diff --git a/chain_capabilities/evm/trigger/trigger.go b/chain_capabilities/evm/trigger/trigger.go index acd11c5df..d9b077b96 100644 --- a/chain_capabilities/evm/trigger/trigger.go +++ b/chain_capabilities/evm/trigger/trigger.go @@ -476,11 +476,11 @@ func (lts *LogTriggerService) sendLogsToWorkflows(ctx context.Context, telemetry workflowExecutionID, execIDErr := workflows.GenerateExecutionIDWithTriggerIndex(telemetryContext.WorkflowID, response.Id, triggerIndex) if execIDErr != nil { - lts.lggr.Errorw("failed to generate execution ID", "err", execIDErr, "isLegacyExecutionID", false, "triggerID", triggerID, "workflowID", telemetryContext.WorkflowID, "eventID", response.Id) + lts.lggr.Errorw("failed to generate execution ID", "err", execIDErr, "triggerID", triggerID, "workflowID", telemetryContext.WorkflowID, "eventID", response.Id) // continue with execution even if we can't generate ID workflowExecutionID = "" } - lts.lggr.Debugw("new log trigger event", "triggerEventID", response.Id, "triggerID", triggerID, "executionID", workflowExecutionID, "isLegacyExecutionID", false) + lts.lggr.Debugw("new log trigger event", "triggerEventID", response.Id, "triggerID", triggerID, "executionID", workflowExecutionID) displayWorkflowName := telemetryContext.DecodedWorkflowName if displayWorkflowName == "" { diff --git a/cron/trigger/trigger.go b/cron/trigger/trigger.go index 28a44b7db..de9091972 100644 --- a/cron/trigger/trigger.go +++ b/cron/trigger/trigger.go @@ -248,7 +248,7 @@ func (s *Service) RegisterTrigger(ctx context.Context, triggerID string, metadat // Send trigger event even if we can't generate execution ID. Here the ID is used only for observability. } - s.lggr.Debugw("sending trigger event", "executionID", workflowExecutionID, "isLegacyExecutionID", false, "triggerID", triggerID, "scheduledExecTimeUTC", scheduledExecutionTimeUTC.Format(time.RFC3339Nano), "actualExecTimeUTC", currentTimeUTC.Format(time.RFC3339Nano)) + s.lggr.Debugw("sending trigger event", "executionID", workflowExecutionID, "triggerID", triggerID, "scheduledExecTimeUTC", scheduledExecutionTimeUTC.Format(time.RFC3339Nano), "actualExecTimeUTC", currentTimeUTC.Format(time.RFC3339Nano)) nextExecutionTime, nextRunErr := job.NextRun() if nextRunErr != nil { diff --git a/cron/trigger/trigger_test.go b/cron/trigger/trigger_test.go index eb269da8f..04e786158 100644 --- a/cron/trigger/trigger_test.go +++ b/cron/trigger/trigger_test.go @@ -1215,11 +1215,8 @@ func TestCronTrigger_ExecutionIDWithTriggerIndex(t *testing.T) { for _, entry := range observedLogs.All() { if entry.Message == "sending trigger event" { for _, field := range entry.Context { - switch field.Key { - case "executionID": + if field.Key == "executionID" { execIDFromLog = field.String - case "isLegacyExecutionID": - isLegacyFromLog = field.Integer == 1 } } found = true diff --git a/http_trigger/README.md b/http_trigger/README.md index 4ff311744..77b8452a9 100644 --- a/http_trigger/README.md +++ b/http_trigger/README.md @@ -119,7 +119,7 @@ All responses follow the JSON-RPC 2.0 specification: #### Workflow Execution ID Generation Execution IDs are generated using: ```go -workflowExecutionID = EncodeExecutionID(workflowID, requestID) +workflowExecutionID = GenerateExecutionIDWithTriggerIndex(workflowID, requestID, 0) ``` This ensures unique execution IDs that can be traced back to their originating request and guarantees uniqueness per workflow. diff --git a/http_trigger/trigger/connector_handler.go b/http_trigger/trigger/connector_handler.go index 18ccb4900..1981b0961 100644 --- a/http_trigger/trigger/connector_handler.go +++ b/http_trigger/trigger/connector_handler.go @@ -13,14 +13,12 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/capabilities" "github.com/smartcontractkit/chainlink-common/pkg/capabilities/v2/triggers/http" - "github.com/smartcontractkit/chainlink-common/pkg/config" "github.com/smartcontractkit/chainlink-common/pkg/contexts" "github.com/smartcontractkit/chainlink-common/pkg/custmsg" jsonrpc "github.com/smartcontractkit/chainlink-common/pkg/jsonrpc2" "github.com/smartcontractkit/chainlink-common/pkg/logger" "github.com/smartcontractkit/chainlink-common/pkg/services" "github.com/smartcontractkit/chainlink-common/pkg/services/orgresolver" - "github.com/smartcontractkit/chainlink-common/pkg/settings/cresettings" "github.com/smartcontractkit/chainlink-common/pkg/settings/limits" "github.com/smartcontractkit/chainlink-common/pkg/types/core" gateway_common "github.com/smartcontractkit/chainlink-common/pkg/types/gateway" @@ -47,18 +45,13 @@ type connectorHandler struct { metrics *Metrics wg sync.WaitGroup stopChan services.StopChan - orgResolver orgresolver.OrgResolver // Optional org resolver for fetching organization IDs - multiTriggerFlag limits.RangeLimiter[config.Timestamp] + orgResolver orgresolver.OrgResolver } func NewConnectorHandler(lggr logger.Logger, gc core.GatewayConnector, config ServiceConfig, capabilityDonID uint32, workflowStore *workflowStore, gatewayMetadataPublisher GatewayMetadataPublisher, requestCache *requestCache, metrics *Metrics, orgResolver orgresolver.OrgResolver, limitsFactory limits.Factory, ) (*connectorHandler, error) { - multiTriggerFlag, err := limits.MakeRangeLimiter(limitsFactory, cresettings.Default.PerWorkflow.FeatureHTTPTriggerNewExecutionIDsActivePeriod) - if err != nil { - return nil, fmt.Errorf("failed to create multi-trigger execution ID flag: %w", err) - } return &connectorHandler{ lggr: logger.Named(lggr, HandlerName), gatewayConnector: gc, @@ -70,7 +63,6 @@ func NewConnectorHandler(lggr logger.Logger, gc core.GatewayConnector, config Se metrics: metrics, stopChan: make(chan struct{}), orgResolver: orgResolver, - multiTriggerFlag: multiTriggerFlag, }, nil } @@ -293,7 +285,7 @@ func (h *connectorHandler) processTrigger(ctx context.Context, gatewayID string, } } - workflowExecutionID, isLegacyExecutionID, err := h.generateWorkflowExecutionID( + workflowExecutionID, err := h.generateWorkflowExecutionID( ctx, workflowMetadata.WorkflowID, workflowMetadata.WorkflowOwner, @@ -350,7 +342,7 @@ func (h *connectorHandler) processTrigger(ctx context.Context, gatewayID string, labeler = labeler.With(events.KeyOrganizationID, orgID) } - l.Debugw("Triggering workflow", "isLegacyExecutionID", isLegacyExecutionID) + l.Debugw("Triggering workflow") input := []byte(triggerReq.Input) err = h.triggerWorkflow(ctx, workflowMetadata.WorkflowID, req.ID, gatewayID, input, triggerReq.Key) if err != nil { @@ -454,7 +446,7 @@ func (h *connectorHandler) generateWorkflowExecutionID( ctx context.Context, workflowID, workflowOwner, orgID, reqID, referenceID string, l logger.Logger, -) (string, bool, error) { +) (string, error) { l = logger.With(l, "referenceID", referenceID, "workflowOwner", workflowOwner, "orgID", orgID) triggerIndex, err := workflows.GetTriggerIndexFromReferenceID(referenceID) @@ -472,27 +464,12 @@ func (h *connectorHandler) generateWorkflowExecutionID( }) strippedWorkflowID := strings.TrimPrefix(workflowID, "0x") - var workflowExecutionID string - var execIDErr error - isLegacyExecutionID := true - // NOTE: Relying on local time is not ideal but we don't have access to DONTime at this stage. - checkErr := h.multiTriggerFlag.Check(ctx, config.NewTimestamp(time.Now())) - if checkErr == nil { - workflowExecutionID, execIDErr = workflows.GenerateExecutionIDWithTriggerIndex(strippedWorkflowID, reqID, triggerIndex) - isLegacyExecutionID = false - } else { - if _, ok := errors.AsType[limits.ErrorRangeLimited[config.Timestamp]](checkErr); ok { - l.Debugw("Multi-trigger execution ID flag not active; using legacy execution ID", "error", checkErr) - } else { - l.Errorw("Multi-trigger execution ID flag check failed; using legacy execution ID", "error", checkErr) - } - workflowExecutionID, execIDErr = workflows.EncodeExecutionID(strippedWorkflowID, reqID) //nolint:staticcheck // SA1019 legacy execution ID path - } + workflowExecutionID, execIDErr := workflows.GenerateExecutionIDWithTriggerIndex(strippedWorkflowID, reqID, triggerIndex) if execIDErr != nil { - l.Errorw("Failed to generate workflow execution ID", "error", execIDErr, "isLegacyExecutionID", isLegacyExecutionID) - return "", isLegacyExecutionID, execIDErr + l.Errorw("Failed to generate workflow execution ID", "error", execIDErr) + return "", execIDErr } - return ensureHexPrefix(workflowExecutionID), isLegacyExecutionID, nil + return ensureHexPrefix(workflowExecutionID), nil } func (h *connectorHandler) handleRequestCaching(ctx context.Context, gatewayID string, req *jsonrpc.Request[json.RawMessage], workflowExecutionID string, l logger.Logger) bool { diff --git a/http_trigger/trigger/connector_handler_test.go b/http_trigger/trigger/connector_handler_test.go index 0531e82fb..c5c964f0a 100644 --- a/http_trigger/trigger/connector_handler_test.go +++ b/http_trigger/trigger/connector_handler_test.go @@ -269,7 +269,7 @@ func requireWorkflowTriggered(t *testing.T, triggerCh <-chan capabilities.Trigge require.NoError(t, err) require.Equal(t, testWorkflowID, triggerResp.WorkflowID) - executionID, err := workflows.EncodeExecutionID(strings.TrimPrefix(testWorkflowID, "0x"), req.ID) //nolint:staticcheck // SA1019 default FeatureMultiTrigger period is inactive + executionID, err := workflows.GenerateExecutionIDWithTriggerIndex(strings.TrimPrefix(testWorkflowID, "0x"), req.ID, 0) require.NoError(t, err) executionID = ensureHexPrefix(executionID) require.Equal(t, executionID, triggerResp.WorkflowExecutionID) @@ -974,7 +974,7 @@ func TestHandleGatewayMessage_TriggerFailureDoesNotCacheThenSameRequestSucceeds( t.Fatal("timed out waiting for trigger delivery") } - executionID, err := workflows.EncodeExecutionID(strings.TrimPrefix(testWorkflowID, "0x"), req.ID) //nolint:staticcheck // SA1019 default FeatureMultiTrigger period is inactive + executionID, err := workflows.GenerateExecutionIDWithTriggerIndex(strings.TrimPrefix(testWorkflowID, "0x"), req.ID, 0) require.NoError(t, err) executionID = ensureHexPrefix(executionID) require.Equal(t, executionID, triggerResp.WorkflowExecutionID) diff --git a/http_trigger/trigger/multi_trigger_flag_test.go b/http_trigger/trigger/multi_trigger_flag_test.go deleted file mode 100644 index 7dd763ab5..000000000 --- a/http_trigger/trigger/multi_trigger_flag_test.go +++ /dev/null @@ -1,153 +0,0 @@ -package trigger - -import ( - "strings" - "testing" - "time" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/smartcontractkit/chainlink-common/pkg/config" - "github.com/smartcontractkit/chainlink-common/pkg/contexts" - "github.com/smartcontractkit/chainlink-common/pkg/logger" - "github.com/smartcontractkit/chainlink-common/pkg/settings" - "github.com/smartcontractkit/chainlink-common/pkg/settings/cresettings" - "github.com/smartcontractkit/chainlink-common/pkg/settings/limits" - "github.com/smartcontractkit/chainlink-common/pkg/workflows" -) - -// TestFeatureMultiTriggerFlagCheckRequiresCRE shows the raw RangeLimiter contract: -// ScopeWorkflow Check fails without CRE on ctx, and succeeds with it. -func TestFeatureMultiTriggerFlagCheckRequiresCRE(t *testing.T) { - t.Parallel() - - now := time.Now().UTC() - period := cresettings.Default.PerWorkflow.FeatureHTTPTriggerNewExecutionIDsActivePeriod - period.DefaultValue = settings.Range[config.Timestamp]{ - Lower: config.Timestamp(now.Add(-time.Hour).Unix()), - Upper: config.Timestamp(now.Add(time.Hour).Unix()), - } - - multiTriggerFlag, err := limits.MakeRangeLimiter(limits.Factory{}, period) - require.NoError(t, err) - t.Cleanup(func() { require.NoError(t, multiTriggerFlag.Close()) }) - - strippedWorkflowID := strings.TrimPrefix(testWorkflowID, "0x") - reqID := "req-multi-trigger-ctx" - ts := config.NewTimestamp(now) - - legacyID, err := workflows.EncodeExecutionID(strippedWorkflowID, reqID) //nolint:staticcheck // intentional legacy path - require.NoError(t, err) - - multiID, err := workflows.GenerateExecutionIDWithTriggerIndex(strippedWorkflowID, reqID, 0) - require.NoError(t, err) - - require.NotEqual(t, legacyID, multiID, "legacy and multi-trigger hashes must differ even for trigger index 0") - - selectExecutionID := func(flagOn bool) string { - if flagOn { - return multiID - } - return legacyID - } - - // without_WithCRE_Check_fails_and_selects_legacy_ID: bare ctx → missing tenant → flag treated as off. - t.Run("without_WithCRE_Check_fails_and_selects_legacy_ID", func(t *testing.T) { - t.Parallel() - - checkErr := multiTriggerFlag.Check(t.Context(), ts) - require.Error(t, checkErr) - var errMissingTenant limits.ErrMissingTenant - require.ErrorAs(t, checkErr, &errMissingTenant) - assert.Equal(t, settings.ScopeWorkflow, errMissingTenant.Scope) - - flagOn := checkErr == nil - assert.False(t, flagOn) - assert.Equal(t, legacyID, selectExecutionID(flagOn)) - }) - - // with_WithCRE_Check_succeeds_and_selects_multi_trigger_ID: Workflow on ctx → Check passes. - t.Run("with_WithCRE_Check_succeeds_and_selects_multi_trigger_ID", func(t *testing.T) { - t.Parallel() - - ctx := contexts.WithCRE(t.Context(), contexts.CRE{Workflow: testWorkflowID}) - require.NoError(t, multiTriggerFlag.Check(ctx, ts)) - - flagOn := multiTriggerFlag.Check(ctx, ts) == nil - assert.True(t, flagOn) - assert.Equal(t, multiID, selectExecutionID(flagOn)) - }) -} - -// TestGenerateWorkflowExecutionID_WithCREUsesMultiTriggerIDs checks that -// generateWorkflowExecutionID itself attaches CRE, so a bare gateway ctx still -// yields multi-trigger IDs when the flag period is active. -func TestGenerateWorkflowExecutionID_WithCREUsesMultiTriggerIDs(t *testing.T) { - t.Parallel() - - lggr := logger.Test(t) - handler, _, _, _ := setup(t, lggr) - - now := time.Now().UTC() - period := cresettings.Default.PerWorkflow.FeatureHTTPTriggerNewExecutionIDsActivePeriod - period.DefaultValue = settings.Range[config.Timestamp]{ - Lower: config.Timestamp(now.Add(-time.Hour).Unix()), - Upper: config.Timestamp(now.Add(time.Hour).Unix()), - } - activeFlag, err := limits.MakeRangeLimiter(limits.Factory{}, period) - require.NoError(t, err) - t.Cleanup(func() { require.NoError(t, activeFlag.Close()) }) - handler.multiTriggerFlag = activeFlag - - reqID := "req-generate-with-cre" - refID := "trigger_0" - wantMulti, err := workflows.GenerateExecutionIDWithTriggerIndex(strings.TrimPrefix(testWorkflowID, "0x"), reqID, 0) - require.NoError(t, err) - wantMulti = ensureHexPrefix(wantMulti) - - // Intentionally pass t.Context() without WithCRE — the method under test must add it. - got, isLegacy, err := handler.generateWorkflowExecutionID( - t.Context(), - testWorkflowID, - testWorkflowOwner, - "test-org", - reqID, - refID, - lggr, - ) - require.NoError(t, err) - assert.False(t, isLegacy) - assert.Equal(t, wantMulti, got) -} - -// TestGenerateWorkflowExecutionID_InactiveFlagUsesLegacyIDs checks that when the -// flag period does not include now, generateWorkflowExecutionID returns a legacy ID -// even though CRE is attached internally. -func TestGenerateWorkflowExecutionID_InactiveFlagUsesLegacyIDs(t *testing.T) { - t.Parallel() - - lggr := logger.Test(t) - handler, _, _, _ := setup(t, lggr) - // Default FeatureMultiTrigger period is in 2100 → inactive with CRE present. - - reqID := "req-generate-legacy" - refID := "trigger_0" - wantLegacy, err := workflows.EncodeExecutionID(strings.TrimPrefix(testWorkflowID, "0x"), reqID) //nolint:staticcheck // SA1019 - require.NoError(t, err) - wantLegacy = ensureHexPrefix(wantLegacy) - - // Intentionally pass t.Context() without WithCRE — the method under test must add it. - got, isLegacy, err := handler.generateWorkflowExecutionID( - t.Context(), - testWorkflowID, - testWorkflowOwner, - "test-org", - reqID, - refID, - lggr, - ) - require.NoError(t, err) - assert.True(t, isLegacy) - assert.Equal(t, wantLegacy, got) -} diff --git a/integration_tests/http/http_trigger_test.go b/integration_tests/http/http_trigger_test.go index 0ad1cf305..1611c87f7 100644 --- a/integration_tests/http/http_trigger_test.go +++ b/integration_tests/http/http_trigger_test.go @@ -339,7 +339,7 @@ func validateHTTPTriggerResponse(t *testing.T, body []byte, requestID string, ex workflowIDFromResponse := respBody.Result.WorkflowID require.Equal(t, expectedWorkflowID, workflowIDFromResponse) - executionID, err := workflows.EncodeExecutionID(strings.TrimPrefix(workflowIDFromResponse, "0x"), requestID) //nolint:staticcheck // SA1019 legacy execution ID path + executionID, err := workflows.GenerateExecutionIDWithTriggerIndex(strings.TrimPrefix(workflowIDFromResponse, "0x"), requestID, 0) require.NoError(t, err) require.Equal(t, "0x"+executionID, respBody.Result.WorkflowExecutionID) }