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
134 changes: 71 additions & 63 deletions cron/trigger/trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,7 @@
return nil
}

func (s *Service) RegisterTrigger(ctx context.Context, triggerID string, metadata capabilities.RequestMetadata, input *crontypedapi.Config) (<-chan capabilities.TriggerAndId[*crontypedapi.Payload], caperrors.Error) {

Check warning on line 184 in cron/trigger/trigger.go

View check run for this annotation

CL-sonarqube-production / SonarQube Code Analysis

Refactor this method to reduce its Cognitive Complexity from 41 to the 30 allowed.

[S3776] Cognitive Complexity of functions should not be too high See more on https://sonarqube.main.prod.cldev.sh/project/issues?id=smartcontractkit_capabilities&pullRequest=784&issues=81c4e7c9-109b-461a-8eb4-098b167a8e8e&open=81c4e7c9-109b-461a-8eb4-098b167a8e8e
ctx = metadata.ContextWithCRE(ctx)
var muCh sync.RWMutex // extra synchronization to prevent the cron task from racing to send on the closed chan and re-register itself
// hold the lock until we call triggers.Write
Expand Down Expand Up @@ -249,50 +249,12 @@
}

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 {
Expand All @@ -301,31 +263,77 @@
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)
}
})

Expand Down
8 changes: 4 additions & 4 deletions cron/trigger/trigger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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":
Expand All @@ -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)

Expand Down
Loading