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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions chain_capabilities/evm/trigger/trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 == "" {
Expand Down
2 changes: 1 addition & 1 deletion cron/trigger/trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
5 changes: 1 addition & 4 deletions cron/trigger/trigger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion http_trigger/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
39 changes: 8 additions & 31 deletions http_trigger/trigger/connector_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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,
Expand All @@ -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
}

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand All @@ -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 {
Expand Down
4 changes: 2 additions & 2 deletions http_trigger/trigger/connector_handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
153 changes: 0 additions & 153 deletions http_trigger/trigger/multi_trigger_flag_test.go

This file was deleted.

2 changes: 1 addition & 1 deletion integration_tests/http/http_trigger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
Loading