From b4637a4800ad80bc983137cec0a6e98d7d0e24f3 Mon Sep 17 00:00:00 2001 From: Puma Date: Sat, 3 Oct 2026 07:13:16 -0700 Subject: [PATCH 1/5] refactor(workflow): box private request and completion payloads --- .github/workflows/quality-gates.yml | 7 +++ .../src/graph/inference_validation_state.rs | 43 ++++++++++++++++++- .../src/workflow/task_execution_worker.rs | 27 +++++++++--- ...6-10-03-workflow-private-payload-layout.md | 7 +++ 4 files changed, 75 insertions(+), 9 deletions(-) create mode 100644 docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-payload-layout.md diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index badccc11a..6cdec4212 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,13 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run private workflow owner regression suites + run: | + for module in graph::inference_validation_state::tests:: workflow::task_execution_worker::tests::; do + cargo test -p pantograph-workflow-service --lib "$module" -- --list > "$RUNNER_TEMP/private-owner-tests.list" + python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith(sys.argv[2]) and line.endswith(": test")]; print(sys.argv[2], "tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/private-owner-tests.list" "$module" + cargo test -p pantograph-workflow-service --lib "$module" + done - name: Run workflow style contract regressions run: | cargo test -p pantograph-workflow-service --lib lint_style_ -- --list > "$RUNNER_TEMP/workflow-style-tests.list" diff --git a/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs b/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs index c0759411e..061f55ae0 100644 --- a/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs +++ b/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs @@ -1158,7 +1158,7 @@ pub(crate) enum DependencyEnvironmentActionIntentStateResolution { Blocked(DependencyEnvironmentActionIntentResult), RequestReady { intent: DependencyEnvironmentActionIntent, - environment_request: ValidatedDependencyEnvironmentRequest, + environment_request: Box, }, } @@ -1846,7 +1846,7 @@ fn request_ready_dependency_environment_action_resolution( ) -> DependencyEnvironmentActionIntentStateResolution { DependencyEnvironmentActionIntentStateResolution::RequestReady { intent, - environment_request, + environment_request: Box::new(environment_request), } } @@ -1877,6 +1877,15 @@ mod tests { use super::*; + #[test] + fn action_resolution_keeps_validated_environment_request_indirect() { + assert!(std::mem::size_of::() <= 384); + assert!( + std::mem::size_of::() + < std::mem::size_of::() + ); + } + #[tokio::test] async fn action_intent_state_rejects_stale_graph_revision() { let store = CurrentInferenceValidationStateStore::new(); @@ -2284,6 +2293,36 @@ mod tests { DependencyEnvironmentActionIntentStatus::RequestReady ); assert!(result.diagnostics.is_empty()); + + let resolution = store + .resolve_dependency_environment_action_request(state_request_with_validation_session( + "graph-session-1", + "aaaaaaaaaaaaaaaa", + "aaaaaaaaaaaaaaaa", + "validation.session.1", + "dependency-node-1", + true, + )) + .await; + let DependencyEnvironmentActionIntentStateResolution::RequestReady { + intent, + environment_request, + } = resolution + else { + panic!("current proof must retain a ready validated request"); + }; + let raw = environment_request.as_request(); + assert_eq!(raw.action, intent.action); + assert_eq!(raw.planning_request.task_id.as_str(), "image_generation"); + assert_eq!( + raw.identity_key, + DependencyPlanningIdentityKey::from_planning_request(&raw.planning_request).unwrap() + ); + raw.validate() + .expect("retained environment request validates"); + let encoded = serde_json::to_value(raw).unwrap(); + let decoded = ValidatedDependencyEnvironmentRequest::try_from(encoded.clone()).unwrap(); + assert_eq!(serde_json::to_value(decoded.as_request()).unwrap(), encoded); } #[tokio::test] diff --git a/crates/pantograph-workflow-service/src/workflow/task_execution_worker.rs b/crates/pantograph-workflow-service/src/workflow/task_execution_worker.rs index 5a15d2c97..e632a73db 100644 --- a/crates/pantograph-workflow-service/src/workflow/task_execution_worker.rs +++ b/crates/pantograph-workflow-service/src/workflow/task_execution_worker.rs @@ -368,7 +368,7 @@ pub(super) enum WorkflowTaskExecutionWorkerShutdownReason { pub(super) enum WorkflowTaskExecutionWorkerOutcome { TaskTerminal(WorkflowTaskExecutionWorkerTerminalOutcome), TaskDeferred(WorkflowTaskExecutionWorkerDeferredOutcome), - RuntimeBranchCompleted(WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome), + RuntimeBranchCompleted(Box), RuntimeBranchFailed(WorkflowTaskExecutionWorkerRuntimeBranchFailedOutcome), RuntimeBranchCancelled(WorkflowTaskExecutionWorkerRuntimeBranchCancelledOutcome), RuntimeBranchDeferred(WorkflowTaskExecutionWorkerRuntimeBranchDeferredOutcome), @@ -1124,12 +1124,12 @@ impl WorkflowTaskExecutionWorkerOutcome { response: WorkflowRunResponse, diagnostics: Vec, ) -> Self { - Self::RuntimeBranchCompleted(WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { + Self::RuntimeBranchCompleted(Box::new(WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { session_id: command.session_id.clone(), workflow_run_id: command.workflow_run_id.clone(), response, diagnostics, - }) + })) } pub(super) fn runtime_branch_failed( @@ -1370,12 +1370,12 @@ async fn claim_and_execute_runtime_branch_event( .await { Ok(response) => WorkflowTaskExecutionWorkerOutcome::RuntimeBranchCompleted( - WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { + Box::new(WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { session_id: command.session_id.clone(), workflow_run_id: command.workflow_run_id.clone(), response, diagnostics: Vec::new(), - }, + }), ), Err(error) => WorkflowTaskExecutionWorkerOutcome::runtime_branch_failed( command, @@ -1886,12 +1886,12 @@ fn runtime_branch_batch_member_completion( unix_timestamp_ms(), proof) { Ok(_record) => match member_outcome.completed_response { Some(response) => WorkflowTaskExecutionWorkerOutcome::RuntimeBranchCompleted( - WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { + Box::new(WorkflowTaskExecutionWorkerRuntimeBranchCompletedOutcome { session_id: member_outcome.session_id.clone(), workflow_run_id: member_outcome.workflow_run_id.clone(), response, diagnostics, - }, + }), ), None => WorkflowTaskExecutionWorkerOutcome::RuntimeBranchFailed( WorkflowTaskExecutionWorkerRuntimeBranchFailedOutcome { @@ -4353,6 +4353,11 @@ mod tests { assert_eq!(outcome.diagnostics, vec![diagnostic]); } + #[test] + fn worker_outcome_keeps_completed_response_indirect() { + assert!(std::mem::size_of::() <= 128); + } + #[test] fn runtime_branch_completed_outcome_preserves_run_scope_and_response() { let command = runtime_branch_command(); @@ -4379,6 +4384,14 @@ mod tests { assert_eq!(outcome.session_id, command.session_id); assert_eq!(outcome.workflow_run_id, command.workflow_run_id); assert_eq!(outcome.response, response); + assert_eq!( + serde_json::to_value(&outcome.response).unwrap(), + serde_json::json!({ + "workflow_run_id": command.workflow_run_id, + "outputs": [], + "timing_ms": 42 + }) + ); assert_eq!(outcome.diagnostics, vec![diagnostic]); } diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-payload-layout.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-payload-layout.md new file mode 100644 index 000000000..bc0ce099b --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-payload-layout.md @@ -0,0 +1,7 @@ +# Private validation request and worker completion payloads + +Fresh PR #37 Clippy retains an 896-byte private RequestReady branch and six large-error reports caused by a 144-byte private RuntimeBranchCompleted outcome. Box only RequestReady.environment_request and the completed worker outcome. The constructor inventory is one validated-request helper and three completed-outcome constructors. The graph session still borrows the validated request for the dependency service; worker response extraction and all Result signatures remain unchanged. Public contracts and the public Ready inference projection are excluded. + +The real ready-proof regression now checks the retained request task, canonical identity, validation and full JSON round trip. Existing blocked/missing/stale proof tests remain. The worker success test retains complete response/scope/diagnostic equality and adds exact response JSON; existing cancellation, unavailable and shutdown behaviors remain. Two private layout bounds provide source-layout evidence, not measured performance claims. CI explicitly discovers and runs both complete owner test modules with nonzero guards. + +Root approved this bounded design. Whitespace and focused test syntax formatting are checked locally; root source review accepted all four files at frozen tree c2a6e2439f8df3a7e4d00461cf36175fbbe19ff4; hosted execution remains pending. No local Rust execution is claimed. From 09ae3182b3efa0e7921763122b25d0f0fc2ada11 Mon Sep 17 00:00:00 2001 From: Puma Date: Sat, 3 Oct 2026 07:22:35 -0700 Subject: [PATCH 2/5] refactor(workflow): box ready inference projection payload --- .github/workflows/quality-gates.yml | 7 +++ CHANGELOG.md | 4 ++ .../src/graph/inference_validation_state.rs | 4 +- .../executable_validation_snapshot.rs | 4 +- .../src/workflow/task_graph.rs | 2 +- .../workflow/tests/task_binding_resolution.rs | 4 +- .../src/workflow/tests/task_graph.rs | 50 ++++++++++++++++++- ...10-03-ready-inference-projection-layout.md | 7 +++ 8 files changed, 73 insertions(+), 9 deletions(-) create mode 100644 docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-ready-inference-projection-layout.md diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index 6cdec4212..a28eea8d2 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,13 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run workflow inference projection regression suites + run: | + for module in workflow::tests::task_graph:: workflow::tests::task_binding_resolution:: workflow::executable_validation_snapshot::tests::; do + cargo test -p pantograph-workflow-service --lib "$module" -- --list > "$RUNNER_TEMP/projection-tests.list" + python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith(sys.argv[2]) and line.endswith(": test")]; print(sys.argv[2], "tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/projection-tests.list" "$module" + cargo test -p pantograph-workflow-service --lib "$module" + done - name: Run private workflow owner regression suites run: | for module in graph::inference_validation_state::tests:: workflow::task_execution_worker::tests::; do diff --git a/CHANGELOG.md b/CHANGELOG.md index 2c8136067..ebc94a76a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -59,3 +59,7 @@ The seven validated runtime-host contract wrappers now implement standard `AsRef ### Diagnostic event artifact payload construction `DiagnosticEventPayload::IoArtifactObserved` now stores `Box`. Rust callers constructing the public variant use `Box::new`, and owned pattern bindings contain a box. Raw DTO fields, validation and tagged JSON stay unchanged; this adds one allocation for this variant without a measured performance claim. The crate remains `publish = false`. + +### Workflow ready inference projection construction + +`WorkflowSchedulerInferenceTaskProjection::Ready` now holds `Box`. Rust constructors use `Box::new`, and owned pattern bindings contain a box; the raw record, equality, borrowed readers and downstream scheduler intent serialization remain unchanged. The projection enum itself has no Serde implementation. This adds one allocation for ready projections without a measured performance claim; the crate remains `publish = false`. diff --git a/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs b/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs index 061f55ae0..000e8c783 100644 --- a/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs +++ b/crates/pantograph-workflow-service/src/graph/inference_validation_state.rs @@ -1468,7 +1468,7 @@ impl CurrentInferenceValidationNodeRecord { }, )?; return Ok(WorkflowSchedulerInferenceTaskProjection::Ready( - WorkflowSchedulerReadyInferenceTaskProjection { + Box::new(WorkflowSchedulerReadyInferenceTaskProjection { node_id: pantograph_scheduler::SchedulerNodeId::parse(self.node_id.as_str()) .map_err(|error| { CurrentInferenceSchedulerProjectionError::IncompleteNodeState { @@ -1532,7 +1532,7 @@ impl CurrentInferenceValidationNodeRecord { .dependency_override_fingerprint .clone(), }, - }, + }), )); } diff --git a/crates/pantograph-workflow-service/src/workflow/executable_validation_snapshot.rs b/crates/pantograph-workflow-service/src/workflow/executable_validation_snapshot.rs index d873efe0a..7c35fa013 100644 --- a/crates/pantograph-workflow-service/src/workflow/executable_validation_snapshot.rs +++ b/crates/pantograph-workflow-service/src/workflow/executable_validation_snapshot.rs @@ -492,7 +492,7 @@ impl WorkflowExecutableValidationSnapshotNode { })?; Ok(WorkflowSchedulerInferenceTaskProjection::Ready( - WorkflowSchedulerReadyInferenceTaskProjection { + Box::new(WorkflowSchedulerReadyInferenceTaskProjection { node_id: scheduler_node_id, descriptor_fingerprint: self.descriptor_fingerprint.clone(), task_type, @@ -504,7 +504,7 @@ impl WorkflowExecutableValidationSnapshotNode { dependency_readiness_source: workflow_scheduler_dependency_readiness_source( snapshot, self, )?, - }, + }), )) } } diff --git a/crates/pantograph-workflow-service/src/workflow/task_graph.rs b/crates/pantograph-workflow-service/src/workflow/task_graph.rs index 23c586a8c..89b9e71f5 100644 --- a/crates/pantograph-workflow-service/src/workflow/task_graph.rs +++ b/crates/pantograph-workflow-service/src/workflow/task_graph.rs @@ -64,7 +64,7 @@ impl WorkflowSchedulerInferenceTaskProjections { #[derive(Debug, Clone, PartialEq, Eq)] pub enum WorkflowSchedulerInferenceTaskProjection { - Ready(WorkflowSchedulerReadyInferenceTaskProjection), + Ready(Box), Blocked(WorkflowSchedulerBlockedInferenceTaskProjection), } diff --git a/crates/pantograph-workflow-service/src/workflow/tests/task_binding_resolution.rs b/crates/pantograph-workflow-service/src/workflow/tests/task_binding_resolution.rs index ea7fa28d8..186d08109 100644 --- a/crates/pantograph-workflow-service/src/workflow/tests/task_binding_resolution.rs +++ b/crates/pantograph-workflow-service/src/workflow/tests/task_binding_resolution.rs @@ -34,7 +34,7 @@ fn workflow_run_id() -> WorkflowRunId { fn inference_projection() -> WorkflowSchedulerInferenceTaskProjections { WorkflowSchedulerInferenceTaskProjections::from_records(vec![ WorkflowSchedulerInferenceTaskProjection::Ready( - WorkflowSchedulerReadyInferenceTaskProjection { + Box::new(WorkflowSchedulerReadyInferenceTaskProjection { node_id: SchedulerNodeId::parse("infer").expect("node id"), descriptor_fingerprint: InferenceInterfaceFingerprint::parse("iface.binding.v1") .expect("fingerprint"), @@ -50,7 +50,7 @@ fn inference_projection() -> WorkflowSchedulerInferenceTaskProjections { trait_settings: Vec::new(), estimate_hints: resource_estimate_hints(), dependency_readiness_source: dependency_readiness_source("iface.binding.v1"), - }, + }), ), ]) .expect("projection") diff --git a/crates/pantograph-workflow-service/src/workflow/tests/task_graph.rs b/crates/pantograph-workflow-service/src/workflow/tests/task_graph.rs index 9660df775..ec9a18eaa 100644 --- a/crates/pantograph-workflow-service/src/workflow/tests/task_graph.rs +++ b/crates/pantograph-workflow-service/src/workflow/tests/task_graph.rs @@ -45,7 +45,7 @@ fn ready_inference_projection( ) -> WorkflowSchedulerInferenceTaskProjections { WorkflowSchedulerInferenceTaskProjections::from_records(vec![ WorkflowSchedulerInferenceTaskProjection::Ready( - WorkflowSchedulerReadyInferenceTaskProjection { + Box::new(WorkflowSchedulerReadyInferenceTaskProjection { node_id: SchedulerNodeId::parse("infer").expect("node id"), descriptor_fingerprint: InferenceInterfaceFingerprint::parse("iface.test.v1") .expect("fingerprint"), @@ -70,7 +70,7 @@ fn ready_inference_projection( }], estimate_hints, dependency_readiness_source: dependency_readiness_source("iface.test.v1"), - }, + }), ), ]) .expect("projection") @@ -200,6 +200,26 @@ fn graph_with_inline_inference_ref() -> WorkflowGraph { } } +#[test] +fn ready_projection_preserves_owned_raw_record_and_compact_layout() { + let projections = inference_projection(); + let projection = projections + .get(&SchedulerNodeId::parse("infer").unwrap()) + .unwrap(); + let WorkflowSchedulerInferenceTaskProjection::Ready(borrowed) = projection else { + panic!("expected ready projection"); + }; + let WorkflowSchedulerInferenceTaskProjection::Ready(owned) = projection.clone() else { + panic!("expected cloned ready projection"); + }; + let raw: WorkflowSchedulerReadyInferenceTaskProjection = *owned; + assert_eq!(&raw, borrowed.as_ref()); + assert!( + std::mem::size_of::() + < std::mem::size_of::() + ); +} + #[test] fn scheduler_task_graph_projects_path_free_inference_intent() { let graph = workflow_scheduler_task_graph_with_inference_projections( @@ -295,6 +315,32 @@ fn scheduler_task_graph_projects_path_free_inference_intent() { SchedulerTraitValue::String("euler_discrete".to_string()) ); + let expected_intent = json!({ + "contract_version": 1, + "workflow_id": "workflow-task-graph", + "workflow_run_id": "run-task-graph", + "node_id": "infer", + "task_id": "infer", + "task_type": "image_generation", + "model_ref": { + "model_id": "image/example/tiny-diffusion", + "revision": "main", + "selected_artifact_id": "diffusers-bundle" + }, + "constraints": { "requested_runtime_id": "pytorch", "requested_device_id": "cuda:0" }, + "trait_settings": [{ "trait_id": "denoising_scheduler", "value": { "kind": "string", "value": "euler_discrete" } }], + "estimate_hints": [ + { "kind": "peak_ram_bytes", "value": 2_147_483_648_u64 }, + { "kind": "peak_vram_bytes", "value": 4_294_967_296_u64 } + ] + }); + assert_eq!(serde_json::to_value(intent).unwrap(), expected_intent); + assert_eq!( + serde_json::from_value::(expected_intent) + .unwrap(), + *intent + ); + let encoded = serde_json::to_string(&graph).expect("encode task graph"); assert!(!encoded.contains("model_path")); assert!(!encoded.contains("/tmp/legacy-model")); diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-ready-inference-projection-layout.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-ready-inference-projection-layout.md new file mode 100644 index 000000000..7e0d24373 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-ready-inference-projection-layout.md @@ -0,0 +1,7 @@ +# Public ready inference projection payload + +Fresh PR #38 Clippy retains the public Ready projection at 528 bytes versus the Blocked record at 80. Box only Ready and migrate four existing constructors: current validation state, executable snapshot projection and two workflow fixture helpers. Raw fields, equality, borrowed readers, freshness/admission checks and blocked behavior remain unchanged. The enum has no Serde implementation; the downstream scheduler intent wire contract is the serialization boundary to preserve. + +Tests cover owned raw-record equality and layout, plus exact downstream intent JSON and deserialization while retaining existing path-free projection assertions. CI explicitly discovers and runs task-graph, task-binding and executable-snapshot modules; inherited validation-owner and Headless tests retain ready/stale/missing-estimate/blocked behavior and the embedded-runtime consumer. + +This is a public Rust construction/owned-binding change despite publish=false, documented with the extra ready-record allocation and no measured performance claim. Root approved the bounded design. Root source review accepted all eight files at frozen tree 42349ac0a39dd876e5fca06ca4b285f3cc6eaca8. Hosted execution remains pending; no local Rust execution is claimed. From 0c085394d1f6864cc712b1e934c1c22d41d57f1a Mon Sep 17 00:00:00 2001 From: Puma Date: Sat, 3 Oct 2026 07:37:35 -0700 Subject: [PATCH 3/5] refactor(workflow): group terminal diagnostic inputs --- .github/workflows/quality-gates.yml | 5 + .../runtime_branch_batch_execution.rs | 27 ++-- .../runtime_branch_run_finalization.rs | 148 ++++++++++++++---- .../src/workflow/session_scheduler_runner.rs | 107 +++++++------ .../2026-10-03-terminal-diagnostic-input.md | 9 ++ 5 files changed, 198 insertions(+), 98 deletions(-) create mode 100644 docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-terminal-diagnostic-input.md diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index a28eea8d2..93fcf95c6 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,11 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run terminal diagnostic regression suite + run: | + cargo test -p pantograph-workflow-service --lib workflow::runtime_branch_run_finalization::tests:: -- --list > "$RUNNER_TEMP/terminal-diagnostic-tests.list" + python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith("workflow::runtime_branch_run_finalization::tests::") and line.endswith(": test")]; print("Terminal diagnostic tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/terminal-diagnostic-tests.list" + cargo test -p pantograph-workflow-service --lib workflow::runtime_branch_run_finalization::tests:: - name: Run workflow inference projection regression suites run: | for module in workflow::tests::task_graph:: workflow::tests::task_binding_resolution:: workflow::executable_validation_snapshot::tests::; do diff --git a/crates/pantograph-workflow-service/src/workflow/runtime_branch_batch_execution.rs b/crates/pantograph-workflow-service/src/workflow/runtime_branch_batch_execution.rs index bb2cf341f..be0d264c8 100644 --- a/crates/pantograph-workflow-service/src/workflow/runtime_branch_batch_execution.rs +++ b/crates/pantograph-workflow-service/src/workflow/runtime_branch_batch_execution.rs @@ -24,6 +24,7 @@ use pantograph_scheduler::{ use super::runtime_branch_run_finalization::{ completed_scheduler_run_response, record_scheduler_task_attempt_terminal, + WorkflowSchedulerTaskAttemptTerminalInput, }; use super::runtime_dispatch_assignment::{ WorkflowRuntimeDispatchAssignmentBatchClaim, @@ -439,18 +440,20 @@ where })?; record_scheduler_task_attempt_terminal( service, - application.started_batch_member.started().task(), - application - .started_batch_member - .started() - .attempt_id() - .as_str(), - application.started_batch_member.started().started_at_ms(), - transition, - reason, - error_summary, - Some(application.started_batch_member.selected_dispatch()), - Some(terminal_mutation), + WorkflowSchedulerTaskAttemptTerminalInput { + task: application.started_batch_member.started().task(), + attempt_id: application + .started_batch_member + .started() + .attempt_id() + .as_str(), + started_at_ms: application.started_batch_member.started().started_at_ms(), + transition, + reason, + error_summary, + selected_dispatch: Some(application.started_batch_member.selected_dispatch()), + terminal_mutation: Some(terminal_mutation), + }, ) .map_err(|error| { WorkflowRuntimeBranchBatchExecutionFailure::active_run_member( diff --git a/crates/pantograph-workflow-service/src/workflow/runtime_branch_run_finalization.rs b/crates/pantograph-workflow-service/src/workflow/runtime_branch_run_finalization.rs index dfcc18b15..0289ed26b 100644 --- a/crates/pantograph-workflow-service/src/workflow/runtime_branch_run_finalization.rs +++ b/crates/pantograph-workflow-service/src/workflow/runtime_branch_run_finalization.rs @@ -37,6 +37,18 @@ pub(super) struct WorkflowSchedulerTaskAttemptDiagnosticAttribution { pub(super) bucket_id: Option, } +#[derive(Debug)] +pub(super) struct WorkflowSchedulerTaskAttemptTerminalInput<'a> { + pub(super) task: &'a WorkflowSchedulerTask, + pub(super) attempt_id: &'a str, + pub(super) started_at_ms: u64, + pub(super) transition: SchedulerTaskAttemptLifecycleTransition, + pub(super) reason: &'a str, + pub(super) error_summary: Option, + pub(super) selected_dispatch: Option<&'a SelectedRuntimeTaskDispatch>, + pub(super) terminal_mutation: Option<&'a WorkflowSchedulerTaskTerminalMutation>, +} + #[derive(Debug)] pub(super) struct WorkflowSchedulerTaskAttemptTerminalDiagnosticRequest<'a> { pub(super) task: &'a WorkflowSchedulerTask, @@ -142,14 +154,16 @@ pub(super) async fn finalize_started_runtime_task_dispatch( scheduler_task_attempt_terminal_transition_from_result(&result); record_scheduler_task_attempt_terminal( service, - started_runtime_task.task(), - started_runtime_task.attempt_id().as_str(), - started_runtime_task.started_at_ms(), - transition, - reason, - error_summary, - Some(selected_dispatch), - Some(&terminal_mutation), + WorkflowSchedulerTaskAttemptTerminalInput { + task: started_runtime_task.task(), + attempt_id: started_runtime_task.attempt_id().as_str(), + started_at_ms: started_runtime_task.started_at_ms(), + transition, + reason, + error_summary, + selected_dispatch: Some(selected_dispatch), + terminal_mutation: Some(&terminal_mutation), + }, )?; service .scheduler_task_orchestrator @@ -191,14 +205,16 @@ pub(super) async fn finalize_started_runtime_task_dispatch( }; record_scheduler_task_attempt_terminal( service, - started_runtime_task.task(), - started_runtime_task.attempt_id().as_str(), - started_runtime_task.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Cancelled, - "scheduler runtime task cancellation observed", - Some(message.clone()), - Some(selected_dispatch), - Some(&terminal_mutation), + WorkflowSchedulerTaskAttemptTerminalInput { + task: started_runtime_task.task(), + attempt_id: started_runtime_task.attempt_id().as_str(), + started_at_ms: started_runtime_task.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Cancelled, + reason: "scheduler runtime task cancellation observed", + error_summary: Some(message.clone()), + selected_dispatch: Some(selected_dispatch), + terminal_mutation: Some(&terminal_mutation), + }, )?; service .scheduler_task_orchestrator @@ -234,14 +250,16 @@ pub(super) async fn finalize_started_runtime_task_dispatch( }; record_scheduler_task_attempt_terminal( service, - started_runtime_task.task(), - started_runtime_task.attempt_id().as_str(), - started_runtime_task.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Failed, - "scheduler runtime task dispatch failed", - Some(error.to_string()), - Some(selected_dispatch), - Some(&terminal_mutation), + WorkflowSchedulerTaskAttemptTerminalInput { + task: started_runtime_task.task(), + attempt_id: started_runtime_task.attempt_id().as_str(), + started_at_ms: started_runtime_task.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Failed, + reason: "scheduler runtime task dispatch failed", + error_summary: Some(error.to_string()), + selected_dispatch: Some(selected_dispatch), + terminal_mutation: Some(&terminal_mutation), + }, )?; service .scheduler_task_orchestrator @@ -358,15 +376,18 @@ fn scheduler_task_attempt_terminal_diagnostic_event_at( pub(super) fn record_scheduler_task_attempt_terminal( service: &WorkflowService, - task: &WorkflowSchedulerTask, - attempt_id: &str, - started_at_ms: u64, - transition: SchedulerTaskAttemptLifecycleTransition, - reason: &str, - error_summary: Option, - selected_dispatch: Option<&SelectedRuntimeTaskDispatch>, - terminal_mutation: Option<&WorkflowSchedulerTaskTerminalMutation>, + input: WorkflowSchedulerTaskAttemptTerminalInput<'_>, ) -> Result<(), WorkflowServiceError> { + let WorkflowSchedulerTaskAttemptTerminalInput { + task, + attempt_id, + started_at_ms, + transition, + reason, + error_summary, + selected_dispatch, + terminal_mutation, + } = input; let attribution = scheduler_task_attempt_diagnostic_attribution(service, task.workflow_run_id.as_str())?; service.workflow_diagnostic_event_record(scheduler_task_attempt_terminal_diagnostic_event( @@ -541,6 +562,69 @@ mod tests { use super::*; + #[test] + fn terminal_input_records_success_failure_and_cancellation_scope() { + use pantograph_diagnostics_ledger::{DiagnosticsLedgerRepository, SqliteDiagnosticsLedger}; + + for transition in [ + SchedulerTaskAttemptLifecycleTransition::Completed, + SchedulerTaskAttemptLifecycleTransition::Failed, + SchedulerTaskAttemptLifecycleTransition::Cancelled, + ] { + let service = WorkflowService::new() + .with_diagnostics_ledger(SqliteDiagnosticsLedger::open_in_memory().unwrap()); + let task = runtime_task(); + let started_at_ms = unix_timestamp_ms(); + let error_summary = (transition != SchedulerTaskAttemptLifecycleTransition::Completed) + .then(|| "terminal fixture detail".to_string()); + record_scheduler_task_attempt_terminal( + &service, + WorkflowSchedulerTaskAttemptTerminalInput { + task: &task, + attempt_id: "attempt.terminal.input", + started_at_ms, + transition, + reason: "terminal input fixture", + error_summary: error_summary.clone(), + selected_dispatch: None, + terminal_mutation: None, + }, + ) + .expect("terminal input records event"); + let records = service + .diagnostics_ledger_guard() + .unwrap() + .diagnostic_events_after(0, 10) + .unwrap(); + assert_eq!(records.len(), 1); + let record = &records[0]; + assert_eq!( + record.workflow_run_id.as_ref().unwrap().as_str(), + task.workflow_run_id.as_str() + ); + assert_eq!(record.node_id.as_deref(), Some(task.node_id.as_str())); + assert!(record.client_id.is_none()); + assert!(record.client_session_id.is_none()); + assert!(record.bucket_id.is_none()); + let DiagnosticEventPayload::SchedulerTaskAttemptLifecycleChanged(payload) = + serde_json::from_str(&record.payload_json).unwrap() + else { + panic!("expected task-attempt terminal event"); + }; + assert_eq!(payload.scheduler_task_id, task.task_id.as_str()); + assert_eq!(payload.scheduler_attempt_id, "attempt.terminal.input"); + assert_eq!(payload.transition, transition); + assert_eq!(payload.started_at_ms, Some(started_at_ms as i64)); + assert_eq!( + payload.duration_ms, + Some((payload.ended_at_ms.unwrap() as u64).saturating_sub(started_at_ms)) + ); + assert_eq!(payload.reason.as_deref(), Some("terminal input fixture")); + assert_eq!(payload.error_summary, error_summary); + assert!(payload.reservation_id.is_none()); + } + } + #[test] fn terminal_diagnostic_event_preserves_scheduler_task_attempt_scope() { let task = runtime_task(); diff --git a/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs b/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs index 74fd40b9a..31456e718 100644 --- a/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs +++ b/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs @@ -24,7 +24,7 @@ use crate::scheduler::task_orchestrator::{ }; use crate::scheduler::{ WorkflowDependencyReadinessLifecycle, WorkflowDependencyReadinessLifecycleError, - WorkflowSchedulerRetryLifecycle, WorkflowSchedulerTaskTerminalMutation, + WorkflowSchedulerRetryLifecycle, }; use super::runtime_branch_run_finalization::{ @@ -33,6 +33,7 @@ use super::runtime_branch_run_finalization::{ WorkflowRuntimeTaskDispatchFinalizationOutcome, WorkflowSchedulerTaskAttemptDiagnosticAttribution, WorkflowSchedulerTaskAttemptTerminalDiagnosticRequest, + WorkflowSchedulerTaskAttemptTerminalInput, }; use super::{ WorkflowHost, WorkflowOutputTarget, WorkflowPortBinding, WorkflowRunResponse, @@ -513,17 +514,16 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { "scheduler non-runtime task completion failed: {error}" )) })?; - self.record_scheduler_task_attempt_terminal( - session_id, - started.task(), - started.attempt_id().as_str(), - started.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Completed, - "scheduler task attempt completed", - None, - None, - None, - )?; + self.record_scheduler_task_attempt_terminal(WorkflowSchedulerTaskAttemptTerminalInput { + task: started.task(), + attempt_id: started.attempt_id().as_str(), + started_at_ms: started.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Completed, + reason: "scheduler task attempt completed", + error_summary: None, + selected_dispatch: None, + terminal_mutation: None, + })?; } Err( crate::scheduler::WorkflowSchedulerTaskOrchestratorError::NonRuntimeTaskAdapter( @@ -542,17 +542,16 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { &error, ); if failed.is_ok() { - self.record_scheduler_task_attempt_terminal( - session_id, - started.task(), - started.attempt_id().as_str(), - started.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Failed, - "scheduler non-runtime task execution failed", - Some(error.to_string()), - None, - None, - )?; + self.record_scheduler_task_attempt_terminal(WorkflowSchedulerTaskAttemptTerminalInput { + task: started.task(), + attempt_id: started.attempt_id().as_str(), + started_at_ms: started.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Failed, + reason: "scheduler non-runtime task execution failed", + error_summary: Some(error.to_string()), + selected_dispatch: None, + terminal_mutation: None, + })?; } return Err(WorkflowServiceError::InvalidRequest(format!( "scheduler non-runtime task execution failed: {error}" @@ -886,17 +885,16 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { )) })? }; - self.record_scheduler_task_attempt_terminal( - session_id, - started_runtime_task.task(), - started_runtime_task.attempt_id().as_str(), - started_runtime_task.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Failed, - "scheduler runtime dispatch selection failed", - Some(scheduler_error.to_string()), - None, - Some(&terminal_mutation), - )?; + self.record_scheduler_task_attempt_terminal(WorkflowSchedulerTaskAttemptTerminalInput { + task: started_runtime_task.task(), + attempt_id: started_runtime_task.attempt_id().as_str(), + started_at_ms: started_runtime_task.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Failed, + reason: "scheduler runtime dispatch selection failed", + error_summary: Some(scheduler_error.to_string()), + selected_dispatch: None, + terminal_mutation: Some(&terminal_mutation), + })?; } else { let terminal_mutation = { let mut store = self.service.session_store_guard()?; @@ -915,17 +913,16 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { )) })? }; - self.record_scheduler_task_attempt_terminal( - session_id, - started_runtime_task.task(), - started_runtime_task.attempt_id().as_str(), - started_runtime_task.started_at_ms(), - SchedulerTaskAttemptLifecycleTransition::Failed, - "scheduler runtime dispatch failed", - Some(scheduler_error.to_string()), - None, - Some(&terminal_mutation), - )?; + self.record_scheduler_task_attempt_terminal(WorkflowSchedulerTaskAttemptTerminalInput { + task: started_runtime_task.task(), + attempt_id: started_runtime_task.attempt_id().as_str(), + started_at_ms: started_runtime_task.started_at_ms(), + transition: SchedulerTaskAttemptLifecycleTransition::Failed, + reason: "scheduler runtime dispatch failed", + error_summary: Some(scheduler_error.to_string()), + selected_dispatch: None, + terminal_mutation: Some(&terminal_mutation), + })?; } return Err(WorkflowServiceError::CapabilityViolation(format!( "runtime scheduler dispatch selection failed closed for {count} runtime inference task(s): {scheduler_error}", @@ -1101,16 +1098,18 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { fn record_scheduler_task_attempt_terminal( &self, - _session_id: &str, - task: &WorkflowSchedulerTask, - attempt_id: &str, - started_at_ms: u64, - transition: SchedulerTaskAttemptLifecycleTransition, - reason: &str, - error_summary: Option, - selected_dispatch: Option<&SelectedRuntimeTaskDispatch>, - terminal_mutation: Option<&WorkflowSchedulerTaskTerminalMutation>, + input: WorkflowSchedulerTaskAttemptTerminalInput<'_>, ) -> Result<(), WorkflowServiceError> { + let WorkflowSchedulerTaskAttemptTerminalInput { + task, + attempt_id, + started_at_ms, + transition, + reason, + error_summary, + selected_dispatch, + terminal_mutation, + } = input; let attribution = self.scheduler_task_attempt_diagnostic_attribution(task.workflow_run_id.as_str())?; self.service.workflow_diagnostic_event_record( diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-terminal-diagnostic-input.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-terminal-diagnostic-input.md new file mode 100644 index 000000000..88a33f02f --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-terminal-diagnostic-input.md @@ -0,0 +1,9 @@ +# Shared private terminal diagnostic input + +Fresh PR #39 Clippy retains high-arity terminal diagnostic recorders in the session runner and runtime finalization owner. Group exactly eight existing values into private WorkflowSchedulerTaskAttemptTerminalInput: task, attempt_id, started_at_ms, transition, reason, owned error_summary, borrowed selected_dispatch and borrowed terminal_mutation. Migrate eight callers: four runner paths, three finalization paths and one batch path. + +Remove the runner's proven-unused _session_id argument. All four removed expressions were a plain session_id local with no side effects, and the parameter performed no validation. Both recorder bodies after destructuring remain byte-identical, including task.workflow_run_id attribution lookup and the unchanged lower event builder. No attribution-policy or public API change is made. + +A real in-memory ledger regression sends the new input through the recorder for success, failure and cancellation, verifying run/node/task/attempt scope, absent attribution, timing consistency, reason and original error detail. Existing precise-timing and incomplete-run tests remain, with nonzero module discovery in hosted CI. Source comparison, changed-call syntax formatting and staged whitespace pass locally. Root approved the design; source review and actual hosted execution remain pending. + +Root initial review confirmed the field/caller mapping and recorder-body preservation. The narrow follow-up removes six redundant field labels and the now-unused runner terminal-mutation import; the new duration assertion uses saturating subtraction. The production event builder remains unchanged and rejects a terminal timestamp preceding the start. Root narrow re-review accepted frozen tree a79c7c43eb66667ec8c0e97c5fc83e5dc9e93ab4. Hosted execution remains pending. From 48b728c5cc73d0227a3c596a7a8eab58308064f9 Mon Sep 17 00:00:00 2001 From: Puma Date: Sat, 3 Oct 2026 07:47:14 -0700 Subject: [PATCH 4/5] refactor(workflow): group private session run context --- .github/workflows/quality-gates.yml | 5 ++ .../src/workflow/session_execution_api.rs | 48 ++++++----- .../src/workflow/session_scheduler_runner.rs | 84 ++++++++++--------- .../src/workflow/task_execution_owner.rs | 42 +++++++--- .../reports/2026-10-03-session-run-context.md | 7 ++ 5 files changed, 112 insertions(+), 74 deletions(-) create mode 100644 docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-session-run-context.md diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index 93fcf95c6..a6074e552 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,11 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run complete session execution regression suite + run: | + cargo test -p pantograph-workflow-service --lib workflow::tests::session_execution:: -- --list > "$RUNNER_TEMP/session-execution-tests.list" + python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith("workflow::tests::session_execution::") and line.endswith(": test")]; print("Session execution tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/session-execution-tests.list" + cargo test -p pantograph-workflow-service --lib workflow::tests::session_execution:: - name: Run terminal diagnostic regression suite run: | cargo test -p pantograph-workflow-service --lib workflow::runtime_branch_run_finalization::tests:: -- --list > "$RUNNER_TEMP/terminal-diagnostic-tests.list" diff --git a/crates/pantograph-workflow-service/src/workflow/session_execution_api.rs b/crates/pantograph-workflow-service/src/workflow/session_execution_api.rs index 5e0216ef5..8b6727ed0 100644 --- a/crates/pantograph-workflow-service/src/workflow/session_execution_api.rs +++ b/crates/pantograph-workflow-service/src/workflow/session_execution_api.rs @@ -48,8 +48,10 @@ use super::runtime_dispatch_assignment::{ WorkflowRuntimeDispatchAssignmentRepository, }; use super::session_io_artifacts::workflow_io_artifact_metadata; -use super::session_scheduler_runner::WorkflowSchedulerSessionRunner; -use super::task_execution_owner::WorkflowTaskExecutionOwner; +use super::session_scheduler_runner::{ + WorkflowSchedulerRunContext, WorkflowSchedulerSessionRunner, +}; +use super::task_execution_owner::{WorkflowNonRuntimeExecutionInput, WorkflowTaskExecutionOwner}; use super::task_execution_runtime::WorkflowTaskExecutionRuntimeOwner; use super::task_execution_worker::{ WorkflowTaskExecutionWorkerOutcome, WorkflowTaskExecutionWorkerRuntimeBranchCommand, @@ -415,12 +417,14 @@ impl WorkflowService { return WorkflowTaskExecutionOwner::run_non_runtime_to_completion( self, host, - &session, - run_snapshot.as_ref(), - &session_id, - &workflow_run_id, - &queued_run, - &scheduler_task_run_summary, + WorkflowNonRuntimeExecutionInput { + session: &session, + run_snapshot: run_snapshot.as_ref(), + session_id: &session_id, + workflow_run_id: &workflow_run_id, + queued_run: &queued_run, + summary: &scheduler_task_run_summary, + }, ) .await; } @@ -565,12 +569,14 @@ impl WorkflowService { Ok(Some(timeout_ms)) => { let run_future = runner.resume_runtime_dependency_readiness( host, - &session_id, - &workflow_run_id, - &active_run.workflow_id, - active_run.output_targets.as_deref(), - &scheduler_task_run_summary, - started_at, + WorkflowSchedulerRunContext { + session_id: &session_id, + workflow_run_id: &workflow_run_id, + workflow_id: &active_run.workflow_id, + output_targets: active_run.output_targets.as_deref(), + summary: &scheduler_task_run_summary, + started_at, + }, attempt_start_transition, ); match tokio::time::timeout(Duration::from_millis(timeout_ms), run_future).await { @@ -585,12 +591,14 @@ impl WorkflowService { runner .resume_runtime_dependency_readiness( host, - &session_id, - &workflow_run_id, - &active_run.workflow_id, - active_run.output_targets.as_deref(), - &scheduler_task_run_summary, - started_at, + WorkflowSchedulerRunContext { + session_id: &session_id, + workflow_run_id: &workflow_run_id, + workflow_id: &active_run.workflow_id, + output_targets: active_run.output_targets.as_deref(), + summary: &scheduler_task_run_summary, + started_at, + }, attempt_start_transition, ) .await diff --git a/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs b/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs index 31456e718..e89b18850 100644 --- a/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs +++ b/crates/pantograph-workflow-service/src/workflow/session_scheduler_runner.rs @@ -231,6 +231,16 @@ impl WorkflowPreDispatchPreparationOutcome { } } +#[derive(Debug, Clone, Copy)] +pub(super) struct WorkflowSchedulerRunContext<'a> { + pub(super) session_id: &'a str, + pub(super) workflow_run_id: &'a str, + pub(super) workflow_id: &'a str, + pub(super) output_targets: Option<&'a [WorkflowOutputTarget]>, + pub(super) summary: &'a WorkflowSchedulerTaskRunSummary, + pub(super) started_at: Instant, +} + struct ReadyRuntimeDispatchContext { task: WorkflowSchedulerTask, ready_record: SchedulerTaskStateRecord, @@ -244,14 +254,17 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { pub(super) async fn run_non_runtime_only( &self, host: &H, - session_id: &str, - workflow_run_id: &str, - workflow_id: &str, inputs: &[WorkflowPortBinding], - output_targets: Option<&[WorkflowOutputTarget]>, - summary: &WorkflowSchedulerTaskRunSummary, - started_at: Instant, + context: WorkflowSchedulerRunContext<'_>, ) -> Result { + let WorkflowSchedulerRunContext { + session_id, + workflow_run_id, + workflow_id, + output_targets, + summary, + started_at, + } = context; if !summary.is_non_runtime_only() || summary.has_runtime_inference() { return Err(WorkflowServiceError::Internal( "scheduler session runner received a runtime-containing run".to_string(), @@ -278,31 +291,22 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { pub(super) async fn resume_runtime_dependency_readiness( &self, host: &(impl WorkflowHost + ?Sized), - session_id: &str, - workflow_run_id: &str, - workflow_id: &str, - output_targets: Option<&[WorkflowOutputTarget]>, - summary: &WorkflowSchedulerTaskRunSummary, - started_at: Instant, + context: WorkflowSchedulerRunContext<'_>, attempt_start_transition: SchedulerTaskAttemptLifecycleTransition, ) -> Result { + let WorkflowSchedulerRunContext { + workflow_run_id, + summary, + .. + } = context; if !summary.has_runtime_inference() { return Err(WorkflowServiceError::InvalidRequest(format!( "workflow run '{}' is not a runtime inference run", workflow_run_id ))); } - self.continue_runtime_dependency_readiness( - host, - session_id, - workflow_run_id, - workflow_id, - output_targets, - summary, - started_at, - attempt_start_transition, - ) - .await + self.continue_runtime_dependency_readiness(host, context, attempt_start_transition) + .await } pub(super) async fn resume_progress_loop( @@ -318,14 +322,14 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { async fn continue_runtime_dependency_readiness( &self, host: &(impl WorkflowHost + ?Sized), - session_id: &str, - workflow_run_id: &str, - workflow_id: &str, - output_targets: Option<&[WorkflowOutputTarget]>, - summary: &WorkflowSchedulerTaskRunSummary, - started_at: Instant, + context: WorkflowSchedulerRunContext<'_>, attempt_start_transition: SchedulerTaskAttemptLifecycleTransition, ) -> Result { + let WorkflowSchedulerRunContext { + session_id, + workflow_run_id, + .. + } = context; let preparation = WorkflowPreDispatchPreparationBoundary::new(self.service) .prepare_runtime_dispatch(session_id, workflow_run_id) .await?; @@ -341,12 +345,7 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { } self.run_runtime_dispatch_ready_tasks( host, - session_id, - workflow_run_id, - workflow_id, - output_targets, - summary, - started_at, + context, preparation.admitted_runtime_readiness(), attempt_start_transition, ) @@ -804,15 +803,18 @@ impl<'a> WorkflowSchedulerSessionRunner<'a> { async fn run_runtime_dispatch_ready_tasks( &self, host: &(impl WorkflowHost + ?Sized), - session_id: &str, - workflow_run_id: &str, - workflow_id: &str, - output_targets: Option<&[WorkflowOutputTarget]>, - summary: &WorkflowSchedulerTaskRunSummary, - started_at: Instant, + context: WorkflowSchedulerRunContext<'_>, admitted_runtime_readiness: &[AdmittedRuntimeTaskReadiness], attempt_start_transition: SchedulerTaskAttemptLifecycleTransition, ) -> Result { + let WorkflowSchedulerRunContext { + session_id, + workflow_run_id, + workflow_id, + output_targets, + summary, + started_at, + } = context; let runtime_task_ids = runtime_task_ids_in_state(self.service, session_id, workflow_run_id, |kind| { kind == SchedulerTaskStateKind::Ready diff --git a/crates/pantograph-workflow-service/src/workflow/task_execution_owner.rs b/crates/pantograph-workflow-service/src/workflow/task_execution_owner.rs index e449be161..9f417c342 100644 --- a/crates/pantograph-workflow-service/src/workflow/task_execution_owner.rs +++ b/crates/pantograph-workflow-service/src/workflow/task_execution_owner.rs @@ -3,7 +3,9 @@ use std::time::Duration; use crate::scheduler::WorkflowExecutionSessionDequeuedRun; use pantograph_runtime_attribution::WorkflowRunSnapshotRecord; -use super::session_scheduler_runner::WorkflowSchedulerSessionRunner; +use super::session_scheduler_runner::{ + WorkflowSchedulerRunContext, WorkflowSchedulerSessionRunner, +}; use super::workflow_run_finalization::{ finalize_admitted_workflow_run, WorkflowRunFinalizationRequest, }; @@ -12,6 +14,15 @@ use super::{ WorkflowSchedulerTaskRunSummary, WorkflowService, WorkflowServiceError, }; +pub(super) struct WorkflowNonRuntimeExecutionInput<'a> { + pub(super) session: &'a WorkflowExecutionSessionSummary, + pub(super) run_snapshot: Option<&'a WorkflowRunSnapshotRecord>, + pub(super) session_id: &'a str, + pub(super) workflow_run_id: &'a str, + pub(super) queued_run: &'a WorkflowExecutionSessionDequeuedRun, + pub(super) summary: &'a WorkflowSchedulerTaskRunSummary, +} + pub(super) struct WorkflowTaskExecutionOwner; impl WorkflowTaskExecutionOwner { @@ -26,13 +37,16 @@ impl WorkflowTaskExecutionOwner { pub(super) async fn run_non_runtime_to_completion( service: &WorkflowService, host: &H, - session: &WorkflowExecutionSessionSummary, - run_snapshot: Option<&WorkflowRunSnapshotRecord>, - session_id: &str, - workflow_run_id: &str, - queued_run: &WorkflowExecutionSessionDequeuedRun, - summary: &WorkflowSchedulerTaskRunSummary, + input: WorkflowNonRuntimeExecutionInput<'_>, ) -> Result { + let WorkflowNonRuntimeExecutionInput { + session, + run_snapshot, + session_id, + workflow_run_id, + queued_run, + summary, + } = input; service.record_run_started_event_if_configured(session, run_snapshot, queued_run)?; let run_started_at = std::time::Instant::now(); let queued_workflow_semantic_version = queued_run.queued.workflow_semantic_version.clone(); @@ -40,13 +54,15 @@ impl WorkflowTaskExecutionOwner { let runner = WorkflowSchedulerSessionRunner::new(service); let run_future = runner.run_non_runtime_only( host, - session_id, - workflow_run_id, - &queued_run.workflow_id, &queued_run.queued.inputs, - queued_run.queued.output_targets.as_deref(), - summary, - run_started_at, + WorkflowSchedulerRunContext { + session_id, + workflow_run_id, + workflow_id: &queued_run.workflow_id, + output_targets: queued_run.queued.output_targets.as_deref(), + summary, + started_at: run_started_at, + }, ); let run_result = if let Some(timeout_ms) = queued_run.queued.timeout_ms { match tokio::time::timeout(Duration::from_millis(timeout_ms), run_future).await { diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-session-run-context.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-session-run-context.md new file mode 100644 index 000000000..dd5d1cd58 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-session-run-context.md @@ -0,0 +1,7 @@ +# Private session runner and non-runtime execution inputs + +Group the six values shared by four high-arity runner methods into private WorkflowSchedulerRunContext: borrowed session_id, workflow_run_id, workflow_id, output_targets and summary, plus the original Instant. Keep host, inputs, admitted readiness and attempt transition separate. There are five runner calls: one non-runtime entry, two resume entries and two internal forwards. The Copy context contains only existing borrowed values and Instant; no allocation, new clock read or cloned run data is introduced. + +The sole WorkflowTaskExecutionOwner caller receives a separate six-field WorkflowNonRuntimeExecutionInput: session, run_snapshot, session_id, workflow_run_id, queued_run and summary. Existing finalization requests carry extra result/policy values and are not reused. All existing service/host parameters, borrowing and result signatures remain. Timeout, cancellation cleanup, finalization, no-load behavior and runtime readiness decisions are unchanged. + +Programmatic source comparisons confirm unchanged non-runtime execution/finalization, ready-dispatch body and owner timeout/finalization block. CI explicitly discovers and runs the entire session-execution module, including non-runtime lifecycle, timeout, readiness defer/resume, bootstrap recovery and no-runtime-load cases, with a nonzero guard. Root approved this bounded design. Whitespace/changed-call syntax checks pass locally; root source review accepted all five files at frozen tree 66a25e87fe897ea1001622c0768a51cc51d07882; hosted execution remains pending. From 533a80b9650cff2c1b4f9a8b9087760139fe89c9 Mon Sep 17 00:00:00 2001 From: Puma Date: Sat, 3 Oct 2026 07:52:07 -0700 Subject: [PATCH 5/5] refactor(workflow): clarify private projection and summary types --- .github/workflows/quality-gates.yml | 5 ++ .../src/workflow/task_graph.rs | 16 ++++--- .../src/workflow/task_run_summary.rs | 46 +++++++++++++++---- .../2026-10-03-workflow-private-type-names.md | 7 +++ 4 files changed, 59 insertions(+), 15 deletions(-) create mode 100644 docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-type-names.md diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index a6074e552..52fb41d4b 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,11 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run task summary contract tests + run: | + cargo test -p pantograph-workflow-service --lib workflow::task_run_summary::tests:: -- --list > "$RUNNER_TEMP/task-summary-tests.list" + python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith("workflow::task_run_summary::tests::") and line.endswith(": test")]; print("Task summary tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/task-summary-tests.list" + cargo test -p pantograph-workflow-service --lib workflow::task_run_summary::tests:: - name: Run complete session execution regression suite run: | cargo test -p pantograph-workflow-service --lib workflow::tests::session_execution:: -- --list > "$RUNNER_TEMP/session-execution-tests.list" diff --git a/crates/pantograph-workflow-service/src/workflow/task_graph.rs b/crates/pantograph-workflow-service/src/workflow/task_graph.rs index 89b9e71f5..267a1169e 100644 --- a/crates/pantograph-workflow-service/src/workflow/task_graph.rs +++ b/crates/pantograph-workflow-service/src/workflow/task_graph.rs @@ -259,6 +259,14 @@ fn dependency_task_ids(bindings: &[WorkflowSchedulerTaskInputBinding]) -> Vec, + Option, + Option, + Option, + Vec, +); + fn schedulable_intent_for_node( workflow_id: &SchedulerWorkflowId, workflow_run_id: &SchedulerWorkflowRunId, @@ -266,13 +274,7 @@ fn schedulable_intent_for_node( task_id: &SchedulerTaskId, execution_class: WorkflowSchedulerTaskExecutionClass, inference_task_projection: Option<&WorkflowSchedulerInferenceTaskProjection>, -) -> ( - Option, - Option, - Option, - Option, - Vec, -) { +) -> SchedulableIntentProjection { if execution_class != WorkflowSchedulerTaskExecutionClass::RuntimeInference { return (None, None, None, None, Vec::new()); } diff --git a/crates/pantograph-workflow-service/src/workflow/task_run_summary.rs b/crates/pantograph-workflow-service/src/workflow/task_run_summary.rs index cbad3add0..e085ba2ae 100644 --- a/crates/pantograph-workflow-service/src/workflow/task_run_summary.rs +++ b/crates/pantograph-workflow-service/src/workflow/task_run_summary.rs @@ -54,7 +54,7 @@ pub(crate) fn workflow_scheduler_task_run_summary( .iter() .find(|record| record.task_id.as_str() == task_id) else { - return Err(WorkflowSchedulerTaskRunSummaryError::MissingTaskState { + return Err(WorkflowSchedulerTaskRunSummaryError::Missing { task_id: task_id.to_string(), }); }; @@ -62,7 +62,7 @@ pub(crate) fn workflow_scheduler_task_run_summary( || record.workflow_run_id.as_str() != task.workflow_run_id.as_str() || record.node_id.as_str() != task.node_id.as_str() { - return Err(WorkflowSchedulerTaskRunSummaryError::MismatchedTaskState { + return Err(WorkflowSchedulerTaskRunSummaryError::Mismatched { task_id: task_id.to_string(), }); } @@ -95,7 +95,7 @@ pub(crate) fn workflow_scheduler_task_run_summary( } if let Some(extra_task_id) = record_task_ids.into_iter().next() { - return Err(WorkflowSchedulerTaskRunSummaryError::UnexpectedTaskState { + return Err(WorkflowSchedulerTaskRunSummaryError::Unexpected { task_id: extra_task_id.to_string(), }); } @@ -107,11 +107,11 @@ pub(crate) fn workflow_scheduler_task_run_summary( #[non_exhaustive] pub(crate) enum WorkflowSchedulerTaskRunSummaryError { #[error("scheduler task '{task_id}' has no active task-state record")] - MissingTaskState { task_id: String }, + Missing { task_id: String }, #[error("active task-state record for scheduler task '{task_id}' has mismatched correlation")] - MismatchedTaskState { task_id: String }, + Mismatched { task_id: String }, #[error("active task-state record exists for unknown scheduler task '{task_id}'")] - UnexpectedTaskState { task_id: String }, + Unexpected { task_id: String }, } #[cfg(test)] @@ -269,10 +269,36 @@ mod tests { )]); let error = workflow_scheduler_task_run_summary(&graph, &[]).expect_err("missing record"); + assert_eq!( + error.to_string(), + "scheduler task 'prompt' has no active task-state record" + ); + + assert_eq!( + error, + WorkflowSchedulerTaskRunSummaryError::Missing { + task_id: "prompt".to_string() + } + ); + } + #[test] + fn rejects_mismatched_task_state_with_unchanged_display() { + let graph = task_graph(&[( + "prompt", + WorkflowSchedulerTaskExecutionClass::NonRuntimeNodeEngine, + )]); + let mut mismatched = record("prompt", awaiting_inputs()); + mismatched.workflow_run_id = SchedulerWorkflowRunId::parse("other-run").unwrap(); + let error = workflow_scheduler_task_run_summary(&graph, &[mismatched]) + .expect_err("mismatched correlation"); + assert_eq!( + error.to_string(), + "active task-state record for scheduler task 'prompt' has mismatched correlation" + ); assert_eq!( error, - WorkflowSchedulerTaskRunSummaryError::MissingTaskState { + WorkflowSchedulerTaskRunSummaryError::Mismatched { task_id: "prompt".to_string() } ); @@ -293,10 +319,14 @@ mod tests { ], ) .expect_err("unexpected record"); + assert_eq!( + error.to_string(), + "active task-state record exists for unknown scheduler task 'extra'" + ); assert_eq!( error, - WorkflowSchedulerTaskRunSummaryError::UnexpectedTaskState { + WorkflowSchedulerTaskRunSummaryError::Unexpected { task_id: "extra".to_string() } ); diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-type-names.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-type-names.md new file mode 100644 index 000000000..292c07e47 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-workflow-private-type-names.md @@ -0,0 +1,7 @@ +# Private projection tuple and summary error names + +Name the existing five-element schedulable projection tuple with the private SchedulableIntentProjection alias; tuple members, order, function body and caller destructuring are unchanged. Shorten only the crate-private summary error variants from MissingTaskState/MismatchedTaskState/UnexpectedTaskState to Missing/Mismatched/Unexpected, retaining every task_id field and exact thiserror Display string. Neither type has a new public or serialized representation. + +The source inventory finds two production summary consumers, both using Display interpolation before admission/resume; all variant matches are local tests. Derived Debug variant labels intentionally change with the private names and may appear in test failures, but no production Debug consumer was found. Real missing/unexpected summary tests now assert exact Display, and a real mismatched-run correlation test covers the third branch. CI runs the complete summary module with nonzero discovery, retaining projection tests for the tuple alias. + +Root approved this bounded cleanup. Root source review accepted all four files at frozen tree c3df7026f434651b5b57876d967110d55f3150c6. Hosted execution remains pending. No local Rust execution is claimed.