| 1 | //! Presentation-only companion. Replaces per-view live world ownership; the |
| 2 | //! existing Engine projection, PetNative world and score remain authoritative. |
| 3 | //! No agent, prompt, provider or execution API is available to this process. |
| 4 | use std::collections::{BTreeMap, VecDeque}; |
| 5 | use std::io::{self, Read}; |
| 6 | use std::path::{Path, PathBuf}; |
| 7 | use std::sync::{Arc, Mutex, mpsc}; |
| 8 | use std::time::{Duration, Instant}; |
| 9 | |
| 10 | use axum::{ |
| 11 | Json, Router, |
| 12 | extract::{DefaultBodyLimit, State}, |
| 13 | http::{HeaderMap, StatusCode}, |
| 14 | response::{Html, IntoResponse, Response}, |
| 15 | routing::{get, post}, |
| 16 | }; |
| 17 | use rquickjs::{Context, Runtime}; |
| 18 | use serde::{Deserialize, Serialize}; |
| 19 | use serde_json::{Value, json}; |
| 20 | |
| 21 | use super::{appearance::Appearance, persistence::Store}; |
| 22 | use crate::fleet::files::{WorkspaceFile, same_file}; |
| 23 | |
| 24 | const LEASE: Duration = Duration::from_secs(2); |
| 25 | const MAX_CLIENTS: usize = 4096; |
| 26 | /// Request body limit for every owner route. Clients size batches to it. |
| 27 | pub(super) const MAX_REQUEST_BYTES: usize = 64 * 1024; |
| 28 | /// How long a route waits for the world thread's reply before answering 503. |
| 29 | pub(super) const WORK_REPLY_TIMEOUT: Duration = Duration::from_secs(5); |
| 30 | |
| 31 | #[derive(Clone, Serialize, Deserialize)] |
| 32 | #[serde(deny_unknown_fields)] |
| 33 | pub struct Descriptor { |
| 34 | pub version: u8, |
| 35 | pub port: u16, |
| 36 | pub token: String, |
| 37 | pub identity: String, |
| 38 | } |
| 39 | |
| 40 | pub fn directory() -> io::Result<PathBuf> { |
| 41 | if let Some(path) = std::env::var_os("CODEWHALE_PET_HOME") { |
| 42 | return Ok(path.into()); |
| 43 | } |
| 44 | Ok(dirs::home_dir() |
| 45 | .ok_or_else(|| io::Error::other("Home directory unavailable"))? |
| 46 | .join(".codewhale/pet-shared")) |
| 47 | } |
| 48 | |
| 49 | pub fn descriptor(root: &Path) -> io::Result<Descriptor> { |
| 50 | let file = WorkspaceFile::open(root, Path::new("connection.json"), false)?.open_file()?; |
| 51 | let mut text = String::new(); |
| 52 | file.take(4097).read_to_string(&mut text)?; |
| 53 | if text.len() > 4096 { |
| 54 | return Err(io::Error::other("Invalid pet connection")); |
| 55 | } |
| 56 | let d: Descriptor = serde_json::from_str(&text)?; |
| 57 | if d.version != 1 |
| 58 | || d.port == 0 |
| 59 | || d.token.len() != 64 |
| 60 | || !d.token.bytes().all(|b| b.is_ascii_hexdigit()) |
| 61 | || uuid::Uuid::parse_str(&d.identity).is_err() |
| 62 | { |
| 63 | return Err(io::Error::other("Invalid pet connection")); |
| 64 | } |
| 65 | Ok(d) |
| 66 | } |
| 67 | |
| 68 | #[derive(Clone, Serialize, Deserialize)] |
| 69 | struct Receipt { |
| 70 | seq: u64, |
| 71 | hash: String, |
| 72 | cursor: u64, |
| 73 | } |
| 74 | |
| 75 | #[derive(Serialize, Deserialize)] |
| 76 | struct Saved { |
| 77 | version: u8, |
| 78 | identity: String, |
| 79 | token: String, |
| 80 | port: u16, |
| 81 | source: String, |
| 82 | source_revision: u64, |
| 83 | cursor: u64, |
| 84 | clients: BTreeMap<String, Receipt>, |
| 85 | recording: Value, |
| 86 | #[serde(default)] |
| 87 | appearance: Appearance, |
| 88 | } |
| 89 | |
| 90 | #[derive(Clone, Deserialize, Serialize)] |
| 91 | #[serde(deny_unknown_fields)] |
| 92 | pub struct Request { |
| 93 | pub identity: String, |
| 94 | pub client: String, |
| 95 | pub seq: u64, |
| 96 | pub source_revision: u64, |
| 97 | pub action: Action, |
| 98 | } |
| 99 | |
| 100 | #[derive(Clone, Deserialize, Serialize)] |
| 101 | #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] |
| 102 | pub enum Action { |
| 103 | Interact { food: bool, x: f64, y: f64 }, |
| 104 | Select { source: String }, |
| 105 | Appearance { appearance: Appearance }, |
| 106 | } |
| 107 | |
| 108 | #[derive(Deserialize)] |
| 109 | #[serde(deny_unknown_fields)] |
| 110 | pub struct Producer { |
| 111 | pub identity: String, |
| 112 | pub epoch: String, |
| 113 | pub client: String, |
| 114 | pub source: String, |
| 115 | pub source_revision: u64, |
| 116 | pub seq: u64, |
| 117 | pub waiting: bool, |
| 118 | pub events: Vec<Value>, |
| 119 | } |
| 120 | |
| 121 | #[derive(Deserialize)] |
| 122 | #[serde(deny_unknown_fields)] |
| 123 | pub struct AudioLease { |
| 124 | pub client: String, |
| 125 | pub enabled: bool, |
| 126 | } |
| 127 | |
| 128 | enum Work { |
| 129 | Action(Request, tokio::sync::oneshot::Sender<Result<Value, String>>), |
| 130 | Producer( |
| 131 | Producer, |
| 132 | tokio::sync::oneshot::Sender<Result<Value, String>>, |
| 133 | ), |
| 134 | Audio( |
| 135 | AudioLease, |
| 136 | tokio::sync::oneshot::Sender<Result<Value, String>>, |
| 137 | ), |
| 138 | Export(tokio::sync::oneshot::Sender<Result<Value, String>>), |
| 139 | } |
| 140 | |
| 141 | #[derive(Clone)] |
| 142 | struct Service { |
| 143 | descriptor: Descriptor, |
| 144 | frames: Arc<Mutex<VecDeque<Value>>>, |
| 145 | tx: mpsc::SyncSender<Work>, |
| 146 | } |
| 147 | |
| 148 | fn valid_client(id: &str) -> bool { |
| 149 | uuid::Uuid::parse_str(id).is_ok() |
| 150 | } |
| 151 | fn valid_source(id: &str) -> bool { |
| 152 | !id.is_empty() |
| 153 | && id.len() <= 128 |
| 154 | && id |
| 155 | .bytes() |
| 156 | .all(|b| b.is_ascii_alphanumeric() || b"-_:./".contains(&b)) |
| 157 | } |
| 158 | |
| 159 | fn authorized(headers: &HeaderMap, state: &Service) -> bool { |
| 160 | let origin = format!("http://127.0.0.1:{}", state.descriptor.port); |
| 161 | let host = format!("127.0.0.1:{}", state.descriptor.port); |
| 162 | if headers.get("host").and_then(|v| v.to_str().ok()) != Some(host.as_str()) { |
| 163 | return false; |
| 164 | } |
| 165 | if headers |
| 166 | .get("origin") |
| 167 | .is_some_and(|o| o.as_bytes() != origin.as_bytes()) |
| 168 | { |
| 169 | return false; |
| 170 | } |
| 171 | let bearer = format!("Bearer {}", state.descriptor.token); |
| 172 | let cookie = format!("cw_pet={}", state.descriptor.token); |
| 173 | headers |
| 174 | .get("authorization") |
| 175 | .is_some_and(|v| v.as_bytes() == bearer.as_bytes()) |
| 176 | || headers |
| 177 | .get("cookie") |
| 178 | .and_then(|v| v.to_str().ok()) |
| 179 | .is_some_and(|v| v.split(';').any(|v| v.trim() == cookie)) |
| 180 | } |
| 181 | |
| 182 | fn answer(result: Result<Value, String>) -> Response { |
| 183 | match result { |
| 184 | Ok(value) => Json(value).into_response(), |
| 185 | Err(error) => (StatusCode::CONFLICT, Json(json!({"error": error}))).into_response(), |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | async fn submit( |
| 190 | state: &Service, |
| 191 | make: impl FnOnce(tokio::sync::oneshot::Sender<Result<Value, String>>) -> Work, |
| 192 | ) -> Response { |
| 193 | let (tx, rx) = tokio::sync::oneshot::channel(); |
| 194 | if state.tx.try_send(make(tx)).is_err() { |
| 195 | return StatusCode::SERVICE_UNAVAILABLE.into_response(); |
| 196 | } |
| 197 | match tokio::time::timeout(WORK_REPLY_TIMEOUT, rx).await { |
| 198 | Ok(Ok(result)) => answer(result), |
| 199 | _ => StatusCode::SERVICE_UNAVAILABLE.into_response(), |
| 200 | } |
| 201 | } |
| 202 | |
| 203 | async fn frame( |
| 204 | State(state): State<Service>, |
| 205 | headers: HeaderMap, |
| 206 | axum::extract::Query(query): axum::extract::Query<BTreeMap<String, String>>, |
| 207 | ) -> Response { |
| 208 | if !authorized(&headers, &state) { |
| 209 | return StatusCode::UNAUTHORIZED.into_response(); |
| 210 | } |
| 211 | let Ok(frames) = state.frames.lock() else { |
| 212 | return StatusCode::SERVICE_UNAVAILABLE.into_response(); |
| 213 | }; |
| 214 | let found = if let Some(tick) = query.get("tick") { |
| 215 | let wanted = tick.parse::<u64>().ok(); |
| 216 | frames |
| 217 | .iter() |
| 218 | .find(|f| wanted.is_some() && f["tick"].as_u64() == wanted) |
| 219 | } else { |
| 220 | frames.back() |
| 221 | }; |
| 222 | found.map_or_else( |
| 223 | || StatusCode::NOT_FOUND.into_response(), |
| 224 | |v| Json(v.clone()).into_response(), |
| 225 | ) |
| 226 | } |
| 227 | |
| 228 | async fn action( |
| 229 | State(state): State<Service>, |
| 230 | headers: HeaderMap, |
| 231 | Json(request): Json<Request>, |
| 232 | ) -> Response { |
| 233 | if !authorized(&headers, &state) { |
| 234 | return StatusCode::UNAUTHORIZED.into_response(); |
| 235 | } |
| 236 | submit(&state, |reply| Work::Action(request, reply)).await |
| 237 | } |
| 238 | async fn producer( |
| 239 | State(state): State<Service>, |
| 240 | headers: HeaderMap, |
| 241 | Json(request): Json<Producer>, |
| 242 | ) -> Response { |
| 243 | if !authorized(&headers, &state) { |
| 244 | return StatusCode::UNAUTHORIZED.into_response(); |
| 245 | } |
| 246 | submit(&state, |reply| Work::Producer(request, reply)).await |
| 247 | } |
| 248 | async fn audio( |
| 249 | State(state): State<Service>, |
| 250 | headers: HeaderMap, |
| 251 | Json(request): Json<AudioLease>, |
| 252 | ) -> Response { |
| 253 | if !authorized(&headers, &state) { |
| 254 | return StatusCode::UNAUTHORIZED.into_response(); |
| 255 | } |
| 256 | submit(&state, |reply| Work::Audio(request, reply)).await |
| 257 | } |
| 258 | async fn export(State(state): State<Service>, headers: HeaderMap) -> Response { |
| 259 | if !authorized(&headers, &state) { |
| 260 | return StatusCode::UNAUTHORIZED.into_response(); |
| 261 | } |
| 262 | submit(&state, Work::Export).await |
| 263 | } |
| 264 | async fn attach(State(state): State<Service>, headers: HeaderMap) -> Response { |
| 265 | if !authorized(&headers, &state) { |
| 266 | return StatusCode::UNAUTHORIZED.into_response(); |
| 267 | } |
| 268 | ( |
| 269 | [ |
| 270 | ( |
| 271 | "set-cookie", |
| 272 | format!( |
| 273 | "cw_pet={}; HttpOnly; SameSite=Strict; Path=/", |
| 274 | state.descriptor.token |
| 275 | ), |
| 276 | ), |
| 277 | ("cache-control", "no-store".into()), |
| 278 | ], |
| 279 | Json(json!({"identity":state.descriptor.identity})), |
| 280 | ) |
| 281 | .into_response() |
| 282 | } |
| 283 | |
| 284 | /// Starts only on an explicit pet command or when a view first attaches. |
| 285 | /// The lifetime lock is held across HTTP serving, ticks and all checkpoints. |
| 286 | // `pet serve` is its own console process, never inside the alt-screen: its |
| 287 | // startup line and stop reason are the operator's only output. |
| 288 | #[allow(clippy::print_stdout, clippy::print_stderr)] |
| 289 | pub fn serve(root: PathBuf, requested_port: u16) -> anyhow::Result<()> { |
| 290 | std::fs::create_dir_all(&root)?; |
| 291 | #[cfg(unix)] |
| 292 | { |
| 293 | use std::os::unix::fs::{MetadataExt, PermissionsExt}; |
| 294 | let metadata = std::fs::symlink_metadata(&root)?; |
| 295 | if !metadata.is_dir() || metadata.uid() != unsafe { libc::geteuid() } { |
| 296 | anyhow::bail!("Pet directory must belong to this user and cannot be a link"); |
| 297 | } |
| 298 | std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700))?; |
| 299 | } |
| 300 | let lock_path = WorkspaceFile::open(&root, Path::new("owner.lock"), true)?; |
| 301 | let original = lock_path.open_update(true, false)?; |
| 302 | let mut lifetime = fd_lock::RwLock::new(original.try_clone()?); |
| 303 | let _guard = lifetime |
| 304 | .try_write() |
| 305 | .map_err(|_| anyhow::anyhow!("Another pet owner is running"))?; |
| 306 | let mut store = Store::at(&root)?; |
| 307 | let previous = store.load()?; |
| 308 | let mut saved = if let Some(text) = previous { |
| 309 | let value: Saved = serde_json::from_str(&text)?; |
| 310 | if value.version != 1 |
| 311 | || uuid::Uuid::parse_str(&value.identity).is_err() |
| 312 | || value.token.len() != 64 |
| 313 | || !value.token.bytes().all(|b| b.is_ascii_hexdigit()) |
| 314 | || !valid_source(&value.source) |
| 315 | || value.clients.len() > MAX_CLIENTS |
| 316 | || !value.appearance.valid() |
| 317 | { |
| 318 | anyhow::bail!("Invalid shared habitat; the existing file was kept"); |
| 319 | } |
| 320 | value |
| 321 | } else { |
| 322 | Saved { |
| 323 | version: 1, |
| 324 | identity: uuid::Uuid::new_v4().to_string(), |
| 325 | token: format!( |
| 326 | "{}{}", |
| 327 | uuid::Uuid::new_v4().simple(), |
| 328 | uuid::Uuid::new_v4().simple() |
| 329 | ), |
| 330 | port: requested_port, |
| 331 | source: "unattached".into(), |
| 332 | source_revision: 0, |
| 333 | cursor: 0, |
| 334 | clients: BTreeMap::new(), |
| 335 | recording: Value::Null, |
| 336 | appearance: Appearance::default(), |
| 337 | } |
| 338 | }; |
| 339 | let listener = std::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, saved.port))?; |
| 340 | saved.port = listener.local_addr()?.port(); |
| 341 | listener.set_nonblocking(true)?; |
| 342 | let descriptor = Descriptor { |
| 343 | version: 1, |
| 344 | port: saved.port, |
| 345 | token: saved.token.clone(), |
| 346 | identity: saved.identity.clone(), |
| 347 | }; |
| 348 | let connection = WorkspaceFile::open(&root, Path::new("connection.json"), true)?; |
| 349 | let epoch = uuid::Uuid::new_v4().to_string(); |
| 350 | let (tx, rx) = mpsc::sync_channel(128); |
| 351 | let frames = Arc::new(Mutex::new(VecDeque::new())); |
| 352 | let service = Service { |
| 353 | descriptor: descriptor.clone(), |
| 354 | frames: frames.clone(), |
| 355 | tx, |
| 356 | }; |
| 357 | let (ready_tx, ready_rx) = mpsc::sync_channel(1); |
| 358 | let (ended_tx, ended_rx) = tokio::sync::oneshot::channel(); |
| 359 | let owner = std::thread::Builder::new() |
| 360 | .name("pet-owner-world".into()) |
| 361 | .spawn(move || { |
| 362 | let result = run_world( |
| 363 | saved, store, rx, frames, &epoch, &lock_path, &original, ready_tx, |
| 364 | ); |
| 365 | if let Err(error) = result { |
| 366 | eprintln!("Shared pet stopped: {error}"); |
| 367 | } |
| 368 | let _ = ended_tx.send(()); |
| 369 | })?; |
| 370 | ready_rx |
| 371 | .recv_timeout(Duration::from_secs(15))? |
| 372 | .map_err(anyhow::Error::msg)?; |
| 373 | connection.replace(&serde_json::to_vec(&descriptor)?)?; |
| 374 | let router = Router::new() |
| 375 | .route("/", get(|| async { Html(include_str!("shared.html")) })) |
| 376 | .route("/pet-native.js", get(|| async { ([("content-type", "text/javascript")], include_str!("pet-native.js")) })) |
| 377 | .route("/whale-points.tsv", get(|| async { include_str!("../ambient_life/whale-points.tsv") })) |
| 378 | .route("/v1/frame", get(frame)).route("/v1/action", post(action)) |
| 379 | .route("/v1/producer", post(producer)).route("/v1/audio", post(audio)) |
| 380 | .route("/v1/export", get(export)).route("/v1/attach", post(attach)) |
| 381 | .layer(DefaultBodyLimit::max(MAX_REQUEST_BYTES)) |
| 382 | .layer(axum::middleware::from_fn(|request: axum::extract::Request, next: axum::middleware::Next| async move { |
| 383 | let mut response = next.run(request).await; |
| 384 | for (key, value) in [("cache-control", "no-store"), ("referrer-policy", "no-referrer"), ("x-content-type-options", "nosniff"), ("content-security-policy", "default-src 'none'; script-src 'self' 'unsafe-inline'; style-src 'unsafe-inline'; connect-src 'self'; img-src 'self' blob:; frame-ancestors 'none'; base-uri 'none'; form-action 'none'")] { |
| 385 | response.headers_mut().insert(axum::http::HeaderName::from_static(key), axum::http::HeaderValue::from_static(value)); |
| 386 | } |
| 387 | response |
| 388 | })).with_state(service); |
| 389 | println!( |
| 390 | "Shared pet {} listening on 127.0.0.1:{}", |
| 391 | descriptor.identity, descriptor.port |
| 392 | ); |
| 393 | let rt = tokio::runtime::Builder::new_multi_thread() |
| 394 | .worker_threads(2) |
| 395 | .enable_all() |
| 396 | .build()?; |
| 397 | rt.block_on(async { |
| 398 | axum::serve(tokio::net::TcpListener::from_std(listener)?, router) |
| 399 | .with_graceful_shutdown(async { |
| 400 | tokio::select! { _=tokio::signal::ctrl_c()=>{}, _=ended_rx=>{} } |
| 401 | }) |
| 402 | .await |
| 403 | })?; |
| 404 | let _ = owner.join(); |
| 405 | Ok(()) |
| 406 | } |
| 407 | |
| 408 | fn run_world( |
| 409 | mut saved: Saved, |
| 410 | mut store: Store, |
| 411 | rx: mpsc::Receiver<Work>, |
| 412 | frames: Arc<Mutex<VecDeque<Value>>>, |
| 413 | epoch: &str, |
| 414 | lock_path: &WorkspaceFile, |
| 415 | original: &std::fs::File, |
| 416 | ready: mpsc::SyncSender<Result<(), String>>, |
| 417 | ) -> anyhow::Result<()> { |
| 418 | let runtime = Runtime::new()?; |
| 419 | runtime.set_memory_limit(64 * 1024 * 1024); |
| 420 | runtime.set_max_stack_size(2 * 1024 * 1024); |
| 421 | let deadline = Arc::new(Mutex::new(Instant::now() + Duration::from_secs(10))); |
| 422 | let limit = deadline.clone(); |
| 423 | runtime.set_interrupt_handler(Some(Box::new(move || { |
| 424 | limit.lock().map_or(true, |d| Instant::now() > *d) |
| 425 | }))); |
| 426 | let context = Context::full(&runtime)?; |
| 427 | context |
| 428 | .with(|ctx| -> rquickjs::Result<()> { |
| 429 | let points: Vec<Vec<f64>> = include_str!("../ambient_life/whale-points.tsv") |
| 430 | .lines() |
| 431 | .map(|line| { |
| 432 | line.split_whitespace() |
| 433 | .filter_map(|n| n.parse().ok()) |
| 434 | .collect() |
| 435 | }) |
| 436 | .collect(); |
| 437 | ctx.globals() |
| 438 | .set("points", serde_json::to_string(&points).unwrap())?; |
| 439 | ctx.eval::<(), _>(include_bytes!("pet-native.js").as_slice())?; |
| 440 | ctx.eval::<(), _>("globalThis.pet = new PetNative(points, '', '[]', true)")?; |
| 441 | if !saved.recording.is_null() { |
| 442 | ctx.globals().set("saved", saved.recording.to_string())?; |
| 443 | ctx.eval::<(), _>( |
| 444 | "pet.restoreRecording(saved); pet.resumeEngine(); delete globalThis.saved", |
| 445 | )?; |
| 446 | } |
| 447 | Ok(()) |
| 448 | }) |
| 449 | .map_err(|_| anyhow::anyhow!("Shared habitat could not restore; its file was kept"))?; |
| 450 | let mut producer: Option<(String, u64, Instant, String)> = None; |
| 451 | let mut audio: Option<(String, Instant)> = None; |
| 452 | let mut playback: Option<(super::audio::Output, super::audio_cursor::AudioCursor)> = None; |
| 453 | let mut audio_error = false; |
| 454 | let mut waiting = false; |
| 455 | let mut last_save = Instant::now(); |
| 456 | let mut storage_error = false; |
| 457 | let origin = Instant::now(); |
| 458 | let initial_time: f64 = context.with(|ctx| ctx.eval("JSON.parse(pet.snapshot()).timeMs"))?; |
| 459 | let mut last = origin; |
| 460 | let mut ticks = 0u64; |
| 461 | let mut measurements = VecDeque::<f64>::new(); |
| 462 | save(&context, &mut saved, &mut store)?; |
| 463 | let _ = ready.send(Ok(())); |
| 464 | loop { |
| 465 | *deadline |
| 466 | .lock() |
| 467 | .map_err(|_| anyhow::anyhow!("Clock lock failed"))? = |
| 468 | Instant::now() + Duration::from_secs(5); |
| 469 | let now = Instant::now(); |
| 470 | if !same_file(&lock_path.open_update(false, false)?, original)? { |
| 471 | anyhow::bail!("Owner lock was replaced"); |
| 472 | } |
| 473 | if producer |
| 474 | .as_ref() |
| 475 | .is_some_and(|(_, _, seen, _)| now.duration_since(*seen) > LEASE) |
| 476 | { |
| 477 | producer = None; |
| 478 | waiting = false; |
| 479 | context.with(|ctx| ctx.eval::<(), _>("pet.disconnectEngine()"))?; |
| 480 | } |
| 481 | if audio |
| 482 | .as_ref() |
| 483 | .is_some_and(|(_, seen)| now.duration_since(*seen) > LEASE) |
| 484 | { |
| 485 | audio = None; |
| 486 | } |
| 487 | let elapsed = now.duration_since(last).as_secs_f64(); |
| 488 | if elapsed >= 1.0 / 30.0 { |
| 489 | // A suspended machine advances a bounded amount and marks a gap; |
| 490 | // offline wall time never invents activity or historical sound. |
| 491 | let count = (elapsed * 30.0).floor().min(3.0) as u64; |
| 492 | if elapsed > 0.25 { |
| 493 | // A delayed clock tick loses observation coverage, not the |
| 494 | // producer's transport sequence. Keep its lease until LEASE |
| 495 | // expires above; ordinary scheduler stalls must not reject |
| 496 | // the next valid packet as an unknown producer. |
| 497 | waiting = false; |
| 498 | context.with(|ctx| ctx.eval::<(), _>("pet.disconnectEngine()"))?; |
| 499 | } |
| 500 | let started = Instant::now(); |
| 501 | context.with(|ctx| -> rquickjs::Result<()> { |
| 502 | ctx.globals() |
| 503 | .set("at", initial_time + (ticks + count) as f64 * 1000.0 / 30.0)?; |
| 504 | ctx.globals().set("waiting", waiting)?; |
| 505 | ctx.eval::<(), _>("pet.advanceEngine(at,true,waiting)") |
| 506 | })?; |
| 507 | ticks += count; |
| 508 | if audio.is_none() { |
| 509 | playback = None; |
| 510 | } |
| 511 | if audio.is_some() && playback.is_none() { |
| 512 | match super::audio::Output::start() { |
| 513 | Ok(output) => { |
| 514 | let cursor = super::audio_cursor::AudioCursor::new( |
| 515 | output.target(), |
| 516 | initial_time + ticks as f64 * 1000.0 / 30.0, |
| 517 | ); |
| 518 | playback = Some((output, cursor)); |
| 519 | audio_error = false; |
| 520 | } |
| 521 | Err(_) => { |
| 522 | audio = None; |
| 523 | audio_error = true; |
| 524 | } |
| 525 | } |
| 526 | } |
| 527 | if let Some((output, cursor)) = &mut playback { |
| 528 | let target = output.target(); |
| 529 | if output.failed() |
| 530 | || context |
| 531 | .with(|ctx| { |
| 532 | cursor.present( |
| 533 | &ctx, |
| 534 | &target, |
| 535 | initial_time + ticks as f64 * 1000.0 / 30.0, |
| 536 | ) |
| 537 | }) |
| 538 | .is_err() |
| 539 | { |
| 540 | context.with(|ctx| { |
| 541 | let _ = ctx.catch(); |
| 542 | }); |
| 543 | playback = None; |
| 544 | audio = None; |
| 545 | audio_error = true; |
| 546 | } |
| 547 | } |
| 548 | last = if elapsed > 0.25 { |
| 549 | now |
| 550 | } else { |
| 551 | last + Duration::from_secs_f64(count as f64 / 30.0) |
| 552 | }; |
| 553 | let text: String = context.with(|ctx| ctx.eval("pet.presentation()"))?; |
| 554 | let mut frame: Value = serde_json::from_str(&text)?; |
| 555 | frame["version"] = json!(1); |
| 556 | frame["identity"] = json!(saved.identity); |
| 557 | frame["epoch"] = json!(epoch); |
| 558 | frame["tick"] = |
| 559 | json!((frame["timeMs"].as_f64().unwrap_or(0.0) * 30.0 / 1000.0).round() as u64); |
| 560 | frame["cursor"] = json!(saved.cursor); |
| 561 | frame["source"] = json!(saved.source); |
| 562 | frame["sourceRevision"] = json!(saved.source_revision); |
| 563 | bind_activity(&mut frame, &saved.source, saved.cursor); |
| 564 | // Presentation material keeps missing coverage legible. The core |
| 565 | // pigment, score, particle digest and recording are unchanged. |
| 566 | for key in ["", "still"] { |
| 567 | let pose = if key.is_empty() { |
| 568 | &mut frame |
| 569 | } else { |
| 570 | &mut frame[key] |
| 571 | }; |
| 572 | let hollow = |
| 573 | producer.is_none() || pose["style"]["hollow"].as_bool().unwrap_or(true); |
| 574 | if hollow { |
| 575 | pose["style"]["hollow"] = json!(true); |
| 576 | pose["style"]["r"] = json!(153); |
| 577 | pose["style"]["g"] = json!(176); |
| 578 | pose["style"]["b"] = json!(184); |
| 579 | pose["style"]["alpha"] = |
| 580 | json!(0.68 * pose["state"]["lit"].as_f64().unwrap_or(1.0).max(0.25)); |
| 581 | } |
| 582 | } |
| 583 | for key in ["", "still"] { |
| 584 | let pose = if key.is_empty() { |
| 585 | &mut frame |
| 586 | } else { |
| 587 | &mut frame[key] |
| 588 | }; |
| 589 | if !saved.appearance.event_colors { |
| 590 | for (key, value) in ["r", "g", "b"].into_iter().zip(saved.appearance.particle) { |
| 591 | pose["style"][key] = json!(value); |
| 592 | } |
| 593 | } |
| 594 | pose["style"]["alpha"] = json!( |
| 595 | (pose["style"]["alpha"].as_f64().unwrap_or(0.5) * saved.appearance.brightness) |
| 596 | .clamp(0.0, 1.0) |
| 597 | ); |
| 598 | } |
| 599 | frame["appearance"] = json!(saved.appearance); |
| 600 | frame["producerConnected"] = json!(producer.is_some()); |
| 601 | frame["storageAvailable"] = json!(!storage_error); |
| 602 | frame["audioOwner"] = json!(audio.as_ref().map(|(id, _)| id)); |
| 603 | frame["audioUnavailable"] = json!(audio_error); |
| 604 | measurements.push_back(started.elapsed().as_secs_f64() * 1000.0); |
| 605 | if measurements.len() > 300 { |
| 606 | measurements.pop_front(); |
| 607 | } |
| 608 | frame["performance"] = json!({"worldHz":30,"frames":ticks,"uptimeSeconds":origin.elapsed().as_secs_f64(),"workMs":measurements.back()}); |
| 609 | let mut output = frames |
| 610 | .lock() |
| 611 | .map_err(|_| anyhow::anyhow!("Frame lock failed"))?; |
| 612 | output.push_back(frame); |
| 613 | if output.len() > 16 { |
| 614 | output.pop_front(); |
| 615 | } |
| 616 | } |
| 617 | if last_save.elapsed() >= Duration::from_secs(1) { |
| 618 | storage_error = save(&context, &mut saved, &mut store).is_err(); |
| 619 | last_save = Instant::now(); |
| 620 | } |
| 621 | let work = match rx.recv_timeout(Duration::from_millis(2)) { |
| 622 | Ok(work) => work, |
| 623 | Err(mpsc::RecvTimeoutError::Timeout) => continue, |
| 624 | Err(mpsc::RecvTimeoutError::Disconnected) => { |
| 625 | save(&context, &mut saved, &mut store)?; |
| 626 | return Ok(()); |
| 627 | } |
| 628 | }; |
| 629 | match work { |
| 630 | Work::Export(reply) => { |
| 631 | let result = context |
| 632 | .with(|ctx| super::persistence::export_recording(&ctx, true)) |
| 633 | .map_err(|_| "Recording export failed".to_owned()) |
| 634 | .and_then(|text| { |
| 635 | serde_json::from_slice(&text).map_err(|_| "Recording export failed".into()) |
| 636 | }); |
| 637 | let _ = reply.send(result); |
| 638 | } |
| 639 | Work::Audio(request, reply) => { |
| 640 | let result = if !valid_client(&request.client) { |
| 641 | Err("Invalid view identity".into()) |
| 642 | } else if !request.enabled { |
| 643 | if audio.as_ref().is_some_and(|(id, _)| *id == request.client) { |
| 644 | audio = None; |
| 645 | } |
| 646 | Ok(json!({"granted":false})) |
| 647 | } else if audio.as_ref().is_none_or(|(id, _)| *id == request.client) { |
| 648 | audio = Some((request.client, Instant::now())); |
| 649 | Ok(json!({"granted":true})) |
| 650 | } else { |
| 651 | Ok(json!({"granted":false})) |
| 652 | }; |
| 653 | let _ = reply.send(result); |
| 654 | } |
| 655 | Work::Producer(request, reply) => { |
| 656 | let result = (|| -> Result<Value, String> { |
| 657 | if request.identity != saved.identity |
| 658 | || request.epoch != epoch |
| 659 | || !valid_client(&request.client) |
| 660 | || request.source != saved.source |
| 661 | || request.source_revision != saved.source_revision |
| 662 | || request.events.len() > 64 |
| 663 | { |
| 664 | return Err( |
| 665 | "Source changed; attach at the current frame and discard stale input" |
| 666 | .into(), |
| 667 | ); |
| 668 | } |
| 669 | use sha2::{Digest, Sha256}; |
| 670 | let hash=Sha256::digest(serde_json::to_vec(&json!({"seq":request.seq,"events":request.events,"waiting":request.waiting})).unwrap()).iter().map(|b|format!("{b:02x}")).collect::<String>(); |
| 671 | if let Some((id, seq, _, prior_hash)) = &producer { |
| 672 | if *id != request.client { |
| 673 | return Err("This source already has a producer".into()); |
| 674 | } |
| 675 | if request.seq == *seq && *prior_hash == hash { |
| 676 | return Ok( |
| 677 | json!({"seq":seq,"cursor":saved.cursor,"duplicate":true,"durable":false}), |
| 678 | ); |
| 679 | } |
| 680 | if request.seq != seq + 1 { |
| 681 | producer = None; |
| 682 | waiting = false; |
| 683 | context |
| 684 | .with(|ctx| ctx.eval::<(), _>("pet.disconnectEngine()")) |
| 685 | .map_err(|_| "Coverage reset failed")?; |
| 686 | return Err("Producer gap; reconnect without historical input".into()); |
| 687 | } |
| 688 | } else if request.seq != 0 || !request.events.is_empty() { |
| 689 | return Err( |
| 690 | "Begin a producer lease with sequence zero and no historical events" |
| 691 | .into(), |
| 692 | ); |
| 693 | } |
| 694 | context |
| 695 | .with(|ctx| -> rquickjs::Result<()> { |
| 696 | ctx.globals() |
| 697 | .set("events", serde_json::to_string(&request.events).unwrap())?; |
| 698 | ctx.eval::<(), _>( |
| 699 | "pet.observeEngineBatch(events, JSON.parse(pet.snapshot()).timeMs)", |
| 700 | ) |
| 701 | }) |
| 702 | .map_err(|_| { |
| 703 | context.with(|ctx| { |
| 704 | let _ = ctx.catch(); |
| 705 | }); |
| 706 | "Invalid Engine metadata" |
| 707 | })?; |
| 708 | if !request.events.is_empty() || waiting != request.waiting { |
| 709 | saved.cursor += 1; |
| 710 | } |
| 711 | producer = Some((request.client, request.seq, Instant::now(), hash)); |
| 712 | waiting = request.waiting; |
| 713 | Ok(json!({"seq":request.seq,"cursor":saved.cursor,"durable":false})) |
| 714 | })(); |
| 715 | let _ = reply.send(result); |
| 716 | } |
| 717 | Work::Action(request, reply) => { |
| 718 | let result = (|| -> Result<Value, String> { |
| 719 | use sha2::{Digest, Sha256}; |
| 720 | if request.identity != saved.identity |
| 721 | || !valid_client(&request.client) |
| 722 | || request.seq == 0 |
| 723 | { |
| 724 | return Err("Invalid pet action identity".into()); |
| 725 | } |
| 726 | let hash = Sha256::digest(serde_json::to_vec(&request).unwrap()) |
| 727 | .iter() |
| 728 | .map(|b| format!("{b:02x}")) |
| 729 | .collect::<String>(); |
| 730 | if let Some(receipt) = saved.clients.get(&request.client) { |
| 731 | if request.seq == receipt.seq && hash == receipt.hash { |
| 732 | return Ok(json!({"cursor":receipt.cursor,"duplicate":true})); |
| 733 | } |
| 734 | if request.seq != receipt.seq + 1 { |
| 735 | return Err("Action sequence is stale or has a gap".into()); |
| 736 | } |
| 737 | } else if request.seq != 1 || saved.clients.len() >= MAX_CLIENTS { |
| 738 | return Err( |
| 739 | "Action client is unknown or the retained client limit was reached" |
| 740 | .into(), |
| 741 | ); |
| 742 | } |
| 743 | if request.source_revision != saved.source_revision { |
| 744 | return Err( |
| 745 | "Source changed; review the current source before interacting".into(), |
| 746 | ); |
| 747 | } |
| 748 | // Check storage before accepting an action. A failed commit |
| 749 | // is rolled back together with its idempotency receipt. |
| 750 | save(&context, &mut saved, &mut store) |
| 751 | .map_err(|_| "Pet storage unavailable; action was not accepted")?; |
| 752 | let before = serde_json::to_string(&saved).unwrap(); |
| 753 | match request.action { |
| 754 | Action::Appearance { appearance } => { |
| 755 | if !appearance.valid() { |
| 756 | return Err("Invalid appearance range".into()); |
| 757 | } |
| 758 | saved.appearance = appearance; |
| 759 | } |
| 760 | Action::Interact { food, x, y } => { |
| 761 | if !x.is_finite() || !y.is_finite() || x.abs() > 1.0 || y.abs() > 1.0 { |
| 762 | return Err("Invalid interaction coordinates".into()); |
| 763 | } |
| 764 | context |
| 765 | .with(|ctx| -> rquickjs::Result<()> { |
| 766 | ctx.globals().set("x", x)?; |
| 767 | ctx.globals().set("y", y)?; |
| 768 | ctx.globals() |
| 769 | .set("kind", if food { "food" } else { "attention" })?; |
| 770 | ctx.eval::<(), _>("pet.interact(kind,x,y)") |
| 771 | }) |
| 772 | .map_err(|_| "Interaction failed")?; |
| 773 | } |
| 774 | Action::Select { source } => { |
| 775 | if !valid_source(&source) { |
| 776 | return Err("Invalid source identity".into()); |
| 777 | } |
| 778 | saved.source = source; |
| 779 | saved.source_revision += 1; |
| 780 | producer = None; |
| 781 | waiting = false; |
| 782 | context |
| 783 | .with(|ctx| ctx.eval::<(), _>("pet.disconnectEngine()")) |
| 784 | .map_err(|_| "Source disconnect failed")?; |
| 785 | } |
| 786 | } |
| 787 | saved.cursor += 1; |
| 788 | saved.clients.insert( |
| 789 | request.client, |
| 790 | Receipt { |
| 791 | seq: request.seq, |
| 792 | hash, |
| 793 | cursor: saved.cursor, |
| 794 | }, |
| 795 | ); |
| 796 | if save(&context, &mut saved, &mut store).is_err() { |
| 797 | saved = serde_json::from_str(&before).unwrap(); |
| 798 | context.with(|ctx| -> rquickjs::Result<()> { ctx.globals().set("rollback",saved.recording.to_string())?;ctx.eval::<(),_>("pet.restoreRecording(rollback); pet.disconnectEngine(); delete globalThis.rollback") }).map_err(|_| "Storage rollback failed")?; |
| 799 | producer = None; |
| 800 | waiting = false; |
| 801 | storage_error = true; |
| 802 | return Err("Pet storage unavailable; action was not accepted".into()); |
| 803 | } |
| 804 | storage_error = false; |
| 805 | last_save = Instant::now(); |
| 806 | Ok(json!({"cursor":saved.cursor,"duplicate":false})) |
| 807 | })(); |
| 808 | let _ = reply.send(result); |
| 809 | } |
| 810 | } |
| 811 | } |
| 812 | } |
| 813 | |
| 814 | fn save(context: &Context, saved: &mut Saved, store: &mut Store) -> anyhow::Result<()> { |
| 815 | let (text, segment) = context.with(|ctx| -> rquickjs::Result<(String, bool)> { |
| 816 | let segment: bool = ctx.eval("pet.needsSegment()")?; |
| 817 | Ok(( |
| 818 | ctx.eval(if segment { |
| 819 | "pet.prepareSegment()" |
| 820 | } else { |
| 821 | "pet.recording(true)" |
| 822 | })?, |
| 823 | segment, |
| 824 | )) |
| 825 | })?; |
| 826 | saved.recording = serde_json::from_str(&text)?; |
| 827 | let archive = |
| 828 | if segment { |
| 829 | Some(context.with(|ctx| { |
| 830 | ctx.eval::<String, _>("JSON.stringify(JSON.parse(pet.recording(true)))") |
| 831 | })?) |
| 832 | } else { |
| 833 | None |
| 834 | }; |
| 835 | let tick = saved.recording["checkpoint"]["tick"].as_u64().unwrap_or(0); |
| 836 | store.save_archived( |
| 837 | &serde_json::to_string(saved)?, |
| 838 | archive.as_ref().map(|a| (a.as_bytes(), tick)), |
| 839 | )?; |
| 840 | if segment { |
| 841 | context.with(|ctx| ctx.eval::<(), _>("pet.commitSegment()"))?; |
| 842 | } |
| 843 | Ok(()) |
| 844 | } |
| 845 | |
| 846 | /// Bind the reducer's account-neutral owner projection to the session and |
| 847 | /// cursor this owner serves. The reducer never knows either; `Scene::valid` |
| 848 | /// rejects a frame whose projection disagrees with them. |
| 849 | fn bind_activity(frame: &mut Value, source: &str, cursor: u64) { |
| 850 | if let Some(activity) = frame["activity"].as_object_mut() { |
| 851 | let session = if source == "unattached" { |
| 852 | Value::Null |
| 853 | } else { |
| 854 | json!(source) |
| 855 | }; |
| 856 | activity.insert("sessionId".into(), session); |
| 857 | activity.insert("cursor".into(), json!(cursor)); |
| 858 | } |
| 859 | } |
| 860 | |
| 861 | #[cfg(test)] |
| 862 | mod tests { |
| 863 | use super::*; |
| 864 | use codewhale_protocol::engine_owner::{ |
| 865 | EngineOwnerProjection, OwnerActivityKind, OwnerFreshness, OwnerPresence, |
| 866 | }; |
| 867 | |
| 868 | /// Drive the committed `pet-native.js` bundle (the one the owner and the |
| 869 | /// worker actually evaluate) and read back its projection the way the |
| 870 | /// live client does: strict schema, then the contract's invariants. |
| 871 | fn projection(ctx: &rquickjs::Ctx<'_>, source: &str, cursor: u64) -> EngineOwnerProjection { |
| 872 | let text: String = ctx.eval("pet.presentation()").unwrap(); |
| 873 | let mut frame: Value = serde_json::from_str(&text).unwrap(); |
| 874 | bind_activity(&mut frame, source, cursor); |
| 875 | let activity: EngineOwnerProjection = |
| 876 | serde_json::from_value(frame["activity"].clone()).expect("bundle matches the contract"); |
| 877 | assert!(activity.is_valid(), "invalid projection: {activity:?}"); |
| 878 | assert_eq!(activity.cursor, cursor); |
| 879 | activity |
| 880 | } |
| 881 | |
| 882 | #[test] |
| 883 | fn committed_bundle_projection_satisfies_the_owner_contract() { |
| 884 | let runtime = Runtime::new().unwrap(); |
| 885 | let context = Context::full(&runtime).unwrap(); |
| 886 | context.with(|ctx| { |
| 887 | let points: Vec<Vec<f64>> = include_str!("../ambient_life/whale-points.tsv") |
| 888 | .lines() |
| 889 | .map(|line| line.split_whitespace().filter_map(|n| n.parse().ok()).collect()) |
| 890 | .collect(); |
| 891 | ctx.globals() |
| 892 | .set("points", serde_json::to_string(&points).unwrap()) |
| 893 | .unwrap(); |
| 894 | ctx.eval::<(), _>(include_bytes!("pet-native.js").as_slice()) |
| 895 | .unwrap(); |
| 896 | ctx.eval::<(), _>("globalThis.pet = new PetNative(points, '', '[]', true)") |
| 897 | .unwrap(); |
| 898 | |
| 899 | let unattached = projection(&ctx, "unattached", 0); |
| 900 | assert_eq!(unattached.freshness, OwnerFreshness::Missing); |
| 901 | assert_eq!(unattached.session_id, None); |
| 902 | |
| 903 | let feed = |events: &str, at: f64| { |
| 904 | ctx.globals().set("events", events).unwrap(); |
| 905 | ctx.globals().set("at", at).unwrap(); |
| 906 | ctx.eval::<(), _>("pet.observeEngineBatch(events, at); pet.advanceEngine(at + 50, true, false)") |
| 907 | .unwrap(); |
| 908 | }; |
| 909 | feed( |
| 910 | r#"[{"event":"turn_started","turn_id":"turn-1"}, |
| 911 | {"event":"operation_activity_started","span_id":"call-1","activity_kind":"computer"}]"#, |
| 912 | 0.0, |
| 913 | ); |
| 914 | let working = projection(&ctx, "session-a", 3); |
| 915 | assert_eq!(working.session_id.as_deref(), Some("session-a")); |
| 916 | assert_eq!(working.freshness, OwnerFreshness::Fresh); |
| 917 | assert_eq!(working.authoritative_presence, OwnerPresence::Working); |
| 918 | assert_eq!(working.activity_kind, Some(OwnerActivityKind::Computer)); |
| 919 | |
| 920 | // The old tool-name feed is refused at the boundary. |
| 921 | ctx.globals() |
| 922 | .set( |
| 923 | "events", |
| 924 | r#"[{"event":"tool_call_started","tool_call_id":"a","tool_name":"exec_shell"}]"#, |
| 925 | ) |
| 926 | .unwrap(); |
| 927 | assert!(ctx.eval::<(), _>("pet.observeEngineBatch(events, 100)").is_err()); |
| 928 | let _ = ctx.catch(); |
| 929 | |
| 930 | feed( |
| 931 | r#"[{"event":"operation_activity_completed","span_id":"call-1","activity_kind":"computer","outcome":"succeeded"}, |
| 932 | {"event":"turn_complete","turn_id":"turn-1","turn_outcome":"completed"}]"#, |
| 933 | 200.0, |
| 934 | ); |
| 935 | let done = projection(&ctx, "session-a", 4); |
| 936 | assert_eq!(done.authoritative_presence, OwnerPresence::Done); |
| 937 | assert_eq!(done.done_effect_id.as_deref(), Some("turn-1")); |
| 938 | |
| 939 | feed( |
| 940 | r#"[{"event":"turn_started","turn_id":"turn-2"}, |
| 941 | {"event":"turn_complete","turn_id":"turn-2","turn_outcome":"interrupted"}]"#, |
| 942 | 400.0, |
| 943 | ); |
| 944 | let interrupted = projection(&ctx, "session-a", 5); |
| 945 | assert_ne!(interrupted.authoritative_presence, OwnerPresence::Done); |
| 946 | assert_eq!(interrupted.done_effect_id, None); |
| 947 | }); |
| 948 | } |
| 949 | } |
| 950 |