From 29bdb92859d2f0d8b5ad8ce107032f60a1e67d14 Mon Sep 17 00:00:00 2001 From: Omry Yadan Date: Sat, 8 Aug 2026 04:11:50 +0800 Subject: [PATCH] Add controlled-session lifecycle state machine Add the concurrency-safe host-owned lifecycle machine for preparing, active, terminating, and terminated sessions. Latch the first accepted termination cause, enforce granted controller requests, preserve explicit completion across later disconnects, permit bounded controller finalization after application exit, and prevent terminal status reports from rewriting observed outcomes. Reject startup failures after controller activation. Ensure host cancellation, runtime-observation loss, and controller-requested termination end a pending controller-finalization wait without rewriting the already-latched cause. Add lifecycle race, authorization immutability, disconnect, timeout, idempotence, abortive-finalization, startup-boundary, and hostile-result coverage. Clarify controller-requested termination as a distinct design cause. --- docs/CONTROLLED_SESSION_DESIGN.md | 46 +- internal/controlledsession/lifecycle.go | 435 ++++++++++++++ internal/controlledsession/lifecycle_test.go | 568 +++++++++++++++++++ internal/controlledsession/model.go | 55 +- internal/controlledsession/protocol_test.go | 79 ++- 5 files changed, 1166 insertions(+), 17 deletions(-) create mode 100644 internal/controlledsession/lifecycle.go create mode 100644 internal/controlledsession/lifecycle_test.go diff --git a/docs/CONTROLLED_SESSION_DESIGN.md b/docs/CONTROLLED_SESSION_DESIGN.md index 1740009..b4df344 100644 --- a/docs/CONTROLLED_SESSION_DESIGN.md +++ b/docs/CONTROLLED_SESSION_DESIGN.md @@ -464,7 +464,16 @@ cleanup and no canceled request is replayed. Host Reploy emits the authoritative lease lifecycle result: - `terminated(cause, workload_status, workload_output_finalization_status, - controller_finalization_status, cleanup_status, recovery_action)`. + runtime_observation_status, controller_finalization_status, cleanup_status, + recovery_action)`. + +`runtime_observation_status` is `maintained` only when Host Reploy retained +authoritative runtime observation through terminal-result creation. It is +`lost` when observation failed at any earlier point, including after workload +outputs were finalized or the controller completed. This independent monotonic +fact prevents a late observation failure from being hidden by an earlier +termination cause or otherwise successful statuses. A `lost` status always +makes the session invalid, regardless of every other terminal field. `controller_finalization_status` reports the controller protocol outcome: `completed`, `lost`, `finalization-timeout`, `not-completed`, or @@ -1003,14 +1012,18 @@ The first accepted termination cause is latched and never rewritten. Causes include controller-requested termination, workload exit, host cancellation, controller loss, Docker-observation loss, and startup failure. Later events remain diagnostic observations. Workload status, workload-output-finalization -status, controller finalization status, and pre-delivery cleanup success are -reported separately in the session result, so a cleanup failure can fail the -operation without hiding its original cause. Controller exit and delivery-tail -cleanup are reported separately by the invoking host operation after teardown. - -Channel closure is never successful completion. The controller must explicitly -send `complete` after receiving `workload_outputs_finalized` and finalizing its -client-owned results; for OmegaFlow these include the recording artifacts. +status, runtime-observation status, controller finalization status, and +pre-delivery cleanup success are reported separately in the session result, so +a late observation or cleanup failure can fail the operation without hiding +its original cause. Controller exit and delivery-tail cleanup are reported +separately by the invoking host operation after teardown. + +Channel closure is never successful completion. A controller granted the +`complete` operation must explicitly send `complete` after receiving +`workload_outputs_finalized` and finalizing its client-owned results; for +OmegaFlow these include the recording artifacts. Host Reploy does not open a +controller-finalization wait when `complete` was not granted and records that +controller as `not-completed` in the terminal result. Repeated terminate or host cancel operations are idempotent. Input and resize are rejected after `terminating` begins. A single `complete` remains valid during termination while Host Reploy is waiting for controller finalization. A @@ -1028,16 +1041,19 @@ Normal completion is: 5. Host Reploy drains and closes every declared workload-output surface under the finite output-finalization deadline, then emits the one ordered `workload_outputs_finalized` outcome. -6. Host Reploy gives the live controller a bounded finalization period in which - to close its client-owned output and send `complete`. A failed output outcome - remains a session failure even when partial client artifacts are finalized. +6. When the live controller was granted `complete`, Host Reploy gives it a + bounded finalization period in which to close its client-owned output and + send `complete`. Without that grant, Host Reploy skips the wait and records + `not-completed`. A failed output outcome remains a session failure even when + partial client artifacts are finalized. 7. Host Reploy removes the workload container, temporary mounts, networks, and every other lease resource not required to deliver the final result. It keeps the controller and private session channel alive. 8. Host Reploy records the original cause, workload status, - workload-output-finalization status, controller-finalization status, and - pre-delivery cleanup result, then emits the one authoritative `terminated` - event. Only successful event delivery arms the acknowledgement wait. + workload-output-finalization status, runtime-observation status, + controller-finalization status, and pre-delivery cleanup result, then emits + the one authoritative `terminated` event. Only successful event delivery + arms the acknowledgement wait. 9. Host Reploy waits for a bounded `acknowledge_terminated` response. Channel closure is not an acknowledgement. Timeout or disconnect does not block teardown. diff --git a/internal/controlledsession/lifecycle.go b/internal/controlledsession/lifecycle.go new file mode 100644 index 0000000..12d8b64 --- /dev/null +++ b/internal/controlledsession/lifecycle.go @@ -0,0 +1,435 @@ +package controlledsession + +import ( + "errors" + "fmt" + "sync" +) + +type StateV1 string + +const ( + StatePreparingV1 StateV1 = "preparing" + StateActiveV1 StateV1 = "active" + StateTerminatingV1 StateV1 = "terminating" + StateTerminatedV1 StateV1 = "terminated" +) + +type FinishV1 struct { + WorkloadStatus ProcessStatusV1 + WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 + ControllerFinalizationStatus ControllerFinalizationStatusV1 + CleanupStatus CleanupStatusV1 + RecoveryAction RecoveryActionV1 +} + +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" +) + +type ObservationV1 struct { + Kind ObservationKindV1 + WorkloadStatus *ProcessStatusV1 + WorkloadOutputFinalizationStatus *WorkloadOutputFinalizationStatusV1 + Reason string + Finish *FinishV1 +} + +type SnapshotV1 struct { + State StateV1 + Cause TerminationCauseV1 + WorkloadStatus ProcessStatusV1 + WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 + RuntimeObservationStatus RuntimeObservationStatusV1 + ControllerFinalizationStatus ControllerFinalizationStatusV1 + 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 +} + +var ErrRequestRejected = errors.New("controlled-session request rejected") +var ErrObservationRejected = errors.New("controlled-session observation rejected") + +// MachineV1 serializes the authoritative lifecycle of one session. Its mutex +// makes the first accepted termination cause deterministic even when host +// observations race. +type MachineV1 struct { + mu sync.Mutex + authorization AuthorizationV1 + state StateV1 + activated bool + cause TerminationCauseV1 + workload ProcessStatusV1 + workloadOutputs WorkloadOutputFinalizationStatusV1 + controller ControllerFinalizationStatusV1 + runtimeObservation RuntimeObservationStatusV1 + waitingFinalize bool + resultDelivered bool + resultAcknowledged bool + result *ResultV1 +} + +func NewMachineV1(authorization AuthorizationV1) (*MachineV1, error) { + if err := ValidateAuthorizationV1(authorization); err != nil { + return nil, fmt.Errorf("create controlled-session lifecycle: %w", err) + } + return &MachineV1{ + authorization: cloneAuthorizationV1(authorization), + state: StatePreparingV1, + workload: ProcessStatusV1{Kind: ProcessStatusUnknownV1}, + runtimeObservation: RuntimeObservationStatusV1{Kind: RuntimeObservationMaintainedV1}, + controller: ControllerFinalizationStatusV1{Kind: ControllerFinalizationUnknownV1}, + }, nil +} + +func (machine *MachineV1) Snapshot() SnapshotV1 { + machine.mu.Lock() + defer machine.mu.Unlock() + return machine.snapshotLocked() +} + +func (machine *MachineV1) Observe(observation ObservationV1) (TransitionV1, error) { + machine.mu.Lock() + defer machine.mu.Unlock() + before := machine.state + transition := TransitionV1{Before: before, After: before, Cause: machine.cause} + if observation.Kind == ObservationResultDeliveredV1 { + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Reason != "" || observation.Finish != nil || machine.state != StateTerminatedV1 || machine.result == nil { + return transition, fmt.Errorf("%w: result delivery is valid only after termination and carries no payload", ErrObservationRejected) + } + machine.resultDelivered = true + transition.AwaitingResultAcknowledgement = !machine.resultAcknowledged + transition.ResultAcknowledged = machine.resultAcknowledged + transition.Result = cloneResultV1(machine.result) + return transition, nil + } + if machine.state == StateTerminatedV1 { + return transition, fmt.Errorf("%w: lifecycle is already terminated", ErrObservationRejected) + } + + switch observation.Kind { + case ObservationActivatedV1: + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Reason != "" || observation.Finish != nil || machine.state != StatePreparingV1 { + return transition, fmt.Errorf("%w: activation is valid only while preparing and carries no payload", ErrObservationRejected) + } + machine.state = StateActiveV1 + machine.activated = true + machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationActiveV1} + case ObservationWorkloadExitV1: + 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 != "" { + return transition, fmt.Errorf("%w: workload exit cannot follow workload output finalization", ErrObservationRejected) + } + if machine.state == StatePreparingV1 { + return transition, fmt.Errorf("%w: workload exit is invalid before activation", ErrObservationRejected) + } + if err := validateProcessStatusV1(*observation.WorkloadStatus, false); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + if machine.workload.Kind != ProcessStatusUnknownV1 && !equalProcessStatusV1(machine.workload, *observation.WorkloadStatus) { + return transition, fmt.Errorf("%w: workload exit conflicts with the status already observed", ErrObservationRejected) + } + machine.workload = cloneProcessStatusV1(*observation.WorkloadStatus) + if machine.state == StateActiveV1 { + machine.latchLocked(CauseWorkloadExitV1, &transition) + } + case ObservationHostCancelV1: + if err := validateCauseObservationV1(observation); err != nil { + return transition, err + } + machine.latchLocked(CauseHostCancelV1, &transition) + case ObservationControllerLostV1: + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Finish != nil { + return transition, fmt.Errorf("%w: controller loss carries only an optional reason", ErrObservationRejected) + } + if err := validateOptionalSafeTextV1("controller-loss reason", observation.Reason); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + if machine.controller.Kind == ControllerFinalizationUnknownV1 || machine.controller.Kind == ControllerFinalizationActiveV1 { + machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationLostV1, Reason: observation.Reason} + } + machine.waitingFinalize = false + machine.latchLocked(CauseControllerLostV1, &transition) + case ObservationRuntimeObservationLostV1: + if err := validateCauseObservationV1(observation); err != nil { + return transition, err + } + if machine.runtimeObservation.Kind != RuntimeObservationLostV1 { + machine.runtimeObservation = RuntimeObservationStatusV1{Kind: RuntimeObservationLostV1, Reason: observation.Reason} + } + machine.latchLocked(CauseRuntimeObservationLostV1, &transition) + 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) + } + if machine.controller.Kind != ControllerFinalizationUnknownV1 { + return transition, fmt.Errorf("%w: startup failure is invalid after controller activation", ErrObservationRejected) + } + 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 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 { + return transition, fmt.Errorf("%w: workload output finalization requires an observed terminal workload status", ErrObservationRejected) + } + if err := validateWorkloadOutputFinalizationStatusV1(*observation.WorkloadOutputFinalizationStatus); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + 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) + 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) + } + machine.waitingFinalize = false + machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationTimeoutV1} + case ObservationFinishedV1: + 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 err := validateFinishV1(*observation.Finish); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + if machine.workload.Kind != ProcessStatusUnknownV1 && !equalProcessStatusV1(machine.workload, observation.Finish.WorkloadStatus) { + return transition, fmt.Errorf("%w: terminal workload status conflicts with the status already observed", ErrObservationRejected) + } + if machine.workloadOutputs != observation.Finish.WorkloadOutputFinalizationStatus { + return transition, fmt.Errorf("%w: terminal workload output finalization status conflicts with the status already observed", ErrObservationRejected) + } + if err := machine.validateControllerFinishLocked(observation.Finish.ControllerFinalizationStatus); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + result := ResultV1{ + Cause: machine.cause, WorkloadStatus: cloneProcessStatusV1(observation.Finish.WorkloadStatus), + WorkloadOutputFinalizationStatus: observation.Finish.WorkloadOutputFinalizationStatus, + RuntimeObservationStatus: machine.runtimeObservation, + ControllerFinalizationStatus: observation.Finish.ControllerFinalizationStatus, CleanupStatus: observation.Finish.CleanupStatus, + RecoveryAction: observation.Finish.RecoveryAction, + } + if err := ValidateResultV1(result); err != nil { + return transition, fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + machine.workload = result.WorkloadStatus + machine.controller = result.ControllerFinalizationStatus + machine.result = &result + machine.state = StateTerminatedV1 + default: + return transition, fmt.Errorf("%w: observation kind %q is unsupported", ErrObservationRejected, observation.Kind) + } + + transition.After = machine.state + transition.Cause = machine.cause + transition.WorkloadOutputFinalizationStatus = machine.workloadOutputs + transition.AwaitingControllerFinalization = machine.waitingFinalize + transition.AwaitingResultAcknowledgement = machine.resultDelivered && !machine.resultAcknowledged + transition.ResultAcknowledged = machine.resultAcknowledged + transition.Result = cloneResultV1(machine.result) + return transition, nil +} + +func (machine *MachineV1) validateControllerFinishLocked(status ControllerFinalizationStatusV1) error { + switch machine.controller.Kind { + case ControllerFinalizationCompletedV1, ControllerFinalizationLostV1, ControllerFinalizationTimeoutV1, ControllerFinalizationStartupFailedV1: + if machine.controller != status { + return fmt.Errorf("terminal controller finalization status conflicts with the status already observed") + } + case ControllerFinalizationActiveV1, ControllerFinalizationUnknownV1: + if status.Kind != ControllerFinalizationNotCompletedV1 { + return fmt.Errorf("controller without explicit completion must finish as not-completed") + } + default: + return fmt.Errorf("recorded controller finalization status %q is invalid before finish", machine.controller.Kind) + } + return nil +} + +func (machine *MachineV1) ApplyRequest(request RequestV1) (TransitionV1, error) { + machine.mu.Lock() + defer machine.mu.Unlock() + transition := TransitionV1{Before: machine.state, After: machine.state, Cause: machine.cause} + if err := ValidateRequestV1(request); err != nil { + return transition, fmt.Errorf("%w: %v", ErrRequestRejected, err) + } + // A terminal acknowledgement is protocol flow control, not a granted + // controller capability. Its lifecycle position is the authorization. + if request.Kind == RequestAcknowledgeTerminatedV1 { + if machine.state != StateTerminatedV1 || machine.result == nil || !machine.resultDelivered { + return transition, fmt.Errorf("%w: terminal result has not been delivered", ErrRequestRejected) + } + machine.resultAcknowledged = true + transition.RequestAccepted = true + transition.AwaitingResultAcknowledgement = false + transition.ResultAcknowledged = true + transition.Result = cloneResultV1(machine.result) + return transition, nil + } + operation := operationForRequestV1(request.Kind) + if !containsOperationV1(machine.authorization.Operations, operation) { + return transition, fmt.Errorf("%w: operation %q was not granted", ErrRequestRejected, operation) + } + + switch machine.state { + case StateActiveV1: + switch request.Kind { + case RequestInputV1, RequestResizeV1: + case RequestTerminateV1: + machine.latchLocked(CauseControllerTerminateV1, &transition) + case RequestCompleteV1: + return transition, fmt.Errorf("%w: completion is valid only after workload outputs are finalized", ErrRequestRejected) + } + case StateTerminatingV1: + switch request.Kind { + case RequestTerminateV1: + // 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) + } + machine.waitingFinalize = false + machine.controller = ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1} + default: + return transition, fmt.Errorf("%w: %s is not accepted while terminating", ErrRequestRejected, request.Kind) + } + case StatePreparingV1, StateTerminatedV1: + return transition, fmt.Errorf("%w: requests are not accepted while %s", ErrRequestRejected, machine.state) + default: + return transition, fmt.Errorf("%w: lifecycle state %q is invalid", ErrRequestRejected, machine.state) + } + + transition.After = machine.state + transition.Cause = machine.cause + transition.WorkloadOutputFinalizationStatus = machine.workloadOutputs + transition.AwaitingControllerFinalization = machine.waitingFinalize + transition.AwaitingResultAcknowledgement = machine.resultDelivered && !machine.resultAcknowledged + transition.RequestAccepted = true + transition.ResultAcknowledged = machine.resultAcknowledged + return transition, nil +} + +func (machine *MachineV1) latchLocked(cause TerminationCauseV1, transition *TransitionV1) { + if machine.cause != "" { + return + } + machine.cause = cause + machine.state = StateTerminatingV1 + transition.CauseLatched = true + transition.BeginTermination = true +} + +func (machine *MachineV1) snapshotLocked() SnapshotV1 { + return SnapshotV1{ + State: machine.state, Cause: machine.cause, WorkloadStatus: cloneProcessStatusV1(machine.workload), + WorkloadOutputFinalizationStatus: machine.workloadOutputs, + RuntimeObservationStatus: machine.runtimeObservation, + ControllerFinalizationStatus: machine.controller, AwaitingControllerFinalization: machine.waitingFinalize, + AwaitingResultAcknowledgement: machine.resultDelivered && !machine.resultAcknowledged, + ResultAcknowledged: machine.resultAcknowledged, Result: cloneResultV1(machine.result), + } +} + +func validateCauseObservationV1(observation ObservationV1) error { + if observation.WorkloadStatus != nil || observation.WorkloadOutputFinalizationStatus != nil || observation.Finish != nil { + return fmt.Errorf("%w: %s carries only an optional reason", ErrObservationRejected, observation.Kind) + } + if err := validateOptionalSafeTextV1(string(observation.Kind)+" reason", observation.Reason); err != nil { + return fmt.Errorf("%w: %v", ErrObservationRejected, err) + } + return nil +} + +func validateFinishV1(finish FinishV1) error { + if err := validateProcessStatusV1(finish.WorkloadStatus, true); err != nil { + return err + } + if err := validateWorkloadOutputFinalizationStatusV1(finish.WorkloadOutputFinalizationStatus); err != nil { + return err + } + if err := validateTerminalControllerFinalizationStatusV1(finish.ControllerFinalizationStatus); err != nil { + return err + } + return validateCleanupResultV1(finish.CleanupStatus, finish.RecoveryAction) +} + +func equalProcessStatusV1(left ProcessStatusV1, right ProcessStatusV1) bool { + if left.Kind != right.Kind || left.Reason != right.Reason || (left.Code == nil) != (right.Code == nil) { + return false + } + return left.Code == nil || *left.Code == *right.Code +} + +func cloneProcessStatusV1(status ProcessStatusV1) ProcessStatusV1 { + result := status + if status.Code != nil { + code := *status.Code + result.Code = &code + } + return result +} + +func cloneResultV1(result *ResultV1) *ResultV1 { + if result == nil { + return nil + } + copy := *result + copy.WorkloadStatus = cloneProcessStatusV1(result.WorkloadStatus) + return © +} + +func containsOperationV1(operations []OperationV1, target OperationV1) bool { + for _, operation := range operations { + if operation == target { + return true + } + } + return false +} diff --git a/internal/controlledsession/lifecycle_test.go b/internal/controlledsession/lifecycle_test.go new file mode 100644 index 0000000..1650c47 --- /dev/null +++ b/internal/controlledsession/lifecycle_test.go @@ -0,0 +1,568 @@ +package controlledsession + +import ( + "errors" + "sync" + "testing" +) + +func activatedMachineV1(t *testing.T) *MachineV1 { + t.Helper() + machine, err := NewMachineV1(testAuthorizationV1()) + if err != nil { + t.Fatalf("NewMachineV1() error = %v", err) + } + transition, err := machine.Observe(ObservationV1{Kind: ObservationActivatedV1}) + if err != nil { + t.Fatalf("Observe(activated) error = %v", err) + } + if transition.Before != StatePreparingV1 || transition.After != StateActiveV1 { + t.Fatalf("activation transition = %#v", transition) + } + return machine +} + +func observeWorkloadExitV1(t *testing.T, machine *MachineV1, code int) { + t.Helper() + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadExitV1, + WorkloadStatus: &ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + }); err != nil { + t.Fatalf("Observe(workload exit) error = %v", err) + } +} + +func observeOutputsFinalizedV1(t *testing.T, machine *MachineV1, status WorkloadOutputFinalizationStatusV1) TransitionV1 { + t.Helper() + transition, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &status, + }) + if err != nil { + t.Fatalf("Observe(workload outputs finalized) error = %v", err) + } + return transition +} + +func finishLifecycleV1(t *testing.T, machine *MachineV1, finish FinishV1) TransitionV1 { + t.Helper() + transition, err := machine.Observe(ObservationV1{Kind: ObservationFinishedV1, Finish: &finish}) + if err != nil { + t.Fatalf("Observe(finished) error = %v", err) + } + return transition +} + +func TestLifecycleOutputBarrierPrecedesControllerFinalization(t *testing.T) { + machine := activatedMachineV1(t) + code := 0 + observeWorkloadExitV1(t, machine, code) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestInputV1, Bytes: []byte("late")}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("late input error = %v", err) + } + + if snapshot := machine.Snapshot(); snapshot.AwaitingControllerFinalization { + t.Fatalf("controller finalization began before output barrier: %#v", snapshot) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("complete before output barrier error = %v", err) + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationFinishedV1, Finish: &FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationNotCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("finish before output barrier error = %v", err) + } + + outputStatus := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + barrier := observeOutputsFinalizedV1(t, machine, outputStatus) + if !barrier.AwaitingControllerFinalization || barrier.WorkloadOutputFinalizationStatus != outputStatus { + t.Fatalf("output-barrier transition = %#v", barrier) + } + complete, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}) + if err != nil || complete.AwaitingControllerFinalization { + t.Fatalf("ApplyRequest(complete) = %#v, %v", complete, err) + } + + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: outputStatus, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + 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) + } +} + +func TestLifecycleSkipsFinalizationWaitWithoutCompleteGrant(t *testing.T) { + authorization := testAuthorizationV1() + authorization.Operations = []OperationV1{OperationInputV1, OperationResizeV1, OperationTerminateV1} + machine, err := NewMachineV1(authorization) + if err != nil { + t.Fatal(err) + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationActivatedV1}); err != nil { + t.Fatal(err) + } + + code := 0 + observeWorkloadExitV1(t, machine, code) + outputStatus := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + barrier := observeOutputsFinalizedV1(t, machine, outputStatus) + if barrier.AwaitingControllerFinalization { + t.Fatalf("output barrier entered an impossible finalization wait: %#v", barrier) + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationControllerFinalizationExpiredV1}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("finalization expiry without a wait error = %v", err) + } + + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: outputStatus, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationNotCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result == nil || finished.Result.ControllerFinalizationStatus.Kind != ControllerFinalizationNotCompletedV1 { + t.Fatalf("finish transition = %#v", finished) + } +} + +func TestLifecycleCompleteNeverTerminatesActiveWorkload(t *testing.T) { + machine := activatedMachineV1(t) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("active complete error = %v", err) + } + if snapshot := machine.Snapshot(); snapshot.State != StateActiveV1 || snapshot.Cause != "" || snapshot.ControllerFinalizationStatus.Kind != ControllerFinalizationActiveV1 { + t.Fatalf("active complete changed lifecycle = %#v", snapshot) + } +} + +func TestLifecycleControllerTerminationUsesOutputBarrier(t *testing.T) { + machine := activatedMachineV1(t) + first, err := machine.ApplyRequest(RequestV1{Kind: RequestTerminateV1}) + if err != nil || !first.CauseLatched || first.Cause != CauseControllerTerminateV1 { + t.Fatalf("first terminate = %#v, %v", first, err) + } + second, err := machine.ApplyRequest(RequestV1{Kind: RequestTerminateV1}) + if err != nil || second.CauseLatched || second.Cause != CauseControllerTerminateV1 { + t.Fatalf("repeated terminate = %#v, %v", second, err) + } + outputStatus := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &outputStatus, + }); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("output barrier before workload status error = %v", err) + } + + code := 0 + observeWorkloadExitV1(t, machine, code) + observeOutputsFinalizedV1(t, machine, outputStatus) + third, err := machine.ApplyRequest(RequestV1{Kind: RequestTerminateV1}) + if err != nil || !third.AwaitingControllerFinalization { + t.Fatalf("repeated terminate during finalization = %#v, %v", third, err) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatalf("complete after output barrier error = %v", err) + } +} + +func TestLifecycleRepeatedTerminationSignalsPreserveFinalizationWait(t *testing.T) { + actions := []struct { + name string + apply func(*MachineV1) (TransitionV1, error) + }{ + {name: "controller terminate", apply: func(machine *MachineV1) (TransitionV1, error) { + return machine.ApplyRequest(RequestV1{Kind: RequestTerminateV1}) + }}, + {name: "host cancel", apply: func(machine *MachineV1) (TransitionV1, error) { + return machine.Observe(ObservationV1{Kind: ObservationHostCancelV1, Reason: "host interrupted"}) + }}, + {name: "runtime observation lost", apply: func(machine *MachineV1) (TransitionV1, error) { + return machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}) + }}, + } + for _, action := range actions { + t.Run(action.name, func(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 0) + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + transition, err := action.apply(machine) + if err != nil || transition.Cause != CauseWorkloadExitV1 || !transition.AwaitingControllerFinalization { + t.Fatalf("repeated termination signal = %#v, %v", transition, err) + } + }) + } +} + +func TestLifecycleFailedOutputFinalizationRemainsAuthoritative(t *testing.T) { + machine := activatedMachineV1(t) + code := 1 + observeWorkloadExitV1(t, machine, code) + outputStatus := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "output deadline expired"} + observeOutputsFinalizedV1(t, machine, outputStatus) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatal(err) + } + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: outputStatus, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result.WorkloadOutputFinalizationStatus != outputStatus { + t.Fatalf("result rewrote failed output finalization = %#v", finished.Result) + } +} + +func TestLifecycleRuntimeObservationLossRejectsDrainedOutput(t *testing.T) { + machine := activatedMachineV1(t) + if _, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}); err != nil { + t.Fatal(err) + } + observeWorkloadExitV1(t, machine, 1) + + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &drained, + }); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("drained output after runtime observation loss error = %v", err) + } + if snapshot := machine.Snapshot(); 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) + } +} + +func TestLifecycleRuntimeObservationLossRejectsDrainedOutputAfterEarlierCause(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 1) + if _, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}); err != nil { + t.Fatal(err) + } + + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadOutputsFinalizedV1, + WorkloadOutputFinalizationStatus: &drained, + }); !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 != "" { + t.Fatalf("rejected output finalization changed lifecycle = %#v", snapshot) + } +} + +func TestLifecycleLateRuntimeObservationLossInvalidatesCompletedSession(t *testing.T) { + machine := activatedMachineV1(t) + code := 0 + observeWorkloadExitV1(t, machine, code) + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + observeOutputsFinalizedV1(t, machine, drained) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatal(err) + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}); err != nil { + t.Fatal(err) + } + + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: drained, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result == nil || finished.Result.Cause != CauseWorkloadExitV1 { + t.Fatalf("finish transition = %#v", finished) + } + if got := finished.Result.RuntimeObservationStatus; got != (RuntimeObservationStatusV1{Kind: RuntimeObservationLostV1, Reason: "docker unavailable"}) { + t.Fatalf("runtime observation status = %#v", got) + } + if finished.Result.WorkloadOutputFinalizationStatus != drained || finished.Result.ControllerFinalizationStatus.Kind != ControllerFinalizationCompletedV1 { + t.Fatalf("late observation loss rewrote earlier terminal facts: %#v", finished.Result) + } +} + +func TestLifecycleOutputFinalizationIsSingleAndImmutable(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 0) + drained := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + observeOutputsFinalizedV1(t, machine, drained) + for _, status := range []WorkloadOutputFinalizationStatusV1{ + drained, + {Kind: WorkloadOutputFinalizationFailedV1, Reason: "late failure"}, + } { + if _, err := machine.Observe(ObservationV1{Kind: ObservationWorkloadOutputsFinalizedV1, WorkloadOutputFinalizationStatus: &status}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("duplicate output finalization %#v error = %v", status, err) + } + } +} + +func TestLifecycleOutputBarrierRejectsLateTerminalFacts(t *testing.T) { + machine, err := NewMachineV1(testAuthorizationV1()) + if err != nil { + t.Fatal(err) + } + if _, err := machine.Observe(ObservationV1{Kind: ObservationHostCancelV1, Reason: "host interrupted startup"}); err != nil { + t.Fatal(err) + } + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + code := 0 + if _, err := machine.Observe(ObservationV1{ + Kind: ObservationWorkloadExitV1, + WorkloadStatus: &ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + }); !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 TestLifecycleConnectionClosureCannotBecomeSuccess(t *testing.T) { + machine := activatedMachineV1(t) + transition, err := machine.Observe(ObservationV1{Kind: ObservationControllerLostV1, Reason: "session channel closed"}) + if err != nil || transition.Cause != CauseControllerLostV1 || !transition.BeginTermination { + t.Fatalf("controller-loss transition = %#v, %v", transition, err) + } + status := ProcessStatusV1{Kind: ProcessStatusTerminatedV1, Reason: "lease teardown"} + if _, err := machine.Observe(ObservationV1{Kind: ObservationWorkloadExitV1, WorkloadStatus: &status}); err != nil { + t.Fatalf("Observe(workload exit after controller loss) error = %v", err) + } + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "controller disconnected"}) + if snapshot := machine.Snapshot(); snapshot.AwaitingControllerFinalization { + t.Fatalf("lost controller entered finalization wait: %#v", snapshot) + } + + _, err = machine.Observe(ObservationV1{Kind: ObservationFinishedV1, Finish: &FinishV1{ + WorkloadStatus: status, + WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "controller disconnected"}, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }}) + if !errors.Is(err, ErrObservationRejected) { + t.Fatalf("forged successful controller finalization status error = %v", err) + } + finished := finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: status, + WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "controller disconnected"}, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationLostV1, Reason: "session channel closed"}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if finished.Result.ControllerFinalizationStatus.Kind != ControllerFinalizationLostV1 { + t.Fatalf("result = %#v", finished.Result) + } +} + +func TestLifecyclePreservesAcceptedCompletionAfterChannelCloses(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 0) + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatal(err) + } + transition, err := machine.Observe(ObservationV1{Kind: ObservationControllerLostV1, Reason: "channel closed during teardown"}) + if err != nil || transition.Cause != CauseWorkloadExitV1 { + t.Fatalf("controller loss after completion = %#v, %v", transition, err) + } + if snapshot := machine.Snapshot(); snapshot.ControllerFinalizationStatus.Kind != ControllerFinalizationCompletedV1 { + t.Fatalf("controller finalization status = %#v", snapshot.ControllerFinalizationStatus) + } +} + +func TestLifecycleFirstAcceptedCauseWinsConcurrentRace(t *testing.T) { + machine := activatedMachineV1(t) + code := 1 + observations := []ObservationV1{ + {Kind: ObservationHostCancelV1, Reason: "host interrupted"}, + {Kind: ObservationControllerLostV1, Reason: "channel closed"}, + {Kind: ObservationRuntimeObservationLostV1, Reason: "docker unavailable"}, + {Kind: ObservationWorkloadExitV1, WorkloadStatus: &ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}}, + } + var wait sync.WaitGroup + latched := make(chan TerminationCauseV1, len(observations)) + for _, observation := range observations { + observation := observation + wait.Add(1) + go func() { + defer wait.Done() + transition, err := machine.Observe(observation) + if err != nil { + t.Errorf("Observe(%s) error = %v", observation.Kind, err) + return + } + if transition.CauseLatched { + latched <- transition.Cause + } + }() + } + wait.Wait() + close(latched) + var causes []TerminationCauseV1 + for cause := range latched { + causes = append(causes, cause) + } + if len(causes) != 1 { + t.Fatalf("latched causes = %v, want exactly one", causes) + } + if snapshot := machine.Snapshot(); snapshot.Cause != causes[0] || snapshot.State != StateTerminatingV1 { + t.Fatalf("snapshot = %#v, latched = %v", snapshot, causes) + } +} + +func TestLifecycleStartupFailureIsLimitedToPreActivation(t *testing.T) { + active := activatedMachineV1(t) + if _, err := active.Observe(ObservationV1{Kind: ObservationStartupFailureV1, Reason: "late startup failure"}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("active startup failure error = %v", err) + } + + preActivation, err := NewMachineV1(testAuthorizationV1()) + if err != nil { + t.Fatal(err) + } + if _, err := preActivation.Observe(ObservationV1{Kind: ObservationHostCancelV1, Reason: "host interrupted"}); err != nil { + t.Fatal(err) + } + transition, err := preActivation.Observe(ObservationV1{Kind: ObservationStartupFailureV1, Reason: "startup stopped"}) + if err != nil || transition.Cause != CauseHostCancelV1 { + t.Fatalf("pre-activation startup failure = %#v, %v", transition, err) + } + if snapshot := preActivation.Snapshot(); snapshot.ControllerFinalizationStatus.Kind != ControllerFinalizationStartupFailedV1 { + t.Fatalf("controller finalization status = %#v", snapshot.ControllerFinalizationStatus) + } +} + +func TestLifecycleEnforcesAuthorizationAndCopiesIt(t *testing.T) { + authorization := testAuthorizationV1() + authorization.Operations = []OperationV1{OperationCompleteV1} + machine, err := NewMachineV1(authorization) + if err != nil { + t.Fatal(err) + } + authorization.Operations[0] = OperationInputV1 + if _, err := machine.Observe(ObservationV1{Kind: ObservationActivatedV1}); err != nil { + t.Fatal(err) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestInputV1, Bytes: []byte("not granted")}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("request enabled through caller mutation error = %v", err) + } +} + +func TestLifecycleFinalizationExpiryCannotBeRewrittenByLateComplete(t *testing.T) { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 1) + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + if _, err := machine.Observe(ObservationV1{Kind: ObservationControllerFinalizationExpiredV1}); err != nil { + t.Fatal(err) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("late complete error = %v", err) + } + if snapshot := machine.Snapshot(); snapshot.Cause != CauseWorkloadExitV1 || snapshot.ControllerFinalizationStatus.Kind != ControllerFinalizationTimeoutV1 { + t.Fatalf("snapshot = %#v", snapshot) + } +} + +func TestLifecycleCompleteAndExpiryHaveOneAuthoritativeOutcome(t *testing.T) { + for attempt := 0; attempt < 50; attempt++ { + machine := activatedMachineV1(t) + observeWorkloadExitV1(t, machine, 0) + observeOutputsFinalizedV1(t, machine, WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}) + var wait sync.WaitGroup + wait.Add(2) + errorsSeen := make(chan error, 2) + go func() { + defer wait.Done() + _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}) + errorsSeen <- err + }() + go func() { + defer wait.Done() + _, err := machine.Observe(ObservationV1{Kind: ObservationControllerFinalizationExpiredV1}) + errorsSeen <- err + }() + wait.Wait() + close(errorsSeen) + accepted := 0 + for err := range errorsSeen { + if err == nil { + accepted++ + } else if !errors.Is(err, ErrRequestRejected) && !errors.Is(err, ErrObservationRejected) { + t.Fatalf("unexpected race error = %v", err) + } + } + if accepted != 1 { + t.Fatalf("accepted outcomes = %d, want 1", accepted) + } + kind := machine.Snapshot().ControllerFinalizationStatus.Kind + if kind != ControllerFinalizationCompletedV1 && kind != ControllerFinalizationTimeoutV1 { + t.Fatalf("controller finalization status = %q", kind) + } + } +} + +func TestLifecycleAcceptsOnlyPostResultAcknowledgement(t *testing.T) { + machine := activatedMachineV1(t) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestAcknowledgeTerminatedV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("early acknowledgement error = %v", err) + } + code := 0 + observeWorkloadExitV1(t, machine, code) + outputStatus := WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + observeOutputsFinalizedV1(t, machine, outputStatus) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestCompleteV1}); err != nil { + t.Fatal(err) + } + finishLifecycleV1(t, machine, FinishV1{ + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: outputStatus, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + }) + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestAcknowledgeTerminatedV1}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("pre-delivery acknowledgement error = %v", err) + } + delivered, err := machine.Observe(ObservationV1{Kind: ObservationResultDeliveredV1}) + if err != nil || !delivered.AwaitingResultAcknowledgement || delivered.Result == nil { + t.Fatalf("result-delivered observation = %#v, %v", delivered, err) + } + for attempt := 0; attempt < 2; attempt++ { + transition, err := machine.ApplyRequest(RequestV1{Kind: RequestAcknowledgeTerminatedV1}) + if err != nil || !transition.RequestAccepted || !transition.ResultAcknowledged || transition.Result == nil { + t.Fatalf("acknowledgement %d = %#v, %v", attempt+1, transition, err) + } + } + if snapshot := machine.Snapshot(); !snapshot.ResultAcknowledged { + t.Fatalf("post-acknowledgement snapshot = %#v", snapshot) + } + if _, err := machine.ApplyRequest(RequestV1{Kind: RequestResizeV1, Columns: 80, Rows: 24}); !errors.Is(err, ErrRequestRejected) { + t.Fatalf("late resize error = %v", err) + } +} + +func TestLifecycleRejectsResultDeliveryBeforeTermination(t *testing.T) { + machine := activatedMachineV1(t) + if _, err := machine.Observe(ObservationV1{Kind: ObservationResultDeliveredV1}); !errors.Is(err, ErrObservationRejected) { + t.Fatalf("early result-delivered observation error = %v", err) + } +} diff --git a/internal/controlledsession/model.go b/internal/controlledsession/model.go index cdf47a9..4dad109 100644 --- a/internal/controlledsession/model.go +++ b/internal/controlledsession/model.go @@ -57,6 +57,18 @@ type WorkloadOutputFinalizationStatusV1 struct { Reason string `json:"reason,omitempty"` } +type RuntimeObservationStatusKindV1 string + +const ( + RuntimeObservationMaintainedV1 RuntimeObservationStatusKindV1 = "maintained" + RuntimeObservationLostV1 RuntimeObservationStatusKindV1 = "lost" +) + +type RuntimeObservationStatusV1 struct { + Kind RuntimeObservationStatusKindV1 `json:"kind"` + Reason string `json:"reason,omitempty"` +} + type CleanupStatusKindV1 string const ( @@ -81,6 +93,7 @@ type ResultV1 struct { Cause TerminationCauseV1 `json:"cause"` WorkloadStatus ProcessStatusV1 `json:"workload_status"` WorkloadOutputFinalizationStatus WorkloadOutputFinalizationStatusV1 `json:"workload_output_finalization_status"` + RuntimeObservationStatus RuntimeObservationStatusV1 `json:"runtime_observation_status"` ControllerFinalizationStatus ControllerFinalizationStatusV1 `json:"controller_finalization_status"` CleanupStatus CleanupStatusV1 `json:"cleanup_status"` RecoveryAction RecoveryActionV1 `json:"recovery_action"` @@ -96,10 +109,50 @@ func ValidateResultV1(result ResultV1) error { if err := validateWorkloadOutputFinalizationStatusV1(result.WorkloadOutputFinalizationStatus); err != nil { return err } + if err := validateRuntimeObservationStatusV1(result.RuntimeObservationStatus); err != nil { + return err + } if err := validateTerminalControllerFinalizationStatusV1(result.ControllerFinalizationStatus); err != nil { return err } - return validateCleanupResultV1(result.CleanupStatus, result.RecoveryAction) + if err := validateCleanupResultV1(result.CleanupStatus, result.RecoveryAction); err != nil { + return err + } + return validateResultConsistencyV1(result) +} + +func validateResultConsistencyV1(result ResultV1) error { + if result.Cause == CauseRuntimeObservationLostV1 { + if result.RuntimeObservationStatus.Kind != RuntimeObservationLostV1 { + return fmt.Errorf("runtime-observation-loss termination requires lost runtime observation status") + } + if result.WorkloadOutputFinalizationStatus.Kind != WorkloadOutputFinalizationFailedV1 { + return fmt.Errorf("runtime-observation-loss termination requires failed workload output finalization") + } + } + if result.Cause == CauseControllerLostV1 && result.ControllerFinalizationStatus.Kind != ControllerFinalizationLostV1 { + return fmt.Errorf("controller-loss termination requires lost controller finalization status") + } + if result.Cause == CauseStartupFailureV1 && result.ControllerFinalizationStatus.Kind != ControllerFinalizationStartupFailedV1 { + return fmt.Errorf("startup-failure termination requires startup-failed controller finalization status") + } + return nil +} + +func validateRuntimeObservationStatusV1(status RuntimeObservationStatusV1) error { + switch status.Kind { + case RuntimeObservationMaintainedV1: + if status.Reason != "" { + return fmt.Errorf("maintained runtime observation must not contain a reason") + } + case RuntimeObservationLostV1: + if err := validateOptionalSafeTextV1("runtime observation loss reason", status.Reason); err != nil { + return err + } + default: + return fmt.Errorf("runtime observation status %q is invalid", status.Kind) + } + return nil } func validateWorkloadOutputFinalizationStatusV1(status WorkloadOutputFinalizationStatusV1) error { diff --git a/internal/controlledsession/protocol_test.go b/internal/controlledsession/protocol_test.go index 6a7aa06..5c68dca 100644 --- a/internal/controlledsession/protocol_test.go +++ b/internal/controlledsession/protocol_test.go @@ -3,6 +3,7 @@ package controlledsession import ( "bytes" "encoding/binary" + "encoding/json" "reflect" "strings" "testing" @@ -50,6 +51,14 @@ func TestEventV1RoundTripsStrictTypedFrames(t *testing.T) { {Kind: EventTerminatedV1, Terminated: &ResultV1{ Cause: CauseWorkloadExitV1, WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}, + RuntimeObservationStatus: RuntimeObservationStatusV1{Kind: RuntimeObservationMaintainedV1}, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, RecoveryAction: RecoveryNoneV1, + }}, + {Kind: EventTerminatedV1, Terminated: &ResultV1{ + Cause: CauseWorkloadExitV1, WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1}, + RuntimeObservationStatus: RuntimeObservationStatusV1{Kind: RuntimeObservationLostV1, Reason: "docker unavailable"}, ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1}, CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, RecoveryAction: RecoveryNoneV1, }}, @@ -157,7 +166,7 @@ func TestFrameV1RejectsBadMagicVersionTruncationAndUnknownJSON(t *testing.T) { t.Fatalf("ReadEventV1(case-variant duplicate JSON) error = %v", err) } - nestedCaseVariant := []byte(`{"cause":"workload-exit","workload_status":{"Kind":"exited","code":0},"workload_output_finalization_status":{"kind":"drained"},"controller_finalization_status":{"kind":"completed"},"cleanup_status":{"kind":"succeeded"},"recovery_action":"none"}`) + nestedCaseVariant := []byte(`{"cause":"workload-exit","workload_status":{"Kind":"exited","code":0},"workload_output_finalization_status":{"kind":"drained"},"runtime_observation_status":{"kind":"maintained"},"controller_finalization_status":{"kind":"completed"},"cleanup_status":{"kind":"succeeded"},"recovery_action":"none"}`) framed.Reset() if err := writeFrameV1(&framed, wireEventTerminatedV1, nestedCaseVariant); err != nil { t.Fatal(err) @@ -224,6 +233,74 @@ func TestValidateEventV1RejectsInvalidOutputFinalizationOutcomes(t *testing.T) { } } +func TestReadEventV1RejectsContradictoryRuntimeObservationLossResults(t *testing.T) { + code := 0 + valid := ResultV1{ + Cause: CauseRuntimeObservationLostV1, + WorkloadStatus: ProcessStatusV1{Kind: ProcessStatusExitedV1, Code: &code}, + WorkloadOutputFinalizationStatus: WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationFailedV1, Reason: "runtime observation lost"}, + RuntimeObservationStatus: RuntimeObservationStatusV1{Kind: RuntimeObservationLostV1, Reason: "docker unavailable"}, + ControllerFinalizationStatus: ControllerFinalizationStatusV1{Kind: ControllerFinalizationLostV1, Reason: "docker unavailable"}, + CleanupStatus: CleanupStatusV1{Kind: CleanupStatusSucceededV1}, + RecoveryAction: RecoveryNoneV1, + } + + tests := []struct { + name string + wantError string + mutate func(*ResultV1) + }{ + { + name: "runtime observation maintained", + wantError: "runtime-observation-loss termination requires lost runtime observation status", + mutate: func(result *ResultV1) { + result.RuntimeObservationStatus = RuntimeObservationStatusV1{Kind: RuntimeObservationMaintainedV1} + }, + }, + { + name: "runtime loss with outputs drained", + wantError: "runtime-observation-loss termination requires failed workload output finalization", + mutate: func(result *ResultV1) { + result.WorkloadOutputFinalizationStatus = WorkloadOutputFinalizationStatusV1{Kind: WorkloadOutputFinalizationDrainedV1} + }, + }, + { + name: "controller loss with completed finalization", + wantError: "controller-loss termination requires lost controller finalization status", + mutate: func(result *ResultV1) { + result.Cause = CauseControllerLostV1 + result.ControllerFinalizationStatus = ControllerFinalizationStatusV1{Kind: ControllerFinalizationCompletedV1} + }, + }, + { + name: "startup failure without startup-failed finalization", + wantError: "startup-failure termination requires startup-failed controller finalization status", + mutate: func(result *ResultV1) { + result.Cause = CauseStartupFailureV1 + result.ControllerFinalizationStatus = ControllerFinalizationStatusV1{Kind: ControllerFinalizationNotCompletedV1} + }, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + result := valid + test.mutate(&result) + payload, err := json.Marshal(result) + if err != nil { + t.Fatal(err) + } + var framed bytes.Buffer + if err := writeFrameV1(&framed, wireEventTerminatedV1, payload); err != nil { + t.Fatal(err) + } + if _, err := ReadEventV1(&framed); err == nil || !strings.Contains(err.Error(), test.wantError) { + t.Fatalf("ReadEventV1(contradictory terminated result) error = %v", err) + } + }) + } +} + func TestValidateRequestV1RejectsAmbiguousUnionPayloads(t *testing.T) { tests := []RequestV1{ {Kind: RequestInputV1},