From 9cc05584aa5b8fc27bea589a35619c7a6f286c51 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 8 Sep 2026 14:30:26 -0600 Subject: [PATCH 01/23] pkg/workflows/dontime: introduce feature flagged fix for sequence number inconsistency --- pkg/settings/cresettings/README.md | 3 + pkg/settings/cresettings/defaults.json | 1 + pkg/settings/cresettings/defaults.toml | 1 + pkg/settings/cresettings/settings.go | 42 ++--- pkg/workflows/dontime/factory.go | 23 ++- pkg/workflows/dontime/pb/dontime.go | 24 +++ pkg/workflows/dontime/pb/dontime.pb.go | 52 ++++-- pkg/workflows/dontime/pb/dontime.proto | 3 +- pkg/workflows/dontime/plugin.go | 244 ++++++++++++++++++++++++- 9 files changed, 339 insertions(+), 54 deletions(-) create mode 100644 pkg/workflows/dontime/pb/dontime.go diff --git a/pkg/settings/cresettings/README.md b/pkg/settings/cresettings/README.md index 0ce55544d0..a334496097 100644 --- a/pkg/settings/cresettings/README.md +++ b/pkg/settings/cresettings/README.md @@ -348,6 +348,9 @@ flowchart %% enclave → gateway → relay DON node is likewise its own entry point HandleGatewayMessage +%% TODO placating test for now since this flowchart no longer renders + DonTimeSequencedTimestampsEnabled + classDef bound stroke:#f00 classDef gate stroke:#0f0 classDef queue stroke:#00f diff --git a/pkg/settings/cresettings/defaults.json b/pkg/settings/cresettings/defaults.json index 4e0505e8e1..d39ae28768 100644 --- a/pkg/settings/cresettings/defaults.json +++ b/pkg/settings/cresettings/defaults.json @@ -50,6 +50,7 @@ "VaultMaxBlobPayloadSizeLimit": "25.6kb", "VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit": "31.45728mb", "VaultMaxPerOracleUnexpiredBlobCount": "1000", + "DonTimeSequencedTimestampsEnabled": "[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]", "ConfidentialCompute": { "GlobalRate": "1000rps:1000", "MaxRetries": "3", diff --git a/pkg/settings/cresettings/defaults.toml b/pkg/settings/cresettings/defaults.toml index 3004605698..a04476a9a9 100644 --- a/pkg/settings/cresettings/defaults.toml +++ b/pkg/settings/cresettings/defaults.toml @@ -49,6 +49,7 @@ VaultMaxKeyValueModifiedKeys = '300' VaultMaxBlobPayloadSizeLimit = '25.6kb' VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit = '31.45728mb' VaultMaxPerOracleUnexpiredBlobCount = '1000' +DonTimeSequencedTimestampsEnabled = '[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]' [ConfidentialCompute] GlobalRate = '1000rps:1000' diff --git a/pkg/settings/cresettings/settings.go b/pkg/settings/cresettings/settings.go index e07cfb558a..0814e503fe 100644 --- a/pkg/settings/cresettings/settings.go +++ b/pkg/settings/cresettings/settings.go @@ -51,6 +51,11 @@ var DefaultGetter Getter // Deprecated: use Default var Config Schema +var ( + year2100 = time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC) + disabledFeatureTimeRange = TimeRange(year2100, time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)) +) + var Default = Schema{ WorkflowLimit: Int(1000), WorkflowExecutionConcurrencyLimit: Int(1000), @@ -149,6 +154,8 @@ var Default = Schema{ VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit: Size(31457280 * config.Byte), VaultMaxPerOracleUnexpiredBlobCount: Int(1000), + DonTimeSequencedTimestampsEnabled: disabledFeatureTimeRange, + // Confidential Compute (San Marino framework) node-level settings. Defaults // mirror the previous hardcoded executor defaults so behavior is unchanged // until explicitly overridden. @@ -324,39 +331,24 @@ var Default = Schema{ RequestTimeout: Duration(30 * time.Second), }, - FeatureMultiTriggerExecutionIDsActiveAt: Time(time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureMultiTriggerExecutionIDsActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureHTTPTriggerNewExecutionIDsActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureUseSingleDONTimeProviderPerExecutionActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureChainCapabilityHashBasedOCRActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureEVMWriteReportL1FeeActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), - FeatureAptosWriteReportBlockTimestampActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), + FeatureMultiTriggerExecutionIDsActiveAt: Time(year2100), + FeatureMultiTriggerExecutionIDsActivePeriod: disabledFeatureTimeRange, + FeatureHTTPTriggerNewExecutionIDsActivePeriod: disabledFeatureTimeRange, + FeatureUseSingleDONTimeProviderPerExecutionActivePeriod: disabledFeatureTimeRange, + FeatureChainCapabilityHashBasedOCRActivePeriod: disabledFeatureTimeRange, + FeatureEVMWriteReportL1FeeActivePeriod: disabledFeatureTimeRange, + FeatureAptosWriteReportBlockTimestampActivePeriod: disabledFeatureTimeRange, // ON by default: covers all possible timestamps including zero time.Time{}, // so WorkflowTag is included in the hash matching current prod behavior. // After rollout, set to far-future window to exclude WorkflowTag. FeatureRequestHashIncludeWorkflowTagActivePeriod: TimeRange( - time.Date(1, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC)), + time.Date(1, 1, 1, 0, 0, 0, 0, time.UTC), year2100), // OFF by default: the workflow_specs_v2.workflow_tag reconcile backfill // is intentionally disabled on a fresh deploy. Ops narrows the range to // cover "now" only after FeatureRequestHashIncludeWorkflowTag is muted // on every DON member, so DBs can heal without producing tag-driven // hash divergence during the fill window. - FeatureWorkflowTagBackfillActivePeriod: TimeRange( - time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC), - time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)), + FeatureWorkflowTagBackfillActivePeriod: disabledFeatureTimeRange, }, } @@ -434,6 +426,8 @@ type Schema struct { VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit Setting[config.Size] VaultMaxPerOracleUnexpiredBlobCount Setting[int] + DonTimeSequencedTimestampsEnabled Setting[Range[config.Timestamp]] + // Confidential Compute (San Marino framework) node-level settings. ConfidentialCompute confidentialCompute diff --git a/pkg/workflows/dontime/factory.go b/pkg/workflows/dontime/factory.go index 76c102547a..a73b2adc60 100644 --- a/pkg/workflows/dontime/factory.go +++ b/pkg/workflows/dontime/factory.go @@ -9,8 +9,11 @@ import ( "github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3types" + "github.com/smartcontractkit/chainlink-common/pkg/config" "github.com/smartcontractkit/chainlink-common/pkg/logger" "github.com/smartcontractkit/chainlink-common/pkg/services" + "github.com/smartcontractkit/chainlink-common/pkg/settings/cresettings" + "github.com/smartcontractkit/chainlink-common/pkg/settings/limits" "github.com/smartcontractkit/chainlink-common/pkg/types/core" "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" ) @@ -26,19 +29,30 @@ const ( var _ core.OCR3ReportingPluginFactory = &Factory{} type Factory struct { - store *Store - lggr logger.Logger + store *Store + lggr logger.Logger + sequencedTSEnabled limits.RangeLimiter[config.Timestamp] services.StateMachine } func NewFactory(s *Store, lggr logger.Logger) (*Factory, error) { return &Factory{ - store: s, - lggr: logger.Named(lggr, "OCR3DonTimeFactory"), + store: s, + lggr: logger.Named(lggr, "OCR3DonTimeFactory"), + sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsEnabled.DefaultValue), }, nil } +func (o *Factory) InitLimits(lf limits.Factory) error { + sequencedTSEnabled, err := limits.MakeRangeLimiter[config.Timestamp](lf, cresettings.Default.DonTimeSequencedTimestampsEnabled) + if err != nil { + return err + } + o.sequencedTSEnabled = sequencedTSEnabled + return nil +} + func (o *Factory) NewReportingPlugin(_ context.Context, config ocr3types.ReportingPluginConfig) (ocr3types.ReportingPlugin[[]byte], ocr3types.ReportingPluginInfo, error) { var configProto pb.Config err := proto.Unmarshal(config.OffchainConfig, &configProto) @@ -75,6 +89,7 @@ func (o *Factory) NewReportingPlugin(_ context.Context, config ocr3types.Reporti if err != nil { return nil, ocr3types.ReportingPluginInfo{}, err } + plugin.setSequencedTSEnabled(o.sequencedTSEnabled) pluginInfo := ocr3types.ReportingPluginInfo{ Name: "DON Time Plugin", Limits: ocr3types.ReportingPluginLimits{ diff --git a/pkg/workflows/dontime/pb/dontime.go b/pkg/workflows/dontime/pb/dontime.go new file mode 100644 index 0000000000..8a30b52ac6 --- /dev/null +++ b/pkg/workflows/dontime/pb/dontime.go @@ -0,0 +1,24 @@ +package pb + +import "math" + +// MaxSeqNum returns the max sequence number from TimestampsBySequence, or 0 if none exist. +func (t *ObservedDonTimes) MaxSeqNum() (maxSeqNum int64) { + for seqNum := range t.TimestampsBySequence { + if seqNum > maxSeqNum { + maxSeqNum = seqNum + } + } + return +} + +// EarliestTS returns the easliest timestamp value from TimestampsBySequence, or math.MaxInt64 if none exist. +func (t *ObservedDonTimes) EarliestTS() (earliestTS int64) { + earliestTS = math.MaxInt64 + for _, ts := range t.TimestampsBySequence { + if ts < earliestTS { + earliestTS = ts + } + } + return +} diff --git a/pkg/workflows/dontime/pb/dontime.pb.go b/pkg/workflows/dontime/pb/dontime.pb.go index 36d662710f..961d02094e 100644 --- a/pkg/workflows/dontime/pb/dontime.pb.go +++ b/pkg/workflows/dontime/pb/dontime.pb.go @@ -82,10 +82,12 @@ func (x *Observation) GetLimitByBatchSizeFlag() bool { } type ObservedDonTimes struct { - state protoimpl.MessageState `protogen:"open.v1"` - Timestamps []int64 `protobuf:"varint,1,rep,packed,name=timestamps,proto3" json:"timestamps,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + state protoimpl.MessageState `protogen:"open.v1"` + // Deprecated: Marked as deprecated in dontime.proto. + Timestamps []int64 `protobuf:"varint,1,rep,packed,name=timestamps,proto3" json:"timestamps,omitempty"` + TimestampsBySequence map[int64]int64 `protobuf:"bytes,2,rep,name=timestampsBySequence,proto3" json:"timestampsBySequence,omitempty" protobuf_key:"varint,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *ObservedDonTimes) Reset() { @@ -118,6 +120,7 @@ func (*ObservedDonTimes) Descriptor() ([]byte, []int) { return file_dontime_proto_rawDescGZIP(), []int{1} } +// Deprecated: Marked as deprecated in dontime.proto. func (x *ObservedDonTimes) GetTimestamps() []int64 { if x != nil { return x.Timestamps @@ -125,6 +128,13 @@ func (x *ObservedDonTimes) GetTimestamps() []int64 { return nil } +func (x *ObservedDonTimes) GetTimestampsBySequence() map[int64]int64 { + if x != nil { + return x.TimestampsBySequence + } + return nil +} + type Outcome struct { state protoimpl.MessageState `protogen:"open.v1"` Timestamp int64 `protobuf:"varint,1,opt,name=timestamp,proto3" json:"timestamp,omitempty"` @@ -188,11 +198,15 @@ const file_dontime_proto_rawDesc = "" + "\x18limit_by_batch_size_flag\x18\x04 \x01(\bR\x14limitByBatchSizeFlag\x1a;\n" + "\rRequestsEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01J\x04\b\x03\x10\x04\"2\n" + - "\x10ObservedDonTimes\x12\x1e\n" + + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01J\x04\b\x03\x10\x04\"\xe0\x01\n" + + "\x10ObservedDonTimes\x12\"\n" + "\n" + - "timestamps\x18\x01 \x03(\x03R\n" + - "timestamps\"\xcd\x01\n" + + "timestamps\x18\x01 \x03(\x03B\x02\x18\x01R\n" + + "timestamps\x12_\n" + + "\x14timestampsBySequence\x18\x02 \x03(\v2+.ObservedDonTimes.TimestampsBySequenceEntryR\x14timestampsBySequence\x1aG\n" + + "\x19TimestampsBySequenceEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\x03R\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"\xcd\x01\n" + "\aOutcome\x12\x1c\n" + "\ttimestamp\x18\x01 \x01(\x03R\ttimestamp\x12L\n" + "\x12observed_don_times\x18\x02 \x03(\v2\x1e.Outcome.ObservedDonTimesEntryR\x10observedDonTimes\x1aV\n" + @@ -212,23 +226,25 @@ func file_dontime_proto_rawDescGZIP() []byte { return file_dontime_proto_rawDescData } -var file_dontime_proto_msgTypes = make([]protoimpl.MessageInfo, 5) +var file_dontime_proto_msgTypes = make([]protoimpl.MessageInfo, 6) var file_dontime_proto_goTypes = []any{ (*Observation)(nil), // 0: Observation (*ObservedDonTimes)(nil), // 1: ObservedDonTimes (*Outcome)(nil), // 2: Outcome nil, // 3: Observation.RequestsEntry - nil, // 4: Outcome.ObservedDonTimesEntry + nil, // 4: ObservedDonTimes.TimestampsBySequenceEntry + nil, // 5: Outcome.ObservedDonTimesEntry } var file_dontime_proto_depIdxs = []int32{ 3, // 0: Observation.requests:type_name -> Observation.RequestsEntry - 4, // 1: Outcome.observed_don_times:type_name -> Outcome.ObservedDonTimesEntry - 1, // 2: Outcome.ObservedDonTimesEntry.value:type_name -> ObservedDonTimes - 3, // [3:3] is the sub-list for method output_type - 3, // [3:3] is the sub-list for method input_type - 3, // [3:3] is the sub-list for extension type_name - 3, // [3:3] is the sub-list for extension extendee - 0, // [0:3] is the sub-list for field type_name + 4, // 1: ObservedDonTimes.timestampsBySequence:type_name -> ObservedDonTimes.TimestampsBySequenceEntry + 5, // 2: Outcome.observed_don_times:type_name -> Outcome.ObservedDonTimesEntry + 1, // 3: Outcome.ObservedDonTimesEntry.value:type_name -> ObservedDonTimes + 4, // [4:4] is the sub-list for method output_type + 4, // [4:4] is the sub-list for method input_type + 4, // [4:4] is the sub-list for extension type_name + 4, // [4:4] is the sub-list for extension extendee + 0, // [0:4] is the sub-list for field type_name } func init() { file_dontime_proto_init() } @@ -242,7 +258,7 @@ func file_dontime_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_dontime_proto_rawDesc), len(file_dontime_proto_rawDesc)), NumEnums: 0, - NumMessages: 5, + NumMessages: 6, NumExtensions: 0, NumServices: 0, }, diff --git a/pkg/workflows/dontime/pb/dontime.proto b/pkg/workflows/dontime/pb/dontime.proto index d0baecf3b3..ca82d41d76 100644 --- a/pkg/workflows/dontime/pb/dontime.proto +++ b/pkg/workflows/dontime/pb/dontime.proto @@ -10,7 +10,8 @@ message Observation { } message ObservedDonTimes { - repeated int64 timestamps = 1; + repeated int64 timestamps = 1 [deprecated = true]; + map timestampsBySequence = 2; } message Outcome { diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 136ecdec18..9228301b61 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -18,7 +18,10 @@ import ( "github.com/smartcontractkit/libocr/quorumhelper" "github.com/smartcontractkit/chainlink-common/pkg/beholder" + "github.com/smartcontractkit/chainlink-common/pkg/config" "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/settings/cresettings" + "github.com/smartcontractkit/chainlink-common/pkg/settings/limits" "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" ) @@ -92,6 +95,8 @@ type Plugin struct { minTimeIncrease int64 metrics pluginMetrics + + sequencedTSEnabled limits.RangeLimiter[config.Timestamp] } var _ ocr3types.ReportingPlugin[[]byte] = (*Plugin)(nil) @@ -113,16 +118,21 @@ func NewPlugin(store *Store, config ocr3types.ReportingPluginConfig, offchainCfg } return &Plugin{ - store: store, - config: config, - offChainConfig: offchainCfg, - lggr: logger.Named(lggr, "DONTimePlugin"), - batchSize: int(offchainCfg.MaxBatchSize), - minTimeIncrease: offchainCfg.MinTimeIncrease / int64(time.Millisecond), - metrics: metrics, + store: store, + config: config, + offChainConfig: offchainCfg, + lggr: logger.Named(lggr, "DONTimePlugin"), + batchSize: int(offchainCfg.MaxBatchSize), + minTimeIncrease: offchainCfg.MinTimeIncrease / int64(time.Millisecond), + metrics: metrics, + sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsEnabled.DefaultValue), }, nil } +func (p *Plugin) setSequencedTSEnabled(enabledRange limits.RangeLimiter[config.Timestamp]) { + p.sequencedTSEnabled = enabledRange +} + func (p *Plugin) Query(_ context.Context, _ ocr3types.OutcomeContext) (types.Query, error) { return nil, nil } @@ -150,7 +160,22 @@ func (p *Plugin) Observation(ctx context.Context, outctx ocr3types.OutcomeContex if err := proto.Unmarshal(outctx.PreviousOutcome, previousOutcome); err != nil { p.lggr.Errorf("failed to unmarshal previous outcome in Observation phase") } + // inspect ObservedDonTimes to determine if the feature flag has taken effect + var observedDonTimes *pb.ObservedDonTimes + for _, observedDonTimes = range previousOutcome.ObservedDonTimes { + break + } + if observedDonTimes != nil && len(observedDonTimes.TimestampsBySequence) > 0 { + // feature flag for sequenced timestamps is in effect + // TODO can len be zero for both Timestamps and TimestampsBySequence? + return p.unsequencedObservation(ctx, previousOutcome, query) + } + // unsequenced timestamps + return p.sequencedObservation(ctx, previousOutcome, query) +} +// unsequencedObservation executes the original Observation logic against the previous outcome's [pb.ObservedDonTimes.Timestamps]. +func (p *Plugin) unsequencedObservation(ctx context.Context, previousOutcome *pb.Outcome, query types.Query) (types.Observation, error) { sortedRequests := sortedRequests(p.store.GetRequests()) requests := map[string]int64{} // Maps executionID --> seqNum removedCount := 0 @@ -204,6 +229,61 @@ func (p *Plugin) Observation(ctx context.Context, outctx ocr3types.OutcomeContex return proto.MarshalOptions{Deterministic: true}.Marshal(observation) } +// sequencedObservation executes the updated Observation logic against the previous outcome's [pb.ObservedDonTimes.TimestampsBySequence]. +func (p *Plugin) sequencedObservation(ctx context.Context, previousOutcome *pb.Outcome, query types.Query) (types.Observation, error) { + sortedRequests := sortedRequests(p.store.GetRequests()) + requests := map[string]int64{} // Maps executionID --> seqNum + removedCount := 0 + for _, req := range sortedRequests { + // Validate request sequence number + var maxSeqNum int64 + times, ok := previousOutcome.ObservedDonTimes[req.WorkflowExecutionID] + if ok { + // We have seen this workflow before so check against the sequence + maxSeqNum = times.MaxSeqNum() + } + + if int64(req.SeqNum) > maxSeqNum+1 { + p.store.RemoveRequest(req.WorkflowExecutionID) + req.SendResponse(Response{ + WorkflowExecutionID: req.WorkflowExecutionID, + SeqNum: req.SeqNum, + Timestamp: 0, + Err: fmt.Errorf("requested seqNum %d for executionID %s is greater than expected based on the max seqNum observed so far %d", + req.SeqNum, req.WorkflowExecutionID, maxSeqNum), + }) + removedCount += 1 + continue + } + + requests[req.WorkflowExecutionID] = int64(req.SeqNum) + if len(requests) >= p.batchSize { + break + } + } + + overflowCount := len(sortedRequests) - len(requests) - removedCount + p.lggr.Debugw("Observation batch processed", + "inputRequests", len(sortedRequests), + "batchSize", p.batchSize, + "includedRequests", len(requests), + "removedRequests", removedCount, + "overflowRequests", overflowCount, + ) + if overflowCount > 0 { + p.lggr.Warnw("Observation batch overflow", "overflowRequests", overflowCount) + } + p.metrics.observationBatchOverflow.Record(ctx, int64(overflowCount)) + + observation := &pb.Observation{ + Timestamp: time.Now().UTC().UnixMilli(), + Requests: requests, + LimitByBatchSizeFlag: true, + } + + return proto.MarshalOptions{Deterministic: true}.Marshal(observation) +} + func (p *Plugin) ValidateObservation(_ context.Context, oc ocr3types.OutcomeContext, _ types.Query, ao types.AttributedObservation) error { return nil } @@ -213,6 +293,15 @@ func (p *Plugin) ObservationQuorum(_ context.Context, _ ocr3types.OutcomeContext } func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { + if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.Now())) == nil { //TODO inspect error + return p.sequencedOutcome(ctx, outctx, nil, aos) + } + + return p.unsequencedOutcome(ctx, outctx, nil, aos) +} + +// unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. +func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed type timestampNodePair struct { Timestamp int64 @@ -341,6 +430,147 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ return outcomeBytes, err } +// sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. +func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { + observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed + type timestampNodePair struct { + Timestamp int64 + NodeID int + OffsetFromMedian int64 + } + var timestampNodePairs []timestampNodePair + limitByBatchSizeFlagEnabled := true + + prevOutcome := &pb.Outcome{} + if err := proto.Unmarshal(outctx.PreviousOutcome, prevOutcome); err != nil { + p.lggr.Errorf("failed to unmarshal previous outcome in Outcome phase") + } + if prevOutcome.ObservedDonTimes == nil { + prevOutcome.ObservedDonTimes = make(map[string]*pb.ObservedDonTimes) + } else { + // At the transition point, we need to convert from the old slice format to maps + for _, observedTimes := range prevOutcome.ObservedDonTimes { + if len(observedTimes.Timestamps) > 0 { + for seqNum, ts := range observedTimes.Timestamps { + observedTimes.TimestampsBySequence[int64(seqNum)] = ts + } + observedTimes.Timestamps = nil + } + } + } + + for idx, ao := range aos { + observation := &pb.Observation{} + if err := proto.Unmarshal(ao.Observation, observation); err != nil { + p.lggr.Errorf("failed to unmarshal observation in Outcome phase") + continue + } + + if !observation.GetLimitByBatchSizeFlag() { + limitByBatchSizeFlagEnabled = false + } + + for id, requestSeqNum := range observation.Requests { + var currSeqNum int64 + if times, ok := prevOutcome.ObservedDonTimes[id]; ok { + currSeqNum = times.MaxSeqNum() + 1 + } + // We only count requests for the next sequence number and ignore all other ones. + if requestSeqNum == currSeqNum { + observationCounts[id]++ + } else if requestSeqNum > currSeqNum { + // This should never happen since we don't include out of sequence requests in the Observation phase + p.lggr.Errorf("request seqNum %d for executionID %s is greater than the current seqNum %d", + requestSeqNum, id, currSeqNum) + } + } + + timestampNodePairs = append(timestampNodePairs, timestampNodePair{Timestamp: observation.Timestamp, NodeID: idx}) + } + if len(timestampNodePairs) == 0 { + return nil, errors.New("no observation contains a valid timestamp") + } + + slices.SortFunc(timestampNodePairs, func(a, b timestampNodePair) int { + return cmp.Compare(a.Timestamp, b.Timestamp) + }) + donTime := timestampNodePairs[len(timestampNodePairs)/2].Timestamp + for i := range timestampNodePairs { + timestampNodePairs[i].OffsetFromMedian = timestampNodePairs[i].Timestamp - donTime + } + p.lggr.Debugw("Observed Node Timestamps", + "timestampNodePairs", timestampNodePairs, + "median", donTime, + "collectedDataPoints", len(timestampNodePairs), + "minOffsetFromMedian", timestampNodePairs[0].OffsetFromMedian, + "maxOffsetFromMedian", timestampNodePairs[len(timestampNodePairs)-1].OffsetFromMedian, + ) + + outcome := prevOutcome + + // Compare with prior outcome to ensure DON time never goes backward. + if donTime < outcome.Timestamp+p.minTimeIncrease { + p.lggr.Infow("DON Time incremented by minimum time increase to ensure time progression", "minTimeIncrease", p.minTimeIncrease) + donTime = outcome.Timestamp + p.minTimeIncrease + } + + p.lggr.Infow("New DON Time", "donTime", donTime) + outcome.Timestamp = donTime + + for id, numRequests := range observationCounts { + if numRequests > int64(p.config.F) { + observedDonTimes, ok := outcome.ObservedDonTimes[id] + if !ok { + observedDonTimes = &pb.ObservedDonTimes{TimestampsBySequence: make(map[int64]int64)} + } + currSeqNum := observedDonTimes.MaxSeqNum() + 1 + observedDonTimes.TimestampsBySequence[currSeqNum] = donTime + outcome.ObservedDonTimes[id] = observedDonTimes + } + } + + // Remove expired and empty workflow executions + for id, observedTimes := range outcome.ObservedDonTimes { + if observedTimes == nil || len(observedTimes.TimestampsBySequence) == 0 { + delete(outcome.ObservedDonTimes, id) + p.store.deleteExecutionID(id) + continue + } + if donTime >= observedTimes.EarliestTS()+p.offChainConfig.ExecutionRemovalTime.AsDuration().Milliseconds() { + delete(outcome.ObservedDonTimes, id) + p.store.deleteExecutionID(id) + } + } + + var outcomeBatchOverflowCount int64 + if len(outcome.ObservedDonTimes) > p.batchSize && limitByBatchSizeFlagEnabled { + ids := make([]string, 0, len(outcome.ObservedDonTimes)) + for id := range outcome.ObservedDonTimes { + ids = append(ids, id) + } + slices.Sort(ids) + outcomeBatchOverflowCount = int64(len(ids) - p.batchSize) + for _, id := range ids[p.batchSize:] { + delete(outcome.ObservedDonTimes, id) + } + p.lggr.Warnw("Trimmed outcome observed don times to batch size", + "batchSize", p.batchSize, + "removedEntries", outcomeBatchOverflowCount, + ) + } + + outcomeBytes, err := proto.MarshalOptions{Deterministic: true}.Marshal(outcome) + p.lggr.Infow("Outcome computed", + "observedDonTimesEntries", len(outcome.ObservedDonTimes), + "outcomeSizeBytes", len(outcomeBytes), + ) + p.metrics.donTime.Record(ctx, outcome.Timestamp) + p.metrics.donTimeEntries.Record(ctx, int64(len(outcome.ObservedDonTimes))) + p.metrics.outcomeBatchOverflow.Record(ctx, outcomeBatchOverflowCount) + p.metrics.outcomeSize.Record(ctx, int64(len(outcomeBytes))) + return outcomeBytes, err +} + func (p *Plugin) Reports(_ context.Context, _ uint64, outcome ocr3types.Outcome) ([]ocr3types.ReportPlus[[]byte], error) { allOraclesTransmitNow := &ocr3types.TransmissionSchedule{ Transmitters: make([]commontypes.OracleID, p.config.N), From acb7626c93355609b69537e2058dace712b3558d Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 10 Sep 2026 08:16:51 -0600 Subject: [PATCH 02/23] use donTime for feature flag --- pkg/workflows/dontime/plugin.go | 101 +++++++++++++------------------- 1 file changed, 42 insertions(+), 59 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 9228301b61..ae6f45fcab 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -293,22 +293,51 @@ func (p *Plugin) ObservationQuorum(_ context.Context, _ ocr3types.OutcomeContext } func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { - if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.Now())) == nil { //TODO inspect error - return p.sequencedOutcome(ctx, outctx, nil, aos) - } - - return p.unsequencedOutcome(ctx, outctx, nil, aos) -} - -// unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. -func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { - observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed type timestampNodePair struct { Timestamp int64 NodeID int OffsetFromMedian int64 } var timestampNodePairs []timestampNodePair + for idx, ao := range aos { + observation := &pb.Observation{} + if err := proto.Unmarshal(ao.Observation, observation); err != nil { + p.lggr.Errorf("failed to unmarshal observation in Outcome phase") + continue + } + + timestampNodePairs = append(timestampNodePairs, timestampNodePair{Timestamp: observation.Timestamp, NodeID: idx}) + } + + if len(timestampNodePairs) == 0 { + return nil, errors.New("no observation contains a valid timestamp") + } + + slices.SortFunc(timestampNodePairs, func(a, b timestampNodePair) int { + return cmp.Compare(a.Timestamp, b.Timestamp) + }) + donTime := timestampNodePairs[len(timestampNodePairs)/2].Timestamp + for i := range timestampNodePairs { + timestampNodePairs[i].OffsetFromMedian = timestampNodePairs[i].Timestamp - donTime + } + p.lggr.Debugw("Observed Node Timestamps", + "timestampNodePairs", timestampNodePairs, + "median", donTime, + "collectedDataPoints", len(timestampNodePairs), + "minOffsetFromMedian", timestampNodePairs[0].OffsetFromMedian, + "maxOffsetFromMedian", timestampNodePairs[len(timestampNodePairs)-1].OffsetFromMedian, + ) + + if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))) == nil { //TODO inspect error + return p.sequencedOutcome(ctx, outctx, nil, aos, donTime) + } + + return p.unsequencedOutcome(ctx, outctx, nil, aos, donTime) +} + +// unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. +func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, donTime int64) (ocr3types.Outcome, error) { + observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed limitByBatchSizeFlagEnabled := true prevOutcome := &pb.Outcome{} @@ -319,7 +348,7 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom prevOutcome.ObservedDonTimes = make(map[string]*pb.ObservedDonTimes) } - for idx, ao := range aos { + for _, ao := range aos { observation := &pb.Observation{} if err := proto.Unmarshal(ao.Observation, observation); err != nil { p.lggr.Errorf("failed to unmarshal observation in Outcome phase") @@ -344,27 +373,7 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom requestSeqNum, id, currSeqNum) } } - - timestampNodePairs = append(timestampNodePairs, timestampNodePair{Timestamp: observation.Timestamp, NodeID: idx}) } - if len(timestampNodePairs) == 0 { - return nil, errors.New("no observation contains a valid timestamp") - } - - slices.SortFunc(timestampNodePairs, func(a, b timestampNodePair) int { - return cmp.Compare(a.Timestamp, b.Timestamp) - }) - donTime := timestampNodePairs[len(timestampNodePairs)/2].Timestamp - for i := range timestampNodePairs { - timestampNodePairs[i].OffsetFromMedian = timestampNodePairs[i].Timestamp - donTime - } - p.lggr.Debugw("Observed Node Timestamps", - "timestampNodePairs", timestampNodePairs, - "median", donTime, - "collectedDataPoints", len(timestampNodePairs), - "minOffsetFromMedian", timestampNodePairs[0].OffsetFromMedian, - "maxOffsetFromMedian", timestampNodePairs[len(timestampNodePairs)-1].OffsetFromMedian, - ) outcome := prevOutcome @@ -431,14 +440,8 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom } // sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. -func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation) (ocr3types.Outcome, error) { +func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, donTime int64) (ocr3types.Outcome, error) { observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed - type timestampNodePair struct { - Timestamp int64 - NodeID int - OffsetFromMedian int64 - } - var timestampNodePairs []timestampNodePair limitByBatchSizeFlagEnabled := true prevOutcome := &pb.Outcome{} @@ -459,7 +462,7 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC } } - for idx, ao := range aos { + for _, ao := range aos { observation := &pb.Observation{} if err := proto.Unmarshal(ao.Observation, observation); err != nil { p.lggr.Errorf("failed to unmarshal observation in Outcome phase") @@ -484,27 +487,7 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC requestSeqNum, id, currSeqNum) } } - - timestampNodePairs = append(timestampNodePairs, timestampNodePair{Timestamp: observation.Timestamp, NodeID: idx}) - } - if len(timestampNodePairs) == 0 { - return nil, errors.New("no observation contains a valid timestamp") - } - - slices.SortFunc(timestampNodePairs, func(a, b timestampNodePair) int { - return cmp.Compare(a.Timestamp, b.Timestamp) - }) - donTime := timestampNodePairs[len(timestampNodePairs)/2].Timestamp - for i := range timestampNodePairs { - timestampNodePairs[i].OffsetFromMedian = timestampNodePairs[i].Timestamp - donTime } - p.lggr.Debugw("Observed Node Timestamps", - "timestampNodePairs", timestampNodePairs, - "median", donTime, - "collectedDataPoints", len(timestampNodePairs), - "minOffsetFromMedian", timestampNodePairs[0].OffsetFromMedian, - "maxOffsetFromMedian", timestampNodePairs[len(timestampNodePairs)-1].OffsetFromMedian, - ) outcome := prevOutcome From 880dc0eabca545a024f544de3311addef2ba7820 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 10 Sep 2026 10:14:30 -0600 Subject: [PATCH 03/23] adapt transmitter --- pkg/workflows/dontime/transmitter.go | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/pkg/workflows/dontime/transmitter.go b/pkg/workflows/dontime/transmitter.go index 8aa2ceb943..0181bb9743 100644 --- a/pkg/workflows/dontime/transmitter.go +++ b/pkg/workflows/dontime/transmitter.go @@ -5,10 +5,11 @@ import ( "google.golang.org/protobuf/proto" - "github.com/smartcontractkit/chainlink-common/pkg/logger" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" "github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3types" "github.com/smartcontractkit/libocr/offchainreporting2plus/types" + + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" ) var _ ocr3types.ContractTransmitter[[]byte] = (*Transmitter)(nil) @@ -50,8 +51,17 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 // Nodes behind on multiple requests may wait one OCR round per request. // Caching future times locally could be added as an optimization. - if len(donTimes.Timestamps) > request.SeqNum { - donTime := donTimes.Timestamps[request.SeqNum] + var donTime int64 + var ok bool + if len(donTimes.TimestampsBySequence) > 0 { + donTime, ok = donTimes.TimestampsBySequence[int64(request.SeqNum)] + } else { + ok = len(donTimes.Timestamps) > request.SeqNum + if ok { + donTime = donTimes.Timestamps[request.SeqNum] + } + } + if ok { t.store.RemoveRequest(executionID) // Make space for next request before delivering request.SendResponse(Response{ WorkflowExecutionID: executionID, From 4a8f4c711d6bc3872b6be9545e16752844d20999 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Mon, 14 Sep 2026 09:34:42 -0600 Subject: [PATCH 04/23] adjust donTime before usage; rm limitByBatchSizeFlag --- pkg/workflows/dontime/plugin.go | 64 ++++++++++++++------------------- 1 file changed, 26 insertions(+), 38 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index ae6f45fcab..6410900b44 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -328,18 +328,6 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ "maxOffsetFromMedian", timestampNodePairs[len(timestampNodePairs)-1].OffsetFromMedian, ) - if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))) == nil { //TODO inspect error - return p.sequencedOutcome(ctx, outctx, nil, aos, donTime) - } - - return p.unsequencedOutcome(ctx, outctx, nil, aos, donTime) -} - -// unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. -func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, donTime int64) (ocr3types.Outcome, error) { - observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed - limitByBatchSizeFlagEnabled := true - prevOutcome := &pb.Outcome{} if err := proto.Unmarshal(outctx.PreviousOutcome, prevOutcome); err != nil { p.lggr.Errorf("failed to unmarshal previous outcome in Outcome phase") @@ -348,6 +336,23 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom prevOutcome.ObservedDonTimes = make(map[string]*pb.ObservedDonTimes) } + // Compare with prior outcome to ensure DON time never goes backward. + if donTime < prevOutcome.Timestamp+p.minTimeIncrease { + p.lggr.Infow("DON Time incremented by minimum time increase to ensure time progression", "minTimeIncrease", p.minTimeIncrease) + donTime = prevOutcome.Timestamp + p.minTimeIncrease + } + + if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))) == nil { //TODO inspect error + return p.sequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) + } + + return p.unsequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) +} + +// unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. +func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { + observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed + for _, ao := range aos { observation := &pb.Observation{} if err := proto.Unmarshal(ao.Observation, observation); err != nil { @@ -355,10 +360,6 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom continue } - if !observation.GetLimitByBatchSizeFlag() { - limitByBatchSizeFlagEnabled = false - } - for id, requestSeqNum := range observation.Requests { var currSeqNum int64 if times, ok := prevOutcome.ObservedDonTimes[id]; ok { @@ -411,7 +412,7 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom } var outcomeBatchOverflowCount int64 - if len(outcome.ObservedDonTimes) > p.batchSize && limitByBatchSizeFlagEnabled { + if len(outcome.ObservedDonTimes) > p.batchSize { ids := make([]string, 0, len(outcome.ObservedDonTimes)) for id := range outcome.ObservedDonTimes { ids = append(ids, id) @@ -440,25 +441,16 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom } // sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. -func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, donTime int64) (ocr3types.Outcome, error) { +func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed - limitByBatchSizeFlagEnabled := true - prevOutcome := &pb.Outcome{} - if err := proto.Unmarshal(outctx.PreviousOutcome, prevOutcome); err != nil { - p.lggr.Errorf("failed to unmarshal previous outcome in Outcome phase") - } - if prevOutcome.ObservedDonTimes == nil { - prevOutcome.ObservedDonTimes = make(map[string]*pb.ObservedDonTimes) - } else { - // At the transition point, we need to convert from the old slice format to maps - for _, observedTimes := range prevOutcome.ObservedDonTimes { - if len(observedTimes.Timestamps) > 0 { - for seqNum, ts := range observedTimes.Timestamps { - observedTimes.TimestampsBySequence[int64(seqNum)] = ts - } - observedTimes.Timestamps = nil + // At the transition point, we need to convert from the old slice format to maps + for _, observedTimes := range prevOutcome.ObservedDonTimes { + if len(observedTimes.Timestamps) > 0 { + for seqNum, ts := range observedTimes.Timestamps { + observedTimes.TimestampsBySequence[int64(seqNum)] = ts } + observedTimes.Timestamps = nil } } @@ -469,10 +461,6 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC continue } - if !observation.GetLimitByBatchSizeFlag() { - limitByBatchSizeFlagEnabled = false - } - for id, requestSeqNum := range observation.Requests { var currSeqNum int64 if times, ok := prevOutcome.ObservedDonTimes[id]; ok { @@ -526,7 +514,7 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC } var outcomeBatchOverflowCount int64 - if len(outcome.ObservedDonTimes) > p.batchSize && limitByBatchSizeFlagEnabled { + if len(outcome.ObservedDonTimes) > p.batchSize { ids := make([]string, 0, len(outcome.ObservedDonTimes)) for id := range outcome.ObservedDonTimes { ids = append(ids, id) From 3a8a077761ef5ed169e521b87fd74dfff3fff08c Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 15 Sep 2026 18:21:39 -0600 Subject: [PATCH 05/23] cleanup; drop error cases --- pkg/workflows/dontime/pb/dontime.go | 2 +- pkg/workflows/dontime/plugin.go | 54 +++-------------------------- 2 files changed, 5 insertions(+), 51 deletions(-) diff --git a/pkg/workflows/dontime/pb/dontime.go b/pkg/workflows/dontime/pb/dontime.go index 8a30b52ac6..dbfd22f708 100644 --- a/pkg/workflows/dontime/pb/dontime.go +++ b/pkg/workflows/dontime/pb/dontime.go @@ -12,7 +12,7 @@ func (t *ObservedDonTimes) MaxSeqNum() (maxSeqNum int64) { return } -// EarliestTS returns the easliest timestamp value from TimestampsBySequence, or math.MaxInt64 if none exist. +// EarliestTS returns the earliest timestamp value from TimestampsBySequence, or math.MaxInt64 if none exist. func (t *ObservedDonTimes) EarliestTS() (earliestTS int64) { earliestTS = math.MaxInt64 for _, ts := range t.TimestampsBySequence { diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 6410900b44..aab6691d26 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -5,6 +5,7 @@ import ( "context" "errors" "fmt" + "maps" "slices" "time" @@ -142,11 +143,7 @@ func sortedRequests(requests map[string]*Request) []*Request { return nil } - ids := make([]string, 0, len(requests)) - for id := range requests { - ids = append(ids, id) - } - slices.Sort(ids) + ids := slices.Sorted(maps.Keys(requests)) sorted := make([]*Request, 0, len(ids)) for _, id := range ids { @@ -235,27 +232,6 @@ func (p *Plugin) sequencedObservation(ctx context.Context, previousOutcome *pb.O requests := map[string]int64{} // Maps executionID --> seqNum removedCount := 0 for _, req := range sortedRequests { - // Validate request sequence number - var maxSeqNum int64 - times, ok := previousOutcome.ObservedDonTimes[req.WorkflowExecutionID] - if ok { - // We have seen this workflow before so check against the sequence - maxSeqNum = times.MaxSeqNum() - } - - if int64(req.SeqNum) > maxSeqNum+1 { - p.store.RemoveRequest(req.WorkflowExecutionID) - req.SendResponse(Response{ - WorkflowExecutionID: req.WorkflowExecutionID, - SeqNum: req.SeqNum, - Timestamp: 0, - Err: fmt.Errorf("requested seqNum %d for executionID %s is greater than expected based on the max seqNum observed so far %d", - req.SeqNum, req.WorkflowExecutionID, maxSeqNum), - }) - removedCount += 1 - continue - } - requests[req.WorkflowExecutionID] = int64(req.SeqNum) if len(requests) >= p.batchSize { break @@ -342,6 +318,8 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ donTime = prevOutcome.Timestamp + p.minTimeIncrease } + p.lggr.Infow("New DON Time", "donTime", donTime) + if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))) == nil { //TODO inspect error return p.sequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) } @@ -368,23 +346,11 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ - } else if requestSeqNum > currSeqNum { - // This should never happen since we don't include out of sequence requests in the Observation phase - p.lggr.Errorf("request seqNum %d for executionID %s is greater than the number of observed don times %d", - requestSeqNum, id, currSeqNum) } } } outcome := prevOutcome - - // Compare with prior outcome to ensure DON time never goes backward. - if donTime < outcome.Timestamp+p.minTimeIncrease { - p.lggr.Infow("DON Time incremented by minimum time increase to ensure time progression", "minTimeIncrease", p.minTimeIncrease) - donTime = outcome.Timestamp + p.minTimeIncrease - } - - p.lggr.Infow("New DON Time", "donTime", donTime) outcome.Timestamp = donTime for id, numRequests := range observationCounts { @@ -469,23 +435,11 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ - } else if requestSeqNum > currSeqNum { - // This should never happen since we don't include out of sequence requests in the Observation phase - p.lggr.Errorf("request seqNum %d for executionID %s is greater than the current seqNum %d", - requestSeqNum, id, currSeqNum) } } } outcome := prevOutcome - - // Compare with prior outcome to ensure DON time never goes backward. - if donTime < outcome.Timestamp+p.minTimeIncrease { - p.lggr.Infow("DON Time incremented by minimum time increase to ensure time progression", "minTimeIncrease", p.minTimeIncrease) - donTime = outcome.Timestamp + p.minTimeIncrease - } - - p.lggr.Infow("New DON Time", "donTime", donTime) outcome.Timestamp = donTime for id, numRequests := range observationCounts { From 9d5202609978586d767d554da8773f14f3fcb3b3 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 15 Sep 2026 18:32:54 -0600 Subject: [PATCH 06/23] document map keys --- pkg/workflows/dontime/plugin.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index aab6691d26..ed0b3a3961 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -329,8 +329,8 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ // unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { - observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed - + // req_id->count - how many nodes reported where a new DON timestamp might be needed + observationCounts := map[string]int64{} for _, ao := range aos { observation := &pb.Observation{} if err := proto.Unmarshal(ao.Observation, observation); err != nil { @@ -408,7 +408,8 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom // sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { - observationCounts := map[string]int64{} // counts how many nodes reported where a new DON timestamp might be needed + // req_id->count - how many nodes reported where a new DON timestamp might be needed + observationCounts := map[string]int64{} // At the transition point, we need to convert from the old slice format to maps for _, observedTimes := range prevOutcome.ObservedDonTimes { From f87a6472e1381ce95b80e64020bb6b7566b93908 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 16 Sep 2026 07:39:39 -0600 Subject: [PATCH 07/23] restore unsequencedOutcome --- pkg/workflows/dontime/plugin.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index ed0b3a3961..c3a6d46637 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -436,6 +436,10 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ + } else if requestSeqNum > currSeqNum { + // This should never happen since we don't include out of sequence requests in the Observation phase + p.lggr.Errorf("request seqNum %d for executionID %s is greater than the number of observed don times %d", + requestSeqNum, id, currSeqNum) } } } From 304581ed1cc9004c1941c0442cf5631be131d3d2 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 16 Sep 2026 09:35:45 -0600 Subject: [PATCH 08/23] cleanup and error handling --- pkg/workflows/dontime/plugin.go | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index c3a6d46637..07f46330ad 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -165,10 +165,10 @@ func (p *Plugin) Observation(ctx context.Context, outctx ocr3types.OutcomeContex if observedDonTimes != nil && len(observedDonTimes.TimestampsBySequence) > 0 { // feature flag for sequenced timestamps is in effect // TODO can len be zero for both Timestamps and TimestampsBySequence? - return p.unsequencedObservation(ctx, previousOutcome, query) + return p.sequencedObservation(ctx, previousOutcome, query) } // unsequenced timestamps - return p.sequencedObservation(ctx, previousOutcome, query) + return p.unsequencedObservation(ctx, previousOutcome, query) } // unsequencedObservation executes the original Observation logic against the previous outcome's [pb.ObservedDonTimes.Timestamps]. @@ -320,11 +320,13 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ p.lggr.Infow("New DON Time", "donTime", donTime) - if p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))) == nil { //TODO inspect error - return p.sequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) + if err := p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))); err != nil { + if !errors.Is(err, limits.ErrorBoundLimited[config.Timestamp]{}) { + p.lggr.Warnw("Failed to check for sequenced timestamp feature flag", "err", err) + } + return p.unsequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) } - - return p.unsequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) + return p.sequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) } // unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. From c59607122003bb9d9fd0344df447fa0a14ec0811 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 16 Sep 2026 16:53:53 -0600 Subject: [PATCH 09/23] remove irrelevant test --- pkg/workflows/dontime/plugin_test.go | 12 +----------- 1 file changed, 1 insertion(+), 11 deletions(-) diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index 23a67759cf..311e903197 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -590,7 +590,7 @@ func TestPlugin_Outcome_TrimByBatchSize(t *testing.T) { prevOutcomeBytes, err := proto.Marshal(prevOutcome) require.NoError(t, err) - t.Run("trims when all observations set batch size flag", func(t *testing.T) { + t.Run("batch size enforced", func(t *testing.T) { outcome, err := plugin.Outcome(ctx, ocr3types.OutcomeContext{PreviousOutcome: prevOutcomeBytes}, query, makeObservations(true)) require.NoError(t, err) @@ -602,16 +602,6 @@ func TestPlugin_Outcome_TrimByBatchSize(t *testing.T) { require.Contains(t, outcomeProto.ObservedDonTimes, "workflow-b") require.NotContains(t, outcomeProto.ObservedDonTimes, "workflow-c") }) - - t.Run("does not trim when batch size flag is missing", func(t *testing.T) { - outcome, err := plugin.Outcome(ctx, ocr3types.OutcomeContext{PreviousOutcome: prevOutcomeBytes}, query, makeObservations(false)) - require.NoError(t, err) - - outcomeProto := &pb.Outcome{} - err = proto.Unmarshal(outcome, outcomeProto) - require.NoError(t, err) - require.Len(t, outcomeProto.ObservedDonTimes, 3) - }) } func TestPlugin_ExpiredRequest(t *testing.T) { From e6db8f1ef1ba4eccd303ef49bf665b62b628f9ff Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 08:45:00 -0600 Subject: [PATCH 10/23] cleanup --- pkg/workflows/dontime/plugin.go | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 07f46330ad..741331aca8 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -230,7 +230,6 @@ func (p *Plugin) unsequencedObservation(ctx context.Context, previousOutcome *pb func (p *Plugin) sequencedObservation(ctx context.Context, previousOutcome *pb.Outcome, query types.Query) (types.Observation, error) { sortedRequests := sortedRequests(p.store.GetRequests()) requests := map[string]int64{} // Maps executionID --> seqNum - removedCount := 0 for _, req := range sortedRequests { requests[req.WorkflowExecutionID] = int64(req.SeqNum) if len(requests) >= p.batchSize { @@ -238,12 +237,11 @@ func (p *Plugin) sequencedObservation(ctx context.Context, previousOutcome *pb.O } } - overflowCount := len(sortedRequests) - len(requests) - removedCount + overflowCount := len(sortedRequests) - len(requests) p.lggr.Debugw("Observation batch processed", "inputRequests", len(sortedRequests), "batchSize", p.batchSize, "includedRequests", len(requests), - "removedRequests", removedCount, "overflowRequests", overflowCount, ) if overflowCount > 0 { @@ -348,6 +346,10 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ + } else if requestSeqNum > currSeqNum { + // This should never happen since we don't include out of sequence requests in the Observation phase + p.lggr.Errorf("request seqNum %d for executionID %s is greater than the number of observed don times %d", + requestSeqNum, id, currSeqNum) } } } @@ -438,10 +440,6 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ - } else if requestSeqNum > currSeqNum { - // This should never happen since we don't include out of sequence requests in the Observation phase - p.lggr.Errorf("request seqNum %d for executionID %s is greater than the number of observed don times %d", - requestSeqNum, id, currSeqNum) } } } From 1cc7d35ad254f7da263f92f15fb610cefc4e8d5b Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 09:21:01 -0600 Subject: [PATCH 11/23] don't assume current sequence number --- pkg/workflows/dontime/plugin.go | 28 +++++++++++++++------------- 1 file changed, 15 insertions(+), 13 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 741331aca8..95bd32b2f1 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -412,8 +412,12 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom // sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { - // req_id->count - how many nodes reported where a new DON timestamp might be needed - observationCounts := map[string]int64{} + type reqSeq struct { + reqID string + seqNum int64 + } + // [req_id+seq_num]->count - how many nodes reported where a new DON timestamp might be needed + observationCounts := map[reqSeq]int64{} // At the transition point, we need to convert from the old slice format to maps for _, observedTimes := range prevOutcome.ObservedDonTimes { @@ -433,29 +437,27 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC } for id, requestSeqNum := range observation.Requests { - var currSeqNum int64 + // We only count requests for future sequence numbers and ignore all other ones. if times, ok := prevOutcome.ObservedDonTimes[id]; ok { - currSeqNum = times.MaxSeqNum() + 1 - } - // We only count requests for the next sequence number and ignore all other ones. - if requestSeqNum == currSeqNum { - observationCounts[id]++ + if requestSeqNum <= times.MaxSeqNum() { + continue + } } + observationCounts[reqSeq{id, requestSeqNum}]++ } } outcome := prevOutcome outcome.Timestamp = donTime - for id, numRequests := range observationCounts { + for key, numRequests := range observationCounts { if numRequests > int64(p.config.F) { - observedDonTimes, ok := outcome.ObservedDonTimes[id] + observedDonTimes, ok := outcome.ObservedDonTimes[key.reqID] if !ok { observedDonTimes = &pb.ObservedDonTimes{TimestampsBySequence: make(map[int64]int64)} } - currSeqNum := observedDonTimes.MaxSeqNum() + 1 - observedDonTimes.TimestampsBySequence[currSeqNum] = donTime - outcome.ObservedDonTimes[id] = observedDonTimes + observedDonTimes.TimestampsBySequence[key.seqNum] = donTime + outcome.ObservedDonTimes[key.reqID] = observedDonTimes } } From fc1cd3b7ba2be9e4983282e9e3b65c60952df275 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 09:43:41 -0600 Subject: [PATCH 12/23] dedupe --- pkg/workflows/dontime/plugin.go | 92 ++++++++++++--------------------- 1 file changed, 32 insertions(+), 60 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 95bd32b2f1..d5b873e88e 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -318,17 +318,43 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ p.lggr.Infow("New DON Time", "donTime", donTime) + var outcome *pb.Outcome if err := p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))); err != nil { if !errors.Is(err, limits.ErrorBoundLimited[config.Timestamp]{}) { p.lggr.Warnw("Failed to check for sequenced timestamp feature flag", "err", err) } - return p.unsequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) + outcome = p.unsequencedOutcome(aos, prevOutcome, donTime) + } else { + outcome = p.sequencedOutcome(aos, prevOutcome, donTime) } - return p.sequencedOutcome(ctx, outctx, nil, aos, prevOutcome, donTime) + + var outcomeBatchOverflowCount int64 + if len(outcome.ObservedDonTimes) > p.batchSize { + ids := slices.Sorted(maps.Keys(outcome.ObservedDonTimes)) + outcomeBatchOverflowCount = int64(len(ids) - p.batchSize) + for _, id := range ids[p.batchSize:] { + delete(outcome.ObservedDonTimes, id) + } + p.lggr.Warnw("Trimmed outcome observed don times to batch size", + "batchSize", p.batchSize, + "removedEntries", outcomeBatchOverflowCount, + ) + } + + outcomeBytes, err := proto.MarshalOptions{Deterministic: true}.Marshal(outcome) + p.lggr.Infow("Outcome computed", + "observedDonTimesEntries", len(outcome.ObservedDonTimes), + "outcomeSizeBytes", len(outcomeBytes), + ) + p.metrics.donTime.Record(ctx, outcome.Timestamp) + p.metrics.donTimeEntries.Record(ctx, int64(len(outcome.ObservedDonTimes))) + p.metrics.outcomeBatchOverflow.Record(ctx, outcomeBatchOverflowCount) + p.metrics.outcomeSize.Record(ctx, int64(len(outcomeBytes))) + return outcomeBytes, err } // unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. -func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { +func (p *Plugin) unsequencedOutcome(aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) *pb.Outcome { // req_id->count - how many nodes reported where a new DON timestamp might be needed observationCounts := map[string]int64{} for _, ao := range aos { @@ -380,38 +406,11 @@ func (p *Plugin) unsequencedOutcome(ctx context.Context, outctx ocr3types.Outcom p.store.deleteExecutionID(id) } } - - var outcomeBatchOverflowCount int64 - if len(outcome.ObservedDonTimes) > p.batchSize { - ids := make([]string, 0, len(outcome.ObservedDonTimes)) - for id := range outcome.ObservedDonTimes { - ids = append(ids, id) - } - slices.Sort(ids) - outcomeBatchOverflowCount = int64(len(ids) - p.batchSize) - for _, id := range ids[p.batchSize:] { - delete(outcome.ObservedDonTimes, id) - } - p.lggr.Warnw("Trimmed outcome observed don times to batch size", - "batchSize", p.batchSize, - "removedEntries", outcomeBatchOverflowCount, - ) - } - - outcomeBytes, err := proto.MarshalOptions{Deterministic: true}.Marshal(outcome) - p.lggr.Infow("Outcome computed", - "observedDonTimesEntries", len(outcome.ObservedDonTimes), - "outcomeSizeBytes", len(outcomeBytes), - ) - p.metrics.donTime.Record(ctx, outcome.Timestamp) - p.metrics.donTimeEntries.Record(ctx, int64(len(outcome.ObservedDonTimes))) - p.metrics.outcomeBatchOverflow.Record(ctx, outcomeBatchOverflowCount) - p.metrics.outcomeSize.Record(ctx, int64(len(outcomeBytes))) - return outcomeBytes, err + return outcome } // sequencedOutcome executed the updated outcome logic to produce a sequenced map of [pb.ObservedDonTimes.TimestampsBySequence]. -func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ types.Query, aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) (ocr3types.Outcome, error) { +func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) *pb.Outcome { type reqSeq struct { reqID string seqNum int64 @@ -473,34 +472,7 @@ func (p *Plugin) sequencedOutcome(ctx context.Context, outctx ocr3types.OutcomeC p.store.deleteExecutionID(id) } } - - var outcomeBatchOverflowCount int64 - if len(outcome.ObservedDonTimes) > p.batchSize { - ids := make([]string, 0, len(outcome.ObservedDonTimes)) - for id := range outcome.ObservedDonTimes { - ids = append(ids, id) - } - slices.Sort(ids) - outcomeBatchOverflowCount = int64(len(ids) - p.batchSize) - for _, id := range ids[p.batchSize:] { - delete(outcome.ObservedDonTimes, id) - } - p.lggr.Warnw("Trimmed outcome observed don times to batch size", - "batchSize", p.batchSize, - "removedEntries", outcomeBatchOverflowCount, - ) - } - - outcomeBytes, err := proto.MarshalOptions{Deterministic: true}.Marshal(outcome) - p.lggr.Infow("Outcome computed", - "observedDonTimesEntries", len(outcome.ObservedDonTimes), - "outcomeSizeBytes", len(outcomeBytes), - ) - p.metrics.donTime.Record(ctx, outcome.Timestamp) - p.metrics.donTimeEntries.Record(ctx, int64(len(outcome.ObservedDonTimes))) - p.metrics.outcomeBatchOverflow.Record(ctx, outcomeBatchOverflowCount) - p.metrics.outcomeSize.Record(ctx, int64(len(outcomeBytes))) - return outcomeBytes, err + return outcome } func (p *Plugin) Reports(_ context.Context, _ uint64, outcome ocr3types.Outcome) ([]ocr3types.ReportPlus[[]byte], error) { From 3039cbdfd19203bf0c38d965fef51f277bc6d306 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 09:53:47 -0600 Subject: [PATCH 13/23] simplify Observation --- pkg/workflows/dontime/plugin.go | 75 --------------------------------- 1 file changed, 75 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index d5b873e88e..063c083d42 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -153,81 +153,6 @@ func sortedRequests(requests map[string]*Request) []*Request { } func (p *Plugin) Observation(ctx context.Context, outctx ocr3types.OutcomeContext, query types.Query) (types.Observation, error) { - previousOutcome := &pb.Outcome{} - if err := proto.Unmarshal(outctx.PreviousOutcome, previousOutcome); err != nil { - p.lggr.Errorf("failed to unmarshal previous outcome in Observation phase") - } - // inspect ObservedDonTimes to determine if the feature flag has taken effect - var observedDonTimes *pb.ObservedDonTimes - for _, observedDonTimes = range previousOutcome.ObservedDonTimes { - break - } - if observedDonTimes != nil && len(observedDonTimes.TimestampsBySequence) > 0 { - // feature flag for sequenced timestamps is in effect - // TODO can len be zero for both Timestamps and TimestampsBySequence? - return p.sequencedObservation(ctx, previousOutcome, query) - } - // unsequenced timestamps - return p.unsequencedObservation(ctx, previousOutcome, query) -} - -// unsequencedObservation executes the original Observation logic against the previous outcome's [pb.ObservedDonTimes.Timestamps]. -func (p *Plugin) unsequencedObservation(ctx context.Context, previousOutcome *pb.Outcome, query types.Query) (types.Observation, error) { - sortedRequests := sortedRequests(p.store.GetRequests()) - requests := map[string]int64{} // Maps executionID --> seqNum - removedCount := 0 - for _, req := range sortedRequests { - // Validate request sequence number - numObservedDonTimes := 0 - times, ok := previousOutcome.ObservedDonTimes[req.WorkflowExecutionID] - if ok { - // We have seen this workflow before so check against the sequence - numObservedDonTimes = len(times.Timestamps) - } - - if req.SeqNum > numObservedDonTimes { - p.store.RemoveRequest(req.WorkflowExecutionID) - req.SendResponse(Response{ - WorkflowExecutionID: req.WorkflowExecutionID, - SeqNum: req.SeqNum, - Timestamp: 0, - Err: fmt.Errorf("requested seqNum %d for executionID %s is greater than the number of observed don times %d", - req.SeqNum, req.WorkflowExecutionID, numObservedDonTimes), - }) - removedCount += 1 - continue - } - - requests[req.WorkflowExecutionID] = int64(req.SeqNum) - if len(requests) >= p.batchSize { - break - } - } - - overflowCount := len(sortedRequests) - len(requests) - removedCount - p.lggr.Debugw("Observation batch processed", - "inputRequests", len(sortedRequests), - "batchSize", p.batchSize, - "includedRequests", len(requests), - "removedRequests", removedCount, - "overflowRequests", overflowCount, - ) - if overflowCount > 0 { - p.lggr.Warnw("Observation batch overflow", "overflowRequests", overflowCount) - } - p.metrics.observationBatchOverflow.Record(ctx, int64(overflowCount)) - - observation := &pb.Observation{ - Timestamp: time.Now().UTC().UnixMilli(), - Requests: requests, - LimitByBatchSizeFlag: true, - } - - return proto.MarshalOptions{Deterministic: true}.Marshal(observation) -} - -// sequencedObservation executes the updated Observation logic against the previous outcome's [pb.ObservedDonTimes.TimestampsBySequence]. -func (p *Plugin) sequencedObservation(ctx context.Context, previousOutcome *pb.Outcome, query types.Query) (types.Observation, error) { sortedRequests := sortedRequests(p.store.GetRequests()) requests := map[string]int64{} // Maps executionID --> seqNum for _, req := range sortedRequests { From 823f8fed97393bd05b4960de93493088142352b9 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 10:15:56 -0600 Subject: [PATCH 14/23] fix test --- pkg/workflows/dontime/plugin_test.go | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index 311e903197..a26e7e6352 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -146,7 +146,8 @@ func TestPlugin_ValidateObservation(t *testing.T) { require.NoError(t, err) }) - t.Run("Invalid sequence number", func(t *testing.T) { + //TODO no such thing any more - or pass an old one? + t.Run("Valid skipped sequence number", func(t *testing.T) { store := NewStore(DefaultRequestTimeout) plugin, err := NewPlugin(store, config, offchainCfg, lggr) require.NoError(t, err) @@ -160,13 +161,18 @@ func TestPlugin_ValidateObservation(t *testing.T) { // Add single request to queue executionID := "workflow-123" - requestCh := store.RequestDonTime(executionID, 1) + _ = store.RequestDonTime(executionID, 1) - _, err = plugin.Observation(ctx, outcomeCtx, query) + observation, err := plugin.Observation(ctx, outcomeCtx, query) require.NoError(t, err) - response := <-requestCh - require.ErrorContains(t, response.Err, "requested seqNum 1 for executionID workflow-123 is greater than the number of observed don times 0") + ao := types.AttributedObservation{ + Observation: observation, + Observer: commontypes.OracleID(1), + } + + err = plugin.ValidateObservation(ctx, outcomeCtx, query, ao) + require.NoError(t, err) }) } From 92a8f3f8afb58369e6fba6d2d8d0b19f83540584 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 17 Sep 2026 16:37:38 -0600 Subject: [PATCH 15/23] remove todo --- pkg/workflows/dontime/plugin_test.go | 1 - 1 file changed, 1 deletion(-) diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index a26e7e6352..1b96cac154 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -146,7 +146,6 @@ func TestPlugin_ValidateObservation(t *testing.T) { require.NoError(t, err) }) - //TODO no such thing any more - or pass an old one? t.Run("Valid skipped sequence number", func(t *testing.T) { store := NewStore(DefaultRequestTimeout) plugin, err := NewPlugin(store, config, offchainCfg, lggr) From 8d54f75bc78938fff5785ac74f77fa7f08b8727e Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 29 Sep 2026 07:13:42 -0500 Subject: [PATCH 16/23] sorted map iteration --- pkg/workflows/dontime/plugin.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 063c083d42..08f0a3656c 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -297,10 +297,6 @@ func (p *Plugin) unsequencedOutcome(aos []types.AttributedObservation, prevOutco // We only count requests for the next sequence number and ignore all other ones. if requestSeqNum == currSeqNum { observationCounts[id]++ - } else if requestSeqNum > currSeqNum { - // This should never happen since we don't include out of sequence requests in the Observation phase - p.lggr.Errorf("request seqNum %d for executionID %s is greater than the number of observed don times %d", - requestSeqNum, id, currSeqNum) } } } @@ -374,7 +370,10 @@ func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome outcome := prevOutcome outcome.Timestamp = donTime - for key, numRequests := range observationCounts { + for _, key := range slices.SortedFunc(maps.Keys(observationCounts), func(a, b reqSeq) int { + return cmp.Or(cmp.Compare(a.reqID, b.reqID), cmp.Compare(a.seqNum, b.seqNum)) + }) { + numRequests := observationCounts[key] if numRequests > int64(p.config.F) { observedDonTimes, ok := outcome.ObservedDonTimes[key.reqID] if !ok { From 2c8926a2db7fe4c586413af963060a038328d8c9 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 6 Oct 2026 13:11:35 -0500 Subject: [PATCH 17/23] tests and fixes --- pkg/workflows/dontime/pb/dontime.go | 7 ++- pkg/workflows/dontime/plugin.go | 3 + pkg/workflows/dontime/plugin_test.go | 31 +++++++++- pkg/workflows/dontime/transmitter_test.go | 70 +++++++++++++---------- 4 files changed, 76 insertions(+), 35 deletions(-) diff --git a/pkg/workflows/dontime/pb/dontime.go b/pkg/workflows/dontime/pb/dontime.go index dbfd22f708..9fb2107ac0 100644 --- a/pkg/workflows/dontime/pb/dontime.go +++ b/pkg/workflows/dontime/pb/dontime.go @@ -2,14 +2,15 @@ package pb import "math" -// MaxSeqNum returns the max sequence number from TimestampsBySequence, or 0 if none exist. -func (t *ObservedDonTimes) MaxSeqNum() (maxSeqNum int64) { +// MaxSeqNum returns the max sequence number from TimestampsBySequence, or -1 if none exist. +func (t *ObservedDonTimes) MaxSeqNum() int64 { + var maxSeqNum int64 = -1 for seqNum := range t.TimestampsBySequence { if seqNum > maxSeqNum { maxSeqNum = seqNum } } - return + return maxSeqNum } // EarliestTS returns the earliest timestamp value from TimestampsBySequence, or math.MaxInt64 if none exist. diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 08f0a3656c..1466e95938 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -340,6 +340,7 @@ func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome observationCounts := map[reqSeq]int64{} // At the transition point, we need to convert from the old slice format to maps + //TODO consider reverse when disabling... for _, observedTimes := range prevOutcome.ObservedDonTimes { if len(observedTimes.Timestamps) > 0 { for seqNum, ts := range observedTimes.Timestamps { @@ -378,6 +379,8 @@ func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome observedDonTimes, ok := outcome.ObservedDonTimes[key.reqID] if !ok { observedDonTimes = &pb.ObservedDonTimes{TimestampsBySequence: make(map[int64]int64)} + } else if observedDonTimes.TimestampsBySequence == nil { + observedDonTimes.TimestampsBySequence = make(map[int64]int64) } observedDonTimes.TimestampsBySequence[key.seqNum] = donTime outcome.ObservedDonTimes[key.reqID] = observedDonTimes diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index 1b96cac154..a586f963f4 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -4,6 +4,7 @@ import ( "testing" "time" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/durationpb" @@ -13,7 +14,10 @@ import ( "github.com/smartcontractkit/libocr/offchainreporting2/types" "github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3types" + "github.com/smartcontractkit/chainlink-common/pkg/config" "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/settings" + "github.com/smartcontractkit/chainlink-common/pkg/settings/limits" "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" ) @@ -176,13 +180,26 @@ func TestPlugin_ValidateObservation(t *testing.T) { } func TestPlugin_Outcome(t *testing.T) { + t.Run("sequenced", func(t *testing.T) { testPlugin_Outcome(t, true) }) + t.Run("unsequenced", func(t *testing.T) { testPlugin_Outcome(t, false) }) +} + +func testPlugin_Outcome(t *testing.T, sequenced bool) { lggr := logger.Test(t) store := NewStore(DefaultRequestTimeout) - config, offchainCfg := newTestPluginConfig(t), newTestPluginOffchainConfig(t) + cfg, offchainCfg := newTestPluginConfig(t), newTestPluginOffchainConfig(t) ctx := t.Context() - plugin, err := NewPlugin(store, config, offchainCfg, lggr) + plugin, err := NewPlugin(store, cfg, offchainCfg, lggr) require.NoError(t, err) + if sequenced { + plugin.setSequencedTSEnabled(limits.NewRangeLimiter( + settings.Range[config.Timestamp]{ + Lower: config.NewTimestamp(time.Now()), + Upper: config.NewTimestamp(time.Now().Add(time.Hour)), + }, + )) + } query, err := plugin.Query(ctx, ocr3types.OutcomeContext{PreviousOutcome: []byte("")}) require.NoError(t, err) @@ -238,7 +255,15 @@ func TestPlugin_Outcome(t *testing.T) { err = proto.Unmarshal(outcome, outcomeProto) require.NoError(t, err) require.Equal(t, timestamp, outcomeProto.Timestamp) - require.Equal(t, []int64{timestamp}, outcomeProto.ObservedDonTimes[executionID].Timestamps) + if observed, ok := outcomeProto.ObservedDonTimes[executionID]; assert.True(t, ok) { + if sequenced { + require.Equal(t, map[int64]int64{0: timestamp}, observed.TimestampsBySequence) + require.Empty(t, observed.Timestamps) + } else { + require.Equal(t, []int64{timestamp}, observed.Timestamps) + require.Empty(t, observed.TimestampsBySequence) + } + } } func TestPlugin_Outcome_SequenceNumberHandling(t *testing.T) { diff --git a/pkg/workflows/dontime/transmitter_test.go b/pkg/workflows/dontime/transmitter_test.go index 18f9023498..31c12eb39e 100644 --- a/pkg/workflows/dontime/transmitter_test.go +++ b/pkg/workflows/dontime/transmitter_test.go @@ -7,10 +7,11 @@ import ( "github.com/stretchr/testify/require" "google.golang.org/protobuf/proto" - "github.com/smartcontractkit/chainlink-common/pkg/logger" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" "github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3types" "github.com/smartcontractkit/libocr/offchainreporting2plus/types" + + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/workflows/dontime/pb" ) func TestTransmitter_TransmitDonTimeRequest(t *testing.T) { @@ -20,34 +21,45 @@ func TestTransmitter_TransmitDonTimeRequest(t *testing.T) { transmitter := NewTransmitter(lggr, store, "") - // Create request for second donTime in sequence - executionID := "workflow-123" - timeRequest := store.RequestDonTime(executionID, 1) - timestamp := time.Now().UnixMilli() - outcome := &pb.Outcome{ - Timestamp: timestamp, - ObservedDonTimes: map[string]*pb.ObservedDonTimes{ - executionID: {Timestamps: []int64{timestamp - int64(time.Second), timestamp}}, - }, - } - r := ocr3types.ReportWithInfo[[]byte]{} - var err error - r.Report, err = proto.Marshal(outcome) - require.NoError(t, err) - err = transmitter.Transmit(ctx, types.ConfigDigest{}, 0, r, []types.AttributedOnchainSignature{}) - require.NoError(t, err) - - select { - case donTimeResp := <-timeRequest: - require.Equal(t, timestamp, donTimeResp.Timestamp) - require.Equal(t, executionID, donTimeResp.WorkflowExecutionID) - require.Equal(t, 1, donTimeResp.SeqNum) - require.NoError(t, donTimeResp.Err) - case <-ctx.Done(): - t.Fatal("failed to retrieve donTime from request channel") - } + for _, tc := range []struct { + name string + observed *pb.ObservedDonTimes + }{ + {"unsequenced", &pb.ObservedDonTimes{Timestamps: []int64{timestamp - int64(time.Second), timestamp}}}, + {"sequenced", &pb.ObservedDonTimes{TimestampsBySequence: map[int64]int64{0: timestamp - int64(time.Second), 1: timestamp}}}, + {"both", &pb.ObservedDonTimes{Timestamps: []int64{timestamp - int64(time.Second), timestamp}, + TimestampsBySequence: map[int64]int64{0: timestamp - int64(time.Second), 1: timestamp}}}, + } { + t.Run(tc.name, func(t *testing.T) { + // Create request for second donTime in sequence + executionID := "workflow-123" + timeRequest := store.RequestDonTime(executionID, 1) + + outcome := &pb.Outcome{ + Timestamp: timestamp, + ObservedDonTimes: map[string]*pb.ObservedDonTimes{executionID: tc.observed}, + } - require.Empty(t, store.GetRequest(executionID)) + r := ocr3types.ReportWithInfo[[]byte]{} + var err error + r.Report, err = proto.Marshal(outcome) + require.NoError(t, err) + err = transmitter.Transmit(ctx, types.ConfigDigest{}, 0, r, []types.AttributedOnchainSignature{}) + require.NoError(t, err) + + select { + case donTimeResp := <-timeRequest: + require.Equal(t, timestamp, donTimeResp.Timestamp) + require.Equal(t, executionID, donTimeResp.WorkflowExecutionID) + require.Equal(t, 1, donTimeResp.SeqNum) + require.NoError(t, donTimeResp.Err) + case <-ctx.Done(): + t.Fatal("failed to retrieve donTime from request channel") + } + + require.Empty(t, store.GetRequest(executionID)) + }) + } } From 02cb2e4296de8fed3f3efd42021ad1e901fc7276 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Tue, 6 Oct 2026 15:28:02 -0500 Subject: [PATCH 18/23] handle flag reversion; fix store & transmitter --- pkg/workflows/dontime/plugin.go | 16 ++++++++++++++- pkg/workflows/dontime/plugin_test.go | 2 +- pkg/workflows/dontime/store.go | 30 ++++++---------------------- pkg/workflows/dontime/store_test.go | 17 ++++++++++++++++ pkg/workflows/dontime/transmitter.go | 12 +++++++++-- 5 files changed, 49 insertions(+), 28 deletions(-) diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 1466e95938..981a7083c9 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -280,6 +280,21 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ // unsequencedOutcome executes the original outcome logic to produce an unsequenced slice of [pb.ObservedDonTimes.Timestamps]. func (p *Plugin) unsequencedOutcome(aos []types.AttributedObservation, prevOutcome *pb.Outcome, donTime int64) *pb.Outcome { + // If disabling the feature flag, then at the transition point, we need to convert from the map format to the slices + for _, observedTimes := range prevOutcome.ObservedDonTimes { + if mapLen := len(observedTimes.TimestampsBySequence); mapLen > 0 { + observedTimes.Timestamps = make([]int64, mapLen) + for seqNum, ts := range observedTimes.TimestampsBySequence { + if sliceLen := len(observedTimes.Timestamps); seqNum > int64(sliceLen) { + // There must have been a gap in the sequence so grow the slice + observedTimes.Timestamps = slices.Grow(observedTimes.Timestamps, int(seqNum)-sliceLen) + } + observedTimes.Timestamps[seqNum] = ts + } + observedTimes.TimestampsBySequence = nil + } + } + // req_id->count - how many nodes reported where a new DON timestamp might be needed observationCounts := map[string]int64{} for _, ao := range aos { @@ -340,7 +355,6 @@ func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome observationCounts := map[reqSeq]int64{} // At the transition point, we need to convert from the old slice format to maps - //TODO consider reverse when disabling... for _, observedTimes := range prevOutcome.ObservedDonTimes { if len(observedTimes.Timestamps) > 0 { for seqNum, ts := range observedTimes.Timestamps { diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index a586f963f4..aab3671657 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -564,7 +564,7 @@ func TestPlugin_FinishedExecutions(t *testing.T) { }) t.Run("Transmit: delete removed executionIDs", func(t *testing.T) { - store.setDonTimes("workflow-123", []int64{time.Now().UnixMilli()}) + store.setDonTimes("workflow-123", map[int64]int64{0: time.Now().UnixMilli()}) r := ocr3types.ReportWithInfo[[]byte]{} r.Report, err = proto.Marshal(outcomeProto) diff --git a/pkg/workflows/dontime/store.go b/pkg/workflows/dontime/store.go index 455a0cf4b4..f5180af027 100644 --- a/pkg/workflows/dontime/store.go +++ b/pkg/workflows/dontime/store.go @@ -19,9 +19,9 @@ type Store struct { requests map[string]*Request // Maps workflow execution ID to request requestTimeout time.Duration - // donTimes holds ordered sequence timestamps generated for consecutive workflow requests - // i.e. ExecutionID --> [timestamp-0, timestamp-1 , ...] - donTimes map[string][]int64 + // donTimes holds sequence timestamps generated for workflow requests + // executionID -> sequenceNumber -> timestamp + donTimes map[string]map[int64]int64 lastObservedDonTime int64 mu sync.Mutex } @@ -30,7 +30,7 @@ func NewStore(requestTimeout time.Duration) *Store { return &Store{ requests: make(map[string]*Request), requestTimeout: requestTimeout, - donTimes: make(map[string][]int64), + donTimes: make(map[string]map[int64]int64), lastObservedDonTime: 0, mu: sync.Mutex{}, } @@ -124,30 +124,12 @@ func (s *Store) GetDonTimeForSeqNum(executionID string, seqNum int) *int64 { s.mu.Lock() defer s.mu.Unlock() if times, ok := s.donTimes[executionID]; ok { - if len(times) > seqNum { - return ×[seqNum] - } + return new(times[int64(seqNum)]) } return nil } -func (s *Store) GetDonTimes(executionID string) ([]int64, error) { - s.mu.Lock() - defer s.mu.Unlock() - - if times, ok := s.donTimes[executionID]; ok { - return times, nil - } - return []int64{}, fmt.Errorf("no don time for executionID %s", executionID) -} - -func (s *Store) setDonTimes(executionID string, donTimes []int64) { - s.mu.Lock() - defer s.mu.Unlock() - s.donTimes[executionID] = donTimes -} - -func (s *Store) replaceDonTimes(donTimes map[string][]int64) { +func (s *Store) replaceDonTimes(donTimes map[string]map[int64]int64) { s.mu.Lock() defer s.mu.Unlock() diff --git a/pkg/workflows/dontime/store_test.go b/pkg/workflows/dontime/store_test.go index b5ff631dfd..a471586dba 100644 --- a/pkg/workflows/dontime/store_test.go +++ b/pkg/workflows/dontime/store_test.go @@ -1,6 +1,7 @@ package dontime import ( + "fmt" "testing" "time" @@ -24,3 +25,19 @@ func TestStore_RequestExpiresWithoutPlugin(t *testing.T) { require.Nil(t, store.GetRequest(executionID)) } + +func (s *Store) GetDonTimes(executionID string) (map[int64]int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + + if times, ok := s.donTimes[executionID]; ok { + return times, nil + } + return map[int64]int64{}, fmt.Errorf("no don time for executionID %s", executionID) +} + +func (s *Store) setDonTimes(executionID string, donTimes map[int64]int64) { + s.mu.Lock() + defer s.mu.Unlock() + s.donTimes[executionID] = donTimes +} diff --git a/pkg/workflows/dontime/transmitter.go b/pkg/workflows/dontime/transmitter.go index 0181bb9743..4a26e9c48e 100644 --- a/pkg/workflows/dontime/transmitter.go +++ b/pkg/workflows/dontime/transmitter.go @@ -34,9 +34,17 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 return err } - currentDonTimes := make(map[string][]int64, len(outcome.ObservedDonTimes)) + currentDonTimes := make(map[string]map[int64]int64, len(outcome.ObservedDonTimes)) for id, observedDonTimes := range outcome.ObservedDonTimes { - currentDonTimes[id] = observedDonTimes.Timestamps + if len(observedDonTimes.Timestamps) > 0 { + m := make(map[int64]int64) + for i, t := range observedDonTimes.Timestamps { + m[int64(i)] = t + } + currentDonTimes[id] = m + } else { + currentDonTimes[id] = observedDonTimes.TimestampsBySequence + } } t.store.replaceDonTimes(currentDonTimes) t.store.setLastObservedDonTime(outcome.Timestamp) From c5ba8299670178fca9b22d71b78d926bfc2652e5 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 7 Oct 2026 15:54:08 -0500 Subject: [PATCH 19/23] feedback --- pkg/settings/cresettings/settings.go | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/pkg/settings/cresettings/settings.go b/pkg/settings/cresettings/settings.go index ab5fa62641..12a2b33995 100644 --- a/pkg/settings/cresettings/settings.go +++ b/pkg/settings/cresettings/settings.go @@ -53,7 +53,8 @@ var Config Schema var ( year2100 = time.Date(2100, 1, 1, 0, 0, 0, 0, time.UTC) - disabledFeatureTimeRange = TimeRange(year2100, time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC)) + year2101 = time.Date(2101, 1, 1, 0, 0, 0, 0, time.UTC) + disabledFeatureTimeRange = TimeRange(year2100, year2101) ) var Default = Schema{ @@ -342,10 +343,10 @@ var Default = Schema{ RequestTimeout: Duration(30 * time.Second), }, - FeatureHTTPTriggerNewExecutionIDsActivePeriod: disabledFeatureTimeRange, - FeatureChainCapabilityHashBasedOCRActivePeriod: disabledFeatureTimeRange, - FeatureEVMWriteReportL1FeeActivePeriod: disabledFeatureTimeRange, - FeatureAptosWriteReportBlockTimestampActivePeriod: disabledFeatureTimeRange, + FeatureHTTPTriggerNewExecutionIDsActivePeriod: disabledFeatureTimeRange, + FeatureChainCapabilityHashBasedOCRActivePeriod: disabledFeatureTimeRange, + FeatureEVMWriteReportL1FeeActivePeriod: disabledFeatureTimeRange, + FeatureAptosWriteReportBlockTimestampActivePeriod: disabledFeatureTimeRange, // ON by default: covers all possible timestamps including zero time.Time{}, // so WorkflowTag is included in the hash matching current prod behavior. // After rollout, set to far-future window to exclude WorkflowTag. @@ -356,7 +357,7 @@ var Default = Schema{ // cover "now" only after FeatureRequestHashIncludeWorkflowTag is muted // on every DON member, so DBs can heal without producing tag-driven // hash divergence during the fill window. - FeatureWorkflowTagBackfillActivePeriod: disabledFeatureTimeRange, + FeatureWorkflowTagBackfillActivePeriod: disabledFeatureTimeRange, FeatureConsensusStricterMedianQuorumActivePeriod: disabledFeatureTimeRange, FeatureConsensusIncludeAllTimestampsActivePeriod: disabledFeatureTimeRange, }, From b236a195f365de3ca66d1c624d649ed01e9fd105 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 7 Oct 2026 19:43:46 -0500 Subject: [PATCH 20/23] fix error handling; log metadata --- pkg/settings/limits/range.go | 6 +++--- pkg/workflows/dontime/plugin.go | 2 +- pkg/workflows/dontime/transmitter.go | 11 ++++++++--- 3 files changed, 12 insertions(+), 7 deletions(-) diff --git a/pkg/settings/limits/range.go b/pkg/settings/limits/range.go index a77e77f008..a21fe9116c 100644 --- a/pkg/settings/limits/range.go +++ b/pkg/settings/limits/range.go @@ -15,14 +15,14 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/settings" ) -// BoundLimiter is a limiter for simple bounds checks. +// RangeLimiter is a limiter for bounded range checks. type RangeLimiter[N Number] interface { Limiter[settings.Range[N]] - // Check returns ErrorBoundLimited if the value is above the limit. + // Check returns ErrorRangeLimited if the value is above the limit. Check(context.Context, N) error } -// NewRangeLimiter returns a RangeLimiter with the given lower bounds. +// NewRangeLimiter returns a RangeLimiter with the given bounds. func NewRangeLimiter[N Number](bounds settings.Range[N]) RangeLimiter[N] { return &simpleRangeLimiter[N]{bounds: bounds} } diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 981a7083c9..1118fe3711 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -245,7 +245,7 @@ func (p *Plugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, _ var outcome *pb.Outcome if err := p.sequencedTSEnabled.Check(ctx, config.NewTimestamp(time.UnixMilli(donTime))); err != nil { - if !errors.Is(err, limits.ErrorBoundLimited[config.Timestamp]{}) { + if !errors.Is(err, limits.ErrorRangeLimited[config.Timestamp]{}) { p.lggr.Warnw("Failed to check for sequenced timestamp feature flag", "err", err) } outcome = p.unsequencedOutcome(aos, prevOutcome, donTime) diff --git a/pkg/workflows/dontime/transmitter.go b/pkg/workflows/dontime/transmitter.go index 4a26e9c48e..8f33c3d2bb 100644 --- a/pkg/workflows/dontime/transmitter.go +++ b/pkg/workflows/dontime/transmitter.go @@ -34,22 +34,24 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 return err } + var total int currentDonTimes := make(map[string]map[int64]int64, len(outcome.ObservedDonTimes)) for id, observedDonTimes := range outcome.ObservedDonTimes { if len(observedDonTimes.Timestamps) > 0 { m := make(map[int64]int64) - for i, t := range observedDonTimes.Timestamps { - m[int64(i)] = t + for i, ts := range observedDonTimes.Timestamps { + m[int64(i)] = ts } currentDonTimes[id] = m } else { currentDonTimes[id] = observedDonTimes.TimestampsBySequence } + total += len(currentDonTimes[id]) } t.store.replaceDonTimes(currentDonTimes) t.store.setLastObservedDonTime(outcome.Timestamp) - t.lggr.Infow("Transmitting timestamps", "lastObservedDonTime", outcome.Timestamp) + t.lggr.Infow("Transmitting timestamps", "lastObservedDonTime", outcome.Timestamp, "executions", len(currentDonTimes), "total", total) for executionID, donTimes := range outcome.ObservedDonTimes { request := t.store.GetRequest(executionID) @@ -67,6 +69,9 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 ok = len(donTimes.Timestamps) > request.SeqNum if ok { donTime = donTimes.Timestamps[request.SeqNum] + if donTime == 0 { // feature flag was disabled, and we had a gap in the sequence + ok = false + } } } if ok { From 6e8ff75d10a12d9eec0b60b6981121553123db1f Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 7 Oct 2026 20:02:36 -0500 Subject: [PATCH 21/23] fix feature flag name --- pkg/settings/cresettings/README.md | 2 +- pkg/settings/cresettings/defaults.json | 2 +- pkg/settings/cresettings/defaults.toml | 2 +- pkg/settings/cresettings/settings.go | 4 ++-- pkg/workflows/dontime/factory.go | 4 ++-- pkg/workflows/dontime/plugin.go | 2 +- 6 files changed, 8 insertions(+), 8 deletions(-) diff --git a/pkg/settings/cresettings/README.md b/pkg/settings/cresettings/README.md index bb04a5269e..1e1b554f0b 100644 --- a/pkg/settings/cresettings/README.md +++ b/pkg/settings/cresettings/README.md @@ -374,7 +374,7 @@ flowchart CentralTriggerQueue.Put %% TODO placating test for now since this flowchart no longer renders - DonTimeSequencedTimestampsEnabled + DonTimeSequencedTimestampsActivePeriod classDef bound stroke:#f00 classDef gate stroke:#0f0 diff --git a/pkg/settings/cresettings/defaults.json b/pkg/settings/cresettings/defaults.json index 539a0e1c3d..8fe7160a9e 100644 --- a/pkg/settings/cresettings/defaults.json +++ b/pkg/settings/cresettings/defaults.json @@ -59,7 +59,7 @@ "VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit": "31.45728mb", "VaultMaxPerOracleUnexpiredBlobCount": "1000", "MissingRequestRecoveryEnabled": "false", - "DonTimeSequencedTimestampsEnabled": "[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]", + "DonTimeSequencedTimestampsActivePeriod": "[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]", "ConfidentialCompute": { "GlobalRate": "1000rps:1000", "MaxRetries": "3", diff --git a/pkg/settings/cresettings/defaults.toml b/pkg/settings/cresettings/defaults.toml index 69f509302d..832b3df2ca 100644 --- a/pkg/settings/cresettings/defaults.toml +++ b/pkg/settings/cresettings/defaults.toml @@ -58,7 +58,7 @@ VaultMaxBlobPayloadSizeLimit = '25.6kb' VaultMaxPerOracleUnexpiredBlobCumulativePayloadSizeLimit = '31.45728mb' VaultMaxPerOracleUnexpiredBlobCount = '1000' MissingRequestRecoveryEnabled = 'false' -DonTimeSequencedTimestampsEnabled = '[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]' +DonTimeSequencedTimestampsActivePeriod = '[2100-01-01 00:00:00 +0000 UTC,2101-01-01 00:00:00 +0000 UTC]' [ConfidentialCompute] GlobalRate = '1000rps:1000' diff --git a/pkg/settings/cresettings/settings.go b/pkg/settings/cresettings/settings.go index 12a2b33995..4477d4ffa1 100644 --- a/pkg/settings/cresettings/settings.go +++ b/pkg/settings/cresettings/settings.go @@ -166,7 +166,7 @@ var Default = Schema{ // MissingRequestRecoveryEnabled MissingRequestRecoveryEnabled: Bool(false), - DonTimeSequencedTimestampsEnabled: disabledFeatureTimeRange, + DonTimeSequencedTimestampsActivePeriod: disabledFeatureTimeRange, // Confidential Compute (San Marino framework) node-level settings. Defaults // mirror the previous hardcoded executor defaults so behavior is unchanged @@ -474,7 +474,7 @@ type Schema struct { MissingRequestRecoveryEnabled Setting[bool] - DonTimeSequencedTimestampsEnabled Setting[Range[config.Timestamp]] + DonTimeSequencedTimestampsActivePeriod Setting[Range[config.Timestamp]] // Confidential Compute (San Marino framework) node-level settings. ConfidentialCompute confidentialCompute diff --git a/pkg/workflows/dontime/factory.go b/pkg/workflows/dontime/factory.go index a73b2adc60..93e973b846 100644 --- a/pkg/workflows/dontime/factory.go +++ b/pkg/workflows/dontime/factory.go @@ -40,12 +40,12 @@ func NewFactory(s *Store, lggr logger.Logger) (*Factory, error) { return &Factory{ store: s, lggr: logger.Named(lggr, "OCR3DonTimeFactory"), - sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsEnabled.DefaultValue), + sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsActivePeriod.DefaultValue), }, nil } func (o *Factory) InitLimits(lf limits.Factory) error { - sequencedTSEnabled, err := limits.MakeRangeLimiter[config.Timestamp](lf, cresettings.Default.DonTimeSequencedTimestampsEnabled) + sequencedTSEnabled, err := limits.MakeRangeLimiter[config.Timestamp](lf, cresettings.Default.DonTimeSequencedTimestampsActivePeriod) if err != nil { return err } diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 1118fe3711..fe5cc8af70 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -126,7 +126,7 @@ func NewPlugin(store *Store, config ocr3types.ReportingPluginConfig, offchainCfg batchSize: int(offchainCfg.MaxBatchSize), minTimeIncrease: offchainCfg.MinTimeIncrease / int64(time.Millisecond), metrics: metrics, - sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsEnabled.DefaultValue), + sequencedTSEnabled: limits.NewRangeLimiter(cresettings.Default.DonTimeSequencedTimestampsActivePeriod.DefaultValue), }, nil } From 0e15c320a3c146896270a3ff711e797ff59750ec Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Wed, 7 Oct 2026 21:22:36 -0500 Subject: [PATCH 22/23] fix EarliestTS --- pkg/workflows/dontime/pb/dontime.go | 11 ++++------- pkg/workflows/dontime/plugin.go | 2 +- 2 files changed, 5 insertions(+), 8 deletions(-) diff --git a/pkg/workflows/dontime/pb/dontime.go b/pkg/workflows/dontime/pb/dontime.go index 9fb2107ac0..87e2b4deaf 100644 --- a/pkg/workflows/dontime/pb/dontime.go +++ b/pkg/workflows/dontime/pb/dontime.go @@ -1,7 +1,5 @@ package pb -import "math" - // MaxSeqNum returns the max sequence number from TimestampsBySequence, or -1 if none exist. func (t *ObservedDonTimes) MaxSeqNum() int64 { var maxSeqNum int64 = -1 @@ -13,12 +11,11 @@ func (t *ObservedDonTimes) MaxSeqNum() int64 { return maxSeqNum } -// EarliestTS returns the earliest timestamp value from TimestampsBySequence, or math.MaxInt64 if none exist. -func (t *ObservedDonTimes) EarliestTS() (earliestTS int64) { - earliestTS = math.MaxInt64 +// EarliestTS returns the earliest timestamp value from TimestampsBySequence or nil if none exist. +func (t *ObservedDonTimes) EarliestTS() (earliestTS *int64) { for _, ts := range t.TimestampsBySequence { - if ts < earliestTS { - earliestTS = ts + if earliestTS == nil || ts < *earliestTS { + earliestTS = &ts } } return diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index fe5cc8af70..74403dde1c 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -408,7 +408,7 @@ func (p *Plugin) sequencedOutcome(aos []types.AttributedObservation, prevOutcome p.store.deleteExecutionID(id) continue } - if donTime >= observedTimes.EarliestTS()+p.offChainConfig.ExecutionRemovalTime.AsDuration().Milliseconds() { + if ts := observedTimes.EarliestTS(); ts != nil && donTime >= *ts+p.offChainConfig.ExecutionRemovalTime.AsDuration().Milliseconds() { delete(outcome.ObservedDonTimes, id) p.store.deleteExecutionID(id) } From 2a3e61bf549ed091bd86417cedc997807159a7a5 Mon Sep 17 00:00:00 2001 From: Jordan Krage Date: Thu, 8 Oct 2026 09:58:03 -0500 Subject: [PATCH 23/23] simplifications and naming --- pkg/workflows/dontime/transmitter.go | 38 ++++++++++------------------ 1 file changed, 14 insertions(+), 24 deletions(-) diff --git a/pkg/workflows/dontime/transmitter.go b/pkg/workflows/dontime/transmitter.go index 8f33c3d2bb..9622e76f2b 100644 --- a/pkg/workflows/dontime/transmitter.go +++ b/pkg/workflows/dontime/transmitter.go @@ -34,26 +34,29 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 return err } - var total int - currentDonTimes := make(map[string]map[int64]int64, len(outcome.ObservedDonTimes)) + var totalCount int + entriesCount := make(map[string]map[int64]int64, len(outcome.ObservedDonTimes)) for id, observedDonTimes := range outcome.ObservedDonTimes { if len(observedDonTimes.Timestamps) > 0 { m := make(map[int64]int64) - for i, ts := range observedDonTimes.Timestamps { - m[int64(i)] = ts + for i, donTime := range observedDonTimes.Timestamps { + if donTime == 0 { // feature flag was disabled, and we had a gap in the sequence + continue + } + m[int64(i)] = donTime } - currentDonTimes[id] = m + entriesCount[id] = m } else { - currentDonTimes[id] = observedDonTimes.TimestampsBySequence + entriesCount[id] = observedDonTimes.TimestampsBySequence } - total += len(currentDonTimes[id]) + totalCount += len(entriesCount[id]) } - t.store.replaceDonTimes(currentDonTimes) + t.store.replaceDonTimes(entriesCount) t.store.setLastObservedDonTime(outcome.Timestamp) - t.lggr.Infow("Transmitting timestamps", "lastObservedDonTime", outcome.Timestamp, "executions", len(currentDonTimes), "total", total) + t.lggr.Infow("Transmitting timestamps", "lastObservedDonTime", outcome.Timestamp, "donTimeEntries", len(entriesCount), "donTimeTotal", totalCount) - for executionID, donTimes := range outcome.ObservedDonTimes { + for executionID, donTimes := range entriesCount { request := t.store.GetRequest(executionID) if request == nil { continue @@ -61,20 +64,7 @@ func (t *Transmitter) Transmit(_ context.Context, _ types.ConfigDigest, _ uint64 // Nodes behind on multiple requests may wait one OCR round per request. // Caching future times locally could be added as an optimization. - var donTime int64 - var ok bool - if len(donTimes.TimestampsBySequence) > 0 { - donTime, ok = donTimes.TimestampsBySequence[int64(request.SeqNum)] - } else { - ok = len(donTimes.Timestamps) > request.SeqNum - if ok { - donTime = donTimes.Timestamps[request.SeqNum] - if donTime == 0 { // feature flag was disabled, and we had a gap in the sequence - ok = false - } - } - } - if ok { + if donTime, ok := donTimes[int64(request.SeqNum)]; ok { t.store.RemoveRequest(executionID) // Make space for next request before delivering request.SendResponse(Response{ WorkflowExecutionID: executionID,