Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 3 additions & 4 deletions chain_capabilities/evm/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
26 changes: 26 additions & 0 deletions chain_capabilities/evm/trigger/don_id_label_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
32 changes: 14 additions & 18 deletions chain_capabilities/evm/trigger/trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)))
}
Expand Down Expand Up @@ -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)))
}
27 changes: 13 additions & 14 deletions http_trigger/trigger/connector_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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)
}
Expand All @@ -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.
Comment thread
tarcisiozf marked this conversation as resolved.
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
Expand Down
26 changes: 26 additions & 0 deletions http_trigger/trigger/don_id_label_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
5 changes: 2 additions & 3 deletions http_trigger/trigger/trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading