From 9916cf3e32374b02ca8ad56bfce6c7fc4aad4d7e Mon Sep 17 00:00:00 2001 From: ilija42 Date: Thu, 10 Sep 2026 19:51:16 +0200 Subject: [PATCH] Preserve the submitted tx hash when a Stellar WriteReport outcome cannot be confirmed, keep known outcomes when the transaction lookup fails, emit early-return telemetry only on an observed terminal state, and validate report signatures with one shared rule set --- .../stellar/actions/cre_forwarder_codec.go | 63 ++++-- .../stellar/actions/write_report.go | 164 ++++++++------ .../stellar/actions/write_report_test.go | 211 +++++++++++++++--- 3 files changed, 319 insertions(+), 119 deletions(-) diff --git a/chain_capabilities/stellar/actions/cre_forwarder_codec.go b/chain_capabilities/stellar/actions/cre_forwarder_codec.go index 5bb51a8d2..48cc2daea 100644 --- a/chain_capabilities/stellar/actions/cre_forwarder_codec.go +++ b/chain_capabilities/stellar/actions/cre_forwarder_codec.go @@ -18,8 +18,45 @@ const ( reportProcessedTopicPrefix = "forwarder_ReportProcessed" ocrReportContextLen = 96 ed25519OCRSigLen = ed25519.PublicKeySize + ed25519.SignatureSize + // maxOCRSignatures is the libocr oracle-count ceiling; a signed report can never + // legitimately carry more distinct signers than that. + maxOCRSignatures = 31 ) +// validateSignatureSet is the single rule set for a signed report's signature list. +// It is applied at the WriteReport entry point and again inside EncodeReport so both +// gates reject the same inputs with the same messages: nonempty, at least minSigs, +// at most maxOCRSignatures, every entry ed25519OCRSigLen bytes, and no signer public +// key repeated. +func validateSignatureSet(sigs []*sdk.AttributedSignature, minSigs int) error { + if len(sigs) == 0 { + return fmt.Errorf("signed report must contain at least one signature") + } + if len(sigs) < minSigs { + return fmt.Errorf("signed report contains too few signatures: got %d, want at least %d", len(sigs), minSigs) + } + if len(sigs) > maxOCRSignatures { + return fmt.Errorf("signed report contains too many signatures: got %d, want at most %d", len(sigs), maxOCRSignatures) + } + seen := make(map[[ed25519.PublicKeySize]byte]struct{}, len(sigs)) + for i, attributedSig := range sigs { + sig := attributedSig.GetSignature() + if len(sig) != ed25519OCRSigLen { + return fmt.Errorf( + "signature %d has invalid length: expected %d bytes (%d-byte public key || %d-byte signature), got %d", + i, ed25519OCRSigLen, ed25519.PublicKeySize, ed25519.SignatureSize, len(sig), + ) + } + var signer [ed25519.PublicKeySize]byte + copy(signer[:], sig[:ed25519.PublicKeySize]) + if _, dup := seen[signer]; dup { + return fmt.Errorf("signature %d: duplicate signer public key", i) + } + seen[signer] = struct{}{} + } + return nil +} + // CREForwarderCodec encodes and decodes Stellar CRE forwarder contract calls. type CREForwarderCodec interface { EncodeReport(transmitter, receiver string, report *sdk.ReportResponse) ([]stellartypes.ScVal, error) @@ -61,37 +98,17 @@ func (creForwarderCodecImpl) EncodeReport(transmitter string, receiver string, r } signatures := report.GetSigs() - if len(signatures) == 0 { - return nil, fmt.Errorf("report contains no signatures") + if err := validateSignatureSet(signatures, 1); err != nil { + return nil, err } rawSignatures := make([][]byte, len(signatures)) for i, attributedSig := range signatures { - sig := attributedSig.GetSignature() - if len(sig) != ed25519OCRSigLen { - return nil, fmt.Errorf( - "signature %d: expected %d bytes (%d-byte public key || %d-byte signature), got %d", - i, - ed25519OCRSigLen, - ed25519.PublicKeySize, - ed25519.SignatureSize, - len(sig), - ) - } - rawSignatures[i] = sig + rawSignatures[i] = attributedSig.GetSignature() } slices.SortFunc(rawSignatures, func(a, b []byte) int { return bytes.Compare(a[:ed25519.PublicKeySize], b[:ed25519.PublicKeySize]) }) - for i := 1; i < len(rawSignatures); i++ { - previous := rawSignatures[i-1][:ed25519.PublicKeySize] - current := rawSignatures[i][:ed25519.PublicKeySize] - - if bytes.Equal(previous, current) { - return nil, fmt.Errorf("signature %d: duplicate signer public key", i) - } - } - signatureVals := make([]*stellartypes.ScVal, len(rawSignatures)) for i, sig := range rawSignatures { signatureVals[i] = encodeEd25519Signature(sig[:ed25519.PublicKeySize], sig[ed25519.PublicKeySize:]) diff --git a/chain_capabilities/stellar/actions/write_report.go b/chain_capabilities/stellar/actions/write_report.go index c9c4b713d..19d721c2f 100644 --- a/chain_capabilities/stellar/actions/write_report.go +++ b/chain_capabilities/stellar/actions/write_report.go @@ -5,6 +5,7 @@ import ( "encoding/hex" "errors" "fmt" + "math" "math/big" "strings" "time" @@ -123,10 +124,6 @@ func (wr *writeReport) execute( if err := wr.reportSizeLimit.Check(ctx, commoncfg.SizeOf(request.Report.RawReport)); err != nil { return nil, capabilities.ResponseMetadata{}, fmt.Errorf("%s report size exceeds limit: %w", capcommon.UserError, err) } - requiredSigs := int(wr.transmissionScheduler.F) + 1 - if len(request.Report.Sigs) < requiredSigs { - return nil, capabilities.ResponseMetadata{}, fmt.Errorf("%s signed report contains too few signatures: got %d, want at least %d", capcommon.UserError, len(request.Report.Sigs), requiredSigs) - } transmissionID, err := getTransmissionID(metadata.WorkflowExecutionID, request) if err != nil { @@ -155,8 +152,7 @@ func (wr *writeReport) execute( wr.lggr.Errorw("Returning without a transmission attempt - prior transmission succeeded, but failed to retrieve its tx hash", "error", hashErr) return nil, capabilities.ResponseMetadata{}, hashErr } - reply, err := wr.buildSuccessReply(ctx, request, telemetryContext, txHash) - return reply, capabilities.ResponseMetadata{}, err + return wr.buildSuccessReply(ctx, request, telemetryContext, txHash), capabilities.ResponseMetadata{}, nil case TransmissionStateInvalidReceiver: txHash, hashErr := txHashRetriever.GetFailedTransmissionHash(ctx) if hashErr != nil { @@ -167,11 +163,7 @@ func (wr *writeReport) execute( } return nil, capabilities.ResponseMetadata{}, hashErr } - reply, err := wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, info, transmissionID) - if err != nil { - return nil, capabilities.ResponseMetadata{}, revertReplyBuildError(info, transmissionID, err) - } - return reply, capabilities.ResponseMetadata{}, nil + return wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, info, transmissionID), capabilities.ResponseMetadata{}, nil case TransmissionStateFailed: txHash, hashErr := txHashRetriever.GetFailedTransmissionHash(ctx) if hashErr != nil { @@ -182,11 +174,7 @@ func (wr *writeReport) execute( } return nil, capabilities.ResponseMetadata{}, hashErr } - reply, err := wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, info, transmissionID) - if err != nil { - return nil, capabilities.ResponseMetadata{}, revertReplyBuildError(info, transmissionID, err) - } - return reply, capabilities.ResponseMetadata{}, nil + return wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, info, transmissionID), capabilities.ResponseMetadata{}, nil case TransmissionStateNotAttempted: case TransmissionStateUnknown: // Unknown state must not authorize spend. @@ -251,15 +239,18 @@ func (wr *writeReport) execute( "localTxHash", submitResp.TxHash, "localTxStatus", submitResp.TxStatus, ) - - return nil, ownMeteringMetadata, errors.New("failed to retrieve transmission outcome after report submission") + // The node paid for a submit whose outcome it could not confirm. Hand the caller + // what is known (this node's tx hash, local status and fee) so the transaction can + // be reconciled later instead of retried blind. The hash is repeated in the error + // text because the transport delivers only the error to the workflow when one is + // present; the reply and metering ride along for hosts that forward them. + return wr.unconfirmedSubmitReply(submitResp), ownMeteringMetadata, unconfirmedSubmitError(submitResp) } switch postInfo.State { case TransmissionStateSucceeded: if submitResp.TxStatus == stellartypes.TxSuccess && submitResp.TxHash != "" { - reply, err := wr.buildSuccessReply(ctx, request, telemetryContext, submitResp.TxHash) - return reply, ownMeteringMetadata, err + return wr.buildSuccessReply(ctx, request, telemetryContext, submitResp.TxHash), ownMeteringMetadata, nil } txHash, err := txHashRetriever.GetSuccessfulTransmissionHash(ctx) @@ -273,16 +264,11 @@ func (wr *writeReport) execute( wr.lggr, wr.beholderProcessor, wr.messageBuilder.BuildWriteReportDuplicateTx(telemetryContext, request, submitResp.TxHash, txHash)) } - reply, err := wr.buildSuccessReply(ctx, request, telemetryContext, txHash) - return reply, ownMeteringMetadata, err + return wr.buildSuccessReply(ctx, request, telemetryContext, txHash), ownMeteringMetadata, nil case TransmissionStateFailed, TransmissionStateInvalidReceiver: if submitResp.TxStatus == stellartypes.TxSuccess && submitResp.TxHash != "" { wr.lggr.Errorw("Made a new transmission attempt - transmission failed", "txHash", submitResp.TxHash, "transmissionState", postInfo.State) - reply, err := wr.buildRevertReplyFromTx(ctx, request, telemetryContext, submitResp.TxHash, postInfo, transmissionID) - if err != nil { - return nil, ownMeteringMetadata, revertReplyBuildError(postInfo, transmissionID, err) - } - return reply, ownMeteringMetadata, nil + return wr.buildRevertReplyFromTx(ctx, request, telemetryContext, submitResp.TxHash, postInfo, transmissionID), ownMeteringMetadata, nil } txHash, err := txHashRetriever.GetFailedTransmissionHash(ctx) @@ -300,11 +286,7 @@ func (wr *writeReport) execute( wr.messageBuilder.BuildWriteReportDuplicateTx(telemetryContext, request, submitResp.TxHash, txHash)) } wr.lggr.Errorw("Made a new transmission attempt - transmission failed", "txHash", txHash, "transmissionState", postInfo.State) - reply, err := wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, postInfo, transmissionID) - if err != nil { - return nil, ownMeteringMetadata, revertReplyBuildError(postInfo, transmissionID, err) - } - return reply, ownMeteringMetadata, nil + return wr.buildRevertReplyFromTx(ctx, request, telemetryContext, txHash, postInfo, transmissionID), ownMeteringMetadata, nil default: wr.lggr.Errorw("Invalid transmission state after submit", "state", postInfo.State, "localTxStatus", submitResp.TxStatus) wr.emitInvalidTransmissionState(ctx, request, telemetryContext, postInfo, transmissionID, "WriteReport invalid transmission state after submit", invalidTransmissionStateError(postInfo.State).Error()) @@ -328,13 +310,8 @@ func (s *Stellar) validateWriteReportInputs(metadata capabilities.RequestMetadat if len(request.Report.ReportContext) != ocrReportContextLen { return fmt.Errorf("%s report context has invalid length: got %d, want %d", capcommon.UserError, len(request.Report.ReportContext), ocrReportContextLen) } - if len(request.Report.Sigs) == 0 { - return fmt.Errorf("%s signed report must contain at least one signature", capcommon.UserError) - } - for i, sig := range request.Report.Sigs { - if len(sig.GetSignature()) != ed25519OCRSigLen { - return fmt.Errorf("%s signature %d has invalid length: got %d, want %d", capcommon.UserError, i, len(sig.GetSignature()), ed25519OCRSigLen) - } + if err := validateSignatureSet(request.Report.Sigs, int(s.transmissionScheduler.F)+1); err != nil { + return fmt.Errorf("%s %w", capcommon.UserError, err) } reportMetadata, err := capcommon.DecodeReportMetadata(request.Report.RawReport) @@ -398,19 +375,11 @@ func (wr *writeReport) pollTransmissionInfo( attempt := 0 stageTimer := time.NewTimer(delay) - deltaStagePassed := false + defer stageTimer.Stop() hadSuccessfulPoll := false // Guard so an unexpected state that persists across multiple poll iterations only // emits one InvalidTransmissionState metric, not one per poll tick. invalidStateEmitted := false - defer func() { - stageTimer.Stop() - if wr.monitoringEnabled() && !deltaStagePassed && hadSuccessfulPoll { - monitoring.LogAndEmitSuccess(ctx, "Transmission found before delta stage has passed", - wr.lggr, wr.beholderProcessor, - wr.messageBuilder.BuildWriteReportSuccessfulEarlyReturn(telemetryContext)) - } - }() for { if info, infoErr := wr.forwarderClient.GetTransmissionInfo(ctx, transmissionID); infoErr != nil { @@ -420,6 +389,14 @@ func (wr *writeReport) pollTransmissionInfo( lastValidInfo = info switch lastValidInfo.State { case TransmissionStateSucceeded, TransmissionStateInvalidReceiver, TransmissionStateFailed: + // The early-return signal means exactly this: a peer's terminal state was + // observed before this node's slot opened, so no fee will be spent here. + // It is emitted only on this path, never on timeout or error exits. + if wr.monitoringEnabled() { + monitoring.LogAndEmitSuccess(ctx, "Transmission found before delta stage has passed", + wr.lggr, wr.beholderProcessor, + wr.messageBuilder.BuildWriteReportSuccessfulEarlyReturn(telemetryContext)) + } return lastValidInfo, nil case TransmissionStateNotAttempted, TransmissionStateUnknown: // Not yet visible or unreadable; keep polling until the delta stage window @@ -442,7 +419,6 @@ func (wr *writeReport) pollTransmissionInfo( case <-ctx.Done(): return TransmissionInfo{}, fmt.Errorf("timed out waiting for transmission info") case <-stageTimer.C: - deltaStagePassed = true if lastValidInfo.State == TransmissionStateNotAttempted { if finalInfo, finalErr := wr.forwarderClient.GetTransmissionInfo(ctx, transmissionID); finalErr == nil { hadSuccessfulPoll = true @@ -510,7 +486,7 @@ func (wr *writeReport) buildSuccessReply( request *stellarcap.WriteReportRequest, telemetryContext monitoring.TelemetryContext, txHash string, -) (*stellarcap.WriteReportReply, error) { +) *stellarcap.WriteReportReply { return wr.replyFromTransaction(ctx, request, telemetryContext, txHash, stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_SUCCESS, nil) } @@ -521,7 +497,7 @@ func (wr *writeReport) buildRevertReplyFromTx( txHash string, transmissionInfo TransmissionInfo, transmissionID TransmissionID, -) (*stellarcap.WriteReportReply, error) { +) *stellarcap.WriteReportReply { errorMessage := revertReason(transmissionInfo, transmissionID) return wr.replyFromTransaction(ctx, request, telemetryContext, txHash, stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_REVERTED, &errorMessage) } @@ -533,10 +509,14 @@ func revertReason(transmissionInfo TransmissionInfo, transmissionID Transmission return unknownIssueExecutingReceiverContractMessage } -func revertReplyBuildError(transmissionInfo TransmissionInfo, transmissionID TransmissionID, err error) error { - return fmt.Errorf("%s %s: this is the root cause, but an additional error occurred while fetching more info: %w", capcommon.UserError, revertReason(transmissionInfo, transmissionID), err) -} +// maxLedgerCloseTimeSeconds is the largest close time that can be converted to +// microseconds without wrapping a uint64. +const maxLedgerCloseTimeSeconds = int64(math.MaxUint64 / 1_000_000) +// replyFromTransaction builds the reply for a transaction whose forwarder outcome is +// already known. The outcome (hash and receiver status) is authoritative and never +// depends on the GetTransaction lookup; that lookup only enriches the reply with fee, +// ledger and close time, and when it fails those fields are left unset. func (wr *writeReport) replyFromTransaction( ctx context.Context, request *stellarcap.WriteReportRequest, @@ -544,18 +524,7 @@ func (wr *writeReport) replyFromTransaction( txHash string, receiverStatus stellarcap.ReceiverContractExecutionStatus, errorMessage *string, -) (*stellarcap.WriteReportReply, error) { - txResp, err := capcommon.WithQuickRetry(ctx, wr.lggr, func(ctx context.Context) (stellartypes.GetTransactionResponse, error) { - return wr.service.GetTransaction(ctx, stellartypes.GetTransactionRequest{TxHash: txHash}) - }) - if err != nil { - if wr.monitoringEnabled() { - monitoring.LogAndEmitError(ctx, wr.lggr, wr.beholderProcessor, - wr.messageBuilder.BuildWriteReportTxInfoRetrievalError(telemetryContext, request, txHash, err.Error())) - } - return nil, fmt.Errorf("failed to get transaction for tx hash %s: %w", txHash, err) - } - +) *stellarcap.WriteReportReply { message := errorMessage if receiverStatus == stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_REVERTED && errorMessage == nil { message = new(unknownIssueExecutingReceiverContractMessage) @@ -571,10 +540,31 @@ func (wr *writeReport) replyFromTransaction( ReceiverContractExecutionStatus: &receiverStatus, ErrorMessage: message, } + + txResp, err := capcommon.WithQuickRetry(ctx, wr.lggr, func(ctx context.Context) (stellartypes.GetTransactionResponse, error) { + return wr.service.GetTransaction(ctx, stellartypes.GetTransactionRequest{TxHash: txHash}) + }) + if err != nil { + if wr.monitoringEnabled() { + monitoring.LogAndEmitError(ctx, wr.lggr, wr.beholderProcessor, + wr.messageBuilder.BuildWriteReportTxInfoRetrievalError(telemetryContext, request, txHash, err.Error())) + } + wr.lggr.Warnw("Returning reply without transaction details; enrichment lookup failed", + "txHash", txHash, "receiverStatus", receiverStatus, "error", err) + return reply + } + + // A zero fee from the lookup means the fee is unknown, not that the transaction was + // free; it is left unset so nothing downstream reads zero as a real charge. Billing + // uses this node's own submit response, never this field. if txResp.FeeStroops > 0 { reply.TransactionFee = new(txResp.FeeStroops) } - if txResp.LedgerCloseTime > 0 { + switch { + case txResp.LedgerCloseTime <= 0: + case txResp.LedgerCloseTime > maxLedgerCloseTimeSeconds: + wr.lggr.Errorw("Ignoring out-of-range ledger close time", "txHash", txHash, "ledgerCloseTime", txResp.LedgerCloseTime) + default: reply.BlockTimestamp = new(uint64(txResp.LedgerCloseTime) * 1_000_000) } if txResp.LedgerSequence > 0 { @@ -598,7 +588,45 @@ func (wr *writeReport) replyFromTransaction( } wr.lggr.Infow("Successfully fetched transaction", logAttrs...) - return reply, nil + return reply +} + +// unconfirmedSubmitReply describes this node's own submit when the forwarder never +// confirmed the transmission outcome. The receiver status is deliberately left unset: +// nothing is known about it. The tx hash, local status, fee and close time are what the +// relayer reported for the submit, so the caller can reconcile the transaction later. +func (wr *writeReport) unconfirmedSubmitReply(submitResp *stellartypes.SubmitTransactionResponse) *stellarcap.WriteReportReply { + if submitResp == nil || submitResp.TxHash == "" { + return nil + } + reply := &stellarcap.WriteReportReply{ + TxHash: new(submitResp.TxHash), + TxStatus: localTxStatusToProto(submitResp.TxStatus), + TransactionFee: submitResp.TransactionFee, + BlockTimestamp: submitResp.BlockTimestamp, + ErrorMessage: new(fmt.Sprintf("transmission outcome could not be confirmed after submitting tx %s; the transaction may still be included", submitResp.TxHash)), + } + return reply +} + +// unconfirmedSubmitError names the submitted transaction in the error text so a caller +// that receives only the error can still locate and reconcile the transaction. +func unconfirmedSubmitError(submitResp *stellartypes.SubmitTransactionResponse) error { + if submitResp == nil || submitResp.TxHash == "" { + return errors.New("failed to retrieve transmission outcome after report submission") + } + return fmt.Errorf("failed to retrieve transmission outcome after report submission (tx %s, local status %d)", submitResp.TxHash, submitResp.TxStatus) +} + +func localTxStatusToProto(status stellartypes.TransactionStatus) stellarcap.TxStatus { + switch status { + case stellartypes.TxSuccess: + return stellarcap.TxStatus_TX_STATUS_SUCCESS + case stellartypes.TxFailed: + return stellarcap.TxStatus_TX_STATUS_REVERTED + default: + return stellarcap.TxStatus_TX_STATUS_FATAL + } } func transmissionDebugID(id TransmissionID) string { diff --git a/chain_capabilities/stellar/actions/write_report_test.go b/chain_capabilities/stellar/actions/write_report_test.go index 11135a9dc..3eaa076c8 100644 --- a/chain_capabilities/stellar/actions/write_report_test.go +++ b/chain_capabilities/stellar/actions/write_report_test.go @@ -469,7 +469,40 @@ func TestWriteReport_Validation(t *testing.T) { _, err := h.stellar.WriteReport(t.Context(), reqMeta, req) require.NotNil(t, err) require.Contains(t, err.Error(), "signature 0 has invalid length") - require.Contains(t, err.Error(), "want 96") + require.Contains(t, err.Error(), "expected 96 bytes") + }) + + t.Run("duplicate signer rejected before any forwarder read", func(t *testing.T) { + t.Parallel() + h := newWriteReportHelper(t) + _, reqMeta, req := newWRReportFixture(t) + dup := make([]byte, ed25519OCRSigLen) + copy(dup, req.Report.Sigs[0].GetSignature()) + dup[ed25519OCRSigLen-1] ^= 0xff // same signer public key, different signature bytes + req.Report.Sigs = append(req.Report.Sigs, &workflowpb.AttributedSignature{Signature: dup}) + + _, err := h.stellar.WriteReport(t.Context(), reqMeta, req) + require.NotNil(t, err) + require.Contains(t, err.Error(), "duplicate signer public key") + h.svc.AssertNotCalled(t, "SimulateTransaction", mock.Anything, mock.Anything) + }) + + t.Run("more signatures than any DON can produce", func(t *testing.T) { + t.Parallel() + h := newWriteReportHelper(t) + _, reqMeta, req := newWRReportFixture(t) + sigs := make([]*workflowpb.AttributedSignature, maxOCRSignatures+1) + for i := range sigs { + sig := make([]byte, ed25519OCRSigLen) + sig[0] = byte(i) // distinct signer public keys + sigs[i] = &workflowpb.AttributedSignature{Signature: sig} + } + req.Report.Sigs = sigs + + _, err := h.stellar.WriteReport(t.Context(), reqMeta, req) + require.NotNil(t, err) + require.Contains(t, err.Error(), "too many signatures") + require.Contains(t, err.Error(), "want at most 31") }) t.Run("report metadata cannot be decoded", func(t *testing.T) { @@ -709,8 +742,9 @@ func TestWriteReport_Submit(t *testing.T) { result, capErr := h.stellar.WriteReport(ctx, reqMeta, req) require.NotNil(t, capErr) require.Contains(t, capErr.Error(), "failed to retrieve transmission outcome after report submission") + require.Contains(t, capErr.Error(), testTxHash, "the workflow receives only the error, so it must name the transaction") require.NotNil(t, result) - require.Nil(t, result.Response) + requireUnconfirmedSubmitReply(t, result.Response, stellarcap.TxStatus_TX_STATUS_SUCCESS) validateWRMetering(t, result.ResponseMetadata, testWRChainSelector, testFee) h.svc.AssertNotCalled(t, "GetEvents", mock.Anything, mock.Anything) }) @@ -797,7 +831,7 @@ func TestWriteReport_Submit(t *testing.T) { require.NotNil(t, capErr) require.Contains(t, capErr.Error(), "failed to retrieve transmission outcome after report submission") require.NotNil(t, result) - require.Nil(t, result.Response) + requireUnconfirmedSubmitReply(t, result.Response, stellarcap.TxStatus_TX_STATUS_REVERTED) validateWRMetering(t, result.ResponseMetadata, testWRChainSelector, 0) h.svc.AssertNotCalled(t, "GetEvents", mock.Anything, mock.Anything) }) @@ -1313,7 +1347,7 @@ func TestWriteReport_TxFatalSubmitWithoutCanonicalOutcomeReturnsError(t *testing require.NotNil(t, capErr) require.Contains(t, capErr.Error(), "failed to retrieve transmission outcome after report submission") require.NotNil(t, result) - require.Nil(t, result.Response) + requireUnconfirmedSubmitReply(t, result.Response, stellarcap.TxStatus_TX_STATUS_FATAL) validateWRMetering(t, result.ResponseMetadata, testWRChainSelector, 0) h.svc.AssertNotCalled(t, "GetEvents", mock.Anything, mock.Anything) } @@ -1335,10 +1369,28 @@ func TestReplyBuilders(t *testing.T) { LedgerCloseTime: int64(testBlockTimestamp / 1_000_000), }, nil).Once() - reply, err := wr.buildSuccessReply(t.Context(), req, monitoring.TelemetryContext{}, testTxHash) - require.NoError(t, err) + reply := wr.buildSuccessReply(t.Context(), req, monitoring.TelemetryContext{}, testTxHash) require.Equal(t, stellarcap.TxStatus_TX_STATUS_SUCCESS, reply.TxStatus) require.Equal(t, testTxHash, *reply.TxHash) + require.NotNil(t, reply.BlockTimestamp) + require.Equal(t, testBlockTimestamp, *reply.BlockTimestamp) + }) + + t.Run("buildSuccessReply drops an overflowing ledger close time", func(t *testing.T) { + t.Parallel() + mockSvc := mocks.NewStellarService(t) + wr := &writeReport{service: mockSvc, lggr: logger.Sugared(logger.Test(t))} + mockSvc.EXPECT().GetTransaction(mock.Anything, stellartypes.GetTransactionRequest{TxHash: testTxHash}). + Return(stellartypes.GetTransactionResponse{ + FeeStroops: testFee, + LedgerSequence: 100, + LedgerCloseTime: maxLedgerCloseTimeSeconds + 1, + }, nil).Once() + + reply := wr.buildSuccessReply(t.Context(), req, monitoring.TelemetryContext{}, testTxHash) + require.Nil(t, reply.BlockTimestamp) + require.NotNil(t, reply.TransactionFee) + require.NotNil(t, reply.LedgerSequence) }) t.Run("buildRevertReplyFromTx invalid receiver", func(t *testing.T) { @@ -1348,25 +1400,15 @@ func TestReplyBuilders(t *testing.T) { mockSvc.EXPECT().GetTransaction(mock.Anything, stellartypes.GetTransactionRequest{TxHash: testTxHash}). Return(stellartypes.GetTransactionResponse{FeeStroops: testFee}, nil).Once() - reply, err := wr.buildRevertReplyFromTx(t.Context(), req, monitoring.TelemetryContext{}, testTxHash, TransmissionInfo{State: TransmissionStateInvalidReceiver}, transmissionID) - require.NoError(t, err) + reply := wr.buildRevertReplyFromTx(t.Context(), req, monitoring.TelemetryContext{}, testTxHash, TransmissionInfo{State: TransmissionStateInvalidReceiver}, transmissionID) require.Equal(t, stellarcap.TxStatus_TX_STATUS_SUCCESS, reply.TxStatus) require.Contains(t, *reply.ErrorMessage, "not a Wasm contract") }) - - t.Run("revertReplyBuildError", func(t *testing.T) { - t.Parallel() - buildErr := revertReplyBuildError( - TransmissionInfo{State: TransmissionStateFailed}, - transmissionID, - errors.New("rpc down"), - ) - require.Error(t, buildErr) - require.Contains(t, buildErr.Error(), unknownIssueExecutingReceiverContractMessage) - }) } -func TestWriteReport_ObservedRevertReplyBuildError(t *testing.T) { +// A prior transmission's outcome is known from the forwarder; a failing GetTransaction +// lookup must not turn that known outcome into an error. +func TestWriteReport_ObservedRevertReplyWithoutTxDetails(t *testing.T) { t.Parallel() h := newWriteReportHelper(t) rm, reqMeta, req := newWRReportFixture(t) @@ -1383,9 +1425,16 @@ func TestWriteReport_ObservedRevertReplyBuildError(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 500*time.Millisecond) defer cancel() - _, capErr := h.stellar.WriteReport(ctx, reqMeta, req) - require.NotNil(t, capErr) - require.Contains(t, capErr.Error(), unknownIssueExecutingReceiverContractMessage) + result, capErr := h.stellar.WriteReport(ctx, reqMeta, req) + require.Nil(t, capErr) + require.NotNil(t, result) + require.NotNil(t, result.Response) + require.Equal(t, testTxHash, *result.Response.TxHash) + require.Equal(t, stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_REVERTED, *result.Response.ReceiverContractExecutionStatus) + require.Contains(t, *result.Response.ErrorMessage, unknownIssueExecutingReceiverContractMessage) + require.Nil(t, result.Response.TransactionFee) + require.Nil(t, result.Response.LedgerSequence) + require.Nil(t, result.Response.BlockTimestamp) } func TestWriteReport_PostSubmitPollFailureDoesNotRecoverFromEvents(t *testing.T) { @@ -1408,7 +1457,7 @@ func TestWriteReport_PostSubmitPollFailureDoesNotRecoverFromEvents(t *testing.T) require.NotNil(t, capErr) require.Contains(t, capErr.Error(), "failed to retrieve transmission outcome after report submission") require.NotNil(t, result) - require.Nil(t, result.Response) + requireUnconfirmedSubmitReply(t, result.Response, stellarcap.TxStatus_TX_STATUS_SUCCESS) validateWRMetering(t, result.ResponseMetadata, testWRChainSelector, testFee) h.svc.AssertNotCalled(t, "GetEvents", mock.Anything, mock.Anything) } @@ -1492,8 +1541,9 @@ func TestWriteReport_EmitsTxInfoRetrievalErrorTelemetry(t *testing.T) { h.svc.EXPECT().GetTransaction(mock.Anything, stellartypes.GetTransactionRequest{TxHash: testTxHash}). Return(stellartypes.GetTransactionResponse{}, errors.New("rpc down")).Maybe() - _, capErr := h.stellar.WriteReport(t.Context(), reqMeta, req) - require.NotNil(t, capErr) + result, capErr := h.stellar.WriteReport(t.Context(), reqMeta, req) + require.Nil(t, capErr) + require.NotNil(t, result.Response) require.True(t, hasTelemetryMessage[*monitoring.WriteReportTxInfoRetrievalError](processor.messages)) } @@ -1842,7 +1892,7 @@ func TestReplyFromTransaction_SkipsTelemetryWhenMonitoringDisabled(t *testing.T) lggr: logger.Sugared(logger.Test(t)), messageBuilder: monitoring.NewMessageBuilder(types.ChainInfo{}, capabilities.CapabilityInfo{}, ""), } - _, err := wr.replyFromTransaction( + reply := wr.replyFromTransaction( t.Context(), req, monitoring.TelemetryContext{}, @@ -1850,8 +1900,113 @@ func TestReplyFromTransaction_SkipsTelemetryWhenMonitoringDisabled(t *testing.T) stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_SUCCESS, nil, ) + require.NotNil(t, reply) + require.Equal(t, testTxHash, *reply.TxHash) + require.Equal(t, stellarcap.TxStatus_TX_STATUS_SUCCESS, reply.TxStatus) + require.Equal(t, stellarcap.ReceiverContractExecutionStatus_RECEIVER_CONTRACT_EXECUTION_STATUS_SUCCESS, *reply.ReceiverContractExecutionStatus) + require.Nil(t, reply.TransactionFee) + require.Nil(t, reply.LedgerSequence) + require.Nil(t, reply.BlockTimestamp) +} + +func requireUnconfirmedSubmitReply(t *testing.T, reply *stellarcap.WriteReportReply, status stellarcap.TxStatus) { + t.Helper() + require.NotNil(t, reply, "a paid submit must surface its tx hash even when the outcome is unconfirmed") + require.NotNil(t, reply.TxHash) + require.Equal(t, testTxHash, *reply.TxHash) + require.Equal(t, status, reply.TxStatus) + require.Nil(t, reply.ReceiverContractExecutionStatus, "receiver outcome is unknown and must not be asserted") + require.NotNil(t, reply.ErrorMessage) + require.Contains(t, *reply.ErrorMessage, "could not be confirmed") + require.Contains(t, *reply.ErrorMessage, testTxHash) +} + +func TestUnconfirmedSubmitReply(t *testing.T) { + t.Parallel() + wr := &writeReport{lggr: logger.Sugared(logger.Test(t))} + + require.Nil(t, wr.unconfirmedSubmitReply(nil)) + require.Nil(t, wr.unconfirmedSubmitReply(&stellartypes.SubmitTransactionResponse{TxStatus: stellartypes.TxFatal}), "no hash means nothing to reconcile") + + reply := wr.unconfirmedSubmitReply(successSubmitResp()) + requireUnconfirmedSubmitReply(t, reply, stellarcap.TxStatus_TX_STATUS_SUCCESS) + require.NotNil(t, reply.TransactionFee) + require.Equal(t, testFee, *reply.TransactionFee) + require.NotNil(t, reply.BlockTimestamp) + require.Equal(t, testBlockTimestamp, *reply.BlockTimestamp) + + failed := wr.unconfirmedSubmitReply(&stellartypes.SubmitTransactionResponse{TxStatus: stellartypes.TxFailed, TxHash: testTxHash}) + require.Equal(t, stellarcap.TxStatus_TX_STATUS_REVERTED, failed.TxStatus) + fatal := wr.unconfirmedSubmitReply(&stellartypes.SubmitTransactionResponse{TxStatus: stellartypes.TxFatal, TxHash: testTxHash}) + require.Equal(t, stellarcap.TxStatus_TX_STATUS_FATAL, fatal.TxStatus) + + // The error text is the only channel that reaches the workflow on failure, so it must + // carry the hash when one exists and stay generic when none does. + require.EqualError(t, unconfirmedSubmitError(nil), "failed to retrieve transmission outcome after report submission") + require.EqualError(t, unconfirmedSubmitError(&stellartypes.SubmitTransactionResponse{}), "failed to retrieve transmission outcome after report submission") + withHash := unconfirmedSubmitError(successSubmitResp()) + require.Contains(t, withHash.Error(), "failed to retrieve transmission outcome after report submission") + require.Contains(t, withHash.Error(), testTxHash) +} + +// The early-return signal claims a peer's terminal state was seen before this node's +// slot. A request that times out after only NotAttempted polls, with nothing +// transmitted, must not emit it. +func TestPollTransmissionInfo_NoEarlyReturnTelemetryOnTimeout(t *testing.T) { + t.Parallel() + lggr := logger.Test(t) + mockSvc := mocks.NewStellarService(t) + processor := &recordingWriteReportProcessor{} + _, reqMeta, req := newWRReportFixture(t) + transmissionID, err := getTransmissionID(reqMeta.WorkflowExecutionID, req) + require.NoError(t, err) + + mockSvc.EXPECT().SimulateTransaction(mock.Anything, mock.Anything). + Return(transmissionResp(notAttemptedXDR(t)), nil) + + wr := &writeReport{ + forwarderClient: newForwarderClient(mockSvc, lggr, testForwarderAddress, 100), + lggr: logger.Sugared(lggr), + transmissionScheduler: ts.NewTransmissionScheduler( + p2ptypes.PeerID{2}, []p2ptypes.PeerID{{1}, {2}, {3}}, 30*time.Second, 0, lggr), + messageBuilder: monitoring.NewMessageBuilder(types.ChainInfo{}, capabilities.CapabilityInfo{}, ""), + beholderProcessor: processor, + } + + ctx, cancel := context.WithTimeout(t.Context(), 300*time.Millisecond) + defer cancel() + _, err = wr.pollTransmissionInfo(ctx, req, monitoring.TelemetryContext{}, transmissionID, 1) require.Error(t, err) - require.Contains(t, err.Error(), "failed to get transaction") + require.Contains(t, err.Error(), "timed out waiting for transmission info") + require.False(t, hasTelemetryMessage[*monitoring.WriteReportSuccessfulEarlyReturn](processor.messages), + "success metric emitted for a request that timed out with nothing transmitted") +} + +func TestPollTransmissionInfo_EarlyReturnTelemetryOnTerminalStateBeforeSlot(t *testing.T) { + t.Parallel() + lggr := logger.Test(t) + mockSvc := mocks.NewStellarService(t) + processor := &recordingWriteReportProcessor{} + _, reqMeta, req := newWRReportFixture(t) + transmissionID, err := getTransmissionID(reqMeta.WorkflowExecutionID, req) + require.NoError(t, err) + + mockSvc.EXPECT().SimulateTransaction(mock.Anything, mock.Anything). + Return(transmissionResp(succeededXDR(t)), nil).Once() + + wr := &writeReport{ + forwarderClient: newForwarderClient(mockSvc, lggr, testForwarderAddress, 100), + lggr: logger.Sugared(lggr), + transmissionScheduler: ts.NewTransmissionScheduler( + p2ptypes.PeerID{2}, []p2ptypes.PeerID{{1}, {2}, {3}}, 30*time.Second, 0, lggr), + messageBuilder: monitoring.NewMessageBuilder(types.ChainInfo{}, capabilities.CapabilityInfo{}, ""), + beholderProcessor: processor, + } + + info, err := wr.pollTransmissionInfo(t.Context(), req, monitoring.TelemetryContext{}, transmissionID, 1) + require.NoError(t, err) + require.Equal(t, TransmissionStateSucceeded, info.State) + require.True(t, hasTelemetryMessage[*monitoring.WriteReportSuccessfulEarlyReturn](processor.messages)) } func TestCheckEstimatedSpendLimit(t *testing.T) {