From 678dc68edde8ef7f4405c75bef981efe82af6b69 Mon Sep 17 00:00:00 2001 From: Krish-vemula Date: Thu, 10 Sep 2026 12:08:59 -0700 Subject: [PATCH] Fix Stellar early-return telemetry on poll timeout --- .../stellar/actions/write_report.go | 12 +++---- .../stellar/actions/write_report_test.go | 36 +++++++++++++++++++ 2 files changed, 41 insertions(+), 7 deletions(-) diff --git a/chain_capabilities/stellar/actions/write_report.go b/chain_capabilities/stellar/actions/write_report.go index c9c4b713d..c13b775b0 100644 --- a/chain_capabilities/stellar/actions/write_report.go +++ b/chain_capabilities/stellar/actions/write_report.go @@ -398,18 +398,12 @@ func (wr *writeReport) pollTransmissionInfo( attempt := 0 stageTimer := time.NewTimer(delay) - deltaStagePassed := false 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 { @@ -420,6 +414,11 @@ func (wr *writeReport) pollTransmissionInfo( lastValidInfo = info switch lastValidInfo.State { case TransmissionStateSucceeded, TransmissionStateInvalidReceiver, TransmissionStateFailed: + 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 +441,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 diff --git a/chain_capabilities/stellar/actions/write_report_test.go b/chain_capabilities/stellar/actions/write_report_test.go index 62d3ccd6f..c7501cb69 100644 --- a/chain_capabilities/stellar/actions/write_report_test.go +++ b/chain_capabilities/stellar/actions/write_report_test.go @@ -1700,6 +1700,42 @@ func TestPollTransmissionInfo_EmitsInvalidTransmissionStateOnlyOnce(t *testing.T require.Equal(t, 1, invalidStateCount, "InvalidTransmissionState should fire exactly once even across multiple poll iterations with a persistent unexpected state") } +func TestPollTransmissionInfo_ContextTimeoutAfterNonterminalPollEmitsNoEarlyReturn(t *testing.T) { + t.Parallel() + lggr := logger.Test(t) + processor := &recordingWriteReportProcessor{} + scheduler := ts.NewTransmissionScheduler( + p2ptypes.PeerID{2}, + []p2ptypes.PeerID{{1}, {2}, {3}}, + 5*time.Second, + 0, + lggr, + ) + stub := &stubForwarderClient{ + transmissionInfoFn: func(int) (TransmissionInfo, error) { + return TransmissionInfo{State: TransmissionStateNotAttempted}, nil + }, + } + wr := &writeReport{ + forwarderClient: stub, + lggr: logger.Sugared(lggr), + transmissionScheduler: scheduler, + messageBuilder: monitoring.NewMessageBuilder(types.ChainInfo{}, capabilities.CapabilityInfo{}, ""), + beholderProcessor: processor, + } + _, reqMeta, req := newWRReportFixture(t) + transmissionID, err := getTransmissionID(reqMeta.WorkflowExecutionID, req) + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(t.Context(), 75*time.Millisecond) + defer cancel() + + _, err = wr.pollTransmissionInfo(ctx, req, monitoring.TelemetryContext{}, transmissionID, 2) + require.Error(t, err) + require.Contains(t, err.Error(), "timed out waiting for transmission info") + require.False(t, hasTelemetryMessage[*monitoring.WriteReportSuccessfulEarlyReturn](processor.messages)) +} + func TestWriteReport_EmitsInvalidTransmissionStateOnPostSubmitUnexpectedSuccess(t *testing.T) { t.Parallel() h := newWriteReportHelper(t)