use std::collections::{BTreeMap, BTreeSet}; use std::fs::File; use std::io::{BufRead, BufReader}; use std::net::{SocketAddr, TcpListener, TcpStream}; #[cfg(target_os = "linux")] use std::os::fd::FromRawFd; use std::path::{Path, PathBuf}; use std::process::Command; use std::sync::{Arc, mpsc}; use std::thread; use std::time::{Duration, Instant, SystemTime}; use crate::DEFAULT_PIPELINE_CACHED_MODEL_FILE; use crate::codecs::register_myelin_actor_codecs; use crate::node_actor::{ NodeAgentMsg, StageEdgeKindWire, StageInboundEdgeWire, StageObjectSpecWire, StageOutboundEdgeWire, StageProvisionWire, StageRingSpecWire, }; use crate::observability::frame_collector::{FrameCollector, StageLoadProgress}; use crate::observability::orch_telemetry::{ DashboardSupport, MYELIN_STAGE_ROUTE, MYELIN_SWIM_MEMBERSHIP, OrchTelemetry, }; use crate::orchestration::actor::{OrchestratorActor, OrchestratorMsg, OrchestratorReport}; use crate::orchestration::config::{DEFAULT_CONFIG_PATH, TomlConfigOverlay}; use crate::gguf_shard::{StageShardPlan, plan_stage_shard}; use crate::node_provisioning::{ProviderKind, provider_kind}; use crate::orchestration::cluster_reconciler::{ProvisionedClusterGuard, ReconcilerNodeBinding}; use crate::orchestration::distribution_stack::{DistributionRuntimeStack, duration_ms_u64}; use crate::orchestration::provider_adapters::relay::{ MYELIN_IROH_RELAY_URL_ENV, RelayRuntimeConfig, SWACTOR_IROH_RELAY_URL_ENV, relay_mode_env_value, relay_runtime_config_from_settings, }; use crate::orchestration::provider_adapters::vastai::{ SshCommandBootstrapLauncher, ToolsVastAiLeaseClient, VastAiProvisioningConfig, VastAiProvisioningPlugin, }; use crate::prompt::rpc::{ PromptEvent, SubmitPrompt, TokenizerEvent, read_submit_prompt, write_json_line, }; use crate::provisioning::{ LocalDockerPlugin, LocalProcessPlugin, NodeProvisionSpec, PluginObservation, PluginObservationSink, PluginSink, ProviderMount, ProvisionEvent, ProvisionEventKind, ProvisionLogLine, ProvisionLogStream, ProvisionPlugin, }; use crate::run_fsm::{RunConfig, RunId}; use crate::run_plan::{self, GgufSource, TokenizerSource}; use ::provisioning::{ BootSpec, ClusterShape, DesiredNodeShape, LogicalNodeId as ReconcilerLogicalNodeId, NodeGroupId, ProviderKind as ReconcilerProviderKind, RetryPolicy, RoleId, RunId as ClusterRunId, RunNodeGroupSpec, SwactorId, SwarmJoinTemplate, }; use data_plane::object_record as ingress; use telemetry::{TelemetryPublisherMsg, TelemetrySubscribe, SubscriptionRequest}; use distribution::node::DistributedNodeConfig; use distribution::swim::telemetry::ObservedTransition; use distribution::types::{MemberState, NodeId as DistNodeId}; use iroh::EndpointAddr; use data_plane::edge_wire::WireEvent; use iroh_driver::{ TELEMETRY_ALPN, EDGE_ALPN, EdgeSendHandle, IrohDriver, IrohDriverConfig, }; use iroh_driver::{EndpointAddrMask, MVP_IROH_ENDPOINT_ADDR_MASK_ENV, advertised_endpoint}; use parking_lot::Mutex; use serde_json::{Value, json}; use swactor::actor::ActorAddress; use swactor_engine::{Engine, EngineHandle, TokioBackend, TokioConfig}; const DEFAULT_IMAGE: &str = "myelin-node:latest"; const MYELIN_RUNTIME_CONFIG_ENV: &str = "MYELIN_RUNTIME_CONFIG"; const CACHED_MODEL_HOST_ENV: &str = "MYELIN_CACHED_MODEL_HOST_PATH"; const MYELIN_WORKER_BIN_ENV: &str = "MYELIN_WORKER_BIN"; const CACHED_MODEL_CONTAINER_DIR: &str = "/models/cached"; const DEFAULT_PIPELINE_MODEL_CACHE_DIR: &str = ".model-cache"; const DEFAULT_RPC_BIND: &str = "127.0.0.1:19777"; const DEFAULT_HF_REPO: &str = "bartowski/Llama-3.2-1B-Instruct-GGUF"; const DEFAULT_HF_FILE: &str = "Llama-3.2-1B-Instruct-Q4_K_M.gguf"; const DEFAULT_MODEL_ID: &str = "llama-3.2-1b-instruct-q4"; const DEFAULT_MAX_TOKENS: u32 = 64; const PUMP_INTERVAL: Duration = Duration::from_millis(10); const RUNTIME_READY_ACK_RETRY_INTERVAL: Duration = Duration::from_millis(250); const STAGE_PROVISION_ACTIVE_RESEND_AFTER: Duration = Duration::from_secs(60); const PIPELINE_PROMPT_WAIT_LOG_INTERVAL: Duration = Duration::from_secs(15); const TELEMETRY_FRAME_LOG_ENV: &str = "MYELIN_TELEMETRY_FRAME_LOG"; pub(crate) fn run_with_options( args: I, capture_stdio: bool, stop_rx: Option>, ) -> Result<(), String> where I: IntoIterator, { let mut config_builder = ConfigBuilder::hardcoded_defaults(); if let Some(overlay) = TomlConfigOverlay::load_optional(Path::new(DEFAULT_CONFIG_PATH))? { config_builder = config_builder.overlay_toml(overlay)?; } let mut config = config_builder .overlay_env()? .overlay_cli(args)? .finalize()?; config.prepare_vastai_ssh_key()?; let orch_stdio_rx = if capture_stdio { OrchStdioCapture::install()? } else { None }; let mut orch_telemetry = OrchTelemetry::new(config.run_id, config.telemetry_frame_log.as_deref())?; let run_id = config.run_id; let node_id = config.node_id; let bootstrap = |ds: &mut OrchTelemetry, dash: Option<&DashboardSupport>, phase: &str, status: &str, detail: Value| { ds.emit_bootstrap(dash, run_id, node_id, phase, status, detail); }; bootstrap( &mut orch_telemetry, None, "config", "ready", json!({ "config_profile":match config.config_profile { RuntimeConfigProfile::Local => "local", RuntimeConfigProfile::Deploy => "deploy", }, "image":&config.image, "provider":config.provider.as_str(), "rpc_bind":config.rpc_bind.to_string(), "model_id":&config.model_id, "stage_index":config.stage_index, "legacy_layer_end_exclusive":config.layer_end_exclusive, "relay_mode":format!("{:?}", config.relay.mode), "endpoint_addr_mask":config.endpoint_addr_mask.as_str(), "pipeline_stages":config.pipeline_stages, "provider_config":config.provider_telemetry_detail(), }), ); bootstrap( &mut orch_telemetry, None, "telemetry_preflight", "configured", json!({ "producer":"myelin-orchestrator", "telemetry_endpoint":{ "role":"orchestrator-frame-archive", "transport":"telemetry-frame-log", "configured":config.telemetry_frame_log.is_some(), "archive_path":config.telemetry_frame_log.as_ref().map(|path| path.to_string_lossy().to_string()), }, "expected_worker_producers":["myelin-worker","tinygrad-worker"], "provider":config.provider.as_str(), "pipeline_stages":config.pipeline_stages, "endpoint_addr_mask":config.endpoint_addr_mask.as_str(), "provider_config":config.provider_telemetry_detail(), }), ); let orch_synthetic_id = format!("myelin-orchestrator-{}-telemetry-preflight", config.run_id); for (phase, status) in [ ("TelemetryProducerConfigured", "configured"), ("TelemetryProducerConnected", "ready"), ("TelemetrySyntheticEventSent", "sent"), ("TelemetrySyntheticEventObserved", "observed"), ] { bootstrap( &mut orch_telemetry, None, phase, status, json!({ "producer":"myelin-orchestrator", "producer_class":"rust-orchestrator", "synthetic_id":orch_synthetic_id, "telemetry_endpoint":{ "role":"orchestrator-frame-archive", "transport":"telemetry-frame-log", "configured":config.telemetry_frame_log.is_some(), "archive_path":config.telemetry_frame_log.as_ref().map(|path| path.to_string_lossy().to_string()), }, }), ); } drain_orch_stdio_capture( orch_stdio_rx.as_ref(), &mut orch_telemetry, None, config.run_id, config.node_id, ); let pipeline_plan = if config.uses_planned_execution() { let plan = config.build_run_plan()?; bootstrap( &mut orch_telemetry, None, "run_plan", "ready", json!({ "stage_count":plan.stages.len(), "edge_count":plan.edges.len(), "model_layers":plan.model.num_layers, "hidden_dim":plan.model.hidden_dim, "max_seq_len":plan.model.max_seq_len, "eos_token_id":plan.model.eos_token_id, }), ); Some(plan) } else { None }; let actors_channel = orch_telemetry.channel_by_name("runtime.actors"); let orch_stats_hook = orch_telemetry.stats_hook_on(actors_channel); // Build the core swactor runtime parts, clone the routing handle needed by // integrations, then hand the workers to the engine. The engine owns both // core progression and the Tokio substrate (it schedules all background // work); components retain only cheap Runtime handles (ENGINE_SPEC.md). let (parts, runtime, codec, transport_router) = DistributionRuntimeStack::build_runtime( |registry| { register_myelin_actor_codecs(registry); telemetry::wire::register_telemetry_codec(registry); }, Some(orch_stats_hook), ); let engine = match TokioBackend::new(TokioConfig::default()) .and_then(|backend| Engine::new(parts, backend)) { Ok(engine) => { bootstrap( &mut orch_telemetry, None, "engine", "ready", json!({"backend":"tokio","owns":"core+substrate"}), ); engine } Err(error) => { bootstrap( &mut orch_telemetry, None, "engine", "failed", json!({"error":error.to_string()}), ); return Err(format!("create engine: {error}")); } }; let mut driver = match IrohDriver::with_engine( engine.handle(), IrohDriverConfig { secret_key: None, relay_mode: config.relay.mode.clone(), node: DistributedNodeConfig::default(), peer_auth: None, additional_alpns: vec![EDGE_ALPN.to_vec(), TELEMETRY_ALPN.to_vec()], }, ) { Ok(driver) => driver, Err(error) => { bootstrap( &mut orch_telemetry, None, "iroh_driver", "failed", json!({"error":error.to_string()}), ); return Err(format!("create iroh driver: {error}")); } }; let coordinator_endpoint = advertised_endpoint(driver.endpoint_addr(), config.endpoint_addr_mask)?; bootstrap( &mut orch_telemetry, None, "iroh_driver", "ready", json!({"endpoint":coordinator_endpoint.clone(),"has_relay":coordinator_endpoint.relay_urls().next().is_some(),"direct_addr_count":coordinator_endpoint.ip_addrs().count(),"relay_mode":format!("{:?}", config.relay.mode),"endpoint_addr_mask":config.endpoint_addr_mask.as_str()}), ); bootstrap( &mut orch_telemetry, None, "endpoint_config_snapshot", "ready", json!({ "producer":"myelin-orchestrator", "coordinator_endpoint":coordinator_endpoint.clone(), "has_relay":coordinator_endpoint.relay_urls().next().is_some(), "direct_addr_count":coordinator_endpoint.ip_addrs().count(), "relay_mode":format!("{:?}", config.relay.mode), "endpoint_addr_mask":config.endpoint_addr_mask.as_str(), "connectivity_preflight":"ready", }), ); let stack = DistributionRuntimeStack::new_from_runtime( runtime, codec, transport_router, driver.node_id(), DistributedNodeConfig::default(), engine.handle(), ); bootstrap( &mut orch_telemetry, None, "distribution_stack", "ready", json!({"actors":"initialized","route_view":"initialized","swim":"initialized"}), ); bootstrap( &mut orch_telemetry, None, "codecs", "ready", json!({"registered":["node_agent","orchestrator","provisioner","prompt_rpc","telemetry"]}), ); driver.enable_actor_bridge( stack.runtime.clone(), stack.codec.clone(), stack.actor_bridge_routes(), stack.actors.swim, stack.relay_mirror.clone(), stack.route_view.clone(), stack.outbox.clone(), ); // Engine owns protocol tick injection and core progression; the application // loop only drains integration-owned queues (ENGINE_SPEC.md). stack.spawn_protocol_ticker(PUMP_INTERVAL); driver.install_actor_bridge_pump(PUMP_INTERVAL); bootstrap( &mut orch_telemetry, None, "actor_bridge", "ready", json!({"transport":"iroh","routes":"attached","protocol_ticker":"engine-hosted"}), ); let collector = FrameCollector::new(); bootstrap( &mut orch_telemetry, None, "telemetry_collector", "ready", json!({"alpn":String::from_utf8_lossy(TELEMETRY_ALPN)}), ); let dashboard = DashboardSupport::start(config.dashboard, &engine.handle())?; bootstrap( &mut orch_telemetry, dashboard.as_ref(), "dashboard", "ready", json!({"enabled":dashboard.is_some()}), ); let orchestrator_reports = match stack.runtime.new_inbox::() { Ok(inbox) => inbox, Err(error) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "orchestrator_report_actor", "failed", json!({"error":error.to_string()}), ); return Err(format!("orchestrator report inbox: {error}")); } }; let orchestrator_report_actor = *orchestrator_reports.addr(); stack.register_local_actor(driver.register_actor(orchestrator_report_actor, 1)); bootstrap( &mut orch_telemetry, dashboard.as_ref(), "orchestrator_report_actor", "ready", json!({"actor":orchestrator_report_actor}), ); let orchestrator_actor = match stack.runtime.spawn(OrchestratorActor::new( RunConfig { run_id: RunId(config.run_id), max_tokens: u64::from(config.default_max_tokens), prompt: Vec::new(), }, Some(orchestrator_report_actor), )) { Ok(actor) => actor, Err(error) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "orchestrator_actor", "failed", json!({"error":error.to_string()}), ); return Err(format!("spawn orchestrator actor: {error}")); } }; stack.register_local_actor(driver.register_actor(orchestrator_actor, 1)); bootstrap( &mut orch_telemetry, dashboard.as_ref(), "orchestrator_actor", "ready", json!({"actor":orchestrator_actor}), ); let prompt_events = match stack.runtime.new_inbox::() { Ok(inbox) => inbox, Err(error) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_reply_actor", "failed", json!({"error":error.to_string()}), ); return Err(format!("prompt event inbox: {error}")); } }; let prompt_reply_actor = *prompt_events.addr(); stack.register_local_actor(driver.register_actor(prompt_reply_actor, 1)); bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_reply_actor", "ready", json!({"actor":prompt_reply_actor}), ); let tokenizer_events = match stack.runtime.new_inbox::() { Ok(inbox) => inbox, Err(error) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "tokenizer_reply_actor", "failed", json!({"error":error.to_string()}), ); return Err(format!("tokenizer event inbox: {error}")); } }; let tokenizer_reply_actor = *tokenizer_events.addr(); stack.register_local_actor(driver.register_actor(tokenizer_reply_actor, 1)); bootstrap( &mut orch_telemetry, dashboard.as_ref(), "tokenizer_reply_actor", "ready", json!({"actor":tokenizer_reply_actor}), ); let (work_tx, work_rx) = mpsc::channel::(); let stop_rx = stop_rx.unwrap_or_else(spawn_stop_listener); let provisioner = config.build_provisioner(stack.runtime.clone())?; bootstrap( &mut orch_telemetry, dashboard.as_ref(), "node_provisioner", "ready", json!({ "provider":config.provider.as_str(), "owner":"myelin-orchestrator", "config":config.provider_telemetry_detail(), }), ); let (obs_tx, obs_rx) = mpsc::channel::(); let sink = PluginSink::new(Arc::new(ChannelObservationSink { tx: Mutex::new(obs_tx), })); let pipeline_coordinator_endpoint = coordinator_endpoint.clone(); let (mut provisioned_nodes, ready, swim_to_node) = start_and_provision_workers( provisioner, &config, pipeline_plan.as_ref(), RuntimeReadyAckLoop { driver: &mut driver, stack: &stack, obs_rx: &obs_rx, collector: &collector, orchestrator_reports: &orchestrator_reports, stop_rx: &stop_rx, dashboard: dashboard.as_ref(), orch_telemetry: &mut orch_telemetry, orch_stdio_rx: orch_stdio_rx.as_ref(), run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, sink, coordinator_endpoint, pipeline_coordinator_endpoint, orchestrator_actor, )?; bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_rpc", "started", json!({ "bind":config.rpc_bind.to_string(), "default_max_tokens":config.default_max_tokens, }), ); let rpc_addr = match spawn_prompt_rpc( &engine.handle(), config.rpc_bind, work_tx, config.default_max_tokens, ) { Ok(addr) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_rpc", "ready", json!({ "addr":addr.to_string(), "default_max_tokens":config.default_max_tokens, }), ); addr } Err(error) => { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_rpc", "failed", json!({"error":error}), ); return Err(error); } }; bootstrap( &mut orch_telemetry, dashboard.as_ref(), "prompt_loop", "ready", json!({"addr":rpc_addr.to_string(),"node_actor":ready.first_stage.node_actor}), ); bootstrap( &mut orch_telemetry, dashboard.as_ref(), "serve_prompts", "started", json!({"mode":"single_active_prompt","poll_interval_ms":PUMP_INTERVAL.as_millis()}), ); let result = serve_prompts( RuntimeReadyAckLoop { driver: &mut driver, stack: &stack, obs_rx: &obs_rx, collector: &collector, orchestrator_reports: &orchestrator_reports, stop_rx: &stop_rx, dashboard: dashboard.as_ref(), orch_telemetry: &mut orch_telemetry, orch_stdio_rx: orch_stdio_rx.as_ref(), run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, &work_rx, &prompt_events, ready.first_stage.node_actor, prompt_reply_actor, &tokenizer_events, ready.first_stage.node_actor, ready.final_stage.node_actor, tokenizer_reply_actor, pipeline_plan.as_ref(), ready.first_stage.endpoint.clone(), &swim_to_node, ); if let Err(error) = &result { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "serve_prompts", "failed", json!({"error":error}), ); } bootstrap( &mut orch_telemetry, dashboard.as_ref(), "provider_stop", "started", json!({"provider":config.provider.as_str(),"node_id":config.node_id}), ); let stop_result = provisioned_nodes.stop(); match &stop_result { Ok(()) => bootstrap( &mut orch_telemetry, dashboard.as_ref(), "provider_stop", "ready", json!({"provider":config.provider.as_str(),"node_id":config.node_id}), ), Err(error) => bootstrap( &mut orch_telemetry, dashboard.as_ref(), "provider_stop", "failed", json!({"provider":config.provider.as_str(),"node_id":config.node_id,"error":error}), ), } if result.is_ok() && stop_result.is_ok() { bootstrap( &mut orch_telemetry, dashboard.as_ref(), "orch_exit", "ready", json!({"result":"ok"}), ); } result.and(stop_result) } #[derive(Clone)] struct VastAiRuntimeConfig { api_key: Option, provisioning: VastAiProvisioningConfig, bootstrap_command: Option, ssh_identity: Option, ssh_public_fingerprint: Option, } impl VastAiRuntimeConfig { fn from_builder(builder: &ConfigBuilder) -> Result { let mut provisioning = VastAiProvisioningConfig::default(); macro_rules! raw_config { ($parser:ident, $env:literal, $raw:expr, $value:expr) => { $raw.as_ref() .map(|value| ConfigBuilder::$parser($env, value)) .transpose()? .or($value) }; } if let Some(disk_gb) = raw_config!( parse_value, "MYELIN_VASTAI_DISK_GB", builder.vastai_disk_gb_raw, builder.vastai_disk_gb ) { provisioning.disk_gb = disk_gb; } if let Some(ssh_user) = &builder.vastai_ssh_user { provisioning.ssh_user = ssh_user.clone(); } if let Some(confirm_lease) = raw_config!( parse_bool, "MYELIN_VASTAI_CONFIRM_LEASE", builder.vastai_confirm_lease_raw, builder.vastai_confirm_lease ) { provisioning.confirm_lease = confirm_lease; } provisioning.onstart = builder.vastai_onstart.clone(); provisioning.selection.gpu_name = builder.vastai_gpu_name.clone(); if let Some(min_gpu_ram_mb) = raw_config!( parse_value, "MYELIN_VASTAI_MIN_GPU_RAM_MB", builder.vastai_min_gpu_ram_mb_raw, builder.vastai_min_gpu_ram_mb ) { provisioning.selection.min_gpu_ram_mb = Some(min_gpu_ram_mb); } if let Some(min_down_mbps) = raw_config!( parse_value, "MYELIN_VASTAI_MIN_DOWN_MBPS", builder.vastai_min_down_mbps_raw, builder.vastai_min_down_mbps ) { provisioning.selection.min_down_mbps = min_down_mbps; } if let Some(max_dph_total) = raw_config!( parse_value, "MYELIN_VASTAI_MAX_DPH_TOTAL", builder.vastai_max_dph_total_raw, builder.vastai_max_dph_total ) { provisioning.selection.max_dph_total = Some(max_dph_total); } if let Some(min_up_mbps) = raw_config!( parse_value, "MYELIN_VASTAI_MIN_UP_MBPS", builder.vastai_min_up_mbps_raw, builder.vastai_min_up_mbps ) { provisioning.selection.min_up_mbps = Some(min_up_mbps); } if let Some(min_reliability) = raw_config!( parse_value, "MYELIN_VASTAI_MIN_RELIABILITY", builder.vastai_min_reliability_raw, builder.vastai_min_reliability ) { provisioning.selection.min_reliability = min_reliability; } if let Some(require_verified) = raw_config!( parse_bool, "MYELIN_VASTAI_REQUIRE_VERIFIED", builder.vastai_require_verified_raw, builder.vastai_require_verified ) { provisioning.selection.require_verified = require_verified; } for host_id in &builder.vastai_blacklist_hosts { if !provisioning.selection.blacklist_hosts.contains(host_id) { provisioning.selection.blacklist_hosts.push(*host_id); } } if let Some(poll_interval_secs) = raw_config!( parse_value, "MYELIN_VASTAI_POLL_INTERVAL_SECS", builder.vastai_poll_interval_secs_raw, builder.vastai_poll_interval_secs ) { provisioning.lifecycle.poll_interval = Duration::from_secs(poll_interval_secs); } let ssh_identity = builder .vastai_ssh_identity_raw .as_ref() .map(|value| expand_home_path(value)) .transpose()?; Ok(Self { api_key: builder.vastai_api_key.clone(), provisioning, bootstrap_command: builder.vastai_bootstrap_command.clone(), ssh_identity, ssh_public_fingerprint: None, }) } fn telemetry_detail(&self) -> Value { json!({ "disk_gb": self.provisioning.disk_gb, "ssh_user": &self.provisioning.ssh_user, "gpu_name": &self.provisioning.selection.gpu_name, "min_gpu_ram_mb": self.provisioning.selection.min_gpu_ram_mb, "min_compute_cap": self.provisioning.selection.min_compute_cap, "min_down_mbps": self.provisioning.selection.min_down_mbps, "min_up_mbps": self.provisioning.selection.min_up_mbps, "max_dph_total": self.provisioning.selection.max_dph_total, "min_reliability": self.provisioning.selection.min_reliability, "require_verified": self.provisioning.selection.require_verified, "blacklist_hosts": &self.provisioning.selection.blacklist_hosts, "state_timeout_secs": self.provisioning.lifecycle.state_timeout.as_secs(), "confirm_lease": self.provisioning.confirm_lease, "has_api_key": self.api_key.is_some(), "has_onstart": self.provisioning.onstart.is_some(), "has_bootstrap_command": self.bootstrap_command.is_some(), "has_ssh_identity": self.ssh_identity.is_some(), "ssh_public_fingerprint": self.ssh_public_fingerprint.as_deref(), }) } } #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum RuntimeConfigProfile { Local, Deploy, } impl RuntimeConfigProfile { fn parse(value: &str) -> Result { match value.trim().to_ascii_lowercase().as_str() { "local" => Ok(Self::Local), "deploy" => Ok(Self::Deploy), other => Err(format!( "unsupported {MYELIN_RUNTIME_CONFIG_ENV}={other:?}; use local or deploy" )), } } } #[derive(Clone)] struct CachedModelConfig { host_path: PathBuf, container_path: String, } impl CachedModelConfig { fn from_host_path(provider: &str, requested: PathBuf) -> Result { if !matches!(provider, "process" | "docker" | "vastai") { return Err(format!( "{CACHED_MODEL_HOST_ENV} is a host-local cache path and is only supported by provider=process, provider=docker, or vastai planning" )); } let host_path = requested.canonicalize().map_err(|e| { format!( "resolve {CACHED_MODEL_HOST_ENV} path {}: {e}", requested.display() ) })?; if !host_path.is_file() { return Err(format!( "{CACHED_MODEL_HOST_ENV} must point at a file: {}", host_path.display() )); } let file_name = host_path .file_name() .and_then(|name| name.to_str()) .filter(|name| !name.is_empty()) .ok_or_else(|| { format!( "cached model path has no file name: {}", host_path.display() ) })?; let container_path = format!("{CACHED_MODEL_CONTAINER_DIR}/{file_name}"); Ok(Self { host_path, container_path, }) } fn telemetry_detail(&self) -> Value { json!({ "host_path_present": true, "file": self.host_path.file_name().and_then(|name| name.to_str()), "container_path": &self.container_path, }) } } fn default_pipeline_cached_model_path() -> PathBuf { let relative = PathBuf::from(".") .join(DEFAULT_PIPELINE_MODEL_CACHE_DIR) .join(DEFAULT_PIPELINE_CACHED_MODEL_FILE); if relative.is_file() { return relative; } PathBuf::from(env!("CARGO_MANIFEST_DIR")) .join("..") .join("..") .join(DEFAULT_PIPELINE_MODEL_CACHE_DIR) .join(DEFAULT_PIPELINE_CACHED_MODEL_FILE) } #[derive(Clone)] struct Config { config_profile: RuntimeConfigProfile, image: String, docker_gpus: String, provider: ProviderKind, rpc_bind: SocketAddr, run_id: u64, node_id: u64, stage_index: u32, layer_end_exclusive: Option, pipeline_stages: u32, model_id: String, gguf_source: GgufSource, tokenizer: TokenizerSource, default_max_tokens: u32, dashboard: bool, max_context: Option, relay: RelayRuntimeConfig, endpoint_addr_mask: EndpointAddrMask, vastai: Option, cached_model: Option, telemetry_frame_log: Option, worker_bin: Option, } #[derive(Clone)] struct ConfigBuilder { config_profile: RuntimeConfigProfile, provider: Option, image: String, toml_vastai_image: Option, image_overridden_after_toml: bool, docker_gpus: String, rpc_bind: String, rpc_bind_label: &'static str, run_id: u64, node_id: u64, stage_index: u32, layer_end_exclusive: Option, pipeline_stages: u32, model_id: String, gguf_source: GgufSource, tokenizer: TokenizerSource, default_max_tokens: u32, dashboard: bool, max_context: Option, relay_mode: Option, relay_url: Option, endpoint_addr_mask: Option, vastai_api_key: Option, vastai_bootstrap_command: Option, vastai_disk_gb: Option, vastai_disk_gb_raw: Option, vastai_ssh_user: Option, vastai_confirm_lease: Option, vastai_confirm_lease_raw: Option, vastai_onstart: Option, vastai_ssh_identity_raw: Option, vastai_gpu_name: Option, vastai_min_gpu_ram_mb: Option, vastai_min_gpu_ram_mb_raw: Option, vastai_min_down_mbps: Option, vastai_min_down_mbps_raw: Option, vastai_min_up_mbps: Option, vastai_max_dph_total: Option, vastai_max_dph_total_raw: Option, vastai_min_up_mbps_raw: Option, vastai_min_reliability: Option, vastai_min_reliability_raw: Option, vastai_require_verified: Option, vastai_require_verified_raw: Option, vastai_poll_interval_secs: Option, vastai_blacklist_hosts: Vec, vastai_poll_interval_secs_raw: Option, cached_model_host_path: Option, telemetry_frame_log: Option, worker_bin: Option, } impl ConfigBuilder { fn hardcoded_defaults() -> Self { Self { config_profile: RuntimeConfigProfile::Local, provider: None, image: DEFAULT_IMAGE.to_owned(), toml_vastai_image: None, image_overridden_after_toml: false, docker_gpus: "all".to_owned(), rpc_bind: DEFAULT_RPC_BIND.to_owned(), rpc_bind_label: "MYELIN_PROMPT_RPC_BIND", run_id: 1, node_id: 1, stage_index: 0, layer_end_exclusive: None, pipeline_stages: 1, model_id: DEFAULT_MODEL_ID.to_owned(), gguf_source: GgufSource::HuggingFaceGguf { repo: DEFAULT_HF_REPO.to_owned(), file: DEFAULT_HF_FILE.to_owned(), revision: None, }, tokenizer: TokenizerSource::EmbeddedGguf, default_max_tokens: DEFAULT_MAX_TOKENS, dashboard: false, max_context: None, relay_mode: None, relay_url: None, endpoint_addr_mask: None, vastai_api_key: None, vastai_bootstrap_command: None, vastai_disk_gb: None, vastai_disk_gb_raw: None, vastai_ssh_user: None, vastai_confirm_lease: None, vastai_confirm_lease_raw: None, vastai_onstart: None, vastai_ssh_identity_raw: None, vastai_gpu_name: None, vastai_min_gpu_ram_mb: None, vastai_min_gpu_ram_mb_raw: None, vastai_min_down_mbps: None, vastai_min_down_mbps_raw: None, vastai_max_dph_total: None, vastai_max_dph_total_raw: None, vastai_min_up_mbps: None, vastai_min_up_mbps_raw: None, vastai_min_reliability: None, vastai_min_reliability_raw: None, vastai_require_verified: None, vastai_require_verified_raw: None, vastai_poll_interval_secs: None, vastai_blacklist_hosts: Vec::new(), vastai_poll_interval_secs_raw: None, cached_model_host_path: None, telemetry_frame_log: None, worker_bin: None, } } fn overlay_toml(mut self, overlay: TomlConfigOverlay) -> Result { macro_rules! apply { ($option:expr, |$value:ident| $body:block) => { if let Some($value) = $option $body }; } apply!(overlay.runtime.profile, |profile| { self.config_profile = RuntimeConfigProfile::parse(&profile)?; }); apply!(overlay.runtime.run_id, |run_id| { self.run_id = run_id }); apply!(overlay.runtime.node_id, |node_id| { self.node_id = node_id }); apply!(overlay.runtime.stage_index, |stage_index| { self.stage_index = stage_index }); apply!(overlay.runtime.layer_end_exclusive, |layer_end_exclusive| { self.layer_end_exclusive = Some(layer_end_exclusive) }); apply!(overlay.runtime.pipeline_stages, |pipeline_stages| { self.pipeline_stages = pipeline_stages }); apply!(overlay.provider.kind, |provider| { self.provider = Some(provider_kind::parse_deploy(&provider)?); }); apply!(overlay.image.node, |image| { self.image = image }); apply!(overlay.relay.mode, |mode| { self.relay_mode = Some(mode) }); apply!(overlay.relay.url, |url| { self.relay_url = Some(url) }); apply!(overlay.prompt.rpc_addr, |rpc_bind| { self.rpc_bind = rpc_bind; self.rpc_bind_label = "[prompt].rpc_addr"; }); apply!(overlay.prompt.max_tokens, |max_tokens| { self.default_max_tokens = max_tokens }); apply!(overlay.prompt.dashboard, |dashboard| { self.dashboard = dashboard }); apply!(overlay.model.id, |model_id| { self.model_id = model_id }); apply!(overlay.model.gguf_local_path, |path| { self.gguf_source = GgufSource::LocalPath(path) }); apply!(overlay.model.gguf_repo, |repo| { self.set_hf_source(Some(repo), None, None) }); apply!(overlay.model.gguf_file, |file| { self.set_hf_source(None, Some(file), None) }); apply!(overlay.model.gguf_revision, |revision| { self.set_hf_source(None, None, Some(Some(revision))) }); apply!(overlay.model.tokenizer_local_path, |path| { self.tokenizer = TokenizerSource::LocalPath(path) }); apply!(overlay.model.max_context, |max_context| { self.max_context = Some(max_context) }); apply!(overlay.docker.gpus, |gpus| { self.docker_gpus = gpus }); apply!(overlay.docker.cached_model_host_path, |path| { self.cached_model_host_path = Some(PathBuf::from(path)) }); apply!(overlay.observability.telemetry_frame_log, |path| { self.telemetry_frame_log = Some(PathBuf::from(path)) }); apply!(overlay.vastai.image, |image| { self.toml_vastai_image = Some(image) }); apply!(overlay.vastai.api_key, |api_key| { self.vastai_api_key = Some(api_key) }); apply!(overlay.vastai.bootstrap_command, |command| { self.vastai_bootstrap_command = Some(command) }); apply!(overlay.vastai.disk_gb, |disk_gb| { self.vastai_disk_gb = Some(disk_gb) }); apply!(overlay.vastai.ssh_user, |ssh_user| { self.vastai_ssh_user = Some(ssh_user) }); apply!(overlay.vastai.confirm_lease, |confirm_lease| { self.vastai_confirm_lease = Some(confirm_lease) }); apply!(overlay.vastai.onstart, |onstart| { self.vastai_onstart = Some(onstart) }); apply!(overlay.vastai.ssh_identity, |identity| { self.vastai_ssh_identity_raw = Some(identity) }); apply!(overlay.vastai.gpu_name, |gpu_name| { self.vastai_gpu_name = Some(gpu_name) }); apply!(overlay.vastai.min_gpu_ram_mb, |min_gpu_ram_mb| { self.vastai_min_gpu_ram_mb = Some(min_gpu_ram_mb) }); apply!(overlay.vastai.min_down_mbps, |min_down_mbps| { self.vastai_min_down_mbps = Some(min_down_mbps) }); apply!(overlay.vastai.max_dph_total, |max_dph_total| { self.vastai_max_dph_total = Some(max_dph_total) }); apply!(overlay.vastai.min_up_mbps, |min_up_mbps| { self.vastai_min_up_mbps = Some(min_up_mbps) }); apply!(overlay.vastai.min_reliability, |min_reliability| { self.vastai_min_reliability = Some(min_reliability) }); apply!(overlay.vastai.require_verified, |require_verified| { self.vastai_require_verified = Some(require_verified) }); for host_id in overlay.vastai.blacklist_hosts { self.push_vastai_blacklist_host(host_id); } apply!(overlay.vastai.poll_interval_secs, |poll_interval_secs| { self.vastai_poll_interval_secs = Some(poll_interval_secs) }); Ok(self) } fn overlay_env(mut self) -> Result { macro_rules! apply { ($option:expr, |$value:ident| $body:block) => { if let Some($value) = $option $body }; } macro_rules! env_apply { ($name:expr, |$value:ident| $body:block) => { apply!(env_optional($name), |$value| $body) }; } macro_rules! env_parse { ($name:expr, |$value:ident| $body:block) => { env_apply!($name, |raw| { let $value = Self::parse_value($name, &raw)?; $body }) }; } env_apply!(MYELIN_RUNTIME_CONFIG_ENV, |profile| { self.config_profile = RuntimeConfigProfile::parse(&profile)?; }); env_parse!("MYELIN_RUN_ID", |run_id| { self.run_id = run_id }); env_parse!("MYELIN_LOGICAL_NODE_ID", |node_id| { self.node_id = node_id }); env_parse!("MYELIN_STAGE_INDEX", |stage_index| { self.stage_index = stage_index }); env_parse!("MYELIN_LAYER_END_EXCLUSIVE", |layer_end_exclusive| { self.layer_end_exclusive = Some(layer_end_exclusive) }); env_parse!("MYELIN_PIPELINE_STAGES", |pipeline_stages| { self.pipeline_stages = pipeline_stages }); apply!( env_optional("MYELIN_NODE_PROVIDER").or_else(|| env_optional("MYELIN_PROVIDER")), |provider| { self.provider = Some(provider_kind::parse_deploy(&provider)?); } ); env_apply!("MYELIN_NODE_IMAGE", |image| { self.image = image; self.image_overridden_after_toml = true; }); env_apply!("MYELIN_DOCKER_GPUS", |gpus| { self.docker_gpus = gpus }); env_apply!(CACHED_MODEL_HOST_ENV, |path| { self.cached_model_host_path = Some(PathBuf::from(path)) }); env_apply!(MYELIN_WORKER_BIN_ENV, |path| { self.worker_bin = Some(PathBuf::from(path)) }); env_apply!("MYELIN_PROMPT_RPC_BIND", |rpc_bind| { self.rpc_bind = rpc_bind; self.rpc_bind_label = "MYELIN_PROMPT_RPC_BIND"; }); env_parse!("MYELIN_PROMPT_MAX_TOKENS", |max_tokens| { self.default_max_tokens = max_tokens }); env_apply!("MYELIN_DASHBOARD", |dashboard| { self.dashboard = Self::parse_bool("MYELIN_DASHBOARD", &dashboard)?; }); env_apply!(TELEMETRY_FRAME_LOG_ENV, |path| { self.telemetry_frame_log = Some(PathBuf::from(path)) }); env_apply!("MYELIN_MODEL_ID", |model_id| { self.model_id = model_id }); env_apply!("MYELIN_GGUF_LOCAL_PATH", |path| { self.gguf_source = GgufSource::LocalPath(path) }); env_apply!("MYELIN_GGUF_REPO", |repo| { self.set_hf_source(Some(repo), None, None) }); env_apply!("MYELIN_GGUF_FILE", |file| { self.set_hf_source(None, Some(file), None) }); env_apply!("MYELIN_GGUF_REVISION", |revision| { self.set_hf_source(None, None, Some(Some(revision))) }); env_apply!("MYELIN_TOKENIZER_LOCAL_PATH", |path| { self.tokenizer = TokenizerSource::LocalPath(path) }); env_parse!("MYELIN_MAX_CONTEXT", |max_context| { self.max_context = Some(max_context) }); env_apply!("MYELIN_IROH_RELAY_MODE", |mode| { self.relay_mode = Some(mode.to_ascii_lowercase()) }); apply!( env_optional(MYELIN_IROH_RELAY_URL_ENV) .or_else(|| env_optional(SWACTOR_IROH_RELAY_URL_ENV)), |url| { self.relay_url = Some(url); } ); env_apply!(MVP_IROH_ENDPOINT_ADDR_MASK_ENV, |mask| { self.endpoint_addr_mask = Some(mask) }); apply!( env_optional("VAST_API_KEY") .or_else(|| env_optional("MYELIN_VASTAI_API_KEY")) .or_else(|| env_optional("VASTAI_API_KEY")), |api_key| { self.vastai_api_key = Some(api_key); } ); env_apply!("MYELIN_VASTAI_BOOTSTRAP_COMMAND", |command| { self.vastai_bootstrap_command = Some(command) }); env_apply!("MYELIN_VASTAI_SSH_IDENTITY", |identity| { self.vastai_ssh_identity_raw = Some(identity) }); env_apply!("MYELIN_VASTAI_DISK_GB", |disk_gb| { self.vastai_disk_gb_raw = Some(disk_gb) }); env_apply!("MYELIN_VASTAI_SSH_USER", |ssh_user| { self.vastai_ssh_user = Some(ssh_user) }); env_apply!("MYELIN_VASTAI_CONFIRM_LEASE", |confirm_lease| { self.vastai_confirm_lease_raw = Some(confirm_lease) }); env_apply!("MYELIN_VASTAI_ONSTART", |onstart| { self.vastai_onstart = Some(onstart) }); env_apply!("MYELIN_VASTAI_GPU_NAME", |gpu_name| { self.vastai_gpu_name = Some(gpu_name) }); env_apply!("MYELIN_VASTAI_MIN_GPU_RAM_MB", |min_gpu_ram_mb| { self.vastai_min_gpu_ram_mb_raw = Some(min_gpu_ram_mb) }); env_apply!("MYELIN_VASTAI_MIN_DOWN_MBPS", |min_down_mbps| { self.vastai_min_down_mbps_raw = Some(min_down_mbps) }); env_apply!("MYELIN_VASTAI_MAX_DPH_TOTAL", |max_dph_total| { self.vastai_max_dph_total_raw = Some(max_dph_total) }); env_apply!("MYELIN_VASTAI_MIN_UP_MBPS", |min_up_mbps| { self.vastai_min_up_mbps_raw = Some(min_up_mbps) }); env_apply!("MYELIN_VASTAI_MIN_RELIABILITY", |min_reliability| { self.vastai_min_reliability_raw = Some(min_reliability) }); env_apply!("MYELIN_VASTAI_REQUIRE_VERIFIED", |require_verified| { self.vastai_require_verified_raw = Some(require_verified) }); env_apply!("MYELIN_VASTAI_BLACKLIST_HOSTS", |blacklist_hosts| { for host_id in Self::parse_list("MYELIN_VASTAI_BLACKLIST_HOSTS", &blacklist_hosts)? { self.push_vastai_blacklist_host(host_id); } }); env_apply!("MYELIN_VASTAI_POLL_INTERVAL_SECS", |poll_interval_secs| { self.vastai_poll_interval_secs_raw = Some(poll_interval_secs) }); Ok(self) } fn overlay_cli(mut self, args: impl IntoIterator) -> Result { let mut args = args.into_iter(); while let Some(arg) = args.next() { match arg.as_str() { "--runtime-config" => { self.config_profile = RuntimeConfigProfile::parse(&next_arg(&mut args, "--runtime-config")?)? } "--provider" => { self.provider = Some(provider_kind::parse_deploy(&next_arg( &mut args, "--provider", )?)?) } "--worker-bin" => { self.worker_bin = Some(PathBuf::from(next_arg(&mut args, "--worker-bin")?)) } "--image" => { self.image = next_arg(&mut args, "--image")?; self.image_overridden_after_toml = true; } "--gpus" => self.docker_gpus = next_arg(&mut args, "--gpus")?, "--rpc-bind" => { self.rpc_bind = next_arg(&mut args, "--rpc-bind")?; self.rpc_bind_label = "--rpc-bind"; } "--run-id" => self.run_id = parse_next(&mut args, "--run-id")?, "--node-id" => self.node_id = parse_next(&mut args, "--node-id")?, "--stage-index" => self.stage_index = parse_next(&mut args, "--stage-index")?, "--layer-end-exclusive" => { self.layer_end_exclusive = Some(parse_next(&mut args, "--layer-end-exclusive")?) } "-N" | "--pipeline-stages" => self.pipeline_stages = parse_next(&mut args, &arg)?, "--max-tokens" => self.default_max_tokens = parse_next(&mut args, "--max-tokens")?, "--dashboard" => self.dashboard = true, "--no-dashboard" => self.dashboard = false, "--telemetry-frame-log" => { self.telemetry_frame_log = Some(PathBuf::from(next_arg( &mut args, "--telemetry-frame-log", )?)); } "--model-id" => self.model_id = next_arg(&mut args, "--model-id")?, "--gguf-local-path" => { self.gguf_source = GgufSource::LocalPath(next_arg(&mut args, "--gguf-local-path")?) } "--gguf-repo" => { self.set_hf_source(Some(next_arg(&mut args, "--gguf-repo")?), None, None) } "--gguf-file" => { self.set_hf_source(None, Some(next_arg(&mut args, "--gguf-file")?), None) } "--gguf-revision" => self.set_hf_source( None, None, Some(Some(next_arg(&mut args, "--gguf-revision")?)), ), "--tokenizer-local-path" => { self.tokenizer = TokenizerSource::LocalPath(next_arg(&mut args, "--tokenizer-local-path")?) } "--max-context" => self.max_context = Some(parse_next(&mut args, "--max-context")?), "--cached-model-host-path" => { self.cached_model_host_path = Some(PathBuf::from(next_arg( &mut args, "--cached-model-host-path", )?)); } "--relay-mode" => self.relay_mode = Some(next_arg(&mut args, "--relay-mode")?), "--relay-url" => self.relay_url = Some(next_arg(&mut args, "--relay-url")?), "--endpoint-addr-mask" => { self.endpoint_addr_mask = Some(next_arg(&mut args, "--endpoint-addr-mask")?) } "--vastai-disk-gb" => { self.vastai_disk_gb = Some(parse_next(&mut args, "--vastai-disk-gb")?); self.vastai_disk_gb_raw = None; } "--vastai-min-gpu-ram-mb" => { self.vastai_min_gpu_ram_mb = Some(parse_next(&mut args, "--vastai-min-gpu-ram-mb")?); self.vastai_min_gpu_ram_mb_raw = None; } "--vastai-min-down-mbps" => { self.vastai_min_down_mbps = Some(parse_next(&mut args, "--vastai-min-down-mbps")?); self.vastai_min_down_mbps_raw = None; } "--vastai-max-dph-total" => { self.vastai_max_dph_total = Some(parse_next(&mut args, "--vastai-max-dph-total")?); self.vastai_max_dph_total_raw = None; } "--vastai-min-up-mbps" => { self.vastai_min_up_mbps = Some(parse_next(&mut args, "--vastai-min-up-mbps")?); self.vastai_min_up_mbps_raw = None; } "--vastai-min-reliability" => { self.vastai_min_reliability = Some(parse_next(&mut args, "--vastai-min-reliability")?); self.vastai_min_reliability_raw = None; } "--vastai-blacklist-host" => { let host_id = parse_next(&mut args, "--vastai-blacklist-host")?; self.push_vastai_blacklist_host(host_id); } "--vastai-poll-interval-secs" => { self.vastai_poll_interval_secs = Some(parse_next(&mut args, "--vastai-poll-interval-secs")?); self.vastai_poll_interval_secs_raw = None; } "--vastai-api-key" => { self.vastai_api_key = Some(next_arg(&mut args, "--vastai-api-key")?) } "--vastai-bootstrap-command" => { self.vastai_bootstrap_command = Some(next_arg(&mut args, "--vastai-bootstrap-command")?) } "--vastai-ssh-identity" => { self.vastai_ssh_identity_raw = Some(next_arg(&mut args, "--vastai-ssh-identity")?); } "--vastai-ssh-user" => { self.vastai_ssh_user = Some(next_arg(&mut args, "--vastai-ssh-user")?) } "--vastai-onstart" => { self.vastai_onstart = Some(next_arg(&mut args, "--vastai-onstart")?) } "--vastai-gpu-name" => { self.vastai_gpu_name = Some(next_arg(&mut args, "--vastai-gpu-name")?) } "--vastai-confirm-lease" => { self.vastai_confirm_lease = Some(true); self.vastai_confirm_lease_raw = None; } "--no-vastai-confirm-lease" => { self.vastai_confirm_lease = Some(false); self.vastai_confirm_lease_raw = None; } "--vastai-require-verified" => { self.vastai_require_verified = Some(true); self.vastai_require_verified_raw = None; } "--no-vastai-require-verified" => { self.vastai_require_verified = Some(false); self.vastai_require_verified_raw = None; } _ => return Err(format!("unknown argument {arg:?}")), } } Ok(self) } fn finalize(self) -> Result { let provider = self .provider .clone() .unwrap_or_else(|| match self.config_profile { RuntimeConfigProfile::Local => provider_kind::process(), RuntimeConfigProfile::Deploy => provider_kind::vastai(), }); let provider_name = provider.as_str(); let mut image = self.image.clone(); if provider_name == "vastai" && !self.image_overridden_after_toml { if let Some(vastai_image) = &self.toml_vastai_image { image = vastai_image.clone(); } } if self.pipeline_stages == 0 { return Err("--pipeline-stages must be greater than 0".to_owned()); } let mut cached_model_host_path = self.cached_model_host_path.clone(); if matches!(provider_name, "process" | "docker") && self.pipeline_stages > 1 && cached_model_host_path.is_none() && (matches!( &self.gguf_source, GgufSource::HuggingFaceGguf { repo, file, revision: None, } if repo == DEFAULT_HF_REPO && file == DEFAULT_HF_FILE ) || matches!( &self.gguf_source, GgufSource::HuggingFaceGguf { file, revision: None, .. } if file == DEFAULT_PIPELINE_CACHED_MODEL_FILE )) { cached_model_host_path = Some(default_pipeline_cached_model_path()); } let cached_model = cached_model_host_path .map(|path| CachedModelConfig::from_host_path(provider_name, path)) .transpose()?; let mut gguf_source = self.gguf_source.clone(); if let Some(cached_model) = &cached_model { if provider_name != "vastai" { gguf_source = GgufSource::LocalPath(if provider_name == "process" { cached_model.host_path.to_string_lossy().to_string() } else { cached_model.container_path.clone() }); } } let relay = relay_runtime_config_from_settings( self.run_id, self.relay_mode.as_deref(), self.relay_url.as_deref(), )?; let endpoint_addr_mask = match self.endpoint_addr_mask.as_deref() { Some(mask) => EndpointAddrMask::parse(mask)?, None => EndpointAddrMask::Full, }; let vastai = (provider_name == "vastai") .then(|| VastAiRuntimeConfig::from_builder(&self)) .transpose()?; Ok(Config { config_profile: self.config_profile, image, docker_gpus: self.docker_gpus, provider, rpc_bind: self .rpc_bind .parse() .map_err(|e| format!("invalid {}: {e}", self.rpc_bind_label))?, run_id: self.run_id, node_id: self.node_id, stage_index: self.stage_index, layer_end_exclusive: self.layer_end_exclusive, pipeline_stages: self.pipeline_stages, model_id: self.model_id, gguf_source, tokenizer: self.tokenizer, default_max_tokens: self.default_max_tokens, dashboard: self.dashboard, max_context: self.max_context, relay, endpoint_addr_mask, vastai, cached_model, worker_bin: self.worker_bin, telemetry_frame_log: self.telemetry_frame_log, }) } fn set_hf_source( &mut self, repo: Option, file: Option, revision: Option>, ) { let (current_repo, current_file, current_revision) = match &self.gguf_source { GgufSource::HuggingFaceGguf { repo, file, revision, } => (repo.clone(), file.clone(), revision.clone()), GgufSource::LocalPath(_) => { (DEFAULT_HF_REPO.to_owned(), DEFAULT_HF_FILE.to_owned(), None) } }; self.gguf_source = GgufSource::HuggingFaceGguf { repo: repo.unwrap_or(current_repo), file: file.unwrap_or(current_file), revision: revision.unwrap_or(current_revision), }; } fn push_vastai_blacklist_host(&mut self, host_id: u64) { if !self.vastai_blacklist_hosts.contains(&host_id) { self.vastai_blacklist_hosts.push(host_id); } } fn parse_list(name: &str, value: &str) -> Result, String> where T: std::str::FromStr, T::Err: std::fmt::Display, { value .split(',') .map(str::trim) .filter(|part| !part.is_empty()) .map(|part| Self::parse_value(name, part)) .collect() } fn parse_value(name: &str, value: &str) -> Result where T: std::str::FromStr, T::Err: std::fmt::Display, { value .parse::() .map_err(|e| format!("invalid {name}={value:?}: {e}")) } fn parse_bool(name: &str, value: &str) -> Result { match value.to_ascii_lowercase().as_str() { "1" | "true" | "yes" | "on" => Ok(true), "0" | "false" | "no" | "off" => Ok(false), _ => Err(format!( "invalid {name}={value:?}; use 1/0, true/false, yes/no, or on/off" )), } } } impl Config { fn uses_planned_execution(&self) -> bool { self.cached_model.is_some() || self.pipeline_stages > 1 } fn provider_telemetry_detail(&self) -> Value { match self.provider.as_str() { "process" => json!({ "worker_bin": self.worker_bin.as_ref().map(|path| path.to_string_lossy().to_string()), "cached_model": self.cached_model.as_ref().map(CachedModelConfig::telemetry_detail), }), "docker" => json!({ "docker_gpus": &self.docker_gpus, "cached_model": self.cached_model.as_ref().map(CachedModelConfig::telemetry_detail), }), "vastai" => self .vastai .as_ref() .map_or_else(|| json!({}), VastAiRuntimeConfig::telemetry_detail), _ => json!({}), } } fn build_run_plan(&self) -> Result { let host_path = self.local_planning_gguf_path()?; let metadata = crate::staging::gguf_metadata::read_gguf_planning_metadata(&host_path)?; let model = metadata.to_model_facts( self.model_id.clone(), self.gguf_source.clone(), self.tokenizer.clone(), self.max_context, )?; if self.pipeline_stages > model.num_layers { return Err(format!( "--pipeline-stages={} exceeds GGUF layer count {}; choose N <= {}", self.pipeline_stages, model.num_layers, model.num_layers )); } let activation_extent = model .max_seq_len .checked_mul(model.hidden_dim) .and_then(|value| value.checked_mul(model.dtype_width_bytes)) .ok_or_else(|| "activation ring size overflow while planning pipeline".to_owned())?; let activation_data_capacity = run_plan::MO01_HEADER_BYTES .checked_add(activation_extent) .ok_or_else(|| { "activation ring data capacity overflow while planning pipeline".to_owned() })?; let token_extent = model .max_seq_len .checked_mul(4) .ok_or_else(|| "token ring size overflow while planning pipeline".to_owned())?; let token_data_capacity = run_plan::MO01_HEADER_BYTES .checked_add(token_extent) .ok_or_else(|| { "token ring data capacity overflow while planning pipeline".to_owned() })?; let candidate_pool = (0..self.pipeline_stages) .map(|stage_index| run_plan::NodeId(self.node_id + 1 + u64::from(stage_index))) .collect::>(); let placement = run_plan::PlacementInput::FixedLinear( candidate_pool .iter() .enumerate() .map(|(stage_index, node_id)| run_plan::StagePlacement { stage_index: stage_index as u32, node_id: *node_id, }) .collect(), ); run_plan::plan_run(run_plan::PlannerInput { run_id: run_plan::RunId(self.run_id), orchestrator_node_id: run_plan::NodeId(self.node_id), model, runtime: run_plan::RuntimeConfig { max_tokens: self.default_max_tokens, sampling: run_plan::SamplingPolicy { temperature_millis: 0, top_k: 1, }, }, candidate_pool, stage_count: self.pipeline_stages, placement, activation_ring: run_plan::RingSpec { data_capacity: activation_data_capacity, alignment: 64, direction: run_plan::RingDirection::Egress, host_pinning: run_plan::HostPinning::Pageable, wake_coalescing: run_plan::WakeCoalescing::PendingBit, }, token_ring: run_plan::RingSpec { data_capacity: token_data_capacity, alignment: 8, direction: run_plan::RingDirection::Egress, host_pinning: run_plan::HostPinning::Pageable, wake_coalescing: run_plan::WakeCoalescing::PendingBit, }, }) .map_err(|e| format!("plan pipeline run: {:?}", e.kind())) } fn local_planning_gguf_path(&self) -> Result { if let Some(cached_model) = &self.cached_model { return Ok(cached_model.host_path.clone()); } match &self.gguf_source { GgufSource::LocalPath(path) => { let host_path = PathBuf::from(path); if host_path.is_file() { Ok(host_path) } else { Err(format!( "local pipeline requires a locally inspectable GGUF before provisioning; {path:?} is not a host file, so use --cached-model-host-path" )) } } GgufSource::HuggingFaceGguf { repo, file, revision: None, } if self.provider.as_str() == "vastai" && file == DEFAULT_PIPELINE_CACHED_MODEL_FILE => { let host_path = default_pipeline_cached_model_path(); if host_path.is_file() { Ok(host_path) } else { Err(format!( "VastAI pipeline planning requires local GGUF metadata at {}; selected remote source is {repo}/{file}", host_path.display() )) } } GgufSource::HuggingFaceGguf { repo, file, .. } => Err(format!( "local pipeline requires a locally inspectable GGUF before provisioning; selected source {repo}/{file} is remote, so use --cached-model-host-path" )), } } fn prepare_vastai_ssh_key(&mut self) -> Result<(), String> { if self.provider.as_str() != "vastai" { return Ok(()); } let api_key = self .vastai .as_ref() .and_then(|vastai| vastai.api_key.as_deref()) .ok_or_else(|| { "VAST_API_KEY, MYELIN_VASTAI_API_KEY, or VASTAI_API_KEY is required when MYELIN_NODE_PROVIDER=vastai" .to_owned() })? .to_owned(); let identity = resolve_vastai_ssh_identity( self.vastai .as_ref() .and_then(|vastai| vastai.ssh_identity.clone()), )?; if !identity.is_file() { return Err(format!( "missing VastAI SSH identity {}; create/register one with vastai create ssh-key or set MYELIN_VASTAI_SSH_IDENTITY", identity.display() )); } let public_key = derive_ssh_public_key(&identity)?; let fingerprint = ssh_public_key_fingerprint(&public_key); ensure_vastai_account_ssh_key(&api_key, &public_key)?; eprintln!( "VastAI SSH identity {} fingerprint {} registered for account", identity.display(), fingerprint ); let vastai = self .vastai .as_mut() .expect("VastAI config exists when provider is vastai"); vastai.ssh_identity = Some(identity); vastai.provisioning.ssh_public_key = Some(public_key); vastai.ssh_public_fingerprint = Some(fingerprint); Ok(()) } fn build_provisioner( &self, bootstrap_runtime: swactor::runtime::Runtime, ) -> Result, String> { match self.provider.as_str() { "process" => { let worker_bin = match &self.worker_bin { Some(worker_bin) => worker_bin.clone(), None => { let mut path = std::env::current_exe().map_err(|e| format!("current exe: {e}"))?; path.set_file_name("myelin-worker"); path } }; if !worker_bin.is_file() { return Err(format!( "local process worker binary does not exist: {}", worker_bin.display() )); } Ok(Box::new(LocalProcessPlugin::new(worker_bin))) } "docker" => Ok(Box::new(LocalDockerPlugin::new( env_optional("MYELIN_DOCKER_CONTAINER_PREFIX") .unwrap_or_else(|| "myelin-orchestrator".to_owned()), ))), "vastai" => { let vastai = self.vastai.as_ref().ok_or_else(|| { "VastAI config was not resolved for provider vastai".to_owned() })?; if vastai.bootstrap_command.is_none() { return Err( "MYELIN_VASTAI_BOOTSTRAP_COMMAND is required when MYELIN_NODE_PROVIDER=vastai" .to_owned(), ); } let api_key = vastai.api_key.clone().ok_or_else(|| { "VAST_API_KEY, MYELIN_VASTAI_API_KEY, or VASTAI_API_KEY is required when MYELIN_NODE_PROVIDER=vastai" .to_owned() })?; let ssh_identity = vastai .ssh_identity .clone() .ok_or_else(|| "VastAI SSH identity was not prepared".to_owned())?; Ok(Box::new(VastAiProvisioningPlugin::new( ToolsVastAiLeaseClient::from_api_key(api_key)?, SshCommandBootstrapLauncher::new(Some(ssh_identity), bootstrap_runtime), vastai.provisioning.clone(), ))) } _ => Err("mock provider cannot build a runtime provisioner".to_owned()), } } fn node_spec_env_keys(&self) -> Vec { let mut keys: Vec = vec![ "MYELIN_RUN_ID", "MYELIN_LOGICAL_NODE_ID", "MYELIN_NODE_ATTEMPT_ID", "MYELIN_NODE_PROVIDER", "MYELIN_STAGE_INDEX", "MYELIN_COORDINATOR_ENDPOINT", "MYELIN_ORCHESTRATOR_ACTOR", "MYELIN_MODEL_ID", "MYELIN_IROH_RELAY_MODE", MVP_IROH_ENDPOINT_ADDR_MASK_ENV, "MYELIN_PIPELINE_STAGES", ] .into_iter() .map(str::to_owned) .collect(); keys.extend(self.extra_worker_env().into_iter().map(|(k, _)| k)); keys } /// Conditional worker env pairs shared by `node_spec_env_keys` and `node_spec_for_stage`. fn extra_worker_env(&self) -> Vec<(String, String)> { let provider_name = self.provider.as_str(); let mut env = Vec::new(); if let Some(url) = &self.relay.url { env.push((MYELIN_IROH_RELAY_URL_ENV.to_owned(), url.clone())); } if provider_name == "docker" { env.push(("MYELIN_DOCKER_GPUS".to_owned(), self.docker_gpus.clone())); } if let Some(value) = env_optional("DEV") { env.push(("DEV".to_owned(), value)); } env.extend(local_tinygrad_worker_env(provider_name)); for name in [ "MYELIN_CPU_LINE_PROFILE", "MYELIN_CPU_LINE_PROFILE_INTERVAL_MS", "MYELIN_TOKEN_PROGRESS_EVERY", "CUDA_DEVICE_SCHEDULE", "MYELIN_MODEL_CACHE_DIR", "HF_TOKEN", ] { if let Some(value) = env_optional(name) { env.push((name.to_owned(), value)); } } match &self.gguf_source { GgufSource::LocalPath(path) => { env.push(("MYELIN_GGUF_LOCAL_PATH".to_owned(), path.clone())) } GgufSource::HuggingFaceGguf { repo, file, revision, } => { env.push(("MYELIN_GGUF_REPO".to_owned(), repo.clone())); env.push(("MYELIN_GGUF_FILE".to_owned(), file.clone())); if let Some(revision) = revision { env.push(("MYELIN_GGUF_REVISION".to_owned(), revision.clone())); } } } if let TokenizerSource::LocalPath(path) = &self.tokenizer { env.push(("MYELIN_TOKENIZER_LOCAL_PATH".to_owned(), path.clone())); } if let Some(max_context) = self.max_context { env.push(("MYELIN_MAX_CONTEXT".to_owned(), max_context.to_string())); } env } fn node_spec_for_stage( &self, coordinator: EndpointAddr, orchestrator_actor: ActorAddress, logical_node_id: u64, stage_index: u32, ) -> Result { let provider_name = self.provider.as_str(); let mut env = vec![ ("MYELIN_RUN_ID".to_owned(), self.run_id.to_string()), ( "MYELIN_LOGICAL_NODE_ID".to_owned(), logical_node_id.to_string(), ), ("MYELIN_STAGE_INDEX".to_owned(), stage_index.to_string()), ( "MYELIN_PIPELINE_STAGES".to_owned(), self.pipeline_stages.to_string(), ), ( MVP_IROH_ENDPOINT_ADDR_MASK_ENV.to_owned(), self.endpoint_addr_mask.as_str().to_owned(), ), ("MYELIN_NODE_PROVIDER".to_owned(), provider_name.to_owned()), ( "MYELIN_COORDINATOR_ENDPOINT".to_owned(), serde_json::to_string(&coordinator) .map_err(|e| format!("serialize coordinator endpoint: {e}"))?, ), ( "MYELIN_ORCHESTRATOR_ACTOR".to_owned(), serde_json::to_string(&orchestrator_actor) .map_err(|e| format!("serialize orchestrator actor: {e}"))?, ), ("MYELIN_MODEL_ID".to_owned(), self.model_id.clone()), ( "MYELIN_IROH_RELAY_MODE".to_owned(), relay_mode_env_value(&self.relay.mode).to_owned(), ), ]; env.extend(self.extra_worker_env()); let args = match provider_name { "vastai" => self .vastai .as_ref() .and_then(|vastai| vastai.bootstrap_command.clone()) .into_iter() .collect(), "process" | "docker" => Vec::new(), _ => return Err("myelin-orchestrator does not support mock provider".to_owned()), }; let mounts = if provider_name == "docker" { self.cached_model .as_ref() .map(|cached_model| { vec![ProviderMount { host_path: cached_model.host_path.to_string_lossy().to_string(), container_path: cached_model.container_path.clone(), readonly: self.uses_planned_execution(), }] }) .unwrap_or_default() } else { Vec::new() }; Ok(NodeProvisionSpec { run_id: self.run_id, node_id: logical_node_id, attempt_id: 0, stage_index: Some(stage_index), image: self.image.clone(), env, args, mounts, }) } } #[derive(Clone)] struct RuntimeReady { endpoint: EndpointAddr, node_actor: ActorAddress, telemetry_publisher: ActorAddress, stage_index: u32, readiness_id: u64, swim_node_id: DistNodeId, } #[derive(Clone)] struct PromptRuntimeReady { first_stage: RuntimeReady, final_stage: RuntimeReady, } #[derive(Clone)] struct RuntimeReadyAckTarget { node_id: u64, ready: RuntimeReady, } fn runtime_ready_barrier_met(stack: &DistributionRuntimeStack, ready: &RuntimeReady) -> bool { stack.member_state(ready.swim_node_id) == Some(MemberState::Alive) && stack.route_owner(ready.node_actor) == Some(ready.swim_node_id) } fn enqueue_runtime_ready_ack( stack: &DistributionRuntimeStack, ready: &RuntimeReady, run_id: u64, node_id: u64, ) -> Result<(), String> { stack .runtime .send_to( ready.node_actor, NodeAgentMsg::RuntimeReadyAck { run_id, node_id, stage_index: ready.stage_index, readiness_id: ready.readiness_id, }, ) .map_err(|e| format!("send runtime ready ack: {e}")) } fn enqueue_telemetry_subscribe( stack: &DistributionRuntimeStack, ready: &RuntimeReady, collector: &EndpointAddr, run_id: u64, node_id: u64, ) -> Result<(), String> { let mut flow_id = [0_u8; 16]; flow_id[..8].copy_from_slice(&run_id.to_le_bytes()); flow_id[8..].copy_from_slice(&node_id.to_le_bytes()); stack .runtime .send_to( ready.telemetry_publisher, TelemetryPublisherMsg::Subscribe(TelemetrySubscribe { collector: collector.clone(), request: SubscriptionRequest::all(), flow_id, token: Vec::new(), }), ) .map_err(|e| format!("send telemetry subscribe: {e}")) } struct RuntimeReadyAckLoop<'a> { driver: &'a mut IrohDriver, stack: &'a DistributionRuntimeStack, obs_rx: &'a mpsc::Receiver, collector: &'a FrameCollector, orchestrator_reports: &'a swactor::runtime::Inbox, stop_rx: &'a mpsc::Receiver<()>, dashboard: Option<&'a DashboardSupport>, orch_telemetry: &'a mut OrchTelemetry, orch_stdio_rx: Option<&'a mpsc::Receiver>, run_id: u64, orchestrator_node_id: u64, provider: &'a ProviderKind, orchestrator_actor: ActorAddress, } // synchronous process-control/orchestration sequencing; the engine drives all background work (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn wait_for_runtime_ready_acks( ctx: RuntimeReadyAckLoop<'_>, targets: &[RuntimeReadyAckTarget], collector_endpoint: &EndpointAddr, cluster: &mut ProvisionedClusterGuard, ) -> Result { let RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id, orchestrator_node_id, provider, .. } = ctx; let bootstrap = |ds: &mut OrchTelemetry, phase: &str, status: &str, detail: Value| { ds.emit_bootstrap( dashboard, run_id, orchestrator_node_id, phase, status, detail, ); }; let mut pending = targets .iter() .cloned() .map(|target| { ( ( target.node_id, target.ready.stage_index, target.ready.readiness_id, ), target, ) }) .collect::>(); let mut attempts = BTreeMap::<(u64, u32, u64), u64>::new(); let mut last_send = None::; while !pending.is_empty() { cluster .poll(SystemTime::now()) .map_err(|error| format!("cluster reconcile while awaiting ready ack: {error}"))?; if targets.iter().any(|target| { cluster.current_attempt(target.node_id) != Some(::provisioning::NodeAttemptId(target.ready.readiness_id)) }) { return Ok(false); } collector.pump(driver); drain_orch_stdio_capture( orch_stdio_rx, orch_telemetry, dashboard, run_id, orchestrator_node_id, ); if stop_requested(stop_rx) { return Err( "shutdown requested while waiting for runtime-ready acknowledgements".to_owned(), ); } while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, provider, &observation); } collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); while let Some(report) = orchestrator_reports.try_recv() { let OrchestratorReport::NodeRuntimeReadyAck { run_id: ack_run_id, node_id, stage_index, readiness_id, } = report else { continue; }; if ack_run_id != run_id { continue; } let key = (node_id, stage_index, readiness_id); let Some(target) = pending.remove(&key) else { continue; }; bootstrap( orch_telemetry, "runtime_ready_ack", "ready", json!({ "node_id":node_id, "node_actor":target.ready.node_actor, "readiness_id":readiness_id, "attempts":attempts.get(&key).copied().unwrap_or(0), }), ); } if pending.is_empty() { return Ok(true); } if last_send.is_none_or(|sent_at| sent_at.elapsed() >= RUNTIME_READY_ACK_RETRY_INTERVAL) { for (key, target) in &pending { if stack.route_owner(target.ready.telemetry_publisher) == Some(target.ready.swim_node_id) && let Err(error) = enqueue_telemetry_subscribe( stack, &target.ready, collector_endpoint, run_id, target.node_id, ) { bootstrap( orch_telemetry, "telemetry_subscribe", "failed", json!({"node_id":target.node_id,"error":error}), ); } enqueue_runtime_ready_ack(stack, &target.ready, run_id, target.node_id)?; let attempt = attempts.entry(*key).or_default(); *attempt += 1; let attempt = *attempt; bootstrap( orch_telemetry, "runtime_ready_ack", "sent", json!({ "node_id":target.node_id, "node_actor":target.ready.node_actor, "readiness_id":target.ready.readiness_id, "attempt":attempt, }), ); } last_send = Some(Instant::now()); } thread::sleep(PUMP_INTERVAL); } Ok(true) } fn reconciler_group(config: &Config, spec: &NodeProvisionSpec) -> RunNodeGroupSpec { let group_id = NodeGroupId(format!("node-{}", spec.node_id)); let ssh_user = config .vastai .as_ref() .map(|vastai| vastai.provisioning.ssh_user.clone()) .unwrap_or_else(|| "root".to_owned()); let disk_gb = config .vastai .as_ref() .map(|vastai| vastai.provisioning.disk_gb) .unwrap_or_default(); let orchestrator = spec .env .iter() .find(|(name, _)| name == "MYELIN_ORCHESTRATOR_ACTOR") .map(|(_, value)| value.clone()) .unwrap_or_default(); RunNodeGroupSpec { run_id: ClusterRunId(config.run_id), group_id, role: RoleId(format!( "stage-{}", spec.stage_index.unwrap_or(config.stage_index) )), count: 1, provider: ReconcilerProviderKind::new(config.provider.as_str()), shape: DesiredNodeShape { image: spec.image.clone(), disk_gb, gpu_name: None, min_gpu_ram_mb: None, min_down_mbps: None, min_up_mbps: None, min_reliability: None, require_verified: false, provider_labels: BTreeMap::from([( "myelin.provider_config".to_owned(), config.provider_telemetry_detail().to_string(), )]), }, boot: BootSpec { ssh_user, verify_commands: Vec::new(), start_swactor_command: spec.args.join(" "), stdout_sources: Vec::new(), stderr_sources: Vec::new(), env: spec.env.clone(), args: spec.args.clone(), mounts: spec.mounts.clone(), }, swarm_join: SwarmJoinTemplate { orch_swactor_addr: orchestrator, join_token_ref: "myelin-runtime-ready".to_owned(), }, } } fn build_reconciled_cluster( provisioner: Box, config: &Config, stage_specs: &[NodeProvisionSpec], runtime: swactor::runtime::Runtime, engine: EngineHandle, sink: PluginSink, ) -> Result { let groups = stage_specs .iter() .map(|spec| reconciler_group(config, spec)) .collect::>(); let desired = ClusterShape { run_id: ClusterRunId(config.run_id), generation: 1, groups, }; let expanded = desired.expand().map_err(|error| error.to_string())?; let mut first = Some(provisioner); let mut bindings = Vec::with_capacity(stage_specs.len()); for spec in stage_specs { let logical_node_id = ReconcilerLogicalNodeId(format!("node-{}-0", spec.node_id)); if !expanded.contains_key(&logical_node_id) { return Err(format!( "reconciler shape did not expand node {}", logical_node_id.0 )); } let plugin = match first.take() { Some(plugin) => plugin, None => config.build_provisioner(runtime.clone())?, }; bindings.push(ReconcilerNodeBinding { logical_node_id, provision: spec.clone(), plugin, }); } ProvisionedClusterGuard::new(desired, bindings, RetryPolicy::default(), engine, sink) } // Synchronous orchestration sequencing; provider work and timers are engine-hosted. #[allow(clippy::disallowed_methods)] fn start_and_provision_workers( provisioner: Box, config: &Config, pipeline_plan: Option<&run_plan::RunPlan>, ctx: RuntimeReadyAckLoop<'_>, sink: PluginSink, coordinator: EndpointAddr, pipeline_coordinator: EndpointAddr, orchestrator_actor: ActorAddress, ) -> Result< ( ProvisionedClusterGuard, PromptRuntimeReady, BTreeMap, ), String, > { let RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, .. } = ctx; let run_id = config.run_id; let node_id = config.node_id; let bootstrap = |ds: &mut OrchTelemetry, phase: &str, status: &str, detail: Value| { ds.emit_bootstrap(dashboard, run_id, node_id, phase, status, detail); }; let stage_specs = stage_node_specs( config, pipeline_plan, coordinator.clone(), orchestrator_actor, )?; let expected_node_ids = stage_specs .iter() .map(|spec| spec.node_id) .collect::>(); bootstrap( orch_telemetry, "node_spec", "ready", json!({ "provider":config.provider.as_str(), "image":&config.image, "relay_mode":relay_mode_env_value(&config.relay.mode), "endpoint_addr_mask":config.endpoint_addr_mask.as_str(), "docker_gpus":if config.provider.as_str() == "docker" { Some(config.docker_gpus.as_str()) } else { None }, "provider_config":config.provider_telemetry_detail(), "env_keys":config.node_spec_env_keys(), "worker_count":stage_specs.len(), "worker_node_ids":expected_node_ids, }), ); let stage_shard_plans = if let Some(plan) = pipeline_plan { let plans = pipeline_stage_shard_plans(config, plan)?; emit_stage_shard_plan_summaries( dashboard, orch_telemetry, config.run_id, config.node_id, &plans, ); plans } else { BTreeMap::new() }; for node_spec in &stage_specs { orch_telemetry.emit_event( dashboard, ProvisionEvent { run_id: config.run_id, node_id: node_spec.node_id, kind: ProvisionEventKind::ProvisionStart, provider: Some(config.provider.as_str().to_owned()), message: Some(format!( "reconciling {} image {}", config.provider.as_str(), config.image )), }, ); bootstrap( orch_telemetry, "provider_start", "started", json!({ "provider":config.provider.as_str(), "image":&config.image, "node_id":node_spec.node_id, "stage_index":node_spec.stage_index, "generation":1, }), ); } let mut provisioned_nodes = build_reconciled_cluster( provisioner, config, &stage_specs, stack.runtime.clone(), stack.engine.clone(), sink, )?; loop { provisioned_nodes .poll(SystemTime::now()) .map_err(|error| format!("cluster reconcile: {error}"))?; if provisioned_nodes.awaiting_runtime() || provisioned_nodes.is_converged() { break; } collector.pump(driver); collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); drain_orch_stdio_capture( orch_stdio_rx, orch_telemetry, dashboard, config.run_id, config.node_id, ); while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, &config.provider, &observation); } if stop_requested(stop_rx) { return Err("shutdown requested while reconciling worker nodes".to_owned()); } thread::sleep(PUMP_INTERVAL); } for node_spec in &stage_specs { bootstrap( orch_telemetry, "provider_start", "ready", json!({ "provider":config.provider.as_str(), "node_id":node_spec.node_id, "stage_index":node_spec.stage_index, "attempt":provisioned_nodes.current_attempt(node_spec.node_id).map(|attempt| attempt.0), }), ); } drain_orch_stdio_capture( orch_stdio_rx, orch_telemetry, dashboard, config.run_id, config.node_id, ); bootstrap( orch_telemetry, "node_runtime_ready", "started", json!({"worker_count":expected_node_ids.len(),"node_ids":expected_node_ids}), ); let readies = loop { let readies = match wait_for_runtime_readies( RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, &expected_node_ids, &mut provisioned_nodes, ) { Ok(readies) => readies, Err(error) => { bootstrap( orch_telemetry, "node_runtime_ready", "failed", json!({"error":error}), ); return Err(error); } }; for (node_id, ready) in &readies { bootstrap( orch_telemetry, "node_runtime_ready", "ready", json!({"endpoint":&ready.endpoint,"node_actor":ready.node_actor,"node_id":node_id,"stage_index":ready.stage_index,"attempt":ready.readiness_id}), ); } let ack_targets = readies .iter() .map(|(node_id, ready)| RuntimeReadyAckTarget { node_id: *node_id, ready: ready.clone(), }) .collect::>(); if wait_for_runtime_ready_acks( RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, &ack_targets, &pipeline_coordinator, &mut provisioned_nodes, )? { break readies; } bootstrap( orch_telemetry, "node_runtime_ready", "retry", json!({"reason":"node attempt changed before ready acknowledgement"}), ); }; for (node_id, ready) in &readies { let attempt = ::provisioning::NodeAttemptId(ready.readiness_id); if !provisioned_nodes.observe_runtime_ready( *node_id, attempt, SwactorId(format!("{:?}", ready.node_actor)), SystemTime::now(), ) { return Err(format!( "stale runtime-ready observation for node {node_id} attempt {}", attempt.0 )); } } while !provisioned_nodes.is_converged() { provisioned_nodes .poll(SystemTime::now()) .map_err(|error| format!("cluster convergence: {error}"))?; collector.pump(driver); collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, &config.provider, &observation); } if stop_requested(stop_rx) { return Err("shutdown requested while converging worker nodes".to_owned()); } thread::sleep(PUMP_INTERVAL); } bootstrap( orch_telemetry, "stage_provision", "started", stage_provision_detail(config, pipeline_plan), ); if pipeline_plan.is_none() { let ready = readies .get(&config.node_id) .ok_or_else(|| "missing runtime-ready node for single-stage run".to_owned())?; provision_stage(stack, ready.node_actor, config)?; } bootstrap( orch_telemetry, "weights_loaded", "started", json!({"model_id":&config.model_id,"expected":expected_node_ids.len()}), ); let weights_result = if pipeline_plan.is_some() { wait_for_weights_loaded_count( RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, expected_node_ids.len(), pipeline_plan.expect("pipeline mode requires plan"), &readies, &pipeline_coordinator, &stage_shard_plans, ) } else { wait_for_weights_loaded( RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id: config.run_id, orchestrator_node_id: config.node_id, provider: &config.provider, orchestrator_actor, }, config.stage_index, ) }; match weights_result { Ok(()) => bootstrap( orch_telemetry, "weights_loaded", "ready", json!({"source":"actor_stage_ready","model_id":&config.model_id,"expected":expected_node_ids.len()}), ), Err(error) => { bootstrap( orch_telemetry, "weights_loaded", "failed", json!({"error":error}), ); drain_orch_stdio_capture( orch_stdio_rx, orch_telemetry, dashboard, config.run_id, config.node_id, ); return Err(error); } } let prompt_ready = if let Some(plan) = pipeline_plan { let first_node_id = plan .stages .iter() .find(|stage| stage.stage_index == 0) .map(|stage| stage.node_id.0) .ok_or_else(|| "pipeline plan missing stage 0".to_owned())?; let final_node_id = plan .stages .iter() .find(|stage| stage.stage_index + 1 == stage.stage_count) .map(|stage| stage.node_id.0) .ok_or_else(|| "pipeline plan missing final stage".to_owned())?; PromptRuntimeReady { first_stage: readies .get(&first_node_id) .cloned() .ok_or_else(|| "missing stage 0 runtime-ready node".to_owned())?, final_stage: readies .get(&final_node_id) .cloned() .ok_or_else(|| "missing final-stage runtime-ready node".to_owned())?, } } else { let ready = readies .get(&config.node_id) .cloned() .ok_or_else(|| "missing single-stage runtime-ready node".to_owned())?; PromptRuntimeReady { first_stage: ready.clone(), final_stage: ready, } }; let swim_to_node = readies .iter() .map(|(node_id, ready)| (ready.swim_node_id, *node_id)) .collect::>(); Ok((provisioned_nodes, prompt_ready, swim_to_node)) } fn stage_node_specs( config: &Config, pipeline_plan: Option<&run_plan::RunPlan>, coordinator: EndpointAddr, orchestrator_actor: ActorAddress, ) -> Result, String> { if let Some(plan) = pipeline_plan { let mut stages = plan.stages.clone(); stages.sort_by_key(|stage| stage.stage_index); stages .into_iter() .map(|stage| { config.node_spec_for_stage( coordinator.clone(), orchestrator_actor, stage.node_id.0, stage.stage_index, ) }) .collect() } else { Ok(vec![config.node_spec_for_stage( coordinator, orchestrator_actor, config.node_id, config.stage_index, )?]) } } fn stage_provision_detail(config: &Config, pipeline_plan: Option<&run_plan::RunPlan>) -> Value { if let Some(plan) = pipeline_plan { json!({ "run_id":config.run_id, "stage_count":plan.stages.len(), "stages":plan.stages.iter().map(|stage| { json!({ "node_id":stage.node_id.0, "stage_index":stage.stage_index, "layer_range":{"start":stage.layer_start,"end_exclusive":stage.layer_end_exclusive}, "inbound_edge_id":stage.inbound_edge.0, "outbound_edge_id":stage.outbound_edge.0, }) }).collect::>(), "model_id":&config.model_id, }) } else { json!({ "run_id":config.run_id, "node_id":config.node_id, "stage_index":config.stage_index, "stage_count":1, "layer_range":{"start":0,"end_exclusive":config.layer_end_exclusive}, "model_id":&config.model_id, }) } } fn stage_provision_wire_from_plan( plan: &run_plan::RunPlan, stage_index: u32, readies: &BTreeMap, coordinator: &EndpointAddr, stage_shard_plans: &BTreeMap, ) -> Result { let provision = run_plan::derive_stage_provision(plan, stage_index) .map_err(|e| format!("derive stage {stage_index} provision: {e:?}"))?; Ok(StageProvisionWire { run_id: provision.run_id.0, authorized_orchestrator: 0, node_id: provision.node_id.0, stage_index: provision.stage_index, stage_count: provision.stage_count, layer_start: provision.layer_start, layer_end_exclusive: provision.layer_end_exclusive, inbound_edge_id: provision.inbound.edge_id.0, outbound_edge_id: provision.outbound.edge_id.0, inbound_edge: Some(StageInboundEdgeWire { edge_id: provision.inbound.edge_id.0, kind: stage_edge_kind_wire(provision.inbound.kind), object_spec: stage_object_spec_wire(provision.inbound.object_spec), ring_spec: stage_ring_spec_wire(provision.inbound.ring_spec), }), outbound_edge: Some(StageOutboundEdgeWire { edge_id: provision.outbound.edge_id.0, kind: stage_edge_kind_wire(provision.outbound.kind), consumer_node_id: provision.outbound.consumer_node_id.0, consumer_endpoint: stage_consumer_endpoint( provision.outbound.consumer_node_id.0, plan, readies, coordinator, )?, object_spec: stage_object_spec_wire(provision.outbound.object_spec), ring_spec: stage_ring_spec_wire(provision.outbound.ring_spec), }), model_id: provision.model.model_id, gguf_source: provision.gguf_source, tokenizer: provision.tokenizer, stage_shard_plan: stage_shard_plans.get(&stage_index).cloned(), }) } fn provision_stage_from_plan( stack: &DistributionRuntimeStack, node_actor: ActorAddress, plan: &run_plan::RunPlan, stage_index: u32, readies: &BTreeMap, coordinator: &EndpointAddr, stage_shard_plans: &BTreeMap, ) -> Result<(), String> { let provision = stage_provision_wire_from_plan(plan, stage_index, readies, coordinator, stage_shard_plans)?; stack .runtime .send_to(node_actor, NodeAgentMsg::ProvisionStage(provision)) .map_err(|e| format!("send stage {stage_index} provision: {e}")) } fn pipeline_stage_shard_plans( config: &Config, plan: &run_plan::RunPlan, ) -> Result, String> { if !matches!(config.gguf_source, GgufSource::HuggingFaceGguf { .. }) { return Ok(BTreeMap::new()); } let planning_gguf = config.local_planning_gguf_path()?; let mut out = BTreeMap::new(); for stage in &plan.stages { let shard_plan = plan_stage_shard( &planning_gguf, stage.gguf_source.clone(), stage.stage_index, stage.stage_count, stage.layer_start, stage.layer_end_exclusive, ) .map_err(|error| { format!( "plan stage {} HF shard ranges from {}: {error}", stage.stage_index, planning_gguf.display() ) })?; out.insert(stage.stage_index, shard_plan); } Ok(out) } fn stage_shard_plan_summary_detail(plan: &StageShardPlan) -> Value { let planned_fetch_bytes = plan.planned_fetch_bytes(); json!({ "stage_index":plan.stage_index, "stage_count":plan.stage_count, "layer_start":plan.layer_start, "layer_end_exclusive":plan.layer_end_exclusive, "planned_fetch_bytes":planned_fetch_bytes, "source_total_bytes":plan.source_total_bytes, "tensor_count":plan.tensors.len(), "range_count":plan.planned_range_count(), "tensor_range_count":plan.merged_tensor_ranges.len(), "metadata_bytes":plan.metadata_end, "planned_fraction":if plan.source_total_bytes == 0 { Value::Null } else { json!(planned_fetch_bytes as f64 / plan.source_total_bytes as f64) }, }) } fn emit_stage_shard_plan_summaries( dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, stage_shard_plans: &BTreeMap, ) { for plan in stage_shard_plans.values() { orch_telemetry.emit_bootstrap( dashboard, run_id, node_id, "stage_shard_plan", "ready", stage_shard_plan_summary_detail(plan), ); } } fn stage_consumer_endpoint( consumer_node_id: u64, plan: &run_plan::RunPlan, readies: &BTreeMap, coordinator: &EndpointAddr, ) -> Result, String> { if consumer_node_id == plan.stages.first().map_or(0, |stage| { plan.edges .iter() .find(|edge| edge.kind == run_plan::EdgeKind::TokenIn) .and_then(|edge| match edge.producer { run_plan::EdgeEndpoint::Orchestrator { node_id } => Some(node_id.0), run_plan::EdgeEndpoint::Stage { .. } => None, }) .unwrap_or(stage.node_id.0) }) { return Ok(Some(coordinator.clone())); } readies .get(&consumer_node_id) .map(|ready| Some(ready.endpoint.clone())) .ok_or_else(|| { format!("missing runtime-ready endpoint for consumer node {consumer_node_id}") }) } fn stage_edge_kind_wire(kind: run_plan::EdgeKind) -> StageEdgeKindWire { match kind { run_plan::EdgeKind::TokenIn => StageEdgeKindWire::TokenIn, run_plan::EdgeKind::Activation => StageEdgeKindWire::Activation, run_plan::EdgeKind::TokenOut => StageEdgeKindWire::TokenOut, } } fn stage_object_spec_wire(spec: run_plan::ObjectSpec) -> StageObjectSpecWire { StageObjectSpecWire { max_extent: spec.max_extent, alignment: spec.alignment, } } fn stage_ring_spec_wire(spec: run_plan::RingSpec) -> StageRingSpecWire { StageRingSpecWire { data_capacity: spec.data_capacity, alignment: spec.alignment, } } // synchronous process-control/orchestration sequencing; the engine drives all background work (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn wait_for_runtime_readies( ctx: RuntimeReadyAckLoop<'_>, expected_node_ids: &[u64], cluster: &mut ProvisionedClusterGuard, ) -> Result, String> { let RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id, provider, .. } = ctx; let expected = expected_node_ids.iter().copied().collect::>(); let mut pending = BTreeMap::::new(); loop { cluster .poll(SystemTime::now()) .map_err(|error| format!("cluster reconcile while awaiting runtime: {error}"))?; pending.retain(|node_id, ready| { cluster.current_attempt(*node_id) == Some(::provisioning::NodeAttemptId(ready.readiness_id)) }); collector.pump(driver); emit_swim_transitions( orch_telemetry, dashboard, run_id, expected_node_ids.first().copied().unwrap_or(0), stack, ); emit_swim_probe_events(orch_telemetry, dashboard, stack, "runtime_ready_wait"); collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); drain_orch_stdio_capture( orch_stdio_rx, orch_telemetry, dashboard, run_id, expected_node_ids.first().copied().unwrap_or(0), ); if stop_requested(stop_rx) { return Err("shutdown requested while waiting for pipeline nodes ready".to_owned()); } while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, provider, &observation); match observation { PluginObservation::TelemetryFrame { .. } | PluginObservation::ProviderLine { .. } | PluginObservation::StdoutLine { .. } | PluginObservation::StderrLine { .. } | PluginObservation::Failed { .. } | PluginObservation::Exited { .. } => {} } } while let Some(report) = orchestrator_reports.try_recv() { if let OrchestratorReport::NodeRuntimeReady { run_id: report_run_id, node_id, stage_index, endpoint, node_actor, telemetry_publisher, readiness_id, } = report && report_run_id == run_id && expected.contains(&node_id) && cluster.current_attempt(node_id) == Some(::provisioning::NodeAttemptId(readiness_id)) { pending.insert( node_id, RuntimeReady { endpoint: endpoint.clone(), node_actor, telemetry_publisher, stage_index, readiness_id, swim_node_id: DistNodeId(*endpoint.id.as_bytes()), }, ); } } if expected.iter().all(|node_id| { pending .get(node_id) .is_some_and(|ready| runtime_ready_barrier_met(stack, ready)) }) { return Ok(pending); } thread::sleep(PUMP_INTERVAL); } } // synchronous process-control/orchestration sequencing; the engine drives all background work (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn wait_for_weights_loaded_count( ctx: RuntimeReadyAckLoop<'_>, expected_count: usize, pipeline_plan: &run_plan::RunPlan, readies: &BTreeMap, pipeline_coordinator: &EndpointAddr, stage_shard_plans: &BTreeMap, ) -> Result<(), String> { let RuntimeReadyAckLoop { driver, stack, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id, orchestrator_node_id: node_id, provider, .. } = ctx; let expected_stages = pipeline_plan .stages .iter() .map(|stage| stage.stage_index) .collect::>(); let mut loaded_stages = BTreeSet::::new(); let mut last_resend = Instant::now() - Duration::from_secs(15); let mut resend_attempt = 0_u64; let mut stage_resend_counts = BTreeMap::::new(); let mut stage_last_sends = BTreeMap::::new(); let mut load_progress = BTreeMap::::new(); loop { collector.pump(driver); emit_swim_transitions(orch_telemetry, dashboard, run_id, node_id, stack); emit_swim_probe_events(orch_telemetry, dashboard, stack, "weights_loaded_wait"); drain_orch_stdio_capture(orch_stdio_rx, orch_telemetry, dashboard, run_id, node_id); if stop_requested(stop_rx) { return Err("shutdown requested while waiting for pipeline weights loaded".to_owned()); } if loaded_stages.len() >= expected_count { return Ok(()); } if last_resend.elapsed() >= Duration::from_secs(15) { resend_attempt += 1; let mut pending: Vec<&run_plan::StagePlan> = pipeline_plan .stages .iter() .filter(|stage| !loaded_stages.contains(&stage.stage_index)) .collect(); pending.sort_by_key(|stage| stage.stage_index); if pending.is_empty() { return Err(format!( "missing unloaded pipeline weight stage; loaded {} of {expected_count}", loaded_stages.len() )); } for stage in pending { let mut provision = PipelineStageProvision { driver: &mut *driver, stack, collector, dashboard, orch_telemetry: &mut *orch_telemetry, run_id, node_id, pipeline_plan, readies, pipeline_coordinator, stage_shard_plans, loaded_stages: &loaded_stages, stage_resend_counts: &mut stage_resend_counts, stage_last_sends: &mut stage_last_sends, load_progress: &load_progress, }; send_pipeline_stage_provision(&mut provision, stage, resend_attempt)?; } last_resend = Instant::now(); } while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, provider, &observation); match observation { PluginObservation::Failed { reason, node_id: failed_node_id, .. } => { orch_telemetry.emit_bootstrap( dashboard, run_id, node_id, "stage_provision_wait", "failed", json!({ "classification":"worker_process_failed", "stage_node_id":failed_node_id, "last_load_progress":load_progress.get(&failed_node_id).map(StageLoadProgress::to_json), "reason":reason, }), ); return Err(reason); } PluginObservation::Exited { node_id: exited_node_id, status, .. } => { let reason = format!("node {exited_node_id} exited while loading weights: {status:?}"); orch_telemetry.emit_bootstrap( dashboard, run_id, node_id, "stage_provision_wait", "failed", json!({ "classification":"worker_process_exited", "stage_node_id":exited_node_id, "last_load_progress":load_progress.get(&exited_node_id).map(StageLoadProgress::to_json), "status":status, "reason":reason, }), ); return Err(reason); } PluginObservation::TelemetryFrame { .. } | PluginObservation::ProviderLine { .. } | PluginObservation::StdoutLine { .. } | PluginObservation::StderrLine { .. } => {} } } collector.drain_with_progress(&mut load_progress, |stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); while let Some(report) = orchestrator_reports.try_recv() { match report { OrchestratorReport::WeightsReady { run_id: report_run_id, node_id: _, stage_index, } if report_run_id == run_id && expected_stages.contains(&stage_index) => { loaded_stages.insert(stage_index); if loaded_stages.len() >= expected_count { return Ok(()); } } OrchestratorReport::StageFault { run_id: report_run_id, stage_index, reason, } if report_run_id == run_id && expected_stages.contains(&stage_index) => { let mut error = format!("stage {stage_index} faulted while loading pipeline weights"); if let Some(reason) = reason { error.push_str(": "); error.push_str(&reason); } return Err(error); } _ => {} } } thread::sleep(PUMP_INTERVAL); } } struct PipelineStageProvision<'a> { driver: &'a mut IrohDriver, stack: &'a DistributionRuntimeStack, collector: &'a FrameCollector, dashboard: Option<&'a DashboardSupport>, orch_telemetry: &'a mut OrchTelemetry, run_id: u64, node_id: u64, pipeline_plan: &'a run_plan::RunPlan, readies: &'a BTreeMap, pipeline_coordinator: &'a EndpointAddr, stage_shard_plans: &'a BTreeMap, loaded_stages: &'a BTreeSet, stage_resend_counts: &'a mut BTreeMap, stage_last_sends: &'a mut BTreeMap, load_progress: &'a BTreeMap, } fn send_pipeline_stage_provision( ctx: &mut PipelineStageProvision<'_>, stage: &run_plan::StagePlan, attempt: u64, ) -> Result<(), String> { let stage_node_id = stage.node_id.0; let current_send_count = ctx .stage_resend_counts .get(&stage.stage_index) .copied() .unwrap_or_default(); let now = Instant::now(); let ready = ctx .readies .get(&stage_node_id) .ok_or_else(|| format!("missing runtime-ready node for stage {}", stage.stage_index))?; let route_owner = ctx.stack.route_owner(ready.node_actor); let telemetry_route_owner = ctx.stack.route_owner(ready.telemetry_publisher); let member_state = ctx.stack.member_state(ready.swim_node_id); let route_matches_ready = route_owner == Some(ready.swim_node_id); let (dashboard, run_id, node_id) = (ctx.dashboard, ctx.run_id, ctx.node_id); let bootstrap = |ds: &mut OrchTelemetry, phase: &str, status: &str, detail: Value| { ds.emit_bootstrap(dashboard, run_id, node_id, phase, status, detail); }; ctx.orch_telemetry.emit_bootstrap_to_channel( ctx.dashboard, MYELIN_STAGE_ROUTE, ctx.run_id, ctx.node_id, "stage_route_check", "observed", json!({ "attempt":attempt, "stage_index":stage.stage_index, "stage_node_id":stage_node_id, "node_actor":ready.node_actor, "telemetry_publisher":ready.telemetry_publisher, "swim_node_id":format!("{:?}", ready.swim_node_id), "member_state":member_state.map(|state| format!("{:?}", state)), "route_owner":route_owner.map(|owner| format!("{:?}", owner)), "telemetry_route_owner":telemetry_route_owner.map(|owner| format!("{:?}", owner)), "route_matches_ready":route_matches_ready, }), ); if member_state == Some(MemberState::Dead) { let reason = format!( "stage {} node {} is dead while loading pipeline weights", stage.stage_index, stage_node_id ); let liveness = stage_load_liveness_detail( ctx.load_progress.get(&stage_node_id), stage.stage_index, stage_node_id, member_state, route_owner, telemetry_route_owner, route_matches_ready, "heartbeat_missed", ); bootstrap( ctx.orch_telemetry, "stage_provision_wait", "failed", json!({ "attempt":attempt, "stage_count":ctx.pipeline_plan.stages.len(), "stage_index":stage.stage_index, "stage_node_id":stage_node_id, "stage_send_count":current_send_count, "loaded_stage_count":ctx.loaded_stages.len(), "member_state":"Dead", "route_owner":route_owner.map(|owner| format!("{:?}", owner)), "telemetry_route_owner":telemetry_route_owner.map(|owner| format!("{:?}", owner)), "classification":"heartbeat_missed", "liveness":liveness, "reason":reason, }), ); return Err(reason); } let (should_send, dispatch_reason) = stage_provision_dispatch( ctx.load_progress.get(&stage_node_id), current_send_count, ctx.stage_last_sends.get(&stage.stage_index).copied(), now, ); if !should_send { bootstrap( ctx.orch_telemetry, "stage_provision_wait", "observed", json!({ "attempt":attempt, "stage_count":ctx.pipeline_plan.stages.len(), "stage_index":stage.stage_index, "stage_node_id":stage_node_id, "stage_send_count":current_send_count, "loaded_stage_count":ctx.loaded_stages.len(), "resend_suppressed":true, "resend_reason":dispatch_reason, "liveness":stage_load_liveness_detail( ctx.load_progress.get(&stage_node_id), stage.stage_index, stage_node_id, member_state, route_owner, telemetry_route_owner, route_matches_ready, "waiting", ), "message":format!( "loaded {} of {}; waiting on stage {}", ctx.loaded_stages.len(), ctx.pipeline_plan.stages.len(), stage.stage_index ) }), ); return Ok(()); } let stage_send_count = { let count = ctx .stage_resend_counts .entry(stage.stage_index) .or_default(); *count += 1; *count }; ctx.stage_last_sends.insert(stage.stage_index, now); bootstrap( ctx.orch_telemetry, "stage_provision_send", "sent", json!({ "attempt":attempt, "stage_count":ctx.pipeline_plan.stages.len(), "stage_index":stage.stage_index, "stage_send_count":stage_send_count, "loaded_stage_count":ctx.loaded_stages.len(), "parallel_weight_acquisition":true, "resend_reason":dispatch_reason, }), ); if stage_send_count == 1 || stage_send_count % 15 == 0 { bootstrap( ctx.orch_telemetry, "stage_provision_wait", "observed", json!({ "attempt":attempt, "stage_count":ctx.pipeline_plan.stages.len(), "stage_index":stage.stage_index, "stage_node_id":stage_node_id, "stage_send_count":stage_send_count, "loaded_stage_count":ctx.loaded_stages.len(), "liveness":stage_load_liveness_detail( ctx.load_progress.get(&stage_node_id), stage.stage_index, stage_node_id, member_state, route_owner, telemetry_route_owner, route_matches_ready, "waiting", ), "message":format!( "loaded {} of {}; waiting on stage {}", ctx.loaded_stages.len(), ctx.pipeline_plan.stages.len(), stage.stage_index ) }), ); } provision_stage_from_plan( ctx.stack, ready.node_actor, ctx.pipeline_plan, stage.stage_index, ctx.readies, ctx.pipeline_coordinator, ctx.stage_shard_plans, )?; ctx.collector.pump(ctx.driver); Ok(()) } struct PromptWork { request: SubmitPrompt, events: mpsc::Sender, } struct ActivePrompt { request: SubmitPrompt, events: mpsc::Sender, } fn stage_load_phase_is_active(phase: Option<&str>) -> bool { matches!( phase, Some( "loading_weights" | "prefetching_model" | "prefetching_stage_shard" | "fetching_stage_shard" | "stage_shard_cache_ready" | "stage_shard_ready" | "cache_ready" | "constructing_stage" | "stage_constructed" | "building_tokenizer" | "tokenizer_ready" ) ) } fn stage_provision_dispatch( progress: Option<&StageLoadProgress>, send_count: u64, last_send: Option, now: Instant, ) -> (bool, &'static str) { if send_count == 0 { return (true, "initial"); } let Some(progress) = progress else { return (true, "no_progress_after_send"); }; if progress.failure_reason.is_some() || progress.phase.as_deref() == Some("failed") { return (false, "worker_load_failed"); } if progress.phase.as_deref() == Some("weights_loaded") { return (false, "weights_loaded_report_pending"); } if !stage_load_phase_is_active(progress.phase.as_deref()) { return (true, "unknown_or_inactive_progress"); } let Some(last_progress) = progress.last_progress else { return (true, "active_phase_without_progress_time"); }; if now.duration_since(last_progress) < STAGE_PROVISION_ACTIVE_RESEND_AFTER { return (false, "active_progress"); } if let Some(last_send) = last_send && now.duration_since(last_send) < STAGE_PROVISION_ACTIVE_RESEND_AFTER { return (false, "recent_stale_progress_resend"); } (true, "stale_progress") } fn stage_load_liveness_detail( progress: Option<&StageLoadProgress>, stage_index: u32, stage_node_id: u64, member_state: Option, route_owner: Option, telemetry_route_owner: Option, route_matches_ready: bool, classification: &str, ) -> Value { json!({ "classification": classification, "stage_index": stage_index, "stage_node_id": stage_node_id, "member_state": member_state.map(|state| format!("{:?}", state)), "route_owner": route_owner.map(|owner| format!("{:?}", owner)), "telemetry_route_owner": telemetry_route_owner.map(|owner| format!("{:?}", owner)), "route_matches_ready": route_matches_ready, "load_progress": progress.map(StageLoadProgress::to_json), "host_gpu_missing": progress.is_none_or(|progress| progress.host_gpu_samples == 0), }) } struct OrchStdioCapture; struct OrchStdioLine { stream: ProvisionLogStream, line: String, } #[cfg(target_os = "linux")] impl OrchStdioCapture { fn install() -> Result>, String> { let stdout_read = Self::redirect_stream(libc::STDOUT_FILENO, "stdout")?; let stderr_read = Self::redirect_stream(libc::STDERR_FILENO, "stderr")?; let (tx, rx) = mpsc::channel(); Self::spawn_reader(stdout_read, ProvisionLogStream::Stdout, tx.clone()); Self::spawn_reader(stderr_read, ProvisionLogStream::Stderr, tx); Ok(Some(rx)) } fn redirect_stream(fd: libc::c_int, name: &str) -> Result { let mut pipe_fds = [0; 2]; let pipe_result = unsafe { libc::pipe(pipe_fds.as_mut_ptr()) }; if pipe_result != 0 { return Err(format!( "create orchestrator {name} capture pipe: {}", std::io::Error::last_os_error() )); } let dup_result = unsafe { libc::dup2(pipe_fds[1], fd) }; let close_write_result = unsafe { libc::close(pipe_fds[1]) }; if dup_result < 0 { let error = std::io::Error::last_os_error(); let _ = unsafe { libc::close(pipe_fds[0]) }; return Err(format!("redirect orchestrator {name}: {error}")); } if close_write_result != 0 { let error = std::io::Error::last_os_error(); let _ = unsafe { libc::close(pipe_fds[0]) }; return Err(format!("close orchestrator {name} duplicate fd: {error}")); } Ok(unsafe { File::from_raw_fd(pipe_fds[0]) }) } // provider log capture is out of scope (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn spawn_reader(file: File, stream: ProvisionLogStream, tx: mpsc::Sender) { thread::spawn(move || { let reader = BufReader::new(file); for line in reader.lines() { let Ok(line) = line else { break; }; if tx.send(OrchStdioLine { stream, line }).is_err() { break; } } }); } } #[cfg(not(target_os = "linux"))] impl OrchStdioCapture { fn install() -> Result>, String> { Ok(None) } } fn drain_orch_stdio_capture( rx: Option<&mpsc::Receiver>, telemetry: &mut OrchTelemetry, dashboard: Option<&DashboardSupport>, run_id: u64, node_id: u64, ) { let Some(rx) = rx else { return; }; while let Ok(line) = rx.try_recv() { telemetry.emit_log( dashboard, ProvisionLogLine { run_id, node_id, stream: line.stream, line: line.line, }, ); } } struct ChannelObservationSink { tx: Mutex>, } impl PluginObservationSink for ChannelObservationSink { fn observe(&self, observation: PluginObservation) { let _ = self.tx.lock().send(observation); } } fn spawn_prompt_rpc( engine: &EngineHandle, bind: SocketAddr, work_tx: mpsc::Sender, default_max_tokens: u32, ) -> Result { // Bind synchronously (there is no ambient runtime at orchestrator startup) // and report the bound address, then drive accept on the orchestrator // engine. Each accepted connection runs on the engine's blocking pool, // reusing the synchronous request/response parser unchanged. There is no // listener thread and no per-connection std thread (ENGINE_SPEC.md); // no raw Tokio handle or second runtime is introduced. let std_listener = TcpListener::bind(bind).map_err(|e| format!("bind prompt RPC {bind}: {e}"))?; let addr = std_listener .local_addr() .map_err(|e| format!("read prompt RPC addr: {e}"))?; std_listener .set_nonblocking(true) .map_err(|e| format!("set prompt RPC nonblocking: {e}"))?; let engine = engine.clone(); engine.clone().spawn(async move { let listener = match tokio::net::TcpListener::from_std(std_listener) { Ok(listener) => listener, Err(_) => return, }; loop { match listener.accept().await { Ok((stream, _)) => { let tx = work_tx.clone(); let engine = engine.clone(); engine.spawn_blocking(move || { if let Ok(std_stream) = stream.into_std() { // The synchronous parser uses blocking I/O. let _ = std_stream.set_nonblocking(false); let _ = handle_prompt_connection(std_stream, tx, default_max_tokens); } }); } Err(_) => break, } } }); Ok(addr) } fn handle_prompt_connection( stream: TcpStream, work_tx: mpsc::Sender, default_max_tokens: u32, ) -> Result<(), String> { let mut reader = BufReader::new( stream .try_clone() .map_err(|e| format!("clone prompt stream: {e}"))?, ); let mut writer = stream; loop { let request = match read_submit_prompt(&mut reader) { Ok(Some(request)) => request, Ok(None) => break, Err(error) if error.contains("expected value at line 1 column 1") => break, Err(error) => return Err(error), }; let request = request.with_defaults(default_max_tokens); let (event_tx, event_rx) = mpsc::channel(); work_tx .send(PromptWork { request, events: event_tx, }) .map_err(|_| "prompt loop stopped".to_owned())?; for event in event_rx { let terminal = event.is_terminal(); write_json_line(&mut writer, &event)?; if terminal { break; } } } Ok(()) } fn provision_stage( stack: &DistributionRuntimeStack, node_actor: ActorAddress, config: &Config, ) -> Result<(), String> { stack .runtime .send_to( node_actor, NodeAgentMsg::ProvisionStage({ let layer_end_exclusive = config.layer_end_exclusive.ok_or_else(|| { "unplanned single-stage execution requires --layer-end-exclusive or MYELIN_LAYER_END_EXCLUSIVE; cached-model runs use metadata-derived planning".to_owned() })?; StageProvisionWire { run_id: config.run_id, authorized_orchestrator: 0, node_id: config.node_id, stage_index: config.stage_index, stage_count: 1, layer_start: 0, layer_end_exclusive, inbound_edge_id: 1, outbound_edge_id: 2, inbound_edge: None, outbound_edge: None, model_id: config.model_id.clone(), gguf_source: config.gguf_source.clone(), tokenizer: config.tokenizer.clone(), stage_shard_plan: None, } }), ) .map_err(|e| format!("send stage provision: {e}")) } // synchronous process-control/orchestration sequencing; the engine drives all background work (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn wait_for_weights_loaded(ctx: RuntimeReadyAckLoop<'_>, stage_index: u32) -> Result<(), String> { let RuntimeReadyAckLoop { driver, stack: _, obs_rx, collector, orchestrator_reports, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id, orchestrator_node_id: node_id, provider, .. } = ctx; loop { collector.pump(driver); drain_orch_stdio_capture(orch_stdio_rx, orch_telemetry, dashboard, run_id, node_id); if stop_requested(stop_rx) { return Err("shutdown requested while waiting for weights loaded".to_owned()); } drain_observations_with_exit(obs_rx, dashboard, orch_telemetry, provider, |_, status| { format!("node exited while loading weights: {status:?}") })?; collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); while let Some(report) = orchestrator_reports.try_recv() { match report { OrchestratorReport::WeightsReady { run_id: report_run_id, node_id: _, stage_index: report_stage_index, } if report_run_id == run_id && report_stage_index == stage_index => { return Ok(()); } OrchestratorReport::StageFault { run_id: report_run_id, stage_index: report_stage_index, reason, } if report_run_id == run_id && report_stage_index == stage_index => { let mut error = format!("stage {report_stage_index} faulted while loading weights"); if let Some(reason) = reason { error.push_str(": "); error.push_str(&reason); } return Err(error); } _ => {} } } thread::sleep(PUMP_INTERVAL); } } #[derive(Debug)] struct PipelineTokenRecord { object_id: u64, sequence: u64, token_id: u32, eos: bool, } type PipelineSendHandle = EdgeSendHandle; struct PendingEncode { request_id: u64, } struct PendingDecode { request_id: u64, token_id: u32, eos: bool, reached_limit: bool, } struct PipelinePromptRuntime { token_in_edge_id: u64, token_out_edge_id: u64, token_spec: run_plan::ObjectSpec, token_out_spec: run_plan::ObjectSpec, token_in_sender: PipelineSendHandle, recv_rx: mpsc::Receiver>, recv_tx: mpsc::Sender>, recv_buffer: Vec, tokenizer_encode_actor: ActorAddress, tokenizer_decode_actor: ActorAddress, tokenizer_reply_to: ActorAddress, pending_encode: Option, pending_decode: Option, next_sequence: u64, generated_tokens: Vec, final_text: String, active: Option, started_at: Option, last_progress_at: Option, next_wait_log_at: Option, } impl PipelinePromptRuntime { fn new( driver: &IrohDriver, plan: &run_plan::RunPlan, first_stage_endpoint: EndpointAddr, tokenizer_encode_actor: ActorAddress, tokenizer_decode_actor: ActorAddress, tokenizer_reply_to: ActorAddress, ) -> Result { let token_in_edge = plan .edges .iter() .find(|edge| edge.kind == run_plan::EdgeKind::TokenIn) .ok_or_else(|| "pipeline plan missing token-in edge".to_owned())?; let token_out_edge = plan .edges .iter() .find(|edge| edge.kind == run_plan::EdgeKind::TokenOut) .ok_or_else(|| "pipeline plan missing token-out edge".to_owned())?; let (recv_tx, recv_rx) = mpsc::channel(); Ok(Self { token_in_edge_id: token_in_edge.edge_id.0, token_out_edge_id: token_out_edge.edge_id.0, token_spec: token_in_edge.object_spec, token_out_spec: token_out_edge.object_spec, token_in_sender: driver .spawn_edge_send_pump(first_stage_endpoint, token_in_edge.edge_id.0)?, recv_rx, recv_tx, tokenizer_encode_actor, tokenizer_decode_actor, tokenizer_reply_to, pending_encode: None, pending_decode: None, recv_buffer: Vec::new(), next_sequence: 0, generated_tokens: Vec::new(), final_text: String::new(), active: None, started_at: None, last_progress_at: None, next_wait_log_at: None, }) } fn note_progress(&mut self) { let now = Instant::now(); self.last_progress_at = Some(now); self.next_wait_log_at = now.checked_add(PIPELINE_PROMPT_WAIT_LOG_INTERVAL); } fn start_prompt( &mut self, request: SubmitPrompt, events: mpsc::Sender, runtime: &swactor::runtime::Runtime, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) -> Result<(), String> { let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; let request_id = request.request_id; if self.active.is_some() || self.pending_encode.is_some() || self.pending_decode.is_some() { let active_request_id = self.active.as_ref().map(|active| active.request.request_id); emit_prompt_evt( orch_telemetry, request_id, "pipeline_prompt_busy", "failed", json!({ "active_request_id":active_request_id, "pending_encode":self.pending_encode.is_some(), "pending_decode":self.pending_decode.is_some(), }), ); let _ = events.send(PromptEvent::Fault { request_id, error: "pipeline prompt runtime is busy".to_owned(), }); return Ok(()); } self.generated_tokens.clear(); self.final_text.clear(); self.recv_buffer.clear(); self.pending_decode = None; self.pending_encode = Some(PendingEncode { request_id }); self.started_at = Some(Instant::now()); self.note_progress(); emit_prompt_evt( orch_telemetry, request_id, "pipeline_tokenizer_encode", "started", json!({"node_actor":self.tokenizer_encode_actor,"reply_to":self.tokenizer_reply_to,"prompt_bytes":request.prompt_text.len()}), ); runtime .send_to( self.tokenizer_encode_actor, NodeAgentMsg::EncodePrompt { request_id, prompt: request.prompt_text.clone(), reply_to: self.tokenizer_reply_to, }, ) .map_err(|e| format!("send tokenizer encode request: {e}"))?; self.active = Some(ActivePrompt { request, events }); Ok(()) } fn emit_wait_progress( &mut self, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) { let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; let Some(active) = self.active.as_ref() else { return; }; let now = Instant::now(); let request_id = active.request.request_id; let elapsed_ms = self .started_at .map(|started| duration_ms_u64(now.saturating_duration_since(started))) .unwrap_or(0); let idle_ms = self .last_progress_at .map(|last| duration_ms_u64(now.saturating_duration_since(last))) .unwrap_or(elapsed_ms); if self.next_wait_log_at.is_some_and(|next| now >= next) { emit_prompt_evt( orch_telemetry, request_id, "pipeline_prompt_wait", "waiting", json!({ "elapsed_ms":elapsed_ms, "idle_ms":idle_ms, "pending_encode":self.pending_encode.is_some(), "pending_decode":self.pending_decode.is_some(), "generated_tokens":self.generated_tokens.len(), "next_sequence":self.next_sequence, }), ); self.next_wait_log_at = now.checked_add(PIPELINE_PROMPT_WAIT_LOG_INTERVAL); } } fn drain_tokenizer_events( &mut self, runtime: &swactor::runtime::Runtime, tokenizer_events: &swactor::runtime::Inbox, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) -> Result<(), String> { while let Some(event) = tokenizer_events.try_recv() { match event { TokenizerEvent::PromptEncoded { request_id, tokens } => self .handle_encoded_prompt( request_id, tokens, dashboard, orch_telemetry, run_id, node_id, )?, TokenizerEvent::TokensDecoded { request_id, text } => self.handle_decoded_tokens( runtime, request_id, text, dashboard, orch_telemetry, run_id, node_id, )?, TokenizerEvent::Fault { request_id, error } => { self.fault_active(request_id, format!("tokenizer request failed: {error}")); } } } Ok(()) } fn handle_encoded_prompt( &mut self, request_id: u64, tokens: Vec, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) -> Result<(), String> { let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; let Some(pending) = self.pending_encode.take() else { return Ok(()); }; if pending.request_id != request_id { self.pending_encode = Some(pending); return Ok(()); } let Some(active) = self.active.as_ref() else { return Ok(()); }; if active.request.request_id != request_id { return Ok(()); } emit_prompt_evt( orch_telemetry, request_id, "pipeline_tokenizer_encode", "ready", json!({"node_actor":self.tokenizer_encode_actor,"reply_to":self.tokenizer_reply_to,"tokens":tokens.len()}), ); let sequence = self.next_sequence; emit_prompt_evt( orch_telemetry, request_id, "pipeline_token_in", "started", json!({"edge_id":self.token_in_edge_id,"sequence":sequence,"tokens":tokens.len(),"begin_sequence":true,"token_count":tokens.len(),"token_ids":&tokens}), ); self.send_token_in(sequence, &tokens, true)?; self.note_progress(); emit_prompt_evt( orch_telemetry, request_id, "pipeline_token_in", "ready", json!({"edge_id":self.token_in_edge_id,"sequence":sequence,"begin_sequence":true,"token_count":tokens.len()}), ); Ok(()) } #[allow(clippy::too_many_arguments)] fn handle_decoded_tokens( &mut self, _runtime: &swactor::runtime::Runtime, request_id: u64, text: String, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) -> Result<(), String> { let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; let Some(pending) = self.pending_decode.take() else { return Ok(()); }; if pending.request_id != request_id { self.pending_decode = Some(pending); return Ok(()); } let Some(active) = self.active.as_ref() else { return Ok(()); }; if active.request.request_id != request_id { return Ok(()); } let events = active.events.clone(); emit_prompt_evt( orch_telemetry, request_id, "pipeline_tokenizer_decode", "ready", json!({"node_actor":self.tokenizer_decode_actor,"reply_to":self.tokenizer_reply_to,"text_bytes":text.len()}), ); self.note_progress(); self.final_text.push_str(&text); if !text.is_empty() { let _ = events.send(PromptEvent::TextDelta { request_id, text }); } if pending.eos || pending.reached_limit { let elapsed_ms = self .started_at .map(|started| started.elapsed().as_millis() as u64) .unwrap_or(0); let final_text = self.final_text.clone(); let tokens_generated = self.generated_tokens.len() as u32; emit_prompt_evt( orch_telemetry, request_id, "prompt_complete", "ready", json!({ "event":"Done", "terminal":true, "tokens_generated":tokens_generated, "elapsed_ms":elapsed_ms, "final_text_bytes":final_text.len(), }), ); let _ = events.send(PromptEvent::Done { request_id, final_text, tokens_generated, elapsed_ms, }); self.last_progress_at = None; self.started_at = None; self.next_wait_log_at = None; self.active = None; return Ok(()); } let sequence = self.next_sequence; emit_prompt_evt( orch_telemetry, request_id, "pipeline_token_in", "started", json!({"edge_id":self.token_in_edge_id,"sequence":sequence,"tokens":1,"begin_sequence":false,"token_count":1,"token_id":pending.token_id}), ); self.send_token_in(sequence, &[pending.token_id], false)?; self.note_progress(); emit_prompt_evt( orch_telemetry, request_id, "pipeline_token_in", "ready", json!({"edge_id":self.token_in_edge_id,"sequence":sequence,"begin_sequence":false,"token_count":1,"token_id":pending.token_id}), ); Ok(()) } fn send_token_in( &self, sequence: u64, tokens: &[u32], begin_sequence: bool, ) -> Result<(), String> { let payload = tokens .iter() .flat_map(|token| token.to_le_bytes()) .collect::>(); let mut flags = ingress::ObjectFlags::default(); flags.begin_sequence = begin_sequence; self.token_in_sender.send( ingress::ObjectRecordBuilder::new(ingress_object_spec_from_plan(self.token_spec)) .object_id(ingress::ObjectId(9000_u64.saturating_add(sequence))) .sequence(sequence) .payload(payload) .flags(flags) .encode(), ) } fn request_decode( &mut self, runtime: &swactor::runtime::Runtime, request_id: u64, token_id: u32, eos: bool, reached_limit: bool, ) -> Result<(), String> { self.pending_decode = Some(PendingDecode { request_id, token_id, eos, reached_limit, }); runtime .send_to( self.tokenizer_decode_actor, NodeAgentMsg::DecodeTokens { request_id, tokens: vec![token_id], reply_to: self.tokenizer_reply_to, }, ) .map_err(|e| format!("send tokenizer decode request: {e}")) } fn fault_active(&mut self, request_id: u64, error: String) { let should_fault = self .active .as_ref() .is_some_and(|active| active.request.request_id == request_id); if !should_fault { return; } if let Some(active) = self.active.take() { let _ = active.events.send(PromptEvent::Fault { request_id, error }); } self.pending_encode = None; self.pending_decode = None; self.started_at = None; self.last_progress_at = None; self.next_wait_log_at = None; } fn poll_driver(&mut self, driver: &mut IrohDriver) { for event in driver.drain_edge_events() { match event { WireEvent::BytesRead { edge_id, bytes, .. } if edge_id.0 == self.token_out_edge_id => { self.note_progress(); let _ = self.recv_tx.send(bytes); } WireEvent::StreamFault { edge_id: Some(edge_id), reason, .. } if edge_id.0 == self.token_out_edge_id => { if let Some(request_id) = self.active.as_ref().map(|active| active.request.request_id) { self.fault_active( request_id, format!("pipeline token-out stream fault: {reason:?}"), ); } } _ => {} } } } fn drain_tokens( &mut self, runtime: &swactor::runtime::Runtime, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, run_id: u64, node_id: u64, ) -> Result<(), String> { let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; if self.pending_decode.is_some() { return Ok(()); } while let Ok(bytes) = self.recv_rx.try_recv() { self.recv_buffer.extend_from_slice(&bytes); while self.pending_decode.is_none() && let Some(record) = take_pipeline_token_record(&mut self.recv_buffer, self.token_out_spec)? { if record.sequence != self.next_sequence { return Err(format!( "pipeline token sequence violation: expected {}, got {}", self.next_sequence, record.sequence )); } self.next_sequence = self.next_sequence.saturating_add(1); let Some(active) = self.active.as_ref() else { continue; }; let request_id = active.request.request_id; emit_prompt_evt( orch_telemetry, request_id, "pipeline_token_out", "observed", json!({"edge_id":self.token_out_edge_id,"object_id":record.object_id,"sequence":record.sequence,"token_id":record.token_id,"eos":record.eos,"generated_index":self.generated_tokens.len() + 1}), ); self.generated_tokens.push(record.token_id); let reached_limit = self.generated_tokens.len() as u32 >= active.request.max_tokens; emit_prompt_evt( orch_telemetry, request_id, "pipeline_tokenizer_decode", "started", json!({"node_actor":self.tokenizer_decode_actor,"reply_to":self.tokenizer_reply_to,"token_id":record.token_id,"sequence":record.sequence,"generated_index":self.generated_tokens.len()}), ); self.request_decode( runtime, request_id, record.token_id, record.eos, reached_limit, )?; } } Ok(()) } } fn ingress_object_spec_from_plan(spec: run_plan::ObjectSpec) -> ingress::ObjectSpec { let extent_alignment = match spec.kind { run_plan::ObjectKind::Token => u64::from(spec.dtype_width_bytes), _ => u64::from(spec.alignment), }; ingress::ObjectSpec { max_extent: spec.max_extent, alignment: extent_alignment, layout: ingress::ObjectLayout::Token, } } fn take_pipeline_token_record( buffer: &mut Vec, spec: run_plan::ObjectSpec, ) -> Result, String> { let record = match ingress::read_object_record(buffer, ingress_object_spec_from_plan(spec), false) .map_err(|reason| format!("invalid token-out record: {reason:?}"))? { ingress::ObjectRecordRead::Incomplete => return Ok(None), ingress::ObjectRecordRead::Complete(record) => record, }; let payload = record .payload(buffer) .ok_or_else(|| "token-out record payload missing".to_owned())?; if payload.len() != 4 { return Err(format!( "token-out payload must be exactly one u32, got {}", payload.len() )); } let token_id = u32::from_le_bytes(payload.try_into().unwrap()); let out = PipelineTokenRecord { object_id: record.object_id.0, sequence: record.sequence, token_id, eos: record.flags.end_of_sequence, }; buffer.drain(..record.total_len); Ok(Some(out)) } // synchronous process-control/orchestration sequencing; the engine drives all background work (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn serve_prompts( ctx: RuntimeReadyAckLoop<'_>, work_rx: &mpsc::Receiver, prompt_events: &swactor::runtime::Inbox, node_actor: ActorAddress, reply_to: ActorAddress, tokenizer_events: &swactor::runtime::Inbox, tokenizer_encode_actor: ActorAddress, tokenizer_decode_actor: ActorAddress, tokenizer_reply_to: ActorAddress, pipeline_plan: Option<&run_plan::RunPlan>, prompt_endpoint: EndpointAddr, swim_to_node: &BTreeMap, ) -> Result<(), String> { let RuntimeReadyAckLoop { driver, stack, obs_rx, collector, stop_rx, dashboard, orch_telemetry, orch_stdio_rx, run_id, orchestrator_node_id: node_id, provider, orchestrator_actor, .. } = ctx; let emit_prompt_evt = |ds: &mut OrchTelemetry, request_id: u64, phase: &str, status: &str, detail: Value| { ds.emit_prompt( dashboard, run_id, node_id, request_id, phase, status, detail, ) }; let mut pipeline_runtime = match pipeline_plan { Some(plan) => Some(PipelinePromptRuntime::new( driver, plan, prompt_endpoint, tokenizer_encode_actor, tokenizer_decode_actor, tokenizer_reply_to, )?), None => None, }; let mut active: Option = None; loop { collector.pump(driver); orch_telemetry.flush(dashboard, "orchestrator"); if let Some(pipeline) = pipeline_runtime.as_mut() { pipeline.poll_driver(driver); pipeline.drain_tokenizer_events( &stack.runtime, tokenizer_events, dashboard, orch_telemetry, run_id, node_id, )?; pipeline.drain_tokens(&stack.runtime, dashboard, orch_telemetry, run_id, node_id)?; pipeline.emit_wait_progress(dashboard, orch_telemetry, run_id, node_id); } drain_observations_with_exit( obs_rx, dashboard, orch_telemetry, &provider, |_, status| format!("node exited: {status:?}"), )?; collector.drain(|stream, channel, frame| { if let Some(d) = dashboard { d.publish_frame(stream, channel, frame); } orch_telemetry.archive_frame("node", stream, channel, frame); }); drain_orch_stdio_capture(orch_stdio_rx, orch_telemetry, dashboard, run_id, node_id); let swim_transitions = emit_swim_transitions(orch_telemetry, dashboard, run_id, node_id, stack); for transition in &swim_transitions { if transition.to == MemberState::Dead && let Some(&lost_node_id) = swim_to_node.get(&transition.peer) { let _ = stack.runtime.send_to( orchestrator_actor, OrchestratorMsg::ObserveMembershipLost { run_id, node_id: lost_node_id, }, ); } } if stop_rx.try_recv().is_ok() { orch_telemetry.emit_bootstrap( dashboard, run_id, node_id, "shutdown", "started", json!({"source":"stdin"}), ); let _ = stack.runtime.send_to( orchestrator_actor, OrchestratorMsg::ObserveOperatorStop { run_id }, ); collector.pump(driver); return Ok(()); } if active.is_none() && pipeline_runtime .as_ref() .is_none_or(|pipeline| pipeline.active.is_none()) && let Ok(work) = work_rx.try_recv() { let request = work.request; let request_id = request.request_id; emit_prompt_evt( orch_telemetry, request_id, "prompt_work", "observed", json!({ "prompt_bytes":request.prompt_text.len(), "max_tokens":request.max_tokens, }), ); if let Some(pipeline) = pipeline_runtime.as_mut() { pipeline.start_prompt( request, work.events, &stack.runtime, dashboard, orch_telemetry, run_id, node_id, )?; continue; } emit_prompt_evt( orch_telemetry, request_id, "node_prompt_send", "started", json!({"node_actor":node_actor,"reply_to":reply_to}), ); match stack.runtime.send_to( node_actor, NodeAgentMsg::InferPrompt { request_id, prompt: request.prompt_text.clone(), max_tokens: request.max_tokens, reply_to, }, ) { Ok(()) => { emit_prompt_evt( orch_telemetry, request_id, "node_prompt_send", "ready", json!({"node_actor":node_actor,"reply_to":reply_to}), ); active = Some(ActivePrompt { request, events: work.events, }); } Err(error) => { emit_prompt_evt( orch_telemetry, request_id, "node_prompt_send", "failed", json!({"node_actor":node_actor,"reply_to":reply_to,"error":error.to_string()}), ); return Err(format!("send prompt request: {error}")); } } } while let Some(event) = prompt_events.try_recv() { let request_id = event.request_id(); let Some(current) = active.as_ref() else { emit_prompt_evt( orch_telemetry, request_id, "node_prompt_event", "dropped", json!({"reason":"no_active_prompt","event":prompt_event_name(&event)}), ); continue; }; if request_id != current.request.request_id { emit_prompt_evt( orch_telemetry, request_id, "node_prompt_event", "dropped", json!({ "reason":"request_mismatch", "event":prompt_event_name(&event), "active_request_id":current.request.request_id, }), ); continue; } let completion = match &event { PromptEvent::Done { .. } => Some(("ready", json!({"event":"Done"}))), PromptEvent::Fault { error, .. } => { Some(("failed", json!({"event":"Fault","error":error}))) } PromptEvent::TextDelta { .. } => None, }; emit_prompt_evt( orch_telemetry, request_id, "node_prompt_event", "observed", prompt_event_detail(&event), ); let _ = current.events.send(event); if let Some((status, detail)) = completion { emit_prompt_evt( orch_telemetry, request_id, "prompt_complete", status, detail, ); active = None; } } thread::sleep(PUMP_INTERVAL); } } fn prompt_event_name(event: &PromptEvent) -> &'static str { match event { PromptEvent::TextDelta { .. } => "TextDelta", PromptEvent::Done { .. } => "Done", PromptEvent::Fault { .. } => "Fault", } } fn prompt_event_detail(event: &PromptEvent) -> Value { match event { PromptEvent::TextDelta { text, .. } => { json!({"event":"TextDelta","terminal":false,"text_bytes":text.len()}) } PromptEvent::Done { final_text, tokens_generated, elapsed_ms, .. } => json!({ "event":"Done", "terminal":true, "tokens_generated":tokens_generated, "elapsed_ms":elapsed_ms, "final_text_bytes":final_text.len(), }), PromptEvent::Fault { error, .. } => { json!({"event":"Fault","terminal":true,"error":error}) } } } fn stop_requested(stop_rx: &mpsc::Receiver<()>) -> bool { stop_rx.try_recv().is_ok() } // top-level OS signal handling is process control, out of scope (ENGINE_SPEC.md §2) #[allow(clippy::disallowed_methods)] fn spawn_stop_listener() -> mpsc::Receiver<()> { let (tx, rx) = mpsc::channel(); #[cfg(target_os = "linux")] { thread::spawn(move || { let Ok(mut signals) = signal_hook::iterator::Signals::new([ signal_hook::consts::signal::SIGINT, signal_hook::consts::signal::SIGTERM, ]) else { return; }; if signals.forever().next().is_some() { let _ = tx.send(()); } }); } #[cfg(not(target_os = "linux"))] { drop(tx); } rx } fn drain_observations_with_exit( obs_rx: &mpsc::Receiver, dashboard: Option<&DashboardSupport>, orch_telemetry: &mut OrchTelemetry, provider: &ProviderKind, exit_message: impl Fn(u64, Option) -> String, ) -> Result<(), String> { while let Ok(observation) = obs_rx.try_recv() { emit_plugin_observation(orch_telemetry, dashboard, provider, &observation); match observation { PluginObservation::Failed { reason, .. } => return Err(reason), PluginObservation::Exited { node_id, status, .. } => return Err(exit_message(node_id, status)), PluginObservation::TelemetryFrame { .. } | PluginObservation::ProviderLine { .. } | PluginObservation::StdoutLine { .. } | PluginObservation::StderrLine { .. } => {} } } Ok(()) } fn emit_plugin_observation( orch_telemetry: &mut OrchTelemetry, dashboard: Option<&DashboardSupport>, provider: &ProviderKind, observation: &PluginObservation, ) { match observation { PluginObservation::StdoutLine { run_id, node_id, line, } => orch_telemetry.emit_log( dashboard, ProvisionLogLine { run_id: *run_id, node_id: *node_id, stream: ProvisionLogStream::Stdout, line: line.clone(), }, ), PluginObservation::StderrLine { run_id, node_id, line, } => orch_telemetry.emit_log( dashboard, ProvisionLogLine { run_id: *run_id, node_id: *node_id, stream: ProvisionLogStream::Stderr, line: line.clone(), }, ), PluginObservation::ProviderLine { run_id, node_id, line, } => orch_telemetry.emit_log( dashboard, ProvisionLogLine { run_id: *run_id, node_id: *node_id, stream: ProvisionLogStream::Provider, line: line.clone(), }, ), PluginObservation::TelemetryFrame { channel, payload, .. } => orch_telemetry.emit_bytes_from( dashboard, channel, payload.as_bytes().to_vec(), "node_bootstrap_stdio", ), PluginObservation::Exited { run_id, node_id, status, } => orch_telemetry.emit_event( dashboard, ProvisionEvent { run_id: *run_id, node_id: *node_id, kind: ProvisionEventKind::NodeStopped, provider: Some(provider.as_str().to_owned()), message: Some(format!("node process exited with {status:?}")), }, ), PluginObservation::Failed { run_id, node_id, reason, } => orch_telemetry.emit_event( dashboard, ProvisionEvent { run_id: *run_id, node_id: *node_id, kind: ProvisionEventKind::ProvisionFailed, provider: Some(provider.as_str().to_owned()), message: Some(reason.clone()), }, ), } } fn emit_swim_transitions( orch_telemetry: &mut OrchTelemetry, dashboard: Option<&DashboardSupport>, run_id: u64, node_id: u64, stack: &DistributionRuntimeStack, ) -> Vec { let transitions = stack.drain_swim_transitions(); for transition in &transitions { let peer = format!("{:?}", transition.peer); let from = transition.from.map(|state| format!("{:?}", state)); let to = format!("{:?}", transition.to); let member_state = stack .member_state(transition.peer) .map(|state| format!("{:?}", state)); let last_ack_age_ms = transition.last_ack_age.map(duration_ms_u64); let consecutive_timeouts = transition.consecutive_timeouts; let recent_probe_targets = stack.swim_recent_probe_targets(); orch_telemetry.emit_bootstrap_to_channel( dashboard, MYELIN_SWIM_MEMBERSHIP, run_id, node_id, "membership_transition", "observed", json!({ "peer":peer.clone(), "from":from.clone(), "to":to.clone(), "reason":transition.reason, "last_ack_age_ms":last_ack_age_ms, "consecutive_timeouts":consecutive_timeouts, "recent_probe_targets":recent_probe_targets.clone(), "member_state":member_state.clone(), }), ); orch_telemetry.emit_record(dashboard, &stack.membership_transition(transition)); } transitions } fn emit_swim_probe_events( orch_telemetry: &mut OrchTelemetry, dashboard: Option<&DashboardSupport>, stack: &DistributionRuntimeStack, local_phase: &str, ) { for event in stack.drain_swim_probe_events() { let record = stack.swim_probe_event_record(event, local_phase); orch_telemetry.emit_record(dashboard, &record); } } pub(crate) fn env_optional(name: &str) -> Option { std::env::var(name) .ok() .map(|value| value.trim().to_owned()) .filter(|value| !value.is_empty()) } fn local_tinygrad_worker_env(provider: &str) -> Option<(String, String)> { env_optional("MYELIN_TINYGRAD_WORKER") .map(|value| ("MYELIN_TINYGRAD_WORKER".to_owned(), value)) .or_else(|| { (provider == "process") .then(default_local_tinygrad_worker_path) .flatten() .map(|path| { ( "MYELIN_TINYGRAD_WORKER".to_owned(), path.to_string_lossy().to_string(), ) }) }) } fn default_local_tinygrad_worker_path() -> Option { // The tinygrad worker script ships in the node-image build context at // `apps/myelin/node-image/tinygrad_worker.py` (see `chat/node_image.rs` and // the node-image Dockerfile). The process provider runs it directly via // `python3`, so resolve that path from the workspace cwd or this crate's // manifest dir. (Previously looked in `apps/myelin-node/`, a path left stale // by the `mvp-system` -> `myelin` app refactor and never present on disk.) let cwd_candidate = std::env::current_dir().ok().map(|cwd| { cwd.join("apps") .join("myelin") .join("node-image") .join("tinygrad_worker.py") }); if let Some(candidate) = cwd_candidate.filter(|path| path.is_file()) { return Some(candidate.canonicalize().unwrap_or(candidate)); } let manifest_candidate = PathBuf::from(env!("CARGO_MANIFEST_DIR")) .join("node-image") .join("tinygrad_worker.py"); manifest_candidate.is_file().then(|| { manifest_candidate .canonicalize() .unwrap_or(manifest_candidate) }) } pub(crate) fn resolve_vastai_ssh_identity(explicit: Option) -> Result { match explicit { Some(path) => Ok(path), None => { let home = std::env::var_os("HOME") .filter(|value| !value.is_empty()) .ok_or_else(|| { "MYELIN_VASTAI_SSH_IDENTITY is required because HOME is unset".to_owned() })?; Ok(PathBuf::from(home).join(".ssh").join("id_ed25519")) } } } pub(crate) fn expand_home_path(value: &str) -> Result { let trimmed = value.trim(); if let Some(rest) = trimmed.strip_prefix("~/") { let home = std::env::var_os("HOME") .filter(|value| !value.is_empty()) .ok_or_else(|| "MYELIN_VASTAI_SSH_IDENTITY uses ~/ but HOME is unset".to_owned())?; return Ok(PathBuf::from(home).join(rest)); } Ok(PathBuf::from(trimmed)) } pub(crate) fn derive_ssh_public_key(identity: &Path) -> Result { let output = Command::new("ssh-keygen") .arg("-y") .arg("-f") .arg(identity) .output() .map_err(|e| { format!( "derive VastAI SSH public key from {}: {e}", identity.display() ) })?; let public_key = String::from_utf8_lossy(&output.stdout) .trim_end_matches(['\r', '\n']) .to_owned(); if !output.status.success() || public_key.trim().is_empty() { return Err(format!( "derive VastAI SSH public key from {}: {}", identity.display(), command_output_failure_detail(&output, None) )); } Ok(public_key) } pub(crate) fn ssh_public_key_fingerprint(public_key: &str) -> String { const UNAVAILABLE: &str = "unavailable"; let path = std::env::temp_dir().join(format!("myelin-vastai-ssh-key-{}.pub", std::process::id())); if std::fs::write(&path, format!("{public_key}\n")).is_err() { return UNAVAILABLE.to_owned(); } let output = Command::new("ssh-keygen") .arg("-l") .arg("-f") .arg(&path) .output() .ok(); let _ = std::fs::remove_file(&path); let Some(output) = output.filter(|output| output.status.success()) else { return UNAVAILABLE.to_owned(); }; let stdout = String::from_utf8_lossy(&output.stdout); let mut fields = stdout.split_whitespace(); fields .next() .zip(fields.next()) .map(|(bits, fingerprint)| format!("{bits} {fingerprint}")) .unwrap_or_else(|| UNAVAILABLE.to_owned()) } fn vastai_account_has_ssh_key(api_key: &str, public_key: &str) -> Result { let output = Command::new("vastai") .args(["show", "ssh-keys", "--raw", "--api-key", api_key]) .output() .map_err(vastai_cli_error)?; if !output.status.success() { return Err(format!( "vastai show ssh-keys failed: {}", command_output_failure_detail(&output, Some(api_key)) )); } let stdout = String::from_utf8_lossy(&output.stdout); Ok(account_ssh_keys_output_contains_public_key( &stdout, public_key, )) } pub(crate) fn ensure_vastai_account_ssh_key(api_key: &str, public_key: &str) -> Result<(), String> { if vastai_account_has_ssh_key(api_key, public_key)? { return Ok(()); } let output = Command::new("vastai") .args(["create", "ssh-key"]) .arg(public_key) .args(["-y", "--api-key", api_key]) .output() .map_err(vastai_cli_error)?; if !output.status.success() { return Err(format!( "vastai create ssh-key failed: {}", command_output_failure_detail(&output, Some(api_key)) )); } if vastai_account_has_ssh_key(api_key, public_key)? { Ok(()) } else { Err("VastAI SSH key registration did not make the selected key visible in vastai show ssh-keys".to_owned()) } } fn account_ssh_keys_output_contains_public_key(output: &str, public_key: &str) -> bool { let public_key = public_key.trim(); !public_key.is_empty() && (output.contains(public_key) || public_key .split_whitespace() .nth(1) .is_some_and(|body| !body.is_empty() && output.contains(body))) } fn vastai_cli_error(error: std::io::Error) -> String { if error.kind() == std::io::ErrorKind::NotFound { "vastai CLI is required to verify/register MYELIN_VASTAI_SSH_IDENTITY; install with pip install vastai".to_owned() } else { format!("run vastai CLI: {error}") } } fn command_output_failure_detail(output: &std::process::Output, secret: Option<&str>) -> String { let stderr = String::from_utf8_lossy(&output.stderr); let mut detail = match stderr.trim() { "" => output.status.to_string(), detail => detail.to_owned(), }; if let Some(secret) = secret.filter(|secret| !secret.is_empty()) { detail = detail.replace(secret, ""); } detail } fn next_arg(args: &mut impl Iterator, name: &str) -> Result { args.next() .ok_or_else(|| format!("missing value after {name}")) } fn parse_next(args: &mut impl Iterator, name: &str) -> Result where T: std::str::FromStr, T::Err: std::fmt::Display, { let value = next_arg(args, name)?; value .parse::() .map_err(|e| format!("invalid {name}={value:?}: {e}")) }