From bdfbbb4eae6a4dd9035dbdae65a39c9184a41fc5 Mon Sep 17 00:00:00 2001 From: Omry Yadan Date: Sat, 8 Aug 2026 15:32:25 +0800 Subject: [PATCH] Enforce controlled-session output finalization --- docs/CONTROLLED_SESSION_DESIGN.md | 18 +- internal/controlledsession/lifecycle.go | 151 ++++++---- internal/controlledsession/lifecycle_test.go | 274 ++++++++++++++++++- 3 files changed, 370 insertions(+), 73 deletions(-) diff --git a/docs/CONTROLLED_SESSION_DESIGN.md b/docs/CONTROLLED_SESSION_DESIGN.md index ec971ba..255d48e 100644 --- a/docs/CONTROLLED_SESSION_DESIGN.md +++ b/docs/CONTROLLED_SESSION_DESIGN.md @@ -10,9 +10,10 @@ summary: Capability-scoped execution sessions that inherit Reploy's global conta - Decision state: Focused review complete; high-level decisions approved - Implementation state: Initial global sandbox prerequisites, trusted - application-startup verification, controlled-session authorization, and the - initial framed protocol are implemented; lifecycle, controlled networking, - and Docker orchestration remain later slices + application-startup verification, controlled-session authorization, the + framed protocol, and lifecycle state machine, including its output-finalization + barrier and timeout outcome, are implemented; bounded output draining, + controlled-session networking, and Docker orchestration remain later slices - Initial runtime: Linux containers under Docker - Motivating clients: OmegaFlow recording, sandboxed AI agents, security inspection, and untrusted-code execution @@ -448,8 +449,8 @@ cleanup and no canceled request is replayed. ### Session Events -- `opened`: reports the effective dimensions, identity, generation, and fixed - session capabilities. +- `opened`: reports the effective dimensions, identity, generation, fixed + session capabilities, and workload-output-finalization timeout. - `output(bytes)`: ordered PTY output bytes. - `workload_exit(status, reason)`: reports host-observed workload-shell exit. @@ -502,10 +503,9 @@ Host Reploy owns workload-output finalization; it never waits indefinitely for workload cooperation. Once termination begins, it rejects new output surfaces, performs bounded graceful shutdown followed by forced container stop, and continues draining the PTY. The immutable session plan carries a finite -output-finalization deadline. The initial implementation may use a fixed -host-owned value, but the effective value is reported by `opened` and applies -to workload shutdown, final buffered-byte delivery, and controller -backpressure. +output-finalization deadline. Protocol v1 defines an initial host-owned default +of 30 seconds; the effective value is reported by `opened` and applies to +workload shutdown, final buffered-byte delivery, and controller backpressure. If every final byte is delivered and every output surface reaches EOF before the deadline, Host Reploy emits `workload_outputs_finalized(drained)` only after diff --git a/internal/controlledsession/lifecycle.go b/internal/controlledsession/lifecycle.go index 12d8b64..72e832a 100644 --- a/internal/controlledsession/lifecycle.go +++ b/internal/controlledsession/lifecycle.go @@ -26,16 +26,17 @@ type FinishV1 struct { type ObservationKindV1 string const ( - ObservationActivatedV1 ObservationKindV1 = "activated" - ObservationWorkloadExitV1 ObservationKindV1 = "workload-exit" - ObservationHostCancelV1 ObservationKindV1 = "host-cancel" - ObservationControllerLostV1 ObservationKindV1 = "controller-lost" - ObservationRuntimeObservationLostV1 ObservationKindV1 = "runtime-observation-lost" - ObservationStartupFailureV1 ObservationKindV1 = "startup-failure" - ObservationWorkloadOutputsFinalizedV1 ObservationKindV1 = "workload-outputs-finalized" - ObservationControllerFinalizationExpiredV1 ObservationKindV1 = "controller-finalization-expired" - ObservationFinishedV1 ObservationKindV1 = "finished" - ObservationResultDeliveredV1 ObservationKindV1 = "result-delivered" + ObservationActivatedV1 ObservationKindV1 = "activated" + ObservationWorkloadExitV1 ObservationKindV1 = "workload-exit" + ObservationHostCancelV1 ObservationKindV1 = "host-cancel" + ObservationControllerLostV1 ObservationKindV1 = "controller-lost" + ObservationRuntimeObservationLostV1 ObservationKindV1 = "runtime-observation-lost" + ObservationStartupFailureV1 ObservationKindV1 = "startup-failure" + ObservationWorkloadOutputsFinalizedV1 ObservationKindV1 = "workload-outputs-finalized" + ObservationWorkloadOutputFinalizationExpiredV1 ObservationKindV1 = "workload-output-finalization-expired" + ObservationControllerFinalizationExpiredV1 ObservationKindV1 = "controller-finalization-expired" + ObservationFinishedV1 ObservationKindV1 = "finished" + ObservationResultDeliveredV1 ObservationKindV1 = "result-delivered" ) type ObservationV1 struct { @@ -47,30 +48,32 @@ type ObservationV1 struct { } type SnapshotV1 struct { - State StateV1 - Cause TerminationCauseV1 - WorkloadStatus ProcessStatusV1 - WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 - RuntimeObservationStatus RuntimeObservationStatusV1 - ControllerFinalizationStatus ControllerFinalizationStatusV1 - AwaitingControllerFinalization bool - AwaitingResultAcknowledgement bool - ResultAcknowledged bool - Result *ResultV1 + State StateV1 + Cause TerminationCauseV1 + WorkloadStatus ProcessStatusV1 + WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 + RuntimeObservationStatus RuntimeObservationStatusV1 + ControllerFinalizationStatus ControllerFinalizationStatusV1 + AwaitingWorkloadOutputFinalization bool + AwaitingControllerFinalization bool + AwaitingResultAcknowledgement bool + ResultAcknowledged bool + Result *ResultV1 } type TransitionV1 struct { - Before StateV1 - After StateV1 - Cause TerminationCauseV1 - CauseLatched bool - BeginTermination bool - WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 - AwaitingControllerFinalization bool - AwaitingResultAcknowledgement bool - RequestAccepted bool - ResultAcknowledged bool - Result *ResultV1 + Before StateV1 + After StateV1 + Cause TerminationCauseV1 + CauseLatched bool + BeginTermination bool + WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 + AwaitingWorkloadOutputFinalization bool + AwaitingControllerFinalization bool + AwaitingResultAcknowledgement bool + RequestAccepted bool + ResultAcknowledged bool + Result *ResultV1 } var ErrRequestRejected = errors.New("controlled-session request rejected") @@ -89,6 +92,7 @@ type MachineV1 struct { workloadOutputs WorkloadOutputFinalizationStatusV1 controller ControllerFinalizationStatusV1 runtimeObservation RuntimeObservationStatusV1 + waitingOutputs bool waitingFinalize bool resultDelivered bool resultAcknowledged bool @@ -145,7 +149,10 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro if observation.WorkloadStatus == nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Reason != "" || observation.Finish != nil { return transition, fmt.Errorf("%w: workload exit requires exactly one workload status", ErrObservationRejected) } - if machine.workloadOutputs.Kind != "" { + if machine.workloadOutputs.Kind != "" && !(machine.activated && + machine.state == StateTerminatingV1 && + machine.workload.Kind == ProcessStatusUnknownV1 && + machine.workloadOutputs.Kind == WorkloadOutputFinalizationFailedV1) { return transition, fmt.Errorf("%w: workload exit cannot follow workload output finalization", ErrObservationRejected) } if machine.state == StatePreparingV1 { @@ -186,13 +193,14 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro machine.runtimeObservation = RuntimeObservationStatusV1{Kind: RuntimeObservationLostV1, Reason: observation.Reason} } machine.latchLocked(CauseRuntimeObservationLostV1, &transition) + // Observation loss does not close an activated workload's output + // surfaces. Keep that barrier pending until the supervisor explicitly + // reports failed closure or its bounded finalization deadline expires. + machine.finalizePreActivationOutputsForRuntimeObservationLossLocked(observation.Reason) case ObservationStartupFailureV1: if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Finish != nil { return transition, fmt.Errorf("%w: startup failure carries only a reason", ErrObservationRejected) } - if machine.workloadOutputs.Kind != "" { - return transition, fmt.Errorf("%w: startup failure cannot follow workload output finalization", ErrObservationRejected) - } if err := validateRequiredSafeTextV1("startup-failure reason", observation.Reason); err != nil { return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) } @@ -202,16 +210,12 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationStartupFailedV1, Reason: observation.Reason} machine.latchLocked(CauseStartupFailureV1, &transition) case ObservationWorkloadOutputsFinalizedV1: - if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus == nil || observation.Reason != "" || observation.Finish != nil { - return transition, fmt.Errorf("%w: workload output finalization requires exactly one status", ErrObservationRejected) + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus == nil || observation.Reason != "" || observation.Finish != nil || !machine.waitingOutputs { + return transition, fmt.Errorf("%w: workload output finalization requires exactly one status while output finalization is pending", ErrObservationRejected) } - if machine.state != StateTerminatingV1 { - return transition, fmt.Errorf("%w: workload output finalization is valid only while terminating", ErrObservationRejected) - } - if machine.workloadOutputs.Kind != "" { - return transition, fmt.Errorf("%w: workload output finalization was already observed", ErrObservationRejected) - } - if machine.activated && machine.workload.Kind == ProcessStatusUnknownV1 { + if machine.activated && machine.workload.Kind == ProcessStatusUnknownV1 && + !(machine.runtimeObservation.Kind == RuntimeObservationLostV1 && + observation.WorkloadOutputFinalizationStatus.Kind == WorkloadOutputFinalizationFailedV1) { return transition, fmt.Errorf("%w: workload output finalization requires an observed terminal workload status", ErrObservationRejected) } if err := validateWorkloadOutputFinalizationStatusV1(*observation.WorkloadOutputFinalizationStatus); err != nil { @@ -220,9 +224,18 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro if machine.runtimeObservation.Kind == RuntimeObservationLostV1 && observation.WorkloadOutputFinalizationStatus.Kind == WorkloadOutputFinalizationDrainedV1 { return transition, fmt.Errorf("%w: runtime observation loss requires failed workload output finalization", ErrObservationRejected) } - machine.workloadOutputs = *observation.WorkloadOutputFinalizationStatus - machine.waitingFinalize = machine.controller.Kind == ControllerFinalizationActiveV1 && - containsOperationV1(machine.authorization.Operations, OperationCompleteV1) + machine.completeOutputFinalizationLocked(*observation.WorkloadOutputFinalizationStatus) + case ObservationWorkloadOutputFinalizationExpiredV1: + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Finish != nil || !machine.waitingOutputs { + return transition, fmt.Errorf("%w: output-finalization expiry carries only a required reason while output finalization is pending", ErrObservationRejected) + } + if err := validateRequiredSafeTextV1("workload output finalization expiry reason", observation.Reason); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + machine.completeOutputFinalizationLocked(WorkloadOutputFinalizationStatusV1{ + Kind: WorkloadOutputFinalizationFailedV1, + Reason: observation.Reason, + }) case ObservationControllerFinalizationExpiredV1: if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Finish != nil || observation.Reason != "" || !machine.waitingFinalize { return transition, fmt.Errorf("%w: finalization expiry requires an active controller-finalization wait", ErrObservationRejected) @@ -233,8 +246,8 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Reason != "" || observation.Finish == nil { return transition, fmt.Errorf("%w: finish requires exactly one terminal status set", ErrObservationRejected) } - if machine.state != StateTerminatingV1 || machine.workloadOutputs.Kind == "" || machine.waitingFinalize { - return transition, fmt.Errorf("%w: finish requires finalized workload output and no pending controller finalization", ErrObservationRejected) + if machine.state != StateTerminatingV1 || machine.waitingOutputs || machine.workloadOutputs.Kind == "" || machine.waitingFinalize { + return transition, fmt.Errorf("%w: finish requires finalized workload output and no pending output or controller finalization", ErrObservationRejected) } if err := validateFinishV1(*observation.Finish); err != nil { return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) @@ -269,6 +282,7 @@ func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, erro transition.After = machine.state transition.Cause = machine.cause transition.WorkloadOutputFinalizationStatus = machine.workloadOutputs + transition.AwaitingWorkloadOutputFinalization = machine.waitingOutputs transition.AwaitingControllerFinalization = machine.waitingFinalize transition.AwaitingResultAcknowledgement = machine.resultDelivered && !machine.resultAcknowledged transition.ResultAcknowledged = machine.resultAcknowledged @@ -332,8 +346,8 @@ func (machine *MachineV1) ApplyRequest(request RequestV1) (TransitionV1, error) // Repeated graceful termination is idempotent and does not alter // output or controller finalization already in progress. case RequestCompleteV1: - if !machine.waitingFinalize || machine.controller.Kind != ControllerFinalizationActiveV1 { - return transition, fmt.Errorf("%w: completion is not awaiting controller finalization", ErrRequestRejected) + if machine.waitingOutputs || !machine.waitingFinalize || machine.controller.Kind != ControllerFinalizationActiveV1 { + return transition, fmt.Errorf("%w: completion requires finalized workload outputs and a pending controller finalization", ErrRequestRejected) } machine.waitingFinalize = false machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1} @@ -349,6 +363,7 @@ func (machine *MachineV1) ApplyRequest(request RequestV1) (TransitionV1, error) transition.After = machine.state transition.Cause = machine.cause transition.WorkloadOutputFinalizationStatus = machine.workloadOutputs + transition.AwaitingWorkloadOutputFinalization = machine.waitingOutputs transition.AwaitingControllerFinalization = machine.waitingFinalize transition.AwaitingResultAcknowledgement = machine.resultDelivered && !machine.resultAcknowledged transition.RequestAccepted = true @@ -362,6 +377,13 @@ func (machine *MachineV1) latchLocked(cause TerminationCauseV1, transition *Tran } machine.cause = cause machine.state = StateTerminatingV1 + if transition.Before == StateActiveV1 { + machine.waitingOutputs = true + } else { + // No workload ran, so there are no workload-originated surfaces to + // drain. Record the barrier as satisfied rather than inventing a wait. + machine.workloadOutputs = WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + } transition.CauseLatched = true transition.BeginTermination = true } @@ -372,8 +394,9 @@ func (machine *MachineV1) snapshotLocked() SnapshotV1 { WorkloadOutputFinalizationStatus: machine.workloadOutputs, RuntimeObservationStatus: machine.runtimeObservation, ControllerFinalizationStatus: machine.controller, AwaitingControllerFinalization: machine.waitingFinalize, - AwaitingResultAcknowledgement: machine.resultDelivered && !machine.resultAcknowledged, - ResultAcknowledged: machine.resultAcknowledged, Result: cloneResultV1(machine.result), + AwaitingWorkloadOutputFinalization: machine.waitingOutputs, + AwaitingResultAcknowledgement: machine.resultDelivered && !machine.resultAcknowledged, + ResultAcknowledged: machine.resultAcknowledged, Result: cloneResultV1(machine.result), } } @@ -400,6 +423,28 @@ func validateFinishV1(finish FinishV1) error { return validateCleanupResultV1(finish.CleanupStatus, finish.RecoveryAction) } +func (machine *MachineV1) finalizePreActivationOutputsForRuntimeObservationLossLocked(reason string) { + if machine.activated || + machine.cause != CauseRuntimeObservationLostV1 || + machine.workloadOutputs.Kind != WorkloadOutputFinalizationDrainedV1 { + return + } + if reason == "" { + reason = "runtime observation was lost before workload output finalization completed" + } + machine.completeOutputFinalizationLocked(WorkloadOutputFinalizationStatusV1{ + Kind: WorkloadOutputFinalizationFailedV1, + Reason: reason, + }) +} + +func (machine *MachineV1) completeOutputFinalizationLocked(status WorkloadOutputFinalizationStatusV1) { + machine.workloadOutputs = status + machine.waitingOutputs = false + machine.waitingFinalize = machine.controller.Kind == ControllerFinalizationActiveV1 && + containsOperationV1(machine.authorization.Operations, OperationCompleteV1) +} + func equalProcessStatusV1(left ProcessStatusV1, right ProcessStatusV1) bool { if left.Kind != right.Kind || left.Reason != right.Reason || (left.Code == nil) != (right.Code == nil) { return false diff --git a/internal/controlledsession/lifecycle_test.go b/internal/controlledsession/lifecycle_test.go index 1650c47..bc10d40 100644 --- a/internal/controlledsession/lifecycle_test.go +++ b/internal/controlledsession/lifecycle_test.go @@ -57,6 +57,9 @@ func TestLifecycleOutputBarrierPrecedesControllerFinalization(t *testing.T) { machine := activatedMachineV1(t) code := 0 observeWorkloadExitV1(t, machine, code) + if snapshot := machine.Snapshot(); !snapshot.AwaitingWorkloadOutputFinalization { + t.Fatalf("workload exit did not open output barrier: %#v", snapshot) + } if _, err := machine.ApplyRequest(RequestV1{Kind: RequestInputV1, Bytes: []byte("late")}); !errors.Is(err, ErrRequestRejected) { t.Fatalf("late input error = %v", err) } @@ -227,10 +230,26 @@ func TestLifecycleFailedOutputFinalizationRemainsAuthoritative(t *testing.T) { func TestLifecycleRuntimeObservationLossRejectsDrainedOutput(t *testing.T) { machine := activatedMachineV1(t) - if _, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}); err != nil { + transition, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}) + if err != nil { t.Fatal(err) } - observeWorkloadExitV1(t, machine, 1) + if !transition.AwaitingWorkloadOutputFinalization || transition.WorkloadOutputFinalizationStatus.Kind != "" { + t.Fatalf("runtime observation loss transition = %#v", transition) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("complete before failed output finalization error = %v", err) + } + failed := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "docker unavailable"} + if _, err := machine.Observe(ObservationV1{Kind: ObservationFinishedV1, Finish: &FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusUnknownV1}, + WorkloadOutputFinalizationStatus: failed, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationNotCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("finish before failed output finalization error = %v", err) + } drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} if _, err := machine.Observe(ObservationV1{ @@ -239,14 +258,48 @@ func TestLifecycleRuntimeObservationLossRejectsDrainedOutput(t *testing.T) { }); !errors.Is(err, ErrObservationRejected) { t.Fatalf("drained output after runtime observation loss error = %v", err) } - if snapshot := machine.Snapshot(); snapshot.WorkloadOutputFinalizationStatus.Kind != "" { + if snapshot := machine.Snapshot(); !snapshot.AwaitingWorkloadOutputFinalization || snapshot.WorkloadOutputFinalizationStatus.Kind != "" { t.Fatalf("rejected output finalization changed lifecycle = %#v", snapshot) } - failed := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "runtime observation lost"} - transition := observeOutputsFinalizedV1(t, machine, failed) - if transition.Cause != CauseRuntimeObservationLostV1 || transition.WorkloadOutputFinalizationStatus != failed { - t.Fatalf("failed output finalization transition = %#v", transition) + barrier := observeOutputsFinalizedV1(t, machine, failed) + if barrier.AwaitingWorkloadOutputFinalization || barrier.WorkloadOutputFinalizationStatus != failed { + t.Fatalf("failed output finalization transition = %#v", barrier) + } + observeWorkloadExitV1(t, machine, 1) +} + +func TestLifecyclePreActivationRuntimeObservationLossCanFinish(t *testing.T) { + machine, err := NewMachineV1(testAuthorizationV1()) + if err != nil { + t.Fatal(err) + } + transition, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1}) + if err != nil { + t.Fatal(err) + } + failed := WorkloadOutputFinalizationStatusV1{ + Kind: WorkloadOutputFinalizationFailedV1, + Reason: "runtime observation was lost before workload output finalization completed", + } + if transition.Cause != CauseRuntimeObservationLostV1 || + transition.AwaitingWorkloadOutputFinalization || + transition.WorkloadOutputFinalizationStatus != failed { + t.Fatalf("pre-activation runtime observation loss = %#v", transition) + } + + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusUnknownV1}, + WorkloadOutputFinalizationStatus: failed, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationNotCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.After != StateTerminatedV1 || finished.Result == nil { + t.Fatalf("finish transition = %#v", finished) + } + if err := ValidateResultV1(*finished.Result); err != nil { + t.Fatalf("terminal result is invalid: %v", err) } } @@ -264,7 +317,9 @@ func TestLifecycleRuntimeObservationLossRejectsDrainedOutputAfterEarlierCause(t }); !errors.Is(err, ErrObservationRejected) { t.Fatalf("drained output after later runtime observation loss error = %v", err) } - if snapshot := machine.Snapshot(); snapshot.Cause != CauseWorkloadExitV1 || snapshot.WorkloadOutputFinalizationStatus.Kind != "" { + failed := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "docker unavailable"} + observeOutputsFinalizedV1(t, machine, failed) + if snapshot := machine.Snapshot(); snapshot.Cause != CauseWorkloadExitV1 || snapshot.WorkloadOutputFinalizationStatus != failed || snapshot.AwaitingWorkloadOutputFinalization { t.Fatalf("rejected output finalization changed lifecycle = %#v", snapshot) } } @@ -323,7 +378,16 @@ func TestLifecycleOutputBarrierRejectsLateTerminalFacts(t *testing.T) { if _, err := machine.Observe(ObservationV1{Kind: ObservationHostCancelV1, Reason: "host interrupted startup"}); err != nil { t.Fatal(err) } - observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + if snapshot := machine.Snapshot(); snapshot.AwaitingWorkloadOutputFinalization || snapshot.WorkloadOutputFinalizationStatus.Kind != WorkloadOutputFinalizationDrainedV1 { + t.Fatalf("pre-activation termination output state = %#v", snapshot) + } + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &drained, + }); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("duplicate pre-activation output finalization error = %v", err) + } code := 0 if _, err := machine.Observe(ObservationV1{ Kind: ObservationWorkloadExitV1, @@ -331,8 +395,196 @@ func TestLifecycleOutputBarrierRejectsLateTerminalFacts(t *testing.T) { }); !errors.Is(err, ErrObservationRejected) { t.Fatalf("workload exit after output barrier error = %v", err) } - if _, err := machine.Observe(ObservationV1{Kind: ObservationStartupFailureV1, Reason: "late startup failure"}); !errors.Is(err, ErrObservationRejected) { - t.Fatalf("startup failure after output barrier error = %v", err) +} + +func TestLifecycleOutputFinalizationExpiryCannotBecomeDrained(t *testing.T) { + machine := activatedMachineV1(t) + code := 0 + observeWorkloadExitV1(t, machine, code) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("complete before output timeout error = %v", err) + } + + transition, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputFinalizationExpiredV1, + Reason: "output finalization exceeded 30s", + }) + if err != nil { + t.Fatal(err) + } + if transition.AwaitingWorkloadOutputFinalization || !transition.AwaitingControllerFinalization { + t.Fatalf("output timeout transition = %#v", transition) + } + + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &drained, + }); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("late drained output outcome error = %v", err) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatal(err) + } + + failed := WorkloadOutputFinalizationStatusV1{ + Kind: WorkloadOutputFinalizationFailedV1, + Reason: "output finalization exceeded 30s", + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationFinishedV1, Finish: &FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: drained, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("finish with rewritten output status error = %v", err) + } + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: failed, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result == nil || finished.Result.WorkloadOutputFinalizationStatus != failed { + t.Fatalf("terminal result = %#v", finished.Result) + } +} + +func TestLifecycleOutputFinalizationExpiryAcceptsLateWorkloadExit(t *testing.T) { + tests := []struct { + name string + cause TerminationCauseV1 + start func(*MachineV1) error + }{ + { + name: "controller terminate", + cause: CauseControllerTerminateV1, + start: func(machine *MachineV1) error { + _, err := machine.ApplyRequest(RequestV1{Kind: RequestTerminateV1}) + return err + }, + }, + { + name: "host cancel", + cause: CauseHostCancelV1, + start: func(machine *MachineV1) error { + _, err := machine.Observe(ObservationV1{Kind: ObservationHostCancelV1, Reason: "host interrupted"}) + return err + }, + }, + { + name: "controller lost", + cause: CauseControllerLostV1, + start: func(machine *MachineV1) error { + _, err := machine.Observe(ObservationV1{Kind: ObservationControllerLostV1, Reason: "controller disconnected"}) + return err + }, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + machine := activatedMachineV1(t) + if err := test.start(machine); err != nil { + t.Fatal(err) + } + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputFinalizationExpiredV1, + Reason: "output finalization exceeded 30s", + }); err != nil { + t.Fatal(err) + } + + code := 137 + status := ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadExitV1, + WorkloadStatus: &status, + }); err != nil { + t.Fatalf("late workload exit error = %v", err) + } + + snapshot := machine.Snapshot() + failed := WorkloadOutputFinalizationStatusV1{ + Kind: WorkloadOutputFinalizationFailedV1, + Reason: "output finalization exceeded 30s", + } + if snapshot.Cause != test.cause || + !equalProcessStatusV1(snapshot.WorkloadStatus, status) || + snapshot.WorkloadOutputFinalizationStatus != failed { + t.Fatalf("late workload exit snapshot = %#v", snapshot) + } + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadExitV1, + WorkloadStatus: &status, + }); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("duplicate late workload exit error = %v", err) + } + + if snapshot.AwaitingControllerFinalization { + if _, err := machine.Observe(ObservationV1{Kind: ObservationControllerFinalizationExpiredV1}); err != nil { + t.Fatal(err) + } + snapshot = machine.Snapshot() + } + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: status, + WorkloadOutputFinalizationStatus: failed, + ControllerFinalizationStatus: snapshot.ControllerFinalizationStatus, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result == nil || !equalProcessStatusV1(finished.Result.WorkloadStatus, status) { + t.Fatalf("terminal result lost late workload status = %#v", finished.Result) + } + }) + } +} + +func TestLifecycleAcceptsExactlyOneConcurrentOutputFinalizationOutcome(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 0) + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + observations := []ObservationV1{ + {Kind: ObservationWorkloadOutputsFinalizedV1, WorkloadOutputFinalizationStatus: &drained}, + {Kind: ObservationWorkloadOutputFinalizationExpiredV1, Reason: "output finalization exceeded 30s"}, + } + + var wait sync.WaitGroup + results := make(chan error, len(observations)) + for _, observation := range observations { + observation := observation + wait.Add(1) + go func() { + defer wait.Done() + _, err := machine.Observe(observation) + results <- err + }() + } + wait.Wait() + close(results) + + accepted := 0 + rejected := 0 + for err := range results { + switch { + case err == nil: + accepted++ + case errors.Is(err, ErrObservationRejected): + rejected++ + default: + t.Fatalf("unexpected finalization error = %v", err) + } + } + if accepted != 1 || rejected != 1 { + t.Fatalf("output finalization outcomes: accepted=%d rejected=%d", accepted, rejected) + } + if snapshot := machine.Snapshot(); snapshot.AwaitingWorkloadOutputFinalization { + t.Fatalf("snapshot still awaits output finalization: %#v", snapshot) + } else if err := validateWorkloadOutputFinalizationStatusV1(snapshot.WorkloadOutputFinalizationStatus); err != nil { + t.Fatalf("latched output finalization status = %#v: %v", snapshot.WorkloadOutputFinalizationStatus, err) } }