diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index 2edc2cb9d..8a0d02108 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -251,6 +251,8 @@ jobs: - name: Run workflow-nodes tests run: cargo test -p workflow-nodes --lib + - name: Run embedded runtime unit tests + run: cargo test -p pantograph-embedded-runtime --lib - name: Run runtime registry contract tests run: cargo test -p pantograph-runtime-registry - name: Run task summary contract tests diff --git a/crates/pantograph-embedded-runtime/src/dependency_inventory.rs b/crates/pantograph-embedded-runtime/src/dependency_inventory.rs index ebb0b0652..d4cfb8a31 100644 --- a/crates/pantograph-embedded-runtime/src/dependency_inventory.rs +++ b/crates/pantograph-embedded-runtime/src/dependency_inventory.rs @@ -48,6 +48,14 @@ use crate::dependency_inventory_system_package_source::SystemPackageProviderSour use crate::package_readiness_provider::PackageReadinessProbeRunner; use crate::python_package_readiness_probe::ProcessPythonPackageReadinessProbeRunner; +/// Provider-owned diagnostic details; request attribution remains with the observer. +#[cfg(any(test, feature = "standalone"))] +pub(crate) struct DependencyInventoryDiagnosticInput { + pub code: pantograph_dependency_planning::DependencyPlanningDiagnosticCode, + pub message: String, + pub field_path: &'static str, +} + /// Request context passed from the snapshot producer to dependency inventory. #[derive(Debug, Clone)] pub(crate) struct DependencyInventoryRequest { diff --git a/crates/pantograph-embedded-runtime/src/dependency_inventory_device_toolchain.rs b/crates/pantograph-embedded-runtime/src/dependency_inventory_device_toolchain.rs index 6455dee2d..e228458eb 100644 --- a/crates/pantograph-embedded-runtime/src/dependency_inventory_device_toolchain.rs +++ b/crates/pantograph-embedded-runtime/src/dependency_inventory_device_toolchain.rs @@ -16,7 +16,8 @@ use pantograph_dependency_planning::{ }; use crate::dependency_inventory::{ - DependencyInventoryObservation, DependencyInventoryProvider, DependencyInventoryRequest, + DependencyInventoryDiagnosticInput, DependencyInventoryObservation, + DependencyInventoryProvider, DependencyInventoryRequest, }; use crate::dependency_inventory_device_toolchain_source::DeviceToolchainProviderSource; @@ -128,9 +129,12 @@ fn observe_device_toolchain_binding( binding.binding_id.clone(), DependencyInventoryObservationState::Missing, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::ArtifactMissing, - "Device-toolchain source facts are missing for the requested toolchain.", - "dependency_environment.device_toolchain.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactMissing, + message: "Device-toolchain source facts are missing for the requested toolchain." + .into(), + field_path: "dependency_environment.device_toolchain.source", + }, ready_alternatives(&snapshot.rows), ); } @@ -225,9 +229,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Unavailable, DependencyEnvironmentValidationState::Stale, - DependencyPlanningDiagnosticCode::ArtifactStale, - "Device-toolchain source facts are stale.", - "dependency_environment.device_toolchain.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactStale, + message: "Device-toolchain source facts are stale.".into(), + field_path: "dependency_environment.device_toolchain.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Ready, _) => row( @@ -242,9 +248,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Missing, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::ArtifactMissing, - "Device-toolchain source facts are missing.", - "dependency_environment.device_toolchain.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactMissing, + message: "Device-toolchain source facts are missing.".into(), + field_path: "dependency_environment.device_toolchain.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Failed, _) => observation_with_diagnostic( @@ -252,9 +260,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Failed, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::RuntimeUnavailable, - "Device-toolchain source reported a failure.", - "dependency_environment.device_toolchain.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::RuntimeUnavailable, + message: "Device-toolchain source reported a failure.".into(), + field_path: "dependency_environment.device_toolchain.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Unsupported, _) @@ -267,9 +277,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Unavailable, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::RuntimeUnavailable, - "Device-toolchain source is not ready for the requested toolchain.", - "dependency_environment.device_toolchain.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::RuntimeUnavailable, + message: "Device-toolchain source is not ready for the requested toolchain.".into(), + field_path: "dependency_environment.device_toolchain.source", + }, source_row.alternatives.clone(), ), } @@ -336,9 +348,11 @@ fn invalid_row( binding_id, DependencyInventoryObservationState::Invalid, DependencyEnvironmentValidationState::Invalid, - DependencyPlanningDiagnosticCode::InvalidRequest, - message, - field_path, + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::InvalidRequest, + message: message.into(), + field_path, + }, Vec::new(), ) } @@ -348,12 +362,15 @@ fn observation_with_diagnostic( binding_id: DependencyBindingId, state: DependencyInventoryObservationState, validation_state: DependencyEnvironmentValidationState, - code: DependencyPlanningDiagnosticCode, - message: impl Into, - field_path: &'static str, + diagnostic_input: DependencyInventoryDiagnosticInput, alternatives: Vec, ) -> DependencyInventoryObservationRow { - let diagnostic = diagnostic(item, code, message.into(), field_path); + let DependencyInventoryDiagnosticInput { + code, + message, + field_path, + } = diagnostic_input; + let diagnostic = diagnostic(item, code, message, field_path); row( binding_id, state, diff --git a/crates/pantograph-embedded-runtime/src/dependency_inventory_dispatch.rs b/crates/pantograph-embedded-runtime/src/dependency_inventory_dispatch.rs index 5e6894a0d..64908d305 100644 --- a/crates/pantograph-embedded-runtime/src/dependency_inventory_dispatch.rs +++ b/crates/pantograph-embedded-runtime/src/dependency_inventory_dispatch.rs @@ -360,7 +360,7 @@ impl DependencyInventoryProvider for NotImplementedDependencyInventoryProvider { let failure = PackageReadinessProbeFailure::new( PackageReadinessProviderDiagnosticCode::ProbeNotImplemented, None, - CapabilityAvailabilityReason::parse(¬_implemented_reason(&request.payload)) + CapabilityAvailabilityReason::parse(not_implemented_reason(&request.payload)) .expect("inventory provider not implemented reason is valid"), ); let (rows, diagnostics) = diff --git a/crates/pantograph-embedded-runtime/src/dependency_inventory_system_package.rs b/crates/pantograph-embedded-runtime/src/dependency_inventory_system_package.rs index 61686a8d6..526b698ea 100644 --- a/crates/pantograph-embedded-runtime/src/dependency_inventory_system_package.rs +++ b/crates/pantograph-embedded-runtime/src/dependency_inventory_system_package.rs @@ -16,7 +16,8 @@ use pantograph_dependency_planning::{ }; use crate::dependency_inventory::{ - DependencyInventoryObservation, DependencyInventoryProvider, DependencyInventoryRequest, + DependencyInventoryDiagnosticInput, DependencyInventoryObservation, + DependencyInventoryProvider, DependencyInventoryRequest, }; use crate::dependency_inventory_system_package_source::{ SystemPackageProviderSource, SystemPackageProviderSourceError, @@ -135,9 +136,12 @@ fn observe_system_package_binding( binding.binding_id.clone(), DependencyInventoryObservationState::Missing, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::ArtifactMissing, - "System-package source facts are missing for the requested package.", - "dependency_environment.system_package.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactMissing, + message: "System-package source facts are missing for the requested package." + .into(), + field_path: "dependency_environment.system_package.source", + }, ready_alternatives(&snapshot.rows), ); } @@ -240,9 +244,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Unavailable, DependencyEnvironmentValidationState::Stale, - DependencyPlanningDiagnosticCode::ArtifactStale, - "System-package source facts are stale.", - "dependency_environment.system_package.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactStale, + message: "System-package source facts are stale.".into(), + field_path: "dependency_environment.system_package.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Ready, _) => row( @@ -257,9 +263,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Missing, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::ArtifactMissing, - "System-package source facts are missing.", - "dependency_environment.system_package.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::ArtifactMissing, + message: "System-package source facts are missing.".into(), + field_path: "dependency_environment.system_package.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Failed, _) => observation_with_diagnostic( @@ -267,9 +275,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Failed, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::RuntimeUnavailable, - "System-package source reported a failure.", - "dependency_environment.system_package.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::RuntimeUnavailable, + message: "System-package source reported a failure.".into(), + field_path: "dependency_environment.system_package.source", + }, source_row.alternatives.clone(), ), (DependencyProviderSourceState::Unsupported, _) @@ -282,9 +292,11 @@ fn observation_from_source_row( binding_id, DependencyInventoryObservationState::Unavailable, DependencyEnvironmentValidationState::Valid, - DependencyPlanningDiagnosticCode::RuntimeUnavailable, - "System-package source is not ready for the requested package.", - "dependency_environment.system_package.source", + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::RuntimeUnavailable, + message: "System-package source is not ready for the requested package.".into(), + field_path: "dependency_environment.system_package.source", + }, source_row.alternatives.clone(), ), } @@ -380,9 +392,11 @@ fn invalid_row( binding_id, DependencyInventoryObservationState::Invalid, DependencyEnvironmentValidationState::Invalid, - DependencyPlanningDiagnosticCode::InvalidRequest, - message, - field_path, + DependencyInventoryDiagnosticInput { + code: DependencyPlanningDiagnosticCode::InvalidRequest, + message: message.into(), + field_path, + }, Vec::new(), ) } @@ -392,12 +406,15 @@ fn observation_with_diagnostic( binding_id: DependencyBindingId, state: DependencyInventoryObservationState, validation_state: DependencyEnvironmentValidationState, - code: DependencyPlanningDiagnosticCode, - message: impl Into, - field_path: &'static str, + diagnostic_input: DependencyInventoryDiagnosticInput, alternatives: Vec, ) -> DependencyInventoryObservationRow { - let diagnostic = diagnostic(item, code, message.into(), field_path); + let DependencyInventoryDiagnosticInput { + code, + message, + field_path, + } = diagnostic_input; + let diagnostic = diagnostic(item, code, message, field_path); row( binding_id, state, diff --git a/crates/pantograph-embedded-runtime/src/dependency_inventory_tests.rs b/crates/pantograph-embedded-runtime/src/dependency_inventory_tests.rs index a452e640f..bb4126025 100644 --- a/crates/pantograph-embedded-runtime/src/dependency_inventory_tests.rs +++ b/crates/pantograph-embedded-runtime/src/dependency_inventory_tests.rs @@ -326,6 +326,31 @@ async fn inventory_service_routes_device_toolchain_payloads_with_alternatives() device_toolchain_status.state, DependencyBindingStatusState::Unavailable ); + assert_eq!( + device_toolchain_status.diagnostics, + vec![ + pantograph_dependency_planning::DependencyPlanningDiagnostic { + code: DependencyPlanningDiagnosticCode::RuntimeUnavailable, + severity: pantograph_dependency_planning::DependencyPlanningSeverity::Error, + message: "Device-toolchain source is not ready for the requested toolchain." + .to_string(), + model_id: Some(request.as_request().identity_key.model_ref.model_id.clone()), + runtime_id: request + .as_request() + .identity_key + .scheduler_intent + .requested_runtime_id + .clone(), + device_id: request + .as_request() + .identity_key + .scheduler_intent + .requested_device_id + .clone(), + field_path: Some("dependency_environment.device_toolchain.source".to_string()), + } + ] + ); assert_eq!(device_toolchain_status.alternatives.len(), 1); assert_eq!( device_toolchain_status.alternatives[0] diff --git a/crates/pantograph-embedded-runtime/src/inference_interface_facts_provider.rs b/crates/pantograph-embedded-runtime/src/inference_interface_facts_provider.rs index d6854450e..69efd748c 100644 --- a/crates/pantograph-embedded-runtime/src/inference_interface_facts_provider.rs +++ b/crates/pantograph-embedded-runtime/src/inference_interface_facts_provider.rs @@ -357,7 +357,7 @@ mod tests { let facts = resolver_facts_from_sources( PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: package, + facts: Box::new(package), diagnostics: Vec::new(), }, &runtime, @@ -388,7 +388,7 @@ mod tests { fn missing_runtime_facts_keep_capability_but_publish_no_runtime_availability() { let facts = resolver_facts_from_sources( PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: projected_package_facts(), + facts: Box::new(projected_package_facts()), diagnostics: Vec::new(), }, &RuntimeDispatchCapabilityFactsOutcome::Unavailable { diff --git a/crates/pantograph-embedded-runtime/src/lib.rs b/crates/pantograph-embedded-runtime/src/lib.rs index 9eb2ad885..4873a50b9 100644 --- a/crates/pantograph-embedded-runtime/src/lib.rs +++ b/crates/pantograph-embedded-runtime/src/lib.rs @@ -163,7 +163,8 @@ pub use task_executor::{runtime_extension_keys, TauriTaskExecutor as PantographT pub(crate) use workflow_scheduler_diagnostics::EmbeddedWorkflowSchedulerDiagnosticsProvider; pub use workflow_service_composition::{ EmbeddedHostedStartupCompositionInput, EmbeddedHostedStartupCompositionOutput, - EmbeddedHostedStartupPumasSelectorSource, EmbeddedWorkflowServiceComposition, + EmbeddedHostedStartupConfig, EmbeddedHostedStartupPumasSelectorSource, + EmbeddedWorkflowServiceComposition, }; pub type SharedWorkflowService = Arc; diff --git a/crates/pantograph-embedded-runtime/src/lib_tests/workflow_run_execution_tests.rs b/crates/pantograph-embedded-runtime/src/lib_tests/workflow_run_execution_tests.rs index 31a01ef9c..40fe05b9c 100644 --- a/crates/pantograph-embedded-runtime/src/lib_tests/workflow_run_execution_tests.rs +++ b/crates/pantograph-embedded-runtime/src/lib_tests/workflow_run_execution_tests.rs @@ -34,9 +34,10 @@ use pantograph_inference_interface_contracts::{ use pantograph_runtime_attribution::WorkflowVersionRecord; use pantograph_runtime_host_contracts::{ ReservationLifecycleApplication, ReservationLifecycleApplicationState, - ReservationLifecycleEvent, ReservationLifecycleOutcome, ReservationLifecyclePort, - ReservationLifecyclePortError, RESERVATION_LIFECYCLE_CONTRACT_VERSION, - RUNTIME_SESSION_LOAD_PROOF_CONTRACT_VERSION, + ReservationLifecycleDiagnostic, ReservationLifecycleDiagnosticCode, + ReservationLifecycleDiagnosticSeverity, ReservationLifecycleEvent, ReservationLifecycleOutcome, + ReservationLifecyclePort, ReservationLifecyclePortError, + RESERVATION_LIFECYCLE_CONTRACT_VERSION, RUNTIME_SESSION_LOAD_PROOF_CONTRACT_VERSION, }; use pantograph_scheduler::{ SchedulerDispatchCandidateId, SchedulerEstimateHint, SchedulerEstimateHintKind, @@ -239,17 +240,41 @@ async fn workflow_execution_session_dispatches_through_production_embedded_image assert_eq!(body.response.media_type, "image/png"); assert_eq!(dependency_readiness_work_queue.len(), 1); assert_eq!(source_refresher.model_refs(), vec![MODEL_ID.to_string()]); + let lifecycle_events = reservation_lifecycle_port.events(); assert_eq!( - reservation_lifecycle_port - .events() + lifecycle_events .iter() .map(|event| &event.outcome) .collect::>(), vec![ &ReservationLifecycleOutcome::DispatchStarted, - &ReservationLifecycleOutcome::RuntimeHostCompleted, + &ReservationLifecycleOutcome::RetryDeferred, ] ); + let deferred = &lifecycle_events[1]; + assert_eq!( + deferred.reservation_lease_id, + lifecycle_events[0].reservation_lease_id + ); + assert_eq!( + deferred.reservation_lease_id.as_str(), + "reservation.embedded_runtime_session_test.infer" + ); + assert_eq!( + deferred.workflow_run_id, + lifecycle_events[0].workflow_run_id + ); + assert_eq!(deferred.task_id, lifecycle_events[0].task_id); + assert_eq!(deferred.candidate_id, lifecycle_events[0].candidate_id); + assert_eq!( + deferred.diagnostics, + vec![ReservationLifecycleDiagnostic { + severity: ReservationLifecycleDiagnosticSeverity::Info, + code: ReservationLifecycleDiagnosticCode::RetryDeferred, + message: "embedded runtime-host image batch member completed".to_string(), + hint: None, + }] + ); assert_eq!(host.runtime_load_attempts.load(Ordering::SeqCst), 0); assert_eq!(host.run_attempts.load(Ordering::SeqCst), 0); } diff --git a/crates/pantograph-embedded-runtime/src/pumas_dispatch_package_facts.rs b/crates/pantograph-embedded-runtime/src/pumas_dispatch_package_facts.rs index 7e561d403..16366135d 100644 --- a/crates/pantograph-embedded-runtime/src/pumas_dispatch_package_facts.rs +++ b/crates/pantograph-embedded-runtime/src/pumas_dispatch_package_facts.rs @@ -84,7 +84,7 @@ pub(crate) enum PumasDispatchPackageFactsDiagnosticCode { #[derive(Debug, Clone, PartialEq)] pub(crate) enum PumasDispatchPackageFactsBridgeOutcome { Projected { - facts: PumasDispatchPackageFactsProjection, + facts: Box, diagnostics: Vec, }, Unavailable { @@ -210,7 +210,7 @@ fn validate_and_project_dispatch_package_facts( } PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: project_dispatch_package_facts(facts), + facts: Box::new(project_dispatch_package_facts(facts)), diagnostics, } } diff --git a/crates/pantograph-embedded-runtime/src/runtime_dispatch_candidate_provider.rs b/crates/pantograph-embedded-runtime/src/runtime_dispatch_candidate_provider.rs index a00456cce..f167e753c 100644 --- a/crates/pantograph-embedded-runtime/src/runtime_dispatch_candidate_provider.rs +++ b/crates/pantograph-embedded-runtime/src/runtime_dispatch_candidate_provider.rs @@ -70,7 +70,7 @@ const INCOMPATIBLE_RUNTIME_BACKEND_HINT: &str = const MISSING_RUNTIME_DISPATCH_EVIDENCE_HINT: &str = "embedded_runtime_dispatch_candidate_provider.missing_runtime_dispatch_evidence"; -#[derive(Clone)] +#[derive(Clone, Default)] pub(crate) struct EmbeddedRuntimeDispatchCandidateProvider { source_snapshot: EmbeddedRuntimeDispatchCandidateSource, resource_facts_source: Option, @@ -78,22 +78,13 @@ pub(crate) struct EmbeddedRuntimeDispatchCandidateProvider { #[derive(Clone)] enum EmbeddedRuntimeDispatchCandidateSource { - Snapshot(EmbeddedRuntimeDispatchCandidateSourceSnapshot), + Snapshot(Box), Store(EmbeddedRuntimeDispatchSourceFactSnapshotStore), } impl Default for EmbeddedRuntimeDispatchCandidateSource { fn default() -> Self { - Self::Snapshot(EmbeddedRuntimeDispatchCandidateSourceSnapshot::default()) - } -} - -impl Default for EmbeddedRuntimeDispatchCandidateProvider { - fn default() -> Self { - Self { - source_snapshot: EmbeddedRuntimeDispatchCandidateSource::default(), - resource_facts_source: None, - } + Self::Snapshot(Box::default()) } } @@ -122,7 +113,9 @@ impl EmbeddedRuntimeDispatchCandidateProvider { source_snapshot: EmbeddedRuntimeDispatchCandidateSourceSnapshot, ) -> Self { Self { - source_snapshot: EmbeddedRuntimeDispatchCandidateSource::Snapshot(source_snapshot), + source_snapshot: EmbeddedRuntimeDispatchCandidateSource::Snapshot(Box::new( + source_snapshot, + )), resource_facts_source: None, } } @@ -181,7 +174,7 @@ impl EmbeddedRuntimeDispatchCandidateSource { model_ref: &PumasModelRef, ) -> EmbeddedRuntimeDispatchCandidateSourceSnapshot { match self { - Self::Snapshot(snapshot) => snapshot.clone(), + Self::Snapshot(snapshot) => snapshot.as_ref().clone(), Self::Store(store) => store.snapshot_for_dispatch(model_ref, current_time_ms()), } } @@ -1011,6 +1004,54 @@ mod tests { RuntimeDispatchCapabilityFactsProjection, RuntimeDispatchRuntimeCapabilityFacts, }; + #[test] + fn private_source_indirections_preserve_raw_facts_and_snapshot_clone_isolation() { + let expected = source_snapshot(Vec::new(), Vec::new()); + let expected_facts = pumas_package_facts(vec![inference::BackendHintLabel::Diffusers]); + let provider = + EmbeddedRuntimeDispatchCandidateProvider::with_source_snapshot(expected.clone()); + let mut returned = provider + .source_snapshot + .snapshot_for_dispatch(&path_free_model_ref()); + assert_eq!(returned, expected); + let Some(PumasDispatchPackageFactsBridgeOutcome::Projected { facts, diagnostics }) = + returned.pumas_package_facts.as_mut() + else { + panic!("snapshot must preserve projected package facts"); + }; + assert_eq!(facts.as_ref(), &expected_facts); + assert!(diagnostics.is_empty()); + facts.model_ref.model_id = "changed-returned-copy".to_string(); + returned.snapshot_version = 99; + assert_eq!( + provider + .source_snapshot + .snapshot_for_dispatch(&path_free_model_ref()), + expected + ); + assert!(std::mem::size_of::() <= 64); + assert!(std::mem::size_of::() <= 64); + } + + #[test] + fn default_provider_retains_empty_snapshot_and_fail_closed_diagnostics() { + let provider = EmbeddedRuntimeDispatchCandidateProvider::default(); + assert!(provider.resource_facts_source.is_none()); + let EmbeddedRuntimeDispatchCandidateSource::Snapshot(snapshot) = provider.source_snapshot + else { + panic!("default provider must retain snapshot mode"); + }; + assert_eq!( + *snapshot, + EmbeddedRuntimeDispatchCandidateSourceSnapshot::default() + ); + let diagnostics = fail_closed_diagnostics(&snapshot, &path_free_model_ref()); + assert_eq!(diagnostics.len(), 4); + assert!(diagnostics.iter().all(|diagnostic| diagnostic.code + == SchedulerDispatchSelectionDiagnosticCode::NoCandidates + && diagnostic.severity == SchedulerDispatchSelectionDiagnosticSeverity::Error)); + } + #[test] fn fail_closed_provider_reports_missing_source_facts() { let diagnostics = fail_closed_diagnostics( @@ -1235,7 +1276,9 @@ mod tests { let provider = EmbeddedRuntimeDispatchCandidateProvider::with_source_snapshot( EmbeddedRuntimeDispatchCandidateSourceSnapshot { pumas_package_facts: Some(PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: pumas_package_facts(vec![inference::BackendHintLabel::Diffusers]), + facts: Box::new(pumas_package_facts(vec![ + inference::BackendHintLabel::Diffusers, + ])), diagnostics: Vec::new(), }), runtime_capability_facts: Some(RuntimeDispatchCapabilityFactsOutcome::Projected { @@ -1395,7 +1438,9 @@ mod tests { let provider = EmbeddedRuntimeDispatchCandidateProvider::with_source_snapshot( EmbeddedRuntimeDispatchCandidateSourceSnapshot { pumas_package_facts: Some(PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: pumas_package_facts(vec![inference::BackendHintLabel::Diffusers]), + facts: Box::new(pumas_package_facts(vec![ + inference::BackendHintLabel::Diffusers, + ])), diagnostics: Vec::new(), }), runtime_capability_facts: Some(RuntimeDispatchCapabilityFactsOutcome::Projected { @@ -1536,7 +1581,9 @@ mod tests { ) -> EmbeddedRuntimeDispatchCandidateSourceSnapshot { EmbeddedRuntimeDispatchCandidateSourceSnapshot { pumas_package_facts: Some(PumasDispatchPackageFactsBridgeOutcome::Projected { - facts: pumas_package_facts(vec![inference::BackendHintLabel::Diffusers]), + facts: Box::new(pumas_package_facts(vec![ + inference::BackendHintLabel::Diffusers, + ])), diagnostics: Vec::new(), }), runtime_capability_facts: Some(RuntimeDispatchCapabilityFactsOutcome::Projected { diff --git a/crates/pantograph-embedded-runtime/src/runtime_host_media_artifact_sink.rs b/crates/pantograph-embedded-runtime/src/runtime_host_media_artifact_sink.rs index c80763974..bc4bb4f13 100644 --- a/crates/pantograph-embedded-runtime/src/runtime_host_media_artifact_sink.rs +++ b/crates/pantograph-embedded-runtime/src/runtime_host_media_artifact_sink.rs @@ -37,14 +37,18 @@ pub(crate) enum RuntimeHostMediaArtifactSinkError { image_index: usize, message: String, }, - #[error("runtime-host image output artifact write failed for {task_id}.{port_id}[{image_index}]: {source}")] - ArtifactWriteFailed { - task_id: String, - port_id: String, - image_index: usize, - #[source] - source: WorkflowServiceError, - }, + #[error(transparent)] + ArtifactWriteFailed(Box), +} + +#[derive(Debug, Error)] +#[error("runtime-host image output artifact write failed for {task_id}.{port_id}[{image_index}]: {source}")] +pub(crate) struct RuntimeHostArtifactWriteFailure { + task_id: String, + port_id: String, + image_index: usize, + #[source] + source: WorkflowServiceError, } #[derive(Clone)] @@ -101,14 +105,16 @@ impl RuntimeHostMediaArtifactSink for WorkflowServiceRuntimeHostMediaArtifactSin revision_index: Some(request.image_index as u64), body, }) - .map_err( - |source| RuntimeHostMediaArtifactSinkError::ArtifactWriteFailed { - task_id: request.task_id.to_string(), - port_id: request.port_id.to_string(), - image_index: request.image_index, - source, - }, - )?; + .map_err(|source| { + RuntimeHostMediaArtifactSinkError::ArtifactWriteFailed(Box::new( + RuntimeHostArtifactWriteFailure { + task_id: request.task_id.to_string(), + port_id: request.port_id.to_string(), + image_index: request.image_index, + source, + }, + )) + })?; Ok(RuntimeHostExecutionMediaArtifactRef { artifact_id: descriptor.artifact_id, @@ -316,11 +322,56 @@ mod tests { }) .expect_err("missing artifact store must fail closed"); - let RuntimeHostMediaArtifactSinkError::ArtifactWriteFailed { source, .. } = error else { + let concrete_source = std::error::Error::source(&error) + .expect("artifact write failure retains its source") + .downcast_ref::() + .expect("source remains the concrete workflow error"); + assert_eq!(concrete_source.code(), WorkflowErrorCode::InternalError); + let RuntimeHostMediaArtifactSinkError::ArtifactWriteFailed(failure) = error else { panic!("expected artifact write failure"); }; + assert_eq!(failure.source.code(), WorkflowErrorCode::InternalError); + assert!(failure + .source + .to_string() + .contains("artifact store io error")); + } + + #[test] + fn boxed_write_failure_preserves_exact_display_and_concrete_source_chain() { + let error = RuntimeHostMediaArtifactSinkError::ArtifactWriteFailed(Box::new( + RuntimeHostArtifactWriteFailure { + task_id: "task.image".to_string(), + port_id: "image".to_string(), + image_index: 2, + source: WorkflowServiceError::Internal("fixture failure".to_string()), + }, + )); + assert_eq!(error.to_string(), "runtime-host image output artifact write failed for task.image.image[2]: internal_error: fixture failure"); + let source = std::error::Error::source(&error).unwrap(); + assert!(source + .downcast_ref::() + .is_none()); + assert!(source.downcast_ref::>().is_none()); + let source = source + .downcast_ref::() + .expect("unboxed workflow source identity"); assert_eq!(source.code(), WorkflowErrorCode::InternalError); - assert!(source.to_string().contains("artifact store io error")); + assert_eq!(source.to_string(), "internal_error: fixture failure"); + assert!(std::error::Error::source(source).is_none()); + assert!(std::mem::size_of::() <= 128); + } + + #[test] + fn invalid_image_failure_retains_message_and_absent_source() { + let error = RuntimeHostMediaArtifactSinkError::InvalidImagePayload { + task_id: "task.image".to_string(), + port_id: "image".to_string(), + image_index: 2, + message: "invalid fixture bytes".to_string(), + }; + assert_eq!(error.to_string(), "runtime-host image output base64 decode failed for task.image.image[2]: invalid fixture bytes"); + assert!(std::error::Error::source(&error).is_none()); } fn artifact_writer(temp: &TempDir) -> WorkflowArtifactWriter { diff --git a/crates/pantograph-embedded-runtime/src/task_executor/dependency_environment/helpers.rs b/crates/pantograph-embedded-runtime/src/task_executor/dependency_environment/helpers.rs index 52a1db68e..72b252f4c 100644 --- a/crates/pantograph-embedded-runtime/src/task_executor/dependency_environment/helpers.rs +++ b/crates/pantograph-embedded-runtime/src/task_executor/dependency_environment/helpers.rs @@ -103,10 +103,7 @@ impl TauriTaskExecutor { } pub(in crate::task_executor) fn python_runtime_handles_node(node_type: &str) -> bool { - match node_type { - "audio-generation" | "onnx-inference" => true, - _ => false, - } + matches!(node_type, "audio-generation" | "onnx-inference") } pub(in crate::task_executor) fn sanitize_key_component(raw: &str) -> String { @@ -130,3 +127,28 @@ impl TauriTaskExecutor { format!("{:016x}", digest) } } + +#[cfg(test)] +mod tests { + use super::TauriTaskExecutor; + + #[test] + fn python_runtime_node_classification_keeps_exact_accepted_set() { + for node_type in ["audio-generation", "onnx-inference"] { + assert!(TauriTaskExecutor::python_runtime_handles_node(node_type)); + } + for node_type in [ + "llm-inference", + "embedding", + "vision", + "", + "Audio-generation", + "onnx-inference ", + ] { + assert!( + !TauriTaskExecutor::python_runtime_handles_node(node_type), + "{node_type}" + ); + } + } +} diff --git a/crates/pantograph-embedded-runtime/src/task_executor/puma_lib.rs b/crates/pantograph-embedded-runtime/src/task_executor/puma_lib.rs index a6da852cd..02135e330 100644 --- a/crates/pantograph-embedded-runtime/src/task_executor/puma_lib.rs +++ b/crates/pantograph-embedded-runtime/src/task_executor/puma_lib.rs @@ -151,7 +151,7 @@ impl TauriTaskExecutor { if let Some(selector_access) = extensions.get::>(PUMAS_SELECTOR_ACCESS) { - match Self::resolve_puma_lib_selected_detail(&selector_access, requested_model_id) + match Self::resolve_puma_lib_selected_detail(selector_access, requested_model_id) .await { Ok(Some(detail)) => { diff --git a/crates/pantograph-embedded-runtime/src/task_executor/stream_artifacts.rs b/crates/pantograph-embedded-runtime/src/task_executor/stream_artifacts.rs index 64f9f9f7c..5c0002629 100644 --- a/crates/pantograph-embedded-runtime/src/task_executor/stream_artifacts.rs +++ b/crates/pantograph-embedded-runtime/src/task_executor/stream_artifacts.rs @@ -57,15 +57,8 @@ impl StreamArtifactizer { let relationship = artifact_relationship_from_chunk(chunk); let key = MediaStreamKey::new(execution_id, task_id, port); - let (artifact_id, stream_handle, byte_range_start) = self.open_or_get_stream( - &key, - task_id, - execution_id, - port, - media_body.kind, - &media_type, - &relationship, - )?; + let (artifact_id, stream_handle, byte_range_start) = + self.open_or_get_stream(&key, media_body.kind, &media_type, &relationship)?; let (byte_range_end_exclusive, next_sequence) = checked_stream_artifact_progress(byte_range_start, body.len(), sequence)?; @@ -167,9 +160,6 @@ impl StreamArtifactizer { fn open_or_get_stream( &self, key: &MediaStreamKey, - task_id: &str, - execution_id: &str, - port: &str, payload_kind: ArtifactPayloadKind, media_type: &str, relationship: &ArtifactRelationship, @@ -184,11 +174,11 @@ impl StreamArtifactizer { media_type: media_type.to_string(), format: Some(format_metadata(payload_kind, media_type)), attribution: ArtifactAttribution { - workflow_run_id: execution_id.to_string(), + workflow_run_id: key.workflow_run_id.clone(), workflow_id: None, workflow_version_id: None, - node_id: Some(task_id.to_string()), - port_id: Some(port.to_string()), + node_id: Some(key.node_id.clone()), + port_id: Some(key.port.clone()), model_id: None, runtime_id: None, }, @@ -414,6 +404,71 @@ mod tests { use pantograph_workflow_service::WorkflowService; use std::sync::Arc; + #[test] + fn stream_key_preserves_attribution_reuse_and_scope_isolation() { + use pantograph_workflow_service::{ + ArtifactDescriptorQueryRequest, ArtifactPolicy, ArtifactReadRequest, ArtifactStore, + WorkflowArtifactWriter, + }; + let temp = tempfile::tempdir().expect("artifact directory"); + let store = ArtifactStore::open( + temp.path().join("artifacts"), + ArtifactPolicy { + policy_id: "stream-key-test".to_string(), + policy_version: 1, + ttl_seconds: None, + max_disk_bytes: None, + max_memory_bytes: None, + max_single_artifact_bytes: None, + spill_threshold_bytes: None, + delete_on_consume: false, + }, + ) + .expect("artifact store"); + let service = Arc::new( + WorkflowService::new().with_artifact_writer(WorkflowArtifactWriter::new(store)), + ); + let artifactizer = StreamArtifactizer::new(service.clone()); + let chunk = serde_json::json!({"audio_base64":"AA==", "media_type":"audio/wav"}); + let first = artifactizer.artifactize_chunk("node-a", "run-a", "audio", chunk.clone()); + let mut final_chunk = chunk.clone(); + final_chunk["is_final"] = serde_json::json!(true); + let second = artifactizer.artifactize_chunk("node-a", "run-a", "audio", final_chunk); + assert_eq!(first["artifact_id"], second["artifact_id"]); + assert_eq!(first["byte_range_start"], 0); + assert_eq!(second["byte_range_start"], 1); + assert_eq!(second["available_byte_length"], 2); + let id = first["artifact_id"].as_str().expect("artifact id"); + let descriptor = service + .artifact_descriptor(ArtifactDescriptorQueryRequest { + artifact_id: id.to_string(), + }) + .expect("descriptor") + .artifact + .expect("stored artifact"); + assert_eq!(descriptor.attribution.workflow_run_id, "run-a"); + assert_eq!(descriptor.attribution.node_id.as_deref(), Some("node-a")); + assert_eq!(descriptor.attribution.port_id.as_deref(), Some("audio")); + let body = service + .read_artifact_body(ArtifactReadRequest { + artifact_id: id.to_string(), + byte_range_start: None, + byte_range_end_exclusive: None, + }) + .expect("body"); + assert_eq!(body.body, vec![0, 0]); + for (node, run, port) in [ + ("node-b", "run-a", "audio"), + ("node-a", "run-b", "audio"), + ("node-a", "run-a", "other"), + ] { + let isolated = artifactizer.artifactize_chunk(node, run, port, chunk.clone()); + assert_ne!(isolated["artifact_id"], first["artifact_id"]); + assert_eq!(isolated["byte_range_start"], 0); + } + assert_eq!(artifactizer.streams.lock().expect("streams").len(), 3); + } + #[test] fn stream_artifact_progress_rejects_byte_range_overflow() { assert_eq!(checked_stream_artifact_progress(u64::MAX, 1, 0), None); diff --git a/crates/pantograph-embedded-runtime/src/workflow_service_composition.rs b/crates/pantograph-embedded-runtime/src/workflow_service_composition.rs index 87b076248..189dc9c09 100644 --- a/crates/pantograph-embedded-runtime/src/workflow_service_composition.rs +++ b/crates/pantograph-embedded-runtime/src/workflow_service_composition.rs @@ -140,6 +140,22 @@ pub enum EmbeddedHostedStartupPumasSelectorSource { SetupPath(Option), } +/// Owned inputs for hosted startup composition, without changing startup defaults. +/// +/// Rust callers construct these named fields and pass the config to +/// `EmbeddedHostedStartupCompositionInput::new`. +pub struct EmbeddedHostedStartupConfig { + pub runtime_registry: SharedRuntimeRegistry, + pub runtime_registry_controller: Arc, + pub gateway: Arc, + pub pumas_selector_source: Option, + pub project_root: PathBuf, + pub kv_cache_dir: PathBuf, + pub dependency_readiness_runtime_handle: tokio::runtime::Handle, + pub max_loaded_sessions: Option, + pub max_dispatch_source_snapshot_age_ms: u64, +} + pub struct EmbeddedHostedStartupCompositionInput { workflow_service: WorkflowService, max_loaded_sessions: Option, @@ -156,17 +172,18 @@ pub struct EmbeddedHostedStartupCompositionInput { impl EmbeddedHostedStartupCompositionInput { #[must_use] - pub fn new( - runtime_registry: SharedRuntimeRegistry, - runtime_registry_controller: Arc, - gateway: Arc, - pumas_selector_source: Option, - project_root: PathBuf, - kv_cache_dir: PathBuf, - dependency_readiness_runtime_handle: tokio::runtime::Handle, - max_loaded_sessions: Option, - max_dispatch_source_snapshot_age_ms: u64, - ) -> Self { + pub fn new(config: EmbeddedHostedStartupConfig) -> Self { + let EmbeddedHostedStartupConfig { + runtime_registry, + runtime_registry_controller, + gateway, + pumas_selector_source, + project_root, + kv_cache_dir, + dependency_readiness_runtime_handle, + max_loaded_sessions, + max_dispatch_source_snapshot_age_ms, + } = config; Self { workflow_service: WorkflowService::new(), max_loaded_sessions, @@ -850,6 +867,12 @@ mod tests { .await .expect("refresh validation summary"); + assert!( + validation.summary.diagnostics.is_empty(), + "validation diagnostics: {:?}", + validation.summary.diagnostics + ); + assert_eq!(validation.node_projections.len(), 1); let summary = validation .summary .summary @@ -946,6 +969,21 @@ mod tests { ) .await .expect("refresh validation summary"); + assert!( + validation.summary.diagnostics.is_empty(), + "validation diagnostics: {:?}", + validation.summary.diagnostics + ); + assert_eq!(validation.node_projections.len(), 1); + assert_eq!( + validation + .summary + .summary + .as_ref() + .expect("current validation summary") + .status, + DraftGraphValidationStatus::Executable + ); let validation_session_id = validation .summary .validation_session_id @@ -1195,6 +1233,52 @@ mod tests { .contains("dependency-readiness snapshot producer poll interval")); } + #[tokio::test] + async fn hosted_startup_named_config_preserves_owned_values_and_defaults() { + let registry: SharedRuntimeRegistry = Arc::new(RuntimeRegistry::new()); + let gateway = Arc::new(inference::InferenceGateway::new()); + let input = EmbeddedHostedStartupCompositionInput::new(EmbeddedHostedStartupConfig { + runtime_registry: registry.clone(), + runtime_registry_controller: gateway.clone(), + gateway: gateway.clone(), + pumas_selector_source: Some(EmbeddedHostedStartupPumasSelectorSource::SetupPath(Some( + PathBuf::from("pumas-root"), + ))), + project_root: PathBuf::from("project-root"), + kv_cache_dir: PathBuf::from("kv-root"), + dependency_readiness_runtime_handle: tokio::runtime::Handle::current(), + max_loaded_sessions: Some(7), + max_dispatch_source_snapshot_age_ms: 2_345, + }); + assert!(Arc::ptr_eq(&input.runtime_registry, ®istry)); + assert!(Arc::ptr_eq(&input.runtime_registry_controller, &gateway)); + assert!(Arc::ptr_eq(&input.gateway, &gateway)); + assert!( + matches!(&input.pumas_selector_source, Some(EmbeddedHostedStartupPumasSelectorSource::SetupPath(Some(path))) if path == std::path::Path::new("pumas-root")) + ); + assert_eq!(input.project_root, PathBuf::from("project-root")); + assert_eq!(input.kv_cache_dir, PathBuf::from("kv-root")); + assert_eq!(input.max_loaded_sessions, Some(7)); + assert_eq!(input.max_dispatch_source_snapshot_age_ms, 2_345); + assert_eq!( + input.dependency_readiness_producer_config, + EmbeddedDependencyReadinessSnapshotProducerConfig::default() + ); + assert_eq!( + input + .dependency_readiness_runtime_handle + .spawn(async { 42 }) + .await + .expect("original runtime handle"), + 42 + ); + let config = EmbeddedDependencyReadinessSnapshotProducerConfig { + poll_interval: std::time::Duration::from_secs(17), + }; + let input = input.with_dependency_readiness_producer_config(config.clone()); + assert_eq!(input.dependency_readiness_producer_config, config); + } + #[tokio::test] async fn hosted_startup_composition_returns_service_extensions_and_lifecycle_handle() { let temp_dir = tempfile::tempdir().expect("tempdir"); @@ -1208,19 +1292,19 @@ mod tests { ); let registry: SharedRuntimeRegistry = Arc::new(RuntimeRegistry::new()); let gateway = Arc::new(inference::InferenceGateway::new()); - let input = EmbeddedHostedStartupCompositionInput::new( - registry, - gateway.clone(), + let input = EmbeddedHostedStartupCompositionInput::new(EmbeddedHostedStartupConfig { + runtime_registry: registry, + runtime_registry_controller: gateway.clone(), gateway, - Some(EmbeddedHostedStartupPumasSelectorSource::Provided( + pumas_selector_source: Some(EmbeddedHostedStartupPumasSelectorSource::Provided( Arc::new(PumasSelectorAccess::Owner(pumas_api.clone())), )), - temp_dir.path().to_path_buf(), - temp_dir.path().join("kv-cache"), - tokio::runtime::Handle::current(), - Some(1), - 1_000, - ) + project_root: temp_dir.path().to_path_buf(), + kv_cache_dir: temp_dir.path().join("kv-cache"), + dependency_readiness_runtime_handle: tokio::runtime::Handle::current(), + max_loaded_sessions: Some(1), + max_dispatch_source_snapshot_age_ms: 1_000, + }) .with_workflow_service(workflow_service_with_artifact_store(&temp_dir)); let output = EmbeddedWorkflowServiceComposition::resource_backed_hosted_startup(input) @@ -1283,17 +1367,17 @@ mod tests { let temp_dir = tempfile::tempdir().expect("tempdir"); let registry: SharedRuntimeRegistry = Arc::new(RuntimeRegistry::new()); let gateway = Arc::new(inference::InferenceGateway::new()); - let input = EmbeddedHostedStartupCompositionInput::new( - registry, - gateway.clone(), + let input = EmbeddedHostedStartupCompositionInput::new(EmbeddedHostedStartupConfig { + runtime_registry: registry, + runtime_registry_controller: gateway.clone(), gateway, - None, - temp_dir.path().to_path_buf(), - temp_dir.path().join("kv-cache"), - tokio::runtime::Handle::current(), - Some(1), - 1_000, - ); + pumas_selector_source: None, + project_root: temp_dir.path().to_path_buf(), + kv_cache_dir: temp_dir.path().join("kv-cache"), + dependency_readiness_runtime_handle: tokio::runtime::Handle::current(), + max_loaded_sessions: Some(1), + max_dispatch_source_snapshot_age_ms: 1_000, + }); let error = match EmbeddedWorkflowServiceComposition::resource_backed_hosted_startup(input).await { @@ -1326,19 +1410,19 @@ mod tests { .expect("read-only pumas"); let registry: SharedRuntimeRegistry = Arc::new(RuntimeRegistry::new()); let gateway = Arc::new(inference::InferenceGateway::new()); - let input = EmbeddedHostedStartupCompositionInput::new( - registry, - gateway.clone(), + let input = EmbeddedHostedStartupCompositionInput::new(EmbeddedHostedStartupConfig { + runtime_registry: registry, + runtime_registry_controller: gateway.clone(), gateway, - Some(EmbeddedHostedStartupPumasSelectorSource::Provided( + pumas_selector_source: Some(EmbeddedHostedStartupPumasSelectorSource::Provided( Arc::new(PumasSelectorAccess::ReadOnly(Arc::new(read_only))), )), - temp_dir.path().to_path_buf(), - temp_dir.path().join("kv-cache"), - tokio::runtime::Handle::current(), - Some(1), - 1_000, - ); + project_root: temp_dir.path().to_path_buf(), + kv_cache_dir: temp_dir.path().join("kv-cache"), + dependency_readiness_runtime_handle: tokio::runtime::Handle::current(), + max_loaded_sessions: Some(1), + max_dispatch_source_snapshot_age_ms: 1_000, + }); let error = match EmbeddedWorkflowServiceComposition::resource_backed_hosted_startup(input).await { @@ -1365,19 +1449,19 @@ mod tests { ); let registry: SharedRuntimeRegistry = Arc::new(RuntimeRegistry::new()); let gateway = Arc::new(inference::InferenceGateway::new()); - let input = EmbeddedHostedStartupCompositionInput::new( - registry, - gateway.clone(), + let input = EmbeddedHostedStartupCompositionInput::new(EmbeddedHostedStartupConfig { + runtime_registry: registry, + runtime_registry_controller: gateway.clone(), gateway, - Some(EmbeddedHostedStartupPumasSelectorSource::Provided( + pumas_selector_source: Some(EmbeddedHostedStartupPumasSelectorSource::Provided( Arc::new(PumasSelectorAccess::Owner(pumas_api)), )), - temp_dir.path().to_path_buf(), - temp_dir.path().join("kv-cache"), - tokio::runtime::Handle::current(), - Some(0), - 1_000, - ) + project_root: temp_dir.path().to_path_buf(), + kv_cache_dir: temp_dir.path().join("kv-cache"), + dependency_readiness_runtime_handle: tokio::runtime::Handle::current(), + max_loaded_sessions: Some(0), + max_dispatch_source_snapshot_age_ms: 1_000, + }) .with_workflow_service(workflow_service_with_store( &temp_dir, WorkflowService::with_capacity_limits(2, 2), @@ -1435,7 +1519,12 @@ mod tests { position: Position { x: 200.0, y: 0.0 }, data: serde_json::json!({ "task_kind": "image_generation", - "runtime": "pytorch" + "runtime": "pytorch", + "runtime_source_context": { + "operation_type": "image-generation.txt2img", + "context_shape_key": "txt2img.1024x1024.steps30", + "cancellation_mode": "per-run-fanout" + } }), }, ], diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-artifact-write-error-payload.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-artifact-write-error-payload.md new file mode 100644 index 000000000..1bc1eb74c --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-artifact-write-error-payload.md @@ -0,0 +1,7 @@ +# Private artifact-write failure payload and source identity + +Fresh embedded-runtime Clippy identifies three large-error returns caused by ArtifactWriteFailed carrying a raw WorkflowServiceError plus attribution. Keep all media sink trait and Result signatures unchanged, but box a private RuntimeHostArtifactWriteFailure payload containing the unchanged raw source. The enum variant is error-transparent. One production constructor and one existing test pattern are migrated. + +Directly boxing the #[source] field was rejected during design: pinned thiserror 1.0.69 calls source.as_dyn_error(), which can expose Box to downcast consumers. The transparent payload delegates Error::source through its own raw #[source] field instead. Tests explicitly require the concrete WorkflowServiceError downcast and classification, reject wrapper/boxed-source downcasts, preserve exact outer/source messages, and retain the invalid-image variant's message and absent source. The real missing-artifact-store failure now also checks concrete source identity. A private enum-size bound targets the reported large errors. + +Root approved this source-chain-preserving design. Full pinned cargo fmt and staged checks passed. Root source review accepted both files at frozen tree 002ee22ba66b1a957906fc530d1ac080f54c8dbe; actual hosted execution remains required. No public DTO/API or failure classification change is intended. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-contract-fixtures.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-contract-fixtures.md new file mode 100644 index 000000000..b7e5ba4ae --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-contract-fixtures.md @@ -0,0 +1,9 @@ +# Embedded runtime integration fixture contracts + +The full default-feature embedded suite was newly enabled on PR42 and exposed three failures. At b73c5d1 it executed 452 tests: 449 passed, three failed. Preserve the full gate and original execution/publication assertions; this milestone changes only fixtures and assertions. + +The two resource-backed composition fixtures omit mandatory runtime_source_context. inference_interface_request.rs rejects that absence before facts resolution; aggregate publication makes request diagnostics Blocked. The dependency-action owner maps that blocked summary to DriftDetected, so this downstream code alone does not establish descriptor fingerprint drift. Add the explicit image-generation operation/context/cancellation values already used by the image execution fixture. Assert no validation diagnostics and exactly one projection before the original Executable/RequestReady/publication assertions. No production admission bypass or inferred default is added. + +The image test already verifies real production runtime-host output, stored image bytes/media type and no legacy host load/run. Its reservation assertion was outdated: the scheduler explicitly requests DeferToScheduler; the embedded batch response reports DeferredToScheduler even when Completed; the scheduler maps that reservation disposition to RetryDeferred. Preserve every output assertion and require the exact deferred diagnostic (Info, RetryDeferred, successful image-member completion message, no hint) and correlated lease/run/task/candidate across start and terminal events. This asserts deferred reservation ownership without confusing it with failed image execution. + +Root source review accepted the complete three-file fixture-only diff at frozen tree d204dd518d905a4d6d1f1fe471678cb511822f5c. Fresh hosted full embedded execution remains pending. Local execution is not claimed; full pinned formatting is checked. The prior three failing assertions remain recorded in hosted logs rather than reclassified without baseline execution. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-inputs.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-inputs.md new file mode 100644 index 000000000..efa88b303 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-inputs.md @@ -0,0 +1,9 @@ +# Private inventory diagnostic input and stream attribution + +Fresh Clippy at 788ce31 reports only four embedded-runtime arities. This milestone repairs three private methods; the public hosted startup constructor remains separate. + +Both device-toolchain and system-package observation helpers had eight arguments and six callers each. Group only code, owned message and static field_path in a shared private DependencyInventoryDiagnosticInput, gated with the existing test/standalone owners. Keep work item, binding identity, observation state, validation state and alternatives distinct. The message conversion moves from helper entry to the named field at each call; diagnostic attribution lookup and row builder remain unchanged. Exact diagnostic equality is added to the existing real inventory-service device-toolchain test, including model/runtime/device and field path; existing ready/unavailable/system-package coverage remains. + +The stream helper has one caller. MediaStreamKey already copies the exact execution/task/port strings, so remove those three duplicated arguments and clone attribution directly from that key. Lock/map/reuse/relationship/progress/failure behavior remains unchanged. A real artifact-store regression verifies two chunks reuse one artifact, finalization retains exact bytes, attribution matches all three key fields, and changing each field isolates a fresh stream. Existing overflow and redaction regressions remain. + +All changes are private; public DTOs and Result signatures are unchanged. Root source review accepted all six files at frozen tree 6d72a526a367dfd9bfd60e7e1d78ad0eb46d265c. Hosted full embedded execution remains pending. Full pinned cargo fmt passes locally; no local Rust execution is claimed. Prior layout head has 450 embedded passes and three fixture failures, with actual clone-isolation test, Headless and Runtime Separation passing. The fixture repair is a separate preceding commit bbce93c whose hosted result remains pending. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-source-layout.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-source-layout.md new file mode 100644 index 000000000..a9b50dbc4 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-private-source-layout.md @@ -0,0 +1,9 @@ +# Private embedded source payload layout + +Fresh PR42 Clippy reports large variants in the package-facts bridge outcome and dispatch candidate source. Box only Projected.facts and Snapshot, keeping raw projection/snapshot DTOs and public APIs unchanged. The package bridge has one production constructor and five fixture constructors across three files. Candidate source has two constructors (default and explicit snapshot) and one raw-cloning reader. Borrowed readers retain deref coercion; the owned reader explicitly clones the raw snapshot rather than returning a box. + +The regression compares complete raw package facts and snapshots, mutates a returned nested fact and snapshot version, then proves provider state remains isolated. Existing source-owner, stale/missing/blocked, resolver, and candidate-selection tests remain in the full embedded unit gate. Size bounds apply only to these private enums; no measured throughput claim is made. + +Prior artifact-error head b73c5d1 has actual hosted concrete source/downcast and artifact-write tests passing. Its full embedded unit run has 449 passed and three failures: image dispatch gets RetryDeferred instead of RuntimeHostCompleted; resource-backed validation gets Blocked instead of Executable; publication dependency readiness gets Blocked instead of RequestReady with DriftDetected. Those failures remain separate and are not waived or classified as baseline without exact comparison. Format and workspace default/all-feature compilation pass at that head. + +Root source review accepted all four files at frozen tree 614eb5f7f61ee8ffa96d4d07cdaa238f7275e1a2. Fresh hosted execution remains required. Local Rust execution is not attempted under the shared disk constraint; full pinned formatting is checked separately. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-runtime-lint-basics.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-runtime-lint-basics.md new file mode 100644 index 000000000..fd1b2ba93 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-embedded-runtime-lint-basics.md @@ -0,0 +1,7 @@ +# Embedded runtime mechanical lint basics + +Fresh PR #40 Clippy clears runtime-registry and exposes thirteen embedded-runtime library findings. This first area milestone addresses only four mechanical findings: derive the candidate provider's identical default, remove a redundant inventory-reason borrow and selector-access borrow, and express the Python runtime node predicate with matches!. Payload/error layout and arity findings remain separate. + +Regressions prove the default still uses an empty snapshot with no resource facts and four fail-closed diagnostics, and the helper accepts exactly audio-generation/onnx-inference while rejecting other, differently cased and padded node identifiers. CI runs the complete embedded-runtime unit suite with its default backend features. Public APIs, runtime policy and errors are unchanged. + +The same checkout is reused on an embedded-runtime branch, preserving prior refs. The actual full pinned cargo fmt --all -- --check passes locally; root source review accepted frozen tree 405aea674557289e585587c1b23b48e381acbb0d; hosted tests remain required. No local Rust execution or heavy build is claimed. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-hosted-startup-named-config.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-hosted-startup-named-config.md new file mode 100644 index 000000000..2e6c22ada --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-hosted-startup-named-config.md @@ -0,0 +1,9 @@ +# Named owned hosted startup configuration + +The last public embedded-runtime arity finding is the nine-argument EmbeddedHostedStartupCompositionInput::new. Introduce exported EmbeddedHostedStartupConfig with exactly the nine existing owned inputs: runtime registry, registry controller, inference gateway, optional selector source, project root, KV cache path, readiness runtime handle, optional loaded-session limit, and source-snapshot age limit. The constructor destructures these fields into its unchanged initializer. WorkflowService::new and default producer configuration remain unchanged, as do both builder methods. + +All five existing calls are migrated: desktop app_setup and four hosted composition tests. Each original expression, clone and ownership transfer stays attached to its corresponding named field; no policy is inferred or optional field silently defaulted. A constructor regression verifies Arc identity, paths, selector setup path, limits, functional runtime handle, default producer configuration and producer builder override. Existing real startup tests retain success/lifecycle, missing selector, wrong owner and invalid service-capacity coverage. Hosted workspace all-feature compilation covers the desktop caller, and the full embedded unit gate exercises startup tests. + +This is a public Rust source signature change: callers must pass the new config instead of nine positional arguments. The crate is publish=false and uses the workspace version, so no independently published crate version is changed. Unseen Rust consumers are not claimed source-compatible. Serialized DTOs, runtime policy and startup order are unchanged. + +Root source review accepted all four files at frozen tree afc4cc26997ba3f61dde38ab77e1927563e1f18e. Hosted execution remains pending. The constructor initializer and both existing builder bodies were also compared byte-for-byte with the published parent and remain identical. Full pinned formatting passes locally; no local Rust execution is claimed under the shared resource constraint. diff --git a/src-tauri/src/app_setup.rs b/src-tauri/src/app_setup.rs index d68f0ed1b..cb5d69ee5 100644 --- a/src-tauri/src/app_setup.rs +++ b/src-tauri/src/app_setup.rs @@ -9,8 +9,8 @@ use crate::llm::{ use crate::project_root::resolve_project_root; use crate::workflow; use pantograph_embedded_runtime::{ - EmbeddedHostedStartupCompositionInput, EmbeddedHostedStartupPumasSelectorSource, - EmbeddedWorkflowServiceComposition, + EmbeddedHostedStartupCompositionInput, EmbeddedHostedStartupConfig, + EmbeddedHostedStartupPumasSelectorSource, EmbeddedWorkflowServiceComposition, }; use std::sync::Arc; use tauri::{Emitter, Manager}; @@ -245,17 +245,20 @@ pub fn run_app() -> AppStartupResult<()> { let pumas_library_path = pumas_release_dir.or(pumas_launcher_root); let startup_input = EmbeddedHostedStartupCompositionInput::new( - runtime_registry.clone(), - gateway.clone(), - gateway.inner_arc(), - Some(EmbeddedHostedStartupPumasSelectorSource::SetupPath( - pumas_library_path, - )), - project_root.clone(), - kv_cache_dir, - runtime_handle.clone(), - max_loaded_sessions, - HOSTED_DISPATCH_SOURCE_SNAPSHOT_MAX_AGE_MS, + EmbeddedHostedStartupConfig { + runtime_registry: runtime_registry.clone(), + runtime_registry_controller: gateway.clone(), + gateway: gateway.inner_arc(), + pumas_selector_source: Some( + EmbeddedHostedStartupPumasSelectorSource::SetupPath(pumas_library_path), + ), + project_root: project_root.clone(), + kv_cache_dir, + dependency_readiness_runtime_handle: runtime_handle.clone(), + max_loaded_sessions, + max_dispatch_source_snapshot_age_ms: + HOSTED_DISPATCH_SOURCE_SNAPSHOT_MAX_AGE_MS, + }, ) .with_workflow_service(workflow_service); let startup_output = tauri::async_runtime::block_on(