diff --git a/docs/canon/workflow-runtime.md b/docs/canon/workflow-runtime.md index 026145133..0583cec16 100644 --- a/docs/canon/workflow-runtime.md +++ b/docs/canon/workflow-runtime.md @@ -67,6 +67,7 @@ owner: eanzhao - 通过依赖推导(`IWorkflowModuleDependencyExpander`)确定所需模块,经 `WorkflowModuleFactory` 创建并安装 - 收到 `ChatRequestEvent` envelope 后发布 `StartWorkflowEvent` - fork/resume-from-step seed 只走 request-level `WorkflowChatRequestEvent.fork_seed -> StartWorkflowEvent.fork_seed`;run bind 只表达 definition/run binding,不携带 seed。 + - run lineage 是 `WorkflowRunGAgent` owned committed fact,不从 route、actor id、graph/topology、workflow name 或 ID 前缀推断。`WorkflowRunLineage` 分离 retry/fork 与 `workflow_call` parent/child 关系,并始终使用可路由 public `runId`;actor address 只作为可选寻址信息保留。legacy 或未携带 lineage 的 run 必须显式返回 unavailable/legacy-unavailable。 - 由 `WorkflowExecutionKernel` 推进 `StepRequestEvent -> StepCompletedEvent -> WorkflowCompletedEvent` ``` diff --git a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto index 73c7b9791..a31a5549c 100644 --- a/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto +++ b/src/workflow/Aevatar.Workflow.Abstractions/workflow_execution_messages.proto @@ -270,6 +270,7 @@ message WorkflowRunForkSeed map variables = 3; int32 attempt = 4; WorkflowStepIdempotencyState start_step_idempotency = 5; + string original_run_id = 6; } message WorkflowStepIdempotencyState @@ -280,6 +281,60 @@ message WorkflowStepIdempotencyState string idempotency_key = 4; } +enum WorkflowRunLineageAvailability +{ + WORKFLOW_RUN_LINEAGE_AVAILABILITY_UNSPECIFIED = 0; + WORKFLOW_RUN_LINEAGE_AVAILABILITY_AVAILABLE = 1; + WORKFLOW_RUN_LINEAGE_AVAILABILITY_UNAVAILABLE = 2; + WORKFLOW_RUN_LINEAGE_AVAILABILITY_LEGACY_UNAVAILABLE = 3; +} + +enum WorkflowRunLineageRelationKind +{ + WORKFLOW_RUN_LINEAGE_RELATION_KIND_UNSPECIFIED = 0; + WORKFLOW_RUN_LINEAGE_RELATION_KIND_RETRY_FORK = 1; + WORKFLOW_RUN_LINEAGE_RELATION_KIND_SUB_WORKFLOW = 2; +} + +message WorkflowRunLineageRunRef +{ + string run_id = 1; + string actor_id = 2; + string relationship_id = 3; + string step_id = 4; + int32 attempt = 5; + WorkflowRunLineageRelationKind relation_kind = 6; +} + +message WorkflowRunRetryForkLineage +{ + WorkflowRunLineageAvailability availability = 1; + string source_run_id = 2; + string original_run_id = 3; + int32 attempt = 4; + string start_at_step_id = 5; + repeated WorkflowRunLineageRunRef child_runs = 6; +} + +message WorkflowRunSubWorkflowLineage +{ + WorkflowRunLineageAvailability availability = 1; + string parent_run_id = 2; + string parent_actor_id = 3; + string parent_step_id = 4; + string root_run_id = 5; + int32 depth = 6; + repeated WorkflowRunLineageRunRef child_runs = 7; +} + +message WorkflowRunLineage +{ + WorkflowRunLineageAvailability availability = 1; + WorkflowRunRetryForkLineage retry_fork = 2; + WorkflowRunSubWorkflowLineage sub_workflow = 3; + string unavailable_reason = 4; +} + message WorkflowRunForkRequestedEvent { string source_run_id = 1; @@ -288,6 +343,17 @@ message WorkflowRunForkRequestedEvent string scope_id = 4; } +message WorkflowRunLineageRecordedEvent +{ + string source_run_id = 1; + string child_run_id = 2; + string child_actor_id = 3; + string start_at_step_id = 4; + int32 attempt = 5; + WorkflowRunLineageRelationKind relation_kind = 6; + string original_run_id = 7; +} + message StartWorkflowEvent { string workflow_name = 1; @@ -324,6 +390,7 @@ message BindWorkflowRunDefinitionEvent string revision_id = 12; // Authoritative committed definition version/watermark used by this run. Zero means unavailable. int64 definition_version = 13; + WorkflowRunLineage initial_lineage = 14; } // Idempotent provisioning command for a stable Run actor identity. The first @@ -380,6 +447,7 @@ message WorkflowRunExecutionStartedEvent string workflow_command_id = 11; string workflow_correlation_id = 12; string current_turn_id = 13; + WorkflowRunLineage lineage = 14; } message WorkflowRunExecutionContextDelta { diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs index e27fc69e0..301a31332 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Observatory/IWorkflowRunObservatoryQueryService.cs @@ -1,3 +1,4 @@ +using Aevatar.Workflow.Abstractions; using Aevatar.Workflow.Application.Abstractions.Queries; namespace Aevatar.Workflow.Application.Abstractions.Observatory; @@ -122,6 +123,8 @@ public sealed class WorkflowActivityRunFeedRow public long StateVersion { get; init; } public WorkflowRunRecoveryCapability RecoveryCapability { get; init; } = new(); + + public WorkflowRunLineage Lineage { get; init; } = new(); } public sealed class WorkflowActivityRunInitiatorSummary @@ -196,6 +199,8 @@ public sealed class ObservatoryRunSummary // 06-23-observatory-run-coverage-filter: canonical run origin/type (draft | member-invoke | ...), // empty for legacy/unstamped runs. Drives the run-type filter + badge. public string RunOrigin { get; init; } = string.Empty; + + public WorkflowRunLineage Lineage { get; init; } = new(); } public sealed class ObservatoryRunDetail @@ -227,6 +232,8 @@ public sealed class ObservatoryRunDetail public ObservatoryUsageTotals UsageTotals { get; init; } = new(); public WorkflowRunRecoveryCapability RecoveryCapability { get; init; } = new(); + + public WorkflowRunLineage Lineage { 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 a0363c1fb..73fa2280d 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 @@ -89,6 +89,7 @@ message WorkflowActorSnapshot { WorkflowRunActivityWaitingSnapshot activity_waiting = 31; string run_id = 32; WorkflowRunRecoveryCapability recovery_capability = 33; + aevatar.workflow.WorkflowRunLineage lineage = 34; } 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 89ee3a0e3..17bbabbf5 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/RunForks/WorkflowForkRunModels.cs @@ -122,4 +122,5 @@ public sealed record WorkflowForkRunAcceptedReceipt( string CommandId, string CorrelationId, DateTimeOffset AckedAt, - string NewRunId = ""); + string NewRunId = "", + string OriginalRunId = ""); diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowChatRunModels.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowChatRunModels.cs index d2262ec35..0cad2b9c7 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowChatRunModels.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowChatRunModels.cs @@ -147,7 +147,8 @@ public sealed record WorkflowChatRunForkSeed( string StartAtStepId, IReadOnlyDictionary Variables, int Attempt = 0, - WorkflowStepIdempotencyView? StartStepIdempotency = null); + WorkflowStepIdempotencyView? StartStepIdempotency = null, + string OriginalRunId = ""); public enum WorkflowChatSourceKind { diff --git a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs index 9e7ba0c4d..357c07efe 100644 --- a/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs +++ b/src/workflow/Aevatar.Workflow.Application.Abstractions/Runs/WorkflowRunPorts.cs @@ -165,7 +165,8 @@ public sealed record WorkflowRunForkSeedView( IReadOnlyDictionary? IdempotencyByStepId = null, WorkflowCapabilityAdmissionPlan? CapabilityAdmissionPlan = null, string RevisionId = "", - long DefinitionVersion = 0) + long DefinitionVersion = 0, + string OriginalRunId = "") { public WorkflowRunForkSeedView() : this( @@ -182,7 +183,8 @@ public WorkflowRunForkSeedView() new Dictionary(StringComparer.Ordinal), null, string.Empty, - 0) + 0, + string.Empty) { } } diff --git a/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs b/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs index 658ce77fe..5e9debaa9 100644 --- a/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs +++ b/src/workflow/Aevatar.Workflow.Application/Observatory/WorkflowRunObservatoryQueryService.cs @@ -231,6 +231,7 @@ private async Task BuildRunDetailAsync( Diagnostics = BuildDiagnostics(snapshot, report: null, steps: [], viewEvents: []), UsageTotals = new ObservatoryUsageTotals(), RecoveryCapability = CloneRecoveryCapability(snapshot), + Lineage = CloneLineage(snapshot), }; } @@ -266,6 +267,7 @@ private async Task BuildRunDetailAsync( Statistics = ToStatistics(report.Summary), UsageTotals = WorkflowRunObservatoryTimelineMapper.ToUsageTotals(report.Usage), RecoveryCapability = CloneRecoveryCapability(snapshot), + Lineage = CloneLineage(snapshot), }; } @@ -378,6 +380,7 @@ private static ObservatoryRunSummary ToRunSummary(WorkflowActorSnapshot snapshot StateVersion = snapshot.StateVersion, ScopeId = snapshot.ScopeId, RunOrigin = snapshot.RunOrigin, + Lineage = CloneLineage(snapshot), }; } @@ -408,12 +411,28 @@ private static WorkflowActivityRunFeedRow ToActivityRunFeedRow(WorkflowActorSnap DurationMs = completedAtUtc == null || !snapshot.HasDurationMs ? null : snapshot.DurationMs, StateVersion = snapshot.StateVersion, RecoveryCapability = CloneRecoveryCapability(snapshot), + Lineage = CloneLineage(snapshot), }; } private static WorkflowRunRecoveryCapability CloneRecoveryCapability(WorkflowActorSnapshot snapshot) => snapshot.RecoveryCapability?.Clone() ?? new WorkflowRunRecoveryCapability(); + private static Aevatar.Workflow.Abstractions.WorkflowRunLineage CloneLineage(WorkflowActorSnapshot snapshot) => + snapshot.Lineage?.Clone() ?? new Aevatar.Workflow.Abstractions.WorkflowRunLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + UnavailableReason = "Run lineage is unavailable for this legacy run.", + RetryFork = new Aevatar.Workflow.Abstractions.WorkflowRunRetryForkLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + }, + SubWorkflow = new Aevatar.Workflow.Abstractions.WorkflowRunSubWorkflowLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + }, + }; + 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 79632d2f9..8644c7284 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunAcceptedReceiptFactory.cs @@ -21,6 +21,7 @@ public WorkflowForkRunAcceptedReceipt Create( context.CommandId, context.CorrelationId, DateTimeOffset.UtcNow, - target.RunId); + target.RunId, + target.OriginalRunId); } } diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs index 8990e0b57..db0974c31 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTarget.cs @@ -10,6 +10,7 @@ internal sealed class WorkflowForkRunCommandTarget : ICommandDispatchTarget, ICo public WorkflowForkRunCommandTarget( string sourceRunId, + string originalRunId, string startAtStepId, string actorId, string runId, @@ -19,6 +20,7 @@ public WorkflowForkRunCommandTarget( IWorkflowRunProvisioningPort runProvisioningPort) { SourceRunId = Normalize(sourceRunId); + OriginalRunId = Normalize(originalRunId); StartAtStepId = Normalize(startAtStepId); RunId = Normalize(runId); PreparedRequest = preparedRequest ?? throw new ArgumentNullException(nameof(preparedRequest)); @@ -31,6 +33,8 @@ public WorkflowForkRunCommandTarget( public string SourceRunId { get; } + public string OriginalRunId { get; } + public string StartAtStepId { get; } public string RunId { get; } diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs index 69bb734c5..6b89be7ce 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowForkRunCommandTargetResolver.cs @@ -105,6 +105,7 @@ public async Task !string.IsNullOrWhiteSpace(status) && TerminalStatuses.Contains(status.Trim()); diff --git a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowRunForkCoordinator.cs b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowRunForkCoordinator.cs index ccf4ffa77..7b7535330 100644 --- a/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowRunForkCoordinator.cs +++ b/src/workflow/Aevatar.Workflow.Application/RunForks/WorkflowRunForkCoordinator.cs @@ -1,7 +1,9 @@ using Aevatar.CQRS.Core.Abstractions.Commands; +using Aevatar.Foundation.Abstractions; using Aevatar.Foundation.Abstractions.EventSourcing; using Aevatar.Workflow.Abstractions; using Aevatar.Workflow.Application.Abstractions.RunForks; +using Google.Protobuf.WellKnownTypes; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; @@ -10,13 +12,16 @@ namespace Aevatar.Workflow.Application.RunForks; internal sealed class WorkflowRunForkCoordinator : ICommittedStatePublicationHook { private readonly Lazy> _forkDispatchService; + private readonly IActorDispatchPort? _dispatchPort; private readonly ILogger _logger; public WorkflowRunForkCoordinator( Lazy> forkDispatchService, + IActorDispatchPort? dispatchPort = null, ILogger? logger = null) { _forkDispatchService = forkDispatchService ?? throw new ArgumentNullException(nameof(forkDispatchService)); + _dispatchPort = dispatchPort; _logger = logger ?? NullLogger.Instance; } @@ -25,7 +30,8 @@ internal WorkflowRunForkCoordinator( ILogger? logger = null) : this(new Lazy>( () => forkDispatchService ?? throw new ArgumentNullException(nameof(forkDispatchService))), - logger) + dispatchPort: null, + logger: logger) { } @@ -68,7 +74,11 @@ public async Task BeforePublishAsync(CommittedStatePublicationContext context, C requested.Attempt, result.Error?.Code, result.Error?.Reason); + return; } + + await RecordAcceptedForkLineageAsync(context, requested, result.Receipt, ct) + .ConfigureAwait(false); } catch (Exception ex) { @@ -80,4 +90,48 @@ public async Task BeforePublishAsync(CommittedStatePublicationContext context, C requested.Attempt); } } + + private async Task RecordAcceptedForkLineageAsync( + CommittedStatePublicationContext context, + WorkflowRunForkRequestedEvent requested, + WorkflowForkRunAcceptedReceipt? receipt, + CancellationToken ct) + { + if (_dispatchPort == null || receipt == null || !receipt.Accepted) + return; + + var sourceActorId = context.ActorId?.Trim() ?? string.Empty; + var childRunId = receipt.NewRunId?.Trim() ?? string.Empty; + if (string.IsNullOrWhiteSpace(sourceActorId) || string.IsNullOrWhiteSpace(childRunId)) + return; + + // Implement (issue #3252): + // Behavior: source runs record accepted fork children with routable runId and separate child actor address. + // Why this shape: the coordinator only relays the accepted child identity back to the source actor; the source actor commits the lineage fact. + await _dispatchPort.DispatchAsync( + sourceActorId, + new EventEnvelope + { + Id = Guid.NewGuid().ToString("N"), + Timestamp = Timestamp.FromDateTimeOffset(DateTimeOffset.UtcNow), + Payload = Any.Pack(new WorkflowRunLineageRecordedEvent + { + SourceRunId = requested.SourceRunId ?? string.Empty, + ChildRunId = childRunId, + ChildActorId = receipt.NewRunActorId ?? string.Empty, + StartAtStepId = requested.StartAtStepId ?? string.Empty, + Attempt = Math.Max(0, requested.Attempt), + RelationKind = WorkflowRunLineageRelationKind.RetryFork, + OriginalRunId = string.IsNullOrWhiteSpace(receipt.OriginalRunId) + ? requested.SourceRunId ?? string.Empty + : receipt.OriginalRunId, + }), + Route = EnvelopeRouteSemantics.CreateTopologyPublication(sourceActorId, TopologyAudience.Self), + Propagation = new EnvelopePropagation + { + CorrelationId = Guid.NewGuid().ToString("N"), + }, + }, + ct).ConfigureAwait(false); + } } diff --git a/src/workflow/Aevatar.Workflow.Application/Runs/WorkflowChatRequestEnvelopeFactory.cs b/src/workflow/Aevatar.Workflow.Application/Runs/WorkflowChatRequestEnvelopeFactory.cs index 608767c6b..c56a54b44 100644 --- a/src/workflow/Aevatar.Workflow.Application/Runs/WorkflowChatRequestEnvelopeFactory.cs +++ b/src/workflow/Aevatar.Workflow.Application/Runs/WorkflowChatRequestEnvelopeFactory.cs @@ -192,6 +192,7 @@ private static Aevatar.Workflow.Abstractions.WorkflowRunForkSeed ToProto( SourceRunId = Normalize(source.SourceRunId), StartAtStepId = Normalize(source.StartAtStepId), Attempt = Math.Max(0, source.Attempt), + OriginalRunId = Normalize(source.OriginalRunId), }; if (source.StartStepIdempotency != null) { diff --git a/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs b/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs index e32c1770f..57e5ba28e 100644 --- a/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs +++ b/src/workflow/Aevatar.Workflow.Core/Primitives/SubWorkflowOrchestrator.cs @@ -743,6 +743,21 @@ public static WorkflowRunState ApplySubWorkflowInvocationRegistered(WorkflowRunS RemovePendingDefinitionResolution(next, invocationId); RemovePendingInvocation(next, invocationId, childRunId); AddPendingInvocation(next, pending); + // Implement (issue #3252): + // Behavior: parent workflow runs expose durable sub-workflow child run identities separately from retry/fork lineage. + // Why this shape: the committed invocation registration is the parent-owned fact; graph/topology edges remain execution-path data. + next.Lineage = WorkflowRunGAgent.EnsureLineage(next.Lineage); + next.Lineage.Availability = WorkflowRunLineageAvailability.Available; + next.Lineage.SubWorkflow ??= new WorkflowRunSubWorkflowLineage(); + next.Lineage.SubWorkflow.Availability = WorkflowRunLineageAvailability.Available; + WorkflowRunGAgent.UpsertLineageChild( + next.Lineage.SubWorkflow.ChildRuns, + childRunId, + pending.ChildActorId, + invocationId, + pending.ParentStepId, + attempt: 0, + WorkflowRunLineageRelationKind.SubWorkflow); return next; } @@ -919,7 +934,7 @@ await AdvancePendingSubWorkflowInvocationHandoffAsync( if (pending.HandoffPhase < SubWorkflowInvocationHandoffPhase.Bound) { - await BindSubWorkflowActorAsync(childActor.Id, definition, pending.ChildRunId, state, ct); + await BindSubWorkflowActorAsync(childActor.Id, definition, pending, state, ct); if (persistBinding) { await PersistBindingUpsertedAsync( @@ -1085,13 +1100,13 @@ private static SubWorkflowInvocationHandoffPhase ToSubWorkflowInvocationHandoffP private Task BindSubWorkflowActorAsync( string actorId, WorkflowDefinitionSnapshot definition, - string runId, + WorkflowRunState.Types.PendingSubWorkflowInvocation pending, WorkflowRunState state, CancellationToken ct) { return _dispatchPort.DispatchAsync( actorId, - CreateWorkflowRunBindEnvelope(definition, runId, state), + CreateWorkflowRunBindEnvelope(definition, pending, state), ct); } @@ -1466,7 +1481,7 @@ private async Task ResolveOrCreateWorkflowActorByIdAsync(string actorId) private EventEnvelope CreateWorkflowRunBindEnvelope( WorkflowDefinitionSnapshot definition, - string runId, + WorkflowRunState.Types.PendingSubWorkflowInvocation pending, WorkflowRunState state) { var inlineWorkflowYamls = new Dictionary(StringComparer.OrdinalIgnoreCase); @@ -1486,7 +1501,7 @@ private EventEnvelope CreateWorkflowRunBindEnvelope( DefinitionActorId = definition.DefinitionActorId ?? string.Empty, WorkflowYaml = definition.WorkflowYaml ?? string.Empty, WorkflowName = definition.WorkflowName ?? string.Empty, - RunId = runId ?? string.Empty, + RunId = pending.ChildRunId ?? string.Empty, ScopeId = string.IsNullOrWhiteSpace(definition.ScopeId) ? state.ScopeId ?? string.Empty : definition.ScopeId, @@ -1494,6 +1509,10 @@ private EventEnvelope CreateWorkflowRunBindEnvelope( RevisionId = definition.RevisionId ?? string.Empty, DefinitionVersion = Math.Max(0, definition.DefinitionVersion), ExpectedExecutionMode = state.ExpectedExecutionMode, + // Implement (issue #3252): + // Behavior: child workflow runs expose their parent/root run lineage as typed bind facts. + // Why this shape: the child actor commits lineage from the call-site handoff instead of deriving it from runtime topology. + InitialLineage = BuildChildInitialLineage(pending, _ownerActorIdAccessor()), }; return new EventEnvelope @@ -1509,6 +1528,30 @@ private EventEnvelope CreateWorkflowRunBindEnvelope( }; } + private static WorkflowRunLineage BuildChildInitialLineage( + WorkflowRunState.Types.PendingSubWorkflowInvocation pending, + string parentActorId) + { + var parentRunId = WorkflowRunIdNormalizer.Normalize(pending.ParentRunId); + return new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Available, + ParentRunId = parentRunId, + ParentActorId = parentActorId?.Trim() ?? string.Empty, + ParentStepId = pending.ParentStepId?.Trim() ?? string.Empty, + RootRunId = NormalizeRootRunId(pending.RootRunId, parentRunId), + Depth = Math.Max(0, pending.Depth), + }, + }; + } + private static bool ShouldEvictBinding( WorkflowRunState state, WorkflowRunState.Types.SubWorkflowBinding binding, diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs index c986024be..14d7e2607 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.IdentityProvisioning.cs @@ -30,7 +30,8 @@ await BindWorkflowRunDefinitionAsync( binding.RevisionId, binding.DefinitionVersion, binding.CapabilityAdmissionPlan, - binding.ExpectedExecutionMode); + binding.ExpectedExecutionMode, + binding.InitialLineage); } else { diff --git a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs index f5f070a6d..c018d5735 100644 --- a/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs +++ b/src/workflow/Aevatar.Workflow.Core/WorkflowRunGAgent.cs @@ -13,6 +13,7 @@ using Aevatar.Workflow.Core.Modules; using Aevatar.Workflow.Core.Primitives; using Google.Protobuf; +using Google.Protobuf.Collections; using Google.Protobuf.WellKnownTypes; using Microsoft.Extensions.Logging; using System.Diagnostics.CodeAnalysis; @@ -426,6 +427,7 @@ public async Task BindWorkflowRunDefinitionAsync( long definitionVersion, WorkflowCapabilityAdmissionPlan? capabilityAdmissionPlan, ExternalCapabilityExecutionMode expectedExecutionMode, + WorkflowRunLineage? initialLineage = null, CancellationToken ct = default) { if (expectedExecutionMode == ExternalCapabilityExecutionMode.Unspecified || @@ -459,6 +461,13 @@ public async Task BindWorkflowRunDefinitionAsync( CapabilityAdmissionPlan = capabilityAdmissionPlan?.Clone(), ExpectedExecutionMode = expectedExecutionMode, }; + if (initialLineage != null) + { + // Fix (review round 1, F1): + // Sub-workflow child lineage was stamped before dispatch but dropped before bind commit. + // Preserve InitialLineage in the authoritative bind event and cover the bind handler path. + bindDefinitionEvent.InitialLineage = initialLineage.Clone(); + } if (inlineWorkflowYamls != null) { foreach (var (key, value) in inlineWorkflowYamls) @@ -487,7 +496,8 @@ public Task HandleBindWorkflowRunDefinition(BindWorkflowRunDefinitionEvent reque request.RevisionId, request.DefinitionVersion, request.CapabilityAdmissionPlan, - request.ExpectedExecutionMode); + request.ExpectedExecutionMode, + request.InitialLineage); public override Task GetDescriptionAsync() { @@ -657,6 +667,7 @@ await HandleWorkflowCompleted(new WorkflowCompletedEvent CurrentTurnId = string.IsNullOrWhiteSpace(request.CurrentTurnId) ? request.ConversationContext?.CurrentTurnId ?? string.Empty : request.CurrentTurnId, + Lineage = BuildExecutionStartLineage(request.ForkSeed, State.Lineage, runId), }; if (request.CompletionNotificationTarget != null) executionStarted.CompletionNotificationTarget = request.CompletionNotificationTarget.Clone(); @@ -2035,6 +2046,7 @@ protected override WorkflowRunState TransitionState(WorkflowRunState current, IM .On(ApplyWorkflowRunExecutionContextCleared) .On(ApplyWorkflowExecutionStateUpserted) .On(ApplyWorkflowExecutionStateCleared) + .On(ApplyWorkflowRunLineageRecorded) .On(ApplyCompensableStepDispatched) .On(ApplyStepCompleted) .On(ApplyCompensationRequest) @@ -2226,6 +2238,7 @@ private WorkflowRunState ApplyBindWorkflowRunDefinition(WorkflowRunState current next.TerminalNotificationAttempt = 0; next.TerminalNotificationDeliveryStatus = WorkflowRunTerminalNotificationDeliveryStatus.Unspecified; next.TerminalNotificationRetryCallbackId = string.Empty; + next.Lineage = evt.InitialLineage?.Clone() ?? CreateUnavailableLineage("Run lineage is unavailable for this run."); next.InlineWorkflowYamls.Clear(); foreach (var (workflowNameKey, workflowYamlValue) in evt.InlineWorkflowYamls) { @@ -2312,6 +2325,7 @@ private static WorkflowRunState ApplyWorkflowRunExecutionStarted(WorkflowRunStat next.TerminalNotificationDeliveryStatus = WorkflowRunTerminalNotificationDeliveryStatus.Unspecified; next.TerminalNotificationRetryCallbackId = string.Empty; next.CurrentTurnId = evt.CurrentTurnId?.Trim() ?? string.Empty; + next.Lineage = evt.Lineage?.Clone() ?? current.Lineage?.Clone() ?? CreateUnavailableLineage("Run lineage is unavailable for this run."); next.InteractiveActionHandoffs.Clear(); next.ExecutionContext ??= new WorkflowRunExecutionContextState(); ApplyExecutionContextDelta(next.ExecutionContext, evt.ExecutionContextDelta); @@ -2325,6 +2339,198 @@ private static WorkflowRunState ApplyWorkflowRunExecutionStarted(WorkflowRunStat return next; } + [EventHandler] + public async Task HandleWorkflowRunLineageRecorded(WorkflowRunLineageRecordedEvent evt) + { + ArgumentNullException.ThrowIfNull(evt); + + var sourceRunId = WorkflowRunIdNormalizer.Normalize(evt.SourceRunId); + var currentRunId = WorkflowRunIdNormalizer.Normalize(RunId); + if (string.IsNullOrWhiteSpace(sourceRunId) || + string.IsNullOrWhiteSpace(currentRunId) || + !string.Equals(sourceRunId, currentRunId, StringComparison.Ordinal)) + { + Logger.LogWarning( + "Reject workflow lineage record for mismatched source run. actor={ActorId} currentRun={CurrentRunId} sourceRun={SourceRunId}", + Id, + currentRunId, + sourceRunId); + return; + } + + if (evt.RelationKind != WorkflowRunLineageRelationKind.RetryFork || + string.IsNullOrWhiteSpace(evt.ChildRunId)) + { + Logger.LogWarning( + "Reject workflow lineage record with unsupported or missing relation. sourceRun={SourceRunId} childRun={ChildRunId} relation={RelationKind}", + sourceRunId, + evt.ChildRunId, + evt.RelationKind); + return; + } + + // Implement (issue #3252): + // Behavior: source/original runs expose retried or forked child run identities from committed facts. + // Why this shape: the source actor records its own lineage instead of queries deriving children from actor IDs or graph topology. + await PersistDomainEventAsync(new WorkflowRunLineageRecordedEvent + { + SourceRunId = sourceRunId, + ChildRunId = WorkflowRunIdNormalizer.Normalize(evt.ChildRunId), + ChildActorId = evt.ChildActorId?.Trim() ?? string.Empty, + StartAtStepId = evt.StartAtStepId?.Trim() ?? string.Empty, + Attempt = Math.Max(0, evt.Attempt), + RelationKind = WorkflowRunLineageRelationKind.RetryFork, + OriginalRunId = WorkflowRunIdNormalizer.Normalize(evt.OriginalRunId), + }); + } + + private static WorkflowRunState ApplyWorkflowRunLineageRecorded( + WorkflowRunState current, + WorkflowRunLineageRecordedEvent evt) + { + if (evt.RelationKind != WorkflowRunLineageRelationKind.RetryFork) + return current; + + var childRunId = WorkflowRunIdNormalizer.Normalize(evt.ChildRunId); + if (string.IsNullOrWhiteSpace(childRunId)) + return current; + + var next = current.Clone(); + next.Lineage = EnsureLineage(next.Lineage); + next.Lineage.Availability = WorkflowRunLineageAvailability.Available; + next.Lineage.RetryFork ??= new WorkflowRunRetryForkLineage(); + next.Lineage.RetryFork.Availability = WorkflowRunLineageAvailability.Available; + if (string.IsNullOrWhiteSpace(next.Lineage.RetryFork.SourceRunId)) + next.Lineage.RetryFork.SourceRunId = WorkflowRunIdNormalizer.Normalize(evt.SourceRunId); + if (string.IsNullOrWhiteSpace(next.Lineage.RetryFork.OriginalRunId)) + { + var originalRunId = WorkflowRunIdNormalizer.Normalize(evt.OriginalRunId); + next.Lineage.RetryFork.OriginalRunId = string.IsNullOrWhiteSpace(originalRunId) + ? next.Lineage.RetryFork.SourceRunId + : originalRunId; + } + + UpsertLineageChild( + next.Lineage.RetryFork.ChildRuns, + childRunId, + evt.ChildActorId, + relationshipId: string.Empty, + stepId: evt.StartAtStepId, + Math.Max(0, evt.Attempt), + WorkflowRunLineageRelationKind.RetryFork); + return next; + } + + internal static WorkflowRunLineage BuildExecutionStartLineage( + WorkflowRunForkSeed? forkSeed, + WorkflowRunLineage? currentLineage, + string runId) + { + if (forkSeed == null || string.IsNullOrWhiteSpace(forkSeed.SourceRunId)) + return currentLineage?.Clone() ?? CreateUnavailableLineage("Run lineage is unavailable for this run."); + + var sourceRunId = WorkflowRunIdNormalizer.Normalize(forkSeed.SourceRunId); + var originalRunId = WorkflowRunIdNormalizer.Normalize(forkSeed.OriginalRunId); + if (string.IsNullOrWhiteSpace(originalRunId)) + originalRunId = sourceRunId; + + // Implement (issue #3252): + // Behavior: retried and forked child runs carry source and original run IDs as typed lineage. + // Why this shape: fork lineage is stamped from the accepted fork seed instead of inferred from routes or run ID strings. + return new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Available, + SourceRunId = sourceRunId, + OriginalRunId = originalRunId, + Attempt = Math.Max(0, forkSeed.Attempt), + StartAtStepId = forkSeed.StartAtStepId?.Trim() ?? string.Empty, + }, + SubWorkflow = currentLineage?.SubWorkflow?.Clone() ?? new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + }; + } + + internal static WorkflowRunLineage CreateUnavailableLineage(string reason) => + new() + { + Availability = WorkflowRunLineageAvailability.Unavailable, + UnavailableReason = string.IsNullOrWhiteSpace(reason) + ? "Run lineage is unavailable for this run." + : reason, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + }; + + internal static WorkflowRunLineage EnsureLineage(WorkflowRunLineage? lineage) + { + var next = lineage?.Clone() ?? CreateUnavailableLineage("Run lineage is unavailable for this run."); + next.RetryFork ??= new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }; + next.SubWorkflow ??= new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }; + return next; + } + + internal static void UpsertLineageChild( + RepeatedField childRuns, + string runId, + string? actorId, + string relationshipId, + string? stepId, + int attempt, + WorkflowRunLineageRelationKind relationKind) + { + var normalizedRunId = WorkflowRunIdNormalizer.Normalize(runId); + if (string.IsNullOrWhiteSpace(normalizedRunId)) + return; + + var normalizedActorId = actorId?.Trim() ?? string.Empty; + for (var i = 0; i < childRuns.Count; i++) + { + if (!string.Equals(childRuns[i].RunId, normalizedRunId, StringComparison.Ordinal) || + childRuns[i].RelationKind != relationKind) + { + continue; + } + + childRuns[i] = new WorkflowRunLineageRunRef + { + RunId = normalizedRunId, + ActorId = string.IsNullOrWhiteSpace(normalizedActorId) ? childRuns[i].ActorId ?? string.Empty : normalizedActorId, + RelationshipId = string.IsNullOrWhiteSpace(relationshipId) ? childRuns[i].RelationshipId : relationshipId, + StepId = string.IsNullOrWhiteSpace(stepId) ? childRuns[i].StepId : stepId.Trim(), + Attempt = Math.Max(0, attempt), + RelationKind = relationKind, + }; + return; + } + + childRuns.Add(new WorkflowRunLineageRunRef + { + RunId = normalizedRunId, + ActorId = normalizedActorId, + RelationshipId = relationshipId ?? string.Empty, + StepId = stepId?.Trim() ?? string.Empty, + Attempt = Math.Max(0, attempt), + RelationKind = relationKind, + }); + } + private static WorkflowRunState ApplyWorkflowRunExecutionContextUpdated( WorkflowRunState current, WorkflowRunExecutionContextUpdatedEvent evt) diff --git a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto index a81dd5767..b58b7a4ab 100644 --- a/src/workflow/Aevatar.Workflow.Core/workflow_state.proto +++ b/src/workflow/Aevatar.Workflow.Core/workflow_state.proto @@ -202,6 +202,7 @@ message WorkflowRunState string revision_id = 52; int64 definition_version = 53; aevatar.workflow.WorkflowRecoveryFailureKind terminal_recovery_failure_kind = 54; + aevatar.workflow.WorkflowRunLineage lineage = 55; } message WorkflowRunInitiatorState diff --git a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs index f096420ec..3d6a2ca8a 100644 --- a/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs +++ b/src/workflow/Aevatar.Workflow.Projection/Projectors/WorkflowExecutionCurrentStateProjector.cs @@ -99,6 +99,7 @@ public WorkflowExecutionCurrentStateProjector( ForkSeedCompletedStepIds = seedSnapshot.CompletedStepIds.ToList(), ForkSeedLastFailedStepId = seedSnapshot.LastFailedStepId, RecoveryCapability = BuildRecoveryCapability(state, seedSnapshot), + Lineage = MapLineage(state.Lineage), ForkSeedIdempotencies = seedSnapshot.IdempotencyByStepId.ToDictionary( x => x.Key, x => MapStepIdempotency(x.Value), @@ -151,6 +152,21 @@ private static WorkflowRunActivityInitiatorReadModel MapInitiator(WorkflowRunIni }; } + private static WorkflowRunLineage MapLineage(WorkflowRunLineage? lineage) => + lineage?.Clone() ?? new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + UnavailableReason = "Run lineage is unavailable for this run.", + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + }; + private static WorkflowRunActivityStepReadModel ResolveCurrentStep(WorkflowRunState state) { foreach (var executionState in state.ExecutionStates.Values) diff --git a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs index e457125e7..439287497 100644 --- a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs +++ b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowExecutionReadModelMapper.cs @@ -41,6 +41,7 @@ public WorkflowActorSnapshot ToActorSnapshot(WorkflowExecutionCurrentStateDocume ActivityFirstFailure = MapActivityFirstFailure(source.ActivityFirstFailure), ActivityWaiting = MapActivityWaiting(source.ActivityWaiting), RecoveryCapability = MapRecoveryCapability(source.RecoveryCapability), + Lineage = MapLineage(source.Lineage), }; if (source.HasDurationMs) snapshot.DurationMs = source.DurationMs; @@ -62,6 +63,22 @@ private static WorkflowRunRecoveryCapability MapRecoveryCapability( }; } + private static Aevatar.Workflow.Abstractions.WorkflowRunLineage MapLineage( + Aevatar.Workflow.Abstractions.WorkflowRunLineage? source) => + source?.Clone() ?? new Aevatar.Workflow.Abstractions.WorkflowRunLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + UnavailableReason = "Run lineage is unavailable for this legacy run.", + RetryFork = new Aevatar.Workflow.Abstractions.WorkflowRunRetryForkLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + }, + SubWorkflow = new Aevatar.Workflow.Abstractions.WorkflowRunSubWorkflowLineage + { + Availability = Aevatar.Workflow.Abstractions.WorkflowRunLineageAvailability.LegacyUnavailable, + }, + }; + private static WorkflowRecoveryActionCapability MapRecoveryActionCapability( WorkflowRecoveryActionCapabilityReadModel? source) { diff --git a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs index eac817345..a9de37a99 100644 --- a/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs +++ b/src/workflow/Aevatar.Workflow.Projection/ReadModels/WorkflowRunForkSeedReadModelMapper.cs @@ -34,7 +34,8 @@ public WorkflowRunForkSeedView ToSeedView(WorkflowExecutionCurrentStateDocument StringComparer.Ordinal), source.CapabilityAdmissionPlan?.Clone(), source.RevisionId ?? string.Empty, - source.DefinitionVersion); + source.DefinitionVersion, + ResolveOriginalRunId(source.Lineage, source.RunId)); } public WorkflowRunForkSeedProjectionSnapshot ToProjectionSnapshot(WorkflowRunState state) @@ -106,6 +107,17 @@ private static WorkflowStepIdempotencyView ToView(WorkflowStepIdempotencyReadMod source.StepId ?? string.Empty, source.LogicalAttempt, source.IdempotencyKey ?? string.Empty); + + private static string ResolveOriginalRunId( + Aevatar.Workflow.Abstractions.WorkflowRunLineage? lineage, + string? runId) + { + var originalRunId = lineage?.RetryFork?.OriginalRunId?.Trim() ?? string.Empty; + if (!string.IsNullOrWhiteSpace(originalRunId)) + return originalRunId; + + return runId?.Trim() ?? string.Empty; + } } public sealed record WorkflowRunForkSeedProjectionSnapshot( diff --git a/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto b/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto index 85194081b..c41a6009f 100644 --- a/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto +++ b/src/workflow/Aevatar.Workflow.Projection/workflow_projection_transport.proto @@ -133,6 +133,7 @@ message WorkflowExecutionCurrentStateDocument { WorkflowRunRecoveryCapabilityReadModel recovery_capability = 44; string revision_id = 45; int64 definition_version = 46; + aevatar.workflow.WorkflowRunLineage lineage = 47; } message WorkflowRunActivityInitiatorReadModel { diff --git a/test/Aevatar.GAgents.ChannelRuntime.Tests/ChannelWorkflowResultDeliveryContractTests.cs b/test/Aevatar.GAgents.ChannelRuntime.Tests/ChannelWorkflowResultDeliveryContractTests.cs index 9daeb2c84..1246fedeb 100644 --- a/test/Aevatar.GAgents.ChannelRuntime.Tests/ChannelWorkflowResultDeliveryContractTests.cs +++ b/test/Aevatar.GAgents.ChannelRuntime.Tests/ChannelWorkflowResultDeliveryContractTests.cs @@ -600,6 +600,8 @@ await Agent.BindWorkflowRunDefinitionAsync( runOrigin: null, scheduleId: null, workflowId: null, + revisionId: null, + definitionVersion: 0, capabilityAdmissionPlan: null, expectedExecutionMode: command.ExpectedExecutionMode, ct: ct); diff --git a/test/Aevatar.Integration.Tests/TestDoubles/WorkflowGAgentTestBase.cs b/test/Aevatar.Integration.Tests/TestDoubles/WorkflowGAgentTestBase.cs index 25ea79cb1..eec1a26bd 100644 --- a/test/Aevatar.Integration.Tests/TestDoubles/WorkflowGAgentTestBase.cs +++ b/test/Aevatar.Integration.Tests/TestDoubles/WorkflowGAgentTestBase.cs @@ -102,9 +102,11 @@ internal static Task BindInteractiveWorkflowRunDefinitionAsync( runOrigin, scheduleId, workflowId: null, + revisionId: null, + definitionVersion: 0, capabilityAdmissionPlan, ExternalCapabilityExecutionMode.Interactive, - ct); + ct: ct); internal static async Task CreateRegisteredDefinitionAgentAsync( RecordingActorRuntime runtime, diff --git a/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkCoordinatorTests.cs b/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkCoordinatorTests.cs index 5642c7d42..af7193d27 100644 --- a/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkCoordinatorTests.cs +++ b/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkCoordinatorTests.cs @@ -49,10 +49,54 @@ public async Task BeforePublishAsync_WhenCommittedForkRequestedEvent_ShouldDispa command.ScopeId.Should().Be("scope-1"); } - private static CommittedStatePublicationContext CreateContext(IMessage evt) => + [Fact] + public async Task BeforePublishAsync_WhenForkAccepted_ShouldDispatchLineageRecordToSourceActor() + { + var forkDispatchService = new RecordingForkDispatchService + { + Receipt = new WorkflowForkRunAcceptedReceipt( + "run-source-gamma", + "actor-child-delta", + "wf", + true, + "cmd", + "corr", + DateTimeOffset.UtcNow, + "run-child-beta", + "run-original-alpha"), + }; + var dispatchPort = new RecordingActorDispatchPort(); + var coordinator = new WorkflowRunForkCoordinator( + new Lazy>( + () => forkDispatchService), + dispatchPort); + + await coordinator.BeforePublishAsync( + CreateContext(new WorkflowRunForkRequestedEvent + { + SourceRunId = "run-source-gamma", + StartAtStepId = "step-retry", + Attempt = 4, + ScopeId = "scope-alpha", + }, actorId: "actor-source-epsilon"), + CancellationToken.None); + + var dispatched = dispatchPort.Dispatched.Should().ContainSingle().Subject; + dispatched.ActorId.Should().Be("actor-source-epsilon"); + var recorded = dispatched.Envelope.Payload.Unpack(); + recorded.SourceRunId.Should().Be("run-source-gamma"); + recorded.ChildRunId.Should().Be("run-child-beta"); + recorded.ChildActorId.Should().Be("actor-child-delta"); + recorded.OriginalRunId.Should().Be("run-original-alpha"); + recorded.StartAtStepId.Should().Be("step-retry"); + recorded.Attempt.Should().Be(4); + recorded.RelationKind.Should().Be(WorkflowRunLineageRelationKind.RetryFork); + } + + private static CommittedStatePublicationContext CreateContext(IMessage evt, string actorId = "run-source") => new() { - ActorId = "run-source", + ActorId = actorId, ActorType = typeof(object), Published = new CommittedStateEventPublished { @@ -73,6 +117,8 @@ private sealed class RecordingForkDispatchService { public List Commands { get; } = []; + public WorkflowForkRunAcceptedReceipt? Receipt { get; init; } + public Task> DispatchAsync( WorkflowForkRunCommand command, CancellationToken ct = default) @@ -80,7 +126,7 @@ public Task.Success( - new WorkflowForkRunAcceptedReceipt( + Receipt ?? new WorkflowForkRunAcceptedReceipt( command.SourceRunId, "new-run", "wf", @@ -90,4 +136,18 @@ public Task Dispatched { get; } = []; + + public Task DispatchAsync(string actorId, EventEnvelope envelope, CancellationToken ct = default) + { + ct.ThrowIfCancellationRequested(); + Dispatched.Add(new RecordedDispatch(actorId, envelope)); + return Task.FromResult(DispatchAdmissionFactory.Create(actorId, envelope)); + } + } + + private sealed record RecordedDispatch(string ActorId, EventEnvelope Envelope); } diff --git a/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkSeedQueryPortTests.cs b/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkSeedQueryPortTests.cs index d3c390a50..526ccc1d6 100644 --- a/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkSeedQueryPortTests.cs +++ b/test/Aevatar.Workflow.Application.Tests/WorkflowRunForkSeedQueryPortTests.cs @@ -60,6 +60,32 @@ public void ForkSeedReadModelMapper_ShouldMapCompletedRunForkSeed() view.IdempotencyByStepId!["step-b"].IdempotencyKey.Should().Be("run-completed:step-b:1"); } + [Fact] + public void ForkSeedReadModelMapper_ShouldCarryOriginalRunIdFromLineage() + { + var mapper = new WorkflowRunForkSeedReadModelMapper(); + + var view = mapper.ToSeedView(new WorkflowExecutionCurrentStateDocument + { + RunId = "run-source-gamma", + Status = "failed", + ScopeId = "scope-alpha", + Lineage = new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Available, + SourceRunId = "run-source-gamma", + OriginalRunId = "run-original-alpha", + }, + }, + }); + + view.SourceRunId.Should().Be("run-source-gamma"); + view.OriginalRunId.Should().Be("run-original-alpha"); + } + [Fact] public async Task GetForkSeedAsync_ShouldReadFailedRunForkSeedThroughCurrentStateReadModel() { diff --git a/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs b/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs index e9291ccc4..373e33e3e 100644 --- a/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs +++ b/test/Aevatar.Workflow.Application.Tests/WorkflowRunObservatoryQueryServiceTests.cs @@ -146,6 +146,22 @@ public async Task ListActivityRunsForScopeAsync_ShouldReturnPagedTypedRows_FromC Availability = "available", }; snapshot.RecoveryCapability = RecoveryCapability(); + snapshot.Lineage = new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Available, + SourceRunId = "run-source-gamma", + OriginalRunId = "run-original-alpha", + Attempt = 2, + StartAtStepId = "step-failed", + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + }; var currentState = new FakeCurrentStateQueryPort { PageResult = new WorkflowActorCurrentStatePage([snapshot], "cursor-next", 42), @@ -195,6 +211,10 @@ public async Task ListActivityRunsForScopeAsync_ShouldReturnPagedTypedRows_FromC row.RecoveryCapability.WorkflowDefinitionRevisionId.Should().Be("rev-recovery"); row.RecoveryCapability.RetryFailedStep.Eligibility.Should().Be(WorkflowRecoveryEligibility.Eligible); row.RecoveryCapability.RetryFailedStep.StartingStepId.Should().Be("step-failed"); + row.Lineage.RetryFork.SourceRunId.Should().Be("run-source-gamma"); + row.Lineage.RetryFork.OriginalRunId.Should().Be("run-original-alpha"); + row.Lineage.RetryFork.StartAtStepId.Should().Be("step-failed"); + row.Lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); } [Fact] @@ -225,6 +245,7 @@ public async Task ListActivityRunsForScopeAsync_ShouldRepresentMissingFactsExpli row.Waiting.Availability.Should().Be("unavailable"); row.CompletedAtUtc.Should().BeNull(); row.DurationMs.Should().BeNull(); + row.Lineage.Availability.Should().Be(WorkflowRunLineageAvailability.LegacyUnavailable); } [Fact] diff --git a/test/Aevatar.Workflow.Core.Tests/Execution/WorkflowExecutionContextAdapterTests.cs b/test/Aevatar.Workflow.Core.Tests/Execution/WorkflowExecutionContextAdapterTests.cs index bb96ee88c..d534e94e9 100644 --- a/test/Aevatar.Workflow.Core.Tests/Execution/WorkflowExecutionContextAdapterTests.cs +++ b/test/Aevatar.Workflow.Core.Tests/Execution/WorkflowExecutionContextAdapterTests.cs @@ -94,6 +94,8 @@ await agent.BindWorkflowRunDefinitionAsync( runOrigin: null, scheduleId: null, workflowId: null, + revisionId: null, + definitionVersion: 0, capabilityAdmissionPlan: null, expectedExecutionMode: ExternalCapabilityExecutionMode.Interactive); var adapter = WorkflowExecutionContextAdapter.Create(new RecordingEventHandlerContext(), stateHost); diff --git a/test/Aevatar.Workflow.Core.Tests/Primitives/SubWorkflowOrchestratorTests.cs b/test/Aevatar.Workflow.Core.Tests/Primitives/SubWorkflowOrchestratorTests.cs index 73872240d..8375defce 100644 --- a/test/Aevatar.Workflow.Core.Tests/Primitives/SubWorkflowOrchestratorTests.cs +++ b/test/Aevatar.Workflow.Core.Tests/Primitives/SubWorkflowOrchestratorTests.cs @@ -1313,6 +1313,40 @@ public void ApplyStateTransitions_ShouldMaintainBindingAndInvocationIndexes() state.PendingChildRunIdsByParentRunId["parent-run"].ChildRunIds.Should().ContainSingle(x => x == "child-run-b"); } + [Fact] + public void ApplySubWorkflowInvocationRegistered_ShouldAppendTypedSubWorkflowChildLineage() + { + var state = SubWorkflowOrchestrator.ApplySubWorkflowInvocationRegistered( + new WorkflowRunState + { + RunId = "run-parent-alpha", + }, + new SubWorkflowInvocationRegisteredEvent + { + InvocationId = "invoke-sub-001", + ParentRunId = "run-parent-alpha", + ParentStepId = "step-call-child", + WorkflowName = "sub_flow", + ChildActorId = "actor-child-delta", + ChildRunId = "run-child-beta", + Lifecycle = WorkflowCallLifecycle.Transient, + DefinitionActorId = "workflow-definition:sub_flow", + DefinitionVersion = 2, + RootRunId = "run-root-omega", + Depth = 1, + }); + + state.Lineage.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + state.Lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + state.Lineage.RetryFork.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + var child = state.Lineage.SubWorkflow.ChildRuns.Should().ContainSingle().Subject; + child.RunId.Should().Be("run-child-beta"); + child.ActorId.Should().Be("actor-child-delta"); + child.RelationshipId.Should().Be("invoke-sub-001"); + child.StepId.Should().Be("step-call-child"); + child.RelationKind.Should().Be(WorkflowRunLineageRelationKind.SubWorkflow); + } + [Fact] public void PruneIdleSubWorkflowBindings_ShouldKeepReferencedAndPendingSingletons() { diff --git a/test/Aevatar.Workflow.Core.Tests/WorkflowRunGAgentForkOnFailureTests.cs b/test/Aevatar.Workflow.Core.Tests/WorkflowRunGAgentForkOnFailureTests.cs index e3cdc8f28..adc3b293a 100644 --- a/test/Aevatar.Workflow.Core.Tests/WorkflowRunGAgentForkOnFailureTests.cs +++ b/test/Aevatar.Workflow.Core.Tests/WorkflowRunGAgentForkOnFailureTests.cs @@ -21,6 +21,100 @@ namespace Aevatar.Workflow.Core.Tests; public sealed class WorkflowRunGAgentForkOnFailureTests { + [Fact] + public void BuildExecutionStartLineage_WithForkSeed_ShouldExposeSourceAndOriginalRunIds() + { + var lineage = WorkflowRunGAgent.BuildExecutionStartLineage( + new WorkflowRunForkSeed + { + SourceRunId = "run-source-gamma", + OriginalRunId = "run-original-alpha", + StartAtStepId = "step-retry", + Attempt = 3, + }, + currentLineage: null, + runId: "run-child-beta"); + + lineage.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + lineage.RetryFork.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + lineage.RetryFork.SourceRunId.Should().Be("run-source-gamma"); + lineage.RetryFork.OriginalRunId.Should().Be("run-original-alpha"); + lineage.RetryFork.StartAtStepId.Should().Be("step-retry"); + lineage.RetryFork.Attempt.Should().Be(3); + lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + } + + [Fact] + public void BuildExecutionStartLineage_WithoutForkSeed_ShouldReturnExplicitUnavailableLineage() + { + var lineage = WorkflowRunGAgent.BuildExecutionStartLineage( + forkSeed: null, + currentLineage: null, + runId: "run-standalone-beta"); + + lineage.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + lineage.RetryFork.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + lineage.UnavailableReason.Should().Contain("unavailable"); + } + + [Fact] + public async Task HandleBindWorkflowRunDefinition_WithInitialSubWorkflowLineage_ShouldCommitAndApplyChildLineage() + { + const string childRunId = "run-child-beta"; + const string parentRunId = "run-parent-alpha"; + const string rootRunId = "run-root-omega"; + var harness = await CreateUnboundRunAsync(childRunId); + + var initialLineage = new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Unavailable, + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Available, + ParentRunId = parentRunId, + ParentActorId = "actor-parent-gamma", + ParentStepId = "step-call-child", + RootRunId = rootRunId, + Depth = 2, + }, + }; + + await harness.Agent.HandleBindWorkflowRunDefinition(new BindWorkflowRunDefinitionEvent + { + DefinitionActorId = "definition-child-delta", + WorkflowName = "wf_child_beta", + WorkflowYaml = WorkflowYaml(onFailure: false), + RunId = childRunId, + ScopeId = "scope-child", + ExpectedExecutionMode = ExternalCapabilityExecutionMode.Interactive, + InitialLineage = initialLineage, + }); + + var committed = CommittedEvents(harness.CommittedPublisher) + .Should() + .ContainSingle() + .Subject; + committed.RunId.Should().Be(childRunId); + committed.InitialLineage.SubWorkflow.ParentRunId.Should().Be(parentRunId); + committed.InitialLineage.SubWorkflow.RootRunId.Should().Be(rootRunId); + committed.InitialLineage.SubWorkflow.ParentStepId.Should().Be("step-call-child"); + + harness.Agent.State.RunId.Should().Be(childRunId); + harness.Agent.State.Lineage.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + harness.Agent.State.Lineage.RetryFork.Availability.Should().Be(WorkflowRunLineageAvailability.Unavailable); + harness.Agent.State.Lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + harness.Agent.State.Lineage.SubWorkflow.ParentRunId.Should().Be(parentRunId); + harness.Agent.State.Lineage.SubWorkflow.ParentActorId.Should().Be("actor-parent-gamma"); + harness.Agent.State.Lineage.SubWorkflow.ParentStepId.Should().Be("step-call-child"); + harness.Agent.State.Lineage.SubWorkflow.RootRunId.Should().Be(rootRunId); + harness.Agent.State.Lineage.SubWorkflow.Depth.Should().Be(2); + } + [Fact] public async Task TerminalFailedRun_WithForkPolicy_ShouldCommitForkRequestedEvent() { @@ -386,6 +480,28 @@ private static async Task CreateStartedRunAsync(string workflowYaml, return harness with { StepExecutionId = stepRequest.ExecutionId }; } + private static async Task CreateUnboundRunAsync(string runId) + { + var eventStore = new RecordingEventStore(); + var committedHook = new RecordingCommittedStatePublicationHook(); + var topologyPublisher = new RecordingEventPublisher(runId); + var agent = new WorkflowRunGAgent( + new UnsupportedActorRuntime(), + new UnsupportedActorRuntime(), + new EmptyEventModuleFactory(), + [new EmptyWorkflowModulePack()]) + { + EventSourcingBehaviorFactory = new DefaultEventSourcingBehaviorFactory(eventStore), + EventPublisher = topologyPublisher, + Services = new TestServiceProvider(new NoopRuntimeCallbackScheduler(), committedHook), + Logger = NullLogger.Instance, + }; + SetAgentId(agent, runId); + topologyPublisher.Agent = agent; + await agent.ActivateAsync(); + return new RunHarness(agent, runId, string.Empty, committedHook, topologyPublisher); + } + private static async Task CreateRunAsync( string runId, string workflowYaml, diff --git a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionQueryPortsCoverageTests.cs b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionQueryPortsCoverageTests.cs index 046817e07..3a705ff5c 100644 --- a/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionQueryPortsCoverageTests.cs +++ b/test/Aevatar.Workflow.Host.Api.Tests/WorkflowExecutionQueryPortsCoverageTests.cs @@ -48,6 +48,72 @@ public void WorkflowExecutionReadModelMapper_ShouldMapCurrentStateStatuses( snapshot.DeadLetterRemainingUncompensated.Should().Be(2); snapshot.DeadLetterError.Should().Be("refund failed"); } + [Fact] + public void WorkflowExecutionReadModelMapper_ShouldExposeTypedLineage() + { + var mapper = new WorkflowExecutionReadModelMapper(); + + var snapshot = mapper.ToActorSnapshot(new WorkflowExecutionCurrentStateDocument + { + RootActorId = "actor-child-delta", + RunId = "run-child-beta", + Status = "running", + Lineage = new WorkflowRunLineage + { + Availability = WorkflowRunLineageAvailability.Available, + RetryFork = new WorkflowRunRetryForkLineage + { + Availability = WorkflowRunLineageAvailability.Available, + SourceRunId = "run-source-gamma", + OriginalRunId = "run-original-alpha", + Attempt = 2, + StartAtStepId = "step-retry", + }, + SubWorkflow = new WorkflowRunSubWorkflowLineage + { + Availability = WorkflowRunLineageAvailability.Available, + ParentRunId = "run-parent-alpha", + ParentActorId = "actor-parent-gamma", + ParentStepId = "step-call-child", + RootRunId = "run-root-omega", + Depth = 2, + }, + }, + }); + + snapshot.RunId.Should().Be("run-child-beta"); + snapshot.ActorId.Should().Be("actor-child-delta"); + snapshot.Lineage.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + snapshot.Lineage.RetryFork.SourceRunId.Should().Be("run-source-gamma"); + snapshot.Lineage.RetryFork.OriginalRunId.Should().Be("run-original-alpha"); + snapshot.Lineage.RetryFork.StartAtStepId.Should().Be("step-retry"); + snapshot.Lineage.RetryFork.Attempt.Should().Be(2); + snapshot.Lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.Available); + snapshot.Lineage.SubWorkflow.ParentRunId.Should().Be("run-parent-alpha"); + snapshot.Lineage.SubWorkflow.ParentActorId.Should().Be("actor-parent-gamma"); + snapshot.Lineage.SubWorkflow.ParentStepId.Should().Be("step-call-child"); + snapshot.Lineage.SubWorkflow.RootRunId.Should().Be("run-root-omega"); + snapshot.Lineage.SubWorkflow.Depth.Should().Be(2); + } + + [Fact] + public void WorkflowExecutionReadModelMapper_WhenLineageMissing_ShouldReturnLegacyUnavailable() + { + var mapper = new WorkflowExecutionReadModelMapper(); + + var snapshot = mapper.ToActorSnapshot(new WorkflowExecutionCurrentStateDocument + { + RootActorId = "actor-legacy-delta", + RunId = "run-legacy-beta", + Status = "completed", + }); + + snapshot.Lineage.Availability.Should().Be(WorkflowRunLineageAvailability.LegacyUnavailable); + snapshot.Lineage.RetryFork.Availability.Should().Be(WorkflowRunLineageAvailability.LegacyUnavailable); + snapshot.Lineage.SubWorkflow.Availability.Should().Be(WorkflowRunLineageAvailability.LegacyUnavailable); + snapshot.Lineage.UnavailableReason.Should().Contain("legacy"); + } + [Fact] public void WorkflowExecutionReadModelMapper_ShouldExposeCurrentStateInputFileDescriptors() {