use std::collections::BTreeMap;
use std::fs;
use std::io::{self, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use anyhow::{Context, anyhow};
use meerkat::{AgentFactory, Config, FactoryAgentBuilder, PersistentSessionService};
use meerkat_mob::ids::AgentIdentity;
use meerkat_mob::{MobDefinition, MobStorage, ProfileName, SpawnMemberSpec};
use meerkat_mobkit::contact_directory::ContactDirectory;
use meerkat_mobkit::{
AuthPolicy, Base64BlobStoreAdapter, BigQueryNaming, BinaryBlobStore, ConsolePolicy,
ConsoleUiConfig, ConventionalPaths, GatewayPeerKeys, MOBKIT_CONTRACT_VERSION,
MobBootstrapOptions, MobBootstrapSpec, MobKitStorageLayout, ObjectStoreBlobStore,
ReleaseMetadata, RuntimeDecisionState, RuntimeOpsPolicy, TrustedOidcRuntimeConfig,
UnifiedRuntime, load_console_ui_config_from_path_for_realm,
mob_handle_runtime::mob_definition_may_use_image_generation,
};
use meerkat_store::SqliteSessionStore;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use tokio::io::{AsyncBufReadExt, BufReader};
const FALLBACK_TEMPLATE_VERSION: &str = "tux-fallback-v2";
type ScheduleHostInputs = (
meerkat::ScheduleService,
meerkat_mobkit::schedule_wiring::ScheduleMobTargetRegistry,
Arc<PersistentSessionService<FactoryAgentBuilder>>,
PathBuf,
);
type WorkGraphParts = (
meerkat::WorkGraphService,
meerkat_mobkit::workgraph_admission::WorkGraphAdmissionSlot,
PathBuf,
);
type PersistentSessionServiceParts = (
Arc<dyn meerkat_mob::MobSessionService>,
Arc<meerkat_runtime::MeerkatMachine>,
Arc<dyn BinaryBlobStore>,
Option<ScheduleHostInputs>,
Option<WorkGraphParts>,
meerkat_mobkit::storage_health::ResolvedStorageSummary,
meerkat_mobkit::mob_handle_runtime::SessionWriteEpochsHandle,
);
#[derive(Debug, Deserialize)]
struct InitParams {
workspace_root: Option<PathBuf>,
project_root: Option<PathBuf>,
context_root: Option<PathBuf>,
runtime_root: Option<PathBuf>,
store_path: Option<PathBuf>,
persistent_sessions: Option<bool>,
realm: Option<String>,
isolated: Option<bool>,
surface: Option<String>,
runtime_profile: Option<String>,
console_read_only: Option<bool>,
identity_first: Option<bool>,
identity_roster: Option<Vec<meerkat_mobkit::identity_first::DurableAgentSpec>>,
}
#[derive(Debug, Serialize, Deserialize, Default)]
struct RuntimeRegistry {
entries: Vec<RuntimeRegistryEntry>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct RuntimeRegistryEntry {
key: String,
runtime_id: String,
http_base_url: String,
pid: u32,
updated_at_ms: u64,
}
fn current_time_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64
}
fn short_hash(value: &str) -> String {
value.chars().take(8).collect()
}
fn load_registry(path: &Path) -> RuntimeRegistry {
fs::read_to_string(path)
.ok()
.and_then(|text| serde_json::from_str(&text).ok())
.unwrap_or_default()
}
fn save_registry(path: &Path, registry: &RuntimeRegistry) -> anyhow::Result<()> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let text = serde_json::to_string_pretty(registry)?;
fs::write(path, text)?;
Ok(())
}
async fn url_is_alive(url: &str) -> bool {
let client = match reqwest::Client::builder()
.timeout(Duration::from_secs(1))
.build()
{
Ok(client) => client,
Err(_) => return false,
};
client
.get(format!("{}/healthz", url.trim_end_matches('/')))
.send()
.await
.map(|response| response.status().is_success())
.unwrap_or(false)
}
fn conventional_paths(workspace_root: &Path) -> ConventionalPaths {
ConventionalPaths::discover(
workspace_root.join("config"),
workspace_root.join("deployment"),
)
}
fn collect_recursive_files(root: &Path, files: &mut Vec<PathBuf>) {
let Ok(entries) = fs::read_dir(root) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
collect_recursive_files(&path, files);
} else if path.is_file() {
files.push(path);
}
}
}
#[allow(clippy::too_many_arguments)]
fn config_fingerprint(
workspace_root: &Path,
realm: Option<&str>,
isolated: bool,
runtime_profile: &str,
persistent_sessions: bool,
console_read_only: bool,
runtime_root: &Path,
store_path: &Path,
project_root: &Path,
context_root: Option<&Path>,
paths: &ConventionalPaths,
) -> anyhow::Result<String> {
let mut hasher = Sha256::new();
let realpath = workspace_root
.canonicalize()
.unwrap_or_else(|_| workspace_root.to_path_buf());
hasher.update(realpath.to_string_lossy().as_bytes());
hasher.update(b"\n");
hasher.update(realm.unwrap_or("").as_bytes());
hasher.update(b"\n");
hasher.update(if isolated { b"1" } else { b"0" });
hasher.update(b"\n");
hasher.update(runtime_profile.as_bytes());
hasher.update(b"\n");
hasher.update(if persistent_sessions { b"1" } else { b"0" });
hasher.update(b"\n");
hasher.update(if console_read_only { b"1" } else { b"0" });
hasher.update(b"\n");
hasher.update(runtime_root.to_string_lossy().as_bytes());
hasher.update(b"\n");
hasher.update(store_path.to_string_lossy().as_bytes());
hasher.update(b"\n");
hasher.update(project_root.to_string_lossy().as_bytes());
hasher.update(b"\n");
if let Some(ctx) = context_root {
hasher.update(ctx.to_string_lossy().as_bytes());
}
hasher.update(b"\n");
hasher.update(env!("CARGO_PKG_VERSION").as_bytes());
let definition_json = workspace_root.join("definition.json");
if paths.mob_toml.is_some() {
hasher.update(b"\nworkspace-config");
} else if definition_json.exists() {
hasher.update(b"\ndefinition-json");
} else {
hasher.update(b"\nfallback-template:");
hasher.update(FALLBACK_TEMPLATE_VERSION.as_bytes());
}
let mut files = Vec::new();
if let Some(path) = &paths.mob_toml {
files.push(path.clone());
}
if let Some(path) = &paths.gating_toml {
files.push(path.clone());
}
if let Some(path) = &paths.console_toml {
files.push(path.clone());
}
if let Some(path) = &paths.routing_toml {
files.push(path.clone());
}
files.extend(paths.schedule_files.clone());
if definition_json.exists() {
files.push(definition_json);
}
let manifest_toml = workspace_root.join("manifest.toml");
if manifest_toml.exists() {
files.push(manifest_toml);
}
let mut scan_roots = vec![workspace_root.to_path_buf()];
if project_root != workspace_root {
scan_roots.push(project_root.to_path_buf());
}
if let Some(ctx) = context_root
&& ctx != workspace_root
{
scan_roots.push(ctx.to_path_buf());
}
for root in &scan_roots {
for extra_dir in ["skills", "hooks", "mcp", "config"] {
collect_recursive_files(&root.join(extra_dir), &mut files);
}
}
files.sort();
for path in files {
hasher.update(b"\nfile:");
hasher.update(path.to_string_lossy().as_bytes());
if let Ok(bytes) = fs::read(&path) {
hasher.update(b"\n");
hasher.update(bytes);
}
}
Ok(format!("{:x}", hasher.finalize()))
}
fn minimal_definition(runtime_id: &str) -> anyhow::Result<MobDefinition> {
MobDefinition::from_toml(&format!(
r#"
[mob]
id = "{runtime_id}"
orchestrator = "alpha"
[profiles.alpha]
model = "gpt-5.5"
skills = ["alpha-role"]
peer_description = "Runtime guide -- expands this runtime into a small mob and coordinates peers"
external_addressable = true
[profiles.alpha.tools]
builtins = true
comms = true
mob = true
mob_tasks = true
[profiles.worker]
model = "gpt-5.5"
skills = ["worker-role"]
peer_description = "General-purpose peer meerkat"
external_addressable = true
[profiles.worker.tools]
builtins = true
comms = true
mob_tasks = true
[wiring]
auto_wire_orchestrator = true
[skills.alpha-role]
source = "inline"
content = """
## Role
You are Alpha, the runtime guide for a lightweight Meerkat workspace.
## What You Can Do
- Answer directly when the job is simple.
- Grow the runtime into a small mob when parallel work helps.
- Spawn classic sub-agents for delegated background work.
- Spawn peer meerkats when a longer-lived collaborator should appear in the runtime.
## Preferred Growth Pattern
- For quick delegated work, use sub-agent tools.
- For visible collaborators inside this runtime, use mob tools to spawn worker peers.
- When you spawn worker peers, they should appear in the shared runtime UI.
## Coordination
- Use mob tools to spawn, list, wire, and retire meerkats.
- Use peers() and send() when peers are available.
- If asked to create a small team, prefer spawning `worker` peers unless the user clearly asks for classic sub-agents.
## Communication Style
Be explicit about whether you used a sub-agent or spawned a peer meerkat.
"""
[skills.worker-role]
source = "inline"
content = """
You are a general-purpose worker meerkat inside a lightweight runtime.
Complete assigned tasks concisely and report status back to Alpha.
If peer messaging is available, use it to report completion or blockers.
"""
"#
))
.map_err(|error| anyhow!("invalid fallback mob definition: {error}"))
}
fn load_definition(
workspace_root: &Path,
fingerprint: &str,
paths: &ConventionalPaths,
) -> anyhow::Result<(MobDefinition, bool)> {
if let Some(path) = &paths.mob_toml {
let text = fs::read_to_string(path)
.with_context(|| format!("failed to read {}", path.display()))?;
let definition = MobDefinition::from_toml(&text)
.with_context(|| format!("failed to parse {}", path.display()))?;
return Ok((definition, true));
}
let definition_json_path = workspace_root.join("definition.json");
if definition_json_path.exists() {
let text = fs::read_to_string(&definition_json_path)
.with_context(|| format!("failed to read {}", definition_json_path.display()))?;
let definition = serde_json::from_str::<MobDefinition>(&text)
.with_context(|| format!("failed to parse {}", definition_json_path.display()))?;
return Ok((definition, true));
}
let runtime_id = format!("tux-{}", short_hash(fingerprint));
Ok((minimal_definition(&runtime_id)?, false))
}
fn build_persistent_session_service(
layout: &MobKitStorageLayout,
runtime_root: PathBuf,
project_root: PathBuf,
context_root: Option<PathBuf>,
image_generation: bool,
realm_id: &str,
) -> anyhow::Result<PersistentSessionServiceParts> {
let store_dir = layout.state_dir().to_path_buf();
fs::create_dir_all(&store_dir)
.with_context(|| format!("failed to create {}", store_dir.display()))?;
let sqlite_path = layout.session_db().map_err(|e| anyhow!("{e}"))?.path;
let session_store = Arc::new(
SqliteSessionStore::open(sqlite_path.clone())
.with_context(|| format!("failed to open {}", sqlite_path.display()))?,
);
let session_store_incremental = meerkat_mobkit::storage_health::probe_session_store_incremental(
&(session_store.clone() as Arc<dyn meerkat::SessionStore>),
"SqliteSessionStore",
);
let binary_blob_store: Arc<dyn BinaryBlobStore> =
Arc::new(ObjectStoreBlobStore::local(layout.blob_root())?);
let blob_store: Arc<dyn meerkat_core::BlobStore> =
Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
let runtime_db_path = layout.runtime_db();
let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> = Arc::new(
meerkat_runtime::store::SqliteRuntimeStore::new(&runtime_db_path).map_err(|err| {
anyhow!(
"{}",
meerkat_mobkit::storage_health::RuntimeStoreResolutionError {
path: runtime_db_path.clone(),
message: err.to_string(),
}
)
})?,
);
let (runtime_store, session_write_epochs) =
meerkat_mobkit::mob_handle_runtime::epoch_tracking_runtime_store(runtime_store);
let adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
Arc::clone(&runtime_store),
Arc::clone(&blob_store),
));
let mut factory = AgentFactory::new(store_dir)
.session_store(session_store.clone())
.runtime_root(runtime_root)
.project_root(project_root)
.builtins(true)
.shell(true)
.mob(true)
.comms(true)
.memory(true);
if image_generation {
factory = factory.with_image_generation_machine(adapter.clone());
}
if let Some(context_root) = context_root {
factory = factory.context_root(context_root);
}
let config = Config::default();
let mut builder = FactoryAgentBuilder::new(factory, config);
builder.default_blob_store = Some(blob_store.clone());
let (schedule_tools, schedule_slot) =
match meerkat_mobkit::schedule_wiring::attach_schedule_tools_with_identity_targets_reporting(
&builder,
layout.state_dir(),
) {
Ok(tools) => (
Some(tools),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"schedule",
"SqliteScheduleStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"schedule",
format!("schedule store failed to open; schedule tools disabled: {error}"),
),
),
};
let workgraph_state_dir = layout.state_dir().to_path_buf();
let (workgraph, workgraph_slot) =
match meerkat_mobkit::workgraph_wiring::attach_workgraph_tools_reporting(
&builder,
&workgraph_state_dir,
realm_id,
) {
Ok((service, admission_slot)) => (
Some((service, admission_slot, workgraph_state_dir)),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"workgraph",
"SqliteWorkGraphStore",
),
),
Err(error) => (
None,
meerkat_mobkit::storage_health::StorageSlotSummary::degraded(
"workgraph",
format!("workgraph store failed to open; workgraph disabled: {error}"),
),
),
};
let service = Arc::new(PersistentSessionService::new(
builder,
64,
session_store,
Arc::clone(&runtime_store),
blob_store,
));
let schedule_host_inputs = schedule_tools.map(|tools| {
(
tools.service,
tools.mob_target_registry,
Arc::clone(&service),
layout.schedule_db(),
)
});
let mut slots = vec![
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"sessions",
"SqliteSessionStore",
),
meerkat_mobkit::storage_health::StorageSlotSummary::persistent(
"runtime",
"SqliteRuntimeStore",
),
meerkat_mobkit::storage_health::blob_slot_summary(
meerkat_mobkit::storage_health::BlobDurability::PersistentDisk,
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"console",
"InMemoryConsoleLogStore",
"declared default of this surface (UnifiedRuntime::bootstrap keeps in-memory \
console/metadata)",
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"metadata",
"InMemoryMetadataStore",
"declared default of this surface (UnifiedRuntime::bootstrap keeps in-memory \
console/metadata)",
),
schedule_slot,
workgraph_slot,
];
slots.extend(meerkat_mobkit::storage_health::scratch_ring_buffer_slots());
Ok((
service,
adapter,
binary_blob_store,
schedule_host_inputs,
workgraph,
meerkat_mobkit::storage_health::ResolvedStorageSummary::new(
meerkat_mobkit::storage_health::BlobDurability::PersistentDisk,
Some(session_store_incremental),
)
.with_slots(slots),
session_write_epochs,
))
}
fn runtime_decision_state(
runtime_id: &str,
console_ui: ConsoleUiConfig,
console_read_only: bool,
) -> RuntimeDecisionState {
RuntimeDecisionState {
bigquery: BigQueryNaming {
dataset: "tux_local".to_string(),
table: "runtime_events".to_string(),
},
modules: Vec::new(),
auth: AuthPolicy::default(),
trusted_oidc: TrustedOidcRuntimeConfig {
discovery_json: r#"{"issuer":"https://noop.example.com","authorization_endpoint":"https://noop.example.com/auth","token_endpoint":"https://noop.example.com/token","jwks_uri":"https://noop.example.com/.well-known/jwks.json","response_types_supported":["code"],"subject_types_supported":["public"],"id_token_signing_alg_values_supported":["RS256"]}"#.to_string(),
jwks_json: r#"{"keys":[]}"#.to_string(),
audience: runtime_id.to_string(),
},
console: ConsolePolicy {
require_app_auth: false,
read_only: console_read_only,
fetch_timeout_ms: None,
ui: console_ui,
},
ops: RuntimeOpsPolicy::default(),
release_metadata: ReleaseMetadata {
targets: vec!["local".to_string()],
support_matrix: "tux".to_string(),
},
}
}
fn print_json_line(value: &Value) {
let line = serde_json::to_string(value)
.unwrap_or_else(|_| r#"{"jsonrpc":"2.0","id":null,"error":{"code":-32603,"message":"serialization failed"}}"#.to_string());
let mut stdout = io::stdout().lock();
let _ = writeln!(stdout, "{line}");
let _ = stdout.flush();
}
fn parse_init_request(line: &str) -> anyhow::Result<(Value, InitParams)> {
let raw: Value = serde_json::from_str(line).context("failed to parse init request")?;
let method = raw
.get("method")
.and_then(Value::as_str)
.unwrap_or_default();
if method != "mobkit/init" {
return Err(anyhow!("expected mobkit/init, got {method}"));
}
let params = raw.get("params").cloned().unwrap_or_else(|| json!({}));
let parsed: InitParams = serde_json::from_value(params).context("invalid init params")?;
Ok((raw.get("id").cloned().unwrap_or(Value::Null), parsed))
}
fn env_bool(name: &str) -> anyhow::Result<Option<bool>> {
let Ok(value) = std::env::var(name) else {
return Ok(None);
};
let normalized = value.trim().to_ascii_lowercase();
match normalized.as_str() {
"" => Ok(None),
"1" | "true" | "yes" | "on" => Ok(Some(true)),
"0" | "false" | "no" | "off" => Ok(Some(false)),
_ => Err(anyhow!("{name} must be a boolean value")),
}
}
fn init_response(
request_id: Value,
runtime_id: &str,
http_base_url: &str,
launch_state: &str,
) -> Value {
json!({
"jsonrpc": "2.0",
"id": request_id,
"result": {
"contract_version": MOBKIT_CONTRACT_VERSION,
"runtime_id": runtime_id,
"http_base_url": http_base_url,
"launch_state": launch_state,
}
})
}
fn init_error(request_id: Value, code: i64, message: String) -> Value {
json!({
"jsonrpc": "2.0",
"id": request_id,
"error": {
"code": code,
"message": message,
}
})
}
const STORAGE_ADOPT_CHECKPOINTS_USAGE: &str = "usage: mobkit_gateway storage-adopt-checkpoints \
(--db <path> | --state-dir <dir>) [--apply] [--json]\n\
Adopt legacy (pre-typed) session documents inside continuity snapshots \
into typed checkpoint authority (storage-unification H3).\n\
Dry-run by default; --apply rewrites legacy rows in place under the \
exclusive maintenance fence.\n\
Exit codes: 0 clean, 1 refusals or fence/database failure, 2 usage error.";
fn run_storage_adopt_checkpoints(args: &[String]) -> i32 {
let mut db: Option<PathBuf> = None;
let mut state_dir: Option<PathBuf> = None;
let mut apply = false;
let mut json = false;
let mut iter = args.iter();
while let Some(arg) = iter.next() {
match arg.as_str() {
"--db" => match iter.next() {
Some(value) => db = Some(PathBuf::from(value)),
None => {
eprintln!("--db requires a path\n{STORAGE_ADOPT_CHECKPOINTS_USAGE}");
return 2;
}
},
"--state-dir" => match iter.next() {
Some(value) => state_dir = Some(PathBuf::from(value)),
None => {
eprintln!(
"--state-dir requires a directory\n{STORAGE_ADOPT_CHECKPOINTS_USAGE}"
);
return 2;
}
},
"--apply" => apply = true,
"--json" => json = true,
other => {
eprintln!("unknown argument {other:?}\n{STORAGE_ADOPT_CHECKPOINTS_USAGE}");
return 2;
}
}
}
let db_path = match (db, state_dir) {
(Some(db), None) => db,
(None, Some(dir)) => {
let layout = MobKitStorageLayout::with_injected_roots(dir, None);
match layout.continuity_db() {
Ok(resolved) => resolved.path,
Err(error) => {
eprintln!("{error}");
return 1;
}
}
}
_ => {
eprintln!(
"exactly one of --db / --state-dir is required\n{STORAGE_ADOPT_CHECKPOINTS_USAGE}"
);
return 2;
}
};
let mode = if apply {
meerkat_mobkit::identity_first::AdoptionMode::Apply
} else {
meerkat_mobkit::identity_first::AdoptionMode::DryRun
};
match meerkat_mobkit::identity_first::adopt_continuity_snapshots_blocking(&db_path, mode) {
Ok(report) => {
if json {
match serde_json::to_string_pretty(&report) {
Ok(text) => println!("{text}"),
Err(error) => {
eprintln!("failed to serialize adoption report: {error}");
return 1;
}
}
} else {
println!(
"continuity checkpoint adoption ({}) at {}",
if apply { "apply" } else { "dry-run" },
db_path.display()
);
println!("{report}");
}
i32::from(!report.is_clean())
}
Err(error) => {
eprintln!("{error}");
1
}
}
}
const STORAGE_MIGRATE_USAGE: &str = "usage: mobkit_gateway storage-migrate --state-dir <dir> \
[--apply] [--adopt <path>] [--json]\n\
Fenced offline migration of one MobKit state directory \
(storage-unification M6): ledger baseline, legacy-spelling renames, \
twin reconciliation, continuity checkpoint adoption, digest-format \
marker stamping (stamp-digest-format-markers), leftover census.\n\
Dry-run by default; --apply mutates under the exclusive maintenance \
fence. --adopt <path> resolves a divergent file-name twin by adopting \
that copy and archiving the rest read-only (requires --apply).\n\
Exit codes: 0 clean, 1 refusals or fence/store failure, 2 usage error.";
fn run_storage_migrate(args: &[String]) -> i32 {
let mut state_dir: Option<PathBuf> = None;
let mut adopt: Option<PathBuf> = None;
let mut apply = false;
let mut json = false;
let mut iter = args.iter();
while let Some(arg) = iter.next() {
match arg.as_str() {
"--state-dir" => match iter.next() {
Some(value) => state_dir = Some(PathBuf::from(value)),
None => {
eprintln!("--state-dir requires a directory\n{STORAGE_MIGRATE_USAGE}");
return 2;
}
},
"--adopt" => match iter.next() {
Some(value) => adopt = Some(PathBuf::from(value)),
None => {
eprintln!("--adopt requires a path\n{STORAGE_MIGRATE_USAGE}");
return 2;
}
},
"--apply" => apply = true,
"--json" => json = true,
other => {
eprintln!("unknown argument {other:?}\n{STORAGE_MIGRATE_USAGE}");
return 2;
}
}
}
let Some(state_dir) = state_dir else {
eprintln!("--state-dir is required\n{STORAGE_MIGRATE_USAGE}");
return 2;
};
if adopt.is_some() && !apply {
eprintln!("--adopt requires --apply\n{STORAGE_MIGRATE_USAGE}");
return 2;
}
let mode = if apply {
meerkat_mobkit::MigrateMode::Apply
} else {
meerkat_mobkit::MigrateMode::DryRun
};
let report = meerkat_mobkit::migrate_state_dir(&state_dir, mode, adopt.as_deref());
if json {
match serde_json::to_string_pretty(&report) {
Ok(text) => println!("{text}"),
Err(error) => {
eprintln!("failed to serialize migrate report: {error}");
return 1;
}
}
} else {
print_migrate_report_text(&report);
}
i32::from(report.has_errors())
}
const STORAGE_DOWNGRADE_USAGE: &str = "usage: mobkit_gateway storage-downgrade --state-dir <dir> \
[--apply] [--json]\n\
Fenced offline ROLLBACK of the head-canonical continuity upgrade \
(ledger mobkit-continuity v2 -> v1). Re-materializes every head+rows \
session back into a whole-document session_snapshots blob, drops the \
head-canonical trio, and rewinds the ledger row — which is what lets a \
previous release open the file again, keeping every post-upgrade turn.\n\
Dry-run by default; a dry run performs the ENTIRE reconstruction, \
including the per-document reader simulation, and then rolls back, so a \
clean dry run is evidence the apply run will work. --apply mutates \
under the exclusive maintenance fence, in ONE transaction whose last \
statement is the ledger rewind.\n\
Sessions whose retained transcript rewrite history cannot be re-inlined \
into a document a reader accepts are written WITHOUT that history and \
named individually in the report; their turn content is intact.\n\
Exit codes: 0 clean, 1 refusals or fence/store failure, 2 usage error.";
fn run_storage_downgrade(args: &[String]) -> i32 {
let mut state_dir: Option<PathBuf> = None;
let mut apply = false;
let mut json = false;
let mut iter = args.iter();
while let Some(arg) = iter.next() {
match arg.as_str() {
"--state-dir" => match iter.next() {
Some(value) => state_dir = Some(PathBuf::from(value)),
None => {
eprintln!("--state-dir requires a directory\n{STORAGE_DOWNGRADE_USAGE}");
return 2;
}
},
"--apply" => apply = true,
"--json" => json = true,
other => {
eprintln!("unknown argument {other:?}\n{STORAGE_DOWNGRADE_USAGE}");
return 2;
}
}
}
let Some(state_dir) = state_dir else {
eprintln!("--state-dir is required\n{STORAGE_DOWNGRADE_USAGE}");
return 2;
};
let mode = if apply {
meerkat_mobkit::MigrateMode::Apply
} else {
meerkat_mobkit::MigrateMode::DryRun
};
let report = meerkat_mobkit::downgrade_state_dir(&state_dir, mode);
if json {
match serde_json::to_string_pretty(&report) {
Ok(text) => println!("{text}"),
Err(error) => {
eprintln!("failed to serialize downgrade report: {error}");
return 1;
}
}
} else {
print!("{}", meerkat_mobkit::render_downgrade_report(&report));
}
i32::from(report.has_errors())
}
fn print_migrate_report_text(report: &meerkat_mobkit::MobKitMigrateReport) {
let mode = match report.mode {
meerkat_mobkit::MigrateMode::Apply => "apply",
_ => "dry-run",
};
println!(
"Storage migrate ({mode}) over {} ({} database(s) fenced):",
report.state_dir.display(),
report.fenced_databases.len()
);
for twin in &report.twins {
println!("twin [{}]:", twin.slot);
for path in &twin.paths {
println!(" copy: {}", path.display());
}
println!(
" rows: {} equal across all copies; byte-identical: {}",
twin.rows_equal, twin.byte_identical
);
for row in &twin.rows {
let status = match &row.status {
meerkat_mobkit::DivergenceStatus::Equal => "equal".to_string(),
meerkat_mobkit::DivergenceStatus::Divergent => "divergent".to_string(),
meerkat_mobkit::DivergenceStatus::OnlyIn { location } => {
format!("only in {}", location.display())
}
_ => "unknown".to_string(),
};
println!(" {}: {status}", row.key);
}
match &twin.resolution {
meerkat_mobkit::TwinResolution::Refused { reason } => {
println!(" resolution: REFUSED — {reason}");
}
meerkat_mobkit::TwinResolution::Deduped { kept, archived } => {
println!(" resolution: deduped, kept {}", kept.display());
for archive in archived {
println!(" archived read-only: {}", archive.display());
}
}
meerkat_mobkit::TwinResolution::Adopted { adopted, archived } => {
println!(" resolution: adopted {}", adopted.display());
for archive in archived {
println!(" archived read-only: {}", archive.display());
}
}
_ => println!(" resolution: unknown"),
}
for note in &twin.notes {
println!(" note: {note}");
}
for error in &twin.errors {
println!(" error: {error}");
}
}
for rename in &report.renames {
let action = match rename.action {
meerkat_mobkit::RenameAction::WouldRename => "would rename",
meerkat_mobkit::RenameAction::Renamed => "renamed",
meerkat_mobkit::RenameAction::Refused => "REFUSED",
_ => "unknown",
};
println!(
"rename [{}]: {} -> {} ({action}, {} sibling(s), wal checkpointed: {})",
rename.slot,
rename.from.display(),
rename.to.display(),
rename.siblings.len(),
rename.wal_checkpointed
);
}
for entry in &report.ledger {
let action = match entry.action {
meerkat_mobkit::LedgerBaselineAction::WouldStamp => "would-stamp",
meerkat_mobkit::LedgerBaselineAction::Recorded => "recorded",
meerkat_mobkit::LedgerBaselineAction::Stamped => "stamped",
meerkat_mobkit::LedgerBaselineAction::AlreadyCurrent => "already-current",
meerkat_mobkit::LedgerBaselineAction::ReportOnly => "report-only",
meerkat_mobkit::LedgerBaselineAction::Exempt => "exempt",
_ => "unknown",
};
let describe =
|version: Option<i64>| version.map_or_else(|| "none".to_string(), |v| v.to_string());
println!(
"ledger {} [{}]: {} -> {} ({action})",
entry.database.display(),
entry.domain,
describe(entry.before),
describe(entry.after)
);
}
if let Some(adoption) = &report.adoption {
match (&adoption.skipped, &adoption.report) {
(Some(skipped), _) => println!("adoption: {skipped}"),
(None, Some(walk)) => {
println!("adoption at {}:\n{walk}", adoption.database.display());
}
(None, None) => {}
}
}
for outcome in &report.marker_stamping {
match (&outcome.skipped, &outcome.report) {
(Some(skipped), _) => {
println!("marker stamping [{}]: {skipped}", outcome.store);
}
(None, Some(walk)) => {
println!(
"marker stamping [{}] at {}:\n{walk}",
outcome.store,
outcome.database.display()
);
}
(None, None) => {}
}
}
for finding in &report.findings {
let path = finding
.path
.as_ref()
.map(|path| format!(" at {}", path.display()))
.unwrap_or_default();
println!("leftover [{}] {}{path}", finding.code, finding.message);
}
for note in &report.notes {
println!("note: {note}");
}
for error in &report.errors {
println!("error: {error}");
}
println!("storage migrate: {} error(s)", report.errors.len());
}
const STORAGE_PRUNE_USAGE: &str = "usage: mobkit_gateway storage-prune --state-dir <dir> \
[--apply] [--older-than-days N] [--json]\n\
Lifecycle of registered maintenance artifacts (`*.pre-*` backups, \
`*.corrupt-*` quarantines) under one MobKit state directory. Never \
touches anything outside those naming patterns.\n\
Dry-run by default; --apply deletes artifacts at least \
--older-than-days old (default 30; 0 = all).\n\
Exit codes: 0 clean, 1 delete failures, 2 usage error.";
fn run_storage_prune(args: &[String]) -> i32 {
let mut state_dir: Option<PathBuf> = None;
let mut older_than_days: u64 = 30;
let mut apply = false;
let mut json = false;
let mut iter = args.iter();
while let Some(arg) = iter.next() {
match arg.as_str() {
"--state-dir" => match iter.next() {
Some(value) => state_dir = Some(PathBuf::from(value)),
None => {
eprintln!("--state-dir requires a directory\n{STORAGE_PRUNE_USAGE}");
return 2;
}
},
"--older-than-days" => match iter.next().map(|value| value.parse::<u64>()) {
Some(Ok(days)) => older_than_days = days,
_ => {
eprintln!("--older-than-days requires a number\n{STORAGE_PRUNE_USAGE}");
return 2;
}
},
"--apply" => apply = true,
"--json" => json = true,
other => {
eprintln!("unknown argument {other:?}\n{STORAGE_PRUNE_USAGE}");
return 2;
}
}
}
let Some(state_dir) = state_dir else {
eprintln!("--state-dir is required\n{STORAGE_PRUNE_USAGE}");
return 2;
};
let mode = if apply {
meerkat_mobkit::MigrateMode::Apply
} else {
meerkat_mobkit::MigrateMode::DryRun
};
let report = meerkat_mobkit::prune_state_dir(&state_dir, older_than_days, mode);
if json {
match serde_json::to_string_pretty(&report) {
Ok(text) => println!("{text}"),
Err(error) => {
eprintln!("failed to serialize prune report: {error}");
return 1;
}
}
} else {
let mode = if apply { "apply" } else { "dry-run" };
println!(
"Storage prune ({mode}, older than {} day(s)) over {}:",
report.older_than_days,
report.state_dir.display()
);
if report.artifacts.is_empty() {
println!("No registered maintenance artifacts found.");
}
for artifact in &report.artifacts {
let action = match artifact.action {
meerkat_mobkit::PruneAction::WouldDelete => "would delete",
meerkat_mobkit::PruneAction::Deleted => "deleted",
meerkat_mobkit::PruneAction::Kept => "kept (younger than threshold)",
meerkat_mobkit::PruneAction::DeleteFailed => "DELETE FAILED",
_ => "unknown",
};
println!(
" {} {} bytes, {} day(s) old — {action}",
artifact.path.display(),
artifact.bytes,
artifact.age_days
);
}
for error in &report.errors {
println!("error: {error}");
}
println!("storage prune: {} error(s)", report.errors.len());
}
i32::from(report.has_errors())
}
fn main() {
let args: Vec<String> = std::env::args().skip(1).collect();
if args.iter().any(|a| a == "--version" || a == "-V") {
println!(
"mobkit_gateway {} (meerkat-mobkit console/HTTP gateway)",
env!("CARGO_PKG_VERSION")
);
return;
}
if args.first().map(String::as_str) == Some("storage-adopt-checkpoints") {
std::process::exit(run_storage_adopt_checkpoints(&args[1..]));
}
if args.first().map(String::as_str) == Some("storage-migrate") {
std::process::exit(run_storage_migrate(&args[1..]));
}
if args.first().map(String::as_str) == Some("storage-prune") {
std::process::exit(run_storage_prune(&args[1..]));
}
if args.first().map(String::as_str) == Some("storage-downgrade") {
std::process::exit(run_storage_downgrade(&args[1..]));
}
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("warn")),
)
.with_writer(std::io::stderr)
.with_ansi(false)
.init();
tracing::info!(
version = env!("CARGO_PKG_VERSION"),
"mobkit_gateway starting (console/HTTP gateway)"
);
let runtime = match tokio::runtime::Builder::new_multi_thread()
.enable_all()
.thread_stack_size(16 * 1024 * 1024)
.build()
{
Ok(runtime) => runtime,
Err(error) => {
let response = init_error(Value::Null, -32603, error.to_string());
print_json_line(&response);
std::process::exit(1);
}
};
if let Err(error) = runtime.block_on(Box::pin(run())) {
let response = init_error(Value::Null, -32603, error.to_string());
print_json_line(&response);
std::process::exit(1);
}
}
async fn run() -> anyhow::Result<()> {
let stdin = tokio::io::stdin();
let mut reader = BufReader::new(stdin);
let mut init_line = String::new();
if reader.read_line(&mut init_line).await? == 0 {
return Err(anyhow!("stdin closed before init request"));
}
let (request_id, params) = match parse_init_request(init_line.trim()) {
Ok(value) => value,
Err(error) => {
print_json_line(&init_error(Value::Null, -32602, error.to_string()));
return Err(error);
}
};
let workspace_root = params
.workspace_root
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
let workspace_root = workspace_root.canonicalize().unwrap_or(workspace_root);
let project_root = params
.project_root
.unwrap_or_else(|| workspace_root.clone());
let project_root = project_root.canonicalize().unwrap_or(project_root);
let context_root = params
.context_root
.map(|path| path.canonicalize().unwrap_or(path))
.or_else(|| Some(project_root.clone()));
let runtime_root = params
.runtime_root
.unwrap_or_else(|| workspace_root.clone());
let runtime_root = runtime_root.canonicalize().unwrap_or(runtime_root);
let store_path = params
.store_path
.unwrap_or_else(|| runtime_root.join("state"));
let store_path = store_path.canonicalize().unwrap_or(store_path);
let gateway_home = meerkat_mobkit::storage_layout::default_gateway_home()
.context("resolve gateway state directory")?;
let layout = MobKitStorageLayout::standalone_from_store_path(&store_path, gateway_home.clone());
let persistent_sessions = params.persistent_sessions.unwrap_or(false);
let identity_first = params.identity_first.unwrap_or(true);
let identity_roster_seed = params.identity_roster.clone().unwrap_or_default();
let realm = params.realm.as_deref();
let isolated = params.isolated.unwrap_or(false);
let _surface = params.surface.unwrap_or_else(|| "tux".to_string());
let runtime_profile = params
.runtime_profile
.unwrap_or_else(|| "tux-auto".to_string());
let console_read_only = match params.console_read_only {
Some(value) => value,
None => env_bool("MOBKIT_CONSOLE_READ_ONLY")?.unwrap_or(false),
};
let paths = conventional_paths(&workspace_root);
let key = config_fingerprint(
&workspace_root,
realm,
isolated,
&runtime_profile,
persistent_sessions,
console_read_only,
&runtime_root,
&store_path,
&project_root,
context_root.as_deref(),
&paths,
)?;
let registry_file = layout
.registry_file()
.ok_or_else(|| anyhow!("gateway storage layout carries no gateway home"))?;
let mut registry = load_registry(®istry_file);
let mut live_entries = Vec::new();
let mut resumed_entry = None;
for entry in registry.entries.drain(..) {
if url_is_alive(&entry.http_base_url).await {
if entry.key == key {
resumed_entry = Some(entry.clone());
}
live_entries.push(entry);
}
}
registry.entries = live_entries;
save_registry(®istry_file, ®istry)?;
if let Some(entry) = resumed_entry {
print_json_line(&init_response(
request_id,
&entry.runtime_id,
&entry.http_base_url,
"resumed",
));
return Ok(());
}
std::env::set_current_dir(&workspace_root).ok();
let (definition, used_workspace_config) = load_definition(&workspace_root, &key, &paths)?;
let console_ui = match &paths.console_toml {
Some(path) => load_console_ui_config_from_path_for_realm(path, realm)
.with_context(|| format!("failed to load {}", path.display()))?,
None => ConsoleUiConfig::default(),
};
let runtime_id = definition.id.to_string();
let image_generation = mob_definition_may_use_image_generation(&definition);
let (session_spec, schedule_host_inputs, workgraph_service) = if persistent_sessions {
let (
service,
adapter,
binary_blob_store,
schedule_host_inputs,
workgraph,
resolved_storage,
session_write_epochs,
) = build_persistent_session_service(
&layout,
runtime_root.clone(),
project_root.clone(),
context_root.clone(),
image_generation,
&runtime_id,
)?;
let schedule_host_inputs = schedule_host_inputs
.map(|(sched, registry, svc, path)| (sched, registry, svc, path, adapter.clone()));
let workgraph_service = workgraph.as_ref().map(|(service, _, _)| service.clone());
let mut spec = MobBootstrapSpec::new(definition, MobStorage::in_memory(), service)
.with_session_write_epochs(&session_write_epochs)
.with_session_runtime_adapter(adapter.clone())
.with_workgraph_service(workgraph_service.clone());
if let Some((_, admission_slot, state_dir)) = &workgraph {
spec = spec
.with_workgraph_admission_slot(admission_slot.clone())
.with_workgraph_admission_sidecar(state_dir);
}
spec.runtime_adapter = Some(adapter);
spec.binary_blob_store = Some(binary_blob_store);
spec.resolved_storage = Some(resolved_storage);
(spec, schedule_host_inputs, workgraph_service)
} else {
let binary_blob_store: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
let blob_store: Arc<dyn meerkat_core::BlobStore> =
Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
Arc::new(meerkat_runtime::InMemoryRuntimeStore::new());
let (runtime_store, session_write_epochs) =
meerkat_mobkit::mob_handle_runtime::epoch_tracking_runtime_store(runtime_store);
let adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
Arc::clone(&runtime_store),
Arc::clone(&blob_store),
));
let mut factory = AgentFactory::new(&runtime_root)
.runtime_root(runtime_root.clone())
.project_root(project_root.clone())
.builtins(true)
.shell(true)
.mob(true)
.comms(true)
.memory(true);
if image_generation {
factory = factory.with_image_generation_machine(adapter.clone());
}
if let Some(ref ctx) = context_root {
factory = factory.context_root(ctx.clone());
}
let config = Config::default();
let mut builder = FactoryAgentBuilder::new(factory, config);
builder.default_blob_store = Some(blob_store);
let (ephemeral_workgraph, workgraph_admission_slot) =
meerkat_mobkit::workgraph_wiring::attach_workgraph_tools_ephemeral(
&builder,
&runtime_id,
);
let workgraph_service = Some(ephemeral_workgraph);
let session_service = Arc::new(meerkat_session::EphemeralSessionService::new(builder, 64));
let mut spec = MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
.with_session_write_epochs(&session_write_epochs)
.with_session_runtime_adapter(adapter.clone())
.with_workgraph_service(workgraph_service.clone())
.with_workgraph_admission_slot(workgraph_admission_slot);
spec.runtime_adapter = Some(adapter);
spec.binary_blob_store = Some(binary_blob_store);
let mut slots = vec![
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"sessions",
"EphemeralSessionService",
"declared by the default ephemeral launch (persistent_sessions = false)",
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"runtime",
"InMemoryRuntimeStore",
"declared by the default ephemeral launch",
),
meerkat_mobkit::storage_health::blob_slot_summary(
meerkat_mobkit::storage_health::BlobDurability::DeclaredEphemeral,
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"console",
"InMemoryConsoleLogStore",
"declared default of this surface",
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"metadata",
"InMemoryMetadataStore",
"declared default of this surface",
),
meerkat_mobkit::storage_health::StorageSlotSummary::declared_ephemeral(
"workgraph",
"MemoryWorkGraphStore",
"declared by the default ephemeral launch",
),
];
slots.extend(meerkat_mobkit::storage_health::scratch_ring_buffer_slots());
spec.resolved_storage = Some(
meerkat_mobkit::storage_health::ResolvedStorageSummary::new(
meerkat_mobkit::storage_health::BlobDurability::DeclaredEphemeral,
None,
)
.with_slots(slots),
);
(spec, None, workgraph_service)
};
let mob_spec = session_spec.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: !persistent_sessions,
notify_orchestrator_on_resume: true,
default_llm_client: None,
});
let mut runtime = Box::pin(UnifiedRuntime::bootstrap(
mob_spec,
meerkat_mobkit::MobKitConfig {
modules: Vec::new(),
discovery: meerkat_mobkit::DiscoverySpec {
namespace: format!("tux.{}", short_hash(&key)),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
},
Duration::from_secs(30),
))
.await
.context("failed to bootstrap local runtime")?;
if let Some(ref contacts_path) = paths.contacts_toml {
let contacts_text = fs::read_to_string(contacts_path)
.with_context(|| format!("failed to read {}", contacts_path.display()))?;
let directory = ContactDirectory::from_toml(&contacts_text)
.with_context(|| format!("failed to parse {}", contacts_path.display()))?;
runtime.set_contact_directory(directory);
}
if let Some(ref access_path) = paths.access_toml {
let controller = meerkat_mobkit::AccessController::load_or_default(access_path)
.map_err(|error| anyhow!("failed to load {}: {error}", access_path.display()))?;
runtime.set_access_controller(controller);
}
let peer_keys = GatewayPeerKeys::load_or_create(&gateway_home).with_context(|| {
format!(
"failed to load or mint gateway peer key under {}",
gateway_home.display()
)
})?;
runtime.set_gateway_peer_keys(peer_keys);
if !used_workspace_config {
let mut labels = BTreeMap::new();
labels.insert("surface".to_string(), "tux".to_string());
labels.insert("ui".to_string(), "meerkat-tux".to_string());
if let Some(realm) = realm {
labels.insert("realm".to_string(), realm.to_string());
}
runtime
.mob_handle()
.ensure_member(
SpawnMemberSpec::new(ProfileName::from("alpha"), AgentIdentity::from("alpha"))
.with_labels(labels),
)
.await
.map_err(|error| anyhow!("failed to spawn fallback alpha meerkat: {error}"))?;
}
let _identity_roster_provider: Option<
Arc<meerkat_mobkit::identity_first::MutableRosterProvider>,
> = if identity_first {
use meerkat_mobkit::identity_first::{
AgentRuntimeServices, DurabilityPolicy, IdentityFirstRuntimeContext, IdentityRuntime,
IdentityRuntimeConfig, MobSessionBridge, MutableRosterProvider, restore_flow,
};
let store_dir = layout.state_dir();
fs::create_dir_all(store_dir)
.with_context(|| format!("failed to create {}", store_dir.display()))?;
let continuity_db = layout.continuity_db().map_err(|e| anyhow!("{e}"))?.path;
let substrate = meerkat_mobkit::gateway_wiring::open_identity_substrate(&continuity_db)
.await
.map_err(|e| anyhow!("{e}"))?;
let mob_handle = runtime.mob_handle();
let bridge: Arc<dyn meerkat_mobkit::identity_first::SessionBridge> =
if let Some(session_service) = runtime.mob_runtime().session_service().cloned() {
Arc::new(MobSessionBridge::with_session_service(
mob_handle.clone(),
session_service,
))
} else {
Arc::new(MobSessionBridge::new(mob_handle.clone()))
};
let irt = IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: substrate.continuity_store,
lease_provider: substrate.lease_provider,
runtime_instance_id: format!("mobkit-gateway-{}", std::process::id()),
has_runtime_store: persistent_sessions,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(bridge),
default_timeout: None,
})
.with_runtime_services(AgentRuntimeServices::new(mob_handle.clone()));
let roster = Arc::new(MutableRosterProvider::new(identity_roster_seed));
let mob_definition = mob_handle.definition().clone();
let irt = Arc::new(irt);
restore_flow(&irt, &roster.snapshot(), None, None)
.await
.context("identity-first restore_flow failed")?;
runtime.set_console_identity_roster(roster.clone());
runtime.attach_identity_first_context(Arc::new(IdentityFirstRuntimeContext::new(
irt,
roster.clone(),
None,
None,
Some(mob_definition),
)));
tracing::info!(
roster = roster.snapshot().len(),
continuity_db = %continuity_db.display(),
"identity-first gateway mode active"
);
Some(roster)
} else {
None
};
let (_schedule_host, _schedule_watchdog) = if let Some((
schedule_service,
mob_target_registry,
service,
schedule_store_path,
adapter,
)) = schedule_host_inputs
{
let mob_state = runtime.mob_runtime().agent_mob_mcp_state();
mob_target_registry.set_mob_state(mob_state.clone());
match meerkat_mobkit::schedule_wiring::repair_resumable_session_targets_to_mob_members(
&schedule_service,
&mob_target_registry,
)
.await
{
Ok(repaired) if repaired > 0 => {
tracing::info!(
repaired,
"repaired persisted resumable-session schedules to identity mob targets"
);
}
Ok(_) => {}
Err(error) => {
tracing::warn!(
error = %error,
"failed to repair persisted resumable-session schedules to identity mob targets",
);
}
}
let watchdog = meerkat_mobkit::schedule_wiring::spawn_schedule_claim_watchdog(
schedule_service.clone(),
schedule_store_path,
Default::default(),
);
(
meerkat_mobkit::schedule_wiring::spawn_schedule_host_with_identity_runtime(
service,
adapter,
schedule_service,
mob_state,
runtime.mob_handle(),
runtime.identity_runtime().cloned(),
None,
workgraph_service.clone(),
runtime_id.clone(),
),
Some(watchdog),
)
} else {
(None, None)
};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.context("failed to bind gateway listener")?;
let http_base_url = format!(
"http://127.0.0.1:{}",
listener.local_addr().context("missing local addr")?.port()
);
registry.entries.retain(|entry| entry.key != key);
registry.entries.push(RuntimeRegistryEntry {
key: key.clone(),
runtime_id: runtime_id.clone(),
http_base_url: http_base_url.clone(),
pid: std::process::id(),
updated_at_ms: current_time_ms(),
});
save_registry(®istry_file, ®istry)?;
print_json_line(&init_response(
request_id,
&runtime_id,
&http_base_url,
"created",
));
let decisions = runtime_decision_state(&runtime_id, console_ui, console_read_only);
let app = runtime.build_reference_app_router(decisions);
let stdin_guard = async move {
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line).await {
Ok(0) | Err(_) => break, Ok(_) => {}
}
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let id = serde_json::from_str::<Value>(trimmed)
.ok()
.and_then(|value| value.get("id").cloned())
.unwrap_or(Value::Null);
print_json_line(&serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": -32601,
"message": "mobkit_gateway serves the console/admin API over HTTP and does not handle SDK stdin JSON-RPC after init. Point your SDK gateway path (MOBKIT_RPC_GATEWAY_BIN) at the 'rpc_gateway' binary instead."
}
}));
}
std::future::pending::<()>().await;
};
tokio::select! {
result = axum::serve(listener, app).with_graceful_shutdown(async {
let _ = tokio::signal::ctrl_c().await;
}) => {
result.context("gateway HTTP server failed")?;
}
() = stdin_guard => {}
}
let mut registry = load_registry(®istry_file);
registry.entries.retain(|entry| entry.key != key);
save_registry(®istry_file, ®istry)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn init_params_parse_console_read_only() -> anyhow::Result<()> {
let line = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "mobkit/init",
"params": {
"console_read_only": true
}
})
.to_string();
let (_id, params) = parse_init_request(&line)?;
assert_eq!(params.console_read_only, Some(true));
Ok(())
}
#[test]
fn runtime_decision_state_projects_console_read_only() {
let state = runtime_decision_state("test-runtime", ConsoleUiConfig::default(), true);
assert!(state.console.read_only);
}
#[test]
fn config_fingerprint_changes_with_console_read_only() -> anyhow::Result<()> {
let temp = tempfile::tempdir()?;
let paths = conventional_paths(temp.path());
let writable = config_fingerprint(
temp.path(),
None,
false,
"tux-auto",
false,
false,
temp.path(),
temp.path(),
temp.path(),
None,
&paths,
)?;
let read_only = config_fingerprint(
temp.path(),
None,
false,
"tux-auto",
false,
true,
temp.path(),
temp.path(),
temp.path(),
None,
&paths,
)?;
assert_ne!(writable, read_only);
Ok(())
}
}