Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 5 additions & 7 deletions chain_capabilities/stellar/actions/write_report.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand All @@ -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
Expand Down
36 changes: 36 additions & 0 deletions chain_capabilities/stellar/actions/write_report_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading