From d2d084a11f0b316dce803d5025babf869ffe9429 Mon Sep 17 00:00:00 2001 From: Bolek Kulbabinski <1416262+bolekk@users.noreply.github.com> Date: Sun, 20 Sep 2026 09:12:19 -0700 Subject: [PATCH 1/2] Move cron trigger org resolution and event emission to last step (CRE-5715) Org ID resolution and TriggerExecutionStarted emission can be slow (org resolver lookup, event bus). Run them after the trigger response is written back and sent, so slow bookkeeping never delays trigger generation or the store write-back. Co-Authored-By: Claude Sonnet 5 --- cron/trigger/trigger.go | 134 +++++++++++++++++++++------------------- 1 file changed, 71 insertions(+), 63 deletions(-) diff --git a/cron/trigger/trigger.go b/cron/trigger/trigger.go index 6a457a062..91df416f2 100644 --- a/cron/trigger/trigger.go +++ b/cron/trigger/trigger.go @@ -249,50 +249,12 @@ func (s *Service) RegisterTrigger(ctx context.Context, triggerID string, metadat } workflowExecutionID, execIDErr := workflows.GenerateExecutionIDWithTriggerIndex(trigger.workflowID, response.Id, triggerIndex) - if execIDErr != nil { s.lggr.Errorw("failed to generate execution ID", "err", execIDErr, "triggerID", triggerID, "workflowID", trigger.workflowID, "triggerEventID", response.Id) - // Continue with execution even if we can't generate ID or emit event - } else { - // Try to fetch organization ID if org resolver is available - var orgID string - if s.orgResolver != nil && metadata.WorkflowOwner != "" { - func() { - defer func() { - if r := recover(); r != nil { - s.lggr.Warnw("Panic while fetching organization ID from org resolver", "workflowOwner", metadata.WorkflowOwner, "panic", r) - } - }() - if fetchedOrgID, orgErr := s.orgResolver.Get(ctx, metadata.WorkflowOwner); orgErr != nil { - s.lggr.Warnw("Failed to fetch organization ID from org resolver", "workflowOwner", metadata.WorkflowOwner, "error", orgErr) - } else if fetchedOrgID != "" { - orgID = fetchedOrgID - s.lggr.Debugw("Successfully fetched organization ID", "workflowOwner", metadata.WorkflowOwner, "orgID", orgID) - } - }() - } - - // Emit TriggerExecutionStarted event - labeler := custmsg.NewLabeler().With( - events.KeyTriggerID, response.Id, - events.KeyWorkflowID, trigger.workflowID, - events.KeyWorkflowExecutionID, workflowExecutionID, - events.KeyWorkflowOwner, metadata.WorkflowOwner, - events.KeyWorkflowName, displayWorkflowName, - events.KeyDonID, strconv.Itoa(int(metadata.WorkflowDonID)), - events.KeyDonVersion, strconv.Itoa(int(metadata.WorkflowDonConfigVersion)), - events.KeyOrganizationID, orgID, - events.KeyWorkflowRegistryChainSelector, metadata.WorkflowRegistryChainSelector, - events.KeyWorkflowRegistryAddress, metadata.WorkflowRegistryAddress, - events.KeyEngineVersion, metadata.EngineVersion, - ) - if emitErr := events.EmitTriggerExecutionStarted(ctx, labeler); emitErr != nil { - s.lggr.Errorw("failed to emit trigger execution started event", "err", emitErr, "triggerID", triggerID, "workflowExecutionID", workflowExecutionID) - // Continue with execution even if event emission fails - } + // Send trigger event even if we can't generate execution ID. Here the ID is used only for observability. } - s.lggr.Debugw("task callback sending trigger response", "executionID", workflowExecutionID, "isLegacyExecutionID", false, "triggerID", triggerID, "scheduledExecTimeUTC", scheduledExecutionTimeUTC.Format(time.RFC3339Nano), "actualExecTimeUTC", currentTimeUTC.Format(time.RFC3339Nano)) + s.lggr.Debugw("sending trigger event", "executionID", workflowExecutionID, "isLegacyExecutionID", false, "triggerID", triggerID, "scheduledExecTimeUTC", scheduledExecutionTimeUTC.Format(time.RFC3339Nano), "actualExecTimeUTC", currentTimeUTC.Format(time.RFC3339Nano)) nextExecutionTime, nextRunErr := job.NextRun() if nextRunErr != nil { @@ -301,31 +263,77 @@ func (s *Service) RegisterTrigger(ctx context.Context, triggerID string, metadat s.lggr.Errorw("task callback failed to schedule next run", "executionID", workflowExecutionID, "triggerID", triggerID) } - muCh.RLock() - defer muCh.RUnlock() - if callbackCh == nil { - return // unregistered already + func() { + muCh.RLock() + defer muCh.RUnlock() + if callbackCh == nil { + return // unregistered already + } + s.triggers.Write(triggerID, cronTrigger{ + job: job, + nextRun: nextExecutionTime, + workflowID: metadata.WorkflowID, + close: closeCh, + }) + + select { + case callbackCh <- response: + default: + s.lggr.Errorw("callback channel full, dropping event", "executionID", workflowExecutionID, "triggerID", triggerID, "eventID", response.Id) + + lblErr := s.labeler.With( + "workflowOwner", metadata.WorkflowOwner, + "workflowName", displayWorkflowName, + "workflowID", metadata.WorkflowID, + ).Emit(ctx, "callback channel full, dropping event") + if lblErr != nil { + s.lggr.Errorw("cannot emit custom event", "executionID", workflowExecutionID, "triggerID", triggerID, "eventID", response.Id, "err", lblErr) + } + } + }() + + if execIDErr != nil { + return // don't emit anything if we couldn't generate execution ID } - s.triggers.Write(triggerID, cronTrigger{ - job: job, - nextRun: nextExecutionTime, - workflowID: metadata.WorkflowID, - close: closeCh, - }) - select { - case callbackCh <- response: - default: - s.lggr.Errorw("callback channel full, dropping event", "executionID", workflowExecutionID, "triggerID", triggerID, "eventID", response.Id) - - lblErr := s.labeler.With( - "workflowOwner", metadata.WorkflowOwner, - "workflowName", displayWorkflowName, - "workflowID", metadata.WorkflowID, - ).Emit(ctx, "callback channel full, dropping event") - if lblErr != nil { - s.lggr.Errorw("cannot emit custom event", "executionID", workflowExecutionID, "triggerID", triggerID, "eventID", response.Id, "err", lblErr) - } + // Org resolution and event emission are done last, after trigger generation and + // bookkeeping above, since they can involve slow network calls (org resolver lookup, + // event bus) that must not delay the actual trigger or the next scheduled run. + + // Try to fetch organization ID if org resolver is available + var orgID string + if s.orgResolver != nil && metadata.WorkflowOwner != "" { + func() { + defer func() { + if r := recover(); r != nil { + s.lggr.Warnw("Panic while fetching organization ID from org resolver", "workflowOwner", metadata.WorkflowOwner, "panic", r) + } + }() + if fetchedOrgID, orgErr := s.orgResolver.Get(ctx, metadata.WorkflowOwner); orgErr != nil { + s.lggr.Warnw("Failed to fetch organization ID from org resolver", "workflowOwner", metadata.WorkflowOwner, "error", orgErr) + } else if fetchedOrgID != "" { + orgID = fetchedOrgID + s.lggr.Debugw("Successfully fetched organization ID", "workflowOwner", metadata.WorkflowOwner, "orgID", orgID) + } + }() + } + + // Emit TriggerExecutionStarted event + labeler := custmsg.NewLabeler().With( + events.KeyTriggerID, response.Id, + events.KeyWorkflowID, trigger.workflowID, + events.KeyWorkflowExecutionID, workflowExecutionID, + events.KeyWorkflowOwner, metadata.WorkflowOwner, + events.KeyWorkflowName, displayWorkflowName, + events.KeyDonID, strconv.Itoa(int(metadata.WorkflowDonID)), + events.KeyDonVersion, strconv.Itoa(int(metadata.WorkflowDonConfigVersion)), + events.KeyOrganizationID, orgID, + events.KeyWorkflowRegistryChainSelector, metadata.WorkflowRegistryChainSelector, + events.KeyWorkflowRegistryAddress, metadata.WorkflowRegistryAddress, + events.KeyEngineVersion, metadata.EngineVersion, + ) + if emitErr := events.EmitTriggerExecutionStarted(ctx, labeler); emitErr != nil { + s.lggr.Errorw("failed to emit trigger execution started event", "err", emitErr, "triggerID", triggerID, "workflowExecutionID", workflowExecutionID) } }) From 946143e88f39204ef5dd5f31c5ea61fc53e399b3 Mon Sep 17 00:00:00 2001 From: Bolek Kulbabinski <1416262+bolekk@users.noreply.github.com> Date: Sun, 20 Sep 2026 13:41:26 -0700 Subject: [PATCH 2/2] Fix TestCronTrigger_ExecutionIDWithTriggerIndex log message mismatch (CRE-5715) The task callback's debug log was renamed to "sending trigger event" but the test still checked for the old "task callback sending trigger response" text. Co-Authored-By: Claude Sonnet 5 --- cron/trigger/trigger_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/cron/trigger/trigger_test.go b/cron/trigger/trigger_test.go index 6c1b8c0db..8bd219cb4 100644 --- a/cron/trigger/trigger_test.go +++ b/cron/trigger/trigger_test.go @@ -1194,13 +1194,13 @@ func TestCronTrigger_ExecutionIDWithTriggerIndex(t *testing.T) { expectedExecID, err := workflows.GenerateExecutionIDWithTriggerIndex(testWorkflowID, msg.Id, testTriggerIndex) require.NoError(t, err) - // The debug log at "task callback sending trigger response" is written - // before the channel send, so it is already present once we receive msg. + // The debug log at "sending trigger event" is written before the channel + // send, so it is already present once we receive msg. var execIDFromLog string var isLegacyFromLog bool var found bool for _, entry := range observedLogs.All() { - if entry.Message == "task callback sending trigger response" { + if entry.Message == "sending trigger event" { for _, field := range entry.Context { switch field.Key { case "executionID": @@ -1213,7 +1213,7 @@ func TestCronTrigger_ExecutionIDWithTriggerIndex(t *testing.T) { break } } - require.True(t, found, "expected log entry 'task callback sending trigger response'") + require.True(t, found, "expected log entry 'sending trigger event'") require.Equal(t, expectedExecID, execIDFromLog, "execution ID should match expected hash function") require.False(t, isLegacyFromLog)