返回 CodeWhale
owner.rs
根目录 / crates / tui / src / tui / pet_watch / owner.rs
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
950 lines RUST