diff --git a/.github/workflows/quality-gates.yml b/.github/workflows/quality-gates.yml index 8a0d02108..e81a5ab8e 100644 --- a/.github/workflows/quality-gates.yml +++ b/.github/workflows/quality-gates.yml @@ -156,6 +156,9 @@ jobs: - name: Check workspace with all features run: cargo check --workspace --all-features + - name: Verify Tauri injected-state IPC contracts + run: bash scripts/check-tauri-command-state-tests.sh + rust-tests: name: Rust focused tests runs-on: ubuntu-latest @@ -245,6 +248,21 @@ jobs: python3 -c 'import pathlib, sys; names = [line for line in pathlib.Path(sys.argv[1]).read_text().splitlines() if line.startswith("core_executor::tests::inference_tests::") and "rejects_contract_only_with_lifecycle" in line and line.endswith(": test")]; print("Contract-only lifecycle tests discovered:", len(names)); sys.exit(0 if names else 1)' "$RUNNER_TEMP/node-validation-lifecycle-tests.list" cargo test -p node-engine --features inference-nodes --lib rejects_contract_only_with_lifecycle + - name: Run typed embedding capture regressions + run: | + set -euo pipefail + tests=( + core_executor::tests::inference_tests::test_canonical_llm_embedding_uses_typed_gateway_boundary + core_executor::tests::inference_tests::test_canonical_llm_embedding_with_package_facts_emits_compatibility_lifecycle + ) + cargo test -p node-engine --features inference-nodes --lib test_canonical_llm_embedding -- --list > "$RUNNER_TEMP/node-embedding-tests.list" + for test_name in "${tests[@]}"; do + python3 -c 'import pathlib, sys; lines = pathlib.Path(sys.argv[1]).read_text().splitlines(); sys.exit(0 if sys.argv[2] + ": test" in lines else 1)' "$RUNNER_TEMP/node-embedding-tests.list" "$test_name" + result_path="$RUNNER_TEMP/node-embedding-${test_name##*::}.log" + cargo test -p node-engine --features inference-nodes --lib "$test_name" -- --exact 2>&1 | tee "$result_path" + python3 -c 'import pathlib, re, sys; text = pathlib.Path(sys.argv[1]).read_text(); sys.exit(0 if re.search(r"^test result: ok\. 1 passed; 0 failed; 0 ignored;", text, re.M) else 1)' "$result_path" + done + - name: Run node-engine tests run: cargo test -p node-engine --lib @@ -296,6 +314,8 @@ jobs: cargo test -p pantograph-workflow-service --lib scheduler::task_orchestrator::tests:: - name: Run workflow-service contract tests run: cargo test -p pantograph-workflow-service --test contract + - name: Run artifact-store integration contracts + run: cargo test -p pantograph-workflow-service --test artifact_store rust-doc-tests: name: Rust doc tests @@ -441,7 +461,7 @@ jobs: - name: Run warning-deny clippy audit id: run-clippy-audit continue-on-error: true - run: cargo clippy --workspace --all-targets --all-features -- -D warnings + run: cargo clippy --workspace --all-targets --all-features --keep-going -- -D warnings - name: Record clippy audit result id: report-clippy-audit diff --git a/Cargo.lock b/Cargo.lock index 4be64790e..1127ac859 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7552,6 +7552,7 @@ name = "pantograph-uniffi" version = "0.1.0" dependencies = [ "async-trait", + "futures-util", "graph-flow", "inference", "node-engine", diff --git a/bindings/csharp/README.md b/bindings/csharp/README.md index ff8566ece..febd8eed0 100644 --- a/bindings/csharp/README.md +++ b/bindings/csharp/README.md @@ -30,6 +30,16 @@ with any redistributed generator source/binary; this repository does not vendor or package the generator itself. Existing generated-binding/artifact licensing obligations remain unchanged. +## Shutdown result + +Await `runtime.Shutdown()` and handle the generated `FfiException` if the +owned backend cannot stop. Its message contains the standard JSON error envelope +with `internal_error` and the original shutdown cause. A failed stop does not +mean the runtime released its residency; callers may retry after resolving the +cause. Existing C# await expressions remain valid, but regenerate bindings and +ship them with the matching native library when adopting this fallible API. +Rust callers now handle `Result<(), FfiError>` explicitly. + ## Usage Run the repository-level smoke script: diff --git a/crates/inference/src/resource_monitor/mod.rs b/crates/inference/src/resource_monitor/mod.rs index 2a9d64f49..8ad12f028 100644 --- a/crates/inference/src/resource_monitor/mod.rs +++ b/crates/inference/src/resource_monitor/mod.rs @@ -135,7 +135,7 @@ mod tests { #[test] fn unsupported_resource_monitor_returns_typed_unavailable_observation() { - let monitor = unsupported::UnsupportedRuntimeResourceMonitor::default(); + let monitor = unsupported::UnsupportedRuntimeResourceMonitor; let guard = monitor .start_process_monitor(std::process::id()) .expect("unsupported monitor starts"); diff --git a/crates/node-engine/src/core_executor/inference_tests.rs b/crates/node-engine/src/core_executor/inference_tests.rs index 6e4660cbe..5261b7ab2 100644 --- a/crates/node-engine/src/core_executor/inference_tests.rs +++ b/crates/node-engine/src/core_executor/inference_tests.rs @@ -3177,9 +3177,12 @@ impl InferenceBackend for MockTypedTextBackend { } } +#[cfg(feature = "inference-nodes")] +type CapturedEmbeddingRequest = (Vec, String); + #[cfg(feature = "inference-nodes")] struct MockTypedEmbeddingBackend { - embedding_requests: Arc, String)>>>, + embedding_requests: Arc>>, } #[cfg(feature = "inference-nodes")] diff --git a/crates/node-engine/src/engine/dependency_inputs.rs b/crates/node-engine/src/engine/dependency_inputs.rs index 89869b4fb..df812fe07 100644 --- a/crates/node-engine/src/engine/dependency_inputs.rs +++ b/crates/node-engine/src/engine/dependency_inputs.rs @@ -421,7 +421,7 @@ mod tests { inputs.get("text"), Some(&serde_json::json!("generated text")) ); - assert!(inputs.get("stream").is_none()); + assert!(!inputs.contains_key("stream")); } #[test] diff --git a/crates/pantograph-diagnostics-ledger/src/tests.rs b/crates/pantograph-diagnostics-ledger/src/tests.rs index 5d8ff19ba..4ddeeafce 100644 --- a/crates/pantograph-diagnostics-ledger/src/tests.rs +++ b/crates/pantograph-diagnostics-ledger/src/tests.rs @@ -1571,12 +1571,9 @@ fn scheduler_timeline_projection_includes_inference_execution_diagnostics() { assert!(detail.contains("cache handle observed")); assert!(detail.contains("artifact refs 1")); assert!(detail.contains("kv cache restore_input hit")); - assert_eq!( - record - .payload_json - .contains("generated text should not appear"), - false - ); + assert!(!record + .payload_json + .contains("generated text should not appear")); } #[test] @@ -5080,13 +5077,10 @@ fn sqlite_column_exists(conn: &Connection, table_name: &str, column_name: &str) let mut stmt = conn .prepare(&format!("PRAGMA table_info({table_name})")) .expect("table info statement prepares"); - let columns = stmt + let mut columns = stmt .query_map([], |row| row.get::<_, String>(1)) .expect("table info query succeeds"); - let exists = columns - .map(|column| column.expect("column row loads")) - .any(|column| column == column_name); - exists + columns.any(|column| column.expect("column row loads") == column_name) } fn assert_columns_exist(conn: &Connection, table_name: &str, column_names: &[&str]) { diff --git a/crates/pantograph-embedded-runtime/src/lib_tests/data_graph_execution_tests.rs b/crates/pantograph-embedded-runtime/src/lib_tests/data_graph_execution_tests.rs index f90e8a9bf..3c3c1b01e 100644 --- a/crates/pantograph-embedded-runtime/src/lib_tests/data_graph_execution_tests.rs +++ b/crates/pantograph-embedded-runtime/src/lib_tests/data_graph_execution_tests.rs @@ -41,7 +41,7 @@ async fn execute_data_graph_retired_onnx_audio_path_does_not_call_python_sidecar .await .expect("data graph execution"); - assert!(outputs.get("audio").is_none()); + assert!(!outputs.contains_key("audio")); assert_eq!( outputs.get("_graph_id"), Some(&serde_json::json!("runtime-onnx-audio-data-graph")) 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 40fe05b9c..89336a945 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 @@ -450,7 +450,7 @@ fn image_runtime_validation_snapshot( workflow_semantic_version: version.semantic_version.clone(), workflow_execution_fingerprint: version.execution_fingerprint.clone(), descriptor_contract_version: INFERENCE_INTERFACE_CONTRACT_VERSION, - graph_revision: WorkflowGraphRevision::parse(&graph.compute_fingerprint()) + graph_revision: WorkflowGraphRevision::parse(graph.compute_fingerprint()) .expect("valid graph revision"), validation_session_id: DraftGraphValidationSessionId::parse( "embedded_runtime_validation_session_1", diff --git a/crates/pantograph-embedded-runtime/src/node_io_artifacts.rs b/crates/pantograph-embedded-runtime/src/node_io_artifacts.rs index d7ddc31e1..d5034f075 100644 --- a/crates/pantograph-embedded-runtime/src/node_io_artifacts.rs +++ b/crates/pantograph-embedded-runtime/src/node_io_artifacts.rs @@ -304,6 +304,73 @@ fn node_io_artifact_format_metadata(media_type: &str) -> ArtifactFormatMetadata } } +fn io_artifact_payload_kind(kind: ArtifactPayloadKind) -> IoArtifactPayloadKind { + match kind { + ArtifactPayloadKind::Text => IoArtifactPayloadKind::Text, + ArtifactPayloadKind::Image => IoArtifactPayloadKind::Image, + ArtifactPayloadKind::Audio => IoArtifactPayloadKind::Audio, + ArtifactPayloadKind::Video => IoArtifactPayloadKind::Video, + ArtifactPayloadKind::ThreeD => IoArtifactPayloadKind::ThreeD, + ArtifactPayloadKind::LargeTable => IoArtifactPayloadKind::LargeTable, + ArtifactPayloadKind::GenericBinary => IoArtifactPayloadKind::GenericBinary, + ArtifactPayloadKind::Structured => IoArtifactPayloadKind::Structured, + } +} + +fn io_artifact_lifecycle_state(state: ArtifactLifecycleState) -> IoArtifactLifecycleState { + match state { + ArtifactLifecycleState::Declared => IoArtifactLifecycleState::Declared, + ArtifactLifecycleState::Writing => IoArtifactLifecycleState::Writing, + ArtifactLifecycleState::Streaming => IoArtifactLifecycleState::Streaming, + ArtifactLifecycleState::Finalizing => IoArtifactLifecycleState::Finalizing, + ArtifactLifecycleState::Retained => IoArtifactLifecycleState::Retained, + ArtifactLifecycleState::Failed => IoArtifactLifecycleState::Failed, + ArtifactLifecycleState::Expired => IoArtifactLifecycleState::Expired, + ArtifactLifecycleState::Deleted => IoArtifactLifecycleState::Deleted, + } +} + +fn io_artifact_access_mode(mode: ArtifactAccessMode) -> IoArtifactAccessMode { + match mode { + ArtifactAccessMode::Read => IoArtifactAccessMode::Read, + ArtifactAccessMode::Download => IoArtifactAccessMode::Download, + ArtifactAccessMode::Stream => IoArtifactAccessMode::Stream, + } +} + +fn io_artifact_format_metadata(format: ArtifactFormatMetadata) -> IoArtifactFormatMetadata { + IoArtifactFormatMetadata { + format_id: format.format_id, + media_type: format.media_type, + codec_id: format.codec_id, + quality_percent: format.quality_percent, + bitrate_kbps: format.bitrate_kbps, + crf: format.crf, + bit_depth: format.bit_depth, + color_profile_id: format.color_profile_id, + converter_id: format.converter_id, + converter_version: format.converter_version, + library_version: format.library_version, + conversion_id: format.conversion_id, + conversion_status: format.conversion_status.map(|status| match status { + ArtifactConversionStatus::Converted => IoArtifactConversionStatus::Converted, + ArtifactConversionStatus::PassedThrough => IoArtifactConversionStatus::PassedThrough, + ArtifactConversionStatus::Failed => IoArtifactConversionStatus::Failed, + }), + conversion_command_id: format.conversion_command_id, + conversion_dependencies: format + .conversion_dependencies + .into_iter() + .map(|dependency| IoArtifactConversionDependency { + dependency_id: dependency.dependency_id, + active_version: dependency.active_version, + lease_id: dependency.lease_id, + lease_holder: dependency.lease_holder, + }) + .collect(), + } +} + #[cfg(test)] mod tests { use super::{ @@ -394,70 +461,3 @@ mod tests { ); } } - -fn io_artifact_payload_kind(kind: ArtifactPayloadKind) -> IoArtifactPayloadKind { - match kind { - ArtifactPayloadKind::Text => IoArtifactPayloadKind::Text, - ArtifactPayloadKind::Image => IoArtifactPayloadKind::Image, - ArtifactPayloadKind::Audio => IoArtifactPayloadKind::Audio, - ArtifactPayloadKind::Video => IoArtifactPayloadKind::Video, - ArtifactPayloadKind::ThreeD => IoArtifactPayloadKind::ThreeD, - ArtifactPayloadKind::LargeTable => IoArtifactPayloadKind::LargeTable, - ArtifactPayloadKind::GenericBinary => IoArtifactPayloadKind::GenericBinary, - ArtifactPayloadKind::Structured => IoArtifactPayloadKind::Structured, - } -} - -fn io_artifact_lifecycle_state(state: ArtifactLifecycleState) -> IoArtifactLifecycleState { - match state { - ArtifactLifecycleState::Declared => IoArtifactLifecycleState::Declared, - ArtifactLifecycleState::Writing => IoArtifactLifecycleState::Writing, - ArtifactLifecycleState::Streaming => IoArtifactLifecycleState::Streaming, - ArtifactLifecycleState::Finalizing => IoArtifactLifecycleState::Finalizing, - ArtifactLifecycleState::Retained => IoArtifactLifecycleState::Retained, - ArtifactLifecycleState::Failed => IoArtifactLifecycleState::Failed, - ArtifactLifecycleState::Expired => IoArtifactLifecycleState::Expired, - ArtifactLifecycleState::Deleted => IoArtifactLifecycleState::Deleted, - } -} - -fn io_artifact_access_mode(mode: ArtifactAccessMode) -> IoArtifactAccessMode { - match mode { - ArtifactAccessMode::Read => IoArtifactAccessMode::Read, - ArtifactAccessMode::Download => IoArtifactAccessMode::Download, - ArtifactAccessMode::Stream => IoArtifactAccessMode::Stream, - } -} - -fn io_artifact_format_metadata(format: ArtifactFormatMetadata) -> IoArtifactFormatMetadata { - IoArtifactFormatMetadata { - format_id: format.format_id, - media_type: format.media_type, - codec_id: format.codec_id, - quality_percent: format.quality_percent, - bitrate_kbps: format.bitrate_kbps, - crf: format.crf, - bit_depth: format.bit_depth, - color_profile_id: format.color_profile_id, - converter_id: format.converter_id, - converter_version: format.converter_version, - library_version: format.library_version, - conversion_id: format.conversion_id, - conversion_status: format.conversion_status.map(|status| match status { - ArtifactConversionStatus::Converted => IoArtifactConversionStatus::Converted, - ArtifactConversionStatus::PassedThrough => IoArtifactConversionStatus::PassedThrough, - ArtifactConversionStatus::Failed => IoArtifactConversionStatus::Failed, - }), - conversion_command_id: format.conversion_command_id, - conversion_dependencies: format - .conversion_dependencies - .into_iter() - .map(|dependency| IoArtifactConversionDependency { - dependency_id: dependency.dependency_id, - active_version: dependency.active_version, - lease_id: dependency.lease_id, - lease_holder: dependency.lease_holder, - }) - .collect(), - } -} diff --git a/crates/pantograph-embedded-runtime/src/task_executor_tests/puma_lib.rs b/crates/pantograph-embedded-runtime/src/task_executor_tests/puma_lib.rs index bce51235e..479ee3fa7 100644 --- a/crates/pantograph-embedded-runtime/src/task_executor_tests/puma_lib.rs +++ b/crates/pantograph-embedded-runtime/src/task_executor_tests/puma_lib.rs @@ -49,7 +49,7 @@ async fn puma_lib_execution_hydrates_model_ref_from_model_id_without_path_output .expect("puma-lib should resolve selector metadata"); assert!( - outputs.get("model_path").is_none(), + !outputs.contains_key("model_path"), "puma-lib must not emit executable path outputs" ); assert_eq!( @@ -74,15 +74,15 @@ async fn puma_lib_execution_hydrates_model_ref_from_model_id_without_path_output "puma-lib must not hide executable paths inside pumas_model_ref" ); assert!( - outputs.get("backend_key").is_none(), + !outputs.contains_key("backend_key"), "puma-lib must not emit graph-visible backend-key aliases" ); assert!( - outputs.get("resolved_model_package_facts").is_none(), + !outputs.contains_key("resolved_model_package_facts"), "puma-lib must not emit hidden package facts" ); assert!( - outputs.get("resolved_model_artifact_load_target").is_none(), + !outputs.contains_key("resolved_model_artifact_load_target"), "puma-lib must not emit hidden artifact load targets" ); } @@ -117,7 +117,7 @@ async fn puma_lib_execution_preserves_explicit_model_ref_without_model_path_outp .expect("puma-lib should preserve explicit model ref"); assert!( - outputs.get("model_path").is_none(), + !outputs.contains_key("model_path"), "puma-lib must not promote selected artifact paths to executable outputs" ); assert_eq!( @@ -133,7 +133,7 @@ async fn puma_lib_execution_preserves_explicit_model_ref_without_model_path_outp )) ); assert!( - outputs.get("backend_key").is_none(), + !outputs.contains_key("backend_key"), "puma-lib must not emit graph-visible backend-key aliases" ); } @@ -182,7 +182,7 @@ async fn puma_lib_execution_does_not_rebind_model_id_from_raw_pumas_api() { .expect("puma-lib should preserve saved data without selector access"); assert!( - outputs.get("model_path").is_none(), + !outputs.contains_key("model_path"), "raw saved paths must not be emitted as executable outputs" ); assert_eq!( @@ -196,15 +196,15 @@ async fn puma_lib_execution_does_not_rebind_model_id_from_raw_pumas_api() { Some(&serde_json::json!(model_id)) ); assert!( - outputs.get("backend_key").is_none(), + !outputs.contains_key("backend_key"), "puma-lib must not emit graph-visible backend-key aliases" ); assert!( - outputs.get("resolved_model_package_facts").is_none(), + !outputs.contains_key("resolved_model_package_facts"), "raw PUMAS_API alone must not rehydrate selected model facts" ); assert!( - outputs.get("resolved_model_artifact_load_target").is_none(), + !outputs.contains_key("resolved_model_artifact_load_target"), "raw PUMAS_API alone must not rehydrate selected artifact load targets" ); } @@ -260,7 +260,7 @@ async fn puma_lib_execution_hydrates_model_ref_from_selector_access_without_puma .expect("puma-lib should resolve selector metadata from selector access"); assert!( - outputs.get("model_path").is_none(), + !outputs.contains_key("model_path"), "puma-lib must not emit executable path outputs" ); assert_eq!(outputs.get("model_id"), Some(&serde_json::json!(model_id))); @@ -282,11 +282,11 @@ async fn puma_lib_execution_hydrates_model_ref_from_selector_access_without_puma "puma-lib must not hide executable paths inside pumas_model_ref" ); assert!( - outputs.get("resolved_model_package_facts").is_none(), + !outputs.contains_key("resolved_model_package_facts"), "read-only selector rows must not be promoted to full package facts" ); assert!( - outputs.get("resolved_model_artifact_load_target").is_none(), + !outputs.contains_key("resolved_model_artifact_load_target"), "read-only selector rows must not be promoted to artifact load targets" ); } @@ -349,7 +349,7 @@ async fn puma_lib_execution_does_not_emit_inference_settings_from_saved_or_selec .expect("puma-lib should resolve selected detail without inference settings"); assert!( - outputs.get("inference_settings").is_none(), + !outputs.contains_key("inference_settings"), "puma-lib must not emit inference settings from saved node data or selected detail" ); } @@ -396,13 +396,13 @@ async fn puma_lib_execution_does_not_resolve_saved_model_name_without_model_id() .expect("puma-lib should execute with saved data only"); assert!( - outputs.get("model_path").is_none(), + !outputs.contains_key("model_path"), "puma-lib must not emit empty path outputs" ); - assert!(outputs.get("model_id").is_none()); - assert!(outputs.get("model_type").is_none()); - assert!(outputs.get("task_type_primary").is_none()); - assert!(outputs.get("resolved_model_package_facts").is_none()); - assert!(outputs.get("resolved_model_artifact_load_target").is_none()); - assert!(outputs.get("pumas_model_ref").is_none()); + assert!(!outputs.contains_key("model_id")); + assert!(!outputs.contains_key("model_type")); + assert!(!outputs.contains_key("task_type_primary")); + assert!(!outputs.contains_key("resolved_model_package_facts")); + assert!(!outputs.contains_key("resolved_model_artifact_load_target")); + assert!(!outputs.contains_key("pumas_model_ref")); } diff --git a/crates/pantograph-embedded-runtime/src/technical_fit.rs b/crates/pantograph-embedded-runtime/src/technical_fit.rs index 74d1534e3..5cf66eef6 100644 --- a/crates/pantograph-embedded-runtime/src/technical_fit.rs +++ b/crates/pantograph-embedded-runtime/src/technical_fit.rs @@ -1837,7 +1837,7 @@ mod tests { kv_cache: inference::BackendFeatureSupport::Supported, }, runtime_variants: vec![inference::RuntimeVariantCapability { - runtime_variant_id: inference::RuntimeVariantId::parse(&format!( + runtime_variant_id: inference::RuntimeVariantId::parse(format!( "{}.cuda", backend_key )) diff --git a/crates/pantograph-managed-dependencies/src/redistributables/paths.rs b/crates/pantograph-managed-dependencies/src/redistributables/paths.rs index 8e89e663a..7f623f046 100644 --- a/crates/pantograph-managed-dependencies/src/redistributables/paths.rs +++ b/crates/pantograph-managed-dependencies/src/redistributables/paths.rs @@ -121,6 +121,28 @@ pub(crate) fn library_path(name: &str) -> String { } } +pub(crate) fn sanitize_path_segment(segment: &str) -> String { + segment + .chars() + .map(|character| { + if character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_') { + character + } else { + '_' + } + }) + .collect() +} + +pub(crate) fn current_unix_timestamp_ms() -> u64 { + use std::time::{SystemTime, UNIX_EPOCH}; + + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64 +} + #[cfg(test)] mod tests { use super::{managed_redistributable_version_dir, ManagedRedistributableId}; @@ -146,25 +168,3 @@ mod tests { ); } } - -pub(crate) fn sanitize_path_segment(segment: &str) -> String { - segment - .chars() - .map(|character| { - if character.is_ascii_alphanumeric() || matches!(character, '.' | '-' | '_') { - character - } else { - '_' - } - }) - .collect() -} - -pub(crate) fn current_unix_timestamp_ms() -> u64 { - use std::time::{SystemTime, UNIX_EPOCH}; - - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_millis() as u64 -} diff --git a/crates/pantograph-scheduler/tests/queue_state.rs b/crates/pantograph-scheduler/tests/queue_state.rs index 766d2814e..ec04122b3 100644 --- a/crates/pantograph-scheduler/tests/queue_state.rs +++ b/crates/pantograph-scheduler/tests/queue_state.rs @@ -34,8 +34,10 @@ fn boxed_runtime_intent_preserves_exact_json_and_borrowed_payload() { #[test] fn boxed_runtime_intent_retains_validation_and_task_correlation() { let valid = task_intent("run.001", "task.001"); - ValidatedSchedulerTaskStateRecord::try_from(task_record_with_state(ready_state(valid.clone()))) + let expected_record = task_record_with_state(ready_state(valid.clone())); + let validated = ValidatedSchedulerTaskStateRecord::try_from(expected_record.clone()) .expect("valid runtime intent remains accepted"); + assert_eq!(validated.as_ref(), &expected_record); let mut invalid = valid; invalid.contract_version = 0; let expected = invalid.validate().expect_err("invalid raw version"); diff --git a/crates/pantograph-uniffi/Cargo.toml b/crates/pantograph-uniffi/Cargo.toml index ddd64f636..68b0e573a 100644 --- a/crates/pantograph-uniffi/Cargo.toml +++ b/crates/pantograph-uniffi/Cargo.toml @@ -54,6 +54,7 @@ required-features = ["cli"] uniffi = { version = "0.28", features = ["build"] } [dev-dependencies] +futures-util.workspace = true uniffi = { version = "0.28", features = ["bindgen-tests"] } [features] diff --git a/crates/pantograph-uniffi/src/runtime.rs b/crates/pantograph-uniffi/src/runtime.rs index 1e4ab0b11..22a0b07d5 100644 --- a/crates/pantograph-uniffi/src/runtime.rs +++ b/crates/pantograph-uniffi/src/runtime.rs @@ -208,8 +208,13 @@ impl FfiPantographRuntime { } /// Stop inference backends owned by this runtime. - pub async fn shutdown(&self) { - self.runtime.shutdown().await; + /// + /// A failed shutdown is returned through the standard error envelope; the + /// embedded owner retains residency and permits a later retry. + pub async fn shutdown(&self) -> Result<(), FfiError> { + self.runtime.shutdown().await.map_err(|error| { + workflow_adapter_error(WorkflowErrorCode::InternalError, error.to_string()) + }) } /// Register an attribution client and return ClientRegistrationResponse JSON. @@ -1116,3 +1121,7 @@ mod runtime_tests; #[cfg(test)] #[path = "runtime_validation_tests.rs"] mod runtime_validation_tests; + +#[cfg(test)] +#[path = "runtime_shutdown_tests.rs"] +mod runtime_shutdown_tests; diff --git a/crates/pantograph-uniffi/src/runtime_shutdown_tests.rs b/crates/pantograph-uniffi/src/runtime_shutdown_tests.rs new file mode 100644 index 000000000..940b3f0c8 --- /dev/null +++ b/crates/pantograph-uniffi/src/runtime_shutdown_tests.rs @@ -0,0 +1,159 @@ +use std::pin::Pin; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::Arc; + +use async_trait::async_trait; +use futures_util::Stream; +use inference::backend::BackendStartOutcome; +use inference::{ + BackendCapabilities, BackendConfig, BackendError, ChatChunk, EmbeddingResult, InferenceBackend, + InferenceGateway, ProcessSpawner, RerankRequest, RerankResponse, +}; +use node_engine::ExecutorExtensions; +use pantograph_embedded_runtime::{EmbeddedRuntime, EmbeddedRuntimeConfig}; +use pantograph_workflow_service::WorkflowService; +use tokio::sync::RwLock; + +use super::FfiPantographRuntime; +use crate::FfiError; + +struct ShutdownBackend { + fail_next_stop: AtomicBool, + ready: AtomicBool, + stop_calls: Arc, +} + +#[async_trait] +impl InferenceBackend for ShutdownBackend { + fn name(&self) -> &'static str { + "shutdown-fixture" + } + fn description(&self) -> &'static str { + "bounded shutdown fixture" + } + fn capabilities(&self) -> BackendCapabilities { + BackendCapabilities::default() + } + async fn start( + &mut self, + _: &BackendConfig, + _: Arc, + ) -> Result { + panic!("shutdown fixture must not start a backend") + } + async fn stop(&mut self) -> Result<(), BackendError> { + self.stop_calls.fetch_add(1, Ordering::SeqCst); + if self.fail_next_stop.swap(false, Ordering::SeqCst) { + return Err(BackendError::Inference("fixture stop failure".to_string())); + } + self.ready.store(false, Ordering::SeqCst); + Ok(()) + } + fn is_ready(&self) -> bool { + self.ready.load(Ordering::SeqCst) + } + async fn health_check(&self) -> bool { + self.is_ready() + } + fn base_url(&self) -> Option { + None + } + async fn chat_completion_stream( + &self, + _: String, + ) -> Result> + Send>>, BackendError> + { + panic!("shutdown fixture must not execute inference") + } + async fn embeddings( + &self, + _: Vec, + _: &str, + ) -> Result, BackendError> { + panic!("shutdown fixture must not execute embeddings") + } + async fn rerank(&self, _: RerankRequest) -> Result { + panic!("shutdown fixture must not execute reranking") + } +} + +fn runtime( + fail_next_stop: bool, +) -> ( + FfiPantographRuntime, + Arc, + Arc, + std::path::PathBuf, +) { + let root = super::runtime_tests::create_temp_root("shutdown-fixture"); + let stop_calls = Arc::new(AtomicUsize::new(0)); + let gateway = Arc::new(InferenceGateway::with_backend( + Box::new(ShutdownBackend { + fail_next_stop: AtomicBool::new(fail_next_stop), + ready: AtomicBool::new(true), + stop_calls: stop_calls.clone(), + }), + "shutdown-fixture", + )); + let extensions = Arc::new(RwLock::new(ExecutorExtensions::new())); + let app_data_dir = root.join("app-data"); + let embedded = EmbeddedRuntime::with_default_python_runtime( + EmbeddedRuntimeConfig::new(app_data_dir.clone(), root.clone()), + gateway.clone(), + extensions.clone(), + Arc::new(WorkflowService::new()), + None, + ); + ( + FfiPantographRuntime { + runtime: Arc::new(embedded), + app_data_dir, + node_registry: Arc::new(node_engine::NodeRegistry::with_builtins()), + extensions, + }, + gateway, + stop_calls, + root, + ) +} + +#[tokio::test] +async fn ffi_shutdown_reports_owner_failure_and_allows_retry() { + let (runtime, gateway, stop_calls, root) = runtime(true); + assert!(gateway.is_ready().await); + let error = runtime + .shutdown() + .await + .expect_err("failed owner stop must cross FFI"); + let FfiError::Other { message } = error else { + panic!("expected standard FFI error envelope"); + }; + assert_eq!( + serde_json::from_str::(&message).unwrap(), + serde_json::json!({ + "code": "internal_error", "message": "Backend error: Inference error: fixture stop failure" + }) + ); + assert_eq!(stop_calls.load(Ordering::SeqCst), 1); + assert!( + gateway.is_ready().await, + "failed stop must preserve backend readiness" + ); + runtime + .shutdown() + .await + .expect("retry observes successful owner stop"); + assert_eq!(stop_calls.load(Ordering::SeqCst), 2); + assert!(!gateway.is_ready().await); + std::fs::remove_dir_all(root).expect("remove fixture files"); +} + +#[tokio::test] +async fn ffi_shutdown_success_and_repeated_success_are_explicit() { + let (runtime, gateway, stop_calls, root) = runtime(false); + runtime.shutdown().await.expect("successful shutdown"); + assert!(!gateway.is_ready().await); + runtime.shutdown().await.expect("repeated shutdown"); + assert_eq!(stop_calls.load(Ordering::SeqCst), 2); + std::fs::remove_dir_all(root).expect("remove fixture files"); +} diff --git a/crates/pantograph-uniffi/src/runtime_tests.rs b/crates/pantograph-uniffi/src/runtime_tests.rs index 9b3492f5e..481f2eade 100644 --- a/crates/pantograph-uniffi/src/runtime_tests.rs +++ b/crates/pantograph-uniffi/src/runtime_tests.rs @@ -371,7 +371,7 @@ async fn direct_runtime_runs_workflow_session_from_json() { let close: serde_json::Value = serde_json::from_str(&close_json).expect("parse close"); assert_eq!(close["ok"], true); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -434,7 +434,7 @@ async fn direct_runtime_rejects_interactive_graph_at_scheduler_boundary() { .contains("scheduler task session runner has no execution path")); assert!(envelope.message.contains("unsupported=1")); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -490,7 +490,7 @@ async fn direct_runtime_exposes_attribution_client_session_json() { serde_json::from_str(&open_session_json).expect("parse open response"); assert!(opened["session"]["client_session_id"].as_str().is_some()); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -610,7 +610,7 @@ async fn direct_runtime_exposes_workflow_graph_persistence_and_edit_session() { ) .await .expect("close graph edit session"); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -701,7 +701,7 @@ async fn direct_runtime_exposes_backend_owned_graph_authoring_discovery() { .message .contains("No options provider for text-input:text")); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -781,7 +781,7 @@ async fn direct_runtime_puma_lib_options_use_selector_access_from_pumas_api() { .as_str() .is_some_and(|cursor| cursor.starts_with("model-library-updates:"))); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -841,7 +841,7 @@ async fn direct_runtime_exposes_artifact_store_contract_surface() { assert_eq!(envelope.code, WorkflowErrorCode::InvalidRequest); assert!(envelope.message.contains("artifact not found")); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -987,7 +987,7 @@ async fn direct_runtime_exposes_artifact_format_settings_and_capabilities_json() "image quality_percent 0 is outside allowed range" ); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } @@ -1184,7 +1184,7 @@ async fn direct_runtime_exposes_managed_media_dependency_statuses_and_actions_js assert!(envelope.message.contains("missing expected file")); } - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); let _ = std::fs::remove_dir_all(root); } diff --git a/crates/pantograph-uniffi/src/runtime_validation_tests.rs b/crates/pantograph-uniffi/src/runtime_validation_tests.rs index b1c09bfa7..96fc1b288 100644 --- a/crates/pantograph-uniffi/src/runtime_validation_tests.rs +++ b/crates/pantograph-uniffi/src/runtime_validation_tests.rs @@ -170,7 +170,7 @@ async fn public_validation_publication_is_required_again_after_runtime_reopen() .expect("published text runs"), ); assert_eq!(response["outputs"][0]["value"], "published text"); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); drop(runtime); let reopened = open_runtime(&root).await; @@ -186,7 +186,7 @@ async fn public_validation_publication_is_required_again_after_runtime_reopen() .expect("republished text runs"), ); assert_eq!(response["outputs"][0]["value"], "published text"); - reopened.shutdown().await; + reopened.shutdown().await.expect("runtime shutdown"); drop(reopened); std::fs::remove_dir_all(root).expect("remove fixture"); } @@ -234,7 +234,7 @@ async fn public_publication_rejects_stale_validation_and_caller_supplied_proof() .await .expect_err("rejected publication stores no executable snapshot"), ); - runtime.shutdown().await; + runtime.shutdown().await.expect("runtime shutdown"); drop(runtime); std::fs::remove_dir_all(root).expect("remove fixture"); } diff --git a/crates/pantograph-workflow-service/src/graph/connection_insert.rs b/crates/pantograph-workflow-service/src/graph/connection_insert.rs index e436aa8f5..bd49a8b39 100644 --- a/crates/pantograph-workflow-service/src/graph/connection_insert.rs +++ b/crates/pantograph-workflow-service/src/graph/connection_insert.rs @@ -383,6 +383,34 @@ pub fn rejected_insert_response( } } +pub fn rejected_edge_insert_preview_response( + graph: &WorkflowGraph, + rejection: ConnectionRejection, +) -> EdgeInsertionPreviewResponse { + EdgeInsertionPreviewResponse { + accepted: false, + graph_revision: graph.compute_fingerprint(), + bridge: None, + rejection: Some(rejection), + } +} + +pub fn rejected_insert_on_edge_response( + graph: &WorkflowGraph, + rejection: ConnectionRejection, +) -> InsertNodeOnEdgeResponse { + InsertNodeOnEdgeResponse { + accepted: false, + graph_revision: graph.compute_fingerprint(), + inserted_node_id: None, + bridge: None, + graph: Some(graph.clone()), + workflow_event: None, + workflow_execution_session_state: None, + rejection: Some(rejection), + } +} + #[cfg(test)] mod tests { use super::*; @@ -419,31 +447,3 @@ mod tests { ); } } - -pub fn rejected_edge_insert_preview_response( - graph: &WorkflowGraph, - rejection: ConnectionRejection, -) -> EdgeInsertionPreviewResponse { - EdgeInsertionPreviewResponse { - accepted: false, - graph_revision: graph.compute_fingerprint(), - bridge: None, - rejection: Some(rejection), - } -} - -pub fn rejected_insert_on_edge_response( - graph: &WorkflowGraph, - rejection: ConnectionRejection, -) -> InsertNodeOnEdgeResponse { - InsertNodeOnEdgeResponse { - accepted: false, - graph_revision: graph.compute_fingerprint(), - inserted_node_id: None, - bridge: None, - graph: Some(graph.clone()), - workflow_event: None, - workflow_execution_session_state: None, - rejection: Some(rejection), - } -} diff --git a/crates/pantograph-workflow-service/src/graph/registry.rs b/crates/pantograph-workflow-service/src/graph/registry.rs index 0850822d0..1828c6433 100644 --- a/crates/pantograph-workflow-service/src/graph/registry.rs +++ b/crates/pantograph-workflow-service/src/graph/registry.rs @@ -166,7 +166,7 @@ mod tests { let payload = &encoded["inference_payloads"][0]; assert_eq!(payload["role"], serde_json::json!("diagnostics")); assert_eq!(payload["task_id"], serde_json::json!("text_generation")); - assert_llm_inference_payloads_do_not_expose_runtime_policy(&definition); + assert_llm_inference_payloads_do_not_expose_runtime_policy(definition); } #[test] 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 be0d264c8..b91ca792c 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 @@ -1981,7 +1981,7 @@ mod tests { plan.runtime_host_request .members .iter() - .map(|member| prompt_text_from_runtime_host_member_request(member)) + .map(prompt_text_from_runtime_host_member_request) .collect::>(), vec![ "prompt owned by run.2026-05-22.001".to_string(), diff --git a/crates/pantograph-workflow-service/src/workflow/service_config.rs b/crates/pantograph-workflow-service/src/workflow/service_config.rs index 1fe7892f8..7baeca144 100644 --- a/crates/pantograph-workflow-service/src/workflow/service_config.rs +++ b/crates/pantograph-workflow-service/src/workflow/service_config.rs @@ -439,6 +439,26 @@ fn dependency_requirements_registry_error( )) } +fn load_artifact_format_settings( + path: &Path, +) -> Result { + if !path.exists() { + return Ok(ArtifactFormatSettings::default()); + } + let content = std::fs::read_to_string(path).map_err(|error| { + WorkflowServiceError::Internal(format!( + "failed to read artifact format settings {:?}: {error}", + path + )) + })?; + serde_json::from_str(&content).map_err(|error| { + WorkflowServiceError::InvalidRequest(format!( + "artifact format settings file {:?} is invalid: {error}", + path + )) + }) +} + #[cfg(test)] mod tests { use std::sync::{Arc, Mutex}; @@ -594,23 +614,3 @@ mod tests { } } } - -fn load_artifact_format_settings( - path: &Path, -) -> Result { - if !path.exists() { - return Ok(ArtifactFormatSettings::default()); - } - let content = std::fs::read_to_string(path).map_err(|error| { - WorkflowServiceError::Internal(format!( - "failed to read artifact format settings {:?}: {error}", - path - )) - })?; - serde_json::from_str(&content).map_err(|error| { - WorkflowServiceError::InvalidRequest(format!( - "artifact format settings file {:?} is invalid: {error}", - path - )) - }) -} diff --git a/crates/pantograph-workflow-service/src/workflow/tests/session_execution.rs b/crates/pantograph-workflow-service/src/workflow/tests/session_execution.rs index fa52cdfee..311f1a1cd 100644 --- a/crates/pantograph-workflow-service/src/workflow/tests/session_execution.rs +++ b/crates/pantograph-workflow-service/src/workflow/tests/session_execution.rs @@ -4801,7 +4801,7 @@ fn runtime_executable_validation_snapshot( workflow_semantic_version: version.semantic_version.clone(), workflow_execution_fingerprint: version.execution_fingerprint.clone(), descriptor_contract_version: INFERENCE_INTERFACE_CONTRACT_VERSION, - graph_revision: WorkflowGraphRevision::parse(&graph.compute_fingerprint()) + graph_revision: WorkflowGraphRevision::parse(graph.compute_fingerprint()) .expect("valid graph revision"), validation_session_id: DraftGraphValidationSessionId::parse("runtime_validation_session_1") .expect("valid validation session id"), diff --git a/crates/pantograph-workflow-service/src/workflow/tests/task_result_contracts.rs b/crates/pantograph-workflow-service/src/workflow/tests/task_result_contracts.rs index fd3b1e8be..bea71a2ee 100644 --- a/crates/pantograph-workflow-service/src/workflow/tests/task_result_contracts.rs +++ b/crates/pantograph-workflow-service/src/workflow/tests/task_result_contracts.rs @@ -61,9 +61,9 @@ fn scheduler_task_result_validates_path_free_typed_outputs() { encoded["outputs"][1]["value"]["value"]["artifact_id"], "artifact-image-1" ); - assert_eq!(encoded.to_string().contains("model_path"), false); - assert_eq!(encoded.to_string().contains("local_load_path"), false); - assert_eq!(encoded.to_string().contains("runtime_handoff"), false); + assert!(!encoded.to_string().contains("model_path")); + assert!(!encoded.to_string().contains("local_load_path")); + assert!(!encoded.to_string().contains("runtime_handoff")); } #[test] diff --git a/crates/pantograph-workflow-service/src/workflow/tests/workflow_version.rs b/crates/pantograph-workflow-service/src/workflow/tests/workflow_version.rs index 762ff0cb0..79bbcb226 100644 --- a/crates/pantograph-workflow-service/src/workflow/tests/workflow_version.rs +++ b/crates/pantograph-workflow-service/src/workflow/tests/workflow_version.rs @@ -423,7 +423,7 @@ fn executable_validation_publication( graph: &WorkflowGraph, ) -> WorkflowGraphInferenceValidationPublication { let graph_revision = - WorkflowGraphRevision::parse(&graph.compute_fingerprint()).expect("valid graph revision"); + WorkflowGraphRevision::parse(graph.compute_fingerprint()).expect("valid graph revision"); let validation_session_id = DraftGraphValidationSessionId::parse("validation_session_publish") .expect("valid validation session id"); let summary = DraftGraphValidationSummary { diff --git a/crates/pantograph-workflow-service/src/workflow/validation.rs b/crates/pantograph-workflow-service/src/workflow/validation.rs index 7ba8508a9..d0588c9e8 100644 --- a/crates/pantograph-workflow-service/src/workflow/validation.rs +++ b/crates/pantograph-workflow-service/src/workflow/validation.rs @@ -154,23 +154,6 @@ pub(super) fn validate_bindings( Ok(()) } -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn stale_graph_remaining_count_rejects_formatter_underflow() { - let error = remaining_stale_graph_diagnostic_count(1, 2) - .expect_err("shown diagnostics cannot exceed total diagnostics"); - - assert!(matches!( - error, - WorkflowServiceError::Internal(message) - if message.contains("formatter showed 2 reasons for 1 diagnostics") - )); - } -} - pub(super) fn validate_host_output_bindings( bindings: &[WorkflowPortBinding], field_name: &str, @@ -340,3 +323,20 @@ pub(super) fn validate_payload_size( Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn stale_graph_remaining_count_rejects_formatter_underflow() { + let error = remaining_stale_graph_diagnostic_count(1, 2) + .expect_err("shown diagnostics cannot exceed total diagnostics"); + + assert!(matches!( + error, + WorkflowServiceError::Internal(message) + if message.contains("formatter showed 2 reasons for 1 diagnostics") + )); + } +} diff --git a/crates/pantograph-workflow-service/tests/artifact_store.rs b/crates/pantograph-workflow-service/tests/artifact_store.rs index 2178eea24..236531ae3 100644 --- a/crates/pantograph-workflow-service/tests/artifact_store.rs +++ b/crates/pantograph-workflow-service/tests/artifact_store.rs @@ -83,7 +83,7 @@ fn artifact_store_writes_descriptor_and_reads_body_without_path_leak() { ); assert!(!serde_json::to_string(&descriptor) .expect("serialize descriptor") - .contains(temp.path().to_string_lossy().as_ref())); + .contains::<&str>(temp.path().to_string_lossy().as_ref())); let read = store .read_body(ArtifactReadRequest { diff --git a/crates/workflow-nodes/src/input/puma_lib.rs b/crates/workflow-nodes/src/input/puma_lib.rs index 6c6680cdb..2d7c09e86 100644 --- a/crates/workflow-nodes/src/input/puma_lib.rs +++ b/crates/workflow-nodes/src/input/puma_lib.rs @@ -271,7 +271,7 @@ mod options_provider { } pub(crate) fn needs_resolution(&self, model_id: &str) -> bool { - self.summaries.get(model_id).map_or(true, |result| { + self.summaries.get(model_id).is_none_or(|result| { result.summary.is_none() || matches!( result.status, diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-aggregate-test-style.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-aggregate-test-style.md new file mode 100644 index 000000000..4f792788e --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-aggregate-test-style.md @@ -0,0 +1,15 @@ +# Aggregate mechanical and test qualification repair + +The first keep-going strict audit at 9a46796 emitted 51 diagnostics at 50 unique source sites across eight crate/desktop groups. Four desktop command/delegate arity sites are excluded for a wire-compatible design. This milestone addresses the remaining 46 sites without changing public signatures, data contracts, errors or validation policy. + +Exact semantic groups: 23 absent-key assertions become !contains_key (a present null value remains present); six redundant borrows are removed; four literal-false comparisons become negated assertions; one unit Default becomes its unit value; one Option::map_or(true, predicate) becomes is_none_or with the same predicate; one redundant closure becomes its function; one iterator result returns directly; one successful must-use scheduler validation now compares the complete validated raw record; and one desktop base64-length divisibility check uses is_multiple_of(4). The artifact-store path-redaction assertion explicitly selects the &str Pattern type, resolving the all-feature Cow AsRef ambiguity caused by typed_path without allocation or weakening the assertion. + +Six existing cfg(test) modules move after production items in node_io_artifacts, connection_insert, service_config, validation, redistributables/paths and desktop llm/commands/agent. Each moved test-module body is byte-identical after formatting; a Rust lexer comparison verifies all production tokens (including literals and comments) are identical before/after these moves. Production item ordering is retained. + +The required focused workflow retains all existing suites and adds the complete artifact_store integration target (13 source tests), so the compile-blocked redaction/storage contract now executes. Existing embedded, workflow-service, scheduler, ledger, inference and workflow-node suites remain. Desktop/test-only source changes still require actual hosted strict all-target compilation; no blanket lint allowance or assertion removal is introduced. + +Root requested this grouped mechanical milestone from the complete diagnostic inventory. Root source review accepted frozen tree 29cb5ce693a0b6a091a807d3cfa3b203a60468d3 after reviewing the non-move changes and independently checking identical nonblank line multisets for all six module moves. Exact-head hosted execution remains pending. Full pinned cargo fmt is checked locally; no heavy local Rust build is performed. The separate four Tauri arity findings remain open and must preserve existing IPC wire keys. + +## Hosted qualification correction + +PR43 head b2490ae reached the full ledger target in Quality run 37149377334, job 111279795852, and exposed E0597 in the test-only sqlite_column_exists helper. Returning the Map iterator expression extended its temporary borrow beyond the statement's lifetime under the crate's Rust 2021 edition. The original direct-return lint repair was therefore not source-compatible despite its equivalent row predicate. Bind the mapped rows mutably and call any with the same expect and equality directly; the named iterator then drops before the statement, with the same short-circuit order and row-error panic. No production code, SQL, assertions or error text changes. Full ledger and remaining focused suites require a fresh hosted run; the failed run is retained as evidence. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-node-engine-test-lints.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-node-engine-test-lints.md new file mode 100644 index 000000000..995b8dbd3 --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-node-engine-test-lints.md @@ -0,0 +1,11 @@ +# Node-engine test-only lint repair + +Fresh strict Clippy at PR43 head e2dbe87 clears the UniFFI library and reaches two node-engine lib-test findings. Both exact source expressions also exist at main 4938e405. This milestone only names the captured embedding request tuple (Vec, String) with a private feature-gated alias and expresses absent stream membership with !contains_key instead of get().is_none(). All tuple ordering, values, test actions and assertions are unchanged; there is no production API or behavior change. + +The embedding mock has two constructors and two existing test readers, plus one push site. Existing tests assert captured input/model and compatibility lifecycle/output behavior. Because the default node-engine suite does not enable inference-nodes, add a feature-enabled gate requiring both exact existing canonical embedding test names and an individual execution receipt of 1 passed / 0 failed / 0 ignored for each, retaining the full default suite and strict aggregate lint. + +Root source review accepted the mechanical repairs and the corrected exact-name/execution-count gate at frozen tree 2b5bc9252887e86bed791b4c759923592ff9435d. Fresh hosted execution remains pending. Full pinned formatting passes locally; no local Rust build is claimed. PR43 shutdown success/failure and regenerated C# qualification continue independently on its exact published head. + +## Complete independent diagnostics without weakening failure + +Root separately reviewed and accepted adding only --keep-going to the existing strict workspace/all-targets/all-features Clippy invocation. Pinned Cargo 1.92 lists this option, and Clippy delegates these options to cargo check. A dependency-free two-crate local probe emitted both independent denied warnings and still exited 101 (56 KiB output directory, 0.13 seconds). The workflow retains -D warnings and its explicit failing audit/aggregate result. This permits diagnostics from independent compilable crates; dependencies that fail compilation can still prevent downstream checks. No lint is suppressed and no warning becomes acceptable. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-tauri-command-state.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-tauri-command-state.md new file mode 100644 index 000000000..91792cf7a --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-tauri-command-state.md @@ -0,0 +1,21 @@ +# Tauri injected-state bundles with unchanged IPC payloads + +The keep-going strict audit reports four nine-argument sites: workflow_run_execution_session and query_port_options command wrappers plus their two sole delegates. Bundle only framework-injected State arguments in two operation-specific CommandArg implementations. No new managed owner, serialized request, client field, default or service lookup policy is introduced. + +WorkflowRunCommandState preserves this exact extraction sequence and macro-generated keys: gateway, runtimeRegistry, extensions, ragManager, workflowService, diagnosticsStore. PortOptionsCommandState preserves registry, extensions, workflowService. Each extraction calls Tauri State::from_command with the original plugin, command name, InvokeMessage and ACL reference, substituting only the original field key. The bundles hold the same borrowed State values and do not implement Deserialize. Client-supplied state fields cannot replace application-managed values. Missing-state errors retain the framework's original key/name text and first-failure order. + +The two wrappers have one delegate caller each; both delegates destructure the corresponding bundle then execute byte-identical previous bodies. The execution command still receives request, AppHandle and channel in their original relative order. The query command still receives flat nodeType, portId, search, limit, offset and context. The three frontend invokers (WorkflowCommandService, portOptionsCache and pumaModelOptionsCache) are unchanged. This is an internal Rust parameter grouping, not a Tauri wire-key rename. + +Five tests use the official Tauri MockRuntime/IPC path. They verify all six/three managed Arc identities, exact missing-state order/keys, and the real query_port_options handler with full and minimal flat payloads, original optional defaults, provider callbacks and managed extension data. Deliberately supplied client state values are ignored. A dormant gateway's test ProcessSpawner panics if unexpectedly used; no model/process/database setup is part of state extraction. Session state extraction is exercised through a test command using the same production bundle; the real session delegate body is separately byte-compared unchanged, rather than claimed executed by that test. + +The existing Tauri dependency gains only its test-only test feature (empty upstream feature, same locked 2.9.5 package). Required hosted workspace verification adds a script requiring all five exact test names and an individual 1 passed / 0 failed / 0 ignored receipt for each, with backend-llamacpp and no default app features. Existing all-feature workspace compilation and strict all-target Clippy remain. Frontend invocation tests remain part of the existing full frontend suite. + +Root source review accepted corrected tree 58ef1d4de278b899c3d2c72553275d7d3697f1bd after inspecting the import placement and exact pinned-version evidence. This is source acceptance, not Rust execution. After the user approved recovery, the restored official Rust 1.92.0 formatter found formatting changes in three IPC files. Applying pinned cargo fmt --all and rerunning the full cargo fmt --all -- --check passed. The formatting delta contains whitespace, optional trailing commas and two unchanged reordered imports; token comparison preserves the remaining source, literals and comments. Bash syntax and the unchanged delegate-body comparison also pass. Rust compilation and all five native mock-IPC executions remain pending hosted qualification. The optional formatting-artifact workflow is omitted because local pinned formatting is now available. + +## Exact pinned-source compatibility + +Cargo.lock pins Tauri 2.9.5 and tauri-macros 2.5.2. The official tauri-v2.9.5 tag resolves to f2e0405dc2d8bae949de6bb5e479bbf3cde401b6, whose two manifests match those versions. [CommandItem and CommandArg](https://github.com/tauri-apps/tauri/blob/f2e0405dc2d8bae949de6bb5e479bbf3cde401b6/crates/tauri/src/ipc/command.rs) expose exactly the fields and method used here. [State extraction](https://github.com/tauri-apps/tauri/blob/f2e0405dc2d8bae949de6bb5e479bbf3cde401b6/crates/tauri/src/state.rs) has the same lifetime/outlives bounds, state-map lookup and missing-state error as the delegated calls. + +The pinned [command wrapper](https://github.com/tauri-apps/tauri/blob/f2e0405dc2d8bae949de6bb5e479bbf3cde401b6/crates/tauri-macros/src/command/wrapper.rs) builds arguments in signature order, defaults keys to camelCase, and calls CommandArg with the original plugin/name/message/ACL. Its async wrapper propagates the first argument error with ?, while its blocking wrapper reports the first match error. Grouping adjacent state arguments preserves this sequence inside the bundle. + +The pinned [mock IPC helpers](https://github.com/tauri-apps/tauri/blob/f2e0405dc2d8bae949de6bb5e479bbf3cde401b6/crates/tauri/src/test/mod.rs), InvokeRequest fields and InvokeBody::Json match the test fixture. The test module is enabled by the empty test feature; no package upgrade or additional dependency package is required. These are source compatibility checks, not compiled macro expansion or executed test evidence. diff --git a/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-uniffi-shutdown-result.md b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-uniffi-shutdown-result.md new file mode 100644 index 000000000..80a09b69e --- /dev/null +++ b/docs/plans/domain-architecture-and-multimodal/reports/2026-10-03-uniffi-shutdown-result.md @@ -0,0 +1,13 @@ +# Preserve fallible shutdown through UniFFI + +Exact PR42 head e7e5277 clears embedded-runtime Clippy and exposes one UniFFI unused Result at runtime.rs shutdown. EmbeddedRuntime::shutdown already awaits owned shutdown, preserves gateway residency on stop failure and only invalidates loaded session runtimes after successful stop. The adapter must propagate that result rather than report success or discard it. + +Return Result<(), FfiError> through the existing workflow_adapter_error internal_error JSON envelope, retaining the original owner Display cause. Twelve existing Rust shutdown calls now require success explicitly. No embedded lifecycle, invalidation, retry or backend policy is changed. A dedicated real adapter-chain regression uses the production FFI, EmbeddedRuntime and InferenceGateway with a failure-once backend: it checks the exact envelope, one stop invocation, retained actual gateway readiness, then successful retry and stopped readiness. A second test covers initial/repeated success. Non-shutdown trait methods panic if unexpectedly invoked. + +The test-only futures-util manifest/lock edge reuses the existing locked workspace version for the backend trait stream signature; no package version changes. The mock remains a backend fixture, not an alternate adapter implementation. + +Generated C# calls remain await runtime.Shutdown() (NativeSmoke has one executed call plus its compile-surface call; Quickstart has one executed call). Fallible metadata now enables the existing generated FfiException path for Rust FfiError; bindings and native library must be regenerated/shipped together, with no old-generated-binary compatibility claim. Headless CI regenerates and executes NativeSmoke and compiles/runs the packaged quickstart. A C# shutdown-failure injection is not available through the public runtime constructor without adding a production testing API or forcing an actual backend failure; the bounded failure test therefore exercises the real Rust adapter chain, while C# success/generated plumbing remains hosted qualification. + +Root source review accepted the complete eight-file checkpoint at frozen tree 18436e4cf375a4bb204f086453991fa07966eca0. Full pinned format passes locally, with no local Rust/native build claimed. PR42 has 455/455 embedded tests, full focused tests, workspace feature checks and format green; this distinct adapter milestone requires new exact-head aggregate and generated binding qualification. + +The pinned generator source confirms the C# error naming: CsCodeOracle::convert_error_suffix maps Error to Exception, so Rust FfiError is generated as FfiException ([official source](https://github.com/NordSecurity/uniffi-bindgen-cs/blob/2f4880f03ed08d960ad0ee14d11cf95444eee540/bindgen/src/gen_cs/mod.rs)). diff --git a/scripts/check-tauri-command-state-tests.sh b/scripts/check-tauri-command-state-tests.sh new file mode 100644 index 000000000..458965b57 --- /dev/null +++ b/scripts/check-tauri-command-state-tests.sh @@ -0,0 +1,21 @@ +#!/usr/bin/env bash +set -euo pipefail + +log_dir="${RUNNER_TEMP:-target}/tauri-command-state-tests" +mkdir -p "$log_dir" +cargo_args=(test -p pantograph --bin pantograph --no-default-features --features backend-llamacpp) +tests=( + workflow::command_state::tests::run_bundle_preserves_managed_identity_without_payload_state + workflow::command_state::tests::run_bundle_preserves_missing_state_order_and_original_camel_case_keys + workflow::command_state::tests::query_bundle_preserves_managed_identity_without_payload_state + workflow::command_state::tests::query_command_preserves_missing_state_order_and_original_keys + workflow::command_state::tests::query_command_preserves_flat_required_optional_and_context_payloads +) + +cargo "${cargo_args[@]}" workflow::command_state::tests:: -- --list > "$log_dir/tests.list" +for test_name in "${tests[@]}"; do + python3 -c 'import pathlib, sys; lines = pathlib.Path(sys.argv[1]).read_text().splitlines(); sys.exit(0 if sys.argv[2] + ": test" in lines else 1)' "$log_dir/tests.list" "$test_name" + result_path="$log_dir/${test_name##*::}.log" + cargo "${cargo_args[@]}" "$test_name" -- --exact 2>&1 | tee "$result_path" + python3 -c 'import pathlib, re, sys; text = pathlib.Path(sys.argv[1]).read_text(); sys.exit(0 if re.search(r"^test result: ok\. 1 passed; 0 failed; 0 ignored;", text, re.M) else 1)' "$result_path" +done diff --git a/src-tauri/Cargo.toml b/src-tauri/Cargo.toml index 929b6e50a..1107c59b9 100644 --- a/src-tauri/Cargo.toml +++ b/src-tauri/Cargo.toml @@ -96,4 +96,5 @@ axum = { workspace = true } tower-http = { workspace = true, optional = true } [dev-dependencies] +tauri = { version = "2.9.0", features = ["test"] } tempfile.workspace = true diff --git a/src-tauri/src/llm/commands/agent.rs b/src-tauri/src/llm/commands/agent.rs index 889879154..65c9605da 100644 --- a/src-tauri/src/llm/commands/agent.rs +++ b/src-tauri/src/llm/commands/agent.rs @@ -508,19 +508,6 @@ fn format_agent_prompt_with_analysis( prompt } -#[cfg(test)] -mod tests { - use super::drawing_vision_contract_only_error; - - #[test] - fn drawing_vision_error_points_to_typed_inference_contracts() { - let error = drawing_vision_contract_only_error(); - - assert!(error.contains("image_understanding")); - assert!(error.contains("canonical typed inference contracts")); - } -} - /// Format the agent prompt for fix/repair mode. /// Skips vision analysis, includes error context and file content directly. fn format_fix_mode_prompt(request: &AgentRequest) -> String { @@ -588,3 +575,16 @@ fn create_component_updates( }) .collect() } + +#[cfg(test)] +mod tests { + use super::drawing_vision_contract_only_error; + + #[test] + fn drawing_vision_error_points_to_typed_inference_contracts() { + let error = drawing_vision_contract_only_error(); + + assert!(error.contains("image_understanding")); + assert!(error.contains("canonical typed inference contracts")); + } +} diff --git a/src-tauri/src/workflow/command_state.rs b/src-tauri/src/workflow/command_state.rs new file mode 100644 index 000000000..689d99010 --- /dev/null +++ b/src-tauri/src/workflow/command_state.rs @@ -0,0 +1,66 @@ +//! Operation-specific bundles of Tauri-injected state; no client payload fields. + +use tauri::ipc::{CommandArg, CommandItem, InvokeError}; +use tauri::{Runtime, State}; + +use super::commands::{ + SharedExtensions, SharedNodeRegistry, SharedWorkflowDiagnosticsStore, SharedWorkflowService, +}; +use crate::agent::rag::SharedRagManager; +use crate::llm::{SharedGateway, SharedRuntimeRegistry}; + +pub struct WorkflowRunCommandState<'r> { + pub gateway: State<'r, SharedGateway>, + pub runtime_registry: State<'r, SharedRuntimeRegistry>, + pub extensions: State<'r, SharedExtensions>, + pub rag_manager: State<'r, SharedRagManager>, + pub workflow_service: State<'r, SharedWorkflowService>, + pub diagnostics_store: State<'r, SharedWorkflowDiagnosticsStore>, +} + +impl<'r, 'de: 'r, R: Runtime> CommandArg<'de, R> for WorkflowRunCommandState<'r> { + fn from_command(command: CommandItem<'de, R>) -> Result { + // Match the original command's extraction order and macro-generated keys. + Ok(Self { + gateway: managed_state(&command, "gateway")?, + runtime_registry: managed_state(&command, "runtimeRegistry")?, + extensions: managed_state(&command, "extensions")?, + rag_manager: managed_state(&command, "ragManager")?, + workflow_service: managed_state(&command, "workflowService")?, + diagnostics_store: managed_state(&command, "diagnosticsStore")?, + }) + } +} + +pub struct PortOptionsCommandState<'r> { + pub registry: State<'r, SharedNodeRegistry>, + pub extensions: State<'r, SharedExtensions>, + pub workflow_service: State<'r, SharedWorkflowService>, +} + +impl<'r, 'de: 'r, R: Runtime> CommandArg<'de, R> for PortOptionsCommandState<'r> { + fn from_command(command: CommandItem<'de, R>) -> Result { + Ok(Self { + registry: managed_state(&command, "registry")?, + extensions: managed_state(&command, "extensions")?, + workflow_service: managed_state(&command, "workflowService")?, + }) + } +} + +fn managed_state<'r, 'de: 'r, T: Send + Sync + 'static, R: Runtime>( + command: &CommandItem<'de, R>, + key: &'static str, +) -> Result, InvokeError> { + State::from_command(CommandItem { + plugin: command.plugin, + name: command.name, + key, + message: command.message, + acl: command.acl, + }) +} + +#[cfg(test)] +#[path = "command_state_tests.rs"] +mod tests; diff --git a/src-tauri/src/workflow/command_state_tests.rs b/src-tauri/src/workflow/command_state_tests.rs new file mode 100644 index 000000000..905774798 --- /dev/null +++ b/src-tauri/src/workflow/command_state_tests.rs @@ -0,0 +1,309 @@ +use std::path::PathBuf; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; + +use async_trait::async_trait; +use serde_json::{json, Value}; +use tauri::test::{mock_builder, mock_context, noop_assets, MockRuntime}; +use tauri::{App, Builder, WebviewWindow, WebviewWindowBuilder}; +use tokio::sync::{mpsc, RwLock}; + +use super::*; +use crate::workflow::commands; + +struct NoSpawn; + +#[async_trait] +impl inference::ProcessSpawner for NoSpawn { + async fn spawn_sidecar( + &self, + _: &str, + _: &[&str], + ) -> Result< + ( + mpsc::Receiver, + Box, + ), + String, + > { + panic!("state extraction must not spawn a backend") + } + + fn app_data_dir(&self) -> Result { + panic!("state extraction must not resolve runtime paths") + } + + fn binaries_dir(&self) -> Result { + panic!("state extraction must not resolve runtime paths") + } +} + +struct ManagedStates { + gateway: SharedGateway, + runtime_registry: SharedRuntimeRegistry, + extensions: SharedExtensions, + rag_manager: SharedRagManager, + workflow_service: SharedWorkflowService, + diagnostics_store: SharedWorkflowDiagnosticsStore, + registry: SharedNodeRegistry, +} + +impl ManagedStates { + fn new(registry: SharedNodeRegistry) -> Self { + let mut extensions = node_engine::ExecutorExtensions::new(); + extensions.set("state-fixture-marker", "managed extension".to_string()); + Self { + gateway: Arc::new(crate::llm::InferenceGateway::new(Arc::new(NoSpawn))), + runtime_registry: Arc::new(pantograph_runtime_registry::RuntimeRegistry::new()), + extensions: Arc::new(RwLock::new(extensions)), + rag_manager: crate::agent::rag::create_rag_manager(PathBuf::from( + "unused-state-fixture", + )), + workflow_service: Arc::new(pantograph_workflow_service::WorkflowService::new()), + diagnostics_store: Arc::new(crate::workflow::WorkflowDiagnosticsStore::default()), + registry, + } + } + + fn run_builder(&self, prefix: usize) -> Builder { + let mut builder = mock_builder(); + if prefix > 0 { + builder = builder.manage(self.gateway.clone()); + } + if prefix > 1 { + builder = builder.manage(self.runtime_registry.clone()); + } + if prefix > 2 { + builder = builder.manage(self.extensions.clone()); + } + if prefix > 3 { + builder = builder.manage(self.rag_manager.clone()); + } + if prefix > 4 { + builder = builder.manage(self.workflow_service.clone()); + } + if prefix > 5 { + builder = builder.manage(self.diagnostics_store.clone()); + } + builder + } + + fn query_builder(&self, prefix: usize) -> Builder { + let mut builder = mock_builder(); + if prefix > 0 { + builder = builder.manage(self.registry.clone()); + } + if prefix > 1 { + builder = builder.manage(self.extensions.clone()); + } + if prefix > 2 { + builder = builder.manage(self.workflow_service.clone()); + } + builder + } +} + +fn identity(value: &Arc) -> String { + format!("{:p}", Arc::as_ptr(value)) +} + +#[tauri::command] +fn inspect_run_state(state: WorkflowRunCommandState<'_>) -> Vec { + vec![ + identity(state.gateway.inner()), + identity(state.runtime_registry.inner()), + identity(state.extensions.inner()), + identity(state.rag_manager.inner()), + identity(state.workflow_service.inner()), + identity(state.diagnostics_store.inner()), + ] +} + +#[tauri::command] +fn inspect_query_state(state: PortOptionsCommandState<'_>) -> Vec { + vec![ + identity(state.registry.inner()), + identity(state.extensions.inner()), + identity(state.workflow_service.inner()), + ] +} + +fn build(builder: Builder) -> (App, WebviewWindow) { + let app = builder + .invoke_handler(tauri::generate_handler![ + inspect_run_state, + inspect_query_state, + commands::query_port_options, + ]) + .build(mock_context(noop_assets())) + .expect("mock app"); + let webview = WebviewWindowBuilder::new(&app, "main", Default::default()) + .build() + .expect("mock webview"); + (app, webview) +} + +fn invoke( + webview: &WebviewWindow, + command: &str, + body: Value, +) -> Result { + tauri::test::get_ipc_response( + webview, + tauri::webview::InvokeRequest { + cmd: command.to_string(), + callback: tauri::ipc::CallbackFn(0), + error: tauri::ipc::CallbackFn(1), + url: if cfg!(any(windows, target_os = "android")) { + "http://tauri.localhost" + } else { + "tauri://localhost" + } + .parse() + .expect("mock IPC URL"), + body: tauri::ipc::InvokeBody::Json(body), + headers: Default::default(), + invoke_key: tauri::test::INVOKE_KEY.to_string(), + }, + ) + .map(|body| body.deserialize().expect("JSON response")) +} + +fn missing_state(command: &str, key: &str) -> Value { + json!(format!("state not managed for field `{key}` on command `{command}`. You must call `.manage()` before using this command")) +} + +#[test] +fn run_bundle_preserves_managed_identity_without_payload_state() { + let states = ManagedStates::new(Arc::new(node_engine::NodeRegistry::new())); + let (_app, webview) = build(states.run_builder(6)); + let actual = invoke( + &webview, + "inspect_run_state", + json!({"state": "ignored client input"}), + ) + .expect("managed state"); + assert_eq!( + actual, + json!([ + identity(&states.gateway), + identity(&states.runtime_registry), + identity(&states.extensions), + identity(&states.rag_manager), + identity(&states.workflow_service), + identity(&states.diagnostics_store), + ]) + ); +} + +#[test] +fn run_bundle_preserves_missing_state_order_and_original_camel_case_keys() { + let states = ManagedStates::new(Arc::new(node_engine::NodeRegistry::new())); + for (prefix, key) in [ + "gateway", + "runtimeRegistry", + "extensions", + "ragManager", + "workflowService", + "diagnosticsStore", + ] + .iter() + .enumerate() + { + let (_app, webview) = build(states.run_builder(prefix)); + assert_eq!( + invoke(&webview, "inspect_run_state", json!({})), + Err(missing_state("inspect_run_state", key)) + ); + } +} + +#[test] +fn query_bundle_preserves_managed_identity_without_payload_state() { + let states = ManagedStates::new(Arc::new(node_engine::NodeRegistry::new())); + let (_app, webview) = build(states.query_builder(3)); + let actual = invoke(&webview, "inspect_query_state", json!({})).expect("managed state"); + assert_eq!( + actual, + json!([ + identity(&states.registry), + identity(&states.extensions), + identity(&states.workflow_service) + ]) + ); +} + +#[test] +fn query_command_preserves_missing_state_order_and_original_keys() { + let states = ManagedStates::new(Arc::new(node_engine::NodeRegistry::new())); + for (prefix, key) in ["registry", "extensions", "workflowService"] + .iter() + .enumerate() + { + let (_app, webview) = build(states.query_builder(prefix)); + assert_eq!( + invoke( + &webview, + "query_port_options", + json!({"nodeType":"fixture", "portId":"choice"}) + ), + Err(missing_state("query_port_options", key)) + ); + } +} + +struct EchoOptionsProvider(Arc); + +#[async_trait] +impl node_engine::PortOptionsProvider for EchoOptionsProvider { + async fn query_options( + &self, + query: &node_engine::PortOptionsQuery, + extensions: &node_engine::ExecutorExtensions, + ) -> node_engine::Result { + self.0.fetch_add(1, Ordering::SeqCst); + Ok(node_engine::PortOptionsResult { + options: Vec::new(), + total_count: 0, + searchable: true, + metadata: Some( + json!({"query":query,"marker":extensions.get::("state-fixture-marker")}), + ), + }) + } +} + +#[test] +fn query_command_preserves_flat_required_optional_and_context_payloads() { + let calls = Arc::new(AtomicUsize::new(0)); + let mut registry = node_engine::NodeRegistry::new(); + registry.register_port_provider( + "fixture", + "choice", + Box::new(EchoOptionsProvider(calls.clone())), + ); + let states = ManagedStates::new(Arc::new(registry)); + let (_app, webview) = build(states.query_builder(3)); + let context = json!({"targetNodeId":"target-1", "taskKind":"embedding"}); + let actual = invoke(&webview, "query_port_options", json!({ + "nodeType":"fixture", "portId":"choice", "search":"needle", "limit":7, "offset":3, + "context":context, "state":"must not deserialize", "registry":"must not replace managed state", + })).expect("original flat command payload"); + assert_eq!( + actual, + json!({"options":[],"totalCount":0,"searchable":true,"metadata":{ + "query":{"search":"needle","limit":7,"offset":3,"context":context},"marker":"managed extension", + }}) + ); + let minimal = invoke( + &webview, + "query_port_options", + json!({"nodeType":"fixture", "portId":"choice"}), + ) + .expect("original optional defaults"); + assert_eq!( + minimal["metadata"]["query"], + json!({"search":null,"limit":null,"offset":null}) + ); + assert_eq!(calls.load(Ordering::SeqCst), 2); +} diff --git a/src-tauri/src/workflow/commands.rs b/src-tauri/src/workflow/commands.rs index 1c23b739f..6b2b64548 100644 --- a/src-tauri/src/workflow/commands.rs +++ b/src-tauri/src/workflow/commands.rs @@ -6,6 +6,8 @@ use std::{path::PathBuf, sync::Arc}; use tauri::{command, AppHandle, Manager, State}; + +pub use super::command_state::{PortOptionsCommandState, WorkflowRunCommandState}; use tokio::sync::RwLock; use crate::agent::rag::SharedRagManager; @@ -199,26 +201,11 @@ pub async fn workflow_create_execution_session( pub async fn workflow_run_execution_session( request: pantograph_workflow_service::WorkflowExecutionSessionRunRequest, app: AppHandle, - gateway: State<'_, SharedGateway>, - runtime_registry: State<'_, SharedRuntimeRegistry>, - extensions: State<'_, SharedExtensions>, - rag_manager: State<'_, SharedRagManager>, - workflow_service: State<'_, SharedWorkflowService>, - diagnostics_store: State<'_, SharedWorkflowDiagnosticsStore>, + state: WorkflowRunCommandState<'_>, channel: tauri::ipc::Channel, ) -> Result { - super::headless_workflow_commands::workflow_run_execution_session( - request, - app, - gateway, - runtime_registry, - extensions, - rag_manager, - workflow_service, - diagnostics_store, - channel, - ) - .await + super::headless_workflow_commands::workflow_run_execution_session(request, app, state, channel) + .await } #[command] @@ -769,9 +756,7 @@ pub async fn workflow_set_execution_session_keep_alive( #[command] pub async fn query_port_options( - registry: State<'_, SharedNodeRegistry>, - extensions: State<'_, SharedExtensions>, - workflow_service: State<'_, SharedWorkflowService>, + state: PortOptionsCommandState<'_>, node_type: String, port_id: String, search: Option, @@ -780,15 +765,7 @@ pub async fn query_port_options( context: Option, ) -> Result { super::workflow_port_query_commands::query_port_options( - registry, - extensions, - workflow_service, - node_type, - port_id, - search, - limit, - offset, - context, + state, node_type, port_id, search, limit, offset, context, ) .await } diff --git a/src-tauri/src/workflow/diagnostics/overlay.rs b/src-tauri/src/workflow/diagnostics/overlay.rs index 369eeee1c..f576c236d 100644 --- a/src-tauri/src/workflow/diagnostics/overlay.rs +++ b/src-tauri/src/workflow/diagnostics/overlay.rs @@ -434,7 +434,7 @@ fn is_inline_media_string_key(key: &str, value: &str) -> bool { fn is_probably_base64_body(value: &str) -> bool { let value = value.trim(); value.len() >= 128 - && value.len() % 4 == 0 + && value.len().is_multiple_of(4) && value .bytes() .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'+' | b'/' | b'=')) diff --git a/src-tauri/src/workflow/event_adapter/tests/translation_projection.rs b/src-tauri/src/workflow/event_adapter/tests/translation_projection.rs index 719fe35ea..5c7403afe 100644 --- a/src-tauri/src/workflow/event_adapter/tests/translation_projection.rs +++ b/src-tauri/src/workflow/event_adapter/tests/translation_projection.rs @@ -143,7 +143,7 @@ fn translated_task_inputs_resolved_event_preserves_inputs_without_trace_noise() TauriWorkflowEvent::DiagnosticsSnapshot { snapshot, .. } => { let trace = snapshot.runs_by_id.get("exec-1").expect("trace"); assert_eq!(trace.event_count, 1); - assert!(trace.nodes.get("node-a").is_none()); + assert!(!trace.nodes.contains_key("node-a")); } other => panic!("unexpected diagnostics event: {other:?}"), } diff --git a/src-tauri/src/workflow/headless_workflow_commands.rs b/src-tauri/src/workflow/headless_workflow_commands.rs index 88b439c11..ef8e2ab8b 100644 --- a/src-tauri/src/workflow/headless_workflow_commands.rs +++ b/src-tauri/src/workflow/headless_workflow_commands.rs @@ -59,7 +59,8 @@ use tauri::{ipc::Channel, AppHandle, State}; use crate::agent::rag::SharedRagManager; use crate::llm::{SharedGateway, SharedRuntimeRegistry}; -use super::commands::{SharedExtensions, SharedWorkflowDiagnosticsStore, SharedWorkflowService}; +use super::command_state::WorkflowRunCommandState; +use super::commands::{SharedExtensions, SharedWorkflowService}; use super::events::WorkflowEvent; use super::headless_diagnostics::workflow_scheduler_snapshot_response; pub(crate) use super::headless_runtime::build_runtime; @@ -188,14 +189,17 @@ pub async fn workflow_create_execution_session( pub async fn workflow_run_execution_session( request: WorkflowExecutionSessionRunRequest, app: AppHandle, - gateway: State<'_, SharedGateway>, - runtime_registry: State<'_, SharedRuntimeRegistry>, - extensions: State<'_, SharedExtensions>, - rag_manager: State<'_, SharedRagManager>, - workflow_service: State<'_, SharedWorkflowService>, - diagnostics_store: State<'_, SharedWorkflowDiagnosticsStore>, + state: WorkflowRunCommandState<'_>, channel: Channel, ) -> Result { + let WorkflowRunCommandState { + gateway, + runtime_registry, + extensions, + rag_manager, + workflow_service, + diagnostics_store, + } = state; let runtime = build_runtime( &app, gateway.inner(), diff --git a/src-tauri/src/workflow/mod.rs b/src-tauri/src/workflow/mod.rs index a338767b3..e6486440e 100644 --- a/src-tauri/src/workflow/mod.rs +++ b/src-tauri/src/workflow/mod.rs @@ -21,6 +21,7 @@ //! └─────────────────────────────────┘ //! ``` +mod command_state; pub mod commands; pub mod diagnostics; pub mod event_adapter; diff --git a/src-tauri/src/workflow/puma_lib_commands.rs b/src-tauri/src/workflow/puma_lib_commands.rs index 3b670f181..ca8498756 100644 --- a/src-tauri/src/workflow/puma_lib_commands.rs +++ b/src-tauri/src/workflow/puma_lib_commands.rs @@ -160,7 +160,7 @@ pub async fn resolve_model_package_facts_summary( let ext = extensions.read().await; pumas_update_feed_access_from_extensions(&ext) }; - resolve_model_package_facts_summary_from_access(selector_access, &model_id).await + resolve_model_package_facts_summary_from_access(selector_access, model_id).await } pub async fn list_model_library_updates_since( diff --git a/src-tauri/src/workflow/workflow_port_query_commands.rs b/src-tauri/src/workflow/workflow_port_query_commands.rs index 269704fed..af16377ba 100644 --- a/src-tauri/src/workflow/workflow_port_query_commands.rs +++ b/src-tauri/src/workflow/workflow_port_query_commands.rs @@ -1,14 +1,13 @@ +use super::command_state::PortOptionsCommandState; use tauri::State; -use super::commands::{SharedExtensions, SharedNodeRegistry, SharedWorkflowService}; +use super::commands::{SharedNodeRegistry, SharedWorkflowService}; const PUMA_LIB_NODE_TYPE: &str = "puma-lib"; const PUMAS_MODEL_REF_PORT_ID: &str = "pumas_model_ref"; pub async fn query_port_options( - registry: State<'_, SharedNodeRegistry>, - extensions: State<'_, SharedExtensions>, - workflow_service: State<'_, SharedWorkflowService>, + state: PortOptionsCommandState<'_>, node_type: String, port_id: String, search: Option, @@ -16,6 +15,11 @@ pub async fn query_port_options( offset: Option, context: Option, ) -> Result { + let PortOptionsCommandState { + registry, + extensions, + workflow_service, + } = state; let ext = extensions.read().await; let query = node_engine::PortOptionsQuery { search: search.clone(),