返回 CodeWhale
live.rs
根目录 / crates / tui / src / tui / pet_watch / live.rs
1 //! Live view transport replacing the app-local QuickJS Worker. Only immutable
2 //! projections cross back to Ratatui; network, raster and encoding stay here.
3 use super::{graphics, owner};
4 use codewhale_protocol::engine_owner::EngineOwnerProjection;
5 use serde::{Deserialize, Serialize};
6 use serde_json::{Value, json};
7 use std::{
8 io::{self, Read},
9 path::PathBuf,
10 sync::{Arc, Mutex, mpsc},
11 time::{Duration, Instant},
12 };
13
14 #[derive(Clone, Deserialize, Serialize)]
15 pub struct Style {
16 pub r: f64,
17 pub g: f64,
18 pub b: f64,
19 pub alpha: f64,
20 pub hollow: bool,
21 pub channel: String,
22 pub arch: String,
23 }
24 #[derive(Clone, Deserialize, Serialize)]
25 pub struct Pose {
26 pub points: Vec<[f64; 2]>,
27 pub style: Style,
28 pub state: Value,
29 }
30 pub type Activity = EngineOwnerProjection;
31 #[derive(Clone, Deserialize, Serialize)]
32 #[serde(rename_all = "camelCase")]
33 pub struct Scene {
34 pub version: u8,
35 pub identity: String,
36 pub epoch: String,
37 pub tick: u64,
38 pub cursor: u64,
39 pub source: String,
40 pub source_revision: u64,
41 pub time_ms: f64,
42 pub digest: String,
43 pub behaviour: String,
44 pub producer_connected: bool,
45 pub storage_available: bool,
46 pub audio_owner: Option<String>,
47 pub audio_unavailable: bool,
48 pub points: Vec<[f64; 2]>,
49 pub style: Style,
50 pub state: Value,
51 pub still: Pose,
52 #[serde(default)]
53 pub appearance: super::appearance::Appearance,
54 /// Owner activity is an overlay on the body, never a reason to drop the
55 /// frame. A long-lived `pet serve` owner can predate this binary (the
56 /// activity shape changed without a frame version bump), so an activity
57 /// this client cannot parse reads as none and the pet keeps rendering.
58 #[serde(default, deserialize_with = "lenient_activity")]
59 pub activity: Option<Activity>,
60 }
61 fn lenient_activity<'de, D: serde::Deserializer<'de>>(
62 deserializer: D,
63 ) -> Result<Option<Activity>, D::Error> {
64 Ok(Option::<Value>::deserialize(deserializer)?
65 .and_then(|value| serde_json::from_value(value).ok()))
66 }
67 impl Scene {
68 /// Drop an activity that is invalid or bound to another cursor or
69 /// session. Rejecting the whole frame instead would clear queued events,
70 /// drop the producer lease and loop every shared view on reconnect.
71 fn settle_activity(&mut self) {
72 let source = (self.source != "unattached").then_some(self.source.as_str());
73 if self.activity.as_ref().is_some_and(|activity| {
74 !activity.is_valid()
75 || activity.cursor != self.cursor
76 || activity.session_id.as_deref() != source
77 }) {
78 self.activity = None;
79 }
80 }
81 fn valid(&self) -> bool {
82 self.version == 1
83 && self.time_ms.is_finite()
84 && self.points.len() == 980
85 && self.still.points.len() == 980
86 && self
87 .points
88 .iter()
89 .chain(&self.still.points)
90 .flatten()
91 .all(|p| p.is_finite() && p.abs() <= 8.0)
92 }
93 }
94 #[derive(Clone)]
95 pub struct Presentation {
96 pub client: String,
97 pub scene: Scene,
98 pub cells: Vec<u8>,
99 pub image: Option<Vec<u8>>,
100 pub width: u16,
101 pub height: u16,
102 pub created: Instant,
103 pub frame_changed: Instant,
104 pub render_ms: f64,
105 pub bytes: usize,
106 }
107 #[derive(Clone, Default)]
108 pub struct View {
109 pub width: u16,
110 pub height: u16,
111 pub cell_width: f64,
112 pub cell_height: f64,
113 pub motion: bool,
114 pub pixels: bool,
115 pub visible: bool,
116 pub waiting: bool,
117 pub sound: bool,
118 }
119 #[derive(Clone)]
120 pub enum Command {
121 Observe(String),
122 Select,
123 Browser,
124 Window,
125 Export,
126 }
127 pub enum Notice {
128 Exported(PathBuf),
129 Message(String),
130 /// The companion cannot be reached: the view thread stopped, or a frame
131 /// fetch failed and it is reconnecting. Only this marks the pet offline.
132 Unreachable(String),
133 }
134 /// Producer events queued for the next `/v1/producer` post. The owner caps a
135 /// batch at 64 events *and* its request body at [`owner::MAX_REQUEST_BYTES`];
136 /// a batch that exceeds either is a coverage gap, never a post the owner must
137 /// refuse (a 64-event batch of 16 KiB events is far over the body limit).
138 #[derive(Default)]
139 struct Batch {
140 events: Vec<Value>,
141 bytes: usize,
142 }
143
144 impl Batch {
145 const MAX_EVENTS: usize = 64;
146 /// Room left for the post's envelope (identity, epoch, source, sequence).
147 const MAX_EVENT_BYTES: usize = owner::MAX_REQUEST_BYTES - 4 * 1024;
148
149 /// Queue one event. `Ok(false)` means it does not fit: the caller drops
150 /// the batch and restarts the lease rather than posting a partial one.
151 fn push(&mut self, text: &str) -> io::Result<bool> {
152 // `+ 1` counts the separating comma in the posted JSON array.
153 let bytes = self.bytes + text.len() + 1;
154 if self.events.len() >= Self::MAX_EVENTS || bytes > Self::MAX_EVENT_BYTES {
155 return Ok(false);
156 }
157 self.events.push(serde_json::from_str(text)?);
158 self.bytes = bytes;
159 Ok(true)
160 }
161
162 fn clear(&mut self) {
163 self.events.clear();
164 self.bytes = 0;
165 }
166 }
167
168 pub struct Worker {
169 pub tx: mpsc::SyncSender<Command>,
170 pub latest: Arc<Mutex<Option<Presentation>>>,
171 pub view: Arc<Mutex<View>>,
172 pub notices: mpsc::Receiver<Notice>,
173 }
174
175 pub struct Client {
176 pub descriptor: owner::Descriptor,
177 http: reqwest::blocking::Client,
178 }
179 impl Client {
180 pub fn connect() -> io::Result<Self> {
181 let root = owner::directory()?;
182 let http = crate::tls::reqwest_blocking_client_builder()
183 .no_proxy()
184 .connect_timeout(Duration::from_millis(500))
185 .timeout(Duration::from_secs(2))
186 .build()
187 .map_err(io::Error::other)?;
188 let attempt = || -> io::Result<Self> {
189 Ok(Self {
190 descriptor: owner::descriptor(&root)?,
191 http: http.clone(),
192 })
193 };
194 if let Ok(client) = attempt()
195 && client.get("/v1/frame").is_ok()
196 {
197 return Ok(client);
198 }
199 #[cfg(not(test))]
200 {
201 use std::process::{Command as Process, Stdio};
202 let mut process = Process::new(std::env::current_exe()?);
203 process
204 .args(["pet", "serve"])
205 .stdin(Stdio::null())
206 .stdout(Stdio::null())
207 .stderr(Stdio::null());
208 #[cfg(unix)]
209 {
210 use std::os::unix::process::CommandExt;
211 unsafe {
212 process.pre_exec(|| {
213 if libc::setsid() < 0 {
214 return Err(io::Error::last_os_error());
215 }
216 Ok(())
217 });
218 }
219 }
220 #[cfg(windows)]
221 {
222 use std::os::windows::process::CommandExt;
223 process.creation_flags(0x08000000 | 0x00000008);
224 }
225 let mut child = process.spawn()?;
226 std::thread::spawn(move || {
227 let _ = child.wait();
228 });
229 }
230 let began = Instant::now();
231 while began.elapsed() < Duration::from_secs(5) {
232 if let Ok(client) = attempt()
233 && client.get("/v1/frame").is_ok()
234 {
235 return Ok(client);
236 }
237 std::thread::sleep(Duration::from_millis(50));
238 }
239 Err(io::Error::other(
240 "Shared pet unavailable. Run codewhale pet serve; existing recordings are preserved.",
241 ))
242 }
243 pub fn get(&self, path: &str) -> io::Result<Value> {
244 self.request(path, None)
245 }
246 pub fn post(&self, path: &str, body: &Value) -> io::Result<Value> {
247 self.request(path, Some(body))
248 }
249 fn request(&self, path: &str, body: Option<&Value>) -> io::Result<Value> {
250 let url = format!("http://127.0.0.1:{}{path}", self.descriptor.port);
251 let mut request = if let Some(body) = body {
252 self.http.post(url).json(body)
253 } else {
254 self.http.get(url)
255 };
256 if path == "/v1/export" {
257 // The owner serializes the recording on its world thread for up
258 // to its own work bound; giving up sooner abandons an export the
259 // owner still completes. Other calls keep the 2 s client bound.
260 //
261 // Known limitation (U06-m4): export runs synchronously in the
262 // owner's single QuickJS world, so the world answers no other work
263 // until it finishes. Keeping it responsive needs an incremental
264 // export in pet-native.js; aligning this wait is only the client
265 // half.
266 request = request.timeout(owner::WORK_REPLY_TIMEOUT + Duration::from_secs(1));
267 }
268 let response = request
269 .bearer_auth(&self.descriptor.token)
270 .send()
271 .map_err(io::Error::other)?;
272 let status = response.status();
273 let success = status.is_success();
274 let bound = if path == "/v1/export" {
275 super::persistence::MAX_EXPORT_BYTES
276 } else {
277 8 * 1024 * 1024
278 };
279 let mut bytes = Vec::new();
280 response.take(bound as u64 + 1).read_to_end(&mut bytes)?;
281 if bytes.len() > bound {
282 return Err(io::Error::other("Pet response exceeds its bound"));
283 }
284 let value: Value = serde_json::from_slice(&bytes)?;
285 if !success {
286 let message = value["error"]
287 .as_str()
288 .unwrap_or("Shared pet connection failed");
289 return Err(io::Error::new(
290 if status == reqwest::StatusCode::CONFLICT && !message.contains("storage") {
291 io::ErrorKind::InvalidInput
292 } else {
293 io::ErrorKind::Other
294 },
295 message,
296 ));
297 }
298 Ok(value)
299 }
300 pub fn open_browser(&self) -> io::Result<()> {
301 let url = format!(
302 "http://127.0.0.1:{}/#{}",
303 self.descriptor.port, self.descriptor.token
304 );
305 #[cfg(all(unix, not(target_os = "macos")))]
306 {
307 use std::process::{Command as Process, Stdio};
308 // `webbrowser` resolves `$BROWSER` and `xdg-open` through the
309 // environment and `PATH`, where a workspace could shadow the
310 // launcher and receive this URL's owner token. Use the launcher
311 // from a trusted system prefix only (U06-06).
312 let launcher = crate::notify::audio::trusted_system_executable("xdg-open")?;
313 let mut child = Process::new(launcher)
314 .arg(&url)
315 .stdin(Stdio::null())
316 .stdout(Stdio::null())
317 .stderr(Stdio::null())
318 .spawn()?;
319 // Some launchers stay in the foreground with the browser; reap it
320 // off this worker thread so the pet view keeps painting.
321 std::thread::Builder::new()
322 .name("pet-browser-launcher".into())
323 .spawn(move || {
324 let _ = child.wait();
325 })?;
326 Ok(())
327 }
328 // macOS (LaunchServices) and Windows (ShellExecute) resolve the
329 // default browser through the OS, not `PATH`.
330 #[cfg(not(all(unix, not(target_os = "macos"))))]
331 {
332 webbrowser::open(&url).map_err(io::Error::other)
333 }
334 }
335 pub fn open_window(&self) -> io::Result<()> {
336 #[cfg(target_os = "macos")]
337 {
338 use std::process::{Command as Process, Stdio};
339 // The system launcher by absolute path: a bare `open` resolves
340 // through `PATH`, where a workspace entry could shadow it.
341 let mut process = Process::new("/usr/bin/open");
342 if let Some(path) = std::env::var_os("CODEWHALE_PET_APP") {
343 process.arg(path);
344 } else {
345 process.args(["-a", "Codewhale Pet"]);
346 }
347 process
348 .args(["--args", "--companion"])
349 .stdin(Stdio::null())
350 .stdout(Stdio::null())
351 .stderr(Stdio::null());
352 if process.status()?.success() {
353 return Ok(());
354 }
355 Err(io::Error::other(
356 "Build or install the Codewhale Pet app to open its companion window",
357 ))
358 }
359 #[cfg(not(target_os = "macos"))]
360 {
361 Err(io::Error::other(
362 "The native companion window is currently available on macOS",
363 ))
364 }
365 }
366 }
367 impl Worker {
368 pub fn start(session: Option<String>) -> io::Result<Self> {
369 let (tx, rx) = mpsc::sync_channel(128);
370 let (latest, (notices_tx, notices)) = (Arc::new(Mutex::new(None)), mpsc::sync_channel(16));
371 let output = latest.clone();
372 let view = Arc::new(Mutex::new(View {
373 width: 40,
374 height: 8,
375 ..View::default()
376 }));
377 let settings = view.clone();
378 std::thread::Builder::new()
379 .name("pet-view".into())
380 .spawn(move || {
381 if let Err(e) = run(rx, &output, &notices_tx, &settings, session) {
382 let _ = notices_tx.try_send(Notice::Unreachable(e.to_string()));
383 }
384 })?;
385 Ok(Self {
386 tx,
387 latest,
388 notices,
389 view,
390 })
391 }
392 }
393 fn run(
394 rx: mpsc::Receiver<Command>,
395 output: &Mutex<Option<Presentation>>,
396 notices: &mpsc::SyncSender<Notice>,
397 settings: &Mutex<View>,
398 session: Option<String>,
399 ) -> io::Result<()> {
400 use sha2::{Digest, Sha256};
401 let source = session.as_ref().map(|s| {
402 format!(
403 "session:{}",
404 Sha256::digest(s.as_bytes())
405 .iter()
406 .take(10)
407 .map(|b| format!("{b:02x}"))
408 .collect::<String>()
409 )
410 });
411 let mut renderer = graphics::Renderer::default();
412 let id = uuid::Uuid::new_v4().to_string();
413 let mut client = Client::connect()?;
414
415 let mut scene: Option<Scene> = None;
416 let mut previous = None;
417 let mut changed = Instant::now();
418 let mut fetched = Instant::now() - Duration::from_secs(1);
419 let mut encoded = Instant::now() - Duration::from_secs(1);
420 let mut events = Batch::default();
421 let mut producer_seq = None;
422 let mut sequence = 0u64;
423 let mut action: Option<Value> = None;
424 let mut last_action = Instant::now() - Duration::from_secs(1);
425 let mut last_produce = Instant::now();
426 let mut last_audio = Instant::now();
427 let mut last_failure = Instant::now() - Duration::from_secs(10);
428 loop {
429 let view = settings
430 .lock()
431 .map_err(|_| io::Error::other("Pet view lock failed"))?
432 .clone();
433 match rx.recv_timeout(Duration::from_millis(2)) {
434 Ok(command) => match command {
435 Command::Observe(text) => {
436 if !events.push(&text)? {
437 // Overflow is a coverage gap, never a truncated batch:
438 // restart the lease from sequence zero.
439 events.clear();
440 producer_seq = None;
441 }
442 }
443 Command::Select => {
444 if let (Some(s), Some(source)) = (&scene, &source)
445 && action.is_none()
446 {
447 action = Some(
448 json!({"identity":s.identity,"client":id,"seq":sequence+1,"source_revision":s.source_revision,"action":{"kind":"select","source":source}}),
449 );
450 } else {
451 let _ = notices.try_send(Notice::Message(
452 "Save this session and wait for the pet connection before selecting its source.".into(),
453 ));
454 }
455 }
456 Command::Browser => {
457 if let Err(e) = client.open_browser() {
458 let _ = notices.try_send(Notice::Message(e.to_string()));
459 }
460 }
461 Command::Window => {
462 if let Err(e) = client.open_window() {
463 let _ = notices.try_send(Notice::Message(e.to_string()));
464 }
465 }
466 Command::Export => {
467 let result = client.get("/v1/export").and_then(|r| {
468 super::persistence::export(
469 session.as_deref().ok_or_else(|| {
470 io::Error::other("Save the terminal session before exporting")
471 })?,
472 &serde_json::to_vec(&r)?,
473 )
474 });
475 let _ = notices.try_send(match result {
476 Ok(path) => Notice::Exported(path),
477 Err(e) => Notice::Message(e.to_string()),
478 });
479 }
480 },
481 Err(mpsc::RecvTimeoutError::Disconnected) => {
482 let _ = client.post("/v1/audio", &json!({"client":id,"enabled":false}));
483 return Ok(());
484 }
485 Err(mpsc::RecvTimeoutError::Timeout) => {}
486 }
487 if fetched.elapsed() >= Duration::from_millis(if view.visible { 30 } else { 400 }) {
488 match client
489 .get("/v1/frame")
490 .and_then(|value| serde_json::from_value::<Scene>(value).map_err(io::Error::other))
491 .map(|mut next| {
492 next.settle_activity();
493 next
494 }) {
495 Ok(next) if next.valid() => {
496 if scene
497 .as_ref()
498 .is_none_or(|s| s.epoch != next.epoch || s.identity != next.identity)
499 {
500 previous = None;
501 changed = Instant::now();
502 producer_seq = None;
503 events.clear();
504 } else if scene.as_ref().is_some_and(|s| s.tick != next.tick) {
505 previous = scene.clone();
506 changed = Instant::now();
507 }
508 if next.source == "unattached"
509 && let Some(source) = &source
510 && action.is_none()
511 {
512 action = Some(
513 json!({"identity":next.identity,"client":id,"seq":sequence+1,"source_revision":next.source_revision,"action":{"kind":"select","source":source}}),
514 );
515 }
516 scene = Some(next);
517 fetched = Instant::now();
518 }
519 _ => {
520 fetched = Instant::now();
521 events.clear();
522 producer_seq = None;
523 if last_failure.elapsed() > Duration::from_secs(3) {
524 last_failure = Instant::now();
525 let _ =
526 notices.try_send(Notice::Unreachable("Shared pet reconnecting".into()));
527 if let Ok(next) = Client::connect() {
528 client = next;
529 }
530 }
531 continue;
532 }
533 }
534 }
535 let Some(s) = &scene else { continue };
536 if last_action.elapsed() >= Duration::from_millis(250)
537 && let Some(pending) = &action
538 {
539 last_action = Instant::now();
540 match client.post("/v1/action", pending) {
541 Ok(_) => {
542 let selected = pending["action"]["kind"] == "select";
543 sequence += 1;
544 action = None;
545 if selected {
546 producer_seq = None;
547 events.clear();
548 }
549 }
550 Err(e) => {
551 if e.kind() == io::ErrorKind::InvalidInput {
552 action = None;
553 }
554 if last_failure.elapsed() > Duration::from_secs(3) {
555 last_failure = Instant::now();
556 let _ = notices.try_send(Notice::Message(e.to_string()));
557 }
558 }
559 }
560 }
561 if source.as_deref() == Some(s.source.as_str())
562 && last_produce.elapsed() >= Duration::from_millis(100)
563 {
564 let seq = producer_seq.map_or(0, |seq| seq + 1);
565 if seq == 0 {
566 events.clear();
567 }
568 let body = json!({"identity":s.identity,"epoch":s.epoch,"client":id,"source":s.source,"source_revision":s.source_revision,"seq":seq,"waiting":view.waiting,"events":events.events});
569 producer_seq = client.post("/v1/producer", &body).ok().map(|_| seq);
570 events.clear();
571 last_produce = Instant::now();
572 } else if source.as_deref() != Some(s.source.as_str()) {
573 events.clear();
574 producer_seq = None;
575 }
576 if last_audio.elapsed() >= Duration::from_millis(500) {
577 let _=client.post("/v1/audio",&json!({"client":id,"enabled":view.sound&&view.visible&&changed.elapsed()<Duration::from_millis(500)}));
578 last_audio = Instant::now();
579 }
580 if view.visible
581 && encoded.elapsed()
582 >= Duration::from_millis(if view.motion && view.pixels {
583 16
584 } else if view.motion {
585 33
586 } else {
587 200
588 })
589 {
590 let began = Instant::now();
591 let mut pose = if view.motion {
592 Pose {
593 points: s.points.clone(),
594 style: s.style.clone(),
595 state: s.state.clone(),
596 }
597 } else {
598 s.still.clone()
599 };
600 if !s.producer_connected {
601 pose.style.hollow = true;
602 }
603 let fraction = if view.motion {
604 (changed.elapsed().as_secs_f64() * 30.0).clamp(0.0, 1.0)
605 } else {
606 1.0
607 };
608 renderer.set_appearance(&s.appearance);
609 let (cells, image) = renderer.render(
610 &pose,
611 previous.as_ref().filter(|_| view.motion),
612 fraction,
613 &view,
614 if view.motion { s.time_ms / 1000.0 } else { 0.0 },
615 )?;
616 let bytes = image.as_ref().map_or(0, Vec::len);
617 if let Ok(mut slot) = output.lock() {
618 *slot = Some(Presentation {
619 client: id.clone(),
620 scene: s.clone(),
621 cells,
622 image,
623 width: view.width,
624 height: view.height,
625 created: Instant::now(),
626 frame_changed: changed,
627 render_ms: began.elapsed().as_secs_f64() * 1000.0,
628 bytes,
629 });
630 }
631 encoded = began;
632 }
633 }
634 }
635
636 #[cfg(test)]
637 mod tests {
638 use super::*;
639
640 /// U06-03: a producer batch is bounded by the owner's request body limit,
641 /// not only by its event count, so the worker never posts a body the
642 /// owner must refuse (64 events of 16 KiB each is ~1 MiB).
643 #[test]
644 fn producer_batch_fits_the_owner_body_limit() {
645 let mut batch = Batch::default();
646 let large = format!("{{\"event\":\"x\",\"id\":\"{}\"}}", "a".repeat(16_000));
647 let mut accepted = 0;
648 while batch.push(&large).unwrap() {
649 accepted += 1;
650 }
651 assert!(accepted > 0 && accepted < Batch::MAX_EVENTS, "{accepted}");
652 let posted = serde_json::to_vec(&json!({"seq": u64::MAX, "events": batch.events})).unwrap();
653 assert!(posted.len() <= owner::MAX_REQUEST_BYTES, "{}", posted.len());
654
655 batch.clear();
656 let small = "{\"event\":\"tool_call_heartbeat\"}";
657 for _ in 0..Batch::MAX_EVENTS {
658 assert!(batch.push(small).unwrap());
659 }
660 assert!(!batch.push(small).unwrap(), "the event count still bounds");
661 }
662 }
663
663 lines RUST