Skip to content

Commit e2f985f

Browse files
SuyashAlphaCadecaro
authored andcommitted
refactor(ttx): resolve envelope metrics via DI and add TypedSession
Addresses @adecaro's review on #1716: - Metrics through dependency injection. NewEnvelopeMetrics is provided in the dig container (token/sdk/dig) and registered as a resolvable service; GetEnvelopeMetrics(sp) looks it up by type from a view context (mirroring the existing reflect-based service resolvers). Replaces the previous package-level RegisterMetrics/sync.Once. - TypedSession. Introduce session.TypedSession (and NewTypedSession / NewTypedSessionFromContext / NewTypedSessionToParty / NewTypedSessionForCaller) which wraps a JSON session and resolves the envelope metrics once from the view context. Its SendTyped / ReceiveTyped / ReceiveTypedWithTimeout methods replace the static helper calls across recipients, withdrawal, upgrade, endorse, accept, auditor, collectactions, collectendorsements, multisig/boolpolicy spend, and interop/htlc, so views no longer pass a session to package functions and metrics are recorded automatically. SendEnvelopeOnSession is removed in favour of TypedSession. - Make endorse_test resolve services by requested type instead of by call order, so the extra metrics lookup does not break it. Signed-off-by: SuyashAlphaC <suyashagrawal862@gmail.com>
1 parent 9c53303 commit e2f985f

16 files changed

Lines changed: 181 additions & 155 deletions

File tree

token/sdk/dig/sdk.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@ import (
5858
"github.com/hyperledger-labs/fabric-token-sdk/token/services/ttx/dep"
5959
auditor2 "github.com/hyperledger-labs/fabric-token-sdk/token/services/ttx/dep/auditor"
6060
wrapper2 "github.com/hyperledger-labs/fabric-token-sdk/token/services/ttx/dep/wrapper"
61+
jsession "github.com/hyperledger-labs/fabric-token-sdk/token/services/utils/json/session"
6162
"go.opentelemetry.io/otel/trace"
6263
"go.uber.org/dig"
6364
)
@@ -206,6 +207,7 @@ func (p *SDK) Install() error {
206207
p.Container().Provide(digutils.Identity[*db.OwnerCheckServiceProvider](), dig.As(new(ttx.CheckServiceProvider))),
207208
p.Container().Provide(ttx.NewServiceManager),
208209
p.Container().Provide(ttx.NewMetrics),
210+
p.Container().Provide(jsession.NewEnvelopeMetrics),
209211
)
210212
if err != nil {
211213
return errors.WithMessagef(err, "failed setting up dig container")
@@ -240,6 +242,7 @@ func (p *SDK) Install() error {
240242
digutils.Register[driver.ConfigService](p.Container()),
241243
digutils.Register[*identity.DBStorageProvider](p.Container()),
242244
digutils.Register[*ttx.Metrics](p.Container()),
245+
digutils.Register[*jsession.EnvelopeMetrics](p.Container()),
243246
digutils.Register[*auditor.ServiceManager](p.Container()),
244247
digutils.Register[*ftsconfig.Service](p.Container()),
245248
digutils.Register[*ttx.ServiceManager](p.Container()),

token/services/interop/htlc/distribute.go

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,7 @@ func (v *DistributeTermsView) Call(context view.Context) (interface{}, error) {
7676
if err != nil {
7777
return nil, err
7878
}
79-
if err := session.SendEnvelopeOnSession(sess, context.Context(), v.terms, TypeHTLCTerms); err != nil {
79+
if err := session.NewTypedSession(context, sess).SendTyped(context.Context(), v.terms, TypeHTLCTerms); err != nil {
8080
return nil, errors.Wrapf(err, "failed sending terms")
8181
}
8282

@@ -97,8 +97,7 @@ func ReceiveTerms(context view.Context) (*Terms, error) {
9797

9898
func (v *termsReceiverView) Call(context view.Context) (interface{}, error) {
9999
terms := &Terms{}
100-
s := session.JSON(context)
101-
if err := session.ReceiveTyped(s, TypeHTLCTerms, terms); err != nil {
100+
if err := session.NewTypedSessionFromContext(context).ReceiveTyped(TypeHTLCTerms, terms); err != nil {
102101
return nil, errors.Wrapf(err, "failed unmarshalling terms")
103102
}
104103

token/services/ttx/accept.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -91,9 +91,8 @@ func (s *AcceptView) ack(context view.Context) error {
9191

9292
// Ack for distribution
9393
// Send the signature back
94-
session := context.Session()
9594
logger.DebugfContext(context.Context(), "ack response: [%s] from [%s]", utils.Hashable(sigma), defaultIdentity)
96-
if err := jsession.SendEnvelopeOnSession(session, context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
95+
if err := jsession.NewTypedSessionFromContext(context).SendTyped(context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
9796
return errors.WithMessagef(err, "failed sending ack")
9897
}
9998

token/services/ttx/auditor.go

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -148,9 +148,8 @@ func (a *AuditingViewInitiator) Call(context view.Context) (interface{}, error)
148148
// Receive signature
149149
logger.DebugfContext(context.Context(), "Receiving signature for [%s]", a.tx.ID())
150150

151-
jsonSession := session2.NewFromSession(context, session)
152151
var signaturePayload SignaturePayload
153-
if err := session2.ReceiveTypedWithTimeout(jsonSession, TypeSignature, &signaturePayload, time.Minute); err != nil {
152+
if err := session2.NewTypedSession(context, session).ReceiveTypedWithTimeout(TypeSignature, &signaturePayload, time.Minute); err != nil {
154153
logger.ErrorfContext(context.Context(), "failed to read audit event: %s", err)
155154

156155
return nil, errors.WithMessagef(err, "failed to read audit event")
@@ -182,7 +181,7 @@ func (a *AuditingViewInitiator) startRemote(context view.Context) (view.Session,
182181
if err != nil {
183182
return nil, err
184183
}
185-
err = session2.SendEnvelopeOnSession(session, context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction)
184+
err = session2.NewTypedSession(context, session).SendTyped(context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction)
186185
if err != nil {
187186
return nil, errors.Wrap(err, "failed sending transaction")
188187
}
@@ -219,7 +218,7 @@ func (a *AuditingViewInitiator) startLocal(context view.Context) (view.Session,
219218
if err != nil {
220219
return nil, err
221220
}
222-
err = session2.SendEnvelopeOnSession(left, context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction)
221+
err = session2.NewTypedSession(context, left).SendTyped(context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction)
223222
if err != nil {
224223
return nil, errors.Wrap(err, "failed sending transaction")
225224
}
@@ -334,7 +333,7 @@ func (a *AuditApproveView) signAndSendBack(context view.Context) error {
334333
}
335334

336335
logger.DebugfContext(context.Context(), "auditor sending sigma back", utils.Hashable(sigma))
337-
if err := session2.SendEnvelopeOnSession(context.Session(), context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
336+
if err := session2.NewTypedSessionFromContext(context).SendTyped(context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
338337
return errors.WithMessagef(err, "failed sending back auditor signature")
339338
}
340339
logger.DebugfContext(context.Context(), "Signing and sending back transaction...done [%s]", a.tx.ID())

token/services/ttx/boolpolicy/spend.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -69,8 +69,8 @@ func NewReceiveSpendRequestView() *ReceiveSpendRequestView {
6969
// Call implements view.View.
7070
func (f *ReceiveSpendRequestView) Call(context view.Context) (interface{}, error) {
7171
tx := &SpendRequest{}
72-
s := session.JSON(context)
73-
if err := session.ReceiveTypedWithTimeout(s, ttx.TypeSpendRequest, tx, time.Minute*4); err != nil {
72+
s := session.NewTypedSessionFromContext(context)
73+
if err := s.ReceiveTypedWithTimeout(ttx.TypeSpendRequest, tx, time.Minute*4); err != nil {
7474
logger.ErrorfContext(context.Context(), "failed receiving request: %s", err)
7575

7676
return nil, err
@@ -171,14 +171,14 @@ func (c *RequestSpendView) collectAnswers(context view.Context, party view.Ident
171171

172172
return
173173
}
174-
s := session.NewFromSession(context, backendSession)
175-
if err = session.SendTyped(s, context.Context(), request, ttx.TypeSpendRequest); err != nil {
174+
s := session.NewTypedSession(context, backendSession)
175+
if err = s.SendTyped(context.Context(), request, ttx.TypeSpendRequest); err != nil {
176176
ch <- &answer{err: errors.Wrapf(err, "failed to send request to [%s]", party), party: party}
177177

178178
return
179179
}
180180
response := &SpendResponse{}
181-
if err := session.ReceiveTyped(s, ttx.TypeSpendResponse, response); err != nil {
181+
if err := s.ReceiveTyped(ttx.TypeSpendResponse, response); err != nil {
182182
ch <- &answer{err: errors.Wrapf(err, "failed to receive response from [%s]", party), party: party}
183183

184184
return
@@ -223,8 +223,8 @@ func ReceiveSpendTx(context view.Context, request *SpendRequest) (*Transaction,
223223
// assembled transaction, and returns it without endorsing. Endorsement is
224224
// the caller's responsibility once any business-logic checks pass.
225225
func (a *ReceiveSpendTxView) Call(context view.Context) (interface{}, error) {
226-
s := session.JSON(context)
227-
if err := session.SendTyped(s, context.Context(), &SpendResponse{}, ttx.TypeSpendResponse); err != nil {
226+
s := session.NewTypedSessionFromContext(context)
227+
if err := s.SendTyped(context.Context(), &SpendResponse{}, ttx.TypeSpendResponse); err != nil {
228228
return nil, errors.Wrap(err, "failed to send spend response")
229229
}
230230
tx, err := ttx.ReceiveTransaction(context)

token/services/ttx/collectactions.go

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -96,7 +96,7 @@ func (c *collectActionsView) collectRemote(context view.Context, actionTransfer
9696
party := actionTransfer.From
9797
logger.DebugfContext(context.Context(), "collect remote from [%s]", party)
9898

99-
session, err := session2.NewJSON(context, context.Initiator(), party)
99+
ts, err := session2.NewTypedSessionForCaller(context, context.Initiator(), party)
100100
if err != nil {
101101
return errors.Wrap(err, "failed getting session")
102102
}
@@ -106,19 +106,19 @@ func (c *collectActionsView) collectRemote(context view.Context, actionTransfer
106106
if err != nil {
107107
return errors.Wrap(err, "failed marshalling transaction")
108108
}
109-
if err := session2.SendTyped(session, context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction); err != nil {
109+
if err := ts.SendTyped(context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction); err != nil {
110110
return errors.Wrap(err, "failed sending transaction")
111111
}
112-
if err := session2.SendTyped(session, context.Context(), c.actions, TypeActions); err != nil {
112+
if err := ts.SendTyped(context.Context(), c.actions, TypeActions); err != nil {
113113
return errors.Wrapf(err, "failed sending actions")
114114
}
115-
if err := session2.SendTyped(session, context.Context(), actionTransfer, TypeActionTransfer); err != nil {
115+
if err := ts.SendTyped(context.Context(), actionTransfer, TypeActionTransfer); err != nil {
116116
return errors.Wrapf(err, "failed sending action")
117117
}
118118

119119
// Wait to receive a content back
120120
var txResponse TransactionPayload
121-
if err := session2.ReceiveTyped(session, TypeTransactionResponse, &txResponse); err != nil {
121+
if err := ts.ReceiveTyped(TypeTransactionResponse, &txResponse); err != nil {
122122
return errors.Wrap(err, "failed reading message")
123123
}
124124
txPayload := &Payload{
@@ -184,15 +184,15 @@ func (r *receiveActionsView) Call(context view.Context) (interface{}, error) {
184184
}
185185

186186
// actions
187-
s := session2.JSON(context)
187+
ts := session2.NewTypedSessionFromContext(context)
188188
actions := &Actions{}
189-
if err := session2.ReceiveTyped(s, TypeActions, actions); err != nil {
189+
if err := ts.ReceiveTyped(TypeActions, actions); err != nil {
190190
return nil, errors.Wrap(err, "failed receiving actions")
191191
}
192192

193193
// action
194194
action := &ActionTransfer{}
195-
if err := session2.ReceiveTyped(s, TypeActionTransfer, action); err != nil {
195+
if err := ts.ReceiveTyped(TypeActionTransfer, action); err != nil {
196196
return nil, errors.Wrap(err, "failed receiving action")
197197
}
198198

@@ -216,7 +216,7 @@ func (s *collectActionsResponderView) Call(context view.Context) (interface{}, e
216216
return nil, errors.Wrap(err, "failed marshalling ephemeral transaction")
217217
}
218218

219-
if err := session2.SendEnvelopeOnSession(context.Session(), context.Context(), &TransactionPayload{Raw: response}, TypeTransactionResponse); err != nil {
219+
if err := session2.NewTypedSessionFromContext(context).SendTyped(context.Context(), &TransactionPayload{Raw: response}, TypeTransactionResponse); err != nil {
220220
return nil, errors.Wrap(err, "failed sending back response")
221221
}
222222

token/services/ttx/collectendorsements.go

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -332,13 +332,13 @@ func (c *CollectEndorsementsView) signRemote(
332332
if err != nil {
333333
return nil, errors.Wrap(err, "failed getting session")
334334
}
335-
jsonSession := session2.NewFromSession(context, session)
336-
if err := session2.SendTyped(jsonSession, context.Context(), signatureRequest, TypeSignatureRequest); err != nil {
335+
ts := session2.NewTypedSession(context, session)
336+
if err := ts.SendTyped(context.Context(), signatureRequest, TypeSignatureRequest); err != nil {
337337
return nil, errors.Wrap(err, "failed sending transaction content")
338338
}
339339

340340
var signaturePayload SignaturePayload
341-
if err := session2.ReceiveTypedWithTimeout(jsonSession, TypeSignature, &signaturePayload, time.Minute); err != nil {
341+
if err := ts.ReceiveTypedWithTimeout(TypeSignature, &signaturePayload, time.Minute); err != nil {
342342
return nil, errors.Wrap(err, "failed reading message")
343343
}
344344
sigma := signaturePayload.Signature
@@ -508,14 +508,14 @@ func (c *CollectEndorsementsView) distributeTxToParty(
508508
}
509509
// Send the content
510510
logger.DebugfContext(context.Context(), "Send transaction content")
511-
jsonSession := session2.NewFromSession(context, session)
512-
if err := session2.SendTyped(jsonSession, context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction); err != nil {
511+
ts := session2.NewTypedSession(context, session)
512+
if err := ts.SendTyped(context.Context(), &TransactionPayload{Raw: txRaw}, TypeTransaction); err != nil {
513513
return errors.Wrap(err, "failed sending transaction content")
514514
}
515515

516516
logger.DebugfContext(context.Context(), "Wait for ack")
517517
var signaturePayload SignaturePayload
518-
if err := session2.ReceiveTypedWithTimeout(jsonSession, TypeSignature, &signaturePayload, time.Minute); err != nil {
518+
if err := ts.ReceiveTypedWithTimeout(TypeSignature, &signaturePayload, time.Minute); err != nil {
519519
return errors.Wrapf(err, "failed reading message on session [%s]", session.Info().ID)
520520
}
521521
sigma := signaturePayload.Signature

token/services/ttx/endorse.go

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ func (s *EndorseView) handleSignatureRequests(context view.Context) error {
9494

9595
logger.DebugfContext(context.Context(), "expect [%d] requests to sign for tx id [%s]", len(requiredSigners), s.tx.ID())
9696

97-
session := context.Session()
97+
typedSession := jsession.NewTypedSessionFromContext(context)
9898

9999
tokenRequestToSign, err := s.tx.TokenRequest.MarshalToSign()
100100
if err != nil {
@@ -108,8 +108,7 @@ func (s *EndorseView) handleSignatureRequests(context view.Context) error {
108108
signatureRequest = s.tx.FromSignatureRequest
109109
} else {
110110
logger.DebugfContext(context.Context(), "receiving signature request...")
111-
jsonSession := jsession.JSON(context)
112-
if err := jsession.ReceiveTypedWithTimeout(jsonSession, TypeSignatureRequest, signatureRequest, time.Minute); err != nil {
111+
if err := typedSession.ReceiveTypedWithTimeout(TypeSignatureRequest, signatureRequest, time.Minute); err != nil {
113112
return errors.Wrap(err, "failed reading signature request")
114113
}
115114
}
@@ -136,7 +135,7 @@ func (s *EndorseView) handleSignatureRequests(context view.Context) error {
136135
return errors.Wrapf(err, "failed signing request")
137136
}
138137
logger.DebugfContext(context.Context(), "Send back signature [%s][%s]", signerIdentity, utils.Hashable(sigma))
139-
err = jsession.SendEnvelopeOnSession(session, context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature)
138+
err = typedSession.SendTyped(context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature)
140139
if err != nil {
141140
return errors.Wrapf(err, "failed sending signature back")
142141
}
@@ -190,7 +189,7 @@ func (s *EndorseView) ack(context view.Context, msg []byte) error {
190189
return errors.WithMessagef(err, "failed to sign ack response")
191190
}
192191
logger.DebugfContext(context.Context(), "ack response: [%s] from [%s]", utils.Hashable(sigma), defaultIdentity)
193-
if err := jsession.SendEnvelopeOnSession(inSession, context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
192+
if err := jsession.NewTypedSession(context, inSession).SendTyped(context.Context(), &SignaturePayload{Signature: sigma}, TypeSignature); err != nil {
194193
return errors.WithMessagef(err, "failed sending ack")
195194
}
196195

token/services/ttx/endorse_test.go

Lines changed: 34 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ package ttx_test
88

99
import (
1010
"encoding/json"
11+
"reflect"
1112
"testing"
1213

1314
"github.com/hyperledger-labs/fabric-smart-client/pkg/utils/errors"
@@ -89,17 +90,6 @@ func newTestEndorseViewContext(t *testing.T, input *TestEndorseViewContextInput)
8990
}
9091
tms.NewRequestReturns(req, nil)
9192

92-
ctx := &mock2.Context{}
93-
ctx.SessionReturns(session)
94-
ctx.ContextReturns(t.Context())
95-
ctx.GetServiceReturnsOnCall(0, tmsp, nil)
96-
ctx.GetServiceReturnsOnCall(1, np, nil)
97-
ctx.GetServiceReturnsOnCall(2, &endpoint.Service{}, nil)
98-
ctx.GetServiceReturnsOnCall(3, np, nil)
99-
ctx.GetServiceReturnsOnCall(4, tmsp, nil)
100-
tx, err := ttx.NewTransaction(ctx, []byte("a_signer"))
101-
require.NoError(t, err)
102-
10393
storage := &mock2.Storage{}
10494
storage.AppendReturns(nil)
10595
storageProvider := &mock2.StorageProvider{}
@@ -110,14 +100,42 @@ func newTestEndorseViewContext(t *testing.T, input *TestEndorseViewContextInput)
110100
nis.SignReturns([]byte("an_ack_signature"), nil)
111101
networkIdentityProvider.GetSignerReturns(nis, nil)
112102

103+
// Resolve services by their requested type rather than by call order, so the
104+
// stub is robust to additional GetService lookups (e.g. envelope metrics).
105+
getService := func(v interface{}) (interface{}, error) {
106+
rt, ok := v.(reflect.Type)
107+
if !ok {
108+
return nil, errors.Errorf("unexpected service request [%T]", v)
109+
}
110+
switch rt.String() {
111+
case "*dep.TokenManagementServiceProvider":
112+
return tmsp, nil
113+
case "*dep.NetworkProvider":
114+
return np, nil
115+
case "*dep.NetworkIdentityProvider":
116+
return networkIdentityProvider, nil
117+
case "*endpoint.Service":
118+
return &endpoint.Service{}, nil
119+
case "*ttx.StorageProvider":
120+
return storageProvider, nil
121+
case "*session.EnvelopeMetrics":
122+
return nil, errors.New("envelope metrics not registered in test")
123+
default:
124+
return nil, errors.Errorf("unexpected service request [%s]", rt.String())
125+
}
126+
}
127+
128+
ctx := &mock2.Context{}
129+
ctx.SessionReturns(session)
130+
ctx.ContextReturns(t.Context())
131+
ctx.GetServiceStub = getService
132+
tx, err := ttx.NewTransaction(ctx, []byte("a_signer"))
133+
require.NoError(t, err)
134+
113135
ctx = &mock2.Context{}
114136
ctx.SessionReturns(session)
115137
ctx.ContextReturns(t.Context())
116-
ctx.GetServiceReturnsOnCall(0, storageProvider, nil)
117-
ctx.GetServiceReturnsOnCall(1, np, nil)
118-
ctx.GetServiceReturnsOnCall(2, tmsp, nil)
119-
ctx.GetServiceReturnsOnCall(3, networkIdentityProvider, nil)
120-
ctx.GetServiceReturnsOnCall(4, storageProvider, nil)
138+
ctx.GetServiceStub = getService
121139

122140
txRaw, err := tx.Bytes()
123141
require.NoError(t, err)

token/services/ttx/manager.go

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@ import (
1919
"github.com/hyperledger-labs/fabric-token-sdk/token/services/storage/ttxdb"
2020
"github.com/hyperledger-labs/fabric-token-sdk/token/services/tokens"
2121
"github.com/hyperledger-labs/fabric-token-sdk/token/services/ttx/dep"
22-
jsession "github.com/hyperledger-labs/fabric-token-sdk/token/services/utils/json/session"
2322
"go.opentelemetry.io/otel/trace"
2423
)
2524

@@ -49,11 +48,6 @@ func NewServiceManager(
4948
metricsProvider metrics.Provider,
5049
checkServiceProvider CheckServiceProvider,
5150
) *ServiceManager {
52-
// Register the interactive-protocol envelope metrics once, using the
53-
// metrics provider already injected here, so the session helpers can record
54-
// them without resolving a provider on the per-message path.
55-
jsession.RegisterMetrics(metricsProvider)
56-
5751
return &ServiceManager{
5852
p: lazy.NewProviderWithKeyMapper(services.Key, func(tmsID token.TMSID) (*Service, error) {
5953
ttxStoreService, err := ttxStoreServiceManager.StoreServiceByTMSId(tmsID)

0 commit comments

Comments
 (0)