diff --git a/rest-api/flow/.gitignore b/rest-api/flow/.gitignore index 446b30493a..c2ea4f26ec 100644 --- a/rest-api/flow/.gitignore +++ b/rest-api/flow/.gitignore @@ -5,3 +5,7 @@ build flow .local_envrc ignore + +# Go event-rule target package; override the repository-wide Rust target ignore. +!internal/eventrule/target/ +!internal/eventrule/target/*.go diff --git a/rest-api/flow/internal/converter/dao/event_action_execution.go b/rest-api/flow/internal/converter/dao/event_action_execution.go index f88599fdbd..ea46f2d4a1 100644 --- a/rest-api/flow/internal/converter/dao/event_action_execution.go +++ b/rest-api/flow/internal/converter/dao/event_action_execution.go @@ -35,7 +35,7 @@ func EventActionExecutionTo( Observations: execution.Observations, Attempts: execution.Attempts, StatusMessage: execution.StatusMessage, - FirstClaimedAt: execution.FirstClaimedAt, + CreatedAt: execution.CreatedAt, UpdatedAt: execution.UpdatedAt, NextAttemptAt: nextAttemptAt, }, nil @@ -55,20 +55,24 @@ func EventActionExecutionFrom( } execution := &eventrule.Execution{ ExecutionState: eventrule.ExecutionState{ - Status: eventrule.ExecutionStatus(persisted.Status), - Reason: eventrule.ExecutionReason(persisted.Reason), - StatusMessage: persisted.StatusMessage, + ExecutionStatusDetails: eventrule.ExecutionStatusDetails{ + Status: eventrule.ExecutionStatus(persisted.Status), + Reason: eventrule.ExecutionReason(persisted.Reason), + StatusMessage: persisted.StatusMessage, + }, NextAttemptAt: nextAttemptAt, }, - ID: persisted.ID, - EventID: persisted.EventID, - RuleID: persisted.RuleID, - ActionID: persisted.ActionID, - CorrelationKey: persisted.CorrelationKey, - Observations: persisted.Observations, - Attempts: persisted.Attempts, - FirstClaimedAt: persisted.FirstClaimedAt, - UpdatedAt: persisted.UpdatedAt, + ExecutionIdentity: eventrule.ExecutionIdentity{ + EventID: persisted.EventID, + RuleID: persisted.RuleID, + ActionID: persisted.ActionID, + CorrelationKey: persisted.CorrelationKey, + }, + ID: persisted.ID, + Observations: persisted.Observations, + Attempts: persisted.Attempts, + CreatedAt: persisted.CreatedAt, + UpdatedAt: persisted.UpdatedAt, } if err := execution.Validate(); err != nil { return nil, fmt.Errorf("%w: %w", eventrule.ErrInvalidPersistedExecution, err) diff --git a/rest-api/flow/internal/converter/dao/event_action_execution_test.go b/rest-api/flow/internal/converter/dao/event_action_execution_test.go index af1820055f..2182ce05ea 100644 --- a/rest-api/flow/internal/converter/dao/event_action_execution_test.go +++ b/rest-api/flow/internal/converter/dao/event_action_execution_test.go @@ -16,13 +16,21 @@ import ( func TestEventActionExecutionRoundTrip(t *testing.T) { now := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) base := eventrule.Execution{ - ExecutionState: eventrule.ExecutionState{Status: eventrule.ExecutionStatusClaimed}, - ID: uuid.New(), EventID: uuid.New(), RuleID: uuid.New(), ActionID: "notify", - CorrelationKey: "incident-1", Observations: 2, Attempts: 1, - FirstClaimedAt: now, UpdatedAt: now.Add(time.Second), + ExecutionState: eventrule.ExecutionState{ExecutionStatusDetails: eventrule.ExecutionStatusDetails{Status: eventrule.ExecutionStatusPending}}, + ExecutionIdentity: eventrule.ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "notify", + CorrelationKey: "incident-1", + }, + ID: uuid.New(), + Observations: 2, + Attempts: 1, + CreatedAt: now, + UpdatedAt: now.Add(time.Second), } tests := map[string]eventrule.Execution{ - "claimed": executionWithStatus(base, eventrule.ExecutionStatusClaimed), + "pending": executionWithStatus(base, eventrule.ExecutionStatusPending), "submitted": executionWithStatus(base, eventrule.ExecutionStatusSubmitted), "completed": executionWithStatus(base, eventrule.ExecutionStatusCompleted), "skipped": func() eventrule.Execution { @@ -60,7 +68,7 @@ func TestEventActionExecutionRoundTrip(t *testing.T) { func TestEventActionExecutionToRejectsInvalidDomain(t *testing.T) { tests := map[string]*eventrule.Execution{ "nil": nil, - "invalid id": {ExecutionState: eventrule.ExecutionState{Status: eventrule.ExecutionStatusClaimed}}, + "invalid id": {ExecutionState: eventrule.ExecutionState{ExecutionStatusDetails: eventrule.ExecutionStatusDetails{Status: eventrule.ExecutionStatusPending}}}, } for name, execution := range tests { t.Run(name, func(t *testing.T) { @@ -74,10 +82,15 @@ func TestEventActionExecutionToRejectsInvalidDomain(t *testing.T) { func TestEventActionExecutionFrom(t *testing.T) { now := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) valid, err := EventActionExecutionTo(&eventrule.Execution{ - ExecutionState: eventrule.ExecutionState{Status: eventrule.ExecutionStatusClaimed}, - ID: uuid.New(), EventID: uuid.New(), RuleID: uuid.New(), ActionID: "notify", + ExecutionState: eventrule.ExecutionState{ExecutionStatusDetails: eventrule.ExecutionStatusDetails{Status: eventrule.ExecutionStatusPending}}, + ExecutionIdentity: eventrule.ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "notify", + }, + ID: uuid.New(), Observations: 1, Attempts: 1, - FirstClaimedAt: now, UpdatedAt: now, + CreatedAt: now, UpdatedAt: now, }) require.NoError(t, err) diff --git a/rest-api/flow/internal/db/model/event_action_execution.go b/rest-api/flow/internal/db/model/event_action_execution.go index 272fc7623c..9ed60ac2b9 100644 --- a/rest-api/flow/internal/db/model/event_action_execution.go +++ b/rest-api/flow/internal/db/model/event_action_execution.go @@ -10,7 +10,8 @@ import ( "github.com/uptrace/bun" ) -// EventActionExecution is the bun model for the event_action_executions table. +// EventActionExecution is the prospective persistence model for an event-rule +// action execution. The database table is introduced in a later phase. type EventActionExecution struct { bun.BaseModel `bun:"table:event_action_executions,alias:eae"` @@ -24,7 +25,7 @@ type EventActionExecution struct { Observations int `bun:"observations,notnull"` Attempts int `bun:"attempts,notnull"` StatusMessage string `bun:"status_message,notnull"` - FirstClaimedAt time.Time `bun:"first_claimed_at,notnull"` + CreatedAt time.Time `bun:"created_at,notnull"` UpdatedAt time.Time `bun:"updated_at,notnull"` NextAttemptAt *time.Time `bun:"next_attempt_at"` } diff --git a/rest-api/flow/internal/eventrule/action.go b/rest-api/flow/internal/eventrule/action.go index 5d41ce74ec..4dfc4584d0 100644 --- a/rest-api/flow/internal/eventrule/action.go +++ b/rest-api/flow/internal/eventrule/action.go @@ -62,9 +62,13 @@ func (c ActionCondition) AppliesTo(envelope Envelope, resource ResolvedResource) return false } - if c.ComponentTypes != nil && - !slices.Contains(c.ComponentTypes, resource.ComponentType) { - return false + if c.ComponentTypes != nil { + if resource.Kind != ResourceKindComponent { + return false + } + if !slices.Contains(c.ComponentTypes, resource.ComponentType) { + return false + } } return true @@ -98,9 +102,10 @@ func (s ConflictStrategy) validate() error { // ActionSpec is the closed set of typed responses supported by an action. // The unexported validation method prevents implementations outside this -// package while allowing processors to identify a specification with Type. +// package while exposing its type and target-resolution behavior. type ActionSpec interface { Type() ActionType + TargetResolutionStrategy() TargetStrategy validate() error } @@ -175,26 +180,36 @@ func CloneActions(actions []Action) []Action { return cloned } -// TargetStrategy identifies how a target-bearing action resolves concrete +// TargetStrategy identifies whether and how an action resolves concrete // operation targets. type TargetStrategy string const ( + TargetStrategyNone TargetStrategy = "none" TargetStrategyComponent TargetStrategy = "component" TargetStrategyRack TargetStrategy = "rack" TargetStrategyAffectedComponents TargetStrategy = "affected_components" ) -// Validate checks that the target strategy is supported by the schema. +// Validate checks that the target strategy is supported by the domain. func (s TargetStrategy) Validate() error { switch s { - case TargetStrategyComponent, TargetStrategyRack, TargetStrategyAffectedComponents: + case TargetStrategyNone, + TargetStrategyComponent, + TargetStrategyRack, + TargetStrategyAffectedComponents: return nil default: return fmt.Errorf("unknown target strategy %q", s) } } +// RequiresResolution reports whether concrete targets must be resolved for +// the strategy. +func (s TargetStrategy) RequiresResolution() bool { + return s != TargetStrategyNone +} + // SubmitTask describes a task submission requested by an event rule. type SubmitTask struct { OperationType taskcommon.TaskType @@ -209,6 +224,11 @@ func (s SubmitTask) Type() ActionType { return ActionTypeSubmitTask } +// TargetResolutionStrategy returns the task's target strategy. +func (s SubmitTask) TargetResolutionStrategy() TargetStrategy { + return s.TargetStrategy +} + func (s SubmitTask) validate() error { if !s.OperationType.IsValid() { return fmt.Errorf("operation_type %q is invalid", s.OperationType) @@ -221,6 +241,9 @@ func (s SubmitTask) validate() error { if err := s.TargetStrategy.Validate(); err != nil { return err } + if !s.TargetStrategy.RequiresResolution() { + return fmt.Errorf("submit task target strategy must require resolution") + } if err := s.ConflictStrategy.validate(); err != nil { return err @@ -240,6 +263,11 @@ func (s SendAlert) Type() ActionType { return ActionTypeSendAlert } +// TargetResolutionStrategy reports that alerts do not resolve targets. +func (SendAlert) TargetResolutionStrategy() TargetStrategy { + return TargetStrategyNone +} + func (s SendAlert) validate() error { if err := s.Severity.Validate(); err != nil { return err @@ -261,6 +289,11 @@ func (Noop) Type() ActionType { return ActionTypeNoop } +// TargetResolutionStrategy reports that no-op actions do not resolve targets. +func (Noop) TargetResolutionStrategy() TargetStrategy { + return TargetStrategyNone +} + func (n Noop) validate() error { return validateOptionalString("noop reason", n.Reason) } diff --git a/rest-api/flow/internal/eventrule/action_test.go b/rest-api/flow/internal/eventrule/action_test.go index 9356bad871..78f77fb7d1 100644 --- a/rest-api/flow/internal/eventrule/action_test.go +++ b/rest-api/flow/internal/eventrule/action_test.go @@ -4,6 +4,7 @@ package eventrule import ( + "slices" "testing" "time" @@ -50,9 +51,19 @@ func TestActionsValidate(t *testing.T) { }), NewAction("noop", ActionCondition{}, Noop{Reason: "record only"}), ) + wantStrategies := append( + slices.Clone(strategies), + TargetStrategyNone, + TargetStrategyNone, + ) for i := range actions { require.NoError(t, actions[i].Validate()) + require.Equal( + t, + wantStrategies[i], + actions[i].Spec.TargetResolutionStrategy(), + ) } } @@ -65,6 +76,8 @@ func TestActionRejectsInvalidDomainValues(t *testing.T) { } unknownStrategySpec := validTaskSpec unknownStrategySpec.TargetStrategy = "unknown" + noneStrategySpec := validTaskSpec + noneStrategySpec.TargetStrategy = TargetStrategyNone mismatchedOperationSpec := validTaskSpec mismatchedOperationSpec.OperationCode = taskcommon.OpCodeFirmwareControlUpgrade tests := map[string]Action{ @@ -80,6 +93,9 @@ func TestActionRejectsInvalidDomainValues(t *testing.T) { "unknown strategy": NewAction( "task", ActionCondition{}, unknownStrategySpec, ), + "task without target resolution": NewAction( + "task", ActionCondition{}, noneStrategySpec, + ), "mismatched operation": NewAction( "task", ActionCondition{}, mismatchedOperationSpec, ), @@ -93,6 +109,25 @@ func TestActionRejectsInvalidDomainValues(t *testing.T) { } } +func TestTargetStrategy_RequiresResolution(t *testing.T) { + tests := map[string]struct { + strategy TargetStrategy + want bool + }{ + "none": {strategy: TargetStrategyNone}, + "component": {strategy: TargetStrategyComponent, want: true}, + "rack": {strategy: TargetStrategyRack, want: true}, + "affected components": {strategy: TargetStrategyAffectedComponents, want: true}, + } + + for name, test := range tests { + t.Run(name, func(t *testing.T) { + require.NoError(t, test.strategy.Validate()) + require.Equal(t, test.want, test.strategy.RequiresResolution()) + }) + } +} + func TestRuleValidatesPolicy(t *testing.T) { action := NewAction("noop", ActionCondition{}, Noop{}) rule := Rule{ @@ -132,18 +167,35 @@ func TestActionConditionAppliesTo(t *testing.T) { "matches severity and component type": { condition: condition, envelope: Envelope{Severity: SeverityCritical}, - resource: ResolvedResource{ComponentType: flowtypes.ComponentTypeCompute}, - want: true, + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ComponentType: flowtypes.ComponentTypeCompute, + }, + want: true, }, "rejects severity": { condition: condition, envelope: Envelope{Severity: SeverityInfo}, - resource: ResolvedResource{ComponentType: flowtypes.ComponentTypeCompute}, + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ComponentType: flowtypes.ComponentTypeCompute, + }, }, "rejects component type": { condition: condition, envelope: Envelope{Severity: SeverityCritical}, - resource: ResolvedResource{ComponentType: flowtypes.ComponentTypeNVSwitch}, + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ComponentType: flowtypes.ComponentTypeNVSwitch, + }, + }, + "component type condition rejects rack": { + condition: condition, + envelope: Envelope{Severity: SeverityCritical}, + resource: ResolvedResource{ + Kind: ResourceKindRack, + ComponentType: flowtypes.ComponentTypeCompute, + }, }, "empty severity set matches nothing": { condition: ActionCondition{Severities: []Severity{}}, diff --git a/rest-api/flow/internal/eventrule/doc.go b/rest-api/flow/internal/eventrule/doc.go index 0ff6958905..e95311bc9f 100644 --- a/rest-api/flow/internal/eventrule/doc.go +++ b/rest-api/flow/internal/eventrule/doc.go @@ -77,9 +77,10 @@ // condition applies to every event. // // Task actions use a named TargetStrategy rather than an arbitrary inventory -// query. Target resolution and side effects occur outside this package. If a -// target strategy resolves no resources, the processor should record the -// action as skipped and must not submit a task. +// query; actions without targets use TargetStrategyNone. Target resolution and +// side effects occur outside this package. If a target strategy resolves no +// resources, the processor should record the action as skipped and must not +// submit a task. // // # Validation boundaries // diff --git a/rest-api/flow/internal/eventrule/event.go b/rest-api/flow/internal/eventrule/event.go index 4506f10584..96d5663058 100644 --- a/rest-api/flow/internal/eventrule/event.go +++ b/rest-api/flow/internal/eventrule/event.go @@ -156,3 +156,27 @@ type ResolvedResource struct { RackID uuid.UUID ComponentType flowtypes.ComponentType } + +// Validate checks the canonical identity and attributes established during +// enrichment. +func (r ResolvedResource) Validate() error { + if err := r.Kind.Validate(); err != nil { + return err + } + if r.ID == uuid.Nil { + return fmt.Errorf("resolved resource id is required") + } + if r.RackID == uuid.Nil { + return fmt.Errorf("resolved resource rack id is required") + } + if r.Kind == ResourceKindComponent { + if err := r.ComponentType.Validate(); err != nil { + return fmt.Errorf("resolved resource component type: %w", err) + } + } else { + if r.ID != r.RackID { + return fmt.Errorf("resolved rack resource id must equal rack id") + } + } + return nil +} diff --git a/rest-api/flow/internal/eventrule/event_test.go b/rest-api/flow/internal/eventrule/event_test.go index 652500abde..23aafe19a1 100644 --- a/rest-api/flow/internal/eventrule/event_test.go +++ b/rest-api/flow/internal/eventrule/event_test.go @@ -7,6 +7,7 @@ import ( "encoding/json" "testing" + flowtypes "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/types" "github.com/google/uuid" "github.com/stretchr/testify/require" ) @@ -89,3 +90,94 @@ func TestResourceIDMayBeUnresolved(t *testing.T) { resource.ID = uuid.New() require.NoError(t, resource.Validate()) } + +func TestResolvedResource_Validate(t *testing.T) { + componentID := uuid.New() + rackID := uuid.New() + tests := []struct { + name string + resource ResolvedResource + wantErr string + }{ + { + name: "component", + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ID: componentID, + RackID: rackID, + ComponentType: flowtypes.ComponentTypeCompute, + }, + }, + { + name: "rack", + resource: ResolvedResource{ + Kind: ResourceKindRack, + ID: rackID, + RackID: rackID, + }, + }, + { + name: "rack ignores component type", + resource: ResolvedResource{ + Kind: ResourceKindRack, + ID: rackID, + RackID: rackID, + ComponentType: flowtypes.ComponentTypeCompute, + }, + }, + { + name: "kind required", + resource: ResolvedResource{ID: componentID, RackID: rackID}, + wantErr: "unknown resource kind", + }, + { + name: "id required", + resource: ResolvedResource{Kind: ResourceKindComponent, RackID: rackID}, + wantErr: "resolved resource id is required", + }, + { + name: "rack id required", + resource: ResolvedResource{Kind: ResourceKindComponent, ID: componentID}, + wantErr: "resolved resource rack id is required", + }, + { + name: "component type required", + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ID: componentID, + RackID: rackID, + }, + wantErr: "resolved resource component type", + }, + { + name: "component type must be supported", + resource: ResolvedResource{ + Kind: ResourceKindComponent, + ID: componentID, + RackID: rackID, + ComponentType: flowtypes.ComponentType("INVALID"), + }, + wantErr: "unknown component type", + }, + { + name: "rack identities must match", + resource: ResolvedResource{ + Kind: ResourceKindRack, + ID: uuid.New(), + RackID: rackID, + }, + wantErr: "resolved rack resource id must equal rack id", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + err := test.resource.Validate() + if test.wantErr != "" { + require.ErrorContains(t, err, test.wantErr) + return + } + require.NoError(t, err) + }) + } +} diff --git a/rest-api/flow/internal/eventrule/execution.go b/rest-api/flow/internal/eventrule/execution.go index c304cbcc2d..701154b7b8 100644 --- a/rest-api/flow/internal/eventrule/execution.go +++ b/rest-api/flow/internal/eventrule/execution.go @@ -4,7 +4,6 @@ package eventrule import ( - "errors" "fmt" "slices" "time" @@ -12,15 +11,11 @@ import ( "github.com/google/uuid" ) -// ErrRetryScheduled indicates that an existing execution is not yet eligible -// to resume. -var ErrRetryScheduled = errors.New("event action retry is scheduled") - -// ExecutionStatus identifies one action execution state. +// ExecutionStatus identifies one execution state. type ExecutionStatus string const ( - ExecutionStatusClaimed ExecutionStatus = "claimed" + ExecutionStatusPending ExecutionStatus = "pending" ExecutionStatusSkipped ExecutionStatus = "skipped" ExecutionStatusDeferred ExecutionStatus = "deferred" ExecutionStatusSubmitted ExecutionStatus = "submitted" @@ -28,10 +23,11 @@ const ( ExecutionStatusFailed ExecutionStatus = "failed" ) -// CanTransitionTo reports whether an execution with this status may transition -// to the target status. +// CanTransitionTo reports whether an execution with this status may accept an +// attempt result. Pending is used by the creator's first attempt; deferred is +// used by scheduler-owned retries. func (s ExecutionStatus) CanTransitionTo(target ExecutionStatus) bool { - if s != ExecutionStatusClaimed { + if s != ExecutionStatusPending && s != ExecutionStatusDeferred { return false } return target == ExecutionStatusSubmitted || @@ -41,214 +37,330 @@ func (s ExecutionStatus) CanTransitionTo(target ExecutionStatus) bool { target == ExecutionStatusFailed } +// RequiresRetryScheduling reports whether the status requires the store to +// calculate a next-attempt time. +func (s ExecutionStatus) RequiresRetryScheduling() bool { + return s == ExecutionStatusDeferred +} + // ExecutionReason identifies the stable reason for an informational execution -// outcome without expanding the status state machine. +// result without expanding the status state machine. type ExecutionReason string const ( - ExecutionReasonNone ExecutionReason = "" - ExecutionReasonNoTargets ExecutionReason = "no_targets" - ExecutionReasonAttemptFailed ExecutionReason = "attempt_failed" + ExecutionReasonNone ExecutionReason = "" + ExecutionReasonNoTargets ExecutionReason = "no_targets" + ExecutionReasonAttemptFailed ExecutionReason = "attempt_failed" + ExecutionReasonAttemptInterrupted ExecutionReason = "attempt_interrupted" ) -type executionStateContract struct { - reasons []ExecutionReason - requiresNextAttemptAt bool +// ExecutionStatusDetails contains the status fields shared by durable state +// and attempt results. +type ExecutionStatusDetails struct { + Status ExecutionStatus + Reason ExecutionReason + StatusMessage string +} + +// Validate checks that the status, reason, and message are internally +// consistent. +func (d ExecutionStatusDetails) Validate() error { + reasons, ok := executionStatusReasons[d.Status] + if !ok { + return fmt.Errorf("unknown execution status %q", d.Status) + } + + if len(reasons) == 0 { + if d.Reason != ExecutionReasonNone { + return fmt.Errorf( + "%s execution cannot have reason %q", + d.Status, + d.Reason, + ) + } + } else if !slices.Contains(reasons, d.Reason) { + return fmt.Errorf( + "%s execution requires one of reasons %q", + d.Status, + reasons, + ) + } + + return validateOptionalString("execution status message", d.StatusMessage) } // ExecutionState contains the status-dependent state of an execution. type ExecutionState struct { - Status ExecutionStatus - Reason ExecutionReason - StatusMessage string + ExecutionStatusDetails NextAttemptAt time.Time } // Validate checks that the execution state is internally consistent. func (s ExecutionState) Validate() error { - contract, ok := executionStateContracts[s.Status] - if !ok { - return fmt.Errorf("unknown execution status %q", s.Status) + if err := s.ExecutionStatusDetails.Validate(); err != nil { + return err } - return contract.validate(s) -} -// ValidateTransition checks that the execution state is a valid transition -// target. Claimed state is established by Claim and cannot be requested through -// Transition. -func (s ExecutionState) ValidateTransition() error { - if s.Status == ExecutionStatusClaimed { - return fmt.Errorf("cannot transition execution to claimed status") + if s.Status.RequiresRetryScheduling() { + if s.NextAttemptAt.IsZero() { + return fmt.Errorf("%s execution requires next attempt time", s.Status) + } + } else { + if !s.NextAttemptAt.IsZero() { + return fmt.Errorf("%s execution cannot have next attempt time", s.Status) + } } - return s.Validate() + + return nil } -// RetryDue reports whether a deferred execution is eligible to be claimed at -// the given time. +// RetryDue reports whether a deferred execution is eligible for the scheduler +// to dispatch at the given time. func (s ExecutionState) RetryDue(now time.Time) bool { - return s.Status == ExecutionStatusDeferred && + return s.Status.RequiresRetryScheduling() && !s.NextAttemptAt.IsZero() && !now.Before(s.NextAttemptAt) } -var executionStateContracts = map[ExecutionStatus]executionStateContract{ - ExecutionStatusClaimed: {}, +var executionStatusReasons = map[ExecutionStatus][]ExecutionReason{ + ExecutionStatusPending: nil, ExecutionStatusSkipped: { - reasons: []ExecutionReason{ExecutionReasonNoTargets}, + ExecutionReasonNoTargets, }, ExecutionStatusDeferred: { - reasons: []ExecutionReason{ExecutionReasonAttemptFailed}, - requiresNextAttemptAt: true, + ExecutionReasonAttemptFailed, + ExecutionReasonAttemptInterrupted, }, - ExecutionStatusSubmitted: {}, - ExecutionStatusCompleted: {}, - ExecutionStatusFailed: {}, + ExecutionStatusSubmitted: nil, + ExecutionStatusCompleted: nil, + ExecutionStatusFailed: nil, } -func (c executionStateContract) validate(state ExecutionState) error { - if len(c.reasons) == 0 { - if state.Reason != ExecutionReasonNone { - return fmt.Errorf( - "%s execution cannot have reason %q", - state.Status, - state.Reason, - ) - } - } else if !slices.Contains(c.reasons, state.Reason) { - return fmt.Errorf( - "%s execution requires one of reasons %q", - state.Status, - c.reasons, - ) +// ExecutionResult describes the result of one dispatch attempt. Deferred +// results carry a relative delay so the store can derive NextAttemptAt from +// its authoritative clock. A zero delay makes the retry immediately eligible. +type ExecutionResult struct { + ExecutionStatusDetails + RetryAfter time.Duration +} + +// SubmittedExecutionResult creates a submitted dispatch result. +func SubmittedExecutionResult() ExecutionResult { + return ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusSubmitted, + }, } - if c.requiresNextAttemptAt { - if state.NextAttemptAt.IsZero() { - return fmt.Errorf("%s execution requires next attempt time", state.Status) +} + +// CompletedExecutionResult creates a completed dispatch result. +func CompletedExecutionResult() ExecutionResult { + return ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusCompleted, + }, + } +} + +// SkippedExecutionResult creates a skipped dispatch result. +func SkippedExecutionResult(reason ExecutionReason) ExecutionResult { + return ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusSkipped, + Reason: reason, + }, + } +} + +// DeferredExecutionResult creates a deferred dispatch result. +func DeferredExecutionResult( + reason ExecutionReason, + statusMessage string, + retryAfter time.Duration, +) ExecutionResult { + return ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: reason, + StatusMessage: statusMessage, + }, + RetryAfter: retryAfter, + } +} + +// FailedExecutionResult creates a failed dispatch result. +func FailedExecutionResult(statusMessage string) ExecutionResult { + return ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusFailed, + StatusMessage: statusMessage, + }, + } +} + +// Validate checks that the dispatch result is internally consistent. +func (r ExecutionResult) Validate() error { + if r.Status == ExecutionStatusPending { + return fmt.Errorf("pending is not an execution result") + } + + if err := r.ExecutionStatusDetails.Validate(); err != nil { + return err + } + + if r.Status.RequiresRetryScheduling() { + if r.RetryAfter < 0 { + return fmt.Errorf("deferred execution retry delay cannot be negative") } } else { - if !state.NextAttemptAt.IsZero() { - return fmt.Errorf("%s execution cannot have next attempt time", state.Status) + if r.RetryAfter != 0 { + return fmt.Errorf("%s execution cannot have retry delay", r.Status) } } - return validateOptionalString("execution status message", state.StatusMessage) + return nil } -// Execution records the durable processing state for one rule action. -type Execution struct { - ExecutionState - ID uuid.UUID +func (r ExecutionResult) stateAt(now time.Time) ExecutionState { + state := ExecutionState{ + ExecutionStatusDetails: r.ExecutionStatusDetails, + } + if r.Status.RequiresRetryScheduling() { + state.NextAttemptAt = now.Add(r.RetryAfter) + } + return state +} + +// ExecutionIdentity contains the source fields used to derive an +// execution's delivery and semantic-deduplication keys. +type ExecutionIdentity struct { EventID uuid.UUID RuleID uuid.UUID ActionID string CorrelationKey string - Observations int - Attempts int - FirstClaimedAt time.Time - UpdatedAt time.Time } -// IsOwned reports whether the current execution attempt is owned for -// processing. -// -// TODO: Treating every claimed execution as owned is sufficient for the current -// single-process in-memory store, which does not need to distinguish competing -// workers. When the database-backed store is introduced, extend ownership with -// a claim token and lease expiration, and apply the same fencing semantics to -// the in-memory store so both implementations satisfy one store contract. -func (e Execution) IsOwned() bool { - return e.Status == ExecutionStatusClaimed +// Validate checks the execution identity. +func (i ExecutionIdentity) Validate() error { + if i.EventID == uuid.Nil { + return fmt.Errorf("event id is required") + } + if i.RuleID == uuid.Nil { + return fmt.Errorf("event rule id is required") + } + if err := validateRequiredString("event rule action id", i.ActionID); err != nil { + return err + } + return validateOptionalString("event correlation_key", i.CorrelationKey) } -func (e *Execution) recordObservation(now time.Time) { - e.Observations++ - e.UpdatedAt = now +// ExecutionDeliveryKey identifies one rule action for one delivered event. +type ExecutionDeliveryKey struct { + EventID uuid.UUID + RuleID uuid.UUID + ActionID string } -// TryDeduplicate reports whether an observation is within the deduplication -// window and records it when it is. -func (e *Execution) TryDeduplicate(dedupe *Dedupe, observedAt time.Time) bool { - if dedupe == nil || - !dedupe.WithinWindow(e.FirstClaimedAt, observedAt) { - return false - } - - e.recordObservation(observedAt) +// ExecutionSemanticKey identifies one correlated rule action +// independently of an individual event delivery. +type ExecutionSemanticKey struct { + RuleID uuid.UUID + ActionID string + CorrelationKey string +} - return true +// DeliveryKey returns the delivery identity. +func (i ExecutionIdentity) DeliveryKey() ExecutionDeliveryKey { + return ExecutionDeliveryKey{ + EventID: i.EventID, + RuleID: i.RuleID, + ActionID: i.ActionID, + } } -// TryClaim records another observation and attempts to acquire an existing -// execution. It returns the execution only when a due deferred execution is -// reclaimed, and returns an error when the receiver is nil. ErrRetryScheduled -// indicates that a deferred execution is not yet due. -func (e *Execution) TryClaim(now time.Time) (*Execution, error) { - if e == nil { - return nil, fmt.Errorf("event action execution is nil") +// SemanticKey returns the semantic-deduplication identity. +func (i ExecutionIdentity) SemanticKey() ExecutionSemanticKey { + return ExecutionSemanticKey{ + RuleID: i.RuleID, + ActionID: i.ActionID, + CorrelationKey: i.CorrelationKey, } +} - e.recordObservation(now) +// Execution records the durable processing state for one rule action. +type Execution struct { + ExecutionState + ExecutionIdentity + ID uuid.UUID + Observations int + Attempts int + CreatedAt time.Time + UpdatedAt time.Time +} - if !e.IsOwned() && e.ExecutionState.RetryDue(now) { - e.ExecutionState = ExecutionState{Status: ExecutionStatusClaimed} - e.Attempts++ - return e, nil +func (e *Execution) recordObservation(now time.Time) { + e.Observations++ + if now.After(e.UpdatedAt) { + e.UpdatedAt = now } +} - if e.Status == ExecutionStatusDeferred { - return nil, fmt.Errorf("%w for %s", ErrRetryScheduled, e.NextAttemptAt) +// TryDeduplicate reports whether an observation is within the deduplication +// window and records it when it is. +func (e *Execution) TryDeduplicate(dedupe *Dedupe, observedAt time.Time) bool { + if dedupe == nil || !dedupe.WithinWindow(e.CreatedAt, observedAt) { + return false } - return nil, nil + e.recordObservation(observedAt) + return true } -// TransitionTo validates and applies an execution state transition at the -// given time. -func (e *Execution) TransitionTo(state ExecutionState, now time.Time) error { +// TransitionTo validates and applies an attempt result at the given time. A +// transition from deferred records the scheduler retry that produced the new +// result. Dispatch ownership is intentionally separate from domain status; +// the future scheduler store adds lease fencing around this transition. +func (e *Execution) TransitionTo(result ExecutionResult, now time.Time) error { if e == nil { - return fmt.Errorf("event action execution is nil") - } - if !e.IsOwned() { - return fmt.Errorf("execution %s is not owned", e.ID) + return fmt.Errorf("execution is nil") } - if err := state.ValidateTransition(); err != nil { + if err := result.Validate(); err != nil { return err } - if !e.Status.CanTransitionTo(state.Status) { + if !e.Status.CanTransitionTo(result.Status) { return fmt.Errorf( "execution %s cannot transition from %q to %q", e.ID, e.Status, - state.Status, + result.Status, ) } if now.IsZero() { return fmt.Errorf("execution transition time is required") } - if now.Before(e.FirstClaimedAt) { - return fmt.Errorf("execution transition time cannot precede first claimed time") + if now.Before(e.CreatedAt) { + return fmt.Errorf("execution transition time cannot precede creation time") } - e.ExecutionState = state - e.UpdatedAt = now + if e.Status == ExecutionStatusDeferred { + e.Attempts++ + } + e.ExecutionState = result.stateAt(now) + if now.After(e.UpdatedAt) { + e.UpdatedAt = now + } return nil } // Validate checks the durable execution aggregate. func (e *Execution) Validate() error { if e == nil { - return fmt.Errorf("event action execution is nil") + return fmt.Errorf("execution is nil") } if e.ID == uuid.Nil { - return fmt.Errorf("event action execution id is required") + return fmt.Errorf("execution id is required") } - if e.EventID == uuid.Nil { - return fmt.Errorf("event id is required") - } - if e.RuleID == uuid.Nil { - return fmt.Errorf("event rule id is required") - } - if err := validateRequiredString("event rule action id", e.ActionID); err != nil { + if err := e.ExecutionIdentity.Validate(); err != nil { return err } if e.Observations <= 0 { @@ -257,140 +369,43 @@ func (e *Execution) Validate() error { if e.Attempts <= 0 { return fmt.Errorf("execution attempts must be positive") } - if e.FirstClaimedAt.IsZero() { - return fmt.Errorf("execution first claimed time is required") + if e.CreatedAt.IsZero() { + return fmt.Errorf("execution creation time is required") } if e.UpdatedAt.IsZero() { return fmt.Errorf("execution updated time is required") } - if e.UpdatedAt.Before(e.FirstClaimedAt) { - return fmt.Errorf("execution updated time cannot precede first claimed time") + if e.UpdatedAt.Before(e.CreatedAt) { + return fmt.Errorf("execution updated time cannot precede creation time") } return e.ExecutionState.Validate() } -// ExecutionClaim identifies one delivery and optional semantic dedupe claim. -type ExecutionClaim struct { - EventID uuid.UUID - RuleID uuid.UUID - ActionID string - CorrelationKey string - Dedupe *Dedupe - Now time.Time -} - -// ExecutionDeliveryKey identifies one rule action for one delivered event. -type ExecutionDeliveryKey struct { - EventID uuid.UUID - RuleID uuid.UUID - ActionID string -} - -// ExecutionSemanticKey identifies one correlated rule action independently of -// an individual event delivery. -type ExecutionSemanticKey struct { - RuleID uuid.UUID - ActionID string - CorrelationKey string -} - -// DeliveryKey returns the delivery identity represented by the claim. -func (c ExecutionClaim) DeliveryKey() ExecutionDeliveryKey { - return ExecutionDeliveryKey{ - EventID: c.EventID, - RuleID: c.RuleID, - ActionID: c.ActionID, - } -} - -// SemanticKey returns the semantic deduplication identity represented by the -// claim. -func (c ExecutionClaim) SemanticKey() ExecutionSemanticKey { - return ExecutionSemanticKey{ - RuleID: c.RuleID, - ActionID: c.ActionID, - CorrelationKey: c.CorrelationKey, - } -} - -// NewExecution validates the claim and constructs its initial execution. -func (c ExecutionClaim) NewExecution() (*Execution, error) { - if err := c.Validate(); err != nil { +// NewExecution constructs a pending execution using the store-provided +// creation time. +func NewExecution( + identity ExecutionIdentity, + now time.Time, +) (*Execution, error) { + if err := identity.Validate(); err != nil { return nil, err } + if now.IsZero() { + return nil, fmt.Errorf("execution creation time is required") + } - return c.NewExecutionUnchecked(), nil -} - -// NewExecutionUnchecked constructs the initial execution without validating -// the claim. Callers must validate the claim before calling this method. -func (c ExecutionClaim) NewExecutionUnchecked() *Execution { return &Execution{ ExecutionState: ExecutionState{ - Status: ExecutionStatusClaimed, + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusPending, + }, }, - ID: uuid.New(), - EventID: c.EventID, - RuleID: c.RuleID, - ActionID: c.ActionID, - CorrelationKey: c.CorrelationKey, - Observations: 1, - Attempts: 1, - FirstClaimedAt: c.Now, - UpdatedAt: c.Now, - } -} - -// Validate checks the delivery identity and optional semantic-deduplication -// fields of an execution claim. -func (c ExecutionClaim) Validate() error { - if c.EventID == uuid.Nil { - return fmt.Errorf("event id is required") - } - if c.RuleID == uuid.Nil { - return fmt.Errorf("event rule id is required") - } - if err := validateRequiredString("event rule action id", c.ActionID); err != nil { - return err - } - if c.Now.IsZero() { - return fmt.Errorf("execution claim time is required") - } - if c.Dedupe != nil { - if err := c.Dedupe.Validate(); err != nil { - return fmt.Errorf("execution claim dedupe: %w", err) - } - } - if c.Dedupe != nil && c.CorrelationKey == "" { - return fmt.Errorf( - "correlation key is required by rule %s dedupe policy", - c.RuleID, - ) - } - return nil -} - -// NewExecutionClaim derives the delivery and optional semantic-deduplication -// identity for one rule action. -func NewExecutionClaim( - envelope Envelope, - rule Rule, - actionID string, - now time.Time, -) (ExecutionClaim, error) { - claim := ExecutionClaim{ - EventID: envelope.ID, - RuleID: rule.ID, - ActionID: actionID, - CorrelationKey: envelope.CorrelationKey, - Dedupe: rule.Dedupe.Clone(), - Now: now, - } - - if err := claim.Validate(); err != nil { - return ExecutionClaim{}, err - } - - return claim, nil + ExecutionIdentity: identity, + ID: uuid.New(), + Observations: 1, + Attempts: 1, + CreatedAt: now, + UpdatedAt: now, + }, nil } diff --git a/rest-api/flow/internal/eventrule/execution_test.go b/rest-api/flow/internal/eventrule/execution_test.go index f2fa49e29a..cbe23ce8df 100644 --- a/rest-api/flow/internal/eventrule/execution_test.go +++ b/rest-api/flow/internal/eventrule/execution_test.go @@ -4,7 +4,6 @@ package eventrule import ( - "errors" "testing" "time" @@ -15,26 +14,29 @@ import ( func TestExecution_Validate(t *testing.T) { now := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) valid := Execution{ - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, - ID: uuid.New(), EventID: uuid.New(), RuleID: uuid.New(), ActionID: "notify", - Observations: 1, Attempts: 1, - FirstClaimedAt: now, UpdatedAt: now, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, + ExecutionIdentity: ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "notify", + }, + ID: uuid.New(), + Observations: 1, + Attempts: 1, + CreatedAt: now, + UpdatedAt: now, } tests := map[string]struct { execution *Execution mutate func(*Execution) wantErr string }{ - "valid claimed": {execution: &valid}, - "nil": {wantErr: "event action execution is nil"}, - "claimed with status message": { - execution: &valid, - mutate: func(execution *Execution) { execution.StatusMessage = "unexpected" }, - }, + "valid pending": {execution: &valid}, + "nil": {wantErr: "execution is nil"}, "missing id": { execution: &valid, mutate: func(execution *Execution) { execution.ID = uuid.Nil }, - wantErr: "event action execution id is required", + wantErr: "execution id is required", }, "missing event id": { execution: &valid, @@ -61,39 +63,22 @@ func TestExecution_Validate(t *testing.T) { mutate: func(execution *Execution) { execution.Attempts = 0 }, wantErr: "execution attempts must be positive", }, - "missing first claimed time": { + "missing creation time": { execution: &valid, - mutate: func(execution *Execution) { execution.FirstClaimedAt = time.Time{} }, - wantErr: "execution first claimed time is required", + mutate: func(execution *Execution) { execution.CreatedAt = time.Time{} }, + wantErr: "execution creation time is required", }, "missing updated time": { execution: &valid, mutate: func(execution *Execution) { execution.UpdatedAt = time.Time{} }, wantErr: "execution updated time is required", }, - "updated before first claim": { + "updated before creation": { execution: &valid, mutate: func(execution *Execution) { - execution.UpdatedAt = execution.FirstClaimedAt.Add(-time.Second) + execution.UpdatedAt = execution.CreatedAt.Add(-time.Second) }, - wantErr: "execution updated time cannot precede first claimed time", - }, - "skipped without reason": { - execution: &valid, - mutate: func(execution *Execution) { execution.Status = ExecutionStatusSkipped }, - wantErr: "skipped execution requires one of reasons", - }, - "deferred without status message": { - execution: &valid, - mutate: func(execution *Execution) { - execution.Status = ExecutionStatusDeferred - execution.Reason = ExecutionReasonAttemptFailed - execution.NextAttemptAt = now.Add(time.Minute) - }, - }, - "failed without status message": { - execution: &valid, - mutate: func(execution *Execution) { execution.Status = ExecutionStatusFailed }, + wantErr: "execution updated time cannot precede creation time", }, "unknown status": { execution: &valid, @@ -106,12 +91,13 @@ func TestExecution_Validate(t *testing.T) { t.Run(name, func(t *testing.T) { var execution *Execution if test.execution != nil { - mutated := *test.execution - execution = &mutated + copy := *test.execution + execution = © if test.mutate != nil { test.mutate(execution) } } + err := execution.Validate() if test.wantErr == "" { require.NoError(t, err) @@ -120,63 +106,35 @@ func TestExecution_Validate(t *testing.T) { require.ErrorContains(t, err, test.wantErr) }) } -} -func TestExecution_IsOwned(t *testing.T) { - tests := map[string]struct { - status ExecutionStatus - want bool - }{ - "claimed": { - status: ExecutionStatusClaimed, - want: true, - }, - "deferred": { - status: ExecutionStatusDeferred, - }, - "completed": { - status: ExecutionStatusCompleted, - }, - } + t.Run("increments repeated deferred attempts", func(t *testing.T) { + execution, err := NewExecution(ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "retry", + }, now) + require.NoError(t, err) - for name, test := range tests { - t.Run(name, func(t *testing.T) { - execution := Execution{ - ExecutionState: ExecutionState{Status: test.status}, - } - require.Equal(t, test.want, execution.IsOwned()) - }) - } -} + result := DeferredExecutionResult( + ExecutionReasonAttemptFailed, + "retry", + 0, + ) + require.NoError(t, execution.TransitionTo(result, now.Add(time.Second))) + require.Equal(t, 1, execution.Attempts) -func TestExecutionState_ValidateTransition(t *testing.T) { - tests := map[string]struct { - state ExecutionState - wantErr string - }{ - "completed": { - state: ExecutionState{Status: ExecutionStatusCompleted}, - }, - "claimed": { - state: ExecutionState{Status: ExecutionStatusClaimed}, - wantErr: "cannot transition execution to claimed status", - }, - "invalid state": { - state: ExecutionState{Status: ExecutionStatusSkipped}, - wantErr: "skipped execution requires one of reasons", - }, - } + require.NoError(t, execution.TransitionTo(result, now.Add(2*time.Second))) + require.Equal(t, 2, execution.Attempts) - for name, test := range tests { - t.Run(name, func(t *testing.T) { - err := test.state.ValidateTransition() - if test.wantErr == "" { - require.NoError(t, err) - return - } - require.ErrorContains(t, err, test.wantErr) - }) - } + require.NoError(t, execution.TransitionTo(result, now.Add(3*time.Second))) + require.Equal(t, 3, execution.Attempts) + + require.NoError(t, execution.TransitionTo( + CompletedExecutionResult(), + now.Add(4*time.Second), + )) + require.Equal(t, 4, execution.Attempts) + }) } func TestExecutionStatus_CanTransitionTo(t *testing.T) { @@ -185,38 +143,28 @@ func TestExecutionStatus_CanTransitionTo(t *testing.T) { to ExecutionStatus want bool }{ - "claimed to submitted": { - from: ExecutionStatusClaimed, - to: ExecutionStatusSubmitted, - want: true, - }, - "claimed to completed": { - from: ExecutionStatusClaimed, + "pending to completed": { + from: ExecutionStatusPending, to: ExecutionStatusCompleted, want: true, }, - "claimed to skipped": { - from: ExecutionStatusClaimed, - to: ExecutionStatusSkipped, - want: true, - }, - "claimed to deferred": { - from: ExecutionStatusClaimed, + "pending to deferred": { + from: ExecutionStatusPending, to: ExecutionStatusDeferred, want: true, }, - "claimed to failed": { - from: ExecutionStatusClaimed, - to: ExecutionStatusFailed, + "deferred to completed": { + from: ExecutionStatusDeferred, + to: ExecutionStatusCompleted, want: true, }, - "claimed to claimed": { - from: ExecutionStatusClaimed, - to: ExecutionStatusClaimed, + "pending to pending": { + from: ExecutionStatusPending, + to: ExecutionStatusPending, }, - "completed to claimed": { + "completed to failed": { from: ExecutionStatusCompleted, - to: ExecutionStatusClaimed, + to: ExecutionStatusFailed, }, } @@ -230,81 +178,84 @@ func TestExecutionStatus_CanTransitionTo(t *testing.T) { func TestExecution_TransitionTo(t *testing.T) { now := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) zero := time.Time{} - beforeFirstClaim := now.Add(-time.Second) + beforeCreation := now.Add(-time.Second) tests := map[string]struct { execution *Execution - state ExecutionState + result ExecutionResult transition *time.Time wantErr string }{ - "completed": { + "pending to completed": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, }, - state: ExecutionState{Status: ExecutionStatusCompleted}, + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, }, - "invalid target state": { + "deferred to completed": { execution: &Execution{ - ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, + ID: uuid.New(), + ExecutionState: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + NextAttemptAt: now, + }, }, - state: ExecutionState{Status: ExecutionStatusSkipped}, - wantErr: "skipped execution requires one of reasons", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, }, - "missing transition time": { + "pending to deferred uses transition time": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, + }, + result: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + RetryAfter: time.Minute, }, - state: ExecutionState{Status: ExecutionStatusCompleted}, - transition: &zero, - wantErr: "execution transition time is required", }, - "transition before first claim": { + "pending is not an result": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, - FirstClaimedAt: now, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, }, - state: ExecutionState{Status: ExecutionStatusCompleted}, - transition: &beforeFirstClaim, - wantErr: "execution transition time cannot precede first claimed time", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, + wantErr: "pending is not an execution result", }, - "target with stale next attempt": { + "terminal source": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, - }, - state: ExecutionState{ - Status: ExecutionStatusCompleted, - NextAttemptAt: now.Add(time.Minute), + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, }, - wantErr: "completed execution cannot have next attempt time", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusFailed}}, + wantErr: "cannot transition", }, - "not owned": { + "missing transition time": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusCompleted}, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, }, - state: ExecutionState{Status: ExecutionStatusFailed}, - wantErr: "is not owned", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, + transition: &zero, + wantErr: "execution transition time is required", }, - "not owned with invalid target": { + "transition before creation": { execution: &Execution{ ID: uuid.New(), - ExecutionState: ExecutionState{Status: ExecutionStatusCompleted}, + ExecutionState: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, + CreatedAt: now, }, - state: ExecutionState{Status: ExecutionStatusClaimed}, - wantErr: "is not owned", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, + transition: &beforeCreation, + wantErr: "cannot precede creation time", }, "nil execution": { - state: ExecutionState{Status: ExecutionStatusCompleted}, - wantErr: "event action execution is nil", - }, - "nil execution with invalid target": { - state: ExecutionState{Status: ExecutionStatusClaimed}, - wantErr: "event action execution is nil", + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}, + wantErr: "execution is nil", }, } @@ -314,472 +265,386 @@ func TestExecution_TransitionTo(t *testing.T) { if test.transition != nil { transitionAt = *test.transition } - err := test.execution.TransitionTo(test.state, transitionAt) + err := test.execution.TransitionTo(test.result, transitionAt) if test.wantErr != "" { require.ErrorContains(t, err, test.wantErr) return } require.NoError(t, err) - require.Equal(t, test.state, test.execution.ExecutionState) + require.Equal(t, test.result.Status, test.execution.Status) + require.Equal(t, test.result.Reason, test.execution.Reason) + require.Equal(t, test.result.StatusMessage, test.execution.StatusMessage) + if test.result.Status == ExecutionStatusDeferred { + require.Equal(t, transitionAt.Add(test.result.RetryAfter), test.execution.NextAttemptAt) + } else { + require.True(t, test.execution.NextAttemptAt.IsZero()) + } require.Equal(t, transitionAt, test.execution.UpdatedAt) }) } } -func TestExecution_TryClaim(t *testing.T) { - now := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) - nextAttemptAt := now.Add(time.Minute) +func TestExecutionState_Validate(t *testing.T) { + nextAttemptAt := time.Now().Add(time.Minute) tests := map[string]struct { - execution *Execution - now time.Time - wantExecution bool - wantErr error + state ExecutionState + wantErr string }{ - "due deferred execution": { - execution: &Execution{ - ExecutionState: ExecutionState{ - Status: ExecutionStatusDeferred, - Reason: ExecutionReasonAttemptFailed, - NextAttemptAt: nextAttemptAt, + "pending": {state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}}, + "skipped": { + state: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusSkipped, + Reason: ExecutionReasonNoTargets, }, - Observations: 1, - Attempts: 1, }, - now: nextAttemptAt, - wantExecution: true, }, - "deferred execution not due": { - execution: &Execution{ - ExecutionState: ExecutionState{ + "deferred": { + state: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ Status: ExecutionStatusDeferred, Reason: ExecutionReasonAttemptFailed, - NextAttemptAt: nextAttemptAt, + StatusMessage: "inventory unavailable", }, - Observations: 1, - Attempts: 1, + NextAttemptAt: nextAttemptAt, }, - now: now, - wantErr: ErrRetryScheduled, }, - "existing claimed execution": { - execution: &Execution{ - ExecutionState: ExecutionState{Status: ExecutionStatusClaimed}, - Observations: 1, - Attempts: 1, + "deferred after interrupted creator attempt": { + state: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptInterrupted, + }, + NextAttemptAt: nextAttemptAt, }, - now: now, }, - "nil execution": { - now: now, - wantErr: errors.New("event action execution is nil"), + "submitted": {state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusSubmitted}}}, + "completed": {state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}}, + "failed": {state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusFailed}}}, + "unknown status": { + state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: "unknown"}}, + wantErr: "unknown execution status", + }, + "skipped without reason": { + state: ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusSkipped}}, + wantErr: "skipped execution requires one of reasons", + }, + "deferred without next attempt": { + state: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + }, + wantErr: "deferred execution requires next attempt time", + }, + "completed with next attempt": { + state: ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusCompleted, + }, + NextAttemptAt: nextAttemptAt, + }, + wantErr: "completed execution cannot have next attempt time", }, } for name, test := range tests { t.Run(name, func(t *testing.T) { - result, err := test.execution.TryClaim(test.now) - if test.wantExecution { - require.Same(t, test.execution, result) - } else { - require.Nil(t, result) - } - if test.wantErr != nil { - require.ErrorContains(t, err, test.wantErr.Error()) - if errors.Is(test.wantErr, ErrRetryScheduled) { - require.ErrorIs(t, err, ErrRetryScheduled) - } - } else { + err := test.state.Validate() + if test.wantErr == "" { require.NoError(t, err) - } - if test.execution == nil { return } - require.Equal(t, 2, test.execution.Observations) - require.Equal(t, test.now, test.execution.UpdatedAt) - if test.wantExecution { - require.Equal(t, ExecutionStatusClaimed, test.execution.Status) - require.Equal(t, 2, test.execution.Attempts) - } + require.ErrorContains(t, err, test.wantErr) }) } } -func TestExecution_TryDeduplicate(t *testing.T) { - firstClaimedAt := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) +func TestExecutionResultConstructors(t *testing.T) { tests := map[string]struct { - dedupe *Dedupe - observedAt time.Time - want bool + result ExecutionResult + want ExecutionResult }{ - "nil deduplication policy": { - observedAt: firstClaimedAt.Add(time.Second), + "submitted": { + result: SubmittedExecutionResult(), + want: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusSubmitted, + }, + }, }, - "within window": { - dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(time.Second), - want: true, + "completed": { + result: CompletedExecutionResult(), + want: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusCompleted, + }, + }, }, - "at window boundary": { - dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(time.Minute), + "skipped": { + result: SkippedExecutionResult(ExecutionReasonNoTargets), + want: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusSkipped, + Reason: ExecutionReasonNoTargets, + }, + }, + }, + "deferred": { + result: DeferredExecutionResult( + ExecutionReasonAttemptFailed, + "downstream unavailable", + time.Second, + ), + want: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + StatusMessage: "downstream unavailable", + }, + RetryAfter: time.Second, + }, + }, + "failed": { + result: FailedExecutionResult("invalid executor result"), + want: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusFailed, + StatusMessage: "invalid executor result", + }, + }, }, } for name, test := range tests { t.Run(name, func(t *testing.T) { - execution := Execution{ - FirstClaimedAt: firstClaimedAt, - UpdatedAt: firstClaimedAt, - Observations: 1, - } - - require.Equal( - t, - test.want, - execution.TryDeduplicate(test.dedupe, test.observedAt), - ) - if test.want { - require.Equal(t, 2, execution.Observations) - require.Equal(t, test.observedAt, execution.UpdatedAt) - return - } - require.Equal(t, 1, execution.Observations) - require.Equal(t, firstClaimedAt, execution.UpdatedAt) + require.Equal(t, test.want, test.result) + require.NoError(t, test.result.Validate()) }) } } -func TestNewExecutionClaim(t *testing.T) { - eventID := uuid.New() - ruleID := uuid.New() - now := time.Date(2026, 8, 4, 12, 0, 0, 0, time.UTC) +func TestExecutionResult_Validate(t *testing.T) { tests := map[string]struct { - envelope Envelope - rule Rule - actionID string - now time.Time - want ExecutionClaim - wantMessage string + result ExecutionResult + wantErr string }{ - "delivery identity without semantic dedupe": { - envelope: Envelope{ID: eventID, CorrelationKey: "unused"}, - rule: Rule{ID: ruleID}, - actionID: "notify", - now: now, - want: ExecutionClaim{ - EventID: eventID, - RuleID: ruleID, - ActionID: "notify", - CorrelationKey: "unused", - Now: now, + "completed": {result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusCompleted}}}, + "deferred": { + result: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + RetryAfter: time.Second, }, }, - "semantic dedupe identity": { - envelope: Envelope{ID: eventID, CorrelationKey: "incident-1"}, - rule: Rule{ - ID: ruleID, - Policy: Policy{Dedupe: &Dedupe{Window: time.Minute}}, - }, - actionID: "notify", - now: now, - want: ExecutionClaim{ - EventID: eventID, - RuleID: ruleID, - ActionID: "notify", - CorrelationKey: "incident-1", - Dedupe: &Dedupe{Window: time.Minute}, - Now: now, + "immediate deferred": { + result: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptInterrupted, + }, }, }, - "missing event id": { - rule: Rule{ID: ruleID}, - actionID: "notify", - now: now, - wantMessage: "event id is required", + "pending": { + result: ExecutionResult{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}, + wantErr: "pending is not an execution result", }, - "missing rule id": { - envelope: Envelope{ID: eventID}, - actionID: "notify", - now: now, - wantMessage: "event rule id is required", + "negative retry delay": { + result: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + RetryAfter: -time.Second, + }, + wantErr: "retry delay cannot be negative", }, - "missing action id": { - envelope: Envelope{ID: eventID}, - rule: Rule{ID: ruleID}, - now: now, - wantMessage: "event rule action id is empty", - }, - "missing claim time": { - envelope: Envelope{ID: eventID}, - rule: Rule{ID: ruleID}, - actionID: "notify", - wantMessage: "execution claim time is required", - }, - "dedupe requires correlation key": { - envelope: Envelope{ID: eventID}, - rule: Rule{ - ID: ruleID, - Policy: Policy{Dedupe: &Dedupe{Window: time.Minute}}, + "terminal retry delay": { + result: ExecutionResult{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusCompleted, + }, + RetryAfter: time.Second, }, - actionID: "notify", - now: now, - wantMessage: "correlation key is required", + wantErr: "completed execution cannot have retry delay", }, } for name, test := range tests { t.Run(name, func(t *testing.T) { - claim, err := NewExecutionClaim( - test.envelope, - test.rule, - test.actionID, - test.now, - ) - if test.wantMessage != "" { - require.ErrorContains(t, err, test.wantMessage) - require.Equal(t, ExecutionClaim{}, claim) + err := test.result.Validate() + if test.wantErr == "" { + require.NoError(t, err) return } - - require.NoError(t, err) - require.Equal(t, test.want, claim) + require.ErrorContains(t, err, test.wantErr) }) } } -func TestExecutionClaim_Validate(t *testing.T) { - valid := ExecutionClaim{ - EventID: uuid.New(), - RuleID: uuid.New(), - ActionID: "notify", - Now: time.Now(), +func TestExecutionState_RetryDue(t *testing.T) { + now := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) + state := ExecutionState{ + ExecutionStatusDetails: ExecutionStatusDetails{ + Status: ExecutionStatusDeferred, + Reason: ExecutionReasonAttemptFailed, + }, + NextAttemptAt: now, } + require.False(t, state.RetryDue(now.Add(-time.Nanosecond))) + require.True(t, state.RetryDue(now)) + require.False(t, (ExecutionState{ExecutionStatusDetails: ExecutionStatusDetails{Status: ExecutionStatusPending}}).RetryDue(now)) +} + +func TestExecution_TryDeduplicate(t *testing.T) { + createdAt := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) tests := map[string]struct { - mutate func(*ExecutionClaim) - wantErr string + dedupe *Dedupe + observedAt time.Time + want bool }{ - "valid delivery claim": {}, - "valid semantic dedupe claim": { - mutate: func(claim *ExecutionClaim) { - claim.CorrelationKey = "incident-1" - claim.Dedupe = &Dedupe{Window: time.Minute} - }, - }, - "negative dedupe window": { - mutate: func(claim *ExecutionClaim) { claim.Dedupe = &Dedupe{Window: -time.Second} }, - wantErr: "dedupe window must be positive", + "nil deduplication policy": {observedAt: createdAt.Add(time.Second)}, + "within window": { + dedupe: &Dedupe{Window: time.Minute}, + observedAt: createdAt.Add(time.Second), + want: true, }, - "dedupe without correlation key": { - mutate: func(claim *ExecutionClaim) { claim.Dedupe = &Dedupe{Window: time.Minute} }, - wantErr: "correlation key is required", + "at window boundary": { + dedupe: &Dedupe{Window: time.Minute}, + observedAt: createdAt.Add(time.Minute), }, } for name, test := range tests { t.Run(name, func(t *testing.T) { - claim := valid - if test.mutate != nil { - test.mutate(&claim) + execution := Execution{ + CreatedAt: createdAt, + UpdatedAt: createdAt, + Observations: 1, } - err := claim.Validate() - if test.wantErr == "" { - require.NoError(t, err) + require.Equal(t, test.want, execution.TryDeduplicate(test.dedupe, test.observedAt)) + if test.want { + require.Equal(t, 2, execution.Observations) + require.Equal(t, test.observedAt, execution.UpdatedAt) return } - require.ErrorContains(t, err, test.wantErr) + require.Equal(t, 1, execution.Observations) + require.Equal(t, createdAt, execution.UpdatedAt) }) } -} -func TestExecutionClaim_DeliveryKey(t *testing.T) { - claim := ExecutionClaim{ - EventID: uuid.New(), - RuleID: uuid.New(), - ActionID: "notify", - } - require.Equal(t, ExecutionDeliveryKey{ - EventID: claim.EventID, - RuleID: claim.RuleID, - ActionID: claim.ActionID, - }, claim.DeliveryKey()) + t.Run("out-of-order observation preserves latest update time", func(t *testing.T) { + updatedAt := createdAt.Add(30 * time.Second) + execution := Execution{ + CreatedAt: createdAt, + UpdatedAt: updatedAt, + Observations: 1, + } + require.True( + t, + execution.TryDeduplicate( + &Dedupe{Window: time.Minute}, + createdAt.Add(time.Second), + ), + ) + require.Equal(t, 2, execution.Observations) + require.Equal(t, updatedAt, execution.UpdatedAt) + }) } -func TestExecutionClaim_SemanticKey(t *testing.T) { - claim := ExecutionClaim{ - RuleID: uuid.New(), +func TestExecutionIdentity(t *testing.T) { + eventID := uuid.New() + ruleID := uuid.New() + identity := ExecutionIdentity{ + EventID: eventID, + RuleID: ruleID, ActionID: "notify", CorrelationKey: "incident-1", } - require.Equal(t, ExecutionSemanticKey{ - RuleID: claim.RuleID, - ActionID: claim.ActionID, - CorrelationKey: claim.CorrelationKey, - }, claim.SemanticKey()) -} - -func TestExecutionClaim_NewExecution(t *testing.T) { - t.Run("invalid claim", func(t *testing.T) { - invalid, err := (ExecutionClaim{}).NewExecution() - require.Error(t, err) - require.Nil(t, invalid) - }) - - t.Run("valid claim", func(t *testing.T) { - now := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) - claim := ExecutionClaim{ - EventID: uuid.New(), - RuleID: uuid.New(), + t.Run("keys", func(t *testing.T) { + require.Equal(t, ExecutionDeliveryKey{ + EventID: eventID, + RuleID: ruleID, + ActionID: "notify", + }, identity.DeliveryKey()) + require.Equal(t, ExecutionSemanticKey{ + RuleID: ruleID, ActionID: "notify", CorrelationKey: "incident-1", - Now: now, - } - - execution, err := claim.NewExecution() - require.NoError(t, err) - require.NotEqual(t, uuid.Nil, execution.ID) - require.Equal(t, claim.EventID, execution.EventID) - require.Equal(t, claim.RuleID, execution.RuleID) - require.Equal(t, claim.ActionID, execution.ActionID) - require.Equal(t, claim.CorrelationKey, execution.CorrelationKey) - require.Equal(t, ExecutionStatusClaimed, execution.Status) - require.Equal(t, 1, execution.Observations) - require.Equal(t, 1, execution.Attempts) - require.Equal(t, now, execution.FirstClaimedAt) - require.Equal(t, now, execution.UpdatedAt) - require.NoError(t, execution.Validate()) + }, identity.SemanticKey()) }) -} -func TestExecutionState_Validate(t *testing.T) { - nextAttemptAt := time.Now().Add(time.Minute) tests := map[string]struct { - state ExecutionState - wantErr string + identity ExecutionIdentity + wantErr string }{ - "skipped": { - state: ExecutionState{ - Status: ExecutionStatusSkipped, - Reason: ExecutionReasonNoTargets, - }, - }, - "deferred": { - state: ExecutionState{ - Status: ExecutionStatusDeferred, - Reason: ExecutionReasonAttemptFailed, - StatusMessage: "inventory unavailable", - NextAttemptAt: nextAttemptAt, + "valid delivery": { + identity: ExecutionIdentity{ + EventID: eventID, + RuleID: ruleID, + ActionID: "notify", }, }, - "submitted": { - state: ExecutionState{Status: ExecutionStatusSubmitted}, - }, - "completed": { - state: ExecutionState{ - Status: ExecutionStatusCompleted, - StatusMessage: "action completed", - }, - }, - "failed": { - state: ExecutionState{ - Status: ExecutionStatusFailed, - StatusMessage: "invalid target", - }, - }, - "unknown status": { - state: ExecutionState{Status: "unknown"}, - wantErr: "unknown execution status", - }, - "skipped without reason": { - state: ExecutionState{Status: ExecutionStatusSkipped}, - wantErr: "skipped execution requires one of reasons", - }, - "deferred without next attempt": { - state: ExecutionState{ - Status: ExecutionStatusDeferred, - Reason: ExecutionReasonAttemptFailed, - StatusMessage: "inventory unavailable", - }, - wantErr: "deferred execution requires next attempt time", - }, - "completed with next attempt": { - state: ExecutionState{ - Status: ExecutionStatusCompleted, - NextAttemptAt: nextAttemptAt, - }, - wantErr: "completed execution cannot have next attempt time", - }, - "submitted with next attempt": { - state: ExecutionState{ - Status: ExecutionStatusSubmitted, - NextAttemptAt: nextAttemptAt, - }, - wantErr: "submitted execution cannot have next attempt time", + "missing event id": { + identity: ExecutionIdentity{RuleID: ruleID, ActionID: "notify"}, + wantErr: "event id is required", }, - "failed with next attempt": { - state: ExecutionState{ - Status: ExecutionStatusFailed, - NextAttemptAt: nextAttemptAt, - }, - wantErr: "failed execution cannot have next attempt time", + "missing rule id": { + identity: ExecutionIdentity{EventID: eventID, ActionID: "notify"}, + wantErr: "event rule id is required", }, - "failed without status message": { - state: ExecutionState{Status: ExecutionStatusFailed}, + "missing action id": { + identity: ExecutionIdentity{EventID: eventID, RuleID: ruleID}, + wantErr: "event rule action id is empty", }, } for name, test := range tests { t.Run(name, func(t *testing.T) { - err := test.state.Validate() - if test.wantErr != "" { - require.ErrorContains(t, err, test.wantErr) + err := test.identity.Validate() + if test.wantErr == "" { + require.NoError(t, err) return } - require.NoError(t, err) + require.ErrorContains(t, err, test.wantErr) }) } } -func TestExecutionState_RetryDue(t *testing.T) { +func TestNewExecution(t *testing.T) { now := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) - tests := map[string]struct { - state ExecutionState - want bool - }{ - "deferred before retry time": { - state: ExecutionState{ - Status: ExecutionStatusDeferred, - NextAttemptAt: now.Add(time.Second), - }, - }, - "deferred at retry time": { - state: ExecutionState{ - Status: ExecutionStatusDeferred, - NextAttemptAt: now, - }, - want: true, - }, - "deferred after retry time": { - state: ExecutionState{ - Status: ExecutionStatusDeferred, - NextAttemptAt: now.Add(-time.Second), - }, - want: true, - }, - "deferred without retry time": { - state: ExecutionState{Status: ExecutionStatusDeferred}, - }, - "completed": { - state: ExecutionState{ - Status: ExecutionStatusCompleted, - NextAttemptAt: now, - }, - }, + identity := ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "notify", } - for name, test := range tests { - t.Run(name, func(t *testing.T) { - require.Equal(t, test.want, test.state.RetryDue(now)) - }) - } + t.Run("new pending execution", func(t *testing.T) { + execution, err := NewExecution(identity, now) + require.NoError(t, err) + require.NotEqual(t, uuid.Nil, execution.ID) + require.Equal(t, identity, execution.ExecutionIdentity) + require.Equal(t, ExecutionStatusPending, execution.Status) + require.Equal(t, 1, execution.Observations) + require.Equal(t, 1, execution.Attempts) + require.Equal(t, now, execution.CreatedAt) + require.Equal(t, now, execution.UpdatedAt) + require.NoError(t, execution.Validate()) + }) + + t.Run("invalid identity", func(t *testing.T) { + execution, err := NewExecution(ExecutionIdentity{}, now) + require.Error(t, err) + require.Nil(t, execution) + }) + + t.Run("missing store time", func(t *testing.T) { + execution, err := NewExecution(identity, time.Time{}) + require.ErrorContains(t, err, "execution creation time is required") + require.Nil(t, execution) + }) } diff --git a/rest-api/flow/internal/eventrule/executor/executor.go b/rest-api/flow/internal/eventrule/executor/executor.go index 64b2588ffb..293766ec94 100644 --- a/rest-api/flow/internal/eventrule/executor/executor.go +++ b/rest-api/flow/internal/eventrule/executor/executor.go @@ -1,64 +1,24 @@ // SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. // SPDX-License-Identifier: Apache-2.0 -// Package executor defines the action-execution boundary used by the event -// processor and implemented by concrete action dispatchers. +// Package executor defines the action-execution boundary used by event workers +// and the future retry scheduler. package executor import ( "context" - "encoding/json" "fmt" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" - "github.com/google/uuid" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" ) -// Target identifies one canonical target for an action. -type Target struct { - Kind eventrule.ResourceKind - ID uuid.UUID -} - -// Validate checks that the target has a supported kind and canonical identity. -func (t Target) Validate() error { - if err := t.Kind.Validate(); err != nil { - return err - } - if t.ID == uuid.Nil { - return fmt.Errorf("target id is required") - } - return nil -} - -// PrepareRequest contains the normalized inputs needed to prepare one action -// execution. -type PrepareRequest struct { - Execution eventrule.Execution - Envelope eventrule.Envelope - Resource eventrule.ResolvedResource - Action eventrule.Action -} - -// TargetRequest contains the event and task information needed to resolve -// concrete targets. -type TargetRequest struct { - EventType eventrule.Type - Payload json.RawMessage - Resource eventrule.ResolvedResource - Task eventrule.SubmitTask -} - -// TargetResolver resolves concrete targets for a task action. -type TargetResolver interface { - ResolveTargets(context.Context, TargetRequest) ([]Target, error) -} - -// ExecutionRequest contains an execution and its prepared action inputs. +// ExecutionRequest contains the execution, action, and resolved targets +// needed for one dispatch attempt. type ExecutionRequest struct { Execution eventrule.Execution Action eventrule.Action - Targets []Target + Targets []target.Target } // Validate checks the execution and action inputs. @@ -77,41 +37,18 @@ func (r ExecutionRequest) Validate() error { return nil } -// PreparationResult contains either a request ready for execution or an outcome -// that completes preparation without execution. Exactly one field is non-nil. -type PreparationResult struct { - Request *ExecutionRequest - Outcome *eventrule.ExecutionState -} - -// Validate checks that preparation produced exactly one next step. -func (p PreparationResult) Validate() error { - if (p.Request == nil) == (p.Outcome == nil) { - return fmt.Errorf("executor preparation requires exactly one request or outcome") - } - - if p.Outcome != nil { - return p.Outcome.ValidateTransition() - } - - return p.Request.Validate() -} - -// Executor owns action-specific preparation and execution outcome decisions. +// Executor performs action side effects and produces execution results. type Executor interface { - // Prepare returns a valid result and a nil error. Operational results, - // including skipped, deferred, and failed states, are represented by - // PreparationResult. A non-nil error means that no valid result was - // produced and indicates an executor contract failure. - Prepare(context.Context, PrepareRequest) (PreparationResult, error) // Execute may be called multiple times for the same Execution.ID after a - // deferred outcome. Implementations that produce external side effects must + // deferred result. Implementations that produce external side effects must // use Execution.ID, or stable keys derived from it for partitioned work, to // make repeated calls idempotent and reconcile an ambiguous prior result - // before submitting again. Attempts and rotating ownership tokens must not + // before submitting again. Attempts and rotating lease tokens must not // be used as downstream idempotency identities. Execute returns a valid - // outcome and a nil error. Operational failures are represented by deferred - // or failed outcomes. A non-nil error means that no valid outcome was - // produced and indicates an executor contract failure. - Execute(context.Context, ExecutionRequest) (eventrule.ExecutionState, error) + // result and a nil error. Operational failures are represented by deferred + // or failed results. A context cancellation or deadline error means the + // attempt was interrupted and is deferred by the dispatcher. Any other + // non-nil error means that no valid result was produced and indicates an + // executor contract failure. + Execute(context.Context, ExecutionRequest) (eventrule.ExecutionResult, error) } diff --git a/rest-api/flow/internal/eventrule/executor/executor_test.go b/rest-api/flow/internal/eventrule/executor/executor_test.go index fff8548cea..918737f3f1 100644 --- a/rest-api/flow/internal/eventrule/executor/executor_test.go +++ b/rest-api/flow/internal/eventrule/executor/executor_test.go @@ -8,68 +8,11 @@ import ( "time" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" "github.com/google/uuid" "github.com/stretchr/testify/require" ) -func TestPreparationResult_Validate(t *testing.T) { - validRequest := newValidExecutionRequest(t) - tests := map[string]struct { - preparation PreparationResult - wantErr string - }{ - "request": { - preparation: PreparationResult{Request: &validRequest}, - }, - "invalid request": { - preparation: PreparationResult{Request: &ExecutionRequest{}}, - wantErr: "execution: event action execution id is required", - }, - "outcome": { - preparation: PreparationResult{Outcome: &eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusSkipped, - Reason: eventrule.ExecutionReasonNoTargets, - }}, - }, - "neither": { - wantErr: "requires exactly one request or outcome", - }, - "both": { - preparation: PreparationResult{ - Request: &ExecutionRequest{}, - Outcome: &eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusSkipped, - Reason: eventrule.ExecutionReasonNoTargets, - }, - }, - wantErr: "requires exactly one request or outcome", - }, - "invalid outcome": { - preparation: PreparationResult{Outcome: &eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusSkipped, - }}, - wantErr: "skipped execution requires one of reasons", - }, - "claimed outcome": { - preparation: PreparationResult{Outcome: &eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusClaimed, - }}, - wantErr: "cannot transition execution to claimed status", - }, - } - - for name, test := range tests { - t.Run(name, func(t *testing.T) { - err := test.preparation.Validate() - if test.wantErr != "" { - require.ErrorContains(t, err, test.wantErr) - return - } - require.NoError(t, err) - }) - } -} - func TestExecutionRequest_Validate(t *testing.T) { valid := newValidExecutionRequest(t) tests := map[string]struct { @@ -81,7 +24,7 @@ func TestExecutionRequest_Validate(t *testing.T) { "invalid execution": { request: valid, mutate: func(request *ExecutionRequest) { request.Execution.ID = uuid.Nil }, - wantErr: "execution: event action execution id is required", + wantErr: "execution: execution id is required", }, "invalid action": { request: valid, @@ -91,14 +34,14 @@ func TestExecutionRequest_Validate(t *testing.T) { "missing target id": { request: valid, mutate: func(request *ExecutionRequest) { - request.Targets = []Target{{Kind: eventrule.ResourceKindComponent}} + request.Targets = []target.Target{{Kind: eventrule.ResourceKindComponent}} }, wantErr: "target 0: target id is required", }, "invalid target kind": { request: valid, mutate: func(request *ExecutionRequest) { - request.Targets = []Target{{Kind: "invalid", ID: uuid.New()}} + request.Targets = []target.Target{{Kind: "invalid", ID: uuid.New()}} }, wantErr: `target 0: unknown resource kind "invalid"`, }, @@ -120,48 +63,14 @@ func TestExecutionRequest_Validate(t *testing.T) { } } -func TestTarget_Validate(t *testing.T) { - tests := map[string]struct { - target Target - wantErr string - }{ - "component": { - target: Target{Kind: eventrule.ResourceKindComponent, ID: uuid.New()}, - }, - "rack": { - target: Target{Kind: eventrule.ResourceKindRack, ID: uuid.New()}, - }, - "invalid kind": { - target: Target{Kind: "invalid", ID: uuid.New()}, - wantErr: `unknown resource kind "invalid"`, - }, - "missing id": { - target: Target{Kind: eventrule.ResourceKindComponent}, - wantErr: "target id is required", - }, - } - - for name, test := range tests { - t.Run(name, func(t *testing.T) { - err := test.target.Validate() - if test.wantErr != "" { - require.ErrorContains(t, err, test.wantErr) - return - } - require.NoError(t, err) - }) - } -} - func newValidExecutionRequest(t *testing.T) ExecutionRequest { t.Helper() now := time.Date(2026, 8, 10, 12, 0, 0, 0, time.UTC) - execution, err := (eventrule.ExecutionClaim{ + execution, err := eventrule.NewExecution(eventrule.ExecutionIdentity{ EventID: uuid.New(), RuleID: uuid.New(), ActionID: "noop", - Now: now, - }).NewExecution() + }, now) require.NoError(t, err) return ExecutionRequest{ Execution: *execution, diff --git a/rest-api/flow/internal/eventrule/leakage/leakage.go b/rest-api/flow/internal/eventrule/leakage/leakage.go index af8dac8074..2d21b8f903 100644 --- a/rest-api/flow/internal/eventrule/leakage/leakage.go +++ b/rest-api/flow/internal/eventrule/leakage/leakage.go @@ -4,8 +4,15 @@ package leakage import ( + "cmp" + "context" + "fmt" + "slices" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" taskcommon "github.com/NVIDIA/infra-controller/rest-api/flow/internal/task/common" + "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/inventoryobjects/rack" "github.com/google/uuid" ) @@ -15,6 +22,124 @@ const TypeHardwareLeakDetected eventrule.Type = "hardware.leak.detected" // defaultRuleID is the stable identity of the immutable leakage fallback. var defaultRuleID = uuid.MustParse("f34b87f7-cb1b-4b08-aa51-30c0b3b58680") +// InventoryResolver is the canonical inventory resolution capability required +// by leakage target strategies. +type InventoryResolver interface { + // RackByID returns a non-nil canonical rack whose ID matches the requested + // ID, or an error when the rack cannot be resolved. + RackByID(context.Context, uuid.UUID, bool) (*rack.Rack, error) +} + +type targetResolver struct { + inventory InventoryResolver +} + +// RegisterTargetResolvers registers leakage-specific target behavior. +func RegisterTargetResolvers( + registry *target.Registry, + inventory InventoryResolver, +) error { + if inventory == nil { + return fmt.Errorf("leakage target inventory is required") + } + + resolver := &targetResolver{inventory: inventory} + return registry.Register( + TypeHardwareLeakDetected, + eventrule.TargetStrategyAffectedComponents, + resolver.resolveAffectedComponents, + ) +} + +func (r *targetResolver) resolveAffectedComponents( + ctx context.Context, + request target.ResolveRequest, +) ([]target.Target, error) { + resource := request.Resource + switch resource.Kind { + case eventrule.ResourceKindComponent: + return r.resolveAffectedComponentsInRack(ctx, resource) + case eventrule.ResourceKindRack: + resolved := target.Target{ + Kind: eventrule.ResourceKindRack, + ID: resource.RackID, + } + return []target.Target{resolved}, nil + default: + return nil, fmt.Errorf( + "%w: "+ + "leakage affected-components strategy does not support resource kind %q", + target.ErrUnresolvable, + resource.Kind, + ) + } +} + +func (r *targetResolver) resolveAffectedComponentsInRack( + ctx context.Context, + resource eventrule.ResolvedResource, +) ([]target.Target, error) { + resolvedRack, err := r.inventory.RackByID(ctx, resource.RackID, true) + if err != nil { + return nil, fmt.Errorf("resolve leakage rack %s: %w", resource.RackID, err) + } + + ids, err := affectedComponentIDs(resolvedRack, resource.ID) + if err != nil { + return nil, err + } + + targets := make([]target.Target, 0, len(ids)) + for _, id := range ids { + resolved := target.Target{ + Kind: eventrule.ResourceKindComponent, + ID: id, + } + targets = append(targets, resolved) + } + + return targets, nil +} + +// affectedComponentIDs isolates topology selection from event resolution and +// target conversion. It models leak impact using rack slot ordering. +// TODO(topology-provider integration): Replace this helper with the topology +// provider's leakage-impact resolution. +func affectedComponentIDs(resolvedRack *rack.Rack, sourceID uuid.UUID) ([]uuid.UUID, error) { + sourceSlot := -1 + for _, candidate := range resolvedRack.Components { + if candidate.Info.ID == sourceID { + sourceSlot = candidate.Position.SlotID + break + } + } + if sourceSlot < 0 { + return nil, fmt.Errorf( + "%w: "+ + "leaking component %s has no valid position in rack %s topology", + target.ErrUnresolvable, + sourceID, + resolvedRack.Info.ID, + ) + } + + ids := make([]uuid.UUID, 0, len(resolvedRack.Components)) + for _, candidate := range resolvedRack.Components { + candidateSlot := candidate.Position.SlotID + if candidateSlot < 0 || + (candidate.Info.ID != sourceID && candidateSlot >= sourceSlot) { + continue + } + + ids = append(ids, candidate.Info.ID) + } + + slices.SortFunc(ids, func(a, b uuid.UUID) int { + return cmp.Compare(a.String(), b.String()) + }) + return ids, nil +} + // DefaultRule returns the immutable safety fallback for leakage events. func DefaultRule() eventrule.Rule { return eventrule.Rule{ diff --git a/rest-api/flow/internal/eventrule/leakage/leakage_test.go b/rest-api/flow/internal/eventrule/leakage/leakage_test.go index 8180994ba3..e770d839b9 100644 --- a/rest-api/flow/internal/eventrule/leakage/leakage_test.go +++ b/rest-api/flow/internal/eventrule/leakage/leakage_test.go @@ -4,15 +4,25 @@ package leakage import ( + "context" + "errors" "testing" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" + inventoryresolver "github.com/NVIDIA/infra-controller/rest-api/flow/internal/inventory/resolver" taskcommon "github.com/NVIDIA/infra-controller/rest-api/flow/internal/task/common" + "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/common/deviceinfo" + "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/inventoryobjects/component" + "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/inventoryobjects/rack" + flowtypes "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/types" "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) +var _ InventoryResolver = (*inventoryresolver.Resolver)(nil) + func TestDefaultRuleValidates(t *testing.T) { rule := DefaultRule() require.NoError(t, rule.Validate()) @@ -26,3 +36,241 @@ func TestDefaultRuleValidates(t *testing.T) { assert.Equal(t, taskcommon.OperationCode(taskcommon.OpCodePowerControlForcePowerOff), spec.OperationCode) assert.Equal(t, eventrule.TargetStrategyAffectedComponents, spec.TargetStrategy) } + +func TestRegisterTargetResolvers(t *testing.T) { + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, &targetInventory{})) + rule := DefaultRule() + require.NoError(t, registry.ValidateRule(&rule)) +} + +func TestTargetResolver_ResolveAffectedComponents(t *testing.T) { + t.Run("source and components below it", testResolveAffectedComponents) + t.Run("whole rack", testResolveAffectedComponentsRack) + t.Run("source only", testResolveAffectedComponentsSourceOnly) + t.Run("inventory error", testResolveAffectedComponentsError) + t.Run("malformed topology", testResolveAffectedComponentsMalformedTopology) +} + +func testResolveAffectedComponentsRack(t *testing.T) { + rackID := uuid.New() + inventory := &targetInventory{} + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, inventory)) + + targets, err := registry.Resolve(context.Background(), target.ResolveRequest{ + Envelope: eventrule.Envelope{Type: TypeHardwareLeakDetected}, + Resource: eventrule.ResolvedResource{ + Kind: eventrule.ResourceKindRack, + ID: rackID, + RackID: rackID, + }, + Strategy: eventrule.TargetStrategyAffectedComponents, + }) + require.NoError(t, err) + require.Equal(t, []target.Target{{ + Kind: eventrule.ResourceKindRack, + ID: rackID, + }}, targets) + require.False(t, inventory.withComponents) +} + +func TestAffectedComponentIDs(t *testing.T) { + rackID := uuid.New() + sourceID := uuid.MustParse("00000000-0000-0000-0000-000000000003") + belowID := uuid.MustParse("00000000-0000-0000-0000-000000000001") + + tests := []struct { + name string + components []component.Component + want []uuid.UUID + wantErr string + }{ + { + name: "selects source and components below it", + components: []component.Component{ + inventoryComponent(uuid.New(), rackID, 20), + inventoryComponent(sourceID, rackID, 10), + inventoryComponent(belowID, rackID, 1), + }, + want: []uuid.UUID{belowID, sourceID}, + }, + { + name: "excludes components with invalid slots", + components: []component.Component{ + inventoryComponent(sourceID, rackID, 10), + inventoryComponent(belowID, rackID, 1), + inventoryComponent(uuid.New(), rackID, -1), + }, + want: []uuid.UUID{belowID, sourceID}, + }, + { + name: "source absent", + components: []component.Component{ + inventoryComponent(belowID, rackID, 1), + }, + wantErr: "has no valid position", + }, + { + name: "source has negative slot", + components: []component.Component{ + inventoryComponent(sourceID, rackID, -1), + }, + wantErr: "has no valid position", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + ids, err := affectedComponentIDs(&rack.Rack{ + Info: deviceinfo.DeviceInfo{ID: rackID}, + Components: test.components, + }, sourceID) + if test.wantErr != "" { + require.ErrorContains(t, err, test.wantErr) + require.Nil(t, ids) + return + } + require.NoError(t, err) + require.Equal(t, test.want, ids) + }) + } +} + +func testResolveAffectedComponents(t *testing.T) { + rackID := uuid.New() + sourceID := uuid.MustParse("00000000-0000-0000-0000-000000000003") + firstBelowID := uuid.MustParse("00000000-0000-0000-0000-000000000001") + secondBelowID := uuid.MustParse("00000000-0000-0000-0000-000000000002") + aboveID := uuid.MustParse("00000000-0000-0000-0000-000000000004") + inventory := &targetInventory{ + rack: &rack.Rack{ + Info: deviceinfo.DeviceInfo{ID: rackID}, + Components: []component.Component{ + inventoryComponent(aboveID, rackID, 20), + inventoryComponent(sourceID, rackID, 10), + inventoryComponent(secondBelowID, rackID, 5), + inventoryComponent(firstBelowID, rackID, 1), + }, + }, + } + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, inventory)) + + targets, err := registry.Resolve(context.Background(), target.ResolveRequest{ + Envelope: eventrule.Envelope{Type: TypeHardwareLeakDetected}, + Resource: eventrule.ResolvedResource{ + Kind: eventrule.ResourceKindComponent, + ID: sourceID, + RackID: rackID, + ComponentType: flowtypes.ComponentTypeCompute, + }, + Strategy: eventrule.TargetStrategyAffectedComponents, + }) + require.NoError(t, err) + require.Equal(t, []target.Target{ + {Kind: eventrule.ResourceKindComponent, ID: firstBelowID}, + {Kind: eventrule.ResourceKindComponent, ID: secondBelowID}, + {Kind: eventrule.ResourceKindComponent, ID: sourceID}, + }, targets) + require.True(t, inventory.withComponents) +} + +func testResolveAffectedComponentsSourceOnly(t *testing.T) { + rackID := uuid.New() + sourceID := uuid.New() + inventory := &targetInventory{ + rack: &rack.Rack{ + Info: deviceinfo.DeviceInfo{ID: rackID}, + Components: []component.Component{inventoryComponent(sourceID, rackID, 1)}, + }, + } + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, inventory)) + + targets, err := registry.Resolve(context.Background(), target.ResolveRequest{ + Envelope: eventrule.Envelope{Type: TypeHardwareLeakDetected}, + Resource: eventrule.ResolvedResource{ + Kind: eventrule.ResourceKindComponent, + ID: sourceID, + RackID: rackID, + ComponentType: flowtypes.ComponentTypeCompute, + }, + Strategy: eventrule.TargetStrategyAffectedComponents, + }) + require.NoError(t, err) + require.Equal(t, []target.Target{{ + Kind: eventrule.ResourceKindComponent, + ID: sourceID, + }}, targets) +} + +func testResolveAffectedComponentsError(t *testing.T) { + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, &targetInventory{ + rackErr: errors.New("inventory unavailable"), + })) + targets, err := registry.Resolve(context.Background(), target.ResolveRequest{ + Envelope: eventrule.Envelope{Type: TypeHardwareLeakDetected}, + Resource: eventrule.ResolvedResource{ + Kind: eventrule.ResourceKindComponent, + ID: uuid.New(), + RackID: uuid.New(), + ComponentType: flowtypes.ComponentTypeCompute, + }, + Strategy: eventrule.TargetStrategyAffectedComponents, + }) + require.ErrorContains(t, err, "inventory unavailable") + require.Nil(t, targets) +} + +func testResolveAffectedComponentsMalformedTopology(t *testing.T) { + rackID := uuid.New() + sourceID := uuid.New() + registry := target.New() + require.NoError(t, RegisterTargetResolvers(registry, &targetInventory{ + rack: &rack.Rack{ + Info: deviceinfo.DeviceInfo{ID: rackID}, + Components: []component.Component{ + inventoryComponent(uuid.New(), rackID, 1), + }, + }, + })) + + targets, err := registry.Resolve(context.Background(), target.ResolveRequest{ + Envelope: eventrule.Envelope{Type: TypeHardwareLeakDetected}, + Resource: eventrule.ResolvedResource{ + Kind: eventrule.ResourceKindComponent, + ID: sourceID, + RackID: rackID, + ComponentType: flowtypes.ComponentTypeCompute, + }, + Strategy: eventrule.TargetStrategyAffectedComponents, + }) + require.ErrorContains(t, err, "has no valid position") + require.ErrorIs(t, err, target.ErrUnresolvable) + require.Nil(t, targets) +} + +type targetInventory struct { + rack *rack.Rack + rackErr error + withComponents bool +} + +func (i *targetInventory) RackByID( + _ context.Context, + _ uuid.UUID, + withComponents bool, +) (*rack.Rack, error) { + i.withComponents = withComponents + return i.rack, i.rackErr +} + +func inventoryComponent(id, rackID uuid.UUID, slot int) component.Component { + return component.Component{ + Info: deviceinfo.DeviceInfo{ID: id}, + RackID: rackID, + Position: component.InRackPosition{SlotID: slot}, + } +} diff --git a/rest-api/flow/internal/eventrule/policy.go b/rest-api/flow/internal/eventrule/policy.go index 52b828a54e..ae92a5d61a 100644 --- a/rest-api/flow/internal/eventrule/policy.go +++ b/rest-api/flow/internal/eventrule/policy.go @@ -30,13 +30,13 @@ func (d *Dedupe) Clone() *Dedupe { } // WithinWindow reports whether an observation falls within the deduplication -// window anchored at the execution's first claim time. -func (d *Dedupe) WithinWindow(firstClaimedAt, observedAt time.Time) bool { +// window anchored at the action execution's creation time. +func (d *Dedupe) WithinWindow(createdAt, observedAt time.Time) bool { if d == nil { return false } - return !observedAt.Before(firstClaimedAt) && - observedAt.Sub(firstClaimedAt) < d.Window + return !observedAt.Before(createdAt) && + observedAt.Sub(createdAt) < d.Window } // Clone returns an independent copy of the policy and its mutable data. diff --git a/rest-api/flow/internal/eventrule/policy_test.go b/rest-api/flow/internal/eventrule/policy_test.go index 5f21ab59e5..e8feadb8a7 100644 --- a/rest-api/flow/internal/eventrule/policy_test.go +++ b/rest-api/flow/internal/eventrule/policy_test.go @@ -25,40 +25,40 @@ func TestDedupe_Clone(t *testing.T) { } func TestDedupe_WithinWindow(t *testing.T) { - firstClaimedAt := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) + createdAt := time.Date(2026, 8, 5, 12, 0, 0, 0, time.UTC) tests := map[string]struct { dedupe *Dedupe observedAt time.Time want bool }{ "nil policy": { - observedAt: firstClaimedAt, + observedAt: createdAt, }, - "at first claim": { + "at creation": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt, + observedAt: createdAt, want: true, }, - "immediately before first claim": { + "immediately before creation": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(-time.Nanosecond), + observedAt: createdAt.Add(-time.Nanosecond), }, - "far before first claim": { + "far before creation": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(-24 * time.Hour), + observedAt: createdAt.Add(-24 * time.Hour), }, "inside window": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(time.Minute - time.Nanosecond), + observedAt: createdAt.Add(time.Minute - time.Nanosecond), want: true, }, "at window boundary": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(time.Minute), + observedAt: createdAt.Add(time.Minute), }, "outside window": { dedupe: &Dedupe{Window: time.Minute}, - observedAt: firstClaimedAt.Add(time.Minute + time.Nanosecond), + observedAt: createdAt.Add(time.Minute + time.Nanosecond), }, } @@ -67,7 +67,7 @@ func TestDedupe_WithinWindow(t *testing.T) { require.Equal( t, test.want, - test.dedupe.WithinWindow(firstClaimedAt, test.observedAt), + test.dedupe.WithinWindow(createdAt, test.observedAt), ) }) } diff --git a/rest-api/flow/internal/eventrule/processor/config.go b/rest-api/flow/internal/eventrule/processor/config.go index 5ce8af34fc..619db63b47 100644 --- a/rest-api/flow/internal/eventrule/processor/config.go +++ b/rest-api/flow/internal/eventrule/processor/config.go @@ -5,25 +5,23 @@ package processor import ( "fmt" - "time" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/executor" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" inventoryresolver "github.com/NVIDIA/infra-controller/rest-api/flow/internal/inventory/resolver" ) -// Config contains the dependencies and runtime settings for a Processor. +// Config contains the dependencies for a Processor. type Config struct { Inventory inventoryresolver.InventoryReader Rules RuleResolver Executions eventrule.ExecutionStore + Targets target.Resolver Executor executor.Executor - - Clock func() time.Time } -// Validate checks that all required processor dependencies and runtime -// settings are valid. +// Validate checks that all required processor dependencies are present. func (c Config) Validate() error { if c.Inventory == nil { return fmt.Errorf("inventory reader is required") @@ -34,16 +32,11 @@ func (c Config) Validate() error { if c.Executions == nil { return fmt.Errorf("execution store is required") } + if c.Targets == nil { + return fmt.Errorf("target resolver is required") + } if c.Executor == nil { return fmt.Errorf("action executor is required") } - return nil } - -func (c Config) clock() func() time.Time { - if c.Clock != nil { - return c.Clock - } - return time.Now -} diff --git a/rest-api/flow/internal/eventrule/processor/config_test.go b/rest-api/flow/internal/eventrule/processor/config_test.go index 18bf0c56cf..34589fa0a8 100644 --- a/rest-api/flow/internal/eventrule/processor/config_test.go +++ b/rest-api/flow/internal/eventrule/processor/config_test.go @@ -23,10 +23,14 @@ func TestConfigValidate(t *testing.T) { mutate: func(config *Config) { config.Rules = nil }, wantErr: "rule resolver is required", }, - "missing execution store": { + "missing action execution store": { mutate: func(config *Config) { config.Executions = nil }, wantErr: "execution store is required", }, + "missing target resolver": { + mutate: func(config *Config) { config.Targets = nil }, + wantErr: "target resolver is required", + }, "missing action executor": { mutate: func(config *Config) { config.Executor = nil }, wantErr: "action executor is required", @@ -57,9 +61,9 @@ func TestNew(t *testing.T) { require.Nil(t, processor) }) - t.Run("defaults clock", func(t *testing.T) { + t.Run("constructs processor", func(t *testing.T) { processor, err := New(validProcessorConfig()) require.NoError(t, err) - require.NotNil(t, processor.now) + require.NotNil(t, processor) }) } diff --git a/rest-api/flow/internal/eventrule/processor/enrichment.go b/rest-api/flow/internal/eventrule/processor/enrichment.go index fdc9c8360b..3d13268e41 100644 --- a/rest-api/flow/internal/eventrule/processor/enrichment.go +++ b/rest-api/flow/internal/eventrule/processor/enrichment.go @@ -16,29 +16,12 @@ import ( "github.com/google/uuid" ) -// enrichment contains the canonical resource information needed by processing. -type enrichment struct { - ResolvedResource eventrule.ResolvedResource -} - -// enrich coordinates all enrichment applied to an event envelope. +// enrich resolves the canonical resource information for an event envelope. func (p *Processor) enrich( ctx context.Context, envelope eventrule.Envelope, -) (enrichment, error) { - resolvedResource, err := p.enrichResource(ctx, envelope.Resource) - if err != nil { - return enrichment{}, err - } - - return enrichment{ResolvedResource: resolvedResource}, nil -} - -// enrichResource resolves a resource's Flow identity, rack, and component type. -func (p *Processor) enrichResource( - ctx context.Context, - resource eventrule.Resource, ) (eventrule.ResolvedResource, error) { + resource := envelope.Resource switch resource.Kind { case eventrule.ResourceKindRack: return p.enrichRackResource(ctx, resource) diff --git a/rest-api/flow/internal/eventrule/processor/enrichment_test.go b/rest-api/flow/internal/eventrule/processor/enrichment_test.go index 461eca31c6..f900856ff5 100644 --- a/rest-api/flow/internal/eventrule/processor/enrichment_test.go +++ b/rest-api/flow/internal/eventrule/processor/enrichment_test.go @@ -48,9 +48,9 @@ func TestEnrichComponent(t *testing.T) { }), ) require.NoError(t, err) - require.Equal(t, componentID, result.ResolvedResource.ID) - require.Equal(t, flowtypes.ComponentTypeCompute, result.ResolvedResource.ComponentType) - require.Equal(t, rackID, result.ResolvedResource.RackID) + require.Equal(t, componentID, result.ID) + require.Equal(t, flowtypes.ComponentTypeCompute, result.ComponentType) + require.Equal(t, rackID, result.RackID) } func TestEnrichClassifiesFailures(t *testing.T) { @@ -167,8 +167,8 @@ func TestEnrichRackUsesResolvedResourceAsRack(t *testing.T) { }), ) require.NoError(t, err) - require.Equal(t, rackID, result.ResolvedResource.ID) - require.Equal(t, rackID, result.ResolvedResource.RackID) + require.Equal(t, rackID, result.ID) + require.Equal(t, rackID, result.RackID) } func validEnvelope(resource eventrule.Resource) eventrule.Envelope { diff --git a/rest-api/flow/internal/eventrule/processor/errors.go b/rest-api/flow/internal/eventrule/processor/errors.go index 92ce0212ec..0027ef338b 100644 --- a/rest-api/flow/internal/eventrule/processor/errors.go +++ b/rest-api/flow/internal/eventrule/processor/errors.go @@ -8,6 +8,7 @@ import ( "fmt" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" inventoryresolver "github.com/NVIDIA/infra-controller/rest-api/flow/internal/inventory/resolver" ) @@ -31,6 +32,11 @@ func classifyRuleError(err error) error { return err } +func isTerminalTargetError(err error) bool { + return errors.Is(err, target.ErrUnresolvable) || + errors.Is(err, inventoryresolver.ErrUnresolvable) +} + func terminalError(err error) error { return fmt.Errorf("%w: %w", ErrTerminal, err) } diff --git a/rest-api/flow/internal/eventrule/processor/execution.go b/rest-api/flow/internal/eventrule/processor/execution.go index 0c1f9e02be..043d4cd31f 100644 --- a/rest-api/flow/internal/eventrule/processor/execution.go +++ b/rest-api/flow/internal/eventrule/processor/execution.go @@ -5,141 +5,153 @@ package processor import ( "context" + "errors" "fmt" + "time" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/executor" + "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" "github.com/google/uuid" ) +// initialRetryDelay prevents a transient creator-attempt failure from becoming +// immediately due. The scheduler owns delay policy after the first attempt. +// It is global because per-rule retry customization is not required and would +// unnecessarily expand the persisted rule contract. +const initialRetryDelay = 5 * time.Second + +// executionPersistTimeout bounds result persistence after detaching it from +// the attempt context so cancellation cannot leave the execution pending. +const executionPersistTimeout = 5 * time.Second + func (p *Processor) processAction( ctx context.Context, prepared preparedEvent, action eventrule.Action, ) error { - // Check action eligibility and construct a validated execution claim. - claim, err := p.precheckExecution(prepared, action) - if err != nil || claim == nil { - return err - } - - // Atomically claim the action's delivery and optional semantic-dedupe - // identity. A nil preparation request is an accepted duplicate. - prepareReq, err := p.claimExecution(ctx, *claim, prepared, action) - if err != nil || prepareReq == nil { - return err + // Check action eligibility before creating durable state. + if !action.Condition.AppliesTo(prepared.Envelope, prepared.Resource) { + return nil } - // Prepare the action-specific request or outcome. - result, err := p.prepareExecution(ctx, *prepareReq) - if err != nil { + // Atomically create or deduplicate the execution. A nil execution is a + // deduplication result and must not be dispatched. + execution, err := p.executions.CreateExecution( + ctx, + eventrule.ExecutionIdentity{ + EventID: prepared.Envelope.ID, + RuleID: prepared.Rule.ID, + ActionID: action.ID, + CorrelationKey: prepared.Envelope.CorrelationKey, + }, + prepared.Rule.Dedupe.Clone(), + ) + if err != nil || execution == nil { return err } - // Execute a prepared request, or carry forward a preparation outcome. - outcome, err := p.performExecution(ctx, result) - if err != nil { - return err + // Resolve targets or persist the resulting target-resolution result. + targets, result := p.resolveTargets(ctx, prepared, action) + if result != nil { + return p.persistExecution(ctx, execution.ID, *result) } - // Persist the resulting execution state. - return p.persistExecution(ctx, prepareReq.Execution.ID, outcome) + // Execute the action and persist the resulting state. + return p.executeAction(ctx, execution, action, targets) } -func (p *Processor) precheckExecution( +func (p *Processor) resolveTargets( + ctx context.Context, prepared preparedEvent, action eventrule.Action, -) (*eventrule.ExecutionClaim, error) { - if !action.Condition.AppliesTo( - prepared.Envelope, - prepared.Enriched.ResolvedResource, - ) { +) ([]target.Target, *eventrule.ExecutionResult) { + strategy := action.Spec.TargetResolutionStrategy() + if !strategy.RequiresResolution() { return nil, nil } - claim, err := eventrule.NewExecutionClaim( - prepared.Envelope, - *prepared.Rule, - action.ID, - p.now(), + targets, err := p.targets.Resolve( + ctx, + target.ResolveRequest{ + Envelope: prepared.Envelope, + Resource: prepared.Resource, + Strategy: strategy, + }, ) if err != nil { - return nil, terminalError(err) - } - - return &claim, nil -} - -func (p *Processor) claimExecution( - ctx context.Context, - claim eventrule.ExecutionClaim, - prepared preparedEvent, - action eventrule.Action, -) (*executor.PrepareRequest, error) { - execution, err := p.executions.Claim(ctx, claim) - if err != nil || execution == nil { - return nil, err - } - - return &executor.PrepareRequest{ - Execution: *execution, - Envelope: prepared.Envelope, - Resource: prepared.Enriched.ResolvedResource, - Action: action, - }, nil -} - -func (p *Processor) prepareExecution( - ctx context.Context, - prepareReq executor.PrepareRequest, -) (executor.PreparationResult, error) { - result, err := p.executor.Prepare(ctx, prepareReq) - if err != nil { - return executor.PreparationResult{}, terminalError( - fmt.Errorf("executor preparation failed: %w", err), + if isTerminalTargetError(err) { + result := eventrule.FailedExecutionResult(err.Error()) + return nil, &result + } + + result := eventrule.DeferredExecutionResult( + eventrule.ExecutionReasonAttemptFailed, + err.Error(), + initialRetryDelay, ) + return nil, &result } - // Defensively validate the result at the executor boundary even though - // implementations are required to return a valid result with a nil error. - if err := result.Validate(); err != nil { - return executor.PreparationResult{}, terminalError(err) + if len(targets) == 0 { + result := eventrule.SkippedExecutionResult( + eventrule.ExecutionReasonNoTargets, + ) + return nil, &result } - return result, nil + return targets, nil } -func (p *Processor) performExecution( +func (p *Processor) executeAction( ctx context.Context, - result executor.PreparationResult, -) (eventrule.ExecutionState, error) { - // A preparation outcome already represents the resulting state, so no - // action execution is needed. - if result.Outcome != nil { - return *result.Outcome, nil - } - - outcome, err := p.executor.Execute(ctx, *result.Request) + execution *eventrule.Execution, + action eventrule.Action, + targets []target.Target, +) error { + result, err := p.executor.Execute(ctx, executor.ExecutionRequest{ + Execution: *execution, + Action: action, + Targets: targets, + }) if err != nil { - return eventrule.ExecutionState{}, terminalError( - fmt.Errorf("executor execution failed: %w", err), + if ctx.Err() != nil || + errors.Is(err, context.Canceled) || + errors.Is(err, context.DeadlineExceeded) { + result = eventrule.DeferredExecutionResult( + eventrule.ExecutionReasonAttemptInterrupted, + fmt.Sprintf("executor execution interrupted: %v", err), + initialRetryDelay, + ) + } else { + result = eventrule.FailedExecutionResult( + fmt.Sprintf("executor execution failed: %v", err), + ) + } + } else if err := result.Validate(); err != nil { + result = eventrule.FailedExecutionResult( + fmt.Sprintf("invalid executor result: %v", err), ) } - // Defensively validate the outcome at the executor boundary even though - // implementations are required to return a valid outcome with a nil error. - if err := outcome.ValidateTransition(); err != nil { - return eventrule.ExecutionState{}, terminalError(err) - } - - return outcome, nil + return p.persistExecution(ctx, execution.ID, result) } func (p *Processor) persistExecution( ctx context.Context, executionID uuid.UUID, - outcome eventrule.ExecutionState, + result eventrule.ExecutionResult, ) error { - _, err := p.executions.Transition(ctx, executionID, outcome, p.now().UTC()) + persistCtx, cancel := context.WithTimeout( + context.WithoutCancel(ctx), + executionPersistTimeout, + ) + defer cancel() + + _, err := p.executions.TransitionExecution( + persistCtx, + executionID, + result, + ) return err } diff --git a/rest-api/flow/internal/eventrule/processor/execution_test.go b/rest-api/flow/internal/eventrule/processor/execution_test.go deleted file mode 100644 index 505b103480..0000000000 --- a/rest-api/flow/internal/eventrule/processor/execution_test.go +++ /dev/null @@ -1,92 +0,0 @@ -// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -// SPDX-License-Identifier: Apache-2.0 - -package processor - -import ( - "context" - "errors" - "testing" - - "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" - "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/executor" - "github.com/stretchr/testify/require" -) - -func TestProcessor_PerformExecution(t *testing.T) { - executorErr := errors.New("executor unavailable") - completed := eventrule.ExecutionState{Status: eventrule.ExecutionStatusCompleted} - claimed := eventrule.ExecutionState{Status: eventrule.ExecutionStatusClaimed} - tests := map[string]struct { - result executor.PreparationResult - executor executor.Executor - wantOutcome eventrule.ExecutionState - wantErr error - wantMessage string - }{ - "preparation outcome pass-through": { - result: executor.PreparationResult{Outcome: &completed}, - executor: executorStub{}, - wantOutcome: completed, - }, - "valid executor outcome": { - result: executionResult(), - executor: executorStub{outcome: completed}, - wantOutcome: completed, - }, - "executor error is terminal": { - result: executionResult(), - executor: executorStub{err: executorErr}, - wantErr: executorErr, - wantMessage: "executor execution failed", - }, - "invalid executor outcome is terminal": { - result: executionResult(), - executor: executorStub{outcome: claimed}, - wantErr: ErrTerminal, - wantMessage: "cannot transition execution to claimed status", - }, - } - - for name, test := range tests { - t.Run(name, func(t *testing.T) { - processor := &Processor{executor: test.executor} - outcome, err := processor.performExecution( - context.Background(), - test.result, - ) - if test.wantErr != nil { - require.ErrorIs(t, err, test.wantErr) - require.ErrorIs(t, err, ErrTerminal) - require.ErrorContains(t, err, test.wantMessage) - require.Equal(t, eventrule.ExecutionState{}, outcome) - return - } - require.NoError(t, err) - require.Equal(t, test.wantOutcome, outcome) - }) - } -} - -func executionResult() executor.PreparationResult { - return executor.PreparationResult{Request: &executor.ExecutionRequest{}} -} - -type executorStub struct { - outcome eventrule.ExecutionState - err error -} - -func (executorStub) Prepare( - context.Context, - executor.PrepareRequest, -) (executor.PreparationResult, error) { - return executor.PreparationResult{}, nil -} - -func (e executorStub) Execute( - context.Context, - executor.ExecutionRequest, -) (eventrule.ExecutionState, error) { - return e.outcome, e.err -} diff --git a/rest-api/flow/internal/eventrule/processor/integration_test.go b/rest-api/flow/internal/eventrule/processor/integration_test.go index d4f19f59b9..4ec6d3973a 100644 --- a/rest-api/flow/internal/eventrule/processor/integration_test.go +++ b/rest-api/flow/internal/eventrule/processor/integration_test.go @@ -55,7 +55,7 @@ func TestProcessorPreparationIntegration(t *testing.T) { prepared, err := processor.prepare(ctx, envelope) require.NoError(t, err) - require.Equal(t, rackID, prepared.Enriched.ResolvedResource.RackID) + require.Equal(t, rackID, prepared.Resource.RackID) require.Equal(t, builtIn.ID, prepared.Rule.ID) require.NoError(t, ruleManager.SetEnabled(ctx, rackRule.ID, true)) diff --git a/rest-api/flow/internal/eventrule/processor/preparation.go b/rest-api/flow/internal/eventrule/processor/preparation.go index 819ee4aec0..a11a090642 100644 --- a/rest-api/flow/internal/eventrule/processor/preparation.go +++ b/rest-api/flow/internal/eventrule/processor/preparation.go @@ -5,6 +5,7 @@ package processor import ( "context" + "fmt" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" "github.com/google/uuid" @@ -19,12 +20,13 @@ type RuleResolver interface { // preparedEvent contains runtime inputs prepared for policy evaluation. type preparedEvent struct { Envelope eventrule.Envelope - Enriched enrichment + Resource eventrule.ResolvedResource Rule *eventrule.Rule } -// prepare enriches an envelope and resolves its effective rule. An absent rule -// is an accepted no-op represented by a nil preparedEvent.Rule. +// prepare enriches an envelope, resolves its effective rule, and validates +// event-rule runtime compatibility. An absent rule is an accepted no-op +// represented by a nil preparedEvent.Rule. func (p *Processor) prepare( ctx context.Context, envelope eventrule.Envelope, @@ -33,7 +35,7 @@ func (p *Processor) prepare( return preparedEvent{}, terminalError(err) } - enriched, err := p.enrich(ctx, envelope) + resource, err := p.enrich(ctx, envelope) if err != nil { return preparedEvent{}, err } @@ -41,15 +43,22 @@ func (p *Processor) prepare( rule, err := p.rules.GetEffective( ctx, envelope.Type, - enriched.ResolvedResource.RackID, + resource.RackID, ) if err != nil { return preparedEvent{}, classifyRuleError(err) } + if rule != nil && rule.Dedupe != nil && envelope.CorrelationKey == "" { + return preparedEvent{}, terminalError(fmt.Errorf( + "correlation key is required by rule %s dedupe policy", + rule.ID, + )) + } + return preparedEvent{ Envelope: envelope, - Enriched: enriched, + Resource: resource, Rule: rule, }, nil } diff --git a/rest-api/flow/internal/eventrule/processor/preparation_test.go b/rest-api/flow/internal/eventrule/processor/preparation_test.go index a89f7aa4f7..cdf728fe7b 100644 --- a/rest-api/flow/internal/eventrule/processor/preparation_test.go +++ b/rest-api/flow/internal/eventrule/processor/preparation_test.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "testing" + "time" "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/common/deviceinfo" @@ -51,6 +52,18 @@ func TestPrepare(t *testing.T) { wantTerminal: true, wantResolved: true, }, + "dedupe without correlation key is terminal": { + rule: &eventrule.Rule{ + ID: uuid.New(), + Policy: eventrule.Policy{ + Dedupe: &eventrule.Dedupe{Window: time.Minute}, + }, + }, + wantErr: ErrTerminal, + wantMessage: "correlation key is required", + wantTerminal: true, + wantResolved: true, + }, "invalid envelope is terminal": { envelope: eventrule.Envelope{}, wantErr: ErrTerminal, @@ -91,7 +104,7 @@ func TestPrepare(t *testing.T) { require.Equal(t, test.wantResolved, resolverCalled) if test.wantErr == nil { require.NoError(t, err) - require.Equal(t, rackID, result.Enriched.ResolvedResource.RackID) + require.Equal(t, rackID, result.Resource.RackID) require.Equal(t, test.rule, result.Rule) return } diff --git a/rest-api/flow/internal/eventrule/processor/process_test.go b/rest-api/flow/internal/eventrule/processor/process_test.go index 3baa881a56..0edf338cc2 100644 --- a/rest-api/flow/internal/eventrule/processor/process_test.go +++ b/rest-api/flow/internal/eventrule/processor/process_test.go @@ -6,6 +6,7 @@ package processor import ( "context" "errors" + "fmt" "sync" "sync/atomic" "testing" @@ -14,6 +15,9 @@ import ( "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule" eventexecutor "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/executor" memorystore "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/store/memory" + eventtarget "github.com/NVIDIA/infra-controller/rest-api/flow/internal/eventrule/target" + inventoryresolver "github.com/NVIDIA/infra-controller/rest-api/flow/internal/inventory/resolver" + taskcommon "github.com/NVIDIA/infra-controller/rest-api/flow/internal/task/common" "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/common/deviceinfo" "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/common/location" "github.com/NVIDIA/infra-controller/rest-api/flow/pkg/inventoryobjects/rack" @@ -22,26 +26,35 @@ import ( "github.com/stretchr/testify/require" ) -const ( - runtimeMaxExecutionAttempts = 3 - runtimeInitialRetryDelay = time.Second -) - func TestProcessor_Process(t *testing.T) { now := time.Date(2026, 8, 3, 12, 0, 0, 0, time.UTC) rackID := uuid.New() tests := map[string]struct { rule *eventrule.Rule - targets []eventexecutor.Target + ruleErr error + invalidEnvelope bool + noTargets bool targetErr error executorErr error + cancelContext bool + invalidResult bool + dedupe *eventrule.Dedupe wantErr error wantStatus eventrule.ExecutionStatus wantReason eventrule.ExecutionReason + wantMessage string wantExecutions int wantExecutorRuns int }{ "no effective rule is accepted": {}, + "invalid envelope is terminal": { + invalidEnvelope: true, + wantErr: ErrTerminal, + }, + "invalid persisted rule is terminal": { + ruleErr: fmt.Errorf("decode rule: %w", eventrule.ErrInvalidPersistedRule), + wantErr: ErrTerminal, + }, "condition skip creates no execution": { rule: processorRuntimeRule(eventrule.NewAction( "skip", @@ -51,29 +64,105 @@ func TestProcessor_Process(t *testing.T) { eventrule.Noop{}, )), }, - "noop completes": { + "dedupe without correlation key fails before condition skip": { + rule: processorRuntimeRule(eventrule.NewAction( + "skip", + eventrule.ActionCondition{ + ComponentTypes: []flowtypes.ComponentType{flowtypes.ComponentTypeNVSwitch}, + }, + eventrule.Noop{}, + )), + dedupe: &eventrule.Dedupe{Window: time.Minute}, + wantErr: ErrTerminal, + }, + "noop completes on creator fast path": { rule: processorRuntimeRule(noopAction("noop")), wantStatus: eventrule.ExecutionStatusCompleted, wantExecutions: 1, wantExecutorRuns: 1, }, - "no targets skips submit": { + "task submits on creator fast path": { + rule: processorRuntimeRule(submitAction("submit")), + wantStatus: eventrule.ExecutionStatusSubmitted, + wantExecutions: 1, + wantExecutorRuns: 1, + }, + "no targets skips task": { rule: processorRuntimeRule(submitAction("submit")), + noTargets: true, wantStatus: eventrule.ExecutionStatusSkipped, wantReason: eventrule.ExecutionReasonNoTargets, wantExecutions: 1, }, - "terminal target error produces failed outcome": { + "unresolvable target fails": { rule: processorRuntimeRule(submitAction("submit")), - targetErr: terminalError(errors.New("invalid topology")), + targetErr: fmt.Errorf("%w: invalid topology", eventtarget.ErrUnresolvable), + wantStatus: eventrule.ExecutionStatusFailed, + wantMessage: "event target cannot be resolved: invalid topology", + wantExecutions: 1, + }, + "unresolvable inventory target fails": { + rule: processorRuntimeRule(submitAction("submit")), + targetErr: fmt.Errorf( + "rack lookup: %w", + inventoryresolver.ErrUnresolvable, + ), wantStatus: eventrule.ExecutionStatusFailed, + wantMessage: "rack lookup: inventory resource cannot be resolved", + wantExecutions: 1, + }, + "transient target failure defers to scheduler": { + rule: processorRuntimeRule(submitAction("submit")), + targetErr: errors.New("inventory unavailable"), + wantStatus: eventrule.ExecutionStatusDeferred, + wantReason: eventrule.ExecutionReasonAttemptFailed, + wantMessage: "inventory unavailable", wantExecutions: 1, }, - "retryable executor error schedules retry": { + "executor contract failure is persisted": { rule: processorRuntimeRule(noopAction("noop")), - executorErr: errors.New("alert service unavailable"), + executorErr: errors.New("invalid executor result"), + wantStatus: eventrule.ExecutionStatusFailed, + wantMessage: "executor execution failed: invalid executor result", + wantExecutions: 1, + wantExecutorRuns: 1, + }, + "canceled executor attempt defers to scheduler": { + rule: processorRuntimeRule(noopAction("noop")), + executorErr: fmt.Errorf( + "worker shutdown: %w", + context.Canceled, + ), + wantStatus: eventrule.ExecutionStatusDeferred, + wantReason: eventrule.ExecutionReasonAttemptInterrupted, + wantMessage: "executor execution interrupted: worker shutdown: context canceled", + wantExecutions: 1, + wantExecutorRuns: 1, + }, + "canceled processing context defers executor error": { + rule: processorRuntimeRule(noopAction("noop")), + executorErr: errors.New("executor stopped"), + cancelContext: true, + wantStatus: eventrule.ExecutionStatusDeferred, + wantReason: eventrule.ExecutionReasonAttemptInterrupted, + wantMessage: "executor execution interrupted: executor stopped", + wantExecutions: 1, + wantExecutorRuns: 1, + }, + "expired executor attempt defers to scheduler": { + rule: processorRuntimeRule(noopAction("noop")), + executorErr: context.DeadlineExceeded, wantStatus: eventrule.ExecutionStatusDeferred, - wantReason: eventrule.ExecutionReasonAttemptFailed, + wantReason: eventrule.ExecutionReasonAttemptInterrupted, + wantMessage: "executor execution interrupted: context deadline exceeded", + wantExecutions: 1, + wantExecutorRuns: 1, + }, + "invalid executor result is persisted": { + rule: processorRuntimeRule(noopAction("noop")), + invalidResult: true, + wantStatus: eventrule.ExecutionStatusFailed, + wantMessage: `invalid executor result: unknown execution status ""`, wantExecutions: 1, wantExecutorRuns: 1, }, @@ -81,53 +170,121 @@ func TestProcessor_Process(t *testing.T) { for name, test := range tests { t.Run(name, func(t *testing.T) { - store := memorystore.New() + store := memorystore.NewWithClock(func() time.Time { return now }) + rule := test.rule + if rule != nil { + cloned := rule.Clone() + cloned.Dedupe = test.dedupe.Clone() + rule = &cloned + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + var executorRuns int + targets := defaultTargetResolver(rackID) + if test.noTargets || test.targetErr != nil { + targets = targetResolverFunc(func( + context.Context, + eventtarget.ResolveRequest, + ) ([]eventtarget.Target, error) { + return nil, test.targetErr + }) + } + execute := executorFunc(func( + _ context.Context, + request eventexecutor.ExecutionRequest, + ) (eventrule.ExecutionResult, error) { + executorRuns++ + if test.cancelContext { + cancel() + } + if test.executorErr != nil { + return eventrule.ExecutionResult{}, test.executorErr + } + if test.invalidResult { + return eventrule.ExecutionResult{}, nil + } + return successResult(request), nil + }) processor := runtimeProcessor( t, rackID, - test.rule, + rule, + test.ruleErr, store, - targetResolverFunc(func( - context.Context, - eventexecutor.TargetRequest, - ) ([]eventexecutor.Target, error) { - return test.targets, test.targetErr - }), - actionExecutorFunc(func( - context.Context, - eventexecutor.ExecutionRequest, - ) (string, error) { - executorRuns++ - return "result-1", test.executorErr - }), - &now, + targets, + execute, ) + envelope := runtimeEnvelope(rackID) + if test.invalidEnvelope { + envelope.ID = uuid.Nil + } - err := processor.Process(context.Background(), runtimeEnvelope(rackID)) + err := processor.Process(ctx, envelope) if test.wantErr == nil { require.NoError(t, err) } else { - require.ErrorContains(t, err, test.wantErr.Error()) - if errors.Is(test.wantErr, ErrTerminal) { - require.ErrorIs(t, err, ErrTerminal) - } + require.ErrorIs(t, err, test.wantErr) } require.Equal(t, test.wantExecutorRuns, executorRuns) + executions, err := store.Executions() require.NoError(t, err) require.Len(t, executions, test.wantExecutions) - if test.wantExecutions > 0 { + if test.wantExecutions == 1 { require.Equal(t, test.wantStatus, executions[0].Status) require.Equal(t, test.wantReason, executions[0].Reason) + require.Equal(t, test.wantMessage, executions[0].StatusMessage) + require.Equal(t, 1, executions[0].Attempts) + require.Equal(t, now, executions[0].CreatedAt) + if test.wantStatus == eventrule.ExecutionStatusDeferred { + require.Equal( + t, + now.Add(initialRetryDelay), + executions[0].NextAttemptAt, + ) + } } }) } t.Run("deduplication", testProcessDeduplication) - t.Run("retries and exhausts", testProcessRetriesAndExhausts) - t.Run("processes actions independently", testProcessProcessesActionsIndependently) - t.Run("concurrent duplicate executes once", testProcessConcurrentDuplicateExecutesOnce) + t.Run("transient target redelivery does not dispatch", testProcessDeferredRedelivery) + t.Run("processes actions independently", testProcessActionsIndependently) + t.Run("concurrent duplicate dispatches once", testProcessConcurrentDuplicateDispatchesOnce) +} + +func TestProcessor_persistExecution(t *testing.T) { + now := time.Date(2026, 8, 18, 12, 0, 0, 0, time.UTC) + store := memorystore.NewWithClock(func() time.Time { return now }) + created, err := store.CreateExecution( + context.Background(), + eventrule.ExecutionIdentity{ + EventID: uuid.New(), + RuleID: uuid.New(), + ActionID: "action", + }, + nil, + ) + require.NoError(t, err) + + attemptCtx, cancelAttempt := context.WithCancel(context.Background()) + cancelAttempt() + require.ErrorIs(t, attemptCtx.Err(), context.Canceled) + + processor := Processor{ + executions: transitionContextStore{ExecutionStore: store}, + } + require.NoError(t, processor.persistExecution( + attemptCtx, + created.ID, + eventrule.CompletedExecutionResult(), + )) + + executions, err := store.Executions() + require.NoError(t, err) + require.Len(t, executions, 1) + require.Equal(t, eventrule.ExecutionStatusCompleted, executions[0].Status) } func testProcessDeduplication(t *testing.T) { @@ -137,7 +294,7 @@ func testProcessDeduplication(t *testing.T) { dedupe *eventrule.Dedupe secondEventID uuid.UUID }{ - "delivery duplicate without semantic dedupe": {}, + "delivery duplicate": {}, "semantic duplicate across event IDs": { dedupe: &eventrule.Dedupe{Window: time.Minute}, secondEventID: uuid.New(), @@ -146,17 +303,19 @@ func testProcessDeduplication(t *testing.T) { for name, test := range tests { t.Run(name, func(t *testing.T) { - store := memorystore.New() + store := memorystore.NewWithClock(func() time.Time { return now }) rule := processorRuntimeRule(noopAction("noop")) rule.Dedupe = test.dedupe var runs int processor := runtimeProcessor( - t, rackID, rule, store, nil, - actionExecutorFunc(func(context.Context, eventexecutor.ExecutionRequest) (string, error) { + t, rackID, rule, nil, store, defaultTargetResolver(rackID), + executorFunc(func( + _ context.Context, + request eventexecutor.ExecutionRequest, + ) (eventrule.ExecutionResult, error) { runs++ - return "", nil + return successResult(request), nil }), - &now, ) first := runtimeEnvelope(rackID) first.CorrelationKey = "incident-1" @@ -166,86 +325,88 @@ func testProcessDeduplication(t *testing.T) { } require.NoError(t, processor.Process(context.Background(), first)) + now = now.Add(time.Second) require.NoError(t, processor.Process(context.Background(), second)) require.Equal(t, 1, runs) + executions, err := store.Executions() require.NoError(t, err) require.Len(t, executions, 1) require.Equal(t, 2, executions[0].Observations) + require.Equal(t, eventrule.ExecutionStatusCompleted, executions[0].Status) }) } } -func testProcessRetriesAndExhausts(t *testing.T) { +func testProcessDeferredRedelivery(t *testing.T) { now := time.Date(2026, 8, 3, 12, 0, 0, 0, time.UTC) rackID := uuid.New() - store := memorystore.New() - executorErr := errors.New("downstream unavailable") + store := memorystore.NewWithClock(func() time.Time { return now }) + var resolverRuns int processor := runtimeProcessor( - t, rackID, processorRuntimeRule(noopAction("noop")), store, nil, - actionExecutorFunc(func(context.Context, eventexecutor.ExecutionRequest) (string, error) { - return "", executorErr + t, + rackID, + processorRuntimeRule(submitAction("submit")), + nil, + store, + targetResolverFunc(func( + context.Context, + eventtarget.ResolveRequest, + ) ([]eventtarget.Target, error) { + resolverRuns++ + return nil, errors.New("inventory unavailable") + }), + executorFunc(func( + context.Context, + eventexecutor.ExecutionRequest, + ) (eventrule.ExecutionResult, error) { + t.Fatal("executor must not run without targets") + return eventrule.ExecutionResult{}, nil }), - &now, ) envelope := runtimeEnvelope(rackID) - for attempt := 1; attempt <= runtimeMaxExecutionAttempts; attempt++ { - err := processor.Process(context.Background(), envelope) - require.NoError(t, err) - executions, snapshotsErr := store.Executions() - require.NoError(t, snapshotsErr) - execution := executions[0] - require.Equal(t, attempt, execution.Attempts) - if attempt < runtimeMaxExecutionAttempts { - require.Equal(t, eventrule.ExecutionStatusDeferred, execution.Status) - require.Equal( - t, - now.Add(runtimeRetryDelay(attempt)), - execution.NextAttemptAt, - ) - earlyErr := processor.Process(context.Background(), envelope) - require.ErrorIs(t, earlyErr, eventrule.ErrRetryScheduled) - executions, snapshotsErr = store.Executions() - require.NoError(t, snapshotsErr) - require.Equal(t, attempt, executions[0].Attempts) - now = execution.NextAttemptAt - } else { - require.Equal(t, eventrule.ExecutionStatusFailed, execution.Status) - } - } + require.NoError(t, processor.Process(context.Background(), envelope)) + now = now.Add(time.Minute) + require.NoError(t, processor.Process(context.Background(), envelope)) + require.Equal(t, 1, resolverRuns) + + executions, err := store.Executions() + require.NoError(t, err) + require.Len(t, executions, 1) + require.Equal(t, eventrule.ExecutionStatusDeferred, executions[0].Status) + require.Equal(t, 2, executions[0].Observations) } -func testProcessProcessesActionsIndependently(t *testing.T) { +func testProcessActionsIndependently(t *testing.T) { now := time.Date(2026, 8, 3, 12, 0, 0, 0, time.UTC) rackID := uuid.New() - store := memorystore.New() - firstErr := errors.New("first action unavailable") - var executed []string + store := memorystore.NewWithClock(func() time.Time { return now }) processor := runtimeProcessor( t, rackID, processorRuntimeRule(noopAction("first"), noopAction("second")), - store, nil, - actionExecutorFunc(func( + store, + defaultTargetResolver(rackID), + executorFunc(func( _ context.Context, request eventexecutor.ExecutionRequest, - ) (string, error) { - executed = append(executed, request.Action.ID) + ) (eventrule.ExecutionResult, error) { if request.Action.ID == "first" { - return "", firstErr + return eventrule.DeferredExecutionResult( + eventrule.ExecutionReasonAttemptFailed, + "downstream unavailable", + 0, + ), nil } - return "", nil + return successResult(request), nil }), - &now, ) - err := processor.Process(context.Background(), runtimeEnvelope(rackID)) + require.NoError(t, processor.Process(context.Background(), runtimeEnvelope(rackID))) + executions, err := store.Executions() require.NoError(t, err) - require.Equal(t, []string{"first", "second"}, executed) - executions, snapshotsErr := store.Executions() - require.NoError(t, snapshotsErr) require.Len(t, executions, 2) statuses := make(map[string]eventrule.ExecutionStatus, len(executions)) for _, execution := range executions { @@ -255,39 +416,43 @@ func testProcessProcessesActionsIndependently(t *testing.T) { require.Equal(t, eventrule.ExecutionStatusCompleted, statuses["second"]) } -func testProcessConcurrentDuplicateExecutesOnce(t *testing.T) { +func testProcessConcurrentDuplicateDispatchesOnce(t *testing.T) { now := time.Date(2026, 8, 3, 12, 0, 0, 0, time.UTC) rackID := uuid.New() - store := memorystore.New() + store := memorystore.NewWithClock(func() time.Time { return now }) entered := make(chan struct{}) release := make(chan struct{}) var runs atomic.Int32 processor := runtimeProcessor( - t, rackID, processorRuntimeRule(noopAction("noop")), store, nil, - actionExecutorFunc(func(context.Context, eventexecutor.ExecutionRequest) (string, error) { + t, + rackID, + processorRuntimeRule(noopAction("noop")), + nil, + store, + defaultTargetResolver(rackID), + executorFunc(func( + _ context.Context, + request eventexecutor.ExecutionRequest, + ) (eventrule.ExecutionResult, error) { if runs.Add(1) == 1 { close(entered) } <-release - return "", nil + return successResult(request), nil }), - &now, ) envelope := runtimeEnvelope(rackID) const deliveries = 20 - start := make(chan struct{}) errs := make(chan error, deliveries) var wg sync.WaitGroup for range deliveries { wg.Add(1) go func() { defer wg.Done() - <-start errs <- processor.Process(context.Background(), envelope) }() } - close(start) <-entered close(release) wg.Wait() @@ -296,35 +461,29 @@ func testProcessConcurrentDuplicateExecutesOnce(t *testing.T) { require.NoError(t, err) } require.Equal(t, int32(1), runs.Load()) + executions, err := store.Executions() require.NoError(t, err) require.Len(t, executions, 1) + require.Equal(t, eventrule.ExecutionStatusCompleted, executions[0].Status) } func runtimeProcessor( t *testing.T, rackID uuid.UUID, rule *eventrule.Rule, + ruleErr error, store eventrule.ExecutionStore, - targets eventexecutor.TargetResolver, - execute actionExecutorFunc, - now *time.Time, + targets eventtarget.Resolver, + execute eventexecutor.Executor, ) *Processor { t.Helper() - if targets == nil { - targets = targetResolverFunc(func( - context.Context, - eventexecutor.TargetRequest, - ) ([]eventexecutor.Target, error) { - return nil, nil - }) - } resolver := ruleResolverFunc(func( context.Context, eventrule.Type, uuid.UUID, ) (*eventrule.Rule, error) { - return rule, nil + return rule, ruleErr }) processor, err := New(Config{ Inventory: &processorInventory{ @@ -332,12 +491,8 @@ func runtimeProcessor( }, Rules: resolver, Executions: store, - Executor: runtimeExecutor{ - targets: targets, - execute: execute, - now: now, - }, - Clock: func() time.Time { return *now }, + Targets: targets, + Executor: execute, }) require.NoError(t, err) return processor @@ -362,17 +517,13 @@ func newTestProcessor( Inventory: inventory, Rules: rules, Executions: memorystore.New(), - Executor: runtimeExecutor{ - targets: targetResolverFunc(func( - context.Context, - eventexecutor.TargetRequest, - ) ([]eventexecutor.Target, error) { - return nil, nil - }), - execute: func(context.Context, eventexecutor.ExecutionRequest) (string, error) { - return "", nil - }, - }, + Targets: defaultTargetResolver(uuid.New()), + Executor: executorFunc(func( + _ context.Context, + request eventexecutor.ExecutionRequest, + ) (eventrule.ExecutionResult, error) { + return successResult(request), nil + }), }) require.NoError(t, err) return processor @@ -399,105 +550,52 @@ func noopAction(id string) eventrule.Action { } func submitAction(id string) eventrule.Action { - return eventrule.NewAction(id, eventrule.ActionCondition{}, eventrule.SubmitTask{}) + return eventrule.NewAction(id, eventrule.ActionCondition{}, eventrule.SubmitTask{ + OperationType: taskcommon.TaskTypePowerControl, + OperationCode: taskcommon.OpCodePowerControlForcePowerOff, + TargetStrategy: eventrule.TargetStrategyRack, + ConflictStrategy: eventrule.ConflictStrategyQueue, + }) +} + +func successResult(request eventexecutor.ExecutionRequest) eventrule.ExecutionResult { + if request.Action.Spec.Type() == eventrule.ActionTypeSubmitTask { + return eventrule.SubmittedExecutionResult() + } + return eventrule.CompletedExecutionResult() } type targetResolverFunc func( context.Context, - eventexecutor.TargetRequest, -) ([]eventexecutor.Target, error) + eventtarget.ResolveRequest, +) ([]eventtarget.Target, error) -func (f targetResolverFunc) ResolveTargets( +func (f targetResolverFunc) Resolve( ctx context.Context, - request eventexecutor.TargetRequest, -) ([]eventexecutor.Target, error) { + request eventtarget.ResolveRequest, +) ([]eventtarget.Target, error) { return f(ctx, request) } -type actionExecutorFunc func(context.Context, eventexecutor.ExecutionRequest) (string, error) - -type runtimeExecutor struct { - targets eventexecutor.TargetResolver - execute actionExecutorFunc - now *time.Time +func defaultTargetResolver(rackID uuid.UUID) eventtarget.Resolver { + return targetResolverFunc(func( + context.Context, + eventtarget.ResolveRequest, + ) ([]eventtarget.Target, error) { + return []eventtarget.Target{{Kind: eventrule.ResourceKindRack, ID: rackID}}, nil + }) } -func (e runtimeExecutor) Prepare( - ctx context.Context, - request eventexecutor.PrepareRequest, -) (eventexecutor.PreparationResult, error) { - var targets []eventexecutor.Target - if request.Action.Spec.Type() == eventrule.ActionTypeSubmitTask { - task := request.Action.Spec.(eventrule.SubmitTask) - var err error - targets, err = e.targets.ResolveTargets(ctx, eventexecutor.TargetRequest{ - EventType: request.Envelope.Type, - Payload: request.Envelope.Payload, - Resource: request.Resource, - Task: task, - }) - if err != nil { - outcome := e.failureOutcome(request.Execution, err) - return eventexecutor.PreparationResult{Outcome: &outcome}, nil - } - if len(targets) == 0 { - outcome := eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusSkipped, - Reason: eventrule.ExecutionReasonNoTargets, - } - return eventexecutor.PreparationResult{Outcome: &outcome}, nil - } - } - - return eventexecutor.PreparationResult{Request: &eventexecutor.ExecutionRequest{ - Execution: request.Execution, - Action: request.Action, - Targets: targets, - }}, nil -} +type executorFunc func( + context.Context, + eventexecutor.ExecutionRequest, +) (eventrule.ExecutionResult, error) -func (e runtimeExecutor) Execute( +func (f executorFunc) Execute( ctx context.Context, request eventexecutor.ExecutionRequest, -) (eventrule.ExecutionState, error) { - _, err := e.execute(ctx, request) - if err != nil { - return e.failureOutcome(request.Execution, err), nil - } - - status := eventrule.ExecutionStatusCompleted - if request.Action.Spec.Type() == eventrule.ActionTypeSubmitTask { - status = eventrule.ExecutionStatusSubmitted - } - return eventrule.ExecutionState{Status: status}, nil -} - -func (e runtimeExecutor) failureOutcome( - execution eventrule.Execution, - cause error, -) eventrule.ExecutionState { - if errors.Is(cause, ErrTerminal) || - execution.Attempts >= runtimeMaxExecutionAttempts { - return eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusFailed, - StatusMessage: cause.Error(), - } - } - - now := time.Now() - if e.now != nil { - now = *e.now - } - return eventrule.ExecutionState{ - Status: eventrule.ExecutionStatusDeferred, - Reason: eventrule.ExecutionReasonAttemptFailed, - StatusMessage: cause.Error(), - NextAttemptAt: now.Add(runtimeRetryDelay(execution.Attempts)), - } -} - -func runtimeRetryDelay(attempt int) time.Duration { - return runtimeInitialRetryDelay * time.Duration(1<