From 38ce9eb5480953c8a1cf39739157905204bc8d14 Mon Sep 17 00:00:00 2001 From: Vladimir Shchukin Date: Thu, 24 Sep 2026 12:08:05 -0400 Subject: [PATCH 1/5] fix in progress logs request --- .../solana/actions/transmission_info_provider.go | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/chain_capabilities/solana/actions/transmission_info_provider.go b/chain_capabilities/solana/actions/transmission_info_provider.go index dba759acc..cf0866c42 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider.go +++ b/chain_capabilities/solana/actions/transmission_info_provider.go @@ -161,6 +161,8 @@ func accountDataBytesFromJSON(asJSON []byte) ([]byte, error) { return nil, fmt.Errorf("could not extract base64 account data from json") } +// Retreive transmission transaction signature from logs deterministically +// Use log with lowest block number and lowest index. func signatureFromInProgressLogs(inProgressLogs []*soltypes.Log) (solana.Signature, error) { if len(inProgressLogs) == 0 { return solana.Signature{}, fmt.Errorf("no in-progress logs") @@ -172,6 +174,9 @@ func signatureFromInProgressLogs(inProgressLogs []*soltypes.Log) (solana.Signatu log = l minBlock = l.BlockNumber } + if l.BlockNumber == minBlock && l.LogIndex < log.LogIndex { + log = l + } } return solana.Signature(log.TxHash), nil } @@ -202,8 +207,10 @@ func (lr *logReader) registerInProgressFilter(ctx context.Context) error { return nil } +const inProgressLogsLimit = 20 + func (lr *logReader) queryInProgress(ctx context.Context, transmissionID [32]byte) ([]*soltypes.Log, error) { - limit := query.NewLimitAndSort(query.CountLimit(1), query.NewSortBySequence(query.Desc)) + limit := query.NewLimitAndSort(query.CountLimit(inProgressLogsLimit), query.NewSortBySequence(query.Asc)) exprs := []query.Expression{ solprimitives.NewEventSigFilter(lr.sigInProgress), solprimitives.NewAddressFilter(soltypes.PublicKey(lr.forwarderProgramID)), From 8cefe13d072a672c964621c9d7b9d72d9783634a Mon Sep 17 00:00:00 2001 From: Vladimir Shchukin Date: Thu, 24 Sep 2026 12:32:35 -0400 Subject: [PATCH 2/5] fix typo --- chain_capabilities/solana/actions/transmission_info_provider.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/chain_capabilities/solana/actions/transmission_info_provider.go b/chain_capabilities/solana/actions/transmission_info_provider.go index cf0866c42..7155f42c5 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider.go +++ b/chain_capabilities/solana/actions/transmission_info_provider.go @@ -161,7 +161,7 @@ func accountDataBytesFromJSON(asJSON []byte) ([]byte, error) { return nil, fmt.Errorf("could not extract base64 account data from json") } -// Retreive transmission transaction signature from logs deterministically +// Retrieve transmission transaction signature from logs deterministically // Use log with lowest block number and lowest index. func signatureFromInProgressLogs(inProgressLogs []*soltypes.Log) (solana.Signature, error) { if len(inProgressLogs) == 0 { From 7f970cd3afc80f3b32c976c67cf88ddc0cedacc8 Mon Sep 17 00:00:00 2001 From: Vladimir Shchukin Date: Tue, 29 Sep 2026 11:07:21 -0400 Subject: [PATCH 3/5] register processed filter --- .../actions/transmission_info_provider.go | 25 +++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/chain_capabilities/solana/actions/transmission_info_provider.go b/chain_capabilities/solana/actions/transmission_info_provider.go index 7155f42c5..679648987 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider.go +++ b/chain_capabilities/solana/actions/transmission_info_provider.go @@ -23,6 +23,7 @@ import ( ) const eventReportInProgress = "ReportInProgress" +const eventReportProcessed = "ReportProcessed" // transmissionLogSubkeyPath indexes ReportInProgress by transmission_id. var transmissionLogSubkeyPath = []string{"TransmissionId"} @@ -35,6 +36,7 @@ type logReader struct { forwarderProgramID solana.PublicKey forwarderState solana.PublicKey sigInProgress soltypes.EventSignature + sigProcessed soltypes.EventSignature } // OnChainTransmissionInfoProvider uses the ExecutionState PDA for success/failure and @@ -57,6 +59,9 @@ func newOnChainTransmissionInfoProvider(ctx context.Context, programID, forwarde } if err := lr.unregisterLegacyInProgressFilter(ctx); err != nil { return nil, fmt.Errorf("failed to unregister legacy ReportInProgress log filter: %w", err) + } + if err := lr.registerProcessedFilter(ctx); err != nil { + } return &OnChainTransmissionInfoProvider{ SolanaService: s, @@ -207,6 +212,26 @@ func (lr *logReader) registerInProgressFilter(ctx context.Context) error { return nil } +func (lr *logReader) registerProcessedFilter(ctx context.Context) error { + idlJSON := []byte(contracts.FetchForwarderIDL()) + sigProcessed := soltypes.EventSignature(lptypes.NewEventSignatureFromName(eventReportProcessed)) + err := lr.RegisterLogTracking(ctx, soltypes.LPFilterQuery{ + Name: eventReportProcessed + "_" + lr.forwarderProgramID.String() + "_v2", + Address: soltypes.PublicKey(lr.forwarderProgramID), + EventName: eventReportInProgress, + EventSig: sigProcessed, + ContractIdlJSON: idlJSON, + SubkeyPaths: [][]string{transmissionLogSubkeyPath, stateSubkeyPath}, + IncludeReverted: true, + }) + if err != nil { + return fmt.Errorf("failed to register ReportInProgress filter for forwarder: %w", err) + } + + lr.sigProcessed = sigProcessed + return nil +} + const inProgressLogsLimit = 20 func (lr *logReader) queryInProgress(ctx context.Context, transmissionID [32]byte) ([]*soltypes.Log, error) { From 3f61226acf521a7a4c107d0978a02147ccf5d718 Mon Sep 17 00:00:00 2001 From: Vladimir Shchukin Date: Fri, 2 Oct 2026 14:53:57 -0400 Subject: [PATCH 4/5] derive signature from logs for successful tx --- .../actions/transmission_info_provider.go | 151 +++------- .../transmission_info_provider_test.go | 277 ++++++++++++++++++ 2 files changed, 318 insertions(+), 110 deletions(-) create mode 100644 chain_capabilities/solana/actions/transmission_info_provider_test.go diff --git a/chain_capabilities/solana/actions/transmission_info_provider.go b/chain_capabilities/solana/actions/transmission_info_provider.go index 679648987..57eaf869a 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider.go +++ b/chain_capabilities/solana/actions/transmission_info_provider.go @@ -1,16 +1,10 @@ package actions import ( - "bytes" "context" - "encoding/base64" - "encoding/json" - "errors" "fmt" - "strings" "github.com/gagliardetto/solana-go" - "github.com/gagliardetto/solana-go/rpc" "github.com/smartcontractkit/chainlink-common/pkg/types" soltypes "github.com/smartcontractkit/chainlink-common/pkg/types/chains/solana" @@ -18,17 +12,16 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/types/query/primitives" solprimitives "github.com/smartcontractkit/chainlink-common/pkg/types/query/primitives/solana" "github.com/smartcontractkit/chainlink-solana/contracts" - ks_forwarder "github.com/smartcontractkit/chainlink-solana/contracts/generated/keystone_forwarder" lptypes "github.com/smartcontractkit/chainlink-solana/pkg/solana/logpoller/types" ) const eventReportInProgress = "ReportInProgress" const eventReportProcessed = "ReportProcessed" -// transmissionLogSubkeyPath indexes ReportInProgress by transmission_id. +// transmissionLogSubkeyPath indexes forwarder events by transmission_id. var transmissionLogSubkeyPath = []string{"TransmissionId"} -// forwarderStateSubkeyPath indexes ReportInProgress by forwarder_state. +// stateSubkeyPath indexes forwarder events by forwarder_state. var stateSubkeyPath = []string{"State"} type logReader struct { @@ -39,8 +32,11 @@ type logReader struct { sigProcessed soltypes.EventSignature } -// OnChainTransmissionInfoProvider uses the ExecutionState PDA for success/failure and -// ReportInProgress logs for transaction signatures (ReportProcessed is not used; it can be truncated). +// OnChainTransmissionInfoProvider derives transmission state from tracked forwarder logs. +// A ReportProcessed log is tracked only from successfully executed transactions, so its +// presence proves success and its tx hash is always a successful tx signature. +// ReportInProgress logs are tracked including reverted transactions: they prove an attempt +// was made and, absent a ReportProcessed log, that every attempt so far failed. type OnChainTransmissionInfoProvider struct { types.SolanaService forwarderProgramID solana.PublicKey @@ -61,7 +57,7 @@ func newOnChainTransmissionInfoProvider(ctx context.Context, programID, forwarde return nil, fmt.Errorf("failed to unregister legacy ReportInProgress log filter: %w", err) } if err := lr.registerProcessedFilter(ctx); err != nil { - + return nil, fmt.Errorf("failed to register ReportProcessed log filter: %w", err) } return &OnChainTransmissionInfoProvider{ SolanaService: s, @@ -76,94 +72,31 @@ func (p *OnChainTransmissionInfoProvider) GetTransmissionInfo(ctx context.Contex if err != nil { return TransmissionInfo{}, fmt.Errorf("failed to request ReportInProgress events: %w", err) } - if len(inProgressLogs) == 0 { return TransmissionInfo{State: TransmissionStateNotAttempted}, nil } - execStateAddr, err := deriveExecutionStatePDA(p.forwarderState, transmissionID, p.forwarderProgramID) + // The ReportProcessed filter excludes reverted transactions, so any tracked + // ReportProcessed log comes from a successfully executed tx. + processedLogs, err := p.lr.queryProcessed(ctx, transmissionID) if err != nil { - return TransmissionInfo{}, fmt.Errorf("failed to derive execution state PDA: %w", err) + return TransmissionInfo{}, fmt.Errorf("failed to request ReportProcessed events: %w", err) } - - reply, err := p.GetAccountInfoWithOpts(ctx, soltypes.GetAccountInfoRequest{ - Account: soltypes.PublicKey(execStateAddr), - Opts: &soltypes.GetAccountInfoOpts{ - Commitment: soltypes.CommitmentProcessed, - }, - }) - if err != nil { - if !isExecutionStateAccountMissing(err) { - return TransmissionInfo{}, fmt.Errorf("failed to get execution state account: %w", err) - } - reply = &soltypes.GetAccountInfoReply{} + if len(processedLogs) > 0 { + return TransmissionInfo{ + State: TransmissionStateSucceeded, + Signature: solana.Signature(processedLogs[0].TxHash), + }, nil } + // ReportInProgress without ReportProcessed: the report was attempted, but every + // attempt so far landed in a reverted tx (a successful tx would have emitted a + // tracked ReportProcessed log). sig, sigErr := signatureFromInProgressLogs(inProgressLogs) if sigErr != nil { return TransmissionInfo{}, sigErr } - - raw, haveBinary := accountDataBytesForTransmission(reply) - if !haveBinary { - // ReportInProgress but no decodable account payload (reverted tx, or missing data on wire). - return TransmissionInfo{State: TransmissionStateFailed, Signature: sig}, nil - } - - execState, err := ks_forwarder.ParseAccount_ExecutionState(raw) - if err != nil { - return TransmissionInfo{}, fmt.Errorf("failed to parse execution state account: %w", err) - } - - if !bytes.Equal(execState.TransmissionId[:], transmissionID[:]) { - return TransmissionInfo{}, fmt.Errorf("execution state transmission id mismatch") - } - - var state TransmissionState - if execState.Success { - state = TransmissionStateSucceeded - } else { - state = TransmissionStateFailed - } - - return TransmissionInfo{ - State: state, - Signature: sig, - }, nil -} - -// accountDataBytesForTransmission returns raw program data for Anchor parsing. Prefers -// AsDecodedBinary; if empty (e.g. jsonParsed-only path), decodes Solana's ["base64","base64"] from AsJSON. -func accountDataBytesForTransmission(reply *soltypes.GetAccountInfoReply) ([]byte, bool) { - if reply == nil || reply.Value == nil || reply.Value.Data == nil { - return nil, false - } - d := reply.Value.Data - if len(d.AsDecodedBinary) > 0 { - return d.AsDecodedBinary, true - } - raw, err := accountDataBytesFromJSON(d.AsJSON) - if err != nil || len(raw) == 0 { - return nil, false - } - return raw, true -} - -func accountDataBytesFromJSON(asJSON []byte) ([]byte, error) { - if len(asJSON) == 0 { - return nil, fmt.Errorf("empty account data json") - } - var arr []string - if err := json.Unmarshal(asJSON, &arr); err == nil && len(arr) >= 2 && arr[1] == "base64" { - return base64.StdEncoding.DecodeString(arr[0]) - } - var wrapped struct { - Data json.RawMessage `json:"data"` - } - if err := json.Unmarshal(asJSON, &wrapped); err == nil && len(wrapped.Data) > 0 { - return accountDataBytesFromJSON(wrapped.Data) - } - return nil, fmt.Errorf("could not extract base64 account data from json") + return TransmissionInfo{State: TransmissionStateFailed, Signature: sig}, nil } // Retrieve transmission transaction signature from logs deterministically @@ -218,14 +151,16 @@ func (lr *logReader) registerProcessedFilter(ctx context.Context) error { err := lr.RegisterLogTracking(ctx, soltypes.LPFilterQuery{ Name: eventReportProcessed + "_" + lr.forwarderProgramID.String() + "_v2", Address: soltypes.PublicKey(lr.forwarderProgramID), - EventName: eventReportInProgress, + EventName: eventReportProcessed, EventSig: sigProcessed, ContractIdlJSON: idlJSON, SubkeyPaths: [][]string{transmissionLogSubkeyPath, stateSubkeyPath}, - IncludeReverted: true, + // Only successfully executed transactions are tracked, so signatures queried + // through this filter always belong to a successful tx. + IncludeReverted: false, }) if err != nil { - return fmt.Errorf("failed to register ReportInProgress filter for forwarder: %w", err) + return fmt.Errorf("failed to register ReportProcessed filter for forwarder: %w", err) } lr.sigProcessed = sigProcessed @@ -255,26 +190,22 @@ func (lr *logReader) queryInProgress(ctx context.Context, transmissionID [32]byt return logs, nil } -func deriveExecutionStatePDA(forwarderState solana.PublicKey, transmissionID [32]byte, programID solana.PublicKey) (solana.PublicKey, error) { - seeds := [][]byte{ - []byte("execution_state"), - forwarderState.Bytes(), - transmissionID[:], +func (lr *logReader) queryProcessed(ctx context.Context, transmissionID [32]byte) ([]*soltypes.Log, error) { + limit := query.NewLimitAndSort(query.CountLimit(1), query.NewSortBySequence(query.Asc)) + exprs := []query.Expression{ + solprimitives.NewEventSigFilter(lr.sigProcessed), + solprimitives.NewAddressFilter(soltypes.PublicKey(lr.forwarderProgramID)), + solprimitives.NewEventBySubkeyFilter(0, []solprimitives.IndexedValueComparator{ + {Value: transmissionID[:], Operator: primitives.Eq}, + }), + solprimitives.NewEventBySubkeyFilter(1, []solprimitives.IndexedValueComparator{ + {Value: lr.forwarderState.Bytes(), Operator: primitives.Eq}, + }), } - ret, _, err := solana.FindProgramAddress(seeds, programID) - return ret, err -} -func isExecutionStateAccountMissing(err error) bool { - if err == nil { - return false - } - if errors.Is(err, rpc.ErrNotFound) { - return true - } - s := strings.ToLower(err.Error()) - if !strings.Contains(s, "not found") { - return false + logs, err := lr.QueryTrackedLogs(ctx, exprs, limit) + if err != nil { + return nil, fmt.Errorf("failed to query tracked logs: %w", err) } - return strings.Contains(s, "account info") || strings.Contains(s, "getaccountinfo") + return logs, nil } diff --git a/chain_capabilities/solana/actions/transmission_info_provider_test.go b/chain_capabilities/solana/actions/transmission_info_provider_test.go new file mode 100644 index 000000000..944b25363 --- /dev/null +++ b/chain_capabilities/solana/actions/transmission_info_provider_test.go @@ -0,0 +1,277 @@ +package actions + +import ( + "bytes" + "errors" + "slices" + "testing" + + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + + soltypes "github.com/smartcontractkit/chainlink-common/pkg/types/chains/solana" + "github.com/smartcontractkit/chainlink-common/pkg/types/mocks" + "github.com/smartcontractkit/chainlink-common/pkg/types/query" + "github.com/smartcontractkit/chainlink-common/pkg/types/query/primitives" + solprimitives "github.com/smartcontractkit/chainlink-common/pkg/types/query/primitives/solana" + + "github.com/smartcontractkit/chainlink-solana/contracts" + lptypes "github.com/smartcontractkit/chainlink-solana/pkg/solana/logpoller/types" +) + +var ( + testSigInProgress = soltypes.EventSignature(lptypes.NewEventSignatureFromName(eventReportInProgress)) + testSigProcessed = soltypes.EventSignature(lptypes.NewEventSignatureFromName(eventReportProcessed)) +) + +// expectProviderFilterRegistrations mocks the log filter lifecycle the provider sets up on creation: +// the ReportInProgress filter (reverted txs included), removal of the legacy filter name and the +// ReportProcessed filter (successful txs only). +func expectProviderFilterRegistrations(svc *mocks.SolanaService, programID solana.PublicKey) { + idlJSON := []byte(contracts.FetchForwarderIDL()) + subkeyPaths := [][]string{transmissionLogSubkeyPath, stateSubkeyPath} + + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.MatchedBy(func(q soltypes.LPFilterQuery) bool { + return q.Name == eventReportInProgress+"_"+programID.String()+"_v2" && + q.EventName == eventReportInProgress && + q.EventSig == testSigInProgress && + bytes.Equal(q.ContractIdlJSON, idlJSON) && + len(q.SubkeyPaths) == 2 && + slices.Equal(q.SubkeyPaths[0], subkeyPaths[0]) && + slices.Equal(q.SubkeyPaths[1], subkeyPaths[1]) && + q.IncludeReverted + })).Return(nil).Once() + + svc.EXPECT().UnregisterLogTracking(mock.Anything, eventReportInProgress+"_"+programID.String()). + Return(nil).Once() + + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.MatchedBy(func(q soltypes.LPFilterQuery) bool { + return q.Name == eventReportProcessed+"_"+programID.String()+"_v2" && + q.EventName == eventReportProcessed && + q.EventSig == testSigProcessed && + bytes.Equal(q.ContractIdlJSON, idlJSON) && + len(q.SubkeyPaths) == 2 && + slices.Equal(q.SubkeyPaths[0], subkeyPaths[0]) && + slices.Equal(q.SubkeyPaths[1], subkeyPaths[1]) && + !q.IncludeReverted + })).Return(nil).Once() +} + +func newTestProvider(t *testing.T) (*mocks.SolanaService, *OnChainTransmissionInfoProvider, solana.PublicKey, solana.PublicKey) { + t.Helper() + + svc := mocks.NewSolanaService(t) + programID := solana.NewWallet().PublicKey() + forwarderState := solana.NewWallet().PublicKey() + expectProviderFilterRegistrations(svc, programID) + + p, err := newOnChainTransmissionInfoProvider(t.Context(), programID, forwarderState, mocks.WrapSolanaService(svc)) + require.NoError(t, err) + provider, ok := p.(*OnChainTransmissionInfoProvider) + require.True(t, ok) + + return svc, provider, programID, forwarderState +} + +// forwarderLogQueryMatcher matches QueryTrackedLogs expressions built for a forwarder event: +// event signature, forwarder program address, transmission id subkey and forwarder state subkey. +func forwarderLogQueryMatcher(sig soltypes.EventSignature, transmissionID [32]byte, programID, forwarderState solana.PublicKey) interface{} { + return mock.MatchedBy(func(exprs []query.Expression) bool { + if len(exprs) != 4 { + return false + } + eventSig, ok := exprs[0].Primitive.(*solprimitives.EventSig) + if !ok || eventSig.Sig != sig { + return false + } + address, ok := exprs[1].Primitive.(*solprimitives.Address) + if !ok || address.PubKey != soltypes.PublicKey(programID) { + return false + } + transmissionFilter, ok := exprs[2].Primitive.(*solprimitives.EventBySubkey) + if !ok || transmissionFilter.SubKeyIndex != 0 || len(transmissionFilter.ValueComparers) != 1 { + return false + } + if !bytes.Equal(transmissionFilter.ValueComparers[0].Value, transmissionID[:]) || + transmissionFilter.ValueComparers[0].Operator != primitives.Eq { + return false + } + stateFilter, ok := exprs[3].Primitive.(*solprimitives.EventBySubkey) + if !ok || stateFilter.SubKeyIndex != 1 || len(stateFilter.ValueComparers) != 1 { + return false + } + if !bytes.Equal(stateFilter.ValueComparers[0].Value, forwarderState.Bytes()) || + stateFilter.ValueComparers[0].Operator != primitives.Eq { + return false + } + return true + }) +} + +func testLog(sig solana.Signature, blockNumber, logIndex int64) *soltypes.Log { + return &soltypes.Log{ + TxHash: soltypes.Signature(sig), + BlockNumber: blockNumber, + LogIndex: logIndex, + } +} + +func TestNewOnChainTransmissionInfoProvider(t *testing.T) { + t.Parallel() + + t.Run("registers forwarder log filters", func(t *testing.T) { + t.Parallel() + // Filter registrations are asserted through the Once() expectations + // and the mock's AssertExpectations cleanup. + _, _, _, _ = newTestProvider(t) + }) + + t.Run("report in progress filter registration fails", func(t *testing.T) { + t.Parallel() + svc := mocks.NewSolanaService(t) + expectedErr := errors.New("registration failed") + + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.Anything).Return(expectedErr).Once() + + _, err := newOnChainTransmissionInfoProvider(t.Context(), + solana.NewWallet().PublicKey(), solana.NewWallet().PublicKey(), mocks.WrapSolanaService(svc)) + require.Error(t, err) + require.ErrorContains(t, err, "failed to register ReportInProgress log filter") + }) + + t.Run("legacy filter unregistration fails", func(t *testing.T) { + t.Parallel() + svc := mocks.NewSolanaService(t) + expectedErr := errors.New("unregistration failed") + + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.Anything).Return(nil).Once() + svc.EXPECT().UnregisterLogTracking(mock.Anything, mock.Anything).Return(expectedErr).Once() + + _, err := newOnChainTransmissionInfoProvider(t.Context(), + solana.NewWallet().PublicKey(), solana.NewWallet().PublicKey(), mocks.WrapSolanaService(svc)) + require.Error(t, err) + require.ErrorContains(t, err, "failed to unregister legacy ReportInProgress log filter") + }) + + t.Run("report processed filter registration fails", func(t *testing.T) { + t.Parallel() + svc := mocks.NewSolanaService(t) + expectedErr := errors.New("registration failed") + + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.Anything).Return(nil).Once() + svc.EXPECT().UnregisterLogTracking(mock.Anything, mock.Anything).Return(nil).Once() + svc.EXPECT().RegisterLogTracking(mock.Anything, mock.Anything).Return(expectedErr).Once() + + _, err := newOnChainTransmissionInfoProvider(t.Context(), + solana.NewWallet().PublicKey(), solana.NewWallet().PublicKey(), mocks.WrapSolanaService(svc)) + require.Error(t, err) + require.ErrorContains(t, err, "failed to register ReportProcessed log filter") + }) +} + +func TestOnChainTransmissionInfoProvider_GetTransmissionInfo(t *testing.T) { + t.Parallel() + + t.Run("no forwarder logs - not attempted", func(t *testing.T) { + t.Parallel() + svc, provider, programID, forwarderState := newTestProvider(t) + transmissionID := [32]byte{1, 2, 3} + + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{}, nil). + Once() + // No ReportProcessed query is expected: the provider short-circuits. + + info, err := provider.GetTransmissionInfo(t.Context(), transmissionID) + require.NoError(t, err) + require.Equal(t, TransmissionInfo{State: TransmissionStateNotAttempted}, info) + }) + + t.Run("report processed - succeeded with successful tx signature", func(t *testing.T) { + t.Parallel() + svc, provider, programID, forwarderState := newTestProvider(t) + transmissionID := [32]byte{9} + + revertedSig := solana.Signature{1} + successfulSig := solana.Signature{2} + + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{testLog(revertedSig, 100, 0)}, nil). + Once() + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigProcessed, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{testLog(successfulSig, 200, 0)}, nil). + Once() + + info, err := provider.GetTransmissionInfo(t.Context(), transmissionID) + require.NoError(t, err) + require.Equal(t, TransmissionStateSucceeded, info.State) + require.Equal(t, successfulSig, info.Signature) + }) + + t.Run("only reverted attempts - failed with earliest attempted signature", func(t *testing.T) { + t.Parallel() + svc, provider, programID, forwarderState := newTestProvider(t) + transmissionID := [32]byte{7} + + laterSig := solana.Signature{3} + earliestSig := solana.Signature{4} + + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{ + testLog(laterSig, 300, 0), + testLog(earliestSig, 200, 5), + testLog(laterSig, 200, 9), + }, nil). + Once() + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigProcessed, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{}, nil). + Once() + + info, err := provider.GetTransmissionInfo(t.Context(), transmissionID) + require.NoError(t, err) + require.Equal(t, TransmissionStateFailed, info.State) + require.Equal(t, earliestSig, info.Signature) + }) + + t.Run("report in progress query fails", func(t *testing.T) { + t.Parallel() + svc, provider, programID, forwarderState := newTestProvider(t) + transmissionID := [32]byte{5} + expectedErr := errors.New("rpc unavailable") + + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). + Return(nil, expectedErr). + Once() + + _, err := provider.GetTransmissionInfo(t.Context(), transmissionID) + require.Error(t, err) + require.ErrorContains(t, err, "failed to request ReportInProgress events") + }) + + t.Run("report processed query fails", func(t *testing.T) { + t.Parallel() + svc, provider, programID, forwarderState := newTestProvider(t) + transmissionID := [32]byte{6} + expectedErr := errors.New("rpc unavailable") + + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). + Return([]*soltypes.Log{testLog(solana.Signature{1}, 100, 0)}, nil). + Once() + svc.EXPECT(). + QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigProcessed, transmissionID, programID, forwarderState), mock.Anything). + Return(nil, expectedErr). + Once() + + _, err := provider.GetTransmissionInfo(t.Context(), transmissionID) + require.Error(t, err) + require.ErrorContains(t, err, "failed to request ReportProcessed events") + }) +} From 2bcac9ae3080f6bae288d46246851ec8c7971021 Mon Sep 17 00:00:00 2001 From: Vladimir Shchukin Date: Mon, 5 Oct 2026 12:55:40 -0400 Subject: [PATCH 5/5] limit in progress logs to 1 --- .../actions/transmission_info_provider.go | 28 ++----------------- .../transmission_info_provider_test.go | 3 -- 2 files changed, 2 insertions(+), 29 deletions(-) diff --git a/chain_capabilities/solana/actions/transmission_info_provider.go b/chain_capabilities/solana/actions/transmission_info_provider.go index 57eaf869a..05fd12cff 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider.go +++ b/chain_capabilities/solana/actions/transmission_info_provider.go @@ -92,31 +92,7 @@ func (p *OnChainTransmissionInfoProvider) GetTransmissionInfo(ctx context.Contex // ReportInProgress without ReportProcessed: the report was attempted, but every // attempt so far landed in a reverted tx (a successful tx would have emitted a // tracked ReportProcessed log). - sig, sigErr := signatureFromInProgressLogs(inProgressLogs) - if sigErr != nil { - return TransmissionInfo{}, sigErr - } - return TransmissionInfo{State: TransmissionStateFailed, Signature: sig}, nil -} - -// Retrieve transmission transaction signature from logs deterministically -// Use log with lowest block number and lowest index. -func signatureFromInProgressLogs(inProgressLogs []*soltypes.Log) (solana.Signature, error) { - if len(inProgressLogs) == 0 { - return solana.Signature{}, fmt.Errorf("no in-progress logs") - } - log := inProgressLogs[0] - minBlock := inProgressLogs[0].BlockNumber - for _, l := range inProgressLogs { - if l.BlockNumber < minBlock { - log = l - minBlock = l.BlockNumber - } - if l.BlockNumber == minBlock && l.LogIndex < log.LogIndex { - log = l - } - } - return solana.Signature(log.TxHash), nil + return TransmissionInfo{State: TransmissionStateFailed, Signature: solana.Signature(inProgressLogs[0].TxHash)}, nil } // Legacy LogTracking doesn't validate against forwarder state used in Event. @@ -167,7 +143,7 @@ func (lr *logReader) registerProcessedFilter(ctx context.Context) error { return nil } -const inProgressLogsLimit = 20 +const inProgressLogsLimit = 1 func (lr *logReader) queryInProgress(ctx context.Context, transmissionID [32]byte) ([]*soltypes.Log, error) { limit := query.NewLimitAndSort(query.CountLimit(inProgressLogsLimit), query.NewSortBySequence(query.Asc)) diff --git a/chain_capabilities/solana/actions/transmission_info_provider_test.go b/chain_capabilities/solana/actions/transmission_info_provider_test.go index 944b25363..64c8ea90a 100644 --- a/chain_capabilities/solana/actions/transmission_info_provider_test.go +++ b/chain_capabilities/solana/actions/transmission_info_provider_test.go @@ -217,15 +217,12 @@ func TestOnChainTransmissionInfoProvider_GetTransmissionInfo(t *testing.T) { svc, provider, programID, forwarderState := newTestProvider(t) transmissionID := [32]byte{7} - laterSig := solana.Signature{3} earliestSig := solana.Signature{4} svc.EXPECT(). QueryTrackedLogs(mock.Anything, forwarderLogQueryMatcher(testSigInProgress, transmissionID, programID, forwarderState), mock.Anything). Return([]*soltypes.Log{ - testLog(laterSig, 300, 0), testLog(earliestSig, 200, 5), - testLog(laterSig, 200, 9), }, nil). Once() svc.EXPECT().