返回 CodeWhale
computer_display.rs
根目录 / crates / tui / src / runtime_api / computer_display.rs
1 //! `/v1/computer/*` — the Engine's view of the computer it runs on.
2 //!
3 //! ARCHITECTURE §3.3 (research/computer): TigerVNC `Xvnc` serves RFB on a
4 //! Unix socket only (`/run/cw/vnc.sock`, group `cw-display`). This module is
5 //! the single external path to it:
6 //!
7 //! - `GET /v1/computer/display` upgrades to a WebSocket that carries raw RFB
8 //! 3.8 bytes as binary frames. The Engine completes the upstream handshake
9 //! itself (security None on the socket) and offers only None downstream,
10 //! because the WebSocket is already authenticated.
11 //! - Server-to-client bytes pass through verbatim.
12 //! - Client-to-server bytes go through [`ClientParser`], a length-tracked,
13 //! fail-closed parser that runs in its **own task**, so a panic in it ends
14 //! one display connection and never a turn. Only message types 0, 2, 3, 4,
15 //! 5, 6, 150 and 251 are allowed; any other type closes the stream.
16 //! Input (4 key, 5 pointer, 6 clipboard, 251 resize) is dropped unless the
17 //! connection's principal holds the control lease.
18 //! - The lease is human-only and lives here, in the Engine. Agents read it
19 //! (`GET /v1/computer`) to refuse input tools while a human drives.
20 //!
21 //! Auth: every route here authenticates itself (it is merged outside the
22 //! `/v1` route layer) because the display WebSocket also accepts a
23 //! single-use `?ticket=` for browser clients that cannot set headers. Tickets
24 //! are redacted by [`redact_query_secrets`] wherever a URI is logged. The
25 //! Engine never trusts the peer address (S0 Q3: `/proxy` peers arrive as
26 //! `10.0.0.2`, not loopback) — only a token.
27 //!
28 //! Human keystrokes are never logged or put in events: events carry time
29 //! spans and counts only. Frames are never events.
30
31 use codewhale_core::secret_eq::constant_time_eq;
32 use std::collections::{HashMap, VecDeque};
33 use std::path::PathBuf;
34 use std::sync::Arc;
35 use std::sync::atomic::{AtomicU64, Ordering};
36 use std::time::{Duration, Instant};
37
38 use axum::extract::ws::{
39 CloseFrame, Message, WebSocket, WebSocketUpgrade, rejection::WebSocketUpgradeRejection,
40 };
41 use axum::extract::{Path, Query, State};
42 use axum::http::{HeaderMap, StatusCode, Uri, header};
43 use axum::response::{IntoResponse, Response};
44 use axum::routing::{delete, get, post};
45 use axum::{Json, Router};
46 use chrono::{DateTime, Utc};
47 use futures_util::{SinkExt, StreamExt};
48 use serde::{Deserialize, Serialize};
49 use serde_json::{Value, json};
50 use sha2::{Digest, Sha256};
51 use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
52
53 const DISPLAY_SOCKET_ENV: &str = "CODEWHALE_COMPUTER_DISPLAY_SOCKET";
54 const DISPLAY_IDLE_ENV: &str = "CODEWHALE_COMPUTER_DISPLAY_IDLE_SECS";
55 const DEFAULT_DISPLAY_SOCKET: &str = "/run/cw/vnc.sock";
56 /// §5: the Engine closes idle displays after 10 minutes.
57 const DEFAULT_IDLE: Duration = Duration::from_secs(600);
58 /// A lease with no human input for this long expires.
59 const LEASE_IDLE_TTL: Duration = Duration::from_secs(300);
60 /// §2.3: client tokens last at most one hour.
61 const CLIENT_TOKEN_MAX_TTL_SECS: u64 = 3600;
62 const CLIENT_TOKEN_MIN_TTL_SECS: u64 = 60;
63 const CLIENT_TOKEN_MAX_ACTIVE: usize = 64;
64 const DISPLAY_TICKET_TTL: Duration = Duration::from_secs(30);
65 const DISPLAY_TICKET_MAX_ACTIVE: usize = 64;
66 const EVENT_LOG_CAP: usize = 512;
67 const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
68 const SUPERVISOR_TICK: Duration = Duration::from_secs(5);
69 const RFB_VERSION_38: &[u8; 12] = b"RFB 003.008\n";
70 const DEVICE_ID_MAX_BYTES: usize = 128;
71
72 /// Query keys whose values are secrets and must never reach a log line.
73 const SECRET_QUERY_KEYS: &[&str] = &[
74 "ticket",
75 super::mobile::MOBILE_STREAM_TICKET_QUERY,
76 "token",
77 "access_token",
78 ];
79
80 // ---------------------------------------------------------------------------
81 // State
82 // ---------------------------------------------------------------------------
83
84 /// Who a request speaks for. `Owner` is the master runtime token (or an
85 /// Engine started with explicit insecure no-auth); `Client` is a device
86 /// token minted through `POST /v1/auth/client-tokens`.
87 #[derive(Debug, Clone, PartialEq, Eq)]
88 pub(super) enum Principal {
89 Owner,
90 Client { token_id: String, device_id: String },
91 }
92
93 impl Principal {
94 fn holder(&self) -> String {
95 match self {
96 Principal::Owner => "owner".to_string(),
97 Principal::Client { device_id, .. } => format!("device:{device_id}"),
98 }
99 }
100
101 fn device_id(&self) -> Option<&str> {
102 match self {
103 Principal::Owner => None,
104 Principal::Client { device_id, .. } => Some(device_id),
105 }
106 }
107 }
108
109 struct ClientToken {
110 id: String,
111 device_id: String,
112 label: Option<String>,
113 created_at: DateTime<Utc>,
114 expires_at: DateTime<Utc>,
115 }
116
117 struct DisplayTicket {
118 principal: Principal,
119 expires: Instant,
120 }
121
122 struct Lease {
123 principal: Principal,
124 acquired_at: DateTime<Utc>,
125 acquired_instant: Instant,
126 last_activity: Instant,
127 input_events: u64,
128 }
129
130 #[derive(Debug, Clone, Serialize)]
131 pub(super) struct ComputerEvent {
132 pub seq: u64,
133 #[serde(rename = "type")]
134 pub kind: String,
135 pub at: DateTime<Utc>,
136 pub data: Value,
137 }
138
139 struct Inner {
140 // Only the Unix display socket is dialed; other platforms report it absent.
141 #[cfg_attr(not(unix), allow(dead_code))]
142 socket_path: PathBuf,
143 idle_close: Duration,
144 lease_ttl: Duration,
145 lease: parking_lot::Mutex<Option<Lease>>,
146 client_tokens: parking_lot::Mutex<HashMap<[u8; 32], ClientToken>>,
147 tickets: parking_lot::Mutex<HashMap<[u8; 32], DisplayTicket>>,
148 events: parking_lot::Mutex<VecDeque<ComputerEvent>>,
149 next_seq: AtomicU64,
150 next_connection: AtomicU64,
151 attached: AtomicU64,
152 }
153
154 /// Engine-side computer state: display socket, control lease, client tokens,
155 /// display tickets and the `computer.*` event log.
156 #[derive(Clone)]
157 pub(crate) struct ComputerState {
158 inner: Arc<Inner>,
159 }
160
161 /// Accept a display socket path only if it is absolute and made of plain
162 /// components: no `.`/`..`, no NUL. The configured value cannot walk the
163 /// Engine out of the directory it names, and the socket-type check at use
164 /// (`display_socket_present`) refuses anything that is not a Unix socket.
165 fn validated_socket_path(raw: &str) -> Option<PathBuf> {
166 let raw = raw.trim();
167 // This is a Unix socket setting even on hosts without Unix transport.
168 // Host-native Path parsing would reject /run/... on Windows, or normalize
169 // away the dot/repeated-separator components this contract must refuse.
170 let relative = raw.strip_prefix('/')?;
171 if raw.contains(['\0', '\\'])
172 || relative
173 .split('/')
174 .any(|part| part.is_empty() || matches!(part, "." | ".."))
175 {
176 return None;
177 }
178 Some(PathBuf::from(raw))
179 }
180
181 impl ComputerState {
182 pub(crate) fn new(socket_path: PathBuf, idle_close: Duration) -> Self {
183 Self {
184 inner: Arc::new(Inner {
185 socket_path,
186 idle_close,
187 lease_ttl: LEASE_IDLE_TTL,
188 lease: parking_lot::Mutex::new(None),
189 client_tokens: parking_lot::Mutex::new(HashMap::new()),
190 tickets: parking_lot::Mutex::new(HashMap::new()),
191 events: parking_lot::Mutex::new(VecDeque::new()),
192 next_seq: AtomicU64::new(1),
193 next_connection: AtomicU64::new(1),
194 attached: AtomicU64::new(0),
195 }),
196 }
197 }
198
199 pub(crate) fn from_env() -> Self {
200 let socket = std::env::var(DISPLAY_SOCKET_ENV)
201 .ok()
202 .and_then(|raw| {
203 let checked = validated_socket_path(&raw);
204 if checked.is_none() && !raw.trim().is_empty() {
205 tracing::warn!(
206 target: "codewhale::computer",
207 "{DISPLAY_SOCKET_ENV} must be an absolute path with no `.`/`..` \
208 components; using {DEFAULT_DISPLAY_SOCKET}"
209 );
210 }
211 checked
212 })
213 .unwrap_or_else(|| PathBuf::from(DEFAULT_DISPLAY_SOCKET));
214 let idle = std::env::var(DISPLAY_IDLE_ENV)
215 .ok()
216 .and_then(|v| v.trim().parse::<u64>().ok())
217 .filter(|secs| *secs > 0)
218 .map(Duration::from_secs)
219 .unwrap_or(DEFAULT_IDLE);
220 Self::new(socket, idle)
221 }
222
223 fn emit(&self, kind: &str, data: Value) {
224 let seq = self.inner.next_seq.fetch_add(1, Ordering::Relaxed);
225 let event = ComputerEvent {
226 seq,
227 kind: kind.to_string(),
228 at: Utc::now(),
229 data,
230 };
231 tracing::info!(target: "codewhale::computer", seq, kind, "computer event");
232 let mut events = self.inner.events.lock();
233 if events.len() >= EVENT_LOG_CAP {
234 events.pop_front();
235 }
236 events.push_back(event);
237 }
238
239 pub(super) fn events_since(&self, since: u64) -> (Vec<ComputerEvent>, u64) {
240 let events = self.inner.events.lock();
241 let list: Vec<_> = events.iter().filter(|e| e.seq > since).cloned().collect();
242 let next = self
243 .inner
244 .next_seq
245 .load(Ordering::Relaxed)
246 .saturating_sub(1);
247 (list, next)
248 }
249
250 // -- tokens --------------------------------------------------------------
251
252 /// Whether `bearer` is a live (unexpired, unrevoked) client token.
253 pub(super) fn client_principal(&self, bearer: &str) -> Option<Principal> {
254 let key = hash(bearer);
255 let now = Utc::now();
256 let tokens = self.inner.client_tokens.lock();
257 let token = tokens.get(&key)?;
258 (token.expires_at > now).then(|| Principal::Client {
259 token_id: token.id.clone(),
260 device_id: token.device_id.clone(),
261 })
262 }
263
264 fn principal_is_live(&self, principal: &Principal) -> bool {
265 match principal {
266 Principal::Owner => true,
267 Principal::Client { token_id, .. } => {
268 let now = Utc::now();
269 self.inner
270 .client_tokens
271 .lock()
272 .values()
273 .any(|t| &t.id == token_id && t.expires_at > now)
274 }
275 }
276 }
277
278 fn mint_client_token(
279 &self,
280 device_id: String,
281 ttl_secs: u64,
282 label: Option<String>,
283 ) -> Result<(String, ClientTokenView), ApiErr> {
284 let now = Utc::now();
285 let mut tokens = self.inner.client_tokens.lock();
286 tokens.retain(|_, t| t.expires_at > now);
287 if tokens.len() >= CLIENT_TOKEN_MAX_ACTIVE {
288 return Err(ApiErr::new(
289 StatusCode::TOO_MANY_REQUESTS,
290 "too many active client tokens; revoke one first",
291 ));
292 }
293 let secret = format!(
294 "cwct_{}{}",
295 uuid::Uuid::new_v4().simple(),
296 uuid::Uuid::new_v4().simple()
297 );
298 let id = format!("ct_{}", uuid::Uuid::new_v4().simple());
299 let token = ClientToken {
300 id,
301 device_id,
302 label,
303 created_at: now,
304 expires_at: now + chrono::Duration::seconds(ttl_secs as i64),
305 };
306 let view = ClientTokenView::from(&token);
307 tokens.insert(hash(&secret), token);
308 Ok((secret, view))
309 }
310
311 fn revoke_client_token(&self, id: &str) -> bool {
312 let mut tokens = self.inner.client_tokens.lock();
313 let before = tokens.len();
314 tokens.retain(|_, t| t.id != id);
315 before != tokens.len()
316 }
317
318 fn list_client_tokens(&self) -> Vec<ClientTokenView> {
319 let now = Utc::now();
320 let mut tokens = self.inner.client_tokens.lock();
321 tokens.retain(|_, t| t.expires_at > now);
322 let mut list: Vec<_> = tokens.values().map(ClientTokenView::from).collect();
323 list.sort_by_key(|a| a.created_at);
324 list
325 }
326
327 // -- display tickets -----------------------------------------------------
328
329 fn mint_ticket(&self, principal: Principal) -> Result<String, ApiErr> {
330 let now = Instant::now();
331 let mut tickets = self.inner.tickets.lock();
332 tickets.retain(|_, t| t.expires > now);
333 if tickets.len() >= DISPLAY_TICKET_MAX_ACTIVE {
334 return Err(ApiErr::new(
335 StatusCode::TOO_MANY_REQUESTS,
336 "too many outstanding display tickets",
337 ));
338 }
339 let secret = format!("cwdt_{}", uuid::Uuid::new_v4().simple());
340 tickets.insert(
341 hash(&secret),
342 DisplayTicket {
343 principal,
344 expires: now + DISPLAY_TICKET_TTL,
345 },
346 );
347 Ok(secret)
348 }
349
350 /// Single use: a ticket is removed on the first redemption attempt,
351 /// whether or not it had expired.
352 fn redeem_ticket(&self, ticket: &str) -> Option<Principal> {
353 let entry = self.inner.tickets.lock().remove(&hash(ticket))?;
354 (entry.expires > Instant::now() && self.principal_is_live(&entry.principal))
355 .then_some(entry.principal)
356 }
357
358 // -- lease ---------------------------------------------------------------
359
360 /// Expire a stale lease (emitting `computer.control.expired`) and return
361 /// a snapshot of whatever lease remains.
362 fn sweep_lease(&self) -> Option<LeaseView> {
363 let mut guard = self.inner.lease.lock();
364 let expired = guard.as_ref().is_some_and(|lease| {
365 lease.last_activity.elapsed() >= self.inner.lease_ttl
366 || !self.principal_is_live(&lease.principal)
367 });
368 if expired {
369 let lease = guard.take().expect("checked above");
370 drop(guard);
371 self.emit(
372 "computer.control.expired",
373 lease_span_data(&lease, "expired"),
374 );
375 return None;
376 }
377 guard.as_ref().map(|lease| self.lease_view(lease))
378 }
379
380 fn lease_view(&self, lease: &Lease) -> LeaseView {
381 let remaining = self
382 .inner
383 .lease_ttl
384 .saturating_sub(lease.last_activity.elapsed());
385 LeaseView {
386 holder: lease.principal.holder(),
387 device_id: lease.principal.device_id().map(str::to_string),
388 acquired_at: lease.acquired_at,
389 expires_at: Utc::now() + chrono::Duration::from_std(remaining).unwrap_or_default(),
390 input_events: lease.input_events,
391 }
392 }
393
394 fn acquire(&self, principal: &Principal, force: bool) -> Result<LeaseView, LeaseView> {
395 self.sweep_lease();
396 let mut guard = self.inner.lease.lock();
397 if let Some(current) = guard.as_mut() {
398 if current.principal == *principal {
399 current.last_activity = Instant::now();
400 return Ok(self.lease_view(current));
401 }
402 if !force {
403 return Err(self.lease_view(current));
404 }
405 let previous = guard.take().expect("checked above");
406 self.emit(
407 "computer.control.released",
408 lease_span_data(&previous, "taken_over"),
409 );
410 }
411 let now = Instant::now();
412 let lease = Lease {
413 principal: principal.clone(),
414 acquired_at: Utc::now(),
415 acquired_instant: now,
416 last_activity: now,
417 input_events: 0,
418 };
419 let view = self.lease_view(&lease);
420 *guard = Some(lease);
421 drop(guard);
422 self.emit(
423 "computer.control.acquired",
424 json!({ "holder": view.holder, "device_id": view.device_id }),
425 );
426 Ok(view)
427 }
428
429 fn release(&self, principal: &Principal) -> bool {
430 let mut guard = self.inner.lease.lock();
431 if guard
432 .as_ref()
433 .is_some_and(|lease| lease.principal == *principal)
434 {
435 let lease = guard.take().expect("checked above");
436 drop(guard);
437 self.emit(
438 "computer.control.released",
439 lease_span_data(&lease, "hand_back"),
440 );
441 true
442 } else {
443 false
444 }
445 }
446
447 fn holds_lease(&self, principal: &Principal) -> bool {
448 self.inner.lease.lock().as_ref().is_some_and(|lease| {
449 lease.principal == *principal && lease.last_activity.elapsed() < self.inner.lease_ttl
450 })
451 }
452
453 fn note_input(&self, principal: &Principal, count: u64) {
454 if let Some(lease) = self.inner.lease.lock().as_mut()
455 && lease.principal == *principal
456 {
457 lease.last_activity = Instant::now();
458 lease.input_events += count;
459 }
460 }
461 }
462
463 fn lease_span_data(lease: &Lease, reason: &str) -> Value {
464 json!({
465 "holder": lease.principal.holder(),
466 "device_id": lease.principal.device_id(),
467 "reason": reason,
468 "held_ms": lease.acquired_instant.elapsed().as_millis() as u64,
469 "input_events": lease.input_events,
470 })
471 }
472
473 fn hash(secret: &str) -> [u8; 32] {
474 Sha256::digest(secret.as_bytes()).into()
475 }
476
477 #[derive(Debug, Clone, Serialize)]
478 struct LeaseView {
479 holder: String,
480 device_id: Option<String>,
481 acquired_at: DateTime<Utc>,
482 expires_at: DateTime<Utc>,
483 input_events: u64,
484 }
485
486 #[derive(Debug, Clone, Serialize)]
487 struct ClientTokenView {
488 id: String,
489 device_id: String,
490 label: Option<String>,
491 created_at: DateTime<Utc>,
492 expires_at: DateTime<Utc>,
493 }
494
495 impl From<&ClientToken> for ClientTokenView {
496 fn from(t: &ClientToken) -> Self {
497 Self {
498 id: t.id.clone(),
499 device_id: t.device_id.clone(),
500 label: t.label.clone(),
501 created_at: t.created_at,
502 expires_at: t.expires_at,
503 }
504 }
505 }
506
507 // ---------------------------------------------------------------------------
508 // Redaction
509 // ---------------------------------------------------------------------------
510
511 /// Replace the value of every secret-bearing query parameter
512 /// (`ticket`, `mobile_stream_ticket`, `token`, `access_token`) with
513 /// `redacted`. Use on any URI before it reaches a log line or a proxy.
514 pub(crate) fn redact_query_secrets(uri: &str) -> String {
515 let Some((path, query)) = uri.split_once('?') else {
516 return uri.to_string();
517 };
518 let (query, fragment) = match query.split_once('#') {
519 Some((q, f)) => (q, Some(f)),
520 None => (query, None),
521 };
522 let redacted: Vec<String> = query
523 .split('&')
524 .map(|pair| match pair.split_once('=') {
525 Some((key, _))
526 if SECRET_QUERY_KEYS
527 .iter()
528 .any(|k| k.eq_ignore_ascii_case(key)) =>
529 {
530 format!("{key}=redacted")
531 }
532 _ => pair.to_string(),
533 })
534 .collect();
535 let mut out = format!("{path}?{}", redacted.join("&"));
536 if fragment.is_some() {
537 // Fragments never reach a server, but a logged client URL could carry
538 // one (the mobile bootstrap redirect does); drop it wholesale.
539 out.push_str("#redacted");
540 }
541 out
542 }
543
544 // ---------------------------------------------------------------------------
545 // Client-to-server RFB parser
546 // ---------------------------------------------------------------------------
547
548 const MAX_ENCODINGS: usize = 64;
549 const MAX_CUT_TEXT: usize = 256 * 1024;
550 const MAX_SCREENS: usize = 16;
551
552 /// Encodings a client may ask Xvnc for. Anything else is stripped from
553 /// `SetEncodings` so the server never starts a sub-protocol (Fence, xvp,
554 /// QEMU keys, extended clipboard) whose client replies this parser would
555 /// refuse.
556 fn encoding_allowed(encoding: i32) -> bool {
557 matches!(
558 encoding,
559 0 | 1 | 2 | 5 | 7 | 16 // Raw, CopyRect, RRE, Hextile, Tight, ZRLE
560 | -223 // DesktopSize
561 | -224 // LastRect
562 | -239 // Cursor
563 | -307 // DesktopName
564 | -308 // ExtendedDesktopSize
565 | -313 // ContinuousUpdates
566 | -32..=-23 // JPEG quality
567 | -256..=-247 // compression level
568 )
569 }
570
571 #[derive(Debug, Clone, PartialEq, Eq)]
572 pub(crate) enum ParseError {
573 UnknownType(u8),
574 TooLarge { message_type: u8, len: usize },
575 }
576
577 impl std::fmt::Display for ParseError {
578 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
579 match self {
580 ParseError::UnknownType(t) => write!(f, "unknown client message type {t}"),
581 ParseError::TooLarge { message_type, len } => {
582 write!(f, "client message type {message_type} too large ({len})")
583 }
584 }
585 }
586 }
587
588 #[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
589 pub(crate) struct FeedStats {
590 pub input_forwarded: u64,
591 pub input_dropped: u64,
592 }
593
594 /// Length-tracked parser for RFB client messages after `ClientInit`.
595 /// Holds partial messages across WebSocket frames.
596 #[derive(Default)]
597 pub(crate) struct ClientParser {
598 buf: Vec<u8>,
599 }
600
601 impl ClientParser {
602 /// Feed client bytes. Complete, allowed messages are appended to `out`;
603 /// input messages are appended only when `input_allowed`. Returns an
604 /// error (and the stream must close) on any unknown or oversized message.
605 pub(crate) fn feed(
606 &mut self,
607 data: &[u8],
608 input_allowed: bool,
609 out: &mut Vec<u8>,
610 ) -> Result<FeedStats, ParseError> {
611 self.buf.extend_from_slice(data);
612 let mut stats = FeedStats::default();
613 let mut offset = 0;
614 loop {
615 let rest = &self.buf[offset..];
616 let Some(&message_type) = rest.first() else {
617 break;
618 };
619 let need = match message_type {
620 0 => Some(20),
621 2 => (rest.len() >= 4)
622 .then(|| {
623 let n = u16::from_be_bytes([rest[2], rest[3]]) as usize;
624 (n, 4 + 4 * n)
625 })
626 .map(|(n, len)| if n > MAX_ENCODINGS { usize::MAX } else { len }),
627 3 => Some(10),
628 4 => Some(8),
629 5 => Some(6),
630 6 => (rest.len() >= 8).then(|| {
631 let n = u32::from_be_bytes([rest[4], rest[5], rest[6], rest[7]]) as usize;
632 if n > MAX_CUT_TEXT { usize::MAX } else { 8 + n }
633 }),
634 150 => Some(10),
635 251 => (rest.len() >= 8).then(|| {
636 let n = rest[6] as usize;
637 if n > MAX_SCREENS {
638 usize::MAX
639 } else {
640 8 + 16 * n
641 }
642 }),
643 other => return Err(ParseError::UnknownType(other)),
644 };
645 let Some(need) = need else { break };
646 if need == usize::MAX {
647 return Err(ParseError::TooLarge {
648 message_type,
649 len: rest.len(),
650 });
651 }
652 if rest.len() < need {
653 break;
654 }
655 let message = &rest[..need];
656 match message_type {
657 2 => {
658 let kept: Vec<[u8; 4]> = message[4..]
659 .as_chunks::<4>()
660 .0
661 .iter()
662 .copied()
663 .filter(|c| encoding_allowed(i32::from_be_bytes(*c)))
664 .collect();
665 out.extend_from_slice(&[2, 0]);
666 out.extend_from_slice(&(kept.len() as u16).to_be_bytes());
667 for c in &kept {
668 out.extend_from_slice(c);
669 }
670 }
671 4 | 5 | 6 | 251 => {
672 if input_allowed {
673 out.extend_from_slice(message);
674 stats.input_forwarded += 1;
675 } else {
676 stats.input_dropped += 1;
677 }
678 }
679 _ => out.extend_from_slice(message),
680 }
681 offset += need;
682 }
683 self.buf.drain(..offset);
684 Ok(stats)
685 }
686 }
687
688 // ---------------------------------------------------------------------------
689 // Handshakes
690 // ---------------------------------------------------------------------------
691
692 #[cfg(unix)]
693 async fn read_reason<S: AsyncRead + Unpin>(s: &mut S) -> String {
694 let Ok(len) = s.read_u32().await else {
695 return String::new();
696 };
697 let mut reason = vec![0u8; (len as usize).min(1024)];
698 let _ = s.read_exact(&mut reason).await;
699 String::from_utf8_lossy(&reason).into_owned()
700 }
701
702 /// Complete the RFB 3.8 handshake with Xvnc as a client (security None) and
703 /// send a shared `ClientInit`, so each viewer gets its own connection without
704 /// disconnecting the others. After this returns, the next upstream bytes are
705 /// `ServerInit`.
706 #[cfg(unix)]
707 pub(crate) async fn upstream_handshake<S: AsyncRead + AsyncWrite + Unpin>(
708 s: &mut S,
709 ) -> Result<(), String> {
710 let mut version = [0u8; 12];
711 s.read_exact(&mut version)
712 .await
713 .map_err(|e| format!("read server version: {e}"))?;
714 if !version.starts_with(b"RFB 003.") {
715 return Err("upstream is not an RFB server".to_string());
716 }
717 s.write_all(RFB_VERSION_38)
718 .await
719 .map_err(|e| format!("write version: {e}"))?;
720 let count = s
721 .read_u8()
722 .await
723 .map_err(|e| format!("read security: {e}"))?;
724 if count == 0 {
725 return Err(format!("upstream refused: {}", read_reason(s).await));
726 }
727 let mut types = vec![0u8; count as usize];
728 s.read_exact(&mut types)
729 .await
730 .map_err(|e| format!("read security types: {e}"))?;
731 if !types.contains(&1) {
732 return Err("upstream does not offer security type None".to_string());
733 }
734 s.write_all(&[1])
735 .await
736 .map_err(|e| format!("write security: {e}"))?;
737 let result = s
738 .read_u32()
739 .await
740 .map_err(|e| format!("read security result: {e}"))?;
741 if result != 0 {
742 return Err(format!(
743 "upstream security failed: {}",
744 read_reason(s).await
745 ));
746 }
747 s.write_all(&[1])
748 .await
749 .map_err(|e| format!("write ClientInit: {e}"))?;
750 Ok(())
751 }
752
753 /// Buffered reader over the client half of the WebSocket.
754 struct WsIn<R> {
755 stream: R,
756 pending: Vec<u8>,
757 }
758
759 impl<R> WsIn<R>
760 where
761 R: futures_util::Stream<Item = Result<Message, axum::Error>> + Unpin,
762 {
763 /// Next chunk of client bytes. `Ok(None)` is a clean close; text frames
764 /// are a protocol violation (RFB is binary).
765 async fn next_chunk(&mut self) -> Result<Option<Vec<u8>>, String> {
766 if !self.pending.is_empty() {
767 return Ok(Some(std::mem::take(&mut self.pending)));
768 }
769 loop {
770 match self.stream.next().await {
771 None | Some(Ok(Message::Close(_))) => return Ok(None),
772 Some(Ok(Message::Binary(bytes))) => return Ok(Some(bytes.to_vec())),
773 Some(Ok(Message::Ping(_) | Message::Pong(_))) => continue,
774 Some(Ok(Message::Text(_))) => return Err("text frame on RFB stream".to_string()),
775 Some(Err(err)) => return Err(format!("websocket: {err}")),
776 }
777 }
778 }
779
780 async fn read_exact(&mut self, n: usize) -> Result<Vec<u8>, String> {
781 let mut acc = Vec::with_capacity(n);
782 while acc.len() < n {
783 let Some(chunk) = self.next_chunk().await? else {
784 return Err("client closed during handshake".to_string());
785 };
786 acc.extend_from_slice(&chunk);
787 }
788 self.pending = acc.split_off(n);
789 Ok(acc)
790 }
791 }
792
793 // ---------------------------------------------------------------------------
794 // Session
795 // ---------------------------------------------------------------------------
796
797 #[derive(Debug, Clone, PartialEq, Eq)]
798 enum ExitReason {
799 ClientClosed,
800 UpstreamClosed,
801 ProtocolViolation(String),
802 ParserPanic,
803 IdleClosed,
804 Revoked,
805 Error(String),
806 }
807
808 impl ExitReason {
809 fn label(&self) -> String {
810 match self {
811 ExitReason::ClientClosed => "client_closed".into(),
812 ExitReason::UpstreamClosed => "upstream_closed".into(),
813 ExitReason::ProtocolViolation(detail) => format!("protocol_violation: {detail}"),
814 ExitReason::ParserPanic => "parser_panic".into(),
815 ExitReason::IdleClosed => "idle_closed".into(),
816 ExitReason::Revoked => "revoked".into(),
817 ExitReason::Error(detail) => format!("error: {detail}"),
818 }
819 }
820
821 fn close_code(&self) -> u16 {
822 match self {
823 ExitReason::ClientClosed | ExitReason::UpstreamClosed | ExitReason::IdleClosed => 1000,
824 ExitReason::ProtocolViolation(_) => 1008,
825 ExitReason::Revoked => 4401,
826 ExitReason::ParserPanic | ExitReason::Error(_) => 1011,
827 }
828 }
829 }
830
831 struct SessionCounters {
832 last_human_input: parking_lot::Mutex<Instant>,
833 last_screen_bytes: parking_lot::Mutex<Instant>,
834 input_forwarded: AtomicU64,
835 input_dropped: AtomicU64,
836 }
837
838 /// Parser task body: client bytes → [`ClientParser`] → upstream writer.
839 async fn parser_loop<R, W>(
840 mut ws_in: WsIn<R>,
841 mut upstream: W,
842 computer: ComputerState,
843 principal: Principal,
844 counters: Arc<SessionCounters>,
845 ) -> ExitReason
846 where
847 R: futures_util::Stream<Item = Result<Message, axum::Error>> + Unpin,
848 W: AsyncWrite + Unpin,
849 {
850 let mut parser = ClientParser::default();
851 let mut out = Vec::with_capacity(4096);
852 loop {
853 let chunk = match ws_in.next_chunk().await {
854 Ok(Some(chunk)) => chunk,
855 Ok(None) => return ExitReason::ClientClosed,
856 Err(detail) => return ExitReason::ProtocolViolation(detail),
857 };
858 out.clear();
859 let allowed = computer.holds_lease(&principal);
860 let stats = match parser.feed(&chunk, allowed, &mut out) {
861 Ok(stats) => stats,
862 Err(err) => return ExitReason::ProtocolViolation(err.to_string()),
863 };
864 if stats.input_forwarded > 0 {
865 computer.note_input(&principal, stats.input_forwarded);
866 *counters.last_human_input.lock() = Instant::now();
867 counters
868 .input_forwarded
869 .fetch_add(stats.input_forwarded, Ordering::Relaxed);
870 }
871 if stats.input_dropped > 0 {
872 counters
873 .input_dropped
874 .fetch_add(stats.input_dropped, Ordering::Relaxed);
875 }
876 if !out.is_empty() && upstream.write_all(&out).await.is_err() {
877 return ExitReason::UpstreamClosed;
878 }
879 }
880 }
881
882 /// Map a parser task's join result to an exit reason. A panic inside the
883 /// parser is contained here: it ends this display connection only.
884 fn parser_exit(result: Result<ExitReason, tokio::task::JoinError>) -> ExitReason {
885 match result {
886 Ok(reason) => reason,
887 Err(err) if err.is_panic() => ExitReason::ParserPanic,
888 Err(err) => ExitReason::Error(err.to_string()),
889 }
890 }
891
892 async fn run_session<U>(
893 socket: WebSocket,
894 upstream: U,
895 computer: ComputerState,
896 principal: Principal,
897 ) where
898 U: AsyncRead + AsyncWrite + Unpin + Send + 'static,
899 {
900 let connection_id = computer
901 .inner
902 .next_connection
903 .fetch_add(1, Ordering::Relaxed);
904 let started = Instant::now();
905 let (mut ws_tx, ws_rx) = socket.split();
906 let mut ws_in = WsIn {
907 stream: ws_rx,
908 pending: Vec::new(),
909 };
910
911 // Downstream handshake: offer RFB 3.8 with security None only.
912 let handshake = async {
913 ws_tx
914 .send(Message::Binary(RFB_VERSION_38.to_vec().into()))
915 .await
916 .map_err(|e| e.to_string())?;
917 let version = ws_in.read_exact(12).await?;
918 if version != RFB_VERSION_38 {
919 return Err("client must speak RFB 003.008".to_string());
920 }
921 ws_tx
922 .send(Message::Binary(vec![1u8, 1].into()))
923 .await
924 .map_err(|e| e.to_string())?;
925 let choice = ws_in.read_exact(1).await?;
926 if choice != [1] {
927 return Err("client chose an unsupported security type".to_string());
928 }
929 ws_tx
930 .send(Message::Binary(vec![0u8, 0, 0, 0].into()))
931 .await
932 .map_err(|e| e.to_string())?;
933 // ClientInit: its shared flag is ignored; upstream is always shared.
934 ws_in.read_exact(1).await?;
935 Ok::<(), String>(())
936 };
937 if let Err(detail) = tokio::time::timeout(HANDSHAKE_TIMEOUT, handshake)
938 .await
939 .unwrap_or_else(|_| Err("handshake timed out".to_string()))
940 {
941 let _ = ws_tx
942 .send(Message::Close(Some(CloseFrame {
943 code: 1008,
944 reason: "rfb handshake failed".into(),
945 })))
946 .await;
947 tracing::info!(target: "codewhale::computer", %detail, "display handshake failed");
948 return;
949 }
950
951 computer.inner.attached.fetch_add(1, Ordering::Relaxed);
952 computer.emit(
953 "computer.display.attached",
954 json!({
955 "connection_id": connection_id,
956 "holder": principal.holder(),
957 "device_id": principal.device_id(),
958 }),
959 );
960
961 let counters = Arc::new(SessionCounters {
962 last_human_input: parking_lot::Mutex::new(Instant::now()),
963 last_screen_bytes: parking_lot::Mutex::new(Instant::now()),
964 input_forwarded: AtomicU64::new(0),
965 input_dropped: AtomicU64::new(0),
966 });
967 let (mut up_r, up_w) = tokio::io::split(upstream);
968
969 // Parser in its own task (§3.3): a panic here must not reach a turn.
970 let mut parser = tokio::spawn(parser_loop(
971 ws_in,
972 up_w,
973 computer.clone(),
974 principal.clone(),
975 counters.clone(),
976 ));
977
978 let mut tick = tokio::time::interval(SUPERVISOR_TICK);
979 tick.tick().await;
980 let mut buf = vec![0u8; 64 * 1024];
981 let reason = loop {
982 tokio::select! {
983 joined = &mut parser => break parser_exit(joined),
984 read = up_r.read(&mut buf) => match read {
985 Ok(0) | Err(_) => break ExitReason::UpstreamClosed,
986 Ok(n) => {
987 *counters.last_screen_bytes.lock() = Instant::now();
988 if ws_tx.send(Message::Binary(buf[..n].to_vec().into())).await.is_err() {
989 break ExitReason::ClientClosed;
990 }
991 }
992 },
993 _ = tick.tick() => {
994 computer.sweep_lease();
995 if !computer.principal_is_live(&principal) {
996 break ExitReason::Revoked;
997 }
998 let idle = computer.inner.idle_close;
999 let human_idle = counters.last_human_input.lock().elapsed() >= idle;
1000 let screen_idle = counters.last_screen_bytes.lock().elapsed() >= idle;
1001 if human_idle && screen_idle {
1002 break ExitReason::IdleClosed;
1003 }
1004 }
1005 }
1006 };
1007 parser.abort();
1008
1009 let _ = ws_tx
1010 .send(Message::Close(Some(CloseFrame {
1011 code: reason.close_code(),
1012 reason: reason.label().into(),
1013 })))
1014 .await;
1015 computer.inner.attached.fetch_sub(1, Ordering::Relaxed);
1016 let span = json!({
1017 "connection_id": connection_id,
1018 "holder": principal.holder(),
1019 "reason": reason.label(),
1020 "attached_ms": started.elapsed().as_millis() as u64,
1021 "input_forwarded": counters.input_forwarded.load(Ordering::Relaxed),
1022 "input_dropped": counters.input_dropped.load(Ordering::Relaxed),
1023 });
1024 if reason == ExitReason::IdleClosed {
1025 computer.emit("computer.display.idle_closed", span.clone());
1026 }
1027 computer.emit("computer.display.detached", span);
1028 }
1029
1030 // ---------------------------------------------------------------------------
1031 // HTTP
1032 // ---------------------------------------------------------------------------
1033
1034 #[derive(Clone)]
1035 struct RouteState {
1036 computer: ComputerState,
1037 runtime_token: Option<String>,
1038 }
1039
1040 struct ApiErr {
1041 status: StatusCode,
1042 message: String,
1043 extra: Option<Value>,
1044 }
1045
1046 impl ApiErr {
1047 fn new(status: StatusCode, message: impl Into<String>) -> Self {
1048 Self {
1049 status,
1050 message: message.into(),
1051 extra: None,
1052 }
1053 }
1054 fn unauthorized() -> Self {
1055 Self::new(
1056 StatusCode::UNAUTHORIZED,
1057 "runtime API bearer token required",
1058 )
1059 }
1060 }
1061
1062 impl IntoResponse for ApiErr {
1063 fn into_response(self) -> Response {
1064 let mut body = json!({
1065 "error": { "message": self.message, "status": self.status.as_u16() }
1066 });
1067 if let Some(extra) = self.extra {
1068 body["error"]["detail"] = extra;
1069 }
1070 (self.status, Json(body)).into_response()
1071 }
1072 }
1073
1074 fn bearer(headers: &HeaderMap) -> Option<&str> {
1075 headers
1076 .get(header::AUTHORIZATION)
1077 .and_then(|v| v.to_str().ok())
1078 .and_then(|raw| raw.strip_prefix("Bearer "))
1079 .or_else(|| {
1080 headers
1081 .get("x-codewhale-runtime-token")
1082 .and_then(|v| v.to_str().ok())
1083 })
1084 }
1085
1086 fn principal_from_headers(state: &RouteState, headers: &HeaderMap) -> Option<Principal> {
1087 let Some(expected) = state.runtime_token.as_deref() else {
1088 return Some(Principal::Owner);
1089 };
1090 let presented = bearer(headers)?;
1091 if constant_time_eq(presented.as_bytes(), expected.as_bytes()) {
1092 return Some(Principal::Owner);
1093 }
1094 state.computer.client_principal(presented)
1095 }
1096
1097 fn require_principal(state: &RouteState, headers: &HeaderMap) -> Result<Principal, ApiErr> {
1098 principal_from_headers(state, headers).ok_or_else(ApiErr::unauthorized)
1099 }
1100
1101 /// Routes for the computer surface, merged into the Runtime API router
1102 /// outside the `/v1` auth layer (each handler authenticates itself).
1103 pub(super) fn router<S>(computer: ComputerState, runtime_token: Option<String>) -> Router<S>
1104 where
1105 S: Clone + Send + Sync + 'static,
1106 {
1107 Router::new()
1108 .route("/v1/computer", get(computer_status))
1109 .route("/v1/computer/events", get(computer_events))
1110 .route("/v1/computer/control/acquire", post(control_acquire))
1111 .route("/v1/computer/control/release", post(control_release))
1112 .route("/v1/computer/display/tickets", post(display_ticket))
1113 .route("/v1/computer/display", get(display_ws))
1114 .route(
1115 "/v1/auth/client-tokens",
1116 get(list_client_tokens).post(create_client_token),
1117 )
1118 .route("/v1/auth/client-tokens/{id}", delete(revoke_client_token))
1119 .with_state(RouteState {
1120 computer,
1121 runtime_token,
1122 })
1123 }
1124
1125 async fn computer_status(State(state): State<RouteState>, headers: HeaderMap) -> Response {
1126 let principal = match require_principal(&state, &headers) {
1127 Ok(p) => p,
1128 Err(e) => return e.into_response(),
1129 };
1130 let lease = state.computer.sweep_lease();
1131 let you_hold = lease
1132 .as_ref()
1133 .is_some_and(|l| l.holder == principal.holder());
1134 let (_, seq) = state.computer.events_since(u64::MAX);
1135 Json(json!({
1136 "display": {
1137 "available": display_socket_present(&state.computer).await,
1138 "attached": state.computer.inner.attached.load(Ordering::Relaxed),
1139 "idle_close_seconds": state.computer.inner.idle_close.as_secs(),
1140 },
1141 "control": {
1142 "lease": lease,
1143 "you_hold_lease": you_hold,
1144 "human_driving": lease.is_some(),
1145 "lease_idle_ttl_seconds": state.computer.inner.lease_ttl.as_secs(),
1146 },
1147 "events_seq": seq,
1148 }))
1149 .into_response()
1150 }
1151
1152 async fn display_socket_present(computer: &ComputerState) -> bool {
1153 #[cfg(unix)]
1154 {
1155 use std::os::unix::fs::FileTypeExt;
1156 tokio::fs::metadata(&computer.inner.socket_path)
1157 .await
1158 .map(|m| m.file_type().is_socket())
1159 .unwrap_or(false)
1160 }
1161 #[cfg(not(unix))]
1162 {
1163 let _ = computer;
1164 false
1165 }
1166 }
1167
1168 #[derive(Deserialize)]
1169 struct EventsQuery {
1170 since: Option<u64>,
1171 }
1172
1173 async fn computer_events(
1174 State(state): State<RouteState>,
1175 headers: HeaderMap,
1176 Query(query): Query<EventsQuery>,
1177 ) -> Response {
1178 if let Err(e) = require_principal(&state, &headers) {
1179 return e.into_response();
1180 }
1181 state.computer.sweep_lease();
1182 let (events, next) = state.computer.events_since(query.since.unwrap_or(0));
1183 Json(json!({ "events": events, "next_since": next })).into_response()
1184 }
1185
1186 #[derive(Deserialize, Default)]
1187 struct AcquireBody {
1188 #[serde(default)]
1189 force: bool,
1190 }
1191
1192 async fn control_acquire(
1193 State(state): State<RouteState>,
1194 headers: HeaderMap,
1195 body: Option<Json<AcquireBody>>,
1196 ) -> Response {
1197 let principal = match require_principal(&state, &headers) {
1198 Ok(p) => p,
1199 Err(e) => return e.into_response(),
1200 };
1201 let force = body.map(|Json(b)| b.force).unwrap_or(false);
1202 match state.computer.acquire(&principal, force) {
1203 Ok(lease) => Json(json!({ "lease": lease })).into_response(),
1204 Err(current) => ApiErr {
1205 status: StatusCode::CONFLICT,
1206 message: "another client holds the control lease".to_string(),
1207 extra: Some(json!({ "lease": current })),
1208 }
1209 .into_response(),
1210 }
1211 }
1212
1213 async fn control_release(State(state): State<RouteState>, headers: HeaderMap) -> Response {
1214 let principal = match require_principal(&state, &headers) {
1215 Ok(p) => p,
1216 Err(e) => return e.into_response(),
1217 };
1218 if state.computer.release(&principal) {
1219 Json(json!({ "released": true })).into_response()
1220 } else {
1221 ApiErr::new(
1222 StatusCode::CONFLICT,
1223 "this client does not hold the control lease",
1224 )
1225 .into_response()
1226 }
1227 }
1228
1229 async fn display_ticket(State(state): State<RouteState>, headers: HeaderMap) -> Response {
1230 let principal = match require_principal(&state, &headers) {
1231 Ok(p) => p,
1232 Err(e) => return e.into_response(),
1233 };
1234 match state.computer.mint_ticket(principal) {
1235 Ok(ticket) => (
1236 StatusCode::CREATED,
1237 Json(json!({
1238 "ticket": ticket,
1239 "expires_in_seconds": DISPLAY_TICKET_TTL.as_secs(),
1240 })),
1241 )
1242 .into_response(),
1243 Err(e) => e.into_response(),
1244 }
1245 }
1246
1247 #[derive(Deserialize)]
1248 struct DisplayQuery {
1249 ticket: Option<String>,
1250 }
1251
1252 async fn display_ws(
1253 State(state): State<RouteState>,
1254 headers: HeaderMap,
1255 uri: Uri,
1256 Query(query): Query<DisplayQuery>,
1257 ws: Result<WebSocketUpgrade, WebSocketUpgradeRejection>,
1258 ) -> Response {
1259 let principal = match principal_from_headers(&state, &headers) {
1260 Some(p) => p,
1261 None => match query
1262 .ticket
1263 .as_deref()
1264 .and_then(|ticket| state.computer.redeem_ticket(ticket))
1265 {
1266 Some(p) => p,
1267 None => return ApiErr::unauthorized().into_response(),
1268 },
1269 };
1270 tracing::info!(
1271 target: "codewhale::computer",
1272 uri = %redact_query_secrets(&uri.to_string()),
1273 holder = %principal.holder(),
1274 "computer display attach"
1275 );
1276 let ws = match ws {
1277 Ok(ws) => ws,
1278 Err(rejection) => return rejection.into_response(),
1279 };
1280 let upstream = match connect_upstream(&state.computer).await {
1281 Ok(upstream) => upstream,
1282 Err(detail) => {
1283 tracing::warn!(target: "codewhale::computer", %detail, "computer display unavailable");
1284 return ApiErr::new(
1285 StatusCode::SERVICE_UNAVAILABLE,
1286 "computer display unavailable",
1287 )
1288 .into_response();
1289 }
1290 };
1291 let computer = state.computer.clone();
1292 ws.on_upgrade(move |socket| run_session(socket, upstream, computer, principal))
1293 }
1294
1295 #[cfg(unix)]
1296 async fn connect_upstream(computer: &ComputerState) -> Result<tokio::net::UnixStream, String> {
1297 let connect = async {
1298 let mut stream = tokio::net::UnixStream::connect(&computer.inner.socket_path)
1299 .await
1300 .map_err(|e| format!("connect display socket: {e}"))?;
1301 upstream_handshake(&mut stream).await?;
1302 Ok::<_, String>(stream)
1303 };
1304 tokio::time::timeout(HANDSHAKE_TIMEOUT, connect)
1305 .await
1306 .unwrap_or_else(|_| Err("display handshake timed out".to_string()))
1307 }
1308
1309 #[cfg(not(unix))]
1310 async fn connect_upstream(_computer: &ComputerState) -> Result<tokio::io::DuplexStream, String> {
1311 Err("the computer display is Unix-only".to_string())
1312 }
1313
1314 #[derive(Deserialize)]
1315 struct CreateClientTokenBody {
1316 device_id: String,
1317 ttl_seconds: Option<u64>,
1318 label: Option<String>,
1319 }
1320
1321 fn valid_device_id(id: &str) -> bool {
1322 !id.is_empty()
1323 && id.len() <= DEVICE_ID_MAX_BYTES
1324 && id
1325 .bytes()
1326 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b':'))
1327 }
1328
1329 /// Minting is master-token only: a client token can never mint another.
1330 fn require_owner(state: &RouteState, headers: &HeaderMap) -> Result<(), ApiErr> {
1331 let Some(expected) = state.runtime_token.as_deref() else {
1332 return Err(ApiErr::new(
1333 StatusCode::CONFLICT,
1334 "client tokens need Runtime API auth; this Engine runs without a token",
1335 ));
1336 };
1337 match bearer(headers) {
1338 Some(presented) if constant_time_eq(presented.as_bytes(), expected.as_bytes()) => Ok(()),
1339 Some(presented) if state.computer.client_principal(presented).is_some() => {
1340 Err(ApiErr::new(
1341 StatusCode::FORBIDDEN,
1342 "client tokens cannot manage client tokens",
1343 ))
1344 }
1345 _ => Err(ApiErr::unauthorized()),
1346 }
1347 }
1348
1349 async fn create_client_token(
1350 State(state): State<RouteState>,
1351 headers: HeaderMap,
1352 Json(body): Json<CreateClientTokenBody>,
1353 ) -> Response {
1354 if let Err(e) = require_owner(&state, &headers) {
1355 return e.into_response();
1356 }
1357 let device_id = body.device_id.trim().to_string();
1358 if !valid_device_id(&device_id) {
1359 return ApiErr::new(
1360 StatusCode::BAD_REQUEST,
1361 "device_id must be 1-128 of [A-Za-z0-9._:-]",
1362 )
1363 .into_response();
1364 }
1365 let ttl = body
1366 .ttl_seconds
1367 .unwrap_or(CLIENT_TOKEN_MAX_TTL_SECS)
1368 .clamp(CLIENT_TOKEN_MIN_TTL_SECS, CLIENT_TOKEN_MAX_TTL_SECS);
1369 let label = body
1370 .label
1371 .map(|l| l.chars().take(128).collect::<String>())
1372 .filter(|l| !l.trim().is_empty());
1373 match state.computer.mint_client_token(device_id, ttl, label) {
1374 Ok((token, view)) => (
1375 StatusCode::CREATED,
1376 Json(json!({
1377 "token": token,
1378 "id": view.id,
1379 "device_id": view.device_id,
1380 "label": view.label,
1381 "created_at": view.created_at,
1382 "expires_at": view.expires_at,
1383 })),
1384 )
1385 .into_response(),
1386 Err(e) => e.into_response(),
1387 }
1388 }
1389
1390 async fn list_client_tokens(State(state): State<RouteState>, headers: HeaderMap) -> Response {
1391 if let Err(e) = require_owner(&state, &headers) {
1392 return e.into_response();
1393 }
1394 Json(json!({ "tokens": state.computer.list_client_tokens() })).into_response()
1395 }
1396
1397 async fn revoke_client_token(
1398 State(state): State<RouteState>,
1399 headers: HeaderMap,
1400 Path(id): Path<String>,
1401 ) -> Response {
1402 if let Err(e) = require_owner(&state, &headers) {
1403 return e.into_response();
1404 }
1405 if state.computer.revoke_client_token(&id) {
1406 state.computer.sweep_lease();
1407 StatusCode::NO_CONTENT.into_response()
1408 } else {
1409 ApiErr::new(StatusCode::NOT_FOUND, "no such client token").into_response()
1410 }
1411 }
1412
1413 #[cfg(test)]
1414 #[path = "computer_display_tests.rs"]
1415 mod tests;
1416
1416 lines RUST