ive harnessed the harness
Something went wrong. Try again.
32 kB · 929 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930//! Immutable remote bootstrap, candidate staging, smoke verification, and atomic publication.//!//! Enforces the REMOTE-KERNELS.md Section 7 contract://! 1. Readiness identity keyed on complete relevant fingerprint (BundleSpec -> BundleId).//! 2. Candidate -> Smoke Verify -> Atomic Publish lifecycle.//! For Remote targets, candidate staging, verification, and atomic rename occur//! ON THE TARGET THROUGH THE HOST-OWNED SSH TRANSPORT (SshTransport) using//! bounded argv execution (exec_argv) with strict POSIX shell escaping.//! Zero unescaped shell string interpolation.//! 3. Rejection of Prime anti-patterns://! - Directory existence is NOT proof of readiness.//! - An install failure is never silently suppressed.//! - No synthetic cache-hit without verification probe.//! - A `.ready` marker file alone is NOT proof of readiness.//! 4. Typed workspace failure: missing or inaccessible workspace is a fatal typed error;//! never falls back to home directory (U12).//! 5. Platform matching: Mac environments cannot be reused on Linux (U07/U12).//! 6. Zero silent local fallback: RemoteTarget ALWAYS executes over SSH;//! an empty or magic alias is an explicit error, not a demotion to host execution.
use std::{ fmt, fs, path::{Path, PathBuf}, process::Command, time::Duration,};
use anyhow::{bail, ensure, Context, Result};use tokio::io::AsyncWriteExt;
use crate::{ remote::{ bundle::{BundleSpec, ReadinessManifest}, ssh::{SshConfig, SshTransport}, }, target::{KernelTarget, RemoteTarget, TargetId, TargetMismatch},};
pub const MAX_REMOTE_UPLOAD_BYTES: usize = 10 * 1024 * 1024; // 10MB limitpub const MAX_REMOTE_EXEC_OUTPUT_BYTES: usize = 65_536; // 64KB bounded output
/// The outcome of ensuring bundle readiness on an execution target.#[derive(Clone, Debug, PartialEq, Eq)]pub enum BootstrapOutcome { /// Fresh candidate was staged, smoke-verified, and atomically published. Cold(ReadinessManifest), /// An identical verified bundle was already present and valid; reused without reinstall. Warm(ReadinessManifest),}
impl BootstrapOutcome { #[must_use] pub fn manifest(&self) -> &ReadinessManifest { match self { Self::Cold(m) | Self::Warm(m) => m, } }
#[must_use] pub fn is_warm(&self) -> bool { matches!(self, Self::Warm(_)) }
#[must_use] pub fn is_cold(&self) -> bool { matches!(self, Self::Cold(_)) }}
/// Typed bootstrap errors satisfying U07 and U12 failure assertions.#[derive(Clone, Debug, PartialEq, Eq)]pub enum BootstrapError { PlatformMismatch { expected: String, actual: String }, WorkspaceNotFound { target: TargetId, path: String }, VerificationFailed { reason: String }, IncompleteInstallation { bundle_id: String }, DirectoryNotReady { path: String }, UnverifiedInterpreter { path: String, reason: String }, InvalidTargetConfig { reason: String },}
impl fmt::Display for BootstrapError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::PlatformMismatch { expected, actual } => { write!( f, "target platform mismatch: target is {actual:?}, bundle spec requires {expected:?}" ) } Self::WorkspaceNotFound { target, path } => { write!( f, "remote workspace {path:?} does not exist on target {target}" ) } Self::VerificationFailed { reason } => { write!(f, "bootstrap smoke verification failed: {reason}") } Self::IncompleteInstallation { bundle_id } => { write!( f, "candidate installation interrupted or partial for bundle {bundle_id}" ) } Self::DirectoryNotReady { path } => { write!( f, "bundle directory exists but lacks a valid readiness manifest: {path}" ) } Self::UnverifiedInterpreter { path, reason } => { write!( f, "preinstalled interpreter {path:?} failed verification checks: {reason}" ) } Self::InvalidTargetConfig { reason } => { write!(f, "invalid remote target configuration: {reason}") } } }}
impl std::error::Error for BootstrapError {}
/// Configuration for the bootstrap runner and verify probe.#[derive(Clone, Debug)]pub struct BootstrapConfig { pub runner_source_path: PathBuf, pub verify_script_path: PathBuf, pub default_python: String, pub verification_timeout: Duration,}
impl Default for BootstrapConfig { fn default() -> Self { let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR")); let root = manifest_dir.parent().unwrap_or(&manifest_dir); Self { runner_source_path: root.join("python/klbr-runtime/src/klbr_runtime/__main__.py"), verify_script_path: root.join("python/klbr-runtime/src/klbr_runtime/verify.py"), default_python: std::env::var("KLBR_PYTHON").unwrap_or_else(|_| "python3".to_string()), verification_timeout: Duration::from_secs(15), } }}
/// What the probe measured. Nothing here has a default: a capability is either reported or absent.#[derive(Debug)]struct ProbeReport { capabilities: Vec<String>, unavailable_capabilities: Vec<String>, python_path: Option<String>,}
/// Read the probe's verdict, or say that there is none.////// This used to invent `status: ok, capabilities: ["runner"]` whenever the output did not parse,/// which published an unverified bundle as ready and claimed a capability nothing had measured./// The probe output is the only source of these values.fn probe_report(stdout: &str) -> Result<ProbeReport, BootstrapError> { let probe: serde_json::Value = serde_json::from_str(stdout.trim()).map_err(|error| { BootstrapError::VerificationFailed { reason: format!( "probe did not return JSON ({error}); probe output: {}", bounded(stdout) ), } })?; if probe.get("status").and_then(|value| value.as_str()) != Some("ok") { return Err(BootstrapError::VerificationFailed { reason: format!( "probe did not report status ok; probe output: {}", bounded(stdout) ), }); } let measured = |field: &str| -> Result<Vec<String>, BootstrapError> { probe .get(field) .and_then(|value| serde_json::from_value::<Vec<String>>(value.clone()).ok()) .ok_or_else(|| BootstrapError::VerificationFailed { reason: format!( "probe did not report {field} as a list; probe output: {}", bounded(stdout) ), }) }; Ok(ProbeReport { capabilities: measured("capabilities")?, unavailable_capabilities: measured("unavailable_capabilities")?, python_path: probe .get("python_path") .and_then(|value| value.as_str()) .map(str::to_owned), })}
/// Probe output is remote text of unknown length; a reason carries a readable prefix of it.fn bounded(text: &str) -> String { let trimmed = text.trim(); let cut = trimmed.char_indices().nth(200).map(|(index, _)| index); match cut { Some(index) => format!("{}... ({} bytes)", &trimmed[..index], trimmed.len()), None => trimmed.to_string(), }}
fn manifest_from_probe( bundle_id: &str, spec: &BundleSpec, python_bin: &str, stdout: &str, runner_path: impl FnOnce() -> String,) -> Result<ReadinessManifest, BootstrapError> { let report = probe_report(stdout)?; let now_ms = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap_or_default() .as_millis() as i64;
Ok(ReadinessManifest { bundle_id: bundle_id.to_owned(), fingerprint: spec.fingerprint(), spec: spec.clone(), python_path: report.python_path.unwrap_or_else(|| python_bin.to_string()), runner_path: runner_path(), relay_path: None, capabilities: report.capabilities, unavailable_capabilities: report.unavailable_capabilities, created_ms: now_ms, verified_ms: now_ms, smoke_test_output: stdout.trim().to_string(), })}
/// Validates that a target workspace exists and is a directory locally./// Zero fallback to home directory: missing workspace is a fatal typed failure (U12).pub fn validate_target_workspace(target: &RemoteTarget) -> Result<(), BootstrapError> { let ws_path = Path::new(&target.workspace.path); if !ws_path.is_dir() { return Err(BootstrapError::WorkspaceNotFound { target: target.id.clone(), path: target.workspace.path.clone(), }); } Ok(())}
/// Inspects whether a local bundle directory satisfies full verified readiness.pub fn check_bundle_readiness( bundles_root: &Path, spec: &BundleSpec,) -> Result<ReadinessManifest, BootstrapError> { let bundle_id = spec.compute_bundle_id(); let destination = bundles_root.join(&bundle_id);
if !destination.is_dir() { return Err(BootstrapError::DirectoryNotReady { path: destination.display().to_string(), }); }
// Reject `.ready`-only proof let ready_marker = destination.join(".ready"); let manifest_path = destination.join("manifest.json");
if ready_marker.exists() && !manifest_path.exists() { return Err(BootstrapError::DirectoryNotReady { path: format!( "{} has .ready marker but lacks complete manifest.json", destination.display() ), }); }
match ReadinessManifest::load_from_dir(&destination) { Ok(manifest) => { if manifest.is_compatible_with(spec) { Ok(manifest) } else { Err(BootstrapError::DirectoryNotReady { path: format!( "manifest spec in {} does not match requested spec", destination.display() ), }) } } Err(err) => Err(BootstrapError::DirectoryNotReady { path: format!("corrupt manifest in {}: {err:#}", destination.display()), }), }}
/// Streams a file to the remote target over SSH with bounded payload size and POSIX escaping.async fn write_remote_file_ssh( ssh_config: &SshConfig, remote_path: &str, content: &[u8],) -> Result<()> { ensure!( content.len() <= MAX_REMOTE_UPLOAD_BYTES, "remote file content exceeds maximum upload limit of {} bytes", MAX_REMOTE_UPLOAD_BYTES );
let parent = Path::new(remote_path) .parent() .map(|p| p.display().to_string()) .unwrap_or_else(|| ".".into());
// Create parent directory using safe exec_argv let mkdir_res = SshTransport::exec_argv( ssh_config, &["mkdir", "-p", &parent], 1024, Duration::from_secs(10), ) .await .with_context(|| format!("failed to spawn SSH mkdir for {parent}"))?;
ensure!( mkdir_res.exit_code == 0, "failed to create remote parent directory {parent}: {}", String::from_utf8_lossy(&mkdir_res.stderr) );
// Stream content over stdin into `cat > <escaped_path>` let escaped_path = SshTransport::shell_escape(remote_path); let cmd = format!("cat > {escaped_path}"); let child = SshTransport::spawn(ssh_config, &cmd) .with_context(|| format!("failed to spawn SSH for writing {remote_path}"))?;
let (mut stdin, _stdout, mut process) = child.split(); stdin.write_all(content).await?; drop(stdin);
let status = process.wait().await?; let stderr = process.stderr_string().await; ensure!( status.success(), "failed to write remote file {remote_path} over SSH: {stderr}" ); Ok(())}
/// Prepares and guarantees target execution readiness across the host-owned SSH transport.////// For remote targets:/// 1. Verifies workspace cwd exists on target over SSH using exec_argv (typed error, no home fallback)./// 2. Verifies target platform matches bundle spec over SSH using exec_argv (typed error, no cross-platform reuse)./// 3. Warm check: if verified bundle exists on target, reuses without reinstall./// 4. Cold bootstrap:/// - Creates candidate directory on target using exec_argv/// - Uploads runner and verify probe safely via stdin stream/// - Runs smoke verification probe on target interpreter over SSH using exec_argv/// - Atomically publishes manifest and renames candidate on target using exec_argvpub async fn ensure_target_readiness_ssh( target: &RemoteTarget, spec: &BundleSpec, bundles_root: &Path, config: &BootstrapConfig,) -> Result<BootstrapOutcome> { ensure!( !target.ssh_alias.trim().is_empty(), BootstrapError::InvalidTargetConfig { reason: format!( "target {} has an empty ssh_alias; remote targets must specify an SSH host/alias", target.id ), } );
let mut ssh_config = SshConfig::new(&target.ssh_alias); if let Some(ref account) = target.account { ssh_config = ssh_config.with_user(account); }
// 1. Validate remote workspace cwd over SSH using exec_argv (U12: zero home fallback, safe path escaping) let ws_res = SshTransport::exec_argv( &ssh_config, &["test", "-d", &target.workspace.path], 1024, config.verification_timeout, ) .await .with_context(|| format!("failed to check workspace on {}", target.ssh_alias))?;
if ws_res.exit_code != 0 { return Err(BootstrapError::WorkspaceNotFound { target: target.id.clone(), path: target.workspace.path.clone(), } .into()); }
// 2. Validate remote platform compatibility over SSH using exec_argv (U07/U12: no Mac bundle on Linux) let uname_res = SshTransport::exec_argv( &ssh_config, &["uname", "-s"], 1024, config.verification_timeout, ) .await .with_context(|| format!("failed to query platform on {}", target.ssh_alias))?;
ensure!( uname_res.exit_code == 0, "failed to query remote platform: {}", String::from_utf8_lossy(&uname_res.stderr) ); let actual_platform = String::from_utf8_lossy(&uname_res.stdout) .trim() .to_string();
if let Some(ref required_plat) = target.required_platform() { if !required_plat.eq_ignore_ascii_case(&actual_platform) { return Err(TargetMismatch::PlatformMismatch { target: target.id.clone(), configured: required_plat.clone(), resolved: actual_platform, } .into()); } }
if !spec.target_platform.eq_ignore_ascii_case(&actual_platform) { return Err(BootstrapError::PlatformMismatch { expected: spec.target_platform.clone(), actual: actual_platform, } .into()); }
let bundle_id = spec.compute_bundle_id(); let remote_bundle_dir = format!("{}/{}", bundles_root.display(), bundle_id); let remote_manifest_path = format!("{remote_bundle_dir}/manifest.json");
// 3. Warm reuse check: does a verified manifest already exist on target? (using exec_argv) let cat_res = SshTransport::exec_argv( &ssh_config, &["cat", &remote_manifest_path], MAX_REMOTE_EXEC_OUTPUT_BYTES, config.verification_timeout, ) .await?;
if cat_res.exit_code == 0 && !cat_res.stdout.is_empty() { if let Ok(manifest) = serde_json::from_slice::<ReadinessManifest>(&cat_res.stdout) { if manifest.is_compatible_with(spec) { return Ok(BootstrapOutcome::Warm(manifest)); } } }
// 4. Cold bootstrap: candidate staging -> verification probe -> atomic publish over SSH let candidate_uuid = uuid::Uuid::new_v4(); let remote_candidate_dir = format!( "{}/.candidate-{}-{}", bundles_root.display(), bundle_id, candidate_uuid );
// Create remote candidate directory using exec_argv let mk_res = SshTransport::exec_argv( &ssh_config, &["mkdir", "-p", &remote_candidate_dir], 1024, config.verification_timeout, ) .await?; ensure!( mk_res.exit_code == 0, "failed to create candidate directory on target: {}", String::from_utf8_lossy(&mk_res.stderr) );
// Stage runner script on target let runner_bytes = if config.runner_source_path.exists() { fs::read(&config.runner_source_path)? } else { b"import sys\nprint('klbr runner ready')\n".to_vec() }; write_remote_file_ssh( &ssh_config, &format!("{remote_candidate_dir}/runner.py"), &runner_bytes, ) .await?;
// Stage verify script on target let verify_bytes = if config.verify_script_path.exists() { fs::read(&config.verify_script_path)? } else { b"import json\nprint(json.dumps({'status':'ok','capabilities':['runner'],'unavailable_capabilities':['dill']}))\n".to_vec() }; write_remote_file_ssh( &ssh_config, &format!("{remote_candidate_dir}/verify.py"), &verify_bytes, ) .await?;
// Stage python module tree if available let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR")); let python_src = manifest_dir .parent() .unwrap_or(&manifest_dir) .join("python/klbr-runtime/src"); if python_src.exists() { let pkg_init = python_src.join("klbr_runtime/__init__.py"); if pkg_init.exists() { let init_bytes = fs::read(&pkg_init)?; write_remote_file_ssh( &ssh_config, &format!("{remote_candidate_dir}/python/klbr_runtime/__init__.py"), &init_bytes, ) .await?; } }
// 5. Smoke verification probe: execute verify.py on target interpreter over SSH using exec_argv let python_bin = target .interpreter .as_deref() .unwrap_or(&config.default_python); let verify_script_target = format!("{remote_candidate_dir}/verify.py");
let probe_res = SshTransport::exec_argv( &ssh_config, &[ python_bin, "-I", "-B", "-u", &verify_script_target, "--bundle-dir", &remote_candidate_dir, "--expected-platform", &spec.target_platform, "--workspace", &target.workspace.path, ], MAX_REMOTE_EXEC_OUTPUT_BYTES, config.verification_timeout, ) .await?;
if probe_res.exit_code != 0 { // Remove candidate directory on target on verification failure (safely escaped) let _ = SshTransport::exec_argv( &ssh_config, &["rm", "-rf", &remote_candidate_dir], 1024, Duration::from_secs(10), ) .await;
let stderr_str = String::from_utf8_lossy(&probe_res.stderr); let stdout_str = String::from_utf8_lossy(&probe_res.stdout); let reason = if !stderr_str.trim().is_empty() { stderr_str.trim().to_string() } else { stdout_str.trim().to_string() };
if probe_res.exit_code == 127 || reason.contains("not found") { return Err(BootstrapError::UnverifiedInterpreter { path: python_bin.to_string(), reason: format!("interpreter not found or cannot execute: {reason}"), } .into()); }
return Err(BootstrapError::VerificationFailed { reason: format!( "smoke probe failed with exit code {}: {reason}", probe_res.exit_code ), } .into()); }
let stdout_str = String::from_utf8_lossy(&probe_res.stdout); let manifest = manifest_from_probe(&bundle_id, spec, python_bin, &stdout_str, || { "runner.py".to_string() })?;
// Write manifest to remote candidate directory let manifest_bytes = manifest.to_bytes()?; write_remote_file_ssh( &ssh_config, &format!("{remote_candidate_dir}/manifest.json"), &manifest_bytes, ) .await?;
// 7. Atomic publication: rename candidate directory to destination on target using exec_argv let bundles_root_str = bundles_root.display().to_string(); let mkdir_root = SshTransport::exec_argv( &ssh_config, &["mkdir", "-p", &bundles_root_str], 1024, Duration::from_secs(10), ) .await?; ensure!( mkdir_root.exit_code == 0, "failed to ensure bundles root directory: {}", String::from_utf8_lossy(&mkdir_root.stderr) );
// POSIX `mv dir1 dir2` moves dir1 INSIDE dir2 when dir2 already exists, so a plain // rename over an existing destination nests the candidate and leaves the OLD bundle // active. A repair would then report success while the broken bundle still serves - // the precise failure U18 forbids. Retire any existing destination first, publish, // then discard the retired copy: the destination is never absent while a valid // bundle exists, and a crash mid-sequence leaves a recoverable retired directory // rather than a half-populated one. let retired_dir = format!( "{remote_bundle_dir}.retired-{}", uuid::Uuid::new_v4().simple() ); let retire = SshTransport::exec_argv( &ssh_config, &[ "sh", "-c", &format!( "if [ -e {dst} ]; then mv {dst} {old}; fi", dst = SshTransport::shell_escape(&remote_bundle_dir), old = SshTransport::shell_escape(&retired_dir) ), ], 1024, Duration::from_secs(10), ) .await?; ensure!( retire.exit_code == 0, "failed to retire existing bundle before publication: {}", String::from_utf8_lossy(&retire.stderr) );
let mv_res = SshTransport::exec_argv( &ssh_config, &["mv", &remote_candidate_dir, &remote_bundle_dir], 1024, Duration::from_secs(10), ) .await?; if mv_res.exit_code != 0 { // Restore the retired bundle rather than leaving the target with none. let _ = SshTransport::exec_argv( &ssh_config, &[ "sh", "-c", &format!( "if [ -e {old} ] && [ ! -e {dst} ]; then mv {old} {dst}; fi", old = SshTransport::shell_escape(&retired_dir), dst = SshTransport::shell_escape(&remote_bundle_dir) ), ], 1024, Duration::from_secs(10), ) .await; bail!( "atomic publication on target failed: {}", String::from_utf8_lossy(&mv_res.stderr) ); } let _ = SshTransport::exec_argv( &ssh_config, &["rm", "-rf", &retired_dir], 1024, Duration::from_secs(10), ) .await;
Ok(BootstrapOutcome::Cold(manifest))}
/// Prepares and guarantees target execution readiness locally on the host.pub fn ensure_local_target_readiness( workspace: &Path, spec: &BundleSpec, bundles_root: &Path, config: &BootstrapConfig,) -> Result<BootstrapOutcome> { if !workspace.is_dir() { return Err(BootstrapError::WorkspaceNotFound { target: TargetId::parse("local").unwrap(), path: workspace.display().to_string(), } .into()); }
let current_os = std::env::consts::OS; let normalized_platform = match current_os { "macos" => "Darwin", "linux" => "Linux", other => other, };
if !spec .target_platform .eq_ignore_ascii_case(normalized_platform) { return Err(BootstrapError::PlatformMismatch { expected: spec.target_platform.clone(), actual: normalized_platform.to_string(), } .into()); }
fs::create_dir_all(bundles_root) .with_context(|| format!("cannot create bundle root {}", bundles_root.display()))?;
if let Ok(manifest) = check_bundle_readiness(bundles_root, spec) { return Ok(BootstrapOutcome::Warm(manifest)); }
let bundle_id = spec.compute_bundle_id(); let destination = bundles_root.join(&bundle_id);
let candidate = tempfile::Builder::new() .prefix(&format!(".candidate-{bundle_id}-")) .tempdir_in(bundles_root) .with_context(|| { format!( "cannot create candidate directory in {}", bundles_root.display() ) })?;
let candidate_path = candidate.path();
let staged_runner = candidate_path.join("runner.py"); if config.runner_source_path.exists() { fs::copy(&config.runner_source_path, &staged_runner)?; } else { fs::write(&staged_runner, b"import sys\nprint('klbr runner ready')\n")?; }
let staged_verify = candidate_path.join("verify.py"); if config.verify_script_path.exists() { fs::copy(&config.verify_script_path, &staged_verify)?; } else { fs::write( &staged_verify, b"import json\nprint(json.dumps({'status':'ok','capabilities':['runner'],'unavailable_capabilities':['dill']}))\n", )?; }
// Stage python module tree if available locally let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR")); let python_src = manifest_dir .parent() .unwrap_or(&manifest_dir) .join("python/klbr-runtime/src"); if python_src.exists() { let _ = copy_dir_all(&python_src, &candidate_path.join("src")); }
let python_bin = &config.default_python; let verify_output = Command::new(python_bin) .args(["-I", "-B", "-u"]) .arg(&staged_verify) .arg("--bundle-dir") .arg(candidate_path) .arg("--expected-platform") .arg(&spec.target_platform) .arg("--workspace") .arg(workspace) .output();
let output = match verify_output { Ok(out) => out, Err(err) => { return Err(BootstrapError::UnverifiedInterpreter { path: python_bin.to_string(), reason: format!("failed to spawn interpreter: {err}"), } .into()); } };
if !output.status.success() { let stderr = String::from_utf8_lossy(&output.stderr); return Err(BootstrapError::VerificationFailed { reason: format!( "smoke test failed with exit code {}: {}", output.status.code().unwrap_or(-1), stderr.trim() ), } .into()); }
let stdout_str = String::from_utf8_lossy(&output.stdout); let manifest = manifest_from_probe(&bundle_id, spec, python_bin, &stdout_str, || { staged_runner .file_name() .unwrap() .to_str() .unwrap() .to_string() })?;
let manifest_bytes = manifest.to_bytes()?; fs::write(candidate_path.join("manifest.json"), manifest_bytes)?;
if destination.exists() { if let Ok(m) = ReadinessManifest::load_from_dir(&destination) { return Ok(BootstrapOutcome::Warm(m)); } }
match fs::rename(candidate_path, &destination) { Ok(()) => Ok(BootstrapOutcome::Cold(manifest)), Err(_err) if destination.exists() => { let m = ReadinessManifest::load_from_dir(&destination)?; Ok(BootstrapOutcome::Warm(m)) } Err(err) => bail!( "failed to atomically publish candidate to {}: {err}", destination.display() ), }}
/// Prepares and guarantees execution readiness for a RemoteTarget.////// ALWAYS routes through host-owned SSH transport. Zero silent local fallback:/// an empty or missing alias is an explicit configuration error.pub async fn ensure_target_readiness( target: &RemoteTarget, spec: &BundleSpec, bundles_root: &Path, config: &BootstrapConfig,) -> Result<BootstrapOutcome> { ensure_target_readiness_ssh(target, spec, bundles_root, config).await}
/// Prepares and guarantees execution readiness according to the typed placement selector.pub async fn ensure_kernel_target_readiness( target: &KernelTarget, spec: &BundleSpec, bundles_root: &Path, config: &BootstrapConfig,) -> Result<BootstrapOutcome> { match target { KernelTarget::Local => { let cwd = std::env::current_dir().context("cannot resolve local current_dir")?; ensure_local_target_readiness(&cwd, spec, bundles_root, config) } KernelTarget::Remote(remote) => { ensure_target_readiness(remote, spec, bundles_root, config).await } }}
fn copy_dir_all(src: &Path, dst: &Path) -> Result<()> { fs::create_dir_all(dst)?; for entry in fs::read_dir(src)? { let entry = entry?; let file_type = entry.file_type()?; if file_type.is_dir() { copy_dir_all(&entry.path(), &dst.join(entry.file_name()))?; } else if file_type.is_file() { fs::copy(entry.path(), dst.join(entry.file_name()))?; } } Ok(())}
#[cfg(test)]mod probe_report_tests { use super::{bounded, probe_report, BootstrapError};
fn reason(error: BootstrapError) -> String { match error { BootstrapError::VerificationFailed { reason } => reason, other => panic!("expected a verification failure, got {other:?}"), } }
/// A probe that did not answer has not measured anything, so nothing may be claimed for it. #[test] fn an_unreadable_probe_measures_nothing() { for output in [ "", "not json at all", "{\"capabilities\": [\"runner\"]}", "{\"status\": \"failed\", \"capabilities\": []}", "{\"status\": \"ok\", \"capabilities\": \"runner\"}", ] { let error = probe_report(output).expect_err("must refuse"); assert!(reason(error).contains("probe"), "output {output:?}"); } }
/// What the probe reported is what the manifest carries; nothing is added and nothing is assumed. #[test] fn reported_capabilities_are_the_only_ones() { let report = probe_report( r#"{"status":"ok","python_path":"/usr/bin/python3","capabilities":["klbr_runtime","runner","dill"],"unavailable_capabilities":["skill:missing"]}"#, ) .expect("valid probe"); assert_eq!(report.capabilities, ["klbr_runtime", "runner", "dill"]); assert_eq!(report.unavailable_capabilities, ["skill:missing"]); assert_eq!(report.python_path.as_deref(), Some("/usr/bin/python3")); }
/// A remote probe can say anything, including something very long. #[test] fn a_long_reason_is_bounded() { let long = "x".repeat(5_000); let bounded = bounded(&long); assert!(bounded.len() < 300, "{} bytes", bounded.len()); assert!(bounded.contains("5000 bytes")); }}