From 09fe7aa2936aa862c71fdbd339a8dbe1e84e4c53 Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Wed, 7 Oct 2026 10:40:03 -0300 Subject: [PATCH 1/2] drop wf DON ID fallback for trigger events --- chain_capabilities/evm/main.go | 7 ++-- .../evm/trigger/don_id_label_test.go | 26 +++++++++++++++ chain_capabilities/evm/trigger/trigger.go | 32 ++++++++----------- http_trigger/trigger/connector_handler.go | 27 ++++++++-------- http_trigger/trigger/don_id_label_test.go | 26 +++++++++++++++ 5 files changed, 82 insertions(+), 36 deletions(-) create mode 100644 chain_capabilities/evm/trigger/don_id_label_test.go create mode 100644 http_trigger/trigger/don_id_label_test.go diff --git a/chain_capabilities/evm/main.go b/chain_capabilities/evm/main.go index 868f0a4be..7c3638d0c 100644 --- a/chain_capabilities/evm/main.go +++ b/chain_capabilities/evm/main.go @@ -137,10 +137,9 @@ func (c *capabilityGRPCService) Initialise(ctx context.Context, dependencies cor // - job-spec boot path: populated when unambiguous, otherwise 0 (e.g. a node // that belongs to multiple DONs running this capability, or a core node // that pre-dates CRE-4409). - // When it is 0 the trigger service falls back to the consumer workflow's DON - // ID (see trigger.NewLogTriggerService). We deliberately do NOT re-resolve it - // from the registry here: that lookup cannot disambiguate multi-DON nodes and - // would emit a guess instead of the safe workflow-DON fallback. See CRE-4409. + // When it is 0 trigger events carry no DON ID. We deliberately do NOT + // re-resolve it from the registry here: that lookup cannot disambiguate + // multi-DON nodes and could emit an incorrect guess. capabilityDonID := dependencies.CapabilityDonID derivedUnknownTTL := capcommon.MaxRequestTimeoutWithMultiplier(ctx, dependencies.CapabilityRegistry, c.id, capabilityDonID, cfg.UnknownRequestsTTL, c.lggr) c.consensusHandler = chainconsensus.NewHandler(c.lggr, c.requestPoller, consensusMetrics, derivedUnknownTTL, cfg.MaxUnknownRequestsCacheSize) diff --git a/chain_capabilities/evm/trigger/don_id_label_test.go b/chain_capabilities/evm/trigger/don_id_label_test.go new file mode 100644 index 000000000..d2f5453e4 --- /dev/null +++ b/chain_capabilities/evm/trigger/don_id_label_test.go @@ -0,0 +1,26 @@ +package trigger + +import ( + "testing" + + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-common/pkg/custmsg" + "github.com/smartcontractkit/chainlink-common/pkg/workflows/events" +) + +func TestWithCapabilityDonID(t *testing.T) { + t.Parallel() + + t.Run("labels events with the capability DON ID", func(t *testing.T) { + t.Parallel() + labels := withCapabilityDonID(custmsg.NewLabeler(), 2).Labels() + require.Equal(t, "2", labels[events.KeyDonID]) + }) + + t.Run("leaves DON ID unset when unknown", func(t *testing.T) { + t.Parallel() + labels := withCapabilityDonID(custmsg.NewLabeler(), 0).Labels() + require.NotContains(t, labels, events.KeyDonID) + }) +} diff --git a/chain_capabilities/evm/trigger/trigger.go b/chain_capabilities/evm/trigger/trigger.go index acd11c5df..00566e139 100644 --- a/chain_capabilities/evm/trigger/trigger.go +++ b/chain_capabilities/evm/trigger/trigger.go @@ -56,10 +56,8 @@ type LogTriggerService struct { beholderProcessor beholder.ProtoProcessor messageBuilder *monitoring.MessageBuilder - // capabilityDonID is the on-chain DON ID of this capability DON. - // Used to label emitted events with the sending DON ID, distinct from the - // consumer workflow's DON ID carried in RequestMetadata.WorkflowDonID. Zero - // means unknown; the labeler then falls back to WorkflowDonID. + // capabilityDonID is the on-chain DON ID of this capability DON, used to label + // emitted events with the sending DON ID. Zero means unknown. capabilityDonID uint32 triggers LogTriggerStore @@ -494,20 +492,7 @@ func (lts *LogTriggerService) sendLogsToWorkflows(ctx context.Context, telemetry events.KeyWorkflowName, displayWorkflowName, ) - // Emit the *sending* capability DON ID. The trigger plugin runs on a capability - // DON (e.g. chain_capabilities_zone-a), separate from the consumer workflow's - // DON carried in RequestMetadata.WorkflowDonID. The workflow service needs the - // sender's DON to resolve on-chain quorum params (N, F). See CRE-4409. - // capabilityDonID is 0 when the host could not resolve it authoritatively - // (a multi-DON job-spec node, or a core node that pre-dates CRE-4409); in - // that case we fall back to WorkflowDonID. This fallback is permanent, not - // transitional, since the job-spec boot path is still supported. - switch { - case lts.capabilityDonID != 0: - labeler = labeler.With(events.KeyDonID, strconv.Itoa(int(lts.capabilityDonID))) - case telemetryContext.WorkflowDonID != 0: - labeler = labeler.With(events.KeyDonID, strconv.Itoa(int(telemetryContext.WorkflowDonID))) - } + labeler = withCapabilityDonID(labeler, lts.capabilityDonID) if telemetryContext.WorkflowDonConfigVersion != 0 { labeler = labeler.With(events.KeyDonVersion, strconv.Itoa(int(telemetryContext.WorkflowDonConfigVersion))) } @@ -781,3 +766,14 @@ func (r realTicker) Stop() { } var defaultTickerFactory tickerFactory = realTickerFactory{} + +// withCapabilityDonID labels trigger events with the *sending* capability DON ID, +// which the workflow service uses to resolve on-chain quorum params (N, F). The +// label is left unset when the DON ID is unknown, otherwise the consumer workflow's +// DON ID would be wrong for a capability DON. +func withCapabilityDonID(labeler custmsg.MessageEmitter, capabilityDonID uint32) custmsg.MessageEmitter { + if capabilityDonID == 0 { + return labeler + } + return labeler.With(events.KeyDonID, strconv.Itoa(int(capabilityDonID))) +} diff --git a/http_trigger/trigger/connector_handler.go b/http_trigger/trigger/connector_handler.go index 18ccb4900..a75dee87b 100644 --- a/http_trigger/trigger/connector_handler.go +++ b/http_trigger/trigger/connector_handler.go @@ -40,7 +40,7 @@ type connectorHandler struct { lggr logger.Logger gatewayConnector core.GatewayConnector config ServiceConfig - capabilityDonID uint32 // authoritative sending DON ID; 0 = unknown, falls back to WorkflowDONID + capabilityDonID uint32 // authoritative sending DON ID; 0 = unknown requestCache *requestCache workflowStore *workflowStore gatewayMetadataPublisher GatewayMetadataPublisher @@ -323,18 +323,6 @@ func (h *connectorHandler) processTrigger(ctx context.Context, gatewayID string, displayWorkflowName = workflowMetadata.WorkflowName } - // Emit the *sending* capability DON ID. The HTTP trigger plugin runs on a - // capability DON, separate from the consumer workflow's DON. The workflow - // service needs the sender's DON to resolve on-chain quorum params (N, F). - // See CRE-4409. capabilityDonID is 0 when the host could not resolve it - // authoritatively (a multi-DON job-spec node, or a core node that pre-dates - // CRE-4409); in that case we fall back to WorkflowDONID. This fallback is - // permanent, not transitional, since the job-spec boot path is still supported. - donIDForEvent := h.capabilityDonID - if donIDForEvent == 0 { - donIDForEvent = workflowMetadata.WorkflowDONID - } - labeler := custmsg.NewLabeler().With( events.KeyTriggerID, req.ID, events.KeyWorkflowID, workflowMetadata.WorkflowID, @@ -344,8 +332,8 @@ func (h *connectorHandler) processTrigger(ctx context.Context, gatewayID string, events.KeyWorkflowRegistryChainSelector, workflowMetadata.WorkflowRegistryChainSelector, events.KeyWorkflowRegistryAddress, workflowMetadata.WorkflowRegistryAddress, events.KeyEngineVersion, workflowMetadata.EngineVersion, - events.KeyDonID, strconv.Itoa(int(donIDForEvent)), ) + labeler = withCapabilityDonID(labeler, h.capabilityDonID) if orgID != "" { labeler = labeler.With(events.KeyOrganizationID, orgID) } @@ -372,6 +360,17 @@ func (h *connectorHandler) processTrigger(ctx context.Context, gatewayID string, l.Debug("Trigger event processed") } +// withCapabilityDonID labels trigger events with the *sending* capability DON ID, +// which the workflow service uses to resolve on-chain quorum params (N, F). The +// label is left unset when the DON ID is unknown, otherwise the consumer workflow's +// DON ID would be wrong for a capability DON. +func withCapabilityDonID(labeler custmsg.MessageEmitter, capabilityDonID uint32) custmsg.MessageEmitter { + if capabilityDonID == 0 { + return labeler + } + return labeler.With(events.KeyDonID, strconv.Itoa(int(capabilityDonID))) +} + type WorkflowMetadata struct { WorkflowID string WorkflowOwner string diff --git a/http_trigger/trigger/don_id_label_test.go b/http_trigger/trigger/don_id_label_test.go new file mode 100644 index 000000000..d2f5453e4 --- /dev/null +++ b/http_trigger/trigger/don_id_label_test.go @@ -0,0 +1,26 @@ +package trigger + +import ( + "testing" + + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-common/pkg/custmsg" + "github.com/smartcontractkit/chainlink-common/pkg/workflows/events" +) + +func TestWithCapabilityDonID(t *testing.T) { + t.Parallel() + + t.Run("labels events with the capability DON ID", func(t *testing.T) { + t.Parallel() + labels := withCapabilityDonID(custmsg.NewLabeler(), 2).Labels() + require.Equal(t, "2", labels[events.KeyDonID]) + }) + + t.Run("leaves DON ID unset when unknown", func(t *testing.T) { + t.Parallel() + labels := withCapabilityDonID(custmsg.NewLabeler(), 0).Labels() + require.NotContains(t, labels, events.KeyDonID) + }) +} From 49da985c6cf211dcd2b4028b2407528ded5269a0 Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Wed, 7 Oct 2026 15:36:28 -0300 Subject: [PATCH 2/2] update comment --- http_trigger/trigger/trigger.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/http_trigger/trigger/trigger.go b/http_trigger/trigger/trigger.go index abb164be9..c8e70357e 100644 --- a/http_trigger/trigger/trigger.go +++ b/http_trigger/trigger/trigger.go @@ -87,9 +87,8 @@ func (s *service) Initialise(ctx context.Context, dependencies core.StandardCapa requestCache := newRequestCache(s.lggr, dependencies.Store, time.Duration(s.cfg.RequestCacheTTL)*time.Second) // dependencies.CapabilityDonID is the on-chain DON ID this plugin process // serves, used to label emitted events with the *sending* DON. Zero means the - // host could not resolve it authoritatively (a multi-DON job-spec node, or a - // core node that pre-dates CRE-4409); the handler then falls back to - // RequestMetadata.WorkflowDONID. See CRE-4409. + // host could not resolve it (e.g. a multi-DON job-spec node); events are then + // emitted without a DON ID label. s.connectorHandler, err = NewConnectorHandler(s.lggr, dependencies.GatewayConnector, s.cfg, dependencies.CapabilityDonID, workflowStore, metadataPublisher, requestCache, s.metrics, s.orgResolver, s.limitsFactory) if err != nil { return err