diff --git a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto index bb7cb21c8c..73c7b97911 100644 --- a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto +++ b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto @@ -320,6 +320,10 @@ message BindWorkflowRunDefinitionEvent // Authoritative draft/workspace workflow identity carried from the definition binding. // Empty means the run is legacy or not backed by a draft workflow identity. string workflow_id = 11; + // Exact definition revision used by this run. Empty for legacy or ad-hoc runs. + string revision_id = 12; + // Authoritative committed definition version/watermark used by this run. Zero means unavailable. + int64 definition_version = 13; } // Idempotent provisioning command for a stable Run actor identity. The first @@ -416,6 +420,13 @@ message WorkflowStoppedEvent { string reason = 3; google.protobuf.Timestamp completed_at_utc = 4; } + +enum WorkflowRecoveryFailureKind { + WORKFLOW_RECOVERY_FAILURE_KIND_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_FAILURE_KIND_AUTHORIZATION_FAILURE = 1; + WORKFLOW_RECOVERY_FAILURE_KIND_CONFIGURATION_FAILURE = 2; +} + message WorkflowCompletedEvent { string workflow_name = 1; bool success = 2; @@ -423,6 +434,7 @@ message WorkflowCompletedEvent { string error = 4; string run_id = 5; google.protobuf.Timestamp completed_at_utc = 6; + WorkflowRecoveryFailureKind recovery_failure_kind = 7; } message WorkflowRunTerminalTimingRecordedEvent { @@ -642,6 +654,7 @@ message StepCompletedEvent { VoteAgreementDecision vote_agreement_decision = 14; WorkflowStepFailureOutcome failure_outcome = 15; WorkflowFileItemResultSet file_item_results = 16; + WorkflowRecoveryFailureKind recovery_failure_kind = 17; } message SubWorkflowInvokeRequestedEvent { @@ -664,6 +677,7 @@ message WorkflowDefinitionSnapshot map inline_workflow_yamls = 4; string scope_id = 5; int32 definition_version = 6; + string revision_id = 7; } message SubWorkflowDefinitionResolveRequestedEvent { @@ -958,6 +972,7 @@ message WorkflowLlmInvocationCompletedEvent WorkflowUsageMetrics usage = 9; WorkflowManagedHandoffOutcome managed_handoff = 10; WorkflowInteractiveAuthorizationRequirement authorization_requirement = 11; + WorkflowRecoveryFailureKind recovery_failure_kind = 12; } message WorkflowInteractiveAuthorizationRequirement diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs index f1aed680a2..e27fc69e0d 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs @@ -1,3 +1,5 @@ +using Aevatar.Workflow.Application.Abstractions.Queries; + namespace Aevatar.Workflow.Application.Abstractions.Observatory; // 06-19-workflow-run-observatory (C2): read-only, scope-gated run viewer query port. @@ -118,6 +120,8 @@ public sealed class WorkflowActivityRunFeedRow public double? DurationMs { get; init; } public long StateVersion { get; init; } + + public WorkflowRunRecoveryCapability RecoveryCapability { get; init; } = new(); } public sealed class WorkflowActivityRunInitiatorSummary @@ -221,6 +225,8 @@ public sealed class ObservatoryRunDetail public ObservatoryRunStatistics Statistics { get; init; } = new(); public ObservatoryUsageTotals UsageTotals { get; init; } = new(); + + public WorkflowRunRecoveryCapability RecoveryCapability { get; init; } = new(); } public sealed class ObservatoryRunDiagnostic diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Queries/workflow_projection_query_models.proto b/src/workflow/Aevatar.Workflow.Application.Abstractions/Queries/workflow_projection_query_models.proto index f899e1ac16..a0363c1fb4 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Queries/workflow_projection_query_models.proto +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Queries/workflow_projection_query_models.proto @@ -9,6 +9,52 @@ import "google/protobuf/wrappers.proto"; import "workflow_execution_messages.proto"; import "workflow_external_action.proto"; +enum WorkflowRecoveryEligibility { + WORKFLOW_RECOVERY_ELIGIBILITY_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_ELIGIBILITY_ELIGIBLE = 1; + WORKFLOW_RECOVERY_ELIGIBILITY_INELIGIBLE = 2; + WORKFLOW_RECOVERY_ELIGIBILITY_UNAVAILABLE = 3; +} + +enum WorkflowRecoveryUnavailableReasonCode { + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_NONE = 1; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_SOURCE_RUN_NOT_TERMINAL = 2; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_MISSING_SOURCE_FACT = 3; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_AUTHORIZATION_FAILURE = 4; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_CONFIGURATION_FAILURE = 5; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_WORKFLOW_DEFINITION_UNAVAILABLE = 6; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_LEGACY_UNAVAILABLE = 7; +} + +enum WorkflowRecoveryRecommendedAction { + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_RETRY = 1; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_RUN_AGAIN = 2; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_FIX_ACCESS = 3; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_CHANGE_CONFIGURATION = 4; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_EDIT_WORKFLOW = 5; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_EDIT_INPUT = 6; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_TECHNICAL_DETAILS = 7; +} + +message WorkflowRecoveryActionCapability { + WorkflowRecoveryEligibility eligibility = 1; + WorkflowRecoveryUnavailableReasonCode unavailable_reason_code = 2; + string unavailable_reason = 3; + repeated WorkflowRecoveryRecommendedAction recommended_actions = 4; + string starting_step_id = 5; + bool reuses_prior_step_outputs = 6; + bool may_incur_model_or_tool_cost = 7; +} + +message WorkflowRunRecoveryCapability { + WorkflowRecoveryActionCapability retry_failed_step = 1; + WorkflowRecoveryActionCapability run_again = 2; + string workflow_definition_revision_id = 3; + int64 workflow_definition_version = 4; +} + message WorkflowActorSnapshot { string actor_id = 1; string workflow_name = 2; @@ -42,6 +88,7 @@ message WorkflowActorSnapshot { WorkflowRunActivityFailureSnapshot activity_first_failure = 30; WorkflowRunActivityWaitingSnapshot activity_waiting = 31; string run_id = 32; + WorkflowRunRecoveryCapability recovery_capability = 33; } message WorkflowRunActivityInitiatorSnapshot { diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs index 127ccdf8a0..89ee3a0e3e 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs @@ -121,4 +121,5 @@ public sealed record WorkflowForkRunAcceptedReceipt( bool Accepted, string CommandId, string CorrelationId, - DateTimeOffset AckedAt); + DateTimeOffset AckedAt, + string NewRunId = ""); diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs index bcfdba71a1..9e7ba0c4d7 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs @@ -87,12 +87,14 @@ public sealed record WorkflowDefinitionBinding( string SourceKind = "", WorkflowCapabilityAdmissionPlan? CapabilityAdmissionPlan = null, string WorkflowId = "", - string RevisionId = ""); + string RevisionId = "", + long DefinitionVersion = 0); public sealed record WorkflowRunCreationReceipt( string ActorId, string DefinitionActorId, - IReadOnlyList CreatedActorIds); + IReadOnlyList CreatedActorIds, + string RunId = ""); public sealed record WorkflowDefinitionProvisioningReceipt( string ActorId, @@ -161,7 +163,9 @@ public sealed record WorkflowRunForkSeedView( string FinalError, string ScopeId = "", IReadOnlyDictionary? IdempotencyByStepId = null, - WorkflowCapabilityAdmissionPlan? CapabilityAdmissionPlan = null) + WorkflowCapabilityAdmissionPlan? CapabilityAdmissionPlan = null, + string RevisionId = "", + long DefinitionVersion = 0) { public WorkflowRunForkSeedView() : this( @@ -176,7 +180,9 @@ public WorkflowRunForkSeedView() string.Empty, string.Empty, new Dictionary(StringComparer.Ordinal), - null) + null, + string.Empty, + 0) { } } diff --git a/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs b/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs index 19d8873930..658ce77fea 100644 --- a/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs +++ b/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs @@ -230,6 +230,7 @@ private async Task BuildRunDetailAsync( Timeline = [], Diagnostics = BuildDiagnostics(snapshot, report: null, steps: [], viewEvents: []), UsageTotals = new ObservatoryUsageTotals(), + RecoveryCapability = CloneRecoveryCapability(snapshot), }; } @@ -264,6 +265,7 @@ private async Task BuildRunDetailAsync( Timeline = viewEvents, Statistics = ToStatistics(report.Summary), UsageTotals = WorkflowRunObservatoryTimelineMapper.ToUsageTotals(report.Usage), + RecoveryCapability = CloneRecoveryCapability(snapshot), }; } @@ -405,9 +407,13 @@ private static WorkflowActivityRunFeedRow ToActivityRunFeedRow(WorkflowActorSnap // Read the optional snapshot field so completed-without-start stays unavailable. DurationMs = completedAtUtc == null || !snapshot.HasDurationMs ? null : snapshot.DurationMs, StateVersion = snapshot.StateVersion, + RecoveryCapability = CloneRecoveryCapability(snapshot), }; } + private static WorkflowRunRecoveryCapability CloneRecoveryCapability(WorkflowActorSnapshot snapshot) => + snapshot.RecoveryCapability?.Clone() ?? new WorkflowRunRecoveryCapability(); + private static WorkflowActivityRunInitiatorSummary ToActivityInitiatorSummary( WorkflowRunActivityInitiatorSnapshot? source) => source == null diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs index dda64a95e9..79632d2f95 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs @@ -20,6 +20,7 @@ public WorkflowForkRunAcceptedReceipt Create( true, context.CommandId, context.CorrelationId, - DateTimeOffset.UtcNow); + DateTimeOffset.UtcNow, + target.RunId); } } diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs index 22979643fa..8990e0b574 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs @@ -12,6 +12,7 @@ public WorkflowForkRunCommandTarget( string sourceRunId, string startAtStepId, string actorId, + string runId, string workflowName, WorkflowChatRunRequest preparedRequest, IReadOnlyList? createdActorIds, @@ -19,6 +20,7 @@ public WorkflowForkRunCommandTarget( { SourceRunId = Normalize(sourceRunId); StartAtStepId = Normalize(startAtStepId); + RunId = Normalize(runId); PreparedRequest = preparedRequest ?? throw new ArgumentNullException(nameof(preparedRequest)); _innerTarget = new WorkflowRunAcceptedCommandTarget( actorId, @@ -31,6 +33,8 @@ public WorkflowForkRunCommandTarget( public string StartAtStepId { get; } + public string RunId { get; } + public WorkflowChatRunRequest PreparedRequest { get; } public string ActorId => _innerTarget.ActorId; diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs index 99758f9ae2..69bb734c56 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs @@ -91,7 +91,9 @@ public async Task ResolveFromSourceActorAsync( SourceKind: sourceBinding.SourceKind, CapabilityAdmissionPlan: sourceBinding.CapabilityAdmissionPlan?.Clone(), WorkflowId: sourceBinding.WorkflowId, - RevisionId: sourceBinding.RevisionId), + RevisionId: sourceBinding.RevisionId, + DefinitionVersion: Math.Max(0, sourceBinding.SourceVersion)), wrapAsFallbackTrigger: true, ct); @@ -373,7 +374,8 @@ private async Task ResolveFromResolvedDefinitionB SourceKind: resolvedDefinitionBinding.SourceKind?.Trim() ?? string.Empty, CapabilityAdmissionPlan: resolvedDefinitionBinding.CapabilityAdmissionPlan?.Clone(), WorkflowId: resolvedDefinitionBinding.WorkflowId?.Trim() ?? string.Empty, - RevisionId: resolvedDefinitionBinding.RevisionId?.Trim() ?? string.Empty), + RevisionId: resolvedDefinitionBinding.RevisionId?.Trim() ?? string.Empty, + DefinitionVersion: Math.Max(0, resolvedDefinitionBinding.DefinitionVersion)), wrapAsFallbackTrigger: false, ct); diff --git a/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs b/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs index ebe25b977f..fed561ebc9 100644 --- a/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs +++ b/src/workflow/Aevatar.Workflow.Core/Execution/WorkflowExecutionKernel.cs @@ -414,6 +414,7 @@ await TryStartCompensationOrPublishTerminalFailureAsync( RunId = runId, Success = false, Error = evt.Error, + RecoveryFailureKind = evt.RecoveryFailureKind, }, state, evt, @@ -442,6 +443,7 @@ await TryStartCompensationOrPublishTerminalFailureAsync( RunId = runId, Success = false, Error = evt.Error, + RecoveryFailureKind = evt.RecoveryFailureKind, }, state, evt, @@ -522,6 +524,7 @@ await TryStartCompensationOrPublishTerminalFailureAsync( RunId = runId, Success = false, Error = WorkflowRuntimeFailureMessages.StepCompletionHandlingFailed(current, evt, ex), + RecoveryFailureKind = evt.RecoveryFailureKind, }, state, evt, diff --git a/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.ApprovalExecution.cs b/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.ApprovalExecution.cs index 5ddbbc85af..27feef0337 100644 --- a/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.ApprovalExecution.cs +++ b/src/workflow/Aevatar.Workflow.Core/Modules/ConnectorCallModule.ApprovalExecution.cs @@ -486,6 +486,12 @@ await WorkflowRuntimeCallbackLeaseSupport.TryCancelAsync( error, durationMs, responseAnnotations); + if (!success && + string.Equals(reasonCode, "connector_authorization_unavailable", StringComparison.Ordinal)) + { + completion.RecoveryFailureKind = WorkflowRecoveryFailureKind.AuthorizationFailure; + } + completion.Annotations["connector.approval.action_id"] = snapshot.Plan.ActionId; completion.Annotations["connector.approval.status"] = snapshot.ApprovalStatus.ToString(); completion.Annotations["connector.approval.execution_status"] = snapshot.ExecutionStatus.ToString(); diff --git a/src/workflow/Aevatar.Workflow.Core/Modules/LLMCallModule.cs b/src/workflow/Aevatar.Workflow.Core/Modules/LLMCallModule.cs index 50fb7b56fe..c11543a6bb 100644 --- a/src/workflow/Aevatar.Workflow.Core/Modules/LLMCallModule.cs +++ b/src/workflow/Aevatar.Workflow.Core/Modules/LLMCallModule.cs @@ -182,7 +182,13 @@ private async Task HandleLlmInvocationCompletedAsync( var publisherActorId = envelope.Route?.PublisherActorId ?? ctx.AgentId; if (!evt.Success) { - await PublishFailedCompletionAsync(pending, string.IsNullOrWhiteSpace(evt.Error) ? "LLM call failed." : evt.Error, publisherActorId, ctx, ct); + await PublishFailedCompletionAsync( + pending, + string.IsNullOrWhiteSpace(evt.Error) ? "LLM call failed." : evt.Error, + publisherActorId, + evt.RecoveryFailureKind, + ctx, + ct); await RemovePendingAsync(sessionId, pending, ctx, ct); return; } @@ -474,11 +480,37 @@ private static Task PublishFailedCompletionAsync( string workerId, IWorkflowExecutionContext ctx, CancellationToken ct) => + PublishFailedCompletionAsync(pending, error, workerId, WorkflowRecoveryFailureKind.Unspecified, ctx, ct); + + private static Task PublishFailedCompletionAsync( + PendingLlmCallState pending, + string error, + string workerId, + WorkflowRecoveryFailureKind recoveryFailureKind, + IWorkflowExecutionContext ctx, + CancellationToken ct) => PublishFailedCompletionAsync( pending.StepId, pending.RunId, error, workerId, + recoveryFailureKind, + ctx, + ct); + + private static Task PublishFailedCompletionAsync( + string stepId, + string runId, + string error, + string workerId, + IWorkflowExecutionContext ctx, + CancellationToken ct) => + PublishFailedCompletionAsync( + stepId, + runId, + error, + workerId, + WorkflowRecoveryFailureKind.Unspecified, ctx, ct); @@ -487,6 +519,7 @@ private static Task PublishFailedCompletionAsync( string runId, string error, string workerId, + WorkflowRecoveryFailureKind recoveryFailureKind, IWorkflowExecutionContext ctx, CancellationToken ct) => ctx.PublishAsync( @@ -497,6 +530,7 @@ private static Task PublishFailedCompletionAsync( Success = false, Error = error, WorkerId = string.IsNullOrWhiteSpace(workerId) ? ctx.AgentId : workerId, + RecoveryFailureKind = recoveryFailureKind, }, TopologyAudience.Self, ct); diff --git a/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs b/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs index cc4ce8358b..e32c1770f4 100644 --- a/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs +++ b/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs @@ -1491,6 +1491,8 @@ private EventEnvelope CreateWorkflowRunBindEnvelope( ? state.ScopeId ?? string.Empty : definition.ScopeId, InlineWorkflowYamls = { inlineWorkflowYamls }, + RevisionId = definition.RevisionId ?? string.Empty, + DefinitionVersion = Math.Max(0, definition.DefinitionVersion), ExpectedExecutionMode = state.ExpectedExecutionMode, }; diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowGAgent.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowGAgent.cs index d1b3630f90..2dbc57f441 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowGAgent.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowGAgent.cs @@ -402,6 +402,7 @@ private WorkflowDefinitionSnapshot BuildDefinitionSnapshot() WorkflowYaml = State.WorkflowYaml ?? string.Empty, ScopeId = State.ScopeId ?? string.Empty, DefinitionVersion = State.Version, + RevisionId = State.RevisionId ?? string.Empty, }; foreach (var (workflowName, workflowYaml) in State.InlineWorkflowYamls) diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs index a0b8ed6f75..c986024bee 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs @@ -27,6 +27,8 @@ await BindWorkflowRunDefinitionAsync( binding.RunOrigin, binding.ScheduleId, binding.WorkflowId, + binding.RevisionId, + binding.DefinitionVersion, binding.CapabilityAdmissionPlan, binding.ExpectedExecutionMode); } diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs index f0b643e005..f5f070a6d1 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs @@ -422,6 +422,8 @@ public async Task BindWorkflowRunDefinitionAsync( string? runOrigin, string? scheduleId, string? workflowId, + string? revisionId, + long definitionVersion, WorkflowCapabilityAdmissionPlan? capabilityAdmissionPlan, ExternalCapabilityExecutionMode expectedExecutionMode, CancellationToken ct = default) @@ -452,6 +454,8 @@ public async Task BindWorkflowRunDefinitionAsync( RunOrigin = runOrigin?.Trim() ?? string.Empty, ScheduleId = scheduleId?.Trim() ?? string.Empty, WorkflowId = workflowId?.Trim() ?? string.Empty, + RevisionId = revisionId?.Trim() ?? string.Empty, + DefinitionVersion = Math.Max(0, definitionVersion), CapabilityAdmissionPlan = capabilityAdmissionPlan?.Clone(), ExpectedExecutionMode = expectedExecutionMode, }; @@ -480,6 +484,8 @@ public Task HandleBindWorkflowRunDefinition(BindWorkflowRunDefinitionEvent reque request.RunOrigin, request.ScheduleId, request.WorkflowId, + request.RevisionId, + request.DefinitionVersion, request.CapabilityAdmissionPlan, request.ExpectedExecutionMode); @@ -624,6 +630,7 @@ await HandleWorkflowCompleted(new WorkflowCompletedEvent WorkflowName = _compiledWorkflow.Name, Success = false, Error = InputFileBindingError, + RecoveryFailureKind = WorkflowRecoveryFailureKind.ConfigurationFailure, }, request.SessionId); return null; } @@ -1041,6 +1048,7 @@ private static WorkflowLlmInvocationCompletedEvent BuildInteractiveTerminalFailu Error = NormalizeInteractiveValue(completed.Error, 512) ?? NormalizeInteractiveValue(requirement.SafeMessage, 512) ?? "The requested service requires authorization.", + RecoveryFailureKind = WorkflowRecoveryFailureKind.AuthorizationFailure, Usage = completed.Usage?.Clone(), }; @@ -2177,12 +2185,19 @@ private WorkflowRunState ApplyBindWorkflowRunDefinition(WorkflowRunState current next.WorkflowId = string.IsNullOrWhiteSpace(evt.WorkflowId) ? current.WorkflowId : evt.WorkflowId.Trim(); + next.RevisionId = string.IsNullOrWhiteSpace(evt.RevisionId) + ? current.RevisionId + : evt.RevisionId.Trim(); + next.DefinitionVersion = evt.DefinitionVersion <= 0 + ? current.DefinitionVersion + : evt.DefinitionVersion; next.CapabilityAdmissionPlan = evt.CapabilityAdmissionPlan?.Clone(); next.ExpectedExecutionMode = evt.ExpectedExecutionMode; next.Status = "bound"; next.Input = string.Empty; next.FinalOutput = string.Empty; next.FinalError = string.Empty; + next.TerminalRecoveryFailureKind = WorkflowRecoveryFailureKind.Unspecified; next.CompletedAtUtc = null; next.ClearDurationMs(); next.Initiator = null; @@ -2273,6 +2288,7 @@ private static WorkflowRunState ApplyWorkflowRunExecutionStarted(WorkflowRunStat next.Status = RunningStatus; next.FinalOutput = string.Empty; next.FinalError = string.Empty; + next.TerminalRecoveryFailureKind = WorkflowRecoveryFailureKind.Unspecified; next.CompletedAtUtc = null; next.ClearDurationMs(); next.Initiator = BuildInitiator(evt.ExecutionContextDelta?.CallerCredential?.NyxIdAuthority); @@ -2598,6 +2614,7 @@ private static WorkflowRunState ApplyWorkflowCompensationFailed( next.Status = FailedStatus; next.FinalOutput = string.Empty; next.FinalError = evt.Error ?? string.Empty; + next.TerminalRecoveryFailureKind = WorkflowRecoveryFailureKind.Unspecified; next.SagaStatus = WorkflowSagaStatus.CompensationDeadLetter; next.CompensationExecutionId = string.Empty; next.DeadLetterFailedCompensationStepId = evt.FailedCompensationStepId ?? string.Empty; @@ -2632,6 +2649,7 @@ private static WorkflowRunState ApplyWorkflowStopped(WorkflowRunState current, W next.FinalOutput = string.Empty; if (!string.IsNullOrWhiteSpace(evt.Reason)) next.FinalError = evt.Reason; + next.TerminalRecoveryFailureKind = WorkflowRecoveryFailureKind.Unspecified; ApplyTerminalTiming(next, evt.CompletedAtUtc); next.ExecutionStates.Clear(); next.ExecutionContext = new WorkflowRunExecutionContextState(); @@ -2649,6 +2667,9 @@ private static WorkflowRunState ApplyWorkflowCompleted(WorkflowRunState current, next.Status = evt.Success ? CompletedStatus : FailedStatus; next.FinalOutput = evt.Output ?? string.Empty; next.FinalError = evt.Error ?? string.Empty; + next.TerminalRecoveryFailureKind = evt.Success + ? WorkflowRecoveryFailureKind.Unspecified + : evt.RecoveryFailureKind; ApplyTerminalTiming(next, evt.CompletedAtUtc); next.TerminalWorkflowCompletionRecorded = true; next.ExecutionContext = new WorkflowRunExecutionContextState(); @@ -3625,6 +3646,8 @@ await PersistDomainEventAsync(new BindWorkflowRunDefinitionEvent RunOrigin = State.RunOrigin ?? string.Empty, ScheduleId = State.ScheduleId ?? string.Empty, WorkflowId = State.WorkflowId ?? string.Empty, + RevisionId = State.RevisionId ?? string.Empty, + DefinitionVersion = Math.Max(0, State.DefinitionVersion), CapabilityAdmissionPlan = State.CapabilityAdmissionPlan?.Clone(), ExpectedExecutionMode = State.ExpectedExecutionMode, InlineWorkflowYamls = { State.InlineWorkflowYamls }, diff --git a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto index 116d851eef..a81dd57671 100644 --- a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto +++ b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto @@ -199,6 +199,9 @@ message WorkflowRunState google.protobuf.Timestamp completed_at_utc = 49; optional double duration_ms = 50; WorkflowRunInitiatorState initiator = 51; + string revision_id = 52; + int64 definition_version = 53; + aevatar.workflow.WorkflowRecoveryFailureKind terminal_recovery_failure_kind = 54; } message WorkflowRunInitiatorState diff --git a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/ChatEndpoints.cs b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/ChatEndpoints.cs index 00c312796b..7adc5a3ae4 100644 --- a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/ChatEndpoints.cs +++ b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/ChatEndpoints.cs @@ -820,11 +820,12 @@ public static async Task HandleForkRun( if (!dispatch.Succeeded || dispatch.Receipt == null) return MapForkRunFailure(dispatch.Error, scope); - var statusUrl = BuildWorkflowRunStatusUrl(dispatch.Receipt.NewRunActorId); + var statusUrl = BuildWorkflowRunStatusUrl(dispatch.Receipt); return Results.Accepted(statusUrl, new { accepted = true, sourceRunId = dispatch.Receipt.SourceRunId, + newRunId = dispatch.Receipt.NewRunId, newRunActorId = dispatch.Receipt.NewRunActorId, workflowName = dispatch.Receipt.WorkflowName, acceptedCommandId = dispatch.Receipt.CommandId, @@ -846,6 +847,14 @@ public static async Task HandleForkRun( private static string BuildWorkflowRunStatusUrl(WorkflowRunControlAcceptedReceipt receipt) => BuildWorkflowRunStatusUrl(receipt.ActorId); + // Implement (issue #3251): + // Behavior: fork receipts expose a routable run id while preserving the technical actor address. + // Why this shape: status links must not treat NewRunActorId as the run identity when NewRunId exists. + private static string BuildWorkflowRunStatusUrl(WorkflowForkRunAcceptedReceipt receipt) => + string.IsNullOrWhiteSpace(receipt.NewRunId) + ? BuildWorkflowRunStatusUrl(receipt.NewRunActorId) + : $"/api/workflow/observatory/runs/{Uri.EscapeDataString(receipt.NewRunId)}"; + // Refactor (iter165/cluster-003-workflow-actor-shaped-query-surface): // Old pattern: accepted status links pointed at /api/actors/{actorId}. // New principle: accepted status links point at the workflow actor current-state readmodel resource. diff --git a/src/workflow/Aevatar.Workflow.Infrastructure/Runs/WorkflowRunActorPort.cs b/src/workflow/Aevatar.Workflow.Infrastructure/Runs/WorkflowRunActorPort.cs index 26b5d10c4a..4aa4e753f0 100644 --- a/src/workflow/Aevatar.Workflow.Infrastructure/Runs/WorkflowRunActorPort.cs +++ b/src/workflow/Aevatar.Workflow.Infrastructure/Runs/WorkflowRunActorPort.cs @@ -115,6 +115,8 @@ await _dispatchPort.DispatchAsync( definition.RunOrigin, definition.ScheduleId, definition.WorkflowId, + definition.RevisionId, + definition.DefinitionVersion, definition.ExpectedExecutionMode, definitionResolution.CapabilityAdmissionPlan), ct); @@ -122,7 +124,11 @@ await _dispatchPort.DispatchAsync( return new WorkflowRunCreationReceipt( runActor.Id, definitionResolution.ActorId, - createdActorIds); + createdActorIds, + // Fix (review round 1, F2): + // CreateRunAsync previously copied the technical actor id into RunId. + // No authoritative routable run id exists here, so leave RunId empty for honest fallback. + RunId: string.Empty); } catch { @@ -208,6 +214,8 @@ await ValidateDefinitionArtifactAsync( definition.RunOrigin, definition.ScheduleId, definition.WorkflowId, + definition.RevisionId, + definition.DefinitionVersion, definition.ExpectedExecutionMode, definitionResolution.CapabilityAdmissionPlan, executionRequest, @@ -223,7 +231,11 @@ await ValidateDefinitionArtifactAsync( return new WorkflowRunCreationReceipt( runActor.Id, definitionResolution.ActorId, - createdActorIds); + createdActorIds, + // Fix (review round 1, F2): + // EnsureRunAsync has a caller-supplied routable run identity. + // Populate RunId from the normalized request, not by reading the actor address back. + RunId: normalizedRunId); } catch { @@ -829,6 +841,8 @@ private static EventEnvelope CreateWorkflowRunBindEnvelope( string? runOrigin, string? scheduleId, string? workflowId, + string? revisionId, + long definitionVersion, ExternalCapabilityExecutionMode expectedExecutionMode, WorkflowCapabilityAdmissionPlan? capabilityAdmissionPlan) => new() @@ -845,6 +859,8 @@ private static EventEnvelope CreateWorkflowRunBindEnvelope( runOrigin, scheduleId, workflowId, + revisionId, + definitionVersion, expectedExecutionMode, capabilityAdmissionPlan)), Route = EnvelopeRouteSemantics.CreateTopologyPublication(WorkflowRunActorPortPublisherId, TopologyAudience.Self), @@ -864,6 +880,8 @@ private static EventEnvelope CreateWorkflowRunEnsureEnvelope( string? runOrigin, string? scheduleId, string? workflowId, + string? revisionId, + long definitionVersion, ExternalCapabilityExecutionMode expectedExecutionMode, WorkflowCapabilityAdmissionPlan? capabilityAdmissionPlan, WorkflowChatRequestEvent? executionRequest = null, @@ -885,6 +903,8 @@ private static EventEnvelope CreateWorkflowRunEnsureEnvelope( runOrigin, scheduleId, workflowId, + revisionId, + definitionVersion, expectedExecutionMode, capabilityAdmissionPlan), }; @@ -972,6 +992,8 @@ private static BindWorkflowRunDefinitionEvent BuildBindWorkflowRunDefinitionEven string? runOrigin, string? scheduleId, string? workflowId, + string? revisionId, + long definitionVersion, ExternalCapabilityExecutionMode expectedExecutionMode, WorkflowCapabilityAdmissionPlan? capabilityAdmissionPlan) { @@ -985,6 +1007,8 @@ private static BindWorkflowRunDefinitionEvent BuildBindWorkflowRunDefinitionEven RunOrigin = runOrigin?.Trim() ?? string.Empty, ScheduleId = scheduleId?.Trim() ?? string.Empty, WorkflowId = workflowId?.Trim() ?? string.Empty, + RevisionId = revisionId?.Trim() ?? string.Empty, + DefinitionVersion = Math.Max(0, definitionVersion), CapabilityAdmissionPlan = capabilityAdmissionPlan?.Clone(), ExpectedExecutionMode = expectedExecutionMode, }; diff --git a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs index d01f37fa8d..f096420ec1 100644 --- a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs +++ b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs @@ -2,9 +2,11 @@ using Aevatar.Workflow.Abstractions; using Aevatar.Workflow.Abstractions.Security; using Aevatar.Workflow.Core; +using Aevatar.Workflow.Core.Primitives; using Aevatar.Workflow.Projection.Observability; using Aevatar.Workflow.Projection.ReadModels; using Google.Protobuf.WellKnownTypes; +using YamlDotNet.Core; namespace Aevatar.Workflow.Projection.Projectors; @@ -14,7 +16,12 @@ public sealed class WorkflowExecutionCurrentStateProjector WorkflowRunState, WorkflowExecutionCurrentStateDocument> { + private const string FailedStatus = "failed"; + private const string CompletedStatus = "completed"; + private const string StoppedStatus = "stopped"; + private readonly WorkflowRunForkSeedReadModelMapper _forkSeedMapper = new(); + private readonly WorkflowParser _workflowParser = new(); public WorkflowExecutionCurrentStateProjector( IProjectionWriteDispatcher writeDispatcher, @@ -52,6 +59,8 @@ public WorkflowExecutionCurrentStateProjector( DefinitionActorId = state.DefinitionActorId ?? string.Empty, RunId = string.IsNullOrWhiteSpace(state.RunId) ? context.RootActorId : state.RunId, WorkflowId = state.WorkflowId ?? string.Empty, + RevisionId = state.RevisionId ?? string.Empty, + DefinitionVersion = Math.Max(0, state.DefinitionVersion), WorkflowName = state.WorkflowName ?? string.Empty, Status = ResolveCurrentStateStatus(state), ScopeId = state.ScopeId ?? string.Empty, @@ -89,6 +98,7 @@ public WorkflowExecutionCurrentStateProjector( StringComparer.Ordinal), ForkSeedCompletedStepIds = seedSnapshot.CompletedStepIds.ToList(), ForkSeedLastFailedStepId = seedSnapshot.LastFailedStepId, + RecoveryCapability = BuildRecoveryCapability(state, seedSnapshot), ForkSeedIdempotencies = seedSnapshot.IdempotencyByStepId.ToDictionary( x => x.Key, x => MapStepIdempotency(x.Value), @@ -359,4 +369,209 @@ private static WorkflowExecutionInputFileRefReadModel MapInputFileRef(WorkflowFi _ => null, }; } + + // Implement (issue #3251): + // Behavior: recovery availability is a typed read-model fact, not UI inference from failure text or actor ids. + // Why this shape: the projector has committed run state plus fork seed facts and can publish conservative capabilities. + private WorkflowRunRecoveryCapabilityReadModel BuildRecoveryCapability( + WorkflowRunState state, + WorkflowRunForkSeedProjectionSnapshot seedSnapshot) + { + var revisionId = state.RevisionId ?? string.Empty; + var definitionVersion = Math.Max(0, state.DefinitionVersion); + return new WorkflowRunRecoveryCapabilityReadModel + { + WorkflowDefinitionRevisionId = revisionId, + WorkflowDefinitionVersion = definitionVersion, + RetryFailedStep = BuildRetryFailedStepCapability(state, seedSnapshot), + RunAgain = BuildRunAgainCapability(state, seedSnapshot), + }; + } + + private WorkflowRecoveryActionCapabilityReadModel BuildRetryFailedStepCapability( + WorkflowRunState state, + WorkflowRunForkSeedProjectionSnapshot seedSnapshot) + { + var status = Normalize(state.Status); + var failureClass = ClassifyFailure(state); + if (failureClass.ReasonCode != WorkflowRecoveryUnavailableReasonCodeReadModel.None) + return Ineligible(failureClass.ReasonCode, failureClass.Reason, failureClass.RecommendedAction); + + if (!string.Equals(status, FailedStatus, StringComparison.OrdinalIgnoreCase)) + { + return Ineligible( + WorkflowRecoveryUnavailableReasonCodeReadModel.SourceRunNotTerminal, + "Retry from failed step is available only after a failed terminal run.", + WorkflowRecoveryRecommendedActionReadModel.TechnicalDetails); + } + + if (string.IsNullOrWhiteSpace(seedSnapshot.WorkflowYaml)) + { + return Unavailable( + WorkflowRecoveryUnavailableReasonCodeReadModel.MissingSourceFact, + "The source workflow definition is not available for retry.", + WorkflowRecoveryRecommendedActionReadModel.EditWorkflow); + } + + if (string.IsNullOrWhiteSpace(seedSnapshot.LastFailedStepId)) + { + return Unavailable( + WorkflowRecoveryUnavailableReasonCodeReadModel.LegacyUnavailable, + "The failed step was not materialized for this run.", + WorkflowRecoveryRecommendedActionReadModel.TechnicalDetails); + } + + return Eligible( + WorkflowRecoveryRecommendedActionReadModel.Retry, + seedSnapshot.LastFailedStepId, + reusesPriorStepOutputs: true, + mayIncurModelOrToolCost: true); + } + + private WorkflowRecoveryActionCapabilityReadModel BuildRunAgainCapability( + WorkflowRunState state, + WorkflowRunForkSeedProjectionSnapshot seedSnapshot) + { + var status = Normalize(state.Status); + var failureClass = ClassifyFailure(state); + if (failureClass.ReasonCode != WorkflowRecoveryUnavailableReasonCodeReadModel.None) + return Ineligible(failureClass.ReasonCode, failureClass.Reason, failureClass.RecommendedAction); + + if (status is not (FailedStatus or CompletedStatus or StoppedStatus)) + { + return Ineligible( + WorkflowRecoveryUnavailableReasonCodeReadModel.SourceRunNotTerminal, + "Run again is available only after the source run reaches a terminal state.", + WorkflowRecoveryRecommendedActionReadModel.TechnicalDetails); + } + + if (string.IsNullOrWhiteSpace(seedSnapshot.WorkflowYaml)) + { + return Unavailable( + WorkflowRecoveryUnavailableReasonCodeReadModel.WorkflowDefinitionUnavailable, + "The source workflow definition is not available for run again.", + WorkflowRecoveryRecommendedActionReadModel.EditWorkflow); + } + + if (!TryResolveEntryStepId(seedSnapshot.WorkflowYaml, out var entryStepId)) + { + return Unavailable( + WorkflowRecoveryUnavailableReasonCodeReadModel.WorkflowDefinitionUnavailable, + "The source workflow definition entry step is not available for run again.", + WorkflowRecoveryRecommendedActionReadModel.EditWorkflow); + } + + return Eligible( + WorkflowRecoveryRecommendedActionReadModel.RunAgain, + entryStepId, + reusesPriorStepOutputs: false, + mayIncurModelOrToolCost: true); + } + + private bool TryResolveEntryStepId(string workflowYaml, out string entryStepId) + { + try + { + entryStepId = _workflowParser.Parse(workflowYaml).EntryStepId?.Trim() ?? string.Empty; + return !string.IsNullOrWhiteSpace(entryStepId); + } + catch (InvalidOperationException) + { + entryStepId = string.Empty; + return false; + } + catch (YamlException) + { + entryStepId = string.Empty; + return false; + } + } + + private static WorkflowRecoveryActionCapabilityReadModel Eligible( + WorkflowRecoveryRecommendedActionReadModel recommendedAction, + string startingStepId, + bool reusesPriorStepOutputs, + bool mayIncurModelOrToolCost) + { + var capability = new WorkflowRecoveryActionCapabilityReadModel + { + Eligibility = WorkflowRecoveryEligibilityReadModel.Eligible, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCodeReadModel.None, + StartingStepId = WorkflowAuditTextSanitizer.Sanitize(startingStepId), + ReusesPriorStepOutputs = reusesPriorStepOutputs, + MayIncurModelOrToolCost = mayIncurModelOrToolCost, + }; + capability.RecommendedActions.Add(recommendedAction); + return capability; + } + + private static WorkflowRecoveryActionCapabilityReadModel Ineligible( + WorkflowRecoveryUnavailableReasonCodeReadModel reasonCode, + string unavailableReason, + WorkflowRecoveryRecommendedActionReadModel recommendedAction) => + NotEligible( + WorkflowRecoveryEligibilityReadModel.Ineligible, + reasonCode, + unavailableReason, + recommendedAction); + + private static WorkflowRecoveryActionCapabilityReadModel Unavailable( + WorkflowRecoveryUnavailableReasonCodeReadModel reasonCode, + string unavailableReason, + WorkflowRecoveryRecommendedActionReadModel recommendedAction) => + NotEligible( + WorkflowRecoveryEligibilityReadModel.Unavailable, + reasonCode, + unavailableReason, + recommendedAction); + + private static WorkflowRecoveryActionCapabilityReadModel NotEligible( + WorkflowRecoveryEligibilityReadModel eligibility, + WorkflowRecoveryUnavailableReasonCodeReadModel reasonCode, + string unavailableReason, + WorkflowRecoveryRecommendedActionReadModel recommendedAction) + { + var capability = new WorkflowRecoveryActionCapabilityReadModel + { + Eligibility = eligibility, + UnavailableReasonCode = reasonCode, + UnavailableReason = WorkflowAuditTextSanitizer.SanitizeForDisplay(unavailableReason, 240), + }; + capability.RecommendedActions.Add(recommendedAction); + return capability; + } + + private static WorkflowFailureClassification ClassifyFailure(WorkflowRunState state) + { + // Fix (review round 1, F1): + // Recovery classification previously inspected FinalError substrings. + // Map only the typed committed failure kind owned by the run state/event path. + return state.TerminalRecoveryFailureKind switch + { + WorkflowRecoveryFailureKind.AuthorizationFailure => new WorkflowFailureClassification( + WorkflowRecoveryUnavailableReasonCodeReadModel.AuthorizationFailure, + "Access must be fixed before this run can be recovered.", + WorkflowRecoveryRecommendedActionReadModel.FixAccess), + WorkflowRecoveryFailureKind.ConfigurationFailure => new WorkflowFailureClassification( + WorkflowRecoveryUnavailableReasonCodeReadModel.ConfigurationFailure, + "Configuration must be changed before this run can be recovered.", + WorkflowRecoveryRecommendedActionReadModel.ChangeConfiguration), + _ => WorkflowFailureClassification.None, + }; + } + + private static string Normalize(string? value) => + value?.Trim() ?? string.Empty; + + private readonly record struct WorkflowFailureClassification( + WorkflowRecoveryUnavailableReasonCodeReadModel ReasonCode, + string Reason, + WorkflowRecoveryRecommendedActionReadModel RecommendedAction) + { + public static WorkflowFailureClassification None => + new( + WorkflowRecoveryUnavailableReasonCodeReadModel.None, + string.Empty, + WorkflowRecoveryRecommendedActionReadModel.Unspecified); + } } diff --git a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs index 507e85e0d9..e457125e7b 100644 --- a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs +++ b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs @@ -40,12 +40,90 @@ public WorkflowActorSnapshot ToActorSnapshot(WorkflowExecutionCurrentStateDocume ActivityCurrentStep = MapActivityCurrentStep(source.ActivityCurrentStep), ActivityFirstFailure = MapActivityFirstFailure(source.ActivityFirstFailure), ActivityWaiting = MapActivityWaiting(source.ActivityWaiting), + RecoveryCapability = MapRecoveryCapability(source.RecoveryCapability), }; if (source.HasDurationMs) snapshot.DurationMs = source.DurationMs; return snapshot; } + private static WorkflowRunRecoveryCapability MapRecoveryCapability( + WorkflowRunRecoveryCapabilityReadModel? source) + { + if (source == null) + return new WorkflowRunRecoveryCapability(); + + return new WorkflowRunRecoveryCapability + { + WorkflowDefinitionRevisionId = source.WorkflowDefinitionRevisionId ?? string.Empty, + WorkflowDefinitionVersion = source.WorkflowDefinitionVersion, + RetryFailedStep = MapRecoveryActionCapability(source.RetryFailedStep), + RunAgain = MapRecoveryActionCapability(source.RunAgain), + }; + } + + private static WorkflowRecoveryActionCapability MapRecoveryActionCapability( + WorkflowRecoveryActionCapabilityReadModel? source) + { + if (source == null) + return new WorkflowRecoveryActionCapability + { + Eligibility = WorkflowRecoveryEligibility.Unavailable, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCode.LegacyUnavailable, + UnavailableReason = "Recovery capability is unavailable for this legacy run.", + }; + + var capability = new WorkflowRecoveryActionCapability + { + Eligibility = MapRecoveryEligibility(source.Eligibility), + UnavailableReasonCode = MapRecoveryUnavailableReasonCode(source.UnavailableReasonCode), + UnavailableReason = source.UnavailableReason ?? string.Empty, + StartingStepId = source.StartingStepId ?? string.Empty, + ReusesPriorStepOutputs = source.ReusesPriorStepOutputs, + MayIncurModelOrToolCost = source.MayIncurModelOrToolCost, + }; + capability.RecommendedActions.Add(source.RecommendedActions.Select(MapRecoveryRecommendedAction)); + return capability; + } + + private static WorkflowRecoveryEligibility MapRecoveryEligibility( + WorkflowRecoveryEligibilityReadModel value) => + value switch + { + WorkflowRecoveryEligibilityReadModel.Eligible => WorkflowRecoveryEligibility.Eligible, + WorkflowRecoveryEligibilityReadModel.Ineligible => WorkflowRecoveryEligibility.Ineligible, + WorkflowRecoveryEligibilityReadModel.Unavailable => WorkflowRecoveryEligibility.Unavailable, + _ => WorkflowRecoveryEligibility.Unspecified, + }; + + private static WorkflowRecoveryUnavailableReasonCode MapRecoveryUnavailableReasonCode( + WorkflowRecoveryUnavailableReasonCodeReadModel value) => + value switch + { + WorkflowRecoveryUnavailableReasonCodeReadModel.None => WorkflowRecoveryUnavailableReasonCode.None, + WorkflowRecoveryUnavailableReasonCodeReadModel.SourceRunNotTerminal => WorkflowRecoveryUnavailableReasonCode.SourceRunNotTerminal, + WorkflowRecoveryUnavailableReasonCodeReadModel.MissingSourceFact => WorkflowRecoveryUnavailableReasonCode.MissingSourceFact, + WorkflowRecoveryUnavailableReasonCodeReadModel.AuthorizationFailure => WorkflowRecoveryUnavailableReasonCode.AuthorizationFailure, + WorkflowRecoveryUnavailableReasonCodeReadModel.ConfigurationFailure => WorkflowRecoveryUnavailableReasonCode.ConfigurationFailure, + WorkflowRecoveryUnavailableReasonCodeReadModel.WorkflowDefinitionUnavailable => WorkflowRecoveryUnavailableReasonCode.WorkflowDefinitionUnavailable, + WorkflowRecoveryUnavailableReasonCodeReadModel.LegacyUnavailable => WorkflowRecoveryUnavailableReasonCode.LegacyUnavailable, + _ => WorkflowRecoveryUnavailableReasonCode.Unspecified, + }; + + private static WorkflowRecoveryRecommendedAction MapRecoveryRecommendedAction( + WorkflowRecoveryRecommendedActionReadModel value) => + value switch + { + WorkflowRecoveryRecommendedActionReadModel.Retry => WorkflowRecoveryRecommendedAction.Retry, + WorkflowRecoveryRecommendedActionReadModel.RunAgain => WorkflowRecoveryRecommendedAction.RunAgain, + WorkflowRecoveryRecommendedActionReadModel.FixAccess => WorkflowRecoveryRecommendedAction.FixAccess, + WorkflowRecoveryRecommendedActionReadModel.ChangeConfiguration => WorkflowRecoveryRecommendedAction.ChangeConfiguration, + WorkflowRecoveryRecommendedActionReadModel.EditWorkflow => WorkflowRecoveryRecommendedAction.EditWorkflow, + WorkflowRecoveryRecommendedActionReadModel.EditInput => WorkflowRecoveryRecommendedAction.EditInput, + WorkflowRecoveryRecommendedActionReadModel.TechnicalDetails => WorkflowRecoveryRecommendedAction.TechnicalDetails, + _ => WorkflowRecoveryRecommendedAction.Unspecified, + }; + public WorkflowActorProjectionState ToActorProjectionState(WorkflowExecutionCurrentStateDocument source) { return new WorkflowActorProjectionState diff --git a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs index 3a115c2fc6..eac817345f 100644 --- a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs +++ b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs @@ -32,7 +32,9 @@ public WorkflowRunForkSeedView ToSeedView(WorkflowExecutionCurrentStateDocument x => x.Key, x => ToView(x.Value), StringComparer.Ordinal), - source.CapabilityAdmissionPlan?.Clone()); + source.CapabilityAdmissionPlan?.Clone(), + source.RevisionId ?? string.Empty, + source.DefinitionVersion); } public WorkflowRunForkSeedProjectionSnapshot ToProjectionSnapshot(WorkflowRunState state) @@ -56,6 +58,8 @@ public WorkflowRunForkSeedProjectionSnapshot ToProjectionSnapshot(WorkflowRunSta completedStepIds, lastFailedStepId, state.ScopeId ?? string.Empty, + state.RevisionId ?? string.Empty, + state.DefinitionVersion, kernelState?.InputFileRefs.Select(static fileRef => fileRef.Clone()).ToList() ?? [], kernelState?.IdempotencyByStepId.ToDictionary( x => x.Key, @@ -112,5 +116,7 @@ public sealed record WorkflowRunForkSeedProjectionSnapshot( IReadOnlyList CompletedStepIds, string LastFailedStepId, string ScopeId, + string RevisionId, + long DefinitionVersion, IReadOnlyList InputFileRefs, IReadOnlyDictionary IdempotencyByStepId); diff --git a/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto b/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto index 57cdc05e9b..85194081b2 100644 --- a/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto +++ b/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto @@ -8,6 +8,52 @@ import "workflow_execution_messages.proto"; import "workflow_external_action.proto"; import "workflow_capability_admission.proto"; +enum WorkflowRecoveryEligibilityReadModel { + WORKFLOW_RECOVERY_ELIGIBILITY_READ_MODEL_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_ELIGIBILITY_READ_MODEL_ELIGIBLE = 1; + WORKFLOW_RECOVERY_ELIGIBILITY_READ_MODEL_INELIGIBLE = 2; + WORKFLOW_RECOVERY_ELIGIBILITY_READ_MODEL_UNAVAILABLE = 3; +} + +enum WorkflowRecoveryUnavailableReasonCodeReadModel { + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_NONE = 1; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_SOURCE_RUN_NOT_TERMINAL = 2; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_MISSING_SOURCE_FACT = 3; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_AUTHORIZATION_FAILURE = 4; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_CONFIGURATION_FAILURE = 5; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_WORKFLOW_DEFINITION_UNAVAILABLE = 6; + WORKFLOW_RECOVERY_UNAVAILABLE_REASON_CODE_READ_MODEL_LEGACY_UNAVAILABLE = 7; +} + +enum WorkflowRecoveryRecommendedActionReadModel { + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_UNSPECIFIED = 0; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_RETRY = 1; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_RUN_AGAIN = 2; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_FIX_ACCESS = 3; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_CHANGE_CONFIGURATION = 4; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_EDIT_WORKFLOW = 5; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_EDIT_INPUT = 6; + WORKFLOW_RECOVERY_RECOMMENDED_ACTION_READ_MODEL_TECHNICAL_DETAILS = 7; +} + +message WorkflowRecoveryActionCapabilityReadModel { + WorkflowRecoveryEligibilityReadModel eligibility = 1; + WorkflowRecoveryUnavailableReasonCodeReadModel unavailable_reason_code = 2; + string unavailable_reason = 3; + repeated WorkflowRecoveryRecommendedActionReadModel recommended_actions = 4; + string starting_step_id = 5; + bool reuses_prior_step_outputs = 6; + bool may_incur_model_or_tool_cost = 7; +} + +message WorkflowRunRecoveryCapabilityReadModel { + WorkflowRecoveryActionCapabilityReadModel retry_failed_step = 1; + WorkflowRecoveryActionCapabilityReadModel run_again = 2; + string workflow_definition_revision_id = 3; + int64 workflow_definition_version = 4; +} + message WorkflowRunInsightReportDocument { string id = 1; int64 state_version = 2; @@ -84,6 +130,9 @@ message WorkflowExecutionCurrentStateDocument { WorkflowRunActivityStepReadModel activity_current_step = 41; WorkflowRunActivityFailureReadModel activity_first_failure = 42; WorkflowRunActivityWaitingReadModel activity_waiting = 43; + WorkflowRunRecoveryCapabilityReadModel recovery_capability = 44; + string revision_id = 45; + int64 definition_version = 46; } message WorkflowRunActivityInitiatorReadModel { diff --git a/test/Aevatar.Workflow.Application.Tests/WorkflowForkRunCommandDispatchTests.cs b/test/Aevatar.Workflow.Application.Tests/WorkflowForkRunCommandDispatchTests.cs index 6a6b97f002..f1115b920d 100644 --- a/test/Aevatar.Workflow.Application.Tests/WorkflowForkRunCommandDispatchTests.cs +++ b/test/Aevatar.Workflow.Application.Tests/WorkflowForkRunCommandDispatchTests.cs @@ -228,7 +228,9 @@ public async Task DispatchAsync_HappyPath_ShouldCreateRunWithChosenYamlAndDispat { ["step-b"] = new("source-run", "step-b", 2, "source-run:step-b:2"), }, - scopeId: "scope-1"), + scopeId: "scope-1", + revisionId: "rev-source", + definitionVersion: 23), }; var runPort = new RecordingRunProvisioningPort(); var dispatchPort = new RecordingActorDispatchPort(); @@ -258,6 +260,7 @@ public async Task DispatchAsync_HappyPath_ShouldCreateRunWithChosenYamlAndDispat { SourceRunId = "source-run", NewRunActorId = "run-created", + NewRunId = "run-routable", WorkflowName = "edited", Accepted = true, CommandId = "cmd-1857", @@ -271,6 +274,8 @@ public async Task DispatchAsync_HappyPath_ShouldCreateRunWithChosenYamlAndDispat binding.WorkflowYaml.Should().Be(editedYaml); binding.ScopeId.Should().Be("scope-1"); binding.InlineWorkflowYamls.Should().Contain("child", childYaml); + binding.RevisionId.Should().Be("rev-source"); + binding.DefinitionVersion.Should().Be(23); dispatchPort.Dispatches.Should().ContainSingle(); dispatchPort.Dispatches.Single().ActorId.Should().Be("run-created"); @@ -394,7 +399,9 @@ private static WorkflowRunForkSeedView CreateSeedView( IReadOnlyDictionary? inlineWorkflowYamls = null, IReadOnlyDictionary? variables = null, string scopeId = "", - IReadOnlyDictionary? idempotencyByStepId = null) => + IReadOnlyDictionary? idempotencyByStepId = null, + string revisionId = "", + long definitionVersion = 0) => new WorkflowRunForkSeedView( SourceRunId: "source-run", Status: status, @@ -410,7 +417,9 @@ private static WorkflowRunForkSeedView CreateSeedView( LastFailedStepId: "step-b", FinalError: status.Equals("failed", StringComparison.OrdinalIgnoreCase) ? "boom" : string.Empty, ScopeId: scopeId, - IdempotencyByStepId: idempotencyByStepId ?? new Dictionary(StringComparer.Ordinal)); + IdempotencyByStepId: idempotencyByStepId ?? new Dictionary(StringComparer.Ordinal), + RevisionId: revisionId, + DefinitionVersion: definitionVersion); private static string WorkflowYaml(string name) => $$""" @@ -461,7 +470,8 @@ public Task CreateRunAsync( return Task.FromResult(new WorkflowRunCreationReceipt( "run-created", "definition-created", - ["definition-created", "run-created"])); + ["definition-created", "run-created"], + "run-routable")); } public Task DestroyAsync(string actorId, CancellationToken ct = default) diff --git a/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs b/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs index 94fdcd16f7..e9291ccc4f 100644 --- a/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs +++ b/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs @@ -145,6 +145,7 @@ public async Task ListActivityRunsForScopeAsync_ShouldReturnPagedTypedRows_FromC Prompt = "Approve?", Availability = "available", }; + snapshot.RecoveryCapability = RecoveryCapability(); var currentState = new FakeCurrentStateQueryPort { PageResult = new WorkflowActorCurrentStatePage([snapshot], "cursor-next", 42), @@ -191,6 +192,9 @@ public async Task ListActivityRunsForScopeAsync_ShouldReturnPagedTypedRows_FromC row.CompletedAtUtc.Should().Be(DateTimeOffset.UnixEpoch.AddSeconds(220)); row.UpdatedAtUtc.Should().Be(DateTimeOffset.UnixEpoch.AddSeconds(300)); row.DurationMs.Should().Be(120_000); + row.RecoveryCapability.WorkflowDefinitionRevisionId.Should().Be("rev-recovery"); + row.RecoveryCapability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibility.Eligible); + row.RecoveryCapability.RetryFailedStep.StartingStepId.Should().Be("step-failed"); } [Fact] @@ -564,9 +568,11 @@ public async Task GetRunForScopeAsync_ShouldSurfaceFinalResultStepsAndStatistics [Fact] public async Task GetRunForScopeAsync_ShouldFallBackToSummary_WhenReportNotYetMaterialized() { + var snapshot = Snapshot("run-1", CallerScope, WorkflowRunCompletionStatus.Running); + snapshot.RecoveryCapability = RecoveryCapability(); var currentState = new FakeCurrentStateQueryPort { - SingleResult = Snapshot("run-1", CallerScope, WorkflowRunCompletionStatus.Running), + SingleResult = snapshot, }; var service = new WorkflowRunObservatoryQueryService(currentState, new FakeArtifactQueryPort { Report = null }); @@ -580,6 +586,26 @@ public async Task GetRunForScopeAsync_ShouldFallBackToSummary_WhenReportNotYetMa detail.Steps.Should().BeEmpty(); detail.Statistics.TotalSteps.Should().Be(0); detail.Diagnostics.Should().BeEmpty(); + detail.RecoveryCapability.WorkflowDefinitionRevisionId.Should().Be("rev-recovery"); + detail.RecoveryCapability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibility.Eligible); + } + + [Fact] + public async Task GetRunForScopeAsync_ShouldExposeRecoveryCapability_WhenReportIsMaterialized() + { + var snapshot = Snapshot("run-1", CallerScope, WorkflowRunCompletionStatus.Failed); + snapshot.RecoveryCapability = RecoveryCapability(); + var service = new WorkflowRunObservatoryQueryService( + new FakeCurrentStateQueryPort { SingleResult = snapshot }, + new FakeArtifactQueryPort { Report = new WorkflowRunReport { FinalError = "boom" } }); + + var detail = await service.GetRunForScopeAsync(CallerScope, "run-1"); + + detail.Should().NotBeNull(); + detail!.RecoveryCapability.WorkflowDefinitionRevisionId.Should().Be("rev-recovery"); + detail.RecoveryCapability.WorkflowDefinitionVersion.Should().Be(12); + detail.RecoveryCapability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedAction.Retry); } [Fact] @@ -921,6 +947,33 @@ private static WorkflowActorSnapshot Snapshot( return snapshot; } + private static WorkflowRunRecoveryCapability RecoveryCapability() + { + var capability = new WorkflowRunRecoveryCapability + { + WorkflowDefinitionRevisionId = "rev-recovery", + WorkflowDefinitionVersion = 12, + RetryFailedStep = new WorkflowRecoveryActionCapability + { + Eligibility = WorkflowRecoveryEligibility.Eligible, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCode.None, + StartingStepId = "step-failed", + ReusesPriorStepOutputs = true, + MayIncurModelOrToolCost = true, + }, + RunAgain = new WorkflowRecoveryActionCapability + { + Eligibility = WorkflowRecoveryEligibility.Eligible, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCode.None, + StartingStepId = "step-a", + MayIncurModelOrToolCost = true, + }, + }; + capability.RetryFailedStep.RecommendedActions.Add(WorkflowRecoveryRecommendedAction.Retry); + capability.RunAgain.RecommendedActions.Add(WorkflowRecoveryRecommendedAction.RunAgain); + return capability; + } + private static WorkflowRunTimelineEvent TimelineEvent(string stage, string message, string? stepId = null) => new() { diff --git a/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs index 640c37e06f..3498c8f4a2 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/ChatEndpointsInternalTests.cs @@ -2551,7 +2551,8 @@ public async Task HandleForkRun_ShouldReturnAcceptedLocationAndDispatchMappedCom true, "cmd-1", "corr-1", - new DateTimeOffset(2026, 6, 8, 0, 0, 0, TimeSpan.Zero))), + new DateTimeOffset(2026, 6, 8, 0, 0, 0, TimeSpan.Zero), + "new-run-routable")), }; var result = await WorkflowCapabilityEndpoints.HandleForkRun( @@ -2580,11 +2581,12 @@ public async Task HandleForkRun_ShouldReturnAcceptedLocationAndDispatchMappedCom await result.ExecuteAsync(http); http.Response.StatusCode.Should().Be(StatusCodes.Status202Accepted); - http.Response.Headers.Location.ToString().Should().Be("/api/workflow-actors/new-run-actor/current-state"); + http.Response.Headers.Location.ToString().Should().Be("/api/workflow/observatory/runs/new-run-routable"); var body = await ReadBodyAsync(http.Response); + body.Should().Contain("\"newRunId\":\"new-run-routable\""); body.Should().Contain("\"newRunActorId\":\"new-run-actor\""); body.Should().Contain("\"acceptedCommandId\":\"cmd-1\""); - body.Should().Contain("\"statusUrl\":\"/api/workflow-actors/new-run-actor/current-state\""); + body.Should().Contain("\"statusUrl\":\"/api/workflow/observatory/runs/new-run-routable\""); service.Commands.Should().ContainSingle(); service.Commands.Single().SourceRunId.Should().Be("source-run"); service.Commands.Single().StartAtStepId.Should().Be("step-b"); diff --git a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowActivityRunFeedQueryPortTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowActivityRunFeedQueryPortTests.cs index 4af6d1233d..088bc7a648 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowActivityRunFeedQueryPortTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowActivityRunFeedQueryPortTests.cs @@ -1,5 +1,6 @@ using Aevatar.CQRS.Projection.Stores.Abstractions; using Aevatar.Workflow.Application.Abstractions.Projections; +using Aevatar.Workflow.Application.Abstractions.Queries; using Aevatar.Workflow.Projection.Configuration; using Aevatar.Workflow.Projection.Orchestration; using Aevatar.Workflow.Projection.ReadModels; @@ -67,6 +68,59 @@ public void WorkflowExecutionReadModelMapper_ShouldExposeActivityRunSummaryField snapshot.ActivityWaiting.WaitingKind.Should().Be("signal"); } + [Fact] + public void WorkflowExecutionReadModelMapper_ShouldExposeTypedRecoveryCapability() + { + var mapper = new WorkflowExecutionReadModelMapper(); + var snapshot = mapper.ToActorSnapshot(new WorkflowExecutionCurrentStateDocument + { + RootActorId = "actor-recovery", + RunId = "run-recovery", + Status = "failed", + RecoveryCapability = new WorkflowRunRecoveryCapabilityReadModel + { + WorkflowDefinitionRevisionId = "rev-recovery", + WorkflowDefinitionVersion = 12, + RetryFailedStep = new WorkflowRecoveryActionCapabilityReadModel + { + Eligibility = WorkflowRecoveryEligibilityReadModel.Eligible, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCodeReadModel.None, + StartingStepId = "step-failed", + ReusesPriorStepOutputs = true, + MayIncurModelOrToolCost = true, + RecommendedActions = + { + WorkflowRecoveryRecommendedActionReadModel.Retry, + }, + }, + RunAgain = new WorkflowRecoveryActionCapabilityReadModel + { + Eligibility = WorkflowRecoveryEligibilityReadModel.Ineligible, + UnavailableReasonCode = WorkflowRecoveryUnavailableReasonCodeReadModel.ConfigurationFailure, + UnavailableReason = "Configuration must be changed before this run can be recovered.", + RecommendedActions = + { + WorkflowRecoveryRecommendedActionReadModel.ChangeConfiguration, + }, + }, + }, + }); + + snapshot.RecoveryCapability.WorkflowDefinitionRevisionId.Should().Be("rev-recovery"); + snapshot.RecoveryCapability.WorkflowDefinitionVersion.Should().Be(12); + snapshot.RecoveryCapability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibility.Eligible); + snapshot.RecoveryCapability.RetryFailedStep.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCode.None); + snapshot.RecoveryCapability.RetryFailedStep.StartingStepId.Should().Be("step-failed"); + snapshot.RecoveryCapability.RetryFailedStep.ReusesPriorStepOutputs.Should().BeTrue(); + snapshot.RecoveryCapability.RetryFailedStep.MayIncurModelOrToolCost.Should().BeTrue(); + snapshot.RecoveryCapability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedAction.Retry); + snapshot.RecoveryCapability.RunAgain.Eligibility.Should().Be(WorkflowRecoveryEligibility.Ineligible); + snapshot.RecoveryCapability.RunAgain.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCode.ConfigurationFailure); + snapshot.RecoveryCapability.RunAgain.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedAction.ChangeConfiguration); + } + [Fact] public void WorkflowExecutionReadModelMapper_ShouldLeaveDurationUnavailable_WhenDocumentHasNoDuration() { diff --git a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionProjectionProjectorTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionProjectionProjectorTests.cs index c2eac40ff0..e2514c5a48 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionProjectionProjectorTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionProjectionProjectorTests.cs @@ -1129,6 +1129,160 @@ await projector.ProjectAsync( .CompletionStatus.Should().Be(WorkflowRunCompletionStatus.AwaitingToolApproval); } + [Fact] + public async Task WorkflowExecutionCurrentStateProjector_ShouldMaterializeEligibleRecoveryCapability() + { + var dispatcher = new RecordingWriteDispatcher(); + var projector = new WorkflowExecutionCurrentStateProjector( + dispatcher, + new FixedProjectionClock(DateTimeOffset.Parse("2026-08-07T03:00:00+00:00"))); + var state = new WorkflowRunState + { + RunId = "run-retry", + Status = "failed", + WorkflowYaml = CurrentStateWorkflowYaml("wf-retry"), + RevisionId = "rev-source", + DefinitionVersion = 17, + FinalError = "tool timeout", + }; + state.ExecutionStates["workflow_execution_kernel"] = Any.Pack(new WorkflowExecutionKernelState + { + CurrentStepId = "step-b", + Variables = + { + ["input"] = "original input", + ["step-a"] = "prior output", + }, + }); + + await projector.ProjectAsync( + CreateContext(), + WrapCommitted( + new WorkflowCompletedEvent { Success = false, Error = "tool timeout" }, + state)); + + var capability = dispatcher.Upserts.Should().ContainSingle().Subject.RecoveryCapability; + capability.WorkflowDefinitionRevisionId.Should().Be("rev-source"); + capability.WorkflowDefinitionVersion.Should().Be(17); + capability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Eligible); + capability.RetryFailedStep.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCodeReadModel.None); + capability.RetryFailedStep.StartingStepId.Should().Be("step-b"); + capability.RetryFailedStep.ReusesPriorStepOutputs.Should().BeTrue(); + capability.RetryFailedStep.MayIncurModelOrToolCost.Should().BeTrue(); + capability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedActionReadModel.Retry); + capability.RunAgain.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Eligible); + capability.RunAgain.StartingStepId.Should().Be("step-a"); + capability.RunAgain.ReusesPriorStepOutputs.Should().BeFalse(); + capability.RunAgain.MayIncurModelOrToolCost.Should().BeTrue(); + capability.RunAgain.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedActionReadModel.RunAgain); + } + + [Theory] + [InlineData(WorkflowRecoveryFailureKind.AuthorizationFailure, WorkflowRecoveryUnavailableReasonCodeReadModel.AuthorizationFailure, WorkflowRecoveryRecommendedActionReadModel.FixAccess)] + [InlineData(WorkflowRecoveryFailureKind.ConfigurationFailure, WorkflowRecoveryUnavailableReasonCodeReadModel.ConfigurationFailure, WorkflowRecoveryRecommendedActionReadModel.ChangeConfiguration)] + public async Task WorkflowExecutionCurrentStateProjector_ShouldNotAdvertiseRetryForTypedBackendClassifiedFailures( + WorkflowRecoveryFailureKind recoveryFailureKind, + WorkflowRecoveryUnavailableReasonCodeReadModel expectedReason, + WorkflowRecoveryRecommendedActionReadModel expectedAction) + { + var dispatcher = new RecordingWriteDispatcher(); + var projector = new WorkflowExecutionCurrentStateProjector( + dispatcher, + new FixedProjectionClock(DateTimeOffset.Parse("2026-08-07T03:10:00+00:00"))); + var state = new WorkflowRunState + { + RunId = "run-classified", + Status = "failed", + WorkflowYaml = CurrentStateWorkflowYaml("wf-classified"), + FinalError = "diagnostic text without recovery control meaning", + TerminalRecoveryFailureKind = recoveryFailureKind, + }; + state.ExecutionStates["workflow_execution_kernel"] = Any.Pack(new WorkflowExecutionKernelState + { + CurrentStepId = "step-b", + }); + + await projector.ProjectAsync( + CreateContext(), + WrapCommitted( + new WorkflowCompletedEvent + { + Success = false, + Error = "diagnostic text without recovery control meaning", + RecoveryFailureKind = recoveryFailureKind, + }, + state)); + + var capability = dispatcher.Upserts.Should().ContainSingle().Subject.RecoveryCapability; + capability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Ineligible); + capability.RetryFailedStep.UnavailableReasonCode.Should().Be(expectedReason); + capability.RetryFailedStep.UnavailableReason.Should().NotBeNullOrWhiteSpace(); + capability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(expectedAction); + capability.RunAgain.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Ineligible); + capability.RunAgain.UnavailableReasonCode.Should().Be(expectedReason); + capability.RunAgain.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(expectedAction); + } + + [Fact] + public async Task WorkflowExecutionCurrentStateProjector_ShouldExposeUnavailableRecoveryReasonsForMissingFacts() + { + var dispatcher = new RecordingWriteDispatcher(); + var projector = new WorkflowExecutionCurrentStateProjector( + dispatcher, + new FixedProjectionClock(DateTimeOffset.Parse("2026-08-07T03:20:00+00:00"))); + + await projector.ProjectAsync( + CreateContext(), + WrapCommitted( + new WorkflowCompletedEvent { Success = false, Error = "boom" }, + new WorkflowRunState + { + RunId = "run-missing-definition", + Status = "failed", + FinalError = "boom", + })); + + var capability = dispatcher.Upserts.Should().ContainSingle().Subject.RecoveryCapability; + capability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Unavailable); + capability.RetryFailedStep.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCodeReadModel.MissingSourceFact); + capability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedActionReadModel.EditWorkflow); + capability.RunAgain.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Unavailable); + capability.RunAgain.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCodeReadModel.WorkflowDefinitionUnavailable); + } + + [Fact] + public async Task WorkflowExecutionCurrentStateProjector_ShouldExposeLegacyUnavailableWhenFailedStepWasNotMaterialized() + { + var dispatcher = new RecordingWriteDispatcher(); + var projector = new WorkflowExecutionCurrentStateProjector( + dispatcher, + new FixedProjectionClock(DateTimeOffset.Parse("2026-08-07T03:30:00+00:00"))); + + await projector.ProjectAsync( + CreateContext(), + WrapCommitted( + new WorkflowCompletedEvent { Success = false, Error = "boom" }, + new WorkflowRunState + { + RunId = "run-legacy-failed-step", + Status = "failed", + WorkflowYaml = CurrentStateWorkflowYaml("wf-legacy"), + FinalError = "boom", + })); + + var capability = dispatcher.Upserts.Should().ContainSingle().Subject.RecoveryCapability; + capability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Unavailable); + capability.RetryFailedStep.UnavailableReasonCode.Should().Be(WorkflowRecoveryUnavailableReasonCodeReadModel.LegacyUnavailable); + capability.RetryFailedStep.RecommendedActions.Should().ContainSingle() + .Which.Should().Be(WorkflowRecoveryRecommendedActionReadModel.TechnicalDetails); + capability.RunAgain.Eligibility.Should().Be(WorkflowRecoveryEligibilityReadModel.Eligible); + } + [Fact] public async Task WorkflowExecutionCurrentStateProjector_WhenEnvelopeIsNotCommittedState_ShouldSkipWrite() { @@ -1570,6 +1724,17 @@ private static WorkflowExecutionMaterializationContext CreateContext() => ProjectionKind = "workflow", }; + private static string CurrentStateWorkflowYaml(string name) => + $$""" + name: {{name}} + roles: [] + steps: + - id: step-a + type: transform + - id: step-b + type: transform + """; + private static StateEvent PackStateEvent( IMessage payload, long version, diff --git a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowRunActorPortBranchTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowRunActorPortBranchTests.cs index 255edf0bc7..fece76de91 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowRunActorPortBranchTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowRunActorPortBranchTests.cs @@ -283,6 +283,34 @@ public async Task CreateRunAsync_WhenExistingDefinitionHasAdmissionPlan_ShouldRe definitionActor.LastHandledEnvelope.Should().BeNull(); } + [Fact] + public async Task CreateRunAsync_ShouldLeaveReceiptRunIdEmpty_WhenNoRoutableRunIdentityExists() + { + const string workflowYaml = "name: direct\nroles: []\nsteps: []\n"; + var runtime = new RecordingActorRuntime(); + var definitionAgent = new WorkflowGAgent(); + definitionAgent.State.WorkflowName = "direct"; + definitionAgent.State.WorkflowYaml = workflowYaml; + definitionAgent.State.CapabilityAdmissionPlan = CreateCapabilityAdmissionPlan(workflowYaml); + definitionAgent.State.ExpectedExecutionMode = ExternalCapabilityExecutionMode.Interactive; + runtime.StoredActors["definition-1"] = new RecordingActor("definition-1", definitionAgent); + runtime.ActorsToCreate.Enqueue(new RecordingActor("run-technical-address", new StubAgent("run"))); + var port = CreatePort(runtime); + + var result = await port.CreateRunAsync( + new WorkflowDefinitionBinding( + "definition-1", + "direct", + workflowYaml, + new Dictionary(StringComparer.OrdinalIgnoreCase), + ExternalCapabilityExecutionMode.Interactive), + CancellationToken.None); + + result.ActorId.Should().Be("run-technical-address"); + result.RunId.Should().BeEmpty(); + result.RunId.Should().NotBe(result.ActorId); + } + [Fact] public async Task EnsureDefinitionAsync_WhenExistingExplicitIdentityDiffers_ShouldRejectWithoutRebinding() {