Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
40 commits
Select commit Hold shift + click to select a range
f1d191a
SPOR-0001 llo/protocol: reject unknown aggregators at admission, skip…
brunotm Sep 15, 2026
854d28f
SPOR-0003 llo/v31: bound blob payload and persisted aggregates agains…
brunotm Sep 16, 2026
b783d43
SPOR-0003 llo/protocol: bound channel opts and total stream entries a…
brunotm Sep 16, 2026
b77c719
SPOR-0003 llo: bound decimal coefficients on observation decode
brunotm Sep 16, 2026
3d847aa
SPOR-0004 llo/dev/v31: isReportable checks through protocol.Effective…
brunotm Sep 16, 2026
3619afd
SPOR-0004 llo/v31: gate reportability on report codec coverage consen…
brunotm Sep 16, 2026
9e25c05
SPOR-0006 llo: do not halt Observation on baseline verification of al…
brunotm Sep 16, 2026
04c3a42
SPOR-0007 llo: do not fail Observation on retirement cache errors
brunotm Sep 17, 2026
2ca024e
SPOR-0009, SPOR-0013 llo/dev/v31: decouple snapshot freshness from bl…
brunotm Sep 17, 2026
4635b16
SPOR-0002 llo/v31: configure the per-stream contribution floor explic…
brunotm Sep 17, 2026
6018b7f
SPOR-0011 llo/v31: verify the definition set an observation advocates
brunotm Sep 18, 2026
efeeabe
SPOR-0015 llo/protocol: state why the channel cache evicts by inserti…
brunotm Sep 18, 2026
7796568
SPOR-0014 llo/dev/v31: bound how long the blob pump waits at Close
brunotm Sep 18, 2026
1659335
SPOR-0016 llo/dev/v31: cover the layout reset that warms nothing
brunotm Sep 18, 2026
96be4c4
SPOR-0018 llo: memoize the per-definition channel verification checks
brunotm Sep 18, 2026
c116f93
SPOR-0019 llo/dev/v31: cut repeated and serialized blob work
brunotm Sep 18, 2026
e1de001
SPOR-REC: enable race detector in test-ci
brunotm Sep 18, 2026
02f3431
SPOR-REC llo/dev/v31: add golden tests for the precursor and KV records
brunotm Sep 18, 2026
0890f82
SPOR-REC llo/dev/v31: fuzz the observation, precursor and KV decoders
brunotm Sep 18, 2026
c506bed
SPOR-REC llo/dev/v31: add tests for warm and cold restarts
brunotm Sep 18, 2026
d251d95
SPOR-0010 llo/dev/v31: agree on the predecessor signer set in c/pred
brunotm Sep 18, 2026
1ea53a8
SPOR-0015 llo/dev/v31: encode observations and blob payloads determin…
brunotm Sep 18, 2026
3b6f0e7
SPOR-0017 llo/dev/v31: cover reporting from an evicted channel genera…
brunotm Sep 18, 2026
811a21d
SPOR-0008 llo/dev/v31: do not read an empty desired set as a vote to …
brunotm Sep 18, 2026
47d41b0
SPOR-0012 llo/dev/v31: saturate the report cadence deadline
brunotm Sep 18, 2026
7afa96c
llo/transmitter/dataengine: wait on the queue instead of sleeping
brunotm Sep 18, 2026
79cb24a
SPOR-0003 llo/protocol: bound stream value nesting limit during unmar…
brunotm Sep 22, 2026
4abc319
SPOR-0004 llo/dev/v31: accumulate report codec support across rounds
brunotm Sep 22, 2026
e13406b
SPOR-0013 llo/dev/v31: use observation timestamp and ensure monotonicity
brunotm Sep 22, 2026
b6fa058
SPOR-0019 llo/dev/v31: bound the blob payload memo to a byte budget
brunotm Sep 22, 2026
b4025e2
llo/dev/v31: colapse the unreportable channel warnings
brunotm Sep 22, 2026
2b4483c
SPOR-0010 llo/dev/v31: agree on the predecessor signer set per round
brunotm Sep 23, 2026
45f59ee
llo/dev/v31: gate the backfill watermark on prevReportable
brunotm Sep 23, 2026
0d067e7
llo/dev/v31: hold advertised report formats as a set, add golden test…
brunotm Sep 23, 2026
5c19fc0
llo/dev/v31: BlombPump.Take waits if a cycle is currently in flight
brunotm Sep 23, 2026
9287cfa
llo: bound the precursor at admission rather than at write time
brunotm Sep 24, 2026
a96c43d
llo/dev/v31: make the blob in-flight wait factor configurable
brunotm Sep 24, 2026
01d29cd
llo/dev/v31: ensure republish of the carried aggregate when aggregati…
brunotm Sep 25, 2026
9b7805b
llo/dev/v31: observableStreams skip backfills
brunotm Sep 25, 2026
020ffab
llo/dev/v31: honour the observation context and bound the round perio…
brunotm Sep 25, 2026
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
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ test:

.PHONY: test-ci
test-ci: testdb
go test ./... -covermode=atomic -coverpkg=./... -coverprofile=./coverage.txt -json | tee output.txt
go test ./... -covermode=atomic -race -coverpkg=./... -coverprofile=./coverage.txt -json | tee output.txt

.PHONY: lint
lint:
Expand Down
12 changes: 12 additions & 0 deletions llo/dev/v31/backfill.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,5 +59,17 @@ func selectBackfillCandidate(defs llotypes.ChannelDefinitions, validAfter map[ll
if !found {
return 0, 0, protocol.HistoryBackfillOpts{}, false
}
// The candidate must also be emittable, not merely selectable. Reports
// needs the target's report-timestamp resolution and the row's stream
// values; if either fails there, the report is skipped while the watermark
// has already advanced past the row in the state transition, losing it
// permanently. Both are pure functions of (target definition, row), so
// checking them here keeps selection and emission on one path.
if _, err := protocol.ReportTimestampResolutionNanos(target); err != nil {
return 0, 0, protocol.HistoryBackfillOpts{}, false
}
if _, err := protocol.BuildBackfillStreamValues(target, o.Observations[bestRaw]); err != nil {
return 0, 0, protocol.HistoryBackfillOpts{}, false
}
return bestNanos, bestRaw, o, true
}
27 changes: 22 additions & 5 deletions llo/dev/v31/blobcompress.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import (

"github.com/klauspost/compress/zstd"

"github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3_1types"

protocol "github.com/smartcontractkit/chainlink-data-streams/llo/protocol"
)

Expand All @@ -24,6 +26,15 @@ const (
// a huge allocation (zstd bomb).
const maxDecompressedBlobPayloadBytes = protocol.MaxDecompressedObservationLength

// MaxBlobPayloadBytes is the size of a broadcast blob payload libocr accepts,
// and is the value the factory declares as MaxBlobPayloadBytes. It bounds the
// payload as framed for broadcast (codec byte plus body), which is a different
// quantity from maxDecompressedBlobPayloadBytes: the latter is the anti-bomb
// bound applied to the decompressed bytes on the read side, and is deliberately
// looser. Enforcing this one on the write side is what stops the pump from
// broadcasting a blob every peer's libocr would reject.
const MaxBlobPayloadBytes = ocr3_1types.MaxMaxBlobPayloadBytes

// zstd Encoder/Decoder are safe for concurrent use via EncodeAll/DecodeAll and
// are expensive to build, so a single pair is shared process-wide. Built lazily
// so a construction failure surfaces at the call site rather than in init.
Expand Down Expand Up @@ -74,12 +85,18 @@ func encodeBlobPayload(raw []byte) ([]byte, error) {
return nil, err
}
compressed := c.encoder.EncodeAll(raw, []byte{blobCodecZstd})
if len(compressed) < len(raw)+1 {
return compressed, nil
out := compressed
if len(compressed) >= len(raw)+1 {
out = make([]byte, 0, len(raw)+1)
out = append(out, blobCodecRaw)
out = append(out, raw...)
}
// The framed payload is what is broadcast, has to fit the declared limit.
// A payload that compresses poorly can pass the check above and still land over it.
if len(out) > MaxBlobPayloadBytes {
return nil, fmt.Errorf("framed blob payload too large: %d > %d bytes", len(out), MaxBlobPayloadBytes)
}
out := make([]byte, 0, len(raw)+1)
out = append(out, blobCodecRaw)
return append(out, raw...), nil
return out, nil
}

// decodeBlobPayload reverses encodeBlobPayload. The payload is untrusted, so
Expand Down
33 changes: 33 additions & 0 deletions llo/dev/v31/blobcompress_fuzz_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package llo

import (
"testing"
)

// FuzzDecodeBlobPayload feeds arbitrary bytes through the blob payload decoder.
// Blob payloads are attacker-controlled, so the contract is that any input
// either returns bytes within the caller's budget or errors. Never a panic,
// and never more bytes than the budget allows (a zstd bomb).
func FuzzDecodeBlobPayload(f *testing.F) {
raw := []byte("stream values would go here")
framed, err := encodeBlobPayload(raw)
if err != nil {
f.Fatal(err)
}
f.Add(framed)
f.Add(append([]byte{blobCodecRaw}, raw...))
f.Add([]byte{blobCodecZstd})
f.Add([]byte{})

f.Fuzz(func(t *testing.T, payload []byte) {
for _, budget := range []int{0, 1, 1 << 10, maxDecompressedBlobPayloadBytes} {
out, err := decodeBlobPayload(payload, budget)
if err != nil {
continue
}
if len(out) > budget {
t.Fatalf("decoded %d bytes against a budget of %d", len(out), budget)
}
}
})
}
126 changes: 126 additions & 0 deletions llo/dev/v31/blobmemo.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
package llo

import (
"sync"

protocol "github.com/smartcontractkit/chainlink-data-streams/llo/protocol"

ocrtypes "github.com/smartcontractkit/libocr/offchainreporting2plus/types"
)

// maxMemoizedPayloadBytes bounds the decompressed bytes one round's memo
// may hold across all of its entries.The budget is a small multiple of what
// one observation may carry, which covers honest traffic while keeping the
// worst case independent of N.
const maxMemoizedPayloadBytes = 4 * maxObservationDecompressedBytes

// maxMemoizedPayloads bounds the entry count, which the byte budget alone does
// not: an empty payload costs no bytes. No round can legitimately reference more
// handles than every oracle naming the most it is allowed.
const maxMemoizedPayloads = ocrtypes.MaxOracles * maxObservationBlobHandles

// blobPayloadCache memoizes the stream values decoded from blob payloads within
// one sequence number. ValidateObservation and StateTransition decode the same
// observations in the same round, and every FetchBlob re-verifies the blob's
// certificate and re-reads its payload before the plugin decompresses and
// unmarshals it again, so the second decode is entirely repeated work.
//
// Entries are keyed by the marshaled blob handle and scoped to a single
// sequence number. That scope is what keeps the memo consistent with the blob
// transport, which refuses a handle that has expired as of the round's sequence
// number: a hit can never resurrect a blob the round itself would have
// rejected. What one round can hold is bounded by what its observations may
// reference, at most maxObservationBlobHandles handles per observation, each
// contributing at most maxObservationDecompressedBytes, and the whole map is
// dropped when the sequence number advances.
//
// Decoded stream values are treated as immutable: a hit copies the map entries
// into the observation rather than handing out the memoized map.
//
// The memo enforces its own budget (maxMemoizedPayloadBytes) and simply declines
// to store beyond it. Declining is safe because memoization is an optimization,
// and a miss decodes and costs exactly what a hit would have.
type blobPayloadCache struct {
mu sync.Mutex
seqNr uint64
// bytes is the sum of entries' sizes, tracked so the budget does not have
// to walk the map on every write.
bytes int
entries map[string]blobPayloadEntry
}

// blobPayloadEntry is one memoized blob payload. size is the decompressed byte
// count, memoized alongside the values because it is charged against the
// observation's decompression budget: a hit must consume exactly what the
// original decode consumed, or the same observation would decode differently on
// a hit than on a miss.
type blobPayloadEntry struct {
values protocol.StreamValues
size int
}

func newBlobPayloadCache() *blobPayloadCache {
return &blobPayloadCache{}
}

// round returns the memo scoped to seqNr, discarding whatever was held for an
// earlier one. Returns nil for a nil cache, which disables memoization.
func (c *blobPayloadCache) round(seqNr uint64) *roundBlobPayloads {
if c == nil {
return nil
}
c.mu.Lock()
defer c.mu.Unlock()
if c.seqNr != seqNr || c.entries == nil {
c.seqNr = seqNr
c.entries = make(map[string]blobPayloadEntry)
c.bytes = 0
}
return &roundBlobPayloads{cache: c, seqNr: seqNr}
}

// roundBlobPayloads is a handle on the memo for one sequence number. Reads and
// writes through a handle whose round has been superseded are dropped, so a
// round can never see another round's payloads.
type roundBlobPayloads struct {
cache *blobPayloadCache
seqNr uint64
}

func (r *roundBlobPayloads) get(handle []byte) (blobPayloadEntry, bool) {
if r == nil {
return blobPayloadEntry{}, false
}
r.cache.mu.Lock()
defer r.cache.mu.Unlock()
if r.cache.seqNr != r.seqNr {
return blobPayloadEntry{}, false
}
entry, ok := r.cache.entries[string(handle)]
return entry, ok
}

func (r *roundBlobPayloads) put(handle []byte, entry blobPayloadEntry) {
if r == nil {
return
}
r.cache.mu.Lock()
defer r.cache.mu.Unlock()
if r.cache.seqNr != r.seqNr {
return
}
key := string(handle)
prev, replacing := r.cache.entries[key]
if !replacing && len(r.cache.entries) >= maxMemoizedPayloads {
return
}
bytes := r.cache.bytes + entry.size
if replacing {
bytes -= prev.size
}
if bytes > maxMemoizedPayloadBytes {
return
}
r.cache.entries[key] = entry
r.cache.bytes = bytes
}
Loading
Loading