Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions .github/workflows/quality-gates.yml
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,35 @@ 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"
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"
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
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
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"
Expand Down
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<IoArtifactObservedPayload>`. 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<WorkflowSchedulerReadyInferenceTaskProjection>`. 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`.
Original file line number Diff line number Diff line change
Expand Up @@ -1158,7 +1158,7 @@ pub(crate) enum DependencyEnvironmentActionIntentStateResolution {
Blocked(DependencyEnvironmentActionIntentResult),
RequestReady {
intent: DependencyEnvironmentActionIntent,
environment_request: ValidatedDependencyEnvironmentRequest,
environment_request: Box<ValidatedDependencyEnvironmentRequest>,
},
}

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -1532,7 +1532,7 @@ impl CurrentInferenceValidationNodeRecord {
.dependency_override_fingerprint
.clone(),
},
},
}),
));
}

Expand Down Expand Up @@ -1846,7 +1846,7 @@ fn request_ready_dependency_environment_action_resolution(
) -> DependencyEnvironmentActionIntentStateResolution {
DependencyEnvironmentActionIntentStateResolution::RequestReady {
intent,
environment_request,
environment_request: Box::new(environment_request),
}
}

Expand Down Expand Up @@ -1877,6 +1877,15 @@ mod tests {

use super::*;

#[test]
fn action_resolution_keeps_validated_environment_request_indirect() {
assert!(std::mem::size_of::<DependencyEnvironmentActionIntentStateResolution>() <= 384);
assert!(
std::mem::size_of::<DependencyEnvironmentActionIntentStateResolution>()
< std::mem::size_of::<ValidatedDependencyEnvironmentRequest>()
);
}

#[tokio::test]
async fn action_intent_state_rejects_stale_graph_revision() {
let store = CurrentInferenceValidationStateStore::new();
Expand Down Expand Up @@ -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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -504,7 +504,7 @@ impl WorkflowExecutableValidationSnapshotNode {
dependency_readiness_source: workflow_scheduler_dependency_readiness_source(
snapshot, self,
)?,
},
}),
))
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,18 @@ pub(super) struct WorkflowSchedulerTaskAttemptDiagnosticAttribution {
pub(super) bucket_id: Option<BucketId>,
}

#[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<String>,
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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<String>,
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(
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading