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
2 changes: 1 addition & 1 deletion core/scripts/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -509,7 +509,7 @@ require (
github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect
github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect
github.com/smartcontractkit/chainlink-confidential-compute v1.3.0 // indirect
github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect
github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect
github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect
github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect
github.com/smartcontractkit/chainlink-framework/chains v0.0.0-20260724153515-bb6a2de39bcb // indirect
Expand Down
4 changes: 2 additions & 2 deletions core/scripts/go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

83 changes: 53 additions & 30 deletions core/services/llo/delegate.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,8 +69,10 @@ type DelegateConfig struct {
JobName null.String
CaptureEATelemetry bool
CaptureObservationTelemetry bool
CaptureOutcomeTelemetry bool
CaptureReportTelemetry bool
// CaptureOutcomeTelemetry and CaptureReportTelemetry apply to v30 instances.
// v31 instances read their plugin telemetry flags from V31Config.
CaptureOutcomeTelemetry bool
CaptureReportTelemetry bool

// LLO
ChannelDefinitionCache llotypes.ChannelDefinitionCache
Expand Down Expand Up @@ -116,6 +118,22 @@ type DelegateConfig struct {
KeyValueDatabaseFactory ocr3_1types.KeyValueDatabaseFactory
}

// telemeterParams creates a telemetry channel when either plugin version needs
// it. Each factory is only handed the channels its own version enables.
func (cfg DelegateConfig) telemeterParams(lggr logger.Logger) telem.TelemeterParams {
return telem.TelemeterParams{
Logger: lggr,
MonitoringEndpoint: cfg.PluginMonitoringEndpoint,
DonID: cfg.DonID,
CaptureEATelemetry: cfg.CaptureEATelemetry,
CaptureObservationTelemetry: cfg.CaptureObservationTelemetry,
CaptureOutcomeTelemetry: cfg.CaptureOutcomeTelemetry || cfg.V31Config.CaptureOutcomeTelemetry,
CaptureReportTelemetry: cfg.CaptureReportTelemetry || cfg.V31Config.CaptureReportTelemetry,
CaptureAttributedObservationTelemetry: cfg.V31Config.CaptureAttributedObservationTelemetry,
SampleTelemetry: cfg.SampleTelemetry,
}
}

// anyV31 reports whether any protocol instance runs the v31 plugin. The
// OCR3.1-only dependencies are per job, so one v31 instance requires them.
func (cfg DelegateConfig) anyV31() bool {
Expand Down Expand Up @@ -174,16 +192,7 @@ func NewDelegate(cfg DelegateConfig) (job.ServiceCtx, error) {
}
reportCodecs := NewReportCodecs(codecLggr, cfg.DonID)

t := telem.NewTelemeterService(telem.TelemeterParams{
Logger: lggr,
MonitoringEndpoint: cfg.PluginMonitoringEndpoint,
DonID: cfg.DonID,
CaptureEATelemetry: cfg.CaptureEATelemetry,
CaptureObservationTelemetry: cfg.CaptureObservationTelemetry,
CaptureOutcomeTelemetry: cfg.CaptureOutcomeTelemetry,
CaptureReportTelemetry: cfg.CaptureReportTelemetry,
SampleTelemetry: cfg.SampleTelemetry,
})
t := telem.NewTelemeterService(cfg.telemeterParams(lggr))

ds := observation.NewDataSource(logger.Named(lggr, "DataSource"), cfg.Registry, t)

Expand Down Expand Up @@ -262,22 +271,7 @@ func (d *delegate) newOracleV30(i int, configTracker ocr2types.ContractConfigTra
OffchainKeyring: d.cfg.OffchainKeyring,
OnchainKeyring: ocr3shims.OnchainKeyringAsOnchainKeyring2(d.cfg.OnchainKeyring),
ReportingPluginFactory: promwrapper.NewReportingPluginFactory(
llov30.NewPluginFactory(
llov30.PluginFactoryParams{
Config: d.cfg.ReportingPluginConfig,
PredecessorRetirementReportCache: psrrc,
ShouldRetireCache: d.src,
RetirementReportCodec: d.cfg.RetirementReportCodec,
ChannelDefinitionCache: d.cfg.ChannelDefinitionCache,
DataSource: d.ds,
Logger: logger.Named(lggr, "ReportingPlugin"),
OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{},
ReportCodecs: d.reportCodecs,
OutcomeTelemetryCh: d.telem.GetOutcomeTelemetryCh(),
ReportTelemetryCh: d.telem.GetReportTelemetryCh(),
DonID: d.cfg.DonID,
},
),
llov30.NewPluginFactory(d.v30FactoryParams(lggr, psrrc)),
lggr,
"",
d.cfg.ChainID,
Expand All @@ -287,12 +281,40 @@ func (d *delegate) newOracleV30(i int, configTracker ocr2types.ContractConfigTra
})
}

// v30FactoryParams assembles the v30 plugin factory params.
func (d *delegate) v30FactoryParams(lggr logger.Logger, psrrc lloprotocol.PredecessorRetirementReportCache) llov30.PluginFactoryParams {
return llov30.PluginFactoryParams{
Config: d.cfg.ReportingPluginConfig,
PredecessorRetirementReportCache: psrrc,
ShouldRetireCache: d.src,
RetirementReportCodec: d.cfg.RetirementReportCodec,
ChannelDefinitionCache: d.cfg.ChannelDefinitionCache,
DataSource: d.ds,
Logger: logger.Named(lggr, "ReportingPlugin"),
OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{},
ReportCodecs: d.reportCodecs,
OutcomeTelemetryCh: enabledCh(d.cfg.CaptureOutcomeTelemetry, d.telem.GetOutcomeTelemetryCh()),
ReportTelemetryCh: enabledCh(d.cfg.CaptureReportTelemetry, d.telem.GetReportTelemetryCh()),
DonID: d.cfg.DonID,
}
}

// enabledCh returns ch when enabled and nil otherwise, so a protocol instance only
// emits the plugin telemetry its own version enables.
func enabledCh[T any](enabled bool, ch chan<- T) chan<- T {
if !enabled {
return nil
}
return ch
}

// v31FactoryParams assembles the v31 plugin factory params, mapping the job's
// V31Config knobs onto it. Knobs left at zero are forwarded as zero, which the
// factory reads as "apply the plugin default".
func (d *delegate) v31FactoryParams(lggr logger.Logger, psrrc lloprotocol.PredecessorRetirementReportCache) llov31.PluginFactoryParams {
return llov31.PluginFactoryParams{
VerboseLogging: d.cfg.ReportingPluginConfig.VerboseLogging || d.cfg.V31Config.VerboseLogging,
CaptureStagingTelemetry: d.cfg.V31Config.CaptureStagingTelemetry,
PredecessorRetirementReportCache: psrrc,
ShouldRetireCache: d.src,
RetirementReportCodec: d.cfg.RetirementReportCodec,
Expand All @@ -301,8 +323,9 @@ func (d *delegate) v31FactoryParams(lggr logger.Logger, psrrc lloprotocol.Predec
Logger: logger.Named(lggr, "ReportingPlugin"),
OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{},
ReportCodecs: d.reportCodecs,
OutcomeTelemetryCh: d.telem.GetOutcomeTelemetryCh(),
ReportTelemetryCh: d.telem.GetReportTelemetryCh(),
OutcomeTelemetryCh: enabledCh(d.cfg.V31Config.CaptureOutcomeTelemetry, d.telem.GetOutcomeTelemetryCh()),
ReportTelemetryCh: enabledCh(d.cfg.V31Config.CaptureReportTelemetry, d.telem.GetReportTelemetryCh()),
AttributedObservationTelemetryCh: enabledCh(d.cfg.V31Config.CaptureAttributedObservationTelemetry, d.telem.GetAttributedObservationTelemetryCh()),
DonID: d.cfg.DonID,
MaxSnapshotRounds: d.cfg.V31Config.MaxSnapshotRounds,
BlobLifetimeRounds: d.cfg.V31Config.BlobLifetimeRounds,
Expand Down
72 changes: 72 additions & 0 deletions core/services/llo/delegate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
lloconfig "github.com/smartcontractkit/chainlink-data-streams/llo/pluginconfig"
llov30 "github.com/smartcontractkit/chainlink-data-streams/llo/v30"
"github.com/smartcontractkit/chainlink/v2/core/services/llo/telem"
"github.com/smartcontractkit/chainlink/v2/core/services/synchronization"
)

// newTestDelegate returns a delegate carrying only what v31FactoryParams reads.
Expand All @@ -37,9 +38,11 @@ func Test_delegate_v31FactoryParams(t *testing.T) {
BlobInFlightWaitFactor: 5,
MaxBlobSnapshotAge: lloconfig.Duration(9 * time.Second),
MaxRoundPeriod: lloconfig.Duration(time.Minute),
CaptureStagingTelemetry: true,
},
}).v31FactoryParams(logger.Test(t), nil)

assert.True(t, params.CaptureStagingTelemetry)
assert.Equal(t, uint64(7), params.MaxSnapshotRounds)
assert.Equal(t, uint64(11), params.BlobLifetimeRounds)
assert.Equal(t, 3*time.Second, params.MaxDurationBlobObservation)
Expand All @@ -60,6 +63,7 @@ func Test_delegate_v31FactoryParams(t *testing.T) {
assert.Zero(t, params.BlobInFlightWaitFactor)
assert.Zero(t, params.MaxBlobSnapshotAge)
assert.Zero(t, params.MaxRoundPeriod)
assert.Nil(t, params.AttributedObservationTelemetryCh)
})

t.Run("forwards a negative MaxBlobSnapshotAge, which disables the age check", func(t *testing.T) {
Expand Down Expand Up @@ -100,6 +104,74 @@ func Test_delegate_v31FactoryParams(t *testing.T) {
})
}

// Test_delegate_pluginTelemetry checks that each plugin version only gets the
// telemetry channels its own config enables: v30 from the node driven flags, v31
// from V31Config alone.
func Test_delegate_pluginTelemetry(t *testing.T) {
t.Parallel()

type channels struct{ outcome, report, attributed bool }
for _, tc := range []struct {
name string
cfg DelegateConfig
v30, v31 channels
}{
{
name: "all off",
},
{
name: "node flags enable v30 only",
cfg: DelegateConfig{
CaptureEATelemetry: true,
CaptureOutcomeTelemetry: true,
CaptureReportTelemetry: true,
},
v30: channels{outcome: true, report: true},
},
{
name: "v31 flags enable v31 only, without the node flag",
cfg: DelegateConfig{
V31Config: lloconfig.V31Config{
CaptureOutcomeTelemetry: true,
CaptureReportTelemetry: true,
CaptureAttributedObservationTelemetry: true,
},
},
v31: channels{outcome: true, report: true, attributed: true},
},
{
name: "each flag independently",
cfg: DelegateConfig{
CaptureEATelemetry: true,
CaptureOutcomeTelemetry: true,
V31Config: lloconfig.V31Config{CaptureReportTelemetry: true},
},
v30: channels{outcome: true},
v31: channels{report: true},
},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

tc.cfg.PluginMonitoringEndpoint = nopTypedEndpoint{}
d := &delegate{cfg: tc.cfg, telem: telem.NewTelemeterService(tc.cfg.telemeterParams(logger.Test(t)))}

v30 := d.v30FactoryParams(logger.Test(t), nil)
assert.Equal(t, tc.v30.outcome, v30.OutcomeTelemetryCh != nil, "v30 outcome")
assert.Equal(t, tc.v30.report, v30.ReportTelemetryCh != nil, "v30 report")

v31 := d.v31FactoryParams(logger.Test(t), nil)
assert.Equal(t, tc.v31.outcome, v31.OutcomeTelemetryCh != nil, "v31 outcome")
assert.Equal(t, tc.v31.report, v31.ReportTelemetryCh != nil, "v31 report")
assert.Equal(t, tc.v31.attributed, v31.AttributedObservationTelemetryCh != nil, "v31 attributed observation")
})
}
}

type nopTypedEndpoint struct{}

func (nopTypedEndpoint) SendTypedLog(synchronization.TelemetryType, []byte) {}

// stubKeyValueDatabaseFactory and stubBinaryNetworkEndpoint2Factory stand in
// for the OCR3.1-only dependencies; validateInstances only checks that they
// are present.
Expand Down
5 changes: 4 additions & 1 deletion core/services/llo/observation/data_source_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -179,7 +179,10 @@ func (m *mockTelemeter) GetOutcomeTelemetryCh() chan<- *lloprotocol.LLOOutcomeTe
return nil
}
func (m *mockTelemeter) GetReportTelemetryCh() chan<- *lloprotocol.LLOReportTelemetry { return nil }
func (m *mockTelemeter) CaptureEATelemetry() bool { return true }
func (m *mockTelemeter) GetAttributedObservationTelemetryCh() chan<- *lloprotocol.LLOAttributedObservationTelemetry {
return nil
}
func (m *mockTelemeter) CaptureEATelemetry() bool { return true }

func (m *mockTelemeter) CaptureObservationTelemetry() bool { return true }

Expand Down
26 changes: 26 additions & 0 deletions core/services/llo/telem/sampling.go
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,20 @@ func fingerprint(typ synchronization.TelemetryType, msg proto.Message) (string,
hex.EncodeToString(m.ConfigDigest),
}
return strings.Join(traits, samplerDelimiter), nanosToSec(int64(m.ObservationTimestampNanoseconds)), nil //nolint:gosec // G115
case synchronization.LLOAttributedObservation:
m, ok := msg.(*lloprotocol.LLOAttributedObservationTelemetry)
if !ok || m == nil {
return "", 0, errors.New("invalid telemetry type, expected LLOAttributedObservationTelemetry")
}
// A large observation is split in parts sharing the observer: key on the
// lowest stream id so each part is sampled on its own.
traits := []string{
strconv.FormatUint(uint64(m.DonId), 10),
strconv.FormatUint(uint64(m.Observer), 10),
hex.EncodeToString(m.ConfigDigest),
strconv.FormatUint(uint64(lowestStreamID(m.StreamValues)), 10),
}
return strings.Join(traits, samplerDelimiter), nanosToSec(int64(m.AgreedObservationTimestampNanoseconds)), nil //nolint:gosec // G115
case synchronization.PipelineBridge:
m, ok := msg.(*LLOBridgeTelemetry)
if !ok || m == nil {
Expand All @@ -171,6 +185,18 @@ func fingerprint(typ synchronization.TelemetryType, msg proto.Message) (string,
}
}

// lowestStreamID returns the lowest key of values, or 0 when empty.
func lowestStreamID(values map[uint32]*lloprotocol.LLOStreamValue) uint32 {
var lowest uint32
first := true
for id := range values {
if first || id < lowest {
lowest, first = id, false
}
}
return lowest
}

func nanosToSec(n int64) int32 {
return int32(n / int64(time.Second)) //nolint:gosec // G115
}
54 changes: 54 additions & 0 deletions core/services/llo/telem/sampling_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,33 @@ func TestFingerprint(t *testing.T) {
ts: int32(ot.Unix()), //nolint:gosec // G115
err: nil,
},
{
name: "successful attributed observation",
msg: &lloprotocol.LLOAttributedObservationTelemetry{
DonId: donID,
Observer: 3,
ConfigDigest: configDigest,
AgreedObservationTimestampNanoseconds: uint64(ot.UnixNano()),
StreamValues: map[uint32]*lloprotocol.LLOStreamValue{streamID: {}, streamID + 7: {}},
},
typ: synchronization.LLOAttributedObservation,
fingerprint: fmt.Sprintf("%d-%d-%x-%d", donID, 3, configDigest, streamID),
ts: int32(ot.Unix()), //nolint:gosec // G115
err: nil,
},
{
name: "attributed observation without values",
msg: &lloprotocol.LLOAttributedObservationTelemetry{
DonId: donID,
Observer: 3,
ConfigDigest: configDigest,
AgreedObservationTimestampNanoseconds: uint64(ot.UnixNano()),
},
typ: synchronization.LLOAttributedObservation,
fingerprint: fmt.Sprintf("%d-%d-%x-%d", donID, 3, configDigest, 0),
ts: int32(ot.Unix()), //nolint:gosec // G115
err: nil,
},
{
name: "successful bridge",
msg: &LLOBridgeTelemetry{
Expand Down Expand Up @@ -139,6 +166,33 @@ func TestSample(t *testing.T) {
assert.False(t, shouldSend)
}

// TestSample_AttributedObservationParts ensures the parts of a split
// observation are sampled independently, and a repeated part is deduplicated.
func TestSample_AttributedObservationParts(t *testing.T) {
t.Parallel()

samplr := newSampler(logger.TestSugared(t), true)
ts := uint64(time.Unix(1600000000, 0).UnixNano())
part := func(observer uint32, ids ...uint32) *lloprotocol.LLOAttributedObservationTelemetry {
m := &lloprotocol.LLOAttributedObservationTelemetry{
DonId: 2,
Observer: observer,
ConfigDigest: []byte("digest"),
AgreedObservationTimestampNanoseconds: ts,
StreamValues: map[uint32]*lloprotocol.LLOStreamValue{},
}
for _, id := range ids {
m.StreamValues[id] = &lloprotocol.LLOStreamValue{}
}
return m
}

assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 1, 2)))
assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 3, 4)), "second part of the same observation")
assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(2, 1, 2)), "another observer")
assert.False(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 1, 2)), "repeated part")
}

// TestPruningLoop ensures the pruning loop works as expected.
func TestPruningLoop(t *testing.T) {
if testing.Short() {
Expand Down
Loading
Loading