Skip to content
Open
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
33 changes: 32 additions & 1 deletion core/services/llo/telem/telemetry.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,14 @@ import (

const adapterLWBAErrorName = "AdapterLWBAError"

// maxBufferedSeqNrsPerDigest bounds telemetryBuffer independently of
// Transmit() activity. Buffered entries are normally evicted by
// sendBufferedTelemetry when a Transmit() occurs for their digest, but a
// digest that transmits rarely (or never, e.g. because it lost transmission
// duty or is winding down) would otherwise accumulate one entry per round
// forever.
const maxBufferedSeqNrsPerDigest = 10

// DSOpts is the shared, version-agnostic LLO data-source options (llo/v30 and
// llo/v31 both use llodatasource.DSOpts). Aliased here so the telemetry and
// observation paths keep referring to telem.DSOpts.
Expand Down Expand Up @@ -143,7 +151,8 @@ type telemeter struct {
}

// Buffer Report and Outcome telemetry to only send
// for transmitting rounds sequence numbers
// for transmitting rounds sequence numbers. Bounded per-digest by
// maxBufferedSeqNrsPerDigest; see evictOldestSeqNrsLocked.
telemetryBufferMu sync.Mutex
telemetryBuffer map[string]map[uint64][]telemetryEntry

Expand Down Expand Up @@ -319,6 +328,26 @@ func (t *telemeter) sendBufferedTelemetry(digest types.ConfigDigest, seqNr uint6
}()
}

// evictOldestSeqNrsLocked drops the oldest buffered seqNrs for digest until
// at most maxBufferedSeqNrsPerDigest remain. Callers must hold
// telemetryBufferMu.
func (t *telemeter) evictOldestSeqNrsLocked(digest string) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Currently, these 2 for loops run in time O(N*N) where N = len(digestMessages), but this does not seem to be necessary. We want to drop smallest digestMessages seq, until we have maxBufferedSeqNrsPerDigest left. This can be done in O(N) time, you first find the maxBufferedSeqNrsPerDigest-smallest element, and then drop all the elements bigger than it.

digestMessages := t.telemetryBuffer[digest]
for len(digestMessages) > maxBufferedSeqNrsPerDigest {
var oldest uint64
first := true
for seqNr := range digestMessages {
if first || seqNr < oldest {
oldest = seqNr
first = false
}
}
delete(digestMessages, oldest)
t.eng.Warnw("Telemetry: evicted buffered telemetry for stale seqNr; digest may not be transmitting",
"digest", digest, "evictedSeqNr", oldest, "bufferSize", len(digestMessages))
}
}

func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchronization.TelemetryType, msg proto.Message) {
switch typ {
case synchronization.PipelineBridge, synchronization.LLOObservation, synchronization.EnhancedEAMercury:
Expand Down Expand Up @@ -349,6 +378,7 @@ func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchroniz
telemType: typ,
msg: msg,
}}
t.evictOldestSeqNrsLocked(digest)
default: // synchronization.LLOReport and other buffered types
// Report telemetry: append, since multiple reports per seqNr is
// expected (one per reportable channel).
Expand All @@ -363,6 +393,7 @@ func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchroniz
telemType: typ,
msg: msg,
})
t.evictOldestSeqNrsLocked(digest)
}
}

Expand Down
68 changes: 66 additions & 2 deletions core/services/llo/telem/telemetry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -824,7 +824,7 @@ func Test_Telemeter_outcomeTelemetry_samplingAtFlushTime(t *testing.T) {
// second transmits — mimicking a DON where DeltaRound << report interval.
const (
secondsCovered = 3
outcomesPerSecond = 5
outcomesPerSecond = 3 // secondsCovered*outcomesPerSecond must stay <= maxBufferedSeqNrsPerDigest
baseObservationUnix = int64(1737936858)
baseSeqNr = uint64(1000)
)
Expand Down Expand Up @@ -981,7 +981,7 @@ func Test_Telemeter_reportTelemetry_samplingAtFlushTime(t *testing.T) {

const (
secondsCovered = 3
seqNrsPerSecond = 5
seqNrsPerSecond = 3 // secondsCovered*seqNrsPerSecond must stay <= maxBufferedSeqNrsPerDigest
baseObservationUnix = int64(1737936858)
baseSeqNr = uint64(2000)
)
Expand Down Expand Up @@ -1193,3 +1193,67 @@ func Test_Telemeter_reportTelemetry_samplingAtFlushTime(t *testing.T) {
"each per-channel report should be admitted (distinct sampler fingerprints)")
})
}

// Test_Telemeter_telemetryBuffer_boundedWithoutTransmit guards against a
// digest that never (or rarely) transmits accumulating an unbounded number of
// buffered seqNrs. Previously telemetryBuffer entries were only evicted by
// sendBufferedTelemetry on TrackSeqNr, so a stuck/inactive digest would grow
// the buffer forever.
func Test_Telemeter_telemetryBuffer_boundedWithoutTransmit(t *testing.T) {
t.Parallel()

lggr := logger.TestLogger(t)
donID := uint32(1)
cd := (&mockOpts{}).ConfigDigest()

t.Run("outcome telemetry", func(t *testing.T) {
t.Parallel()
m := &mockMonitoringEndpoint{chTypedLogs: make(chan typedLog, 100)}
tm := newTelemeter(TelemeterParams{
Logger: lggr,
MonitoringEndpoint: m,
DonID: donID,
CaptureOutcomeTelemetry: true,
})

const rounds = maxBufferedSeqNrsPerDigest * 5
for i := range rounds {
seqNr := uint64(i)
tm.enqueueTelemetry(cd.Hex(), seqNr, synchronization.LLOOutcome, &lloprotocol.LLOOutcomeTelemetry{
SeqNr: seqNr,
ConfigDigest: cd[:],
})
}

tm.telemetryBufferMu.Lock()
defer tm.telemetryBufferMu.Unlock()
assert.LessOrEqual(t, len(tm.telemetryBuffer[cd.Hex()]), maxBufferedSeqNrsPerDigest)
// the most recent seqNr must always survive eviction
assert.Contains(t, tm.telemetryBuffer[cd.Hex()], uint64(rounds-1))
})

t.Run("report telemetry", func(t *testing.T) {
t.Parallel()
m := &mockMonitoringEndpoint{chTypedLogs: make(chan typedLog, 100)}
tm := newTelemeter(TelemeterParams{
Logger: lggr,
MonitoringEndpoint: m,
DonID: donID,
CaptureReportTelemetry: true,
})

const rounds = maxBufferedSeqNrsPerDigest * 5
for i := range rounds {
seqNr := uint64(i)
tm.enqueueTelemetry(cd.Hex(), seqNr, synchronization.LLOReport, &lloprotocol.LLOReportTelemetry{
SeqNr: seqNr,
ConfigDigest: cd[:],
})
}

tm.telemetryBufferMu.Lock()
defer tm.telemetryBufferMu.Unlock()
assert.LessOrEqual(t, len(tm.telemetryBuffer[cd.Hex()]), maxBufferedSeqNrsPerDigest)
assert.Contains(t, tm.telemetryBuffer[cd.Hex()], uint64(rounds-1))
})
}
Loading