1use std::collections::{BTreeMap, HashMap, VecDeque};
67#[cfg(feature = "adapter-api")]
68use std::net::SocketAddr;
69use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
70use std::sync::{Arc, Mutex as StdMutex};
71
72use async_trait::async_trait;
73#[cfg(feature = "adapter-api")]
74use futures::{SinkExt, StreamExt};
75use serde_json::{json, Value};
76use tokio::io::AsyncBufRead;
77#[cfg(feature = "adapter-api")]
78use tokio::io::{AsyncRead, AsyncReadExt};
79#[cfg(feature = "adapter-api")]
80use tokio::io::{AsyncWrite, AsyncWriteExt};
81#[cfg(feature = "adapter-api")]
82use tokio::net::TcpListener;
83#[cfg(feature = "adapter-api")]
84use tokio::sync::mpsc;
85use tokio::sync::{broadcast, Mutex, Notify, RwLock};
86
87const EVENT_STREAM_KEEPALIVE: std::time::Duration = std::time::Duration::from_secs(15);
90
91use crate::agent::SteerInbox;
92use crate::frontend::{
93 FrontendActions, FrontendApprovalDecision, FrontendAttachSnapshot, FrontendAttachment,
94 FrontendCommandDescriptor, FrontendConnectionState, FrontendDisplayCapabilities, FrontendEvent,
95 FrontendOperationDescriptor, FrontendOperationInvocation, FrontendOperationKind,
96 FrontendOperationResult, FrontendProjectionState, FrontendRequest, FrontendRequestKind,
97 FrontendResponse, FrontendRuntime, FrontendRuntimeDescriptor, FrontendRuntimeError,
98 FrontendRuntimeMetadata, FrontendTurnState, FRONTEND_EVENT_SCHEMA_VERSION,
99 FRONTEND_REPLAY_CAPACITY, FRONTEND_RUNTIME_SCHEMA_VERSION,
100};
101use crate::mcp::{
102 ElicitationAction, ElicitationRequest, ElicitationResponse, McpElicitationHandler,
103};
104use crate::permissions::{ApprovalOutcome, ApprovalRequest, PermissionsApprovalHandler};
105pub use crate::sdk::RuntimeSubmitError;
106use crate::sdk::SdkAgent;
107#[cfg(feature = "adapter-api")]
108use crate::{CoordinatedRuntime, CoordinatedRuntimeClient, RuntimeAuthorization, RuntimeClientId};
109use supercode_interchange::ChatMessage;
110
111pub const SERVER_EVENT_CHANNEL_CAPACITY: usize = 1024;
118
119pub(crate) const SERVER_HISTORY_CAPACITY: usize = 200;
123
124pub const SERVER_MAX_LINE_BYTES: usize = 16 * 1024 * 1024;
131
132#[cfg(feature = "adapter-api")]
137const MAX_HEADER_LINES: usize = 200;
138
139#[derive(Debug, Clone, serde::Deserialize)]
144pub struct RpcRequest {
145 pub id: Value,
149 pub method: String,
151 #[serde(default)]
153 pub params: Value,
154}
155
156fn rpc_ok(id: Value, result: Value) -> Value {
158 json!({"id": id, "result": result})
159}
160
161fn rpc_error(id: Value, code: i32, message: impl Into<String>) -> Value {
168 json!({"id": id, "error": {"code": code, "message": message.into()}})
169}
170
171fn sdk_runtime_rpc_error(id: Value, code: i32, error: &FrontendRuntimeError) -> Value {
172 let code = match error.code() {
173 crate::SdkErrorCode::Unauthenticated => -32030,
174 crate::SdkErrorCode::Unauthorized => -32031,
175 crate::SdkErrorCode::ControllerRequired => -32032,
176 crate::SdkErrorCode::LeaseExpired => -32033,
177 _ => code,
178 };
179 let mut envelope = json!({
180 "id": id,
181 "error": {
182 "code": code,
183 "name": error.code(),
184 "operation": error.operation(),
185 "message": error.to_string(),
186 }
187 });
188 if let Some(detail) = envelope.get_mut("error").and_then(Value::as_object_mut) {
189 match error {
190 FrontendRuntimeError::Unauthorized { permission } => {
191 detail.insert("permission".into(), Value::String(permission.clone()));
192 }
193 FrontendRuntimeError::ControllerRequired {
194 holder,
195 expires_at_ms,
196 } => {
197 if let Some(holder) = holder {
198 detail.insert("holder".into(), Value::String(holder.clone()));
199 }
200 if let Some(expires_at_ms) = expires_at_ms {
201 detail.insert("expiresAtMs".into(), json!(expires_at_ms));
202 }
203 }
204 _ => {}
205 }
206 }
207 envelope
208}
209
210async fn read_bounded_line<R>(reader: &mut R, cap: usize) -> std::io::Result<Option<String>>
218where
219 R: AsyncBufRead + Unpin,
220{
221 use tokio::io::AsyncBufReadExt;
222 let mut out: Vec<u8> = Vec::new();
223 loop {
224 let buf = reader.fill_buf().await?;
225 if buf.is_empty() {
226 return Ok(if out.is_empty() {
227 None
228 } else {
229 Some(strip_crlf(out))
230 });
231 }
232 if let Some(pos) = buf.iter().position(|&b| b == b'\n') {
233 if out.len() + pos > cap {
234 reader.consume(pos + 1);
235 return Err(std::io::Error::new(
236 std::io::ErrorKind::InvalidData,
237 format!("line exceeded {cap} byte cap"),
238 ));
239 }
240 out.extend_from_slice(&buf[..pos]);
241 reader.consume(pos + 1);
242 return Ok(Some(strip_crlf(out)));
243 }
244 let take = buf.len();
245 if out.len() + take > cap {
246 reader.consume(take);
247 loop {
250 let b = reader.fill_buf().await?;
251 if b.is_empty() {
252 break;
253 }
254 if let Some(p) = b.iter().position(|&x| x == b'\n') {
255 reader.consume(p + 1);
256 break;
257 }
258 let n = b.len();
259 reader.consume(n);
260 }
261 return Err(std::io::Error::new(
262 std::io::ErrorKind::InvalidData,
263 format!("line exceeded {cap} byte cap"),
264 ));
265 }
266 out.extend_from_slice(buf);
267 reader.consume(take);
268 }
269}
270
271fn strip_crlf(mut v: Vec<u8>) -> String {
272 if v.last() == Some(&b'\r') {
273 v.pop();
274 }
275 String::from_utf8_lossy(&v).into_owned()
276}
277
278#[cfg(feature = "adapter-api")]
282fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
283 if a.len() != b.len() {
284 return false;
285 }
286 let mut diff = 0u8;
287 for (x, y) in a.iter().zip(b.iter()) {
288 diff |= x ^ y;
289 }
290 diff == 0
291}
292
293pub fn generate_token() -> String {
300 let mut bytes = [0u8; 32];
301 getrandom::getrandom(&mut bytes).expect("OS entropy source for the server bearer token");
307 bytes.iter().map(|b| format!("{b:02x}")).collect()
308}
309
310type TurnCompleteHook = Box<dyn Fn(&SdkAgent) + Send + Sync>;
314
315#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
320pub struct RuntimeStatus {
321 pub session_id: String,
323 pub model: String,
325 pub busy: bool,
327 pub shutting_down: bool,
329}
330
331type PendingFrontendResponses = StdMutex<
332 HashMap<
333 u64,
334 (
335 FrontendRequestKind,
336 std::sync::mpsc::Sender<AcceptedFrontendResponse>,
337 ),
338 >,
339>;
340
341struct AcceptedFrontendResponse {
342 response: FrontendResponse,
343 published: std::sync::mpsc::Receiver<()>,
346}
347
348struct FrontendRequestBroker {
351 next_id: std::sync::atomic::AtomicU64,
352 pending: PendingFrontendResponses,
353 transport: StdMutex<Option<FrontendRequestTransport>>,
354}
355
356#[derive(Clone)]
357struct FrontendRequestTransport {
358 events: broadcast::Sender<FrontendEvent>,
359 state: Arc<StdMutex<FrontendProjectionState>>,
360}
361
362impl FrontendRequestBroker {
363 fn new() -> Arc<Self> {
364 Arc::new(Self {
365 next_id: std::sync::atomic::AtomicU64::new(1),
366 pending: StdMutex::new(HashMap::new()),
367 transport: StdMutex::new(None),
368 })
369 }
370
371 fn bind(
372 &self,
373 events: broadcast::Sender<FrontendEvent>,
374 state: Arc<StdMutex<FrontendProjectionState>>,
375 ) {
376 *self
377 .transport
378 .lock()
379 .unwrap_or_else(std::sync::PoisonError::into_inner) =
380 Some(FrontendRequestTransport { events, state });
381 }
382
383 fn transport(&self) -> Option<FrontendRequestTransport> {
384 self.transport
385 .lock()
386 .unwrap_or_else(std::sync::PoisonError::into_inner)
387 .clone()
388 }
389
390 fn publish(&self, request: &FrontendRequest) -> bool {
391 self.publish_payload(json!({"type": "request", "request": request}))
392 }
393
394 fn publish_payload(&self, payload: Value) -> bool {
395 let Some(transport) = self.transport() else {
396 return false;
397 };
398 let event = {
399 let mut state = transport
400 .state
401 .lock()
402 .unwrap_or_else(std::sync::PoisonError::into_inner);
403 let event = FrontendEvent::new(state.next_sequence, payload);
404 state.next_sequence = state.next_sequence.saturating_add(1);
405 state.replay.push_back(event.clone());
406 while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
407 state.replay.pop_front();
408 }
409 event
410 };
411 transport.events.send(event).is_ok()
412 }
413
414 fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
415 let request_id = response.request_id();
416 let response_kind = match &response {
417 FrontendResponse::Approval { .. } => FrontendRequestKind::Approval,
418 FrontendResponse::Elicitation { .. } => FrontendRequestKind::Elicitation,
419 FrontendResponse::Other { .. } => FrontendRequestKind::Other,
420 };
421 let mut pending = self
422 .pending
423 .lock()
424 .unwrap_or_else(std::sync::PoisonError::into_inner);
425 let expected = pending
426 .get(&request_id)
427 .map(|(kind, _)| *kind)
428 .ok_or(FrontendRuntimeError::UnknownRequest(request_id))?;
429 if expected != response_kind {
430 return Err(FrontendRuntimeError::InvalidResponse(format!(
431 "request {request_id} expects {expected:?}, got {response_kind:?}"
432 )));
433 }
434 let (_, sender) = pending
435 .remove(&request_id)
436 .ok_or(FrontendRuntimeError::UnknownRequest(request_id))?;
437 drop(pending);
438 let payload = json!({
439 "type": "request_resolved",
440 "request_id": request_id,
441 "response": &response,
442 });
443 let (published_tx, published_rx) = std::sync::mpsc::channel();
444 sender
445 .send(AcceptedFrontendResponse {
446 response,
447 published: published_rx,
448 })
449 .map_err(|_| FrontendRuntimeError::UnknownRequest(request_id))?;
450 self.publish_payload(payload);
451 let _ = published_tx.send(());
452 Ok(())
453 }
454
455 fn ask_approval(
456 &self,
457 req: &ApprovalRequest<'_>,
458 child: Option<(&str, &Arc<StdMutex<Vec<crate::subagents::QueuedApproval>>>)>,
459 ) -> ApprovalOutcome {
460 let queued = child.and_then(|(child_agent_id, queue)| {
464 crate::subagents::queue_approval(
465 queue,
466 crate::subagents::QueuedApproval {
467 child_agent_id: child_agent_id.to_string(),
468 tool: req.tool.to_string(),
469 subject: req.subject.map(String::from),
470 queued_at_ms: std::time::SystemTime::now()
471 .duration_since(std::time::UNIX_EPOCH)
472 .map(|duration| duration.as_millis() as i64)
473 .unwrap_or_default(),
474 outcome: None,
475 },
476 )
477 .map(|index| (queue.clone(), index))
478 });
479 let outcome = self.decide_approval(req, child.map(|(id, _)| id));
480 if let Some((queue, index)) = queued {
481 crate::subagents::record_queued_outcome(&queue, index, outcome.into());
482 }
483 outcome
484 }
485
486 fn decide_approval(
487 &self,
488 req: &ApprovalRequest<'_>,
489 child_agent_id: Option<&str>,
490 ) -> ApprovalOutcome {
491 let Some(transport) = self.transport() else {
494 return ApprovalOutcome::Deny;
495 };
496 if transport.events.receiver_count() == 0 {
497 return ApprovalOutcome::Deny;
498 }
499 let id = self.next_id.fetch_add(1, Ordering::SeqCst);
500 let mut payload = json!({
501 "tool": req.tool,
502 "subject": req.subject,
503 "raw_args": req.raw_args,
504 });
505 if let Some(child_agent_id) = child_agent_id {
506 payload["child_agent_id"] = Value::String(child_agent_id.to_string());
507 }
508 let request = FrontendRequest {
509 id,
510 kind: FrontendRequestKind::Approval,
511 payload,
512 };
513 let (tx, rx) = std::sync::mpsc::channel();
514 self.pending
515 .lock()
516 .unwrap_or_else(std::sync::PoisonError::into_inner)
517 .insert(id, (FrontendRequestKind::Approval, tx));
518 if !self.publish(&request) {
519 self.pending
520 .lock()
521 .unwrap_or_else(std::sync::PoisonError::into_inner)
522 .remove(&id);
523 return ApprovalOutcome::Deny;
524 }
525 let wait_for_response = || loop {
526 match rx.recv_timeout(std::time::Duration::from_millis(100)) {
527 Ok(accepted) => {
528 let _ = accepted.published.recv();
529 let FrontendResponse::Approval { decision, .. } = accepted.response else {
530 return ApprovalOutcome::Deny;
531 };
532 return match decision {
533 FrontendApprovalDecision::Deny => ApprovalOutcome::Deny,
534 FrontendApprovalDecision::Allow => ApprovalOutcome::Allow,
535 FrontendApprovalDecision::AllowForSession => {
536 ApprovalOutcome::AllowForSession
537 }
538 };
539 }
540 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
541 return ApprovalOutcome::Deny;
542 }
543 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
544 if self
545 .transport()
546 .map(|transport| transport.events.receiver_count() == 0)
547 .unwrap_or(true)
548 {
549 self.pending
550 .lock()
551 .unwrap_or_else(std::sync::PoisonError::into_inner)
552 .remove(&id);
553 return ApprovalOutcome::Deny;
554 }
555 }
556 }
557 };
558 if tokio::runtime::Handle::try_current()
559 .map(|handle| handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
560 .unwrap_or(false)
561 {
562 tokio::task::block_in_place(wait_for_response)
563 } else {
564 wait_for_response()
565 }
566 }
567
568 async fn ask_elicitation(self: Arc<Self>, req: &ElicitationRequest) -> ElicitationResponse {
569 let cancel = || ElicitationResponse {
570 action: ElicitationAction::Cancel,
571 content: None,
572 };
573 let Some(transport) = self.transport() else {
574 return cancel();
575 };
576 if transport.events.receiver_count() == 0 {
577 return cancel();
578 }
579 let id = self.next_id.fetch_add(1, Ordering::SeqCst);
580 let request = FrontendRequest {
581 id,
582 kind: FrontendRequestKind::Elicitation,
583 payload: json!({
584 "message": req.message,
585 "requested_schema": req.requested_schema,
586 }),
587 };
588 let (tx, rx) = std::sync::mpsc::channel();
589 self.pending
590 .lock()
591 .unwrap_or_else(std::sync::PoisonError::into_inner)
592 .insert(id, (FrontendRequestKind::Elicitation, tx));
593 if !self.publish(&request) {
594 self.pending
595 .lock()
596 .unwrap_or_else(std::sync::PoisonError::into_inner)
597 .remove(&id);
598 return cancel();
599 }
600 let broker = self.clone();
601 tokio::task::spawn_blocking(move || loop {
602 match rx.recv_timeout(std::time::Duration::from_millis(100)) {
603 Ok(accepted) => {
604 let _ = accepted.published.recv();
605 let FrontendResponse::Elicitation {
606 action, content, ..
607 } = accepted.response
608 else {
609 return cancel();
610 };
611 return ElicitationResponse {
612 action: match action {
613 crate::frontend::FrontendElicitationAction::Accept => {
614 ElicitationAction::Accept
615 }
616 crate::frontend::FrontendElicitationAction::Decline => {
617 ElicitationAction::Decline
618 }
619 crate::frontend::FrontendElicitationAction::Cancel => {
620 ElicitationAction::Cancel
621 }
622 },
623 content,
624 };
625 }
626 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return cancel(),
627 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
628 if broker
629 .transport()
630 .map(|transport| transport.events.receiver_count() == 0)
631 .unwrap_or(true)
632 {
633 broker
634 .pending
635 .lock()
636 .unwrap_or_else(std::sync::PoisonError::into_inner)
637 .remove(&id);
638 return cancel();
639 }
640 }
641 }
642 })
643 .await
644 .unwrap_or_else(|_| cancel())
645 }
646}
647
648struct FrontendApprovalHandler(Arc<FrontendRequestBroker>);
649
650impl PermissionsApprovalHandler for FrontendApprovalHandler {
651 fn ask(&self, req: &ApprovalRequest<'_>) -> ApprovalOutcome {
652 self.0.ask_approval(req, None)
653 }
654}
655
656struct FrontendChildApprovalHandler {
657 broker: Arc<FrontendRequestBroker>,
658 child_agent_id: String,
659 queue: Arc<StdMutex<Vec<crate::subagents::QueuedApproval>>>,
660}
661
662impl PermissionsApprovalHandler for FrontendChildApprovalHandler {
663 fn ask(&self, req: &ApprovalRequest<'_>) -> ApprovalOutcome {
664 self.broker
665 .ask_approval(req, Some((&self.child_agent_id, &self.queue)))
666 }
667}
668
669#[derive(Clone)]
672pub struct FrontendRequestBridge {
673 broker: Arc<FrontendRequestBroker>,
674}
675
676impl FrontendRequestBridge {
677 pub fn new() -> Self {
680 Self {
681 broker: FrontendRequestBroker::new(),
682 }
683 }
684
685 pub fn elicitation_handler(&self) -> Arc<dyn McpElicitationHandler> {
687 Arc::new(FrontendElicitationHandler(self.broker.clone()))
688 }
689}
690
691impl Default for FrontendRequestBridge {
692 fn default() -> Self {
693 Self::new()
694 }
695}
696
697struct FrontendElicitationHandler(Arc<FrontendRequestBroker>);
698
699#[async_trait]
700impl McpElicitationHandler for FrontendElicitationHandler {
701 async fn handle(&self, request: &ElicitationRequest) -> ElicitationResponse {
702 self.0.clone().ask_elicitation(request).await
703 }
704}
705
706pub struct RpcEngine {
712 agent: Mutex<SdkAgent>,
713 history_snapshot: RwLock<Vec<ChatMessage>>,
717 session_id: String,
718 model: String,
722 busy: Arc<AtomicBool>,
723 current_cancel: Arc<StdMutex<Option<Arc<Notify>>>>,
729 turn_finished: Arc<Notify>,
732 steer_queue: Arc<StdMutex<SteerInbox>>,
735 events: broadcast::Sender<Value>,
736 frontend_events: broadcast::Sender<FrontendEvent>,
738 frontend_state: Arc<StdMutex<FrontendProjectionState>>,
740 frontend_metadata: FrontendRuntimeMetadata,
741 frontend_active_modules: Vec<String>,
742 frontend_commands: Vec<FrontendCommandDescriptor>,
743 frontend_operations: Vec<FrontendOperationDescriptor>,
744 frontend_requests: Option<Arc<FrontendRequestBroker>>,
745 shutdown: Notify,
746 shutting_down: AtomicBool,
747 accepting_submits: AtomicBool,
751 shutdown_barrier: Mutex<()>,
754 on_turn_complete: Option<TurnCompleteHook>,
760}
761
762struct SdkSubmitClaim {
769 inbox: Arc<StdMutex<SteerInbox>>,
770 busy: Arc<AtomicBool>,
771 cancel: Arc<Notify>,
772 current_cancel: Arc<StdMutex<Option<Arc<Notify>>>>,
773 turn_finished: Arc<Notify>,
774 frontend_events: broadcast::Sender<FrontendEvent>,
775 frontend_state: Arc<StdMutex<FrontendProjectionState>>,
776 lifecycle_started: bool,
777}
778
779impl SdkSubmitClaim {
780 fn mark_lifecycle_started(&mut self) {
781 self.lifecycle_started = true;
782 }
783
784 fn mark_lifecycle_finished(&mut self) {
785 self.lifecycle_started = false;
786 }
787}
788
789impl Drop for SdkSubmitClaim {
790 fn drop(&mut self) {
791 if self.lifecycle_started {
799 let event = {
800 let mut state = self
801 .frontend_state
802 .lock()
803 .unwrap_or_else(std::sync::PoisonError::into_inner);
804 let event = FrontendEvent::new(
805 state.next_sequence,
806 json!({
807 "type": "turn_interrupted",
808 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
809 }),
810 );
811 state.next_sequence = state.next_sequence.saturating_add(1);
812 state.replay.push_back(event.clone());
813 while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
814 state.replay.pop_front();
815 }
816 event
817 };
818 let _ = self.frontend_events.send(event);
819 }
820 self.inbox
821 .lock()
822 .unwrap_or_else(std::sync::PoisonError::into_inner)
823 .close();
824 *self
825 .current_cancel
826 .lock()
827 .unwrap_or_else(std::sync::PoisonError::into_inner) = None;
828 self.busy.store(false, Ordering::SeqCst);
829 self.turn_finished.notify_waiters();
830 self.turn_finished.notify_one();
831 }
832}
833
834impl RpcEngine {
835 pub fn new(
843 agent: impl Into<SdkAgent>,
844 on_turn_complete: Option<TurnCompleteHook>,
845 ) -> Arc<Self> {
846 let agent = agent.into();
847 let session_id = agent
848 .session_name()
849 .map(str::to_owned)
850 .unwrap_or_else(|| format!("supercode-{}", std::process::id()));
851 Self::new_named(agent, session_id, on_turn_complete)
852 }
853
854 pub fn new_named(
858 agent: impl Into<SdkAgent>,
859 session_id: impl Into<String>,
860 on_turn_complete: Option<TurnCompleteHook>,
861 ) -> Arc<Self> {
862 Self::new_named_with_frontend_metadata(
863 agent.into(),
864 session_id,
865 FrontendRuntimeMetadata::default(),
866 on_turn_complete,
867 )
868 }
869
870 pub fn new_named_with_frontend_metadata(
873 agent: impl Into<SdkAgent>,
874 session_id: impl Into<String>,
875 frontend_metadata: FrontendRuntimeMetadata,
876 on_turn_complete: Option<TurnCompleteHook>,
877 ) -> Arc<Self> {
878 Self::build(
879 agent.into(),
880 session_id.into(),
881 frontend_metadata,
882 None,
883 on_turn_complete,
884 )
885 }
886
887 pub fn new_named_with_frontend_requests(
891 agent: impl Into<SdkAgent>,
892 session_id: impl Into<String>,
893 frontend_metadata: FrontendRuntimeMetadata,
894 on_turn_complete: Option<TurnCompleteHook>,
895 ) -> Arc<Self> {
896 let bridge = FrontendRequestBridge::new();
897 Self::new_named_with_frontend_bridge(
898 agent.into(),
899 session_id,
900 frontend_metadata,
901 bridge,
902 on_turn_complete,
903 )
904 }
905
906 pub fn new_named_with_frontend_bridge(
909 agent: impl Into<SdkAgent>,
910 session_id: impl Into<String>,
911 frontend_metadata: FrontendRuntimeMetadata,
912 bridge: FrontendRequestBridge,
913 on_turn_complete: Option<TurnCompleteHook>,
914 ) -> Arc<Self> {
915 Self::build(
916 agent.into(),
917 session_id.into(),
918 frontend_metadata,
919 Some(bridge.broker),
920 on_turn_complete,
921 )
922 }
923
924 fn build(
925 mut agent: SdkAgent,
926 session_id: String,
927 frontend_metadata: FrontendRuntimeMetadata,
928 frontend_requests: Option<Arc<FrontendRequestBroker>>,
929 on_turn_complete: Option<TurnCompleteHook>,
930 ) -> Arc<Self> {
931 let (tx, _rx) = broadcast::channel(SERVER_EVENT_CHANNEL_CAPACITY);
932 let events_tx = tx.clone();
933 let (frontend_tx, _frontend_rx) = broadcast::channel(SERVER_EVENT_CHANNEL_CAPACITY);
934 let frontend_events_tx = frontend_tx.clone();
935 let model = agent.config().model.clone();
936 let steer_queue = agent.inner().steer_queue_handle();
937 let history_snapshot = bounded_history_snapshot(agent.history());
938 let frontend_state = Arc::new(StdMutex::new(FrontendProjectionState {
939 history: history_snapshot.clone(),
940 history_cursor: 0,
941 next_sequence: 1,
942 replay: VecDeque::new(),
943 }));
944 if let Some(broker) = &frontend_requests {
945 broker.bind(frontend_tx.clone(), frontend_state.clone());
946 let legacy_broker = broker.clone();
947 agent
948 .inner_mut()
949 .set_legacy_approval_handler(Box::new(move |call| {
950 let Ok(raw_args) = call.function.parsed_arguments() else {
951 return false;
952 };
953 let subject = raw_args
954 .get("command")
955 .or_else(|| raw_args.get("path"))
956 .or_else(|| raw_args.get("file_path"))
957 .or_else(|| raw_args.get("patch"))
958 .and_then(Value::as_str);
959 matches!(
960 legacy_broker.ask_approval(
961 &ApprovalRequest {
962 tool: &call.function.name,
963 subject,
964 raw_args: &raw_args,
965 },
966 None,
967 ),
968 ApprovalOutcome::Allow | ApprovalOutcome::AllowForSession
969 )
970 }));
971 agent
972 .inner_mut()
973 .set_permissions_approval_handler(FrontendApprovalHandler(broker.clone()));
974 agent
981 .inner_mut()
982 .set_user_question_handler(Arc::new(FrontendElicitationHandler(broker.clone())));
983 let broker = broker.clone();
984 agent
985 .inner_mut()
986 .set_child_approval_handler_factory(move |child_agent_id, queue| {
987 Arc::new(FrontendChildApprovalHandler {
988 broker: broker.clone(),
989 child_agent_id,
990 queue,
991 }) as Arc<dyn PermissionsApprovalHandler>
992 });
993 }
994 let event_frontend_state = frontend_state.clone();
995 let frontend_active_modules = agent
996 .config()
997 .module_activation
998 .iter()
999 .map(ToString::to_string)
1000 .collect();
1001 let mut frontend_operations = agent
1002 .config()
1003 .prompts
1004 .keys()
1005 .filter(|name| valid_frontend_command_name(name))
1006 .map(|name| FrontendOperationDescriptor {
1007 id: format!("prompt:{name}"),
1008 kind: FrontendOperationKind::Prompt,
1009 command: Some(FrontendCommandDescriptor {
1010 name: name.clone(),
1011 description: None,
1012 argument_hint: Some("[arguments]".into()),
1013 }),
1014 })
1015 .collect::<Vec<_>>();
1016 frontend_operations.sort_by(|left, right| left.id.cmp(&right.id));
1017 if agent.config().model_switch_allow_switch {
1030 frontend_operations.push(FrontendOperationDescriptor {
1031 id: "model:switch".into(),
1032 kind: FrontendOperationKind::Model,
1033 command: Some(FrontendCommandDescriptor {
1034 name: "model".into(),
1035 description: Some("show or switch this session's model".into()),
1036 argument_hint: Some("[model]".into()),
1037 }),
1038 });
1039 }
1040 frontend_operations.push(FrontendOperationDescriptor {
1041 id: "context:usage".into(),
1042 kind: FrontendOperationKind::Context,
1043 command: Some(FrontendCommandDescriptor {
1044 name: "context".into(),
1045 description: Some("context-window usage for this session".into()),
1046 argument_hint: None,
1047 }),
1048 });
1049 let frontend_commands = frontend_operations
1052 .iter()
1053 .filter_map(|operation| operation.command.as_ref())
1054 .map(|command| FrontendCommandDescriptor {
1055 name: command.name.clone(),
1056 description: command.description.clone(),
1057 argument_hint: None,
1058 })
1059 .collect();
1060 agent.inner_mut().set_event_sink(Box::new(move |event| {
1061 let payload = event.to_json();
1065 let _ = events_tx.send(payload.clone());
1066 let sequenced = {
1067 let mut state = event_frontend_state
1068 .lock()
1069 .unwrap_or_else(std::sync::PoisonError::into_inner);
1070 let event = FrontendEvent::new(state.next_sequence, payload);
1071 state.next_sequence = state.next_sequence.saturating_add(1);
1072 state.replay.push_back(event.clone());
1073 while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1074 state.replay.pop_front();
1075 }
1076 event
1077 };
1078 let _ = frontend_events_tx.send(sequenced);
1079 }));
1080 Arc::new(RpcEngine {
1081 agent: Mutex::new(agent),
1082 history_snapshot: RwLock::new(history_snapshot),
1083 session_id,
1084 model,
1085 busy: Arc::new(AtomicBool::new(false)),
1086 current_cancel: Arc::new(StdMutex::new(None)),
1087 turn_finished: Arc::new(Notify::new()),
1088 steer_queue,
1089 events: tx,
1090 frontend_events: frontend_tx,
1091 frontend_state,
1092 frontend_metadata,
1093 frontend_active_modules,
1094 frontend_commands,
1095 frontend_operations,
1096 frontend_requests,
1097 shutdown: Notify::new(),
1098 shutting_down: AtomicBool::new(false),
1099 accepting_submits: AtomicBool::new(true),
1100 shutdown_barrier: Mutex::new(()),
1101 on_turn_complete,
1102 })
1103 }
1104
1105 pub fn subscribe(&self) -> broadcast::Receiver<Value> {
1109 self.events.subscribe()
1110 }
1111
1112 pub fn frontend_descriptor(&self) -> FrontendRuntimeDescriptor {
1114 FrontendRuntimeDescriptor {
1115 schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
1116 session_id: self.session_id.clone(),
1117 source_harness: self.frontend_metadata.source_harness.clone(),
1118 emulation_profile: self.frontend_metadata.emulation_profile.clone(),
1119 active_modules: self.frontend_active_modules.clone(),
1120 commands: self.frontend_commands.clone(),
1121 operations: self.frontend_operations.clone(),
1122 actions: FrontendActions {
1123 submit: true,
1124 interrupt: true,
1125 steer: true,
1126 respond: self.frontend_requests.is_some(),
1127 detach: true,
1128 close: true,
1132 },
1133 display: FrontendDisplayCapabilities {
1134 event_kinds: vec![
1135 "user_message".into(),
1136 "turn_started".into(),
1137 "turn_succeeded".into(),
1138 "turn_interrupted".into(),
1139 "turn_failed".into(),
1140 "text_delta".into(),
1141 "turn_completed".into(),
1142 "tool_call_started".into(),
1143 "tool_call_completed".into(),
1144 "cache_warning".into(),
1145 "usage".into(),
1146 "background_output".into(),
1147 "request".into(),
1148 "request_resolved".into(),
1149 "scheduled_prompt_started".into(),
1150 "scheduled_prompt_deferred".into(),
1151 "scheduled_prompt_completed".into(),
1152 "scheduler_error".into(),
1153 ],
1154 opaque_fallback: true,
1155 },
1156 model: self.model.clone(),
1157 turn_state: if self.busy.load(Ordering::SeqCst) {
1158 FrontendTurnState::Busy
1159 } else {
1160 FrontendTurnState::Idle
1161 },
1162 connection_state: if self.is_shutting_down() {
1163 FrontendConnectionState::ShuttingDown
1164 } else {
1165 FrontendConnectionState::Connected
1166 },
1167 extensions: Default::default(),
1168 }
1169 }
1170
1171 pub fn frontend_attach(
1176 &self,
1177 history_limit: usize,
1178 ) -> Result<FrontendAttachment, FrontendRuntimeError> {
1179 let live = self.frontend_subscribe();
1180 let snapshot = self.frontend_snapshot(history_limit)?;
1181 Ok(FrontendAttachment::new(
1182 snapshot.descriptor,
1183 snapshot.history,
1184 snapshot.history_cursor,
1185 snapshot.replay,
1186 live,
1187 None,
1188 ))
1189 }
1190
1191 pub fn frontend_subscribe(&self) -> broadcast::Receiver<FrontendEvent> {
1195 self.frontend_events.subscribe()
1196 }
1197
1198 pub fn frontend_snapshot(
1200 &self,
1201 history_limit: usize,
1202 ) -> Result<FrontendAttachSnapshot, FrontendRuntimeError> {
1203 let state = self
1204 .frontend_state
1205 .lock()
1206 .unwrap_or_else(std::sync::PoisonError::into_inner);
1207 let limit = history_limit.min(SERVER_HISTORY_CAPACITY);
1208 let start = state.history.len().saturating_sub(limit);
1209 let replay = state
1210 .replay
1211 .iter()
1212 .filter(|event| event.sequence > state.history_cursor)
1213 .cloned()
1214 .collect::<VecDeque<_>>();
1215 if let Some(first) = replay.front() {
1216 let expected = state.history_cursor.saturating_add(1);
1217 if first.sequence > expected {
1218 return Err(FrontendRuntimeError::ReplayGap(first.sequence - expected));
1219 }
1220 }
1221 Ok(FrontendAttachSnapshot {
1222 descriptor: self.frontend_descriptor(),
1223 history: state.history[start..].to_vec(),
1224 history_cursor: state.history_cursor,
1225 replay,
1226 })
1227 }
1228
1229 fn publish_frontend_payload(&self, payload: Value) {
1230 let event = {
1231 let mut state = self
1232 .frontend_state
1233 .lock()
1234 .unwrap_or_else(std::sync::PoisonError::into_inner);
1235 let event = FrontendEvent::new(state.next_sequence, payload);
1236 state.next_sequence = state.next_sequence.saturating_add(1);
1237 state.replay.push_back(event.clone());
1238 while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1239 state.replay.pop_front();
1240 }
1241 event
1242 };
1243 let _ = self.frontend_events.send(event);
1244 }
1245
1246 pub fn session_id(&self) -> &str {
1248 &self.session_id
1249 }
1250
1251 fn claim_submit(&self) -> Result<SdkSubmitClaim, RuntimeSubmitError> {
1252 let mut current_cancel = self
1257 .current_cancel
1258 .lock()
1259 .unwrap_or_else(std::sync::PoisonError::into_inner);
1260 if !self.accepting_submits.load(Ordering::SeqCst) {
1261 return Err(RuntimeSubmitError::Interrupted);
1262 }
1263 if self.busy.swap(true, Ordering::SeqCst) {
1264 return Err(RuntimeSubmitError::Busy);
1265 }
1266 if !self.accepting_submits.load(Ordering::SeqCst) {
1270 self.busy.store(false, Ordering::SeqCst);
1271 self.turn_finished.notify_waiters();
1272 return Err(RuntimeSubmitError::Interrupted);
1273 }
1274 let cancel = Arc::new(Notify::new());
1275 *current_cancel = Some(cancel.clone());
1276 drop(current_cancel);
1277 self.steer_queue
1278 .lock()
1279 .unwrap_or_else(std::sync::PoisonError::into_inner)
1280 .open();
1281 Ok(SdkSubmitClaim {
1282 inbox: self.steer_queue.clone(),
1283 busy: self.busy.clone(),
1284 cancel,
1285 current_cancel: self.current_cancel.clone(),
1286 turn_finished: self.turn_finished.clone(),
1287 frontend_events: self.frontend_events.clone(),
1288 frontend_state: self.frontend_state.clone(),
1289 lifecycle_started: false,
1290 })
1291 }
1292
1293 async fn submit_claimed(
1294 &self,
1295 prompt: String,
1296 image_urls: Vec<String>,
1297 mut submit_claim: SdkSubmitClaim,
1298 ) -> Result<String, RuntimeSubmitError> {
1299 let cancel = submit_claim.cancel.clone();
1300 self.publish_frontend_payload(json!({"type": "user_message", "text": &prompt}));
1301 self.publish_frontend_payload(json!({
1302 "type": "turn_started",
1303 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
1304 }));
1305 submit_claim.mark_lifecycle_started();
1306 let outcome = {
1307 let mut agent = self.agent.lock().await;
1308 let result = tokio::select! {
1309 biased;
1310 _ = cancel.notified() => Err(RuntimeSubmitError::Interrupted),
1311 result = async {
1312 if image_urls.is_empty() {
1313 agent.inner_mut().send(&prompt).await
1314 } else {
1315 agent.inner_mut().send_with_images(&prompt, &image_urls).await
1316 }
1317 } => result.map_err(|error| RuntimeSubmitError::Agent(error.to_string())),
1318 };
1319 if result.is_ok() {
1320 if let Some(hook) = &self.on_turn_complete {
1321 hook(&agent);
1322 }
1323 }
1324 let history = bounded_history_snapshot(agent.history());
1328 *self.history_snapshot.write().await = history.clone();
1329 let mut state = self
1330 .frontend_state
1331 .lock()
1332 .unwrap_or_else(std::sync::PoisonError::into_inner);
1333 state.history = history;
1334 state.history_cursor = state.next_sequence.saturating_sub(1);
1335 let request_history = compact_frontend_request_history(&state.replay);
1342 state.replay.clear();
1343 for payload in request_history {
1344 let event = FrontendEvent::new(state.next_sequence, payload);
1345 state.next_sequence = state.next_sequence.saturating_add(1);
1346 state.replay.push_back(event);
1347 while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1348 state.replay.pop_front();
1349 }
1350 }
1351 result
1352 };
1353 *self
1354 .current_cancel
1355 .lock()
1356 .unwrap_or_else(std::sync::PoisonError::into_inner) = None;
1357 let lifecycle = match &outcome {
1358 Ok(reply) => json!({
1359 "type": "turn_succeeded",
1360 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION,
1361 "reply": reply,
1362 }),
1363 Err(RuntimeSubmitError::Interrupted) => json!({
1364 "type": "turn_interrupted",
1365 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
1366 }),
1367 Err(error) => json!({
1368 "type": "turn_failed",
1369 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION,
1370 "message": error.to_string()
1371 }),
1372 };
1373 self.publish_frontend_payload(lifecycle);
1374 submit_claim.mark_lifecycle_finished();
1375 drop(submit_claim);
1379 outcome
1380 }
1381
1382 pub async fn submit(&self, prompt: impl Into<String>) -> Result<String, RuntimeSubmitError> {
1384 let submit_claim = self.claim_submit()?;
1385 self.submit_claimed(prompt.into(), Vec::new(), submit_claim)
1386 .await
1387 }
1388
1389 pub async fn submit_with_images(
1391 &self,
1392 prompt: impl Into<String>,
1393 image_urls: Vec<String>,
1394 ) -> Result<String, RuntimeSubmitError> {
1395 let submit_claim = self.claim_submit()?;
1396 self.submit_claimed(prompt.into(), image_urls, submit_claim)
1397 .await
1398 }
1399
1400 pub fn send_input(self: &Arc<Self>, prompt: String) -> Result<(), RuntimeSubmitError> {
1403 self.send_input_with_images(prompt, Vec::new())
1404 }
1405
1406 pub fn send_input_with_images(
1409 self: &Arc<Self>,
1410 prompt: String,
1411 image_urls: Vec<String>,
1412 ) -> Result<(), RuntimeSubmitError> {
1413 let submit_claim = self.claim_submit()?;
1414 let runtime = self.clone();
1415 tokio::spawn(async move {
1416 let _ = runtime
1417 .submit_claimed(prompt, image_urls, submit_claim)
1418 .await;
1419 });
1420 Ok(())
1421 }
1422
1423 pub fn steer(&self, prompt: impl Into<String>) -> Result<(), FrontendRuntimeError> {
1427 let accepted = self
1428 .steer_queue
1429 .lock()
1430 .unwrap_or_else(std::sync::PoisonError::into_inner)
1431 .enqueue(prompt.into());
1432 if accepted {
1433 Ok(())
1434 } else {
1435 Err(FrontendRuntimeError::UnsupportedAction("steer"))
1436 }
1437 }
1438
1439 pub fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
1441 self.frontend_requests
1442 .as_ref()
1443 .ok_or(FrontendRuntimeError::UnsupportedAction("respond"))?
1444 .respond(response)
1445 }
1446
1447 pub async fn invoke(
1450 &self,
1451 operation: FrontendOperationInvocation,
1452 ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
1453 match operation {
1454 FrontendOperationInvocation::Prompt {
1455 operation_id,
1456 arguments,
1457 } => {
1458 let prompt_name = self
1459 .frontend_operations
1460 .iter()
1461 .find(|descriptor| {
1462 descriptor.id == operation_id
1463 && descriptor.kind == FrontendOperationKind::Prompt
1464 })
1465 .and_then(|descriptor| descriptor.command.as_ref())
1466 .map(|command| command.name.as_str())
1467 .ok_or_else(|| {
1468 FrontendRuntimeError::UnsupportedOperation(operation_id.clone())
1469 })?;
1470 let prompt = if arguments.is_empty() {
1471 format!("/{prompt_name}")
1472 } else {
1473 format!("/{prompt_name} {arguments}")
1474 };
1475 let reply = self.submit(prompt).await?;
1476 Ok(FrontendOperationResult::Prompt { reply })
1477 }
1478 FrontendOperationInvocation::Context { operation_id } => {
1479 if !self.frontend_operations.iter().any(|descriptor| {
1480 descriptor.id == operation_id
1481 && descriptor.kind == FrontendOperationKind::Context
1482 }) {
1483 return Err(FrontendRuntimeError::UnsupportedOperation(operation_id));
1484 }
1485 let usage = self.agent.lock().await.context_usage();
1489 Ok(FrontendOperationResult::Context { usage })
1490 }
1491 FrontendOperationInvocation::Model {
1492 operation_id,
1493 model,
1494 } => {
1495 if !self.frontend_operations.iter().any(|descriptor| {
1496 descriptor.id == operation_id && descriptor.kind == FrontendOperationKind::Model
1497 }) {
1498 return Err(FrontendRuntimeError::UnsupportedOperation(operation_id));
1499 }
1500 let mut agent = self.agent.lock().await;
1501 let previous = agent.model().to_string();
1502 if model.trim().is_empty() {
1503 return Ok(FrontendOperationResult::Model {
1504 model: previous.clone(),
1505 previous,
1506 });
1507 }
1508 let resolved = agent
1512 .config()
1513 .model_routing
1514 .resolve_alias(model.trim())
1515 .to_string();
1516 if let Some(refusal) = agent.config().model_routing.refusal(&resolved) {
1517 return Err(FrontendRuntimeError::UnsupportedOperation(refusal));
1518 }
1519 agent.switch_model(resolved.clone());
1524 Ok(FrontendOperationResult::Model {
1525 model: resolved,
1526 previous,
1527 })
1528 }
1529 }
1530 }
1531
1532 pub async fn interrupt(&self) -> bool {
1534 let cancel = self
1535 .current_cancel
1536 .lock()
1537 .unwrap_or_else(std::sync::PoisonError::into_inner)
1538 .clone();
1539 match cancel {
1540 Some(cancel) => {
1541 cancel.notify_one();
1542 true
1543 }
1544 None => false,
1545 }
1546 }
1547
1548 pub fn status(&self) -> RuntimeStatus {
1550 RuntimeStatus {
1551 session_id: self.session_id.clone(),
1552 model: self.model.clone(),
1553 busy: self.busy.load(Ordering::SeqCst),
1554 shutting_down: self.is_shutting_down(),
1555 }
1556 }
1557
1558 pub async fn history(&self, limit: usize) -> Vec<ChatMessage> {
1562 let history = self.history_snapshot.read().await;
1563 let start = history.len().saturating_sub(limit);
1564 history[start..].to_vec()
1565 }
1566
1567 pub async fn finalize_with<R>(&self, finalize: impl FnOnce(&SdkAgent) -> R) -> R {
1572 let agent = self.agent.lock().await;
1573 finalize(&agent)
1574 }
1575
1576 pub async fn shutdown(&self) {
1579 let _barrier = self.shutdown_barrier.lock().await;
1580 self.accepting_submits.store(false, Ordering::SeqCst);
1585 self.interrupt().await;
1586 loop {
1587 let finished = self.turn_finished.notified();
1588 if !self.busy.load(Ordering::SeqCst) {
1589 break;
1590 }
1591 finished.await;
1592 }
1593 self.signal_shutdown();
1594 }
1595
1596 pub fn is_shutting_down(&self) -> bool {
1599 self.shutting_down.load(Ordering::SeqCst)
1600 }
1601
1602 pub async fn wait_for_shutdown(&self) {
1605 if self.is_shutting_down() {
1610 return;
1611 }
1612 self.shutdown.notified().await;
1613 }
1614
1615 fn signal_shutdown(&self) {
1625 self.shutting_down.store(true, Ordering::SeqCst);
1626 self.shutdown.notify_waiters();
1627 }
1628
1629 pub async fn handle_request(self: &Arc<Self>, req: RpcRequest) -> Value {
1633 match req.method.as_str() {
1634 "submit" => self.handle_submit(req).await,
1635 "frontend.send_input" => self.handle_frontend_send_input(req),
1636 "interrupt" => self.handle_interrupt(req).await,
1637 "steer" => self.handle_steer(req),
1638 "respond" => self.handle_respond(req),
1639 "status" => self.handle_status(req).await,
1640 "history" => self.handle_history(req).await,
1641 "frontend.describe" => self.handle_frontend_describe(req),
1642 "frontend.attach" => self.handle_frontend_attach(req),
1643 "frontend.invoke" => self.handle_frontend_invoke(req).await,
1644 "shutdown" => self.handle_shutdown(req).await,
1645 other => rpc_error(req.id, -32601, format!("unknown method `{other}`")),
1646 }
1647 }
1648
1649 fn handle_frontend_send_input(self: &Arc<Self>, req: RpcRequest) -> Value {
1650 let Some(prompt) = req.params.get("prompt").and_then(Value::as_str) else {
1651 return rpc_error(
1652 req.id,
1653 -32602,
1654 "frontend.send_input requires a string `params.prompt`",
1655 );
1656 };
1657 let image_urls = match parse_image_urls(&req.params, "frontend.send_input") {
1658 Ok(image_urls) => image_urls,
1659 Err(message) => return rpc_error(req.id, -32602, message),
1660 };
1661 match self.send_input_with_images(prompt.to_string(), image_urls) {
1662 Ok(()) => rpc_ok(req.id, json!({"accepted": true})),
1663 Err(RuntimeSubmitError::Busy) => {
1664 rpc_error(req.id, -32000, "a turn is already in progress")
1665 }
1666 Err(RuntimeSubmitError::Interrupted) => rpc_error(req.id, -32001, "turn interrupted"),
1667 Err(RuntimeSubmitError::Agent(error)) => rpc_error(req.id, -32002, error),
1668 }
1669 }
1670
1671 async fn handle_submit(&self, req: RpcRequest) -> Value {
1682 let Some(prompt) = req.params.get("prompt").and_then(|v| v.as_str()) else {
1683 return rpc_error(req.id, -32602, "submit requires a string `params.prompt`");
1684 };
1685 match self.submit(prompt).await {
1686 Ok(reply) => rpc_ok(req.id, json!({"reply": reply})),
1687 Err(RuntimeSubmitError::Busy) => rpc_error(
1688 req.id,
1689 -32000,
1690 "a turn is already in progress; `interrupt` it or wait for its response before submitting another",
1691 ),
1692 Err(RuntimeSubmitError::Interrupted) => {
1693 rpc_error(req.id, -32001, "turn interrupted")
1694 }
1695 Err(RuntimeSubmitError::Agent(error)) => rpc_error(req.id, -32002, error),
1696 }
1697 }
1698
1699 async fn handle_interrupt(&self, req: RpcRequest) -> Value {
1705 if self.interrupt().await {
1706 rpc_ok(req.id, json!({"interrupted": true}))
1707 } else {
1708 rpc_ok(
1709 req.id,
1710 json!({"interrupted": false, "reason": "no turn in progress"}),
1711 )
1712 }
1713 }
1714
1715 fn handle_steer(&self, req: RpcRequest) -> Value {
1718 let Some(prompt) = req.params.get("prompt").and_then(Value::as_str) else {
1719 return rpc_error(req.id, -32602, "steer requires a string `params.prompt`");
1720 };
1721 match self.steer(prompt) {
1722 Ok(()) => rpc_ok(req.id, json!({"queued": true})),
1723 Err(error) => rpc_error(req.id, -32020, error.to_string()),
1724 }
1725 }
1726
1727 fn handle_respond(&self, req: RpcRequest) -> Value {
1728 let response = match req.params.get("response").cloned() {
1729 Some(value) => match serde_json::from_value::<FrontendResponse>(value) {
1730 Ok(response) => response,
1731 Err(error) => return rpc_error(req.id, -32602, error.to_string()),
1732 },
1733 None => return rpc_error(req.id, -32602, "respond requires `params.response`"),
1734 };
1735 match self.respond(response) {
1736 Ok(()) => rpc_ok(req.id, json!({"accepted": true})),
1737 Err(FrontendRuntimeError::UnsupportedAction(_)) => {
1738 rpc_error(req.id, -32020, "frontend respond is not enabled")
1739 }
1740 Err(FrontendRuntimeError::UnknownRequest(id)) => rpc_error(
1741 req.id,
1742 -32021,
1743 format!("frontend request {id} is not pending"),
1744 ),
1745 Err(error) => rpc_error(req.id, -32022, error.to_string()),
1746 }
1747 }
1748
1749 async fn handle_status(&self, req: RpcRequest) -> Value {
1753 rpc_ok(
1754 req.id,
1755 serde_json::to_value(self.status()).unwrap_or_default(),
1756 )
1757 }
1758
1759 async fn handle_history(&self, req: RpcRequest) -> Value {
1762 let limit = req
1763 .params
1764 .get("limit")
1765 .and_then(Value::as_u64)
1766 .unwrap_or(50)
1767 .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
1768 rpc_ok(req.id, json!({"messages": self.history(limit).await}))
1769 }
1770
1771 fn handle_frontend_describe(&self, req: RpcRequest) -> Value {
1772 rpc_ok(
1773 req.id,
1774 serde_json::to_value(self.frontend_descriptor()).unwrap_or_default(),
1775 )
1776 }
1777
1778 fn handle_frontend_attach(&self, req: RpcRequest) -> Value {
1779 let limit = req
1780 .params
1781 .get("limit")
1782 .and_then(Value::as_u64)
1783 .unwrap_or(50)
1784 .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
1785 match self.frontend_snapshot(limit) {
1786 Ok(snapshot) => rpc_ok(req.id, serde_json::to_value(snapshot).unwrap_or_default()),
1787 Err(error) => rpc_error(req.id, -32010, error.to_string()),
1788 }
1789 }
1790
1791 async fn handle_frontend_invoke(&self, req: RpcRequest) -> Value {
1792 let operation = match req.params.get("operation").cloned() {
1793 Some(value) => match serde_json::from_value::<FrontendOperationInvocation>(value) {
1794 Ok(operation) => operation,
1795 Err(error) => return rpc_error(req.id, -32602, error.to_string()),
1796 },
1797 None => {
1798 return rpc_error(
1799 req.id,
1800 -32602,
1801 "frontend.invoke requires `params.operation`",
1802 )
1803 }
1804 };
1805 match self.invoke(operation).await {
1806 Ok(result) => rpc_ok(req.id, serde_json::to_value(result).unwrap_or_default()),
1807 Err(FrontendRuntimeError::UnsupportedOperation(id)) => rpc_error(
1808 req.id,
1809 -32023,
1810 FrontendRuntimeError::UnsupportedOperation(id).to_string(),
1811 ),
1812 Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
1813 rpc_error(req.id, -32000, "a turn is already in progress")
1814 }
1815 Err(error) => rpc_error(req.id, -32022, error.to_string()),
1816 }
1817 }
1818
1819 async fn handle_shutdown(&self, req: RpcRequest) -> Value {
1825 self.shutdown().await;
1826 rpc_ok(req.id, json!({"shutting_down": true}))
1827 }
1828}
1829
1830fn compact_frontend_request_history(replay: &VecDeque<FrontendEvent>) -> Vec<Value> {
1837 let mut by_id: BTreeMap<u64, (Option<Value>, Option<Value>)> = BTreeMap::new();
1838 for event in replay {
1839 let (request_id, resolved) = match event.kind.as_str() {
1840 "request" => (
1841 event.payload.pointer("/request/id").and_then(Value::as_u64),
1842 false,
1843 ),
1844 "request_resolved" => (
1845 event.payload.get("request_id").and_then(Value::as_u64),
1846 true,
1847 ),
1848 _ => continue,
1849 };
1850 let Some(request_id) = request_id else {
1851 continue;
1852 };
1853 let entry = by_id.entry(request_id).or_default();
1854 let slot = if resolved { &mut entry.1 } else { &mut entry.0 };
1855 slot.get_or_insert_with(|| event.payload.clone());
1856 }
1857 by_id
1858 .into_values()
1859 .flat_map(|(request, resolution)| request.into_iter().chain(resolution))
1860 .collect()
1861}
1862
1863fn valid_frontend_command_name(name: &str) -> bool {
1864 !name.is_empty()
1865 && !name.starts_with('/')
1866 && name
1867 .chars()
1868 .all(|character| !character.is_whitespace() && !character.is_control())
1869}
1870
1871fn bounded_history_snapshot(history: &[ChatMessage]) -> Vec<ChatMessage> {
1872 let start = history.len().saturating_sub(SERVER_HISTORY_CAPACITY);
1873 history[start..].to_vec()
1874}
1875
1876#[cfg(feature = "adapter-api")]
1901pub async fn run_stdio<R, W>(engine: Arc<RpcEngine>, reader: R, writer: W) -> std::io::Result<()>
1902where
1903 R: AsyncBufRead + Unpin + Send + 'static,
1904 W: AsyncWrite + Unpin + Send + 'static,
1905{
1906 let (out_tx, mut out_rx) = mpsc::unbounded_channel::<Value>();
1907
1908 let writer_task = tokio::spawn(async move {
1909 let mut writer = writer;
1910 while let Some(v) = out_rx.recv().await {
1911 let line = format!("{v}\n");
1912 if writer.write_all(line.as_bytes()).await.is_err() {
1913 break;
1914 }
1915 if writer.flush().await.is_err() {
1916 break;
1917 }
1918 }
1919 });
1920
1921 let mut events = engine.subscribe();
1922 let evt_tx = out_tx.clone();
1923 let evt_engine = engine.clone();
1924 let event_task = tokio::spawn(async move {
1925 loop {
1926 tokio::select! {
1927 biased;
1928 _ = evt_engine.wait_for_shutdown() => break,
1929 recv = events.recv() => {
1930 match recv {
1931 Ok(v) => {
1932 if evt_tx.send(json!({"event": v})).is_err() {
1933 break;
1934 }
1935 }
1936 Err(broadcast::error::RecvError::Lagged(_)) => continue,
1937 Err(broadcast::error::RecvError::Closed) => break,
1938 }
1939 }
1940 }
1941 }
1942 });
1943
1944 let mut reader = reader;
1945 loop {
1946 if engine.is_shutting_down() {
1947 break;
1948 }
1949 tokio::select! {
1950 biased;
1951 _ = engine.wait_for_shutdown() => break,
1952 line = read_bounded_line(&mut reader, SERVER_MAX_LINE_BYTES) => {
1953 match line {
1954 Ok(None) => {
1955 engine.signal_shutdown();
1963 break;
1964 }
1965 Ok(Some(text)) => {
1966 let text = text.trim();
1967 if text.is_empty() {
1968 continue;
1969 }
1970 match serde_json::from_str::<RpcRequest>(text) {
1971 Ok(req) => {
1972 let engine = engine.clone();
1973 let out_tx = out_tx.clone();
1974 tokio::spawn(async move {
1975 let resp = engine.handle_request(req).await;
1976 let _ = out_tx.send(resp);
1977 });
1978 }
1979 Err(e) => {
1980 let _ = out_tx.send(rpc_error(Value::Null, -32700, format!("parse error: {e}")));
1981 }
1982 }
1983 }
1984 Err(e) => {
1985 let _ = out_tx.send(rpc_error(Value::Null, -32700, format!("{e}")));
1986 }
1987 }
1988 }
1989 }
1990 }
1991 drop(out_tx);
2001 let _ = event_task.await;
2002 let _ = writer_task.await;
2003 Ok(())
2004}
2005
2006#[cfg(feature = "adapter-api")]
2009pub(crate) struct HttpRequest {
2010 pub(crate) method: String,
2011 pub(crate) path: String,
2013 pub(crate) query: String,
2014 pub(crate) headers: HashMap<String, String>,
2015 pub(crate) body: Vec<u8>,
2016}
2017
2018#[cfg(feature = "adapter-api")]
2026pub(crate) async fn read_http_request<R>(reader: &mut R) -> std::io::Result<Option<HttpRequest>>
2027where
2028 R: AsyncBufRead + AsyncRead + Unpin,
2029{
2030 const HEAD_LINE_CAP: usize = 8 * 1024;
2031 let Some(request_line) = read_bounded_line(reader, HEAD_LINE_CAP).await? else {
2032 return Ok(None);
2033 };
2034 let mut parts = request_line.split_whitespace();
2035 let method = parts.next().unwrap_or("").to_string();
2036 let target = parts.next().unwrap_or("").to_string();
2037 if method.is_empty() || target.is_empty() {
2038 return Err(std::io::Error::new(
2039 std::io::ErrorKind::InvalidData,
2040 "malformed request line",
2041 ));
2042 }
2043 let (path, query) = match target.split_once('?') {
2044 Some((p, q)) => (p.to_string(), q.to_string()),
2045 None => (target, String::new()),
2046 };
2047
2048 let mut headers = HashMap::new();
2049 let mut content_length: usize = 0;
2050 for _ in 0..MAX_HEADER_LINES {
2051 let Some(line) = read_bounded_line(reader, HEAD_LINE_CAP).await? else {
2052 return Ok(None);
2053 };
2054 if line.is_empty() {
2055 break;
2056 }
2057 if let Some((k, v)) = line.split_once(':') {
2058 let k = k.trim().to_ascii_lowercase();
2059 let v = v.trim().to_string();
2060 if k == "content-length" {
2061 content_length = v.parse().unwrap_or(0);
2062 }
2063 headers.insert(k, v);
2064 }
2065 }
2066 if content_length > SERVER_MAX_LINE_BYTES {
2067 return Err(std::io::Error::new(
2068 std::io::ErrorKind::InvalidData,
2069 format!("request body exceeded {SERVER_MAX_LINE_BYTES} byte cap"),
2070 ));
2071 }
2072 let mut body = vec![0u8; content_length];
2073 if content_length > 0 {
2074 reader.read_exact(&mut body).await?;
2075 }
2076 Ok(Some(HttpRequest {
2077 method,
2078 path,
2079 query,
2080 headers,
2081 body,
2082 }))
2083}
2084
2085#[cfg(feature = "adapter-api")]
2086pub(crate) async fn write_http_response<W: AsyncWrite + Unpin>(
2087 writer: &mut W,
2088 status: u16,
2089 reason: &str,
2090 content_type: &str,
2091 body: &[u8],
2092) -> std::io::Result<()> {
2093 let head = format!(
2094 "HTTP/1.1 {status} {reason}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2095 body.len()
2096 );
2097 writer.write_all(head.as_bytes()).await?;
2098 writer.write_all(body).await?;
2099 writer.flush().await
2100}
2101
2102#[cfg(feature = "adapter-api")]
2103fn browser_observer_asset(path: &str) -> Option<(&'static str, &'static [u8])> {
2104 match path {
2105 "/observer" | "/observer/" => Some((
2106 "text/html; charset=utf-8",
2107 include_bytes!("../embedded/frontend-browser/index.html"),
2108 )),
2109 "/observer/app.mjs" => Some((
2110 "text/javascript; charset=utf-8",
2111 include_bytes!("../embedded/frontend-browser/app.mjs"),
2112 )),
2113 "/observer/client.mjs" => Some((
2114 "text/javascript; charset=utf-8",
2115 include_bytes!("../embedded/frontend-browser/client.mjs"),
2116 )),
2117 "/observer/view.mjs" => Some((
2118 "text/javascript; charset=utf-8",
2119 include_bytes!("../embedded/frontend-browser/view.mjs"),
2120 )),
2121 "/observer/style.css" => Some((
2122 "text/css; charset=utf-8",
2123 include_bytes!("../embedded/frontend-browser/style.css"),
2124 )),
2125 "/frontend/client.mjs" => Some((
2126 "text/javascript; charset=utf-8",
2127 include_bytes!("../embedded/frontend/client.mjs"),
2128 )),
2129 "/frontend/generated-client.mjs" => Some((
2130 "text/javascript; charset=utf-8",
2131 include_bytes!("../embedded/frontend/generated-client.mjs"),
2132 )),
2133 "/frontend/generated.mjs" => Some((
2134 "text/javascript; charset=utf-8",
2135 include_bytes!("../embedded/frontend/generated.mjs"),
2136 )),
2137 _ => None,
2138 }
2139}
2140
2141#[cfg(feature = "adapter-api")]
2142async fn write_browser_observer_asset<W: AsyncWrite + Unpin>(
2143 writer: &mut W,
2144 content_type: &str,
2145 body: &[u8],
2146) -> std::io::Result<()> {
2147 let head = format!(
2148 "HTTP/1.1 200 OK\r\n\
2149 Content-Type: {content_type}\r\n\
2150 Content-Length: {}\r\n\
2151 Cache-Control: no-store\r\n\
2152 Content-Security-Policy: default-src 'none'; script-src 'self'; style-src 'self'; connect-src 'self'; img-src 'self' https://brand.volter.ai; base-uri 'none'; form-action 'self'; frame-ancestors 'none'\r\n\
2153 Referrer-Policy: no-referrer\r\n\
2154 X-Content-Type-Options: nosniff\r\n\
2155 Connection: close\r\n\r\n",
2156 body.len()
2157 );
2158 writer.write_all(head.as_bytes()).await?;
2159 writer.write_all(body).await?;
2160 writer.flush().await
2161}
2162
2163#[cfg(feature = "adapter-api")]
2165#[derive(Clone)]
2166pub struct RuntimeHttpCredential {
2167 token: Arc<str>,
2168 authorization: RuntimeAuthorization,
2169 client_id: Option<crate::RuntimeClientId>,
2170 bootstrap: bool,
2171 runtime_id: Option<Arc<str>>,
2172 generation: Option<[u8; 16]>,
2173 revocation: Option<Arc<RuntimeCredentialRevocation>>,
2174}
2175
2176#[cfg(feature = "adapter-api")]
2177impl RuntimeHttpCredential {
2178 pub fn new(token: impl Into<Arc<str>>, authorization: RuntimeAuthorization) -> Self {
2181 Self {
2182 token: token.into(),
2183 authorization,
2184 client_id: None,
2185 bootstrap: false,
2186 runtime_id: None,
2187 generation: None,
2188 revocation: None,
2189 }
2190 }
2191
2192 pub fn owner(token: impl Into<Arc<str>>) -> Self {
2197 Self {
2198 token: token.into(),
2199 authorization: RuntimeAuthorization::owner(),
2200 client_id: None,
2201 bootstrap: true,
2202 runtime_id: None,
2203 generation: None,
2204 revocation: None,
2205 }
2206 }
2207
2208 pub fn observer(token: impl Into<Arc<str>>) -> Self {
2210 Self::new(token, RuntimeAuthorization::observer())
2211 }
2212
2213 fn frontend(
2214 token: impl Into<Arc<str>>,
2215 client_id: crate::RuntimeClientId,
2216 authorization: RuntimeAuthorization,
2217 runtime_id: impl Into<Arc<str>>,
2218 generation: [u8; 16],
2219 ) -> Self {
2220 Self {
2221 token: token.into(),
2222 authorization,
2223 client_id: Some(client_id),
2224 bootstrap: false,
2225 runtime_id: Some(runtime_id.into()),
2226 generation: Some(generation),
2227 revocation: Some(Arc::new(RuntimeCredentialRevocation::new())),
2228 }
2229 }
2230}
2231
2232#[cfg(feature = "adapter-api")]
2233struct AuthenticatedRuntimeHttpCredential {
2234 authorization: RuntimeAuthorization,
2235 client_id: Option<crate::RuntimeClientId>,
2236 bootstrap: bool,
2237 revocation: Option<tokio::sync::watch::Receiver<bool>>,
2238 attachment: Option<RuntimeCredentialAttachment>,
2239 via_bearer_header: bool,
2240}
2241
2242#[cfg(feature = "adapter-api")]
2243struct RuntimeCredentialRevocation {
2244 signal: tokio::sync::watch::Sender<bool>,
2245 active_attachments: AtomicUsize,
2246 drained: tokio::sync::Notify,
2247}
2248
2249#[cfg(feature = "adapter-api")]
2250impl RuntimeCredentialRevocation {
2251 fn new() -> Self {
2252 let (signal, _) = tokio::sync::watch::channel(false);
2253 Self {
2254 signal,
2255 active_attachments: AtomicUsize::new(0),
2256 drained: tokio::sync::Notify::new(),
2257 }
2258 }
2259
2260 fn register(self: &Arc<Self>) -> RuntimeCredentialAttachment {
2261 self.active_attachments.fetch_add(1, Ordering::AcqRel);
2262 RuntimeCredentialAttachment {
2263 revocation: self.clone(),
2264 }
2265 }
2266
2267 async fn revoke_and_wait(&self) {
2268 let _ = self.signal.send(true);
2269 loop {
2270 let drained = self.drained.notified();
2271 if self.active_attachments.load(Ordering::Acquire) == 0 {
2272 return;
2273 }
2274 drained.await;
2275 }
2276 }
2277}
2278
2279#[cfg(feature = "adapter-api")]
2280struct RuntimeCredentialAttachment {
2281 revocation: Arc<RuntimeCredentialRevocation>,
2282}
2283
2284#[cfg(feature = "adapter-api")]
2285impl Drop for RuntimeCredentialAttachment {
2286 fn drop(&mut self) {
2287 if self
2288 .revocation
2289 .active_attachments
2290 .fetch_sub(1, Ordering::AcqRel)
2291 == 1
2292 {
2293 self.revocation.drained.notify_one();
2298 }
2299 }
2300}
2301
2302#[cfg(feature = "adapter-api")]
2303struct IssuedRuntimeHttpCredential(Arc<str>);
2304
2305#[cfg(feature = "adapter-api")]
2306impl IssuedRuntimeHttpCredential {
2307 fn as_bytes(&self) -> &[u8] {
2308 self.0.as_bytes()
2309 }
2310}
2311
2312#[cfg(feature = "adapter-api")]
2313struct RuntimeHttpCredentialRegistry {
2314 credentials: StdMutex<Vec<RuntimeHttpCredential>>,
2315 runtime_id: String,
2316 generation: [u8; 16],
2317}
2318
2319#[cfg(feature = "adapter-api")]
2320impl RuntimeHttpCredentialRegistry {
2321 fn new(
2322 runtime_id: impl Into<String>,
2323 credentials: Vec<RuntimeHttpCredential>,
2324 ) -> std::io::Result<Arc<Self>> {
2325 let mut generation = [0_u8; 16];
2326 getrandom::getrandom(&mut generation).map_err(|error| {
2327 std::io::Error::other(format!(
2328 "cannot create runtime credential generation: {error}"
2329 ))
2330 })?;
2331 Ok(Arc::new(Self {
2332 credentials: StdMutex::new(credentials),
2333 runtime_id: runtime_id.into(),
2334 generation,
2335 }))
2336 }
2337
2338 fn authenticate(&self, request: &HttpRequest) -> Option<AuthenticatedRuntimeHttpCredential> {
2339 let credentials = self
2340 .credentials
2341 .lock()
2342 .unwrap_or_else(std::sync::PoisonError::into_inner);
2343 debug_assert!(credentials.iter().all(|credential| {
2344 credential.client_id.is_none()
2345 || (credential.runtime_id.as_deref() == Some(self.runtime_id.as_str())
2346 && credential.generation == Some(self.generation))
2347 }));
2348 check_auth(request, &credentials)
2349 }
2350
2351 fn issue_frontend(
2352 &self,
2353 client_id: crate::RuntimeClientId,
2354 observer: bool,
2355 ) -> std::io::Result<IssuedRuntimeHttpCredential> {
2356 let authorization = if observer {
2357 RuntimeAuthorization::observer()
2358 } else {
2359 RuntimeAuthorization::interactive()
2360 };
2361 for _ in 0..3 {
2362 let mut secret = [0_u8; 32];
2363 getrandom::getrandom(&mut secret).map_err(|error| {
2364 std::io::Error::other(format!("cannot mint frontend credential: {error}"))
2365 })?;
2366 let token: Arc<str> = encode_credential(&secret).into();
2367 secret.fill(0);
2368 let mut credentials = self
2369 .credentials
2370 .lock()
2371 .unwrap_or_else(std::sync::PoisonError::into_inner);
2372 if credentials
2373 .iter()
2374 .any(|credential| constant_time_eq(token.as_bytes(), credential.token.as_bytes()))
2375 {
2376 continue;
2377 }
2378 credentials.push(RuntimeHttpCredential::frontend(
2379 token.clone(),
2380 client_id,
2381 authorization,
2382 self.runtime_id.clone(),
2383 self.generation,
2384 ));
2385 return Ok(IssuedRuntimeHttpCredential(token));
2386 }
2387 Err(std::io::Error::new(
2388 std::io::ErrorKind::AlreadyExists,
2389 "frontend credential collision limit exceeded",
2390 ))
2391 }
2392
2393 async fn revoke_client(&self, client_id: &crate::RuntimeClientId) -> bool {
2394 let revocations = {
2397 let mut credentials = self
2398 .credentials
2399 .lock()
2400 .unwrap_or_else(std::sync::PoisonError::into_inner);
2401 let mut revocations = Vec::new();
2402 credentials.retain(|credential| {
2403 if credential.client_id.as_ref() == Some(client_id) {
2404 if let Some(revocation) = &credential.revocation {
2405 revocations.push(revocation.clone());
2406 }
2407 false
2408 } else {
2409 true
2410 }
2411 });
2412 revocations
2413 };
2414 let revoked = !revocations.is_empty();
2415 for revocation in revocations {
2416 revocation.revoke_and_wait().await;
2417 }
2418 revoked
2419 }
2420}
2421
2422#[cfg(feature = "adapter-api")]
2423fn encode_credential(secret: &[u8; 32]) -> String {
2424 const HEX: &[u8; 16] = b"0123456789abcdef";
2425 let mut encoded = String::with_capacity(64);
2426 for byte in secret {
2427 encoded.push(HEX[(byte >> 4) as usize] as char);
2428 encoded.push(HEX[(byte & 0x0f) as usize] as char);
2429 }
2430 encoded
2431}
2432
2433#[cfg(feature = "adapter-api")]
2440fn check_auth(
2441 req: &HttpRequest,
2442 credentials: &[RuntimeHttpCredential],
2443) -> Option<AuthenticatedRuntimeHttpCredential> {
2444 if let Some(auth) = req.headers.get("authorization") {
2445 if let Some(t) = auth.strip_prefix("Bearer ") {
2446 for credential in credentials {
2447 if constant_time_eq(t.as_bytes(), credential.token.as_bytes()) {
2448 return Some(AuthenticatedRuntimeHttpCredential {
2449 authorization: credential.authorization.clone(),
2450 client_id: credential.client_id.clone(),
2451 bootstrap: credential.bootstrap,
2452 revocation: credential
2453 .revocation
2454 .as_ref()
2455 .map(|revocation| revocation.signal.subscribe()),
2456 attachment: credential.revocation.as_ref().and_then(|revocation| {
2457 matches!(req.path.as_str(), "/events" | "/frontend/events")
2458 .then(|| revocation.register())
2459 }),
2460 via_bearer_header: true,
2461 });
2462 }
2463 }
2464 }
2465 }
2466 for pair in req.query.split('&') {
2467 if let Some((k, v)) = pair.split_once('=') {
2468 if k == "token" {
2469 for credential in credentials
2470 .iter()
2471 .filter(|credential| credential.client_id.is_none())
2472 {
2473 if constant_time_eq(v.as_bytes(), credential.token.as_bytes()) {
2474 return Some(AuthenticatedRuntimeHttpCredential {
2475 authorization: credential.authorization.clone(),
2476 client_id: credential.client_id.clone(),
2477 bootstrap: credential.bootstrap,
2478 revocation: None,
2479 attachment: None,
2480 via_bearer_header: false,
2481 });
2482 }
2483 }
2484 }
2485 }
2486 }
2487 None
2488}
2489
2490#[cfg(feature = "adapter-api")]
2491fn coordinated_http_client(
2492 request: &HttpRequest,
2493 coordinator: &Arc<CoordinatedRuntime>,
2494 credential: AuthenticatedRuntimeHttpCredential,
2495) -> Result<Arc<CoordinatedRuntimeClient>, crate::RuntimeLeaseError> {
2496 let supplied_client_id = request
2500 .headers
2501 .get("x-supercode-client-id")
2502 .map(String::as_str);
2503 let client_id = match credential.client_id.as_ref() {
2504 Some(bound) if supplied_client_id == Some(bound.as_str()) => bound.as_str(),
2505 Some(_) => return Err(crate::RuntimeLeaseError::InvalidClientId),
2506 None => supplied_client_id.unwrap_or("legacy-owner"),
2507 };
2508 let mut authorization = credential.authorization;
2509 if let Some(requested) = request.headers.get("x-supercode-permissions") {
2510 authorization = authorization.restrict_to(&RuntimeAuthorization::parse_header(requested)?);
2511 }
2512 Ok(coordinator.client(RuntimeClientId::parse(client_id)?, authorization))
2513}
2514
2515#[cfg(feature = "adapter-api")]
2516async fn coordinated_runtime_rpc(
2517 client: Arc<CoordinatedRuntimeClient>,
2518 request: RpcRequest,
2519) -> Value {
2520 let id = request.id.clone();
2521 let method = crate::FrontendFacadeMethod::from_wire_name(&request.method);
2522 let result = match method {
2523 Some(crate::FrontendFacadeMethod::AcquireControl) => client
2524 .acquire_control()
2525 .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2526 Some(crate::FrontendFacadeMethod::TakeControl) => client
2527 .take_control()
2528 .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2529 Some(crate::FrontendFacadeMethod::Heartbeat) => client
2530 .heartbeat()
2531 .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2532 Some(crate::FrontendFacadeMethod::Lease) => client
2533 .lease_snapshot()
2534 .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2535 Some(crate::FrontendFacadeMethod::Detach) => {
2536 serde_json::to_value(client.detach()).map_err(json_sdk_error)
2537 }
2538 Some(crate::FrontendFacadeMethod::Close) => match client.close().await {
2539 Ok(()) => Ok(json!({"closed":true})),
2540 Err(error) => Err(error),
2541 },
2542 None if request.method == "shutdown" => match client.close().await {
2543 Ok(()) => Ok(json!({"shutting_down":true})),
2544 Err(error) => Err(error),
2545 },
2546 _ => return frontend_http_rpc(client, request).await,
2547 };
2548 match result {
2549 Ok(value) => rpc_ok(id, value),
2550 Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
2551 }
2552}
2553
2554#[cfg(feature = "adapter-api")]
2555fn json_sdk_error(error: serde_json::Error) -> FrontendRuntimeError {
2556 FrontendRuntimeError::Transport(error.to_string())
2557}
2558
2559#[cfg(feature = "adapter-api")]
2571async fn serve_frontend_credential_door<W: AsyncWrite + Unpin>(
2572 req: &HttpRequest,
2573 write_half: &mut W,
2574 peer_is_loopback: bool,
2575 credential: &AuthenticatedRuntimeHttpCredential,
2576 credentials: &RuntimeHttpCredentialRegistry,
2577 coordinator: &Arc<CoordinatedRuntime>,
2578) -> Option<std::io::Result<()>> {
2579 if matches!(
2580 req.path.as_str(),
2581 "/_supercode/frontend-credentials/mint" | "/_supercode/frontend-credentials/revoke"
2582 ) && !credential.via_bearer_header
2583 {
2584 let body =
2585 sdk_runtime_rpc_error(Value::Null, -32030, &FrontendRuntimeError::Unauthenticated)
2586 .to_string();
2587 return Some(
2588 write_http_response(
2589 write_half,
2590 401,
2591 "Unauthorized",
2592 "application/json",
2593 body.as_bytes(),
2594 )
2595 .await,
2596 );
2597 }
2598 if req.path == "/_supercode/frontend-credentials/mint" {
2599 if req.method != "POST" || !peer_is_loopback || !credential.bootstrap {
2600 let body = sdk_runtime_rpc_error(
2601 Value::Null,
2602 -32031,
2603 &FrontendRuntimeError::Unauthorized {
2604 permission: "bootstrap".into(),
2605 },
2606 )
2607 .to_string();
2608 return Some(
2609 write_http_response(
2610 write_half,
2611 403,
2612 "Forbidden",
2613 "application/json",
2614 body.as_bytes(),
2615 )
2616 .await,
2617 );
2618 }
2619 let request: Value = match serde_json::from_slice(&req.body) {
2620 Ok(request) => request,
2621 Err(error) => {
2622 return Some(
2623 write_http_response(
2624 write_half,
2625 400,
2626 "Bad Request",
2627 "application/json",
2628 format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2629 )
2630 .await,
2631 );
2632 }
2633 };
2634 let Some(client_id) = request.get("clientId").and_then(Value::as_str) else {
2635 return Some(
2636 write_http_response(
2637 write_half,
2638 400,
2639 "Bad Request",
2640 "application/json",
2641 b"{\"error\":\"mint request omitted clientId\"}",
2642 )
2643 .await,
2644 );
2645 };
2646 let client_id = match crate::RuntimeClientId::parse(client_id) {
2647 Ok(client_id) => client_id,
2648 Err(error) => {
2649 return Some(
2650 write_http_response(
2651 write_half,
2652 400,
2653 "Bad Request",
2654 "application/json",
2655 format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2656 )
2657 .await,
2658 );
2659 }
2660 };
2661 let observer = match request.get("grant").and_then(Value::as_str) {
2662 Some("observer") => true,
2663 Some("interactive") => false,
2664 _ => {
2665 return Some(
2666 write_http_response(
2667 write_half,
2668 400,
2669 "Bad Request",
2670 "application/json",
2671 b"{\"error\":\"grant must be observer or interactive\"}",
2672 )
2673 .await,
2674 );
2675 }
2676 };
2677 let token = match credentials.issue_frontend(client_id, observer) {
2678 Ok(token) => token,
2679 Err(error) => return Some(Err(error)),
2680 };
2681 return Some(
2682 write_http_response(
2683 write_half,
2684 200,
2685 "OK",
2686 "application/octet-stream",
2687 token.as_bytes(),
2688 )
2689 .await,
2690 );
2691 }
2692
2693 if req.path == "/_supercode/frontend-credentials/revoke" {
2694 if req.method != "POST" || !peer_is_loopback || !credential.bootstrap {
2695 let body = sdk_runtime_rpc_error(
2696 Value::Null,
2697 -32031,
2698 &FrontendRuntimeError::Unauthorized {
2699 permission: "bootstrap".into(),
2700 },
2701 )
2702 .to_string();
2703 return Some(
2704 write_http_response(
2705 write_half,
2706 403,
2707 "Forbidden",
2708 "application/json",
2709 body.as_bytes(),
2710 )
2711 .await,
2712 );
2713 }
2714 let request: Value = match serde_json::from_slice(&req.body) {
2715 Ok(request) => request,
2716 Err(error) => {
2717 return Some(
2718 write_http_response(
2719 write_half,
2720 400,
2721 "Bad Request",
2722 "application/json",
2723 format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2724 )
2725 .await,
2726 );
2727 }
2728 };
2729 let Some(client_id) = request.get("clientId").and_then(Value::as_str) else {
2730 return Some(
2731 write_http_response(
2732 write_half,
2733 400,
2734 "Bad Request",
2735 "application/json",
2736 b"{\"error\":\"revoke request omitted clientId\"}",
2737 )
2738 .await,
2739 );
2740 };
2741 let client_id = match crate::RuntimeClientId::parse(client_id) {
2742 Ok(client_id) => client_id,
2743 Err(error) => {
2744 return Some(
2745 write_http_response(
2746 write_half,
2747 400,
2748 "Bad Request",
2749 "application/json",
2750 format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2751 )
2752 .await,
2753 );
2754 }
2755 };
2756 let revoked = credentials.revoke_client(&client_id).await;
2757 if revoked {
2758 coordinator
2759 .client(client_id, RuntimeAuthorization::observer())
2760 .detach();
2761 }
2762 return Some(
2763 write_http_response(
2764 write_half,
2765 200,
2766 "OK",
2767 "application/json",
2768 if revoked {
2769 b"{\"revoked\":true}"
2770 } else {
2771 b"{\"revoked\":false}"
2772 },
2773 )
2774 .await,
2775 );
2776 }
2777
2778 None
2779}
2780
2781#[cfg(feature = "adapter-api")]
2782async fn handle_http_conn(
2783 stream: tokio::net::TcpStream,
2784 engine: Arc<RpcEngine>,
2785 coordinator: Arc<CoordinatedRuntime>,
2786 credentials: Arc<RuntimeHttpCredentialRegistry>,
2787) -> std::io::Result<()> {
2788 let peer_is_loopback = stream.peer_addr()?.ip().is_loopback();
2789 let (read_half, mut write_half) = stream.into_split();
2790 let mut reader = tokio::io::BufReader::new(read_half);
2791 let Some(req) = read_http_request(&mut reader).await? else {
2792 return Ok(());
2793 };
2794
2795 if req.method == "GET" {
2799 if let Some((content_type, body)) = browser_observer_asset(&req.path) {
2800 return write_browser_observer_asset(&mut write_half, content_type, body).await;
2801 }
2802 }
2803
2804 let Some(credential) = credentials.authenticate(&req) else {
2805 let body =
2806 sdk_runtime_rpc_error(Value::Null, -32030, &FrontendRuntimeError::Unauthenticated)
2807 .to_string();
2808 return write_http_response(
2809 &mut write_half,
2810 401,
2811 "Unauthorized",
2812 "application/json",
2813 body.as_bytes(),
2814 )
2815 .await;
2816 };
2817
2818 if let Some(result) = serve_frontend_credential_door(
2819 &req,
2820 &mut write_half,
2821 peer_is_loopback,
2822 &credential,
2823 &credentials,
2824 &coordinator,
2825 )
2826 .await
2827 {
2828 return result;
2829 }
2830
2831 let mut revocation = credential.revocation.clone();
2832 let mut attachment = credential.attachment;
2833 let credential = AuthenticatedRuntimeHttpCredential {
2834 authorization: credential.authorization,
2835 client_id: credential.client_id,
2836 bootstrap: credential.bootstrap,
2837 revocation: None,
2838 attachment: None,
2839 via_bearer_header: credential.via_bearer_header,
2840 };
2841 let client = match coordinated_http_client(&req, &coordinator, credential) {
2842 Ok(client) => client,
2843 Err(error) => {
2844 let permission = match error {
2845 crate::RuntimeLeaseError::InvalidClientId => "client_id",
2846 crate::RuntimeLeaseError::InvalidAuthorization => "authorization",
2847 _ => "runtime",
2848 };
2849 let body = sdk_runtime_rpc_error(
2850 Value::Null,
2851 -32031,
2852 &FrontendRuntimeError::Unauthorized {
2853 permission: permission.into(),
2854 },
2855 )
2856 .to_string();
2857 return write_http_response(
2858 &mut write_half,
2859 403,
2860 "Forbidden",
2861 "application/json",
2862 body.as_bytes(),
2863 )
2864 .await;
2865 }
2866 };
2867
2868 match (req.method.as_str(), req.path.as_str()) {
2869 ("POST", "/rpc") => {
2870 let body_text = String::from_utf8_lossy(&req.body);
2871 let resp = match serde_json::from_str::<RpcRequest>(&body_text) {
2872 Ok(rpc_req) if matches!(rpc_req.method.as_str(), "status" | "history") => {
2873 engine.handle_request(rpc_req).await
2874 }
2875 Ok(rpc_req) => coordinated_runtime_rpc(client.clone(), rpc_req).await,
2876 Err(e) => rpc_error(Value::Null, -32700, format!("parse error: {e}")),
2877 };
2878 let body = resp.to_string();
2879 write_http_response(
2880 &mut write_half,
2881 200,
2882 "OK",
2883 "application/json",
2884 body.as_bytes(),
2885 )
2886 .await
2887 }
2888 ("GET", "/events") => {
2889 if let Err(error) = client.observe() {
2890 let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
2891 return write_http_response(
2892 &mut write_half,
2893 403,
2894 "Forbidden",
2895 "application/json",
2896 body.as_bytes(),
2897 )
2898 .await;
2899 }
2900 let mut events = engine.subscribe();
2904 let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
2905 if write_half.write_all(head.as_bytes()).await.is_err() {
2906 client.detach();
2907 return Ok(());
2908 }
2909 let _ = write_half.flush().await;
2910 let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
2914 keepalive.tick().await;
2915 loop {
2916 tokio::select! {
2917 biased;
2918 recv = events.recv() => {
2919 match recv {
2920 Ok(v) => {
2921 let line = format!("data: {v}\n\n");
2922 if write_half.write_all(line.as_bytes()).await.is_err() {
2923 break;
2924 }
2925 if write_half.flush().await.is_err() {
2926 break;
2927 }
2928 }
2929 Err(broadcast::error::RecvError::Lagged(_)) => continue,
2930 Err(broadcast::error::RecvError::Closed) => break,
2931 }
2932 }
2933 _ = engine.wait_for_shutdown() => break,
2938 _ = wait_for_credential_revocation(&mut revocation) => break,
2939 _ = reader.read_u8() => break,
2940 _ = keepalive.tick() => {
2941 if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
2942 break;
2943 }
2944 }
2945 }
2946 }
2947 let _ = write_half.shutdown().await;
2948 client.detach();
2949 drop(attachment.take());
2950 Ok(())
2951 }
2952 ("GET", "/frontend/events") => {
2953 if let Err(error) = client.observe() {
2954 let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
2955 return write_http_response(
2956 &mut write_half,
2957 403,
2958 "Forbidden",
2959 "application/json",
2960 body.as_bytes(),
2961 )
2962 .await;
2963 }
2964 let mut events = engine.frontend_subscribe();
2969 let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
2970 if write_half.write_all(head.as_bytes()).await.is_err() {
2971 client.detach();
2972 return Ok(());
2973 }
2974 let _ = write_half.flush().await;
2975 let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
2979 keepalive.tick().await;
2980 loop {
2981 tokio::select! {
2982 biased;
2983 recv = events.recv() => {
2984 match recv {
2985 Ok(event) => {
2986 let value = serde_json::to_string(&event).unwrap_or_default();
2987 let line = format!("data: {value}\n\n");
2988 if write_half.write_all(line.as_bytes()).await.is_err() {
2989 break;
2990 }
2991 if write_half.flush().await.is_err() {
2992 break;
2993 }
2994 }
2995 Err(broadcast::error::RecvError::Lagged(_)) => break,
2996 Err(broadcast::error::RecvError::Closed) => break,
2997 }
2998 }
2999 _ = engine.wait_for_shutdown() => break,
3003 _ = wait_for_credential_revocation(&mut revocation) => break,
3004 _ = reader.read_u8() => break,
3005 _ = keepalive.tick() => {
3006 if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
3007 break;
3008 }
3009 }
3010 }
3011 }
3012 let _ = write_half.shutdown().await;
3013 client.detach();
3014 drop(attachment.take());
3015 Ok(())
3016 }
3017 _ => {
3018 write_http_response(
3019 &mut write_half,
3020 404,
3021 "Not Found",
3022 "application/json",
3023 b"{\"error\":\"not found\"}",
3024 )
3025 .await
3026 }
3027 }
3028}
3029
3030#[cfg(feature = "adapter-api")]
3031async fn wait_for_credential_revocation(receiver: &mut Option<tokio::sync::watch::Receiver<bool>>) {
3032 let Some(receiver) = receiver else {
3033 std::future::pending::<()>().await;
3034 return;
3035 };
3036 if *receiver.borrow() {
3037 return;
3038 }
3039 while receiver.changed().await.is_ok() {
3040 if *receiver.borrow() {
3041 return;
3042 }
3043 }
3044}
3045
3046#[cfg(feature = "adapter-api")]
3047async fn handle_frontend_http_conn(
3048 stream: tokio::net::TcpStream,
3049 coordinator: Arc<CoordinatedRuntime>,
3050 events: broadcast::Sender<FrontendEvent>,
3051 credentials: Arc<RuntimeHttpCredentialRegistry>,
3052) -> std::io::Result<()> {
3053 let peer_is_loopback = stream.peer_addr()?.ip().is_loopback();
3054 let (read_half, mut write_half) = stream.into_split();
3055 let mut reader = tokio::io::BufReader::new(read_half);
3056 let Some(req) = read_http_request(&mut reader).await? else {
3057 return Ok(());
3058 };
3059 if req.method == "GET" {
3060 if let Some((content_type, body)) = browser_observer_asset(&req.path) {
3061 return write_browser_observer_asset(&mut write_half, content_type, body).await;
3062 }
3063 }
3064 let Some(credential) = credentials.authenticate(&req) else {
3065 return write_http_response(
3066 &mut write_half,
3067 401,
3068 "Unauthorized",
3069 "application/json",
3070 b"{\"error\":\"missing or invalid bearer token\"}",
3071 )
3072 .await;
3073 };
3074 if let Some(result) = serve_frontend_credential_door(
3077 &req,
3078 &mut write_half,
3079 peer_is_loopback,
3080 &credential,
3081 &credentials,
3082 &coordinator,
3083 )
3084 .await
3085 {
3086 return result;
3087 }
3088 let client = match coordinated_http_client(&req, &coordinator, credential) {
3089 Ok(client) => client,
3090 Err(error) => {
3091 let body = json!({"error":error.to_string()}).to_string();
3092 return write_http_response(
3093 &mut write_half,
3094 400,
3095 "Bad Request",
3096 "application/json",
3097 body.as_bytes(),
3098 )
3099 .await;
3100 }
3101 };
3102 match (req.method.as_str(), req.path.as_str()) {
3103 ("POST", "/rpc") => {
3104 let body_text = String::from_utf8_lossy(&req.body);
3105 let response = match serde_json::from_str::<RpcRequest>(&body_text) {
3106 Ok(request) => coordinated_runtime_rpc(client.clone(), request).await,
3107 Err(error) => rpc_error(Value::Null, -32700, format!("parse error: {error}")),
3108 };
3109 let body = response.to_string();
3110 write_http_response(
3111 &mut write_half,
3112 200,
3113 "OK",
3114 "application/json",
3115 body.as_bytes(),
3116 )
3117 .await
3118 }
3119 ("GET", "/frontend/events") => {
3120 if let Err(error) = client.observe() {
3121 let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
3122 return write_http_response(
3123 &mut write_half,
3124 403,
3125 "Forbidden",
3126 "application/json",
3127 body.as_bytes(),
3128 )
3129 .await;
3130 }
3131 let mut receiver = events.subscribe();
3132 let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
3133 if write_half.write_all(head.as_bytes()).await.is_err() {
3134 client.detach();
3135 return Ok(());
3136 }
3137 let _ = write_half.flush().await;
3138 let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
3139 keepalive.tick().await;
3140 loop {
3141 tokio::select! {
3142 _ = reader.read_u8() => break,
3143 _ = keepalive.tick() => {
3144 if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
3145 break;
3146 }
3147 }
3148 event = receiver.recv() => match event {
3149 Ok(event) => {
3150 let value = serde_json::to_string(&event).unwrap_or_default();
3151 let line = format!("data: {value}\n\n");
3152 if write_half.write_all(line.as_bytes()).await.is_err() || write_half.flush().await.is_err() {
3153 break;
3154 }
3155 }
3156 Err(broadcast::error::RecvError::Lagged(_)) => break,
3157 Err(broadcast::error::RecvError::Closed) => break,
3158 }
3159 }
3160 }
3161 client.detach();
3162 Ok(())
3163 }
3164 _ => {
3165 write_http_response(
3166 &mut write_half,
3167 404,
3168 "Not Found",
3169 "application/json",
3170 b"{\"error\":\"not found\"}",
3171 )
3172 .await
3173 }
3174 }
3175}
3176
3177#[cfg(feature = "adapter-api")]
3178async fn frontend_http_rpc(runtime: Arc<dyn FrontendRuntime>, request: RpcRequest) -> Value {
3179 let id = request.id;
3180 let Some(method) = crate::FrontendFacadeMethod::from_wire_name(&request.method) else {
3181 return rpc_error(id, -32601, format!("unknown method `{}`", request.method));
3182 };
3183 match method {
3184 crate::FrontendFacadeMethod::Describe => match runtime.describe().await {
3185 Ok(descriptor) => rpc_ok(id, serde_json::to_value(descriptor).unwrap_or_default()),
3186 Err(error) => sdk_runtime_rpc_error(id, -32010, &error),
3187 },
3188 crate::FrontendFacadeMethod::Attach => {
3189 let limit = request
3190 .params
3191 .get("limit")
3192 .and_then(Value::as_u64)
3193 .unwrap_or(50)
3194 .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
3195 match runtime.attach(limit).await {
3196 Ok(attachment) => rpc_ok(
3197 id,
3198 serde_json::to_value(FrontendAttachSnapshot {
3199 descriptor: attachment.descriptor,
3200 history: attachment.history,
3201 history_cursor: attachment.history_cursor,
3202 replay: attachment.replay,
3203 })
3204 .unwrap_or_default(),
3205 ),
3206 Err(error) => sdk_runtime_rpc_error(id, -32010, &error),
3207 }
3208 }
3209 crate::FrontendFacadeMethod::SendInput => {
3210 let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3211 return rpc_error(
3212 id,
3213 -32602,
3214 "frontend.send_input requires a string `params.prompt`",
3215 );
3216 };
3217 let image_urls = match parse_image_urls(&request.params, "frontend.send_input") {
3218 Ok(image_urls) => image_urls,
3219 Err(message) => return rpc_error(id, -32602, message),
3220 };
3221 match runtime
3222 .clone()
3223 .send_input_with_images(prompt.to_string(), image_urls)
3224 .await
3225 {
3226 Ok(()) => rpc_ok(id, json!({"accepted": true})),
3227 Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3228 sdk_runtime_rpc_error(id, -32000, &error)
3229 }
3230 Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3231 }
3232 }
3233 crate::FrontendFacadeMethod::Invoke => {
3234 let operation = request
3235 .params
3236 .get("operation")
3237 .cloned()
3238 .ok_or("frontend.invoke requires `params.operation`")
3239 .and_then(|value| {
3240 serde_json::from_value(value).map_err(|_| "invalid frontend operation")
3241 });
3242 match operation {
3243 Ok(operation) => match runtime.invoke(operation).await {
3244 Ok(result) => rpc_ok(id, serde_json::to_value(result).unwrap_or_default()),
3245 Err(error @ FrontendRuntimeError::UnsupportedOperation(_)) => {
3246 sdk_runtime_rpc_error(id, -32023, &error)
3247 }
3248 Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3249 sdk_runtime_rpc_error(id, -32000, &error)
3250 }
3251 Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
3252 sdk_runtime_rpc_error(id, -32001, &error)
3253 }
3254 Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3255 },
3256 Err(message) => rpc_error(id, -32602, message),
3257 }
3258 }
3259 crate::FrontendFacadeMethod::Submit => {
3260 let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3261 return rpc_error(id, -32602, "submit requires a string `params.prompt`");
3262 };
3263 let image_urls = match request.params.get("image_urls") {
3264 None => Vec::new(),
3265 Some(Value::Array(values)) => {
3266 let Some(urls) = values.iter().map(Value::as_str).collect::<Option<Vec<_>>>()
3267 else {
3268 return rpc_error(
3269 id,
3270 -32602,
3271 "submit requires string entries in `params.image_urls`",
3272 );
3273 };
3274 urls.into_iter().map(str::to_owned).collect()
3275 }
3276 Some(_) => {
3277 return rpc_error(id, -32602, "submit requires array `params.image_urls`")
3278 }
3279 };
3280 match runtime
3281 .submit_with_images(prompt.to_string(), image_urls)
3282 .await
3283 {
3284 Ok(reply) => rpc_ok(id, json!({"reply":reply})),
3285 Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3286 sdk_runtime_rpc_error(id, -32000, &error)
3287 }
3288 Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
3289 sdk_runtime_rpc_error(id, -32001, &error)
3290 }
3291 Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3292 }
3293 }
3294 crate::FrontendFacadeMethod::Interrupt => match runtime.interrupt().await {
3295 Ok(interrupted) => rpc_ok(id, json!({"interrupted":interrupted})),
3296 Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3297 },
3298 crate::FrontendFacadeMethod::Steer => {
3299 let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3300 return rpc_error(id, -32602, "steer requires a string `params.prompt`");
3301 };
3302 match runtime.steer(prompt.to_string()).await {
3303 Ok(()) => rpc_ok(id, json!({"queued":true})),
3304 Err(error @ FrontendRuntimeError::UnsupportedAction(_)) => {
3305 sdk_runtime_rpc_error(id, -32020, &error)
3306 }
3307 Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3308 }
3309 }
3310 crate::FrontendFacadeMethod::Respond => {
3311 let response = request
3312 .params
3313 .get("response")
3314 .cloned()
3315 .ok_or("respond requires `params.response`")
3316 .and_then(|value| serde_json::from_value(value).map_err(|_| "invalid response"));
3317 match response {
3318 Ok(response) => match runtime.respond(response).await {
3319 Ok(()) => rpc_ok(id, json!({"accepted":true})),
3320 Err(error @ FrontendRuntimeError::UnsupportedAction(_)) => {
3321 sdk_runtime_rpc_error(id, -32020, &error)
3322 }
3323 Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3324 },
3325 Err(message) => rpc_error(id, -32602, message),
3326 }
3327 }
3328 crate::FrontendFacadeMethod::Lease
3329 | crate::FrontendFacadeMethod::AcquireControl
3330 | crate::FrontendFacadeMethod::TakeControl
3331 | crate::FrontendFacadeMethod::Heartbeat
3332 | crate::FrontendFacadeMethod::Detach
3333 | crate::FrontendFacadeMethod::Close => rpc_error(
3334 id,
3335 -32020,
3336 format!(
3337 "frontend action `{}` requires a coordinated runtime",
3338 method.id()
3339 ),
3340 ),
3341 }
3342}
3343
3344fn parse_image_urls(params: &Value, operation: &str) -> std::result::Result<Vec<String>, String> {
3345 let urls = match params.get("image_urls") {
3346 None => Ok(Vec::new()),
3347 Some(Value::Array(values)) => values
3348 .iter()
3349 .map(|value| {
3350 value.as_str().map(str::to_owned).ok_or_else(|| {
3351 format!("{operation} requires string entries in `params.image_urls`")
3352 })
3353 })
3354 .collect(),
3355 Some(_) => Err(format!("{operation} requires array `params.image_urls`")),
3356 }?;
3357 validate_frontend_image_urls(urls, operation)
3358}
3359
3360fn validate_frontend_image_urls(
3361 urls: Vec<String>,
3362 operation: &str,
3363) -> std::result::Result<Vec<String>, String> {
3364 if urls.len() > 4 {
3365 return Err(format!("{operation} accepts at most 4 images"));
3366 }
3367 let mut total = 0usize;
3368 for url in &urls {
3369 if !(url.starts_with("data:image/")
3370 || url.starts_with("https://")
3371 || url.starts_with("http://"))
3372 {
3373 return Err(format!(
3374 "{operation} images must be image data URLs or HTTP(S) URLs"
3375 ));
3376 }
3377 if url.len() > 12 * 1024 * 1024 {
3378 return Err(format!("{operation} image exceeds the encoded size limit"));
3379 }
3380 total = total.saturating_add(url.len());
3381 }
3382 if total > 32 * 1024 * 1024 {
3383 return Err(format!(
3384 "{operation} images exceed the encoded total size limit"
3385 ));
3386 }
3387 Ok(urls)
3388}
3389
3390#[cfg(feature = "adapter-api")]
3392pub(crate) struct FrontendHttpServer {
3393 address: SocketAddr,
3394 task: tokio::task::JoinHandle<()>,
3395}
3396
3397#[cfg(feature = "adapter-api")]
3398impl FrontendHttpServer {
3399 pub(crate) fn address(&self) -> SocketAddr {
3401 self.address
3402 }
3403}
3404
3405#[cfg(feature = "adapter-api")]
3406impl Drop for FrontendHttpServer {
3407 fn drop(&mut self) {
3408 self.task.abort();
3409 }
3410}
3411
3412#[cfg(feature = "adapter-api")]
3415pub(crate) async fn run_frontend_http(
3416 runtime: Arc<dyn FrontendRuntime>,
3417 events: broadcast::Sender<FrontendEvent>,
3418 bind: &str,
3419 token: Arc<str>,
3420 runtime_id: impl Into<String>,
3421) -> std::io::Result<FrontendHttpServer> {
3422 let listener = TcpListener::bind(bind).await?;
3423 let address = listener.local_addr()?;
3424 let coordinator = CoordinatedRuntime::new(runtime);
3425 let credentials =
3429 RuntimeHttpCredentialRegistry::new(runtime_id, vec![RuntimeHttpCredential::owner(token)])?;
3430 let task = tokio::spawn(async move {
3431 while let Ok((stream, _)) = listener.accept().await {
3432 let coordinator = coordinator.clone();
3433 let events = events.clone();
3434 let credentials = credentials.clone();
3435 tokio::spawn(async move {
3436 let _ = handle_frontend_http_conn(stream, coordinator, events, credentials).await;
3437 });
3438 }
3439 });
3440 Ok(FrontendHttpServer { address, task })
3441}
3442
3443#[cfg(feature = "adapter-api")]
3446pub struct FrontendWebSocketServer {
3447 address: SocketAddr,
3448 task: tokio::task::JoinHandle<()>,
3449}
3450
3451#[cfg(feature = "adapter-api")]
3452impl FrontendWebSocketServer {
3453 pub fn address(&self) -> SocketAddr {
3455 self.address
3456 }
3457}
3458
3459#[cfg(feature = "adapter-api")]
3460impl Drop for FrontendWebSocketServer {
3461 fn drop(&mut self) {
3462 self.task.abort();
3463 }
3464}
3465
3466#[cfg(feature = "adapter-api")]
3470pub async fn run_frontend_websocket(
3471 engine: Arc<RpcEngine>,
3472 bind: &str,
3473 credentials: Vec<RuntimeHttpCredential>,
3474) -> std::io::Result<FrontendWebSocketServer> {
3475 let runtime: Arc<dyn FrontendRuntime> = engine.clone();
3476 let events = engine.frontend_events.clone();
3477 run_frontend_websocket_runtime_inner(runtime, events, Some(engine), bind, credentials).await
3478}
3479
3480#[cfg(all(feature = "adapter-api", test))]
3485pub(crate) async fn run_frontend_websocket_runtime(
3486 runtime: Arc<dyn FrontendRuntime>,
3487 events: broadcast::Sender<FrontendEvent>,
3488 bind: &str,
3489 credentials: Vec<RuntimeHttpCredential>,
3490) -> std::io::Result<FrontendWebSocketServer> {
3491 run_frontend_websocket_runtime_inner(runtime, events, None, bind, credentials).await
3492}
3493
3494#[cfg(feature = "adapter-api")]
3495async fn run_frontend_websocket_runtime_inner(
3496 runtime: Arc<dyn FrontendRuntime>,
3497 events: broadcast::Sender<FrontendEvent>,
3498 shutdown_engine: Option<Arc<RpcEngine>>,
3499 bind: &str,
3500 credentials: Vec<RuntimeHttpCredential>,
3501) -> std::io::Result<FrontendWebSocketServer> {
3502 if credentials.is_empty()
3503 || credentials
3504 .iter()
3505 .any(|credential| credential.token.is_empty())
3506 {
3507 return Err(std::io::Error::new(
3508 std::io::ErrorKind::InvalidInput,
3509 "at least one non-empty runtime WebSocket credential is required",
3510 ));
3511 }
3512 let listener = TcpListener::bind(bind).await?;
3513 let address = listener.local_addr()?;
3514 let coordinator = CoordinatedRuntime::new(runtime);
3515 let credentials: Arc<[RuntimeHttpCredential]> = credentials.into();
3516 let task = tokio::spawn(async move {
3517 loop {
3518 tokio::select! {
3519 biased;
3520 _ = wait_for_optional_runtime_shutdown(shutdown_engine.as_ref()) => break,
3521 accepted = listener.accept() => {
3522 let Ok((stream, _)) = accepted else { continue };
3523 let coordinator = coordinator.clone();
3524 let credentials = credentials.clone();
3525 let events = events.clone();
3526 let shutdown_engine = shutdown_engine.clone();
3527 tokio::spawn(async move {
3528 let _ = handle_frontend_websocket(stream, events, shutdown_engine, coordinator, credentials).await;
3529 });
3530 }
3531 }
3532 }
3533 });
3534 Ok(FrontendWebSocketServer { address, task })
3535}
3536
3537#[cfg(feature = "adapter-api")]
3538async fn wait_for_optional_runtime_shutdown(engine: Option<&Arc<RpcEngine>>) {
3539 match engine {
3540 Some(engine) => engine.wait_for_shutdown().await,
3541 None => std::future::pending().await,
3542 }
3543}
3544
3545#[cfg(feature = "adapter-api")]
3546#[allow(clippy::result_large_err)] async fn handle_frontend_websocket(
3548 stream: tokio::net::TcpStream,
3549 events: broadcast::Sender<FrontendEvent>,
3550 shutdown_engine: Option<Arc<RpcEngine>>,
3551 coordinator: Arc<CoordinatedRuntime>,
3552 credentials: Arc<[RuntimeHttpCredential]>,
3553) -> Result<(), tokio_tungstenite::tungstenite::Error> {
3554 use std::sync::Mutex as SyncMutex;
3555 use tokio_tungstenite::tungstenite::handshake::server::{ErrorResponse, Request, Response};
3556
3557 let selected = Arc::new(SyncMutex::new(None::<Arc<CoordinatedRuntimeClient>>));
3558 let selected_by_callback = selected.clone();
3559 let socket = tokio_tungstenite::accept_hdr_async(
3560 stream,
3561 move |request: &Request, response: Response| -> Result<Response, ErrorResponse> {
3562 let reject = |status, message: &str| {
3563 tokio_tungstenite::tungstenite::http::Response::builder()
3564 .status(status)
3565 .body(Some(message.to_string()))
3566 .expect("static WebSocket rejection is valid")
3567 };
3568 if request.uri().path() != "/frontend/v2" {
3569 return Err(reject(404, "frontend WebSocket route not found"));
3570 }
3571 let token = request
3572 .headers()
3573 .get("authorization")
3574 .and_then(|value| value.to_str().ok())
3575 .and_then(|value| value.strip_prefix("Bearer "));
3576 let Some(credential) = token.and_then(|token| {
3577 credentials.iter().find(|credential| {
3578 constant_time_eq(token.as_bytes(), credential.token.as_bytes())
3579 })
3580 }) else {
3581 return Err(reject(401, "missing or invalid bearer token"));
3582 };
3583 let client_id = request
3584 .headers()
3585 .get("x-supercode-client-id")
3586 .and_then(|value| value.to_str().ok())
3587 .unwrap_or("legacy-websocket-owner");
3588 let Ok(client_id) = RuntimeClientId::parse(client_id) else {
3589 return Err(reject(400, "invalid runtime client id"));
3590 };
3591 let mut authorization = credential.authorization.clone();
3592 if let Some(requested) = request
3593 .headers()
3594 .get("x-supercode-permissions")
3595 .and_then(|value| value.to_str().ok())
3596 {
3597 let Ok(requested) = RuntimeAuthorization::parse_header(requested) else {
3598 return Err(reject(400, "invalid runtime authorization grant"));
3599 };
3600 authorization = authorization.restrict_to(&requested);
3601 }
3602 *selected_by_callback
3603 .lock()
3604 .unwrap_or_else(std::sync::PoisonError::into_inner) =
3605 Some(coordinator.client(client_id, authorization));
3606 Ok(response)
3607 },
3608 )
3609 .await?;
3610 let client = selected
3611 .lock()
3612 .unwrap_or_else(std::sync::PoisonError::into_inner)
3613 .take()
3614 .expect("successful WebSocket handshake selects a runtime client");
3615 if let Err(error) = client.observe() {
3616 let mut socket = socket;
3617 let value = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
3618 socket
3619 .send(tokio_tungstenite::tungstenite::Message::Text(value.into()))
3620 .await?;
3621 socket.close(None).await?;
3622 return Ok(());
3623 }
3624
3625 let mut events = events.subscribe();
3626 let (mut writer, mut reader) = socket.split();
3627 loop {
3628 tokio::select! {
3629 biased;
3630 incoming = reader.next() => match incoming {
3631 Some(Ok(tokio_tungstenite::tungstenite::Message::Text(text))) => {
3632 let response = match serde_json::from_str::<RpcRequest>(&text) {
3633 Ok(request) => coordinated_runtime_rpc(client.clone(), request).await,
3634 Err(error) => rpc_error(Value::Null, -32700, format!("parse error: {error}")),
3635 };
3636 writer.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await?;
3637 }
3638 Some(Ok(tokio_tungstenite::tungstenite::Message::Ping(payload))) => {
3639 writer.send(tokio_tungstenite::tungstenite::Message::Pong(payload)).await?;
3640 }
3641 Some(Ok(tokio_tungstenite::tungstenite::Message::Close(_))) | None => break,
3642 Some(Ok(_)) => {}
3643 Some(Err(error)) => {
3644 client.detach();
3645 return Err(error);
3646 }
3647 },
3648 event = events.recv() => match event {
3649 Ok(event) => {
3650 let notification = json!({
3651 "jsonrpc":"2.0",
3652 "method":"frontend.v2.event",
3653 "params":{"event":event},
3654 });
3655 writer.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await?;
3656 }
3657 Err(broadcast::error::RecvError::Lagged(count)) => {
3658 let notification = json!({
3659 "jsonrpc":"2.0",
3660 "method":"frontend.v2.event",
3661 "params":{"error":{"name":"transport","message":format!("event replay gap: {count}")}},
3662 });
3663 writer.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await?;
3664 break;
3665 }
3666 Err(broadcast::error::RecvError::Closed) => break,
3667 },
3668 _ = wait_for_optional_runtime_shutdown(shutdown_engine.as_ref()) => break,
3669 }
3670 }
3671 client.detach();
3672 Ok(())
3673}
3674
3675#[cfg(feature = "adapter-api")]
3683pub async fn run_http(
3684 engine: Arc<RpcEngine>,
3685 bind: &str,
3686 token: Arc<str>,
3687) -> std::io::Result<SocketAddr> {
3688 run_http_authorized(engine, bind, vec![RuntimeHttpCredential::owner(token)]).await
3689}
3690
3691#[cfg(feature = "adapter-api")]
3696pub async fn run_http_authorized(
3697 engine: Arc<RpcEngine>,
3698 bind: &str,
3699 credentials: Vec<RuntimeHttpCredential>,
3700) -> std::io::Result<SocketAddr> {
3701 run_http_authorized_with_lease_ttl(
3702 engine,
3703 bind,
3704 credentials,
3705 crate::DEFAULT_RUNTIME_LEASE_TTL_MS,
3706 )
3707 .await
3708}
3709
3710#[cfg(feature = "adapter-api")]
3713pub async fn run_http_authorized_with_lease_ttl(
3714 engine: Arc<RpcEngine>,
3715 bind: &str,
3716 credentials: Vec<RuntimeHttpCredential>,
3717 lease_ttl_ms: u64,
3718) -> std::io::Result<SocketAddr> {
3719 if credentials.is_empty()
3720 || credentials
3721 .iter()
3722 .any(|credential| credential.token.is_empty())
3723 {
3724 return Err(std::io::Error::new(
3725 std::io::ErrorKind::InvalidInput,
3726 "at least one non-empty runtime HTTP credential is required",
3727 ));
3728 }
3729 if lease_ttl_ms == 0 {
3730 return Err(std::io::Error::new(
3731 std::io::ErrorKind::InvalidInput,
3732 "runtime lease TTL must be non-zero",
3733 ));
3734 }
3735 let listener = TcpListener::bind(bind).await?;
3736 let local_addr = listener.local_addr()?;
3737 let credentials = RuntimeHttpCredentialRegistry::new(engine.session_id(), credentials)?;
3738 let eng = engine;
3739 let runtime: Arc<dyn FrontendRuntime> = eng.clone();
3740 let coordinator = CoordinatedRuntime::with_lease_ttl(runtime, lease_ttl_ms);
3741 tokio::spawn(async move {
3742 loop {
3743 tokio::select! {
3744 biased;
3745 _ = eng.wait_for_shutdown() => break,
3746 accepted = listener.accept() => {
3747 let Ok((stream, _addr)) = accepted else { continue };
3748 let eng = eng.clone();
3749 let coordinator = coordinator.clone();
3750 let credentials = credentials.clone();
3751 tokio::spawn(async move {
3752 let _ = handle_http_conn(stream, eng, coordinator, credentials).await;
3753 });
3754 }
3755 }
3756 }
3757 });
3758 Ok(local_addr)
3759}
3760
3761#[cfg(test)]
3762mod frontend_binding_conformance_tests;
3763
3764#[cfg(test)]
3765mod tests {
3766 use super::*;
3767 use tokio::io::BufReader;
3768
3769 fn cursor(data: &[u8]) -> BufReader<std::io::Cursor<Vec<u8>>> {
3770 BufReader::new(std::io::Cursor::new(data.to_vec()))
3771 }
3772
3773 #[cfg(all(feature = "adapter-api", supercode_workspace_assets))]
3774 #[test]
3775 fn packaged_observer_assets_match_the_sdk_sources() {
3776 let pairs: &[(&str, &[u8], &[u8])] = &[
3777 (
3778 "frontend-browser/index.html",
3779 include_bytes!("../embedded/frontend-browser/index.html"),
3780 include_bytes!("../../../sdk/frontend-browser/index.html"),
3781 ),
3782 (
3783 "frontend-browser/app.mjs",
3784 include_bytes!("../embedded/frontend-browser/app.mjs"),
3785 include_bytes!("../../../sdk/frontend-browser/app.mjs"),
3786 ),
3787 (
3788 "frontend-browser/client.mjs",
3789 include_bytes!("../embedded/frontend-browser/client.mjs"),
3790 include_bytes!("../../../sdk/frontend-browser/client.mjs"),
3791 ),
3792 (
3793 "frontend-browser/view.mjs",
3794 include_bytes!("../embedded/frontend-browser/view.mjs"),
3795 include_bytes!("../../../sdk/frontend-browser/view.mjs"),
3796 ),
3797 (
3798 "frontend-browser/style.css",
3799 include_bytes!("../embedded/frontend-browser/style.css"),
3800 include_bytes!("../../../sdk/frontend-browser/style.css"),
3801 ),
3802 (
3803 "frontend/client.mjs",
3804 include_bytes!("../embedded/frontend/client.mjs"),
3805 include_bytes!("../../../sdk/frontend/client.mjs"),
3806 ),
3807 (
3808 "frontend/generated-client.mjs",
3809 include_bytes!("../embedded/frontend/generated-client.mjs"),
3810 include_bytes!("../../../sdk/frontend/generated-client.mjs"),
3811 ),
3812 (
3813 "frontend/generated.mjs",
3814 include_bytes!("../embedded/frontend/generated.mjs"),
3815 include_bytes!("../../../sdk/frontend/generated.mjs"),
3816 ),
3817 ];
3818 for (name, packaged, source) in pairs {
3819 assert_eq!(packaged, source, "packaged observer asset drifted: {name}");
3820 }
3821 }
3822
3823 #[tokio::test]
3824 async fn admitted_submit_has_a_cancel_token_before_shutdown_observes_busy() {
3825 let agent =
3826 crate::Agent::new(crate::Config::builder().api_key("test-only-key").build()).unwrap();
3827 let engine = RpcEngine::new(agent, None);
3828 let claim = engine.claim_submit().unwrap();
3829 assert!(engine.busy.load(Ordering::SeqCst));
3830 assert!(engine
3831 .current_cancel
3832 .lock()
3833 .unwrap_or_else(std::sync::PoisonError::into_inner)
3834 .is_some());
3835
3836 let cancel = claim.cancel.clone();
3837 let shutdown_engine = engine.clone();
3838 let shutdown = tokio::spawn(async move { shutdown_engine.shutdown().await });
3839 tokio::time::timeout(std::time::Duration::from_secs(1), cancel.notified())
3840 .await
3841 .expect("shutdown must interrupt an admitted claim before its future starts");
3842 assert!(
3843 !shutdown.is_finished(),
3844 "shutdown must retain the barrier until the admitted claim drains"
3845 );
3846 drop(claim);
3847 tokio::time::timeout(std::time::Duration::from_secs(1), shutdown)
3848 .await
3849 .expect("claim drain must release shutdown")
3850 .unwrap();
3851 }
3852
3853 #[tokio::test]
3854 async fn read_bounded_line_reads_a_normal_line() {
3855 let mut r = cursor(b"hello\nworld\n");
3856 assert_eq!(
3857 read_bounded_line(&mut r, 1024).await.unwrap(),
3858 Some("hello".to_string())
3859 );
3860 assert_eq!(
3861 read_bounded_line(&mut r, 1024).await.unwrap(),
3862 Some("world".to_string())
3863 );
3864 assert_eq!(read_bounded_line(&mut r, 1024).await.unwrap(), None);
3865 }
3866
3867 #[tokio::test]
3868 async fn read_bounded_line_strips_trailing_cr() {
3869 let mut r = cursor(b"hello\r\n");
3870 assert_eq!(
3871 read_bounded_line(&mut r, 1024).await.unwrap(),
3872 Some("hello".to_string())
3873 );
3874 }
3875
3876 #[tokio::test]
3877 async fn read_bounded_line_returns_final_line_without_trailing_newline() {
3878 let mut r = cursor(b"no newline at eof");
3879 assert_eq!(
3880 read_bounded_line(&mut r, 1024).await.unwrap(),
3881 Some("no newline at eof".to_string())
3882 );
3883 assert_eq!(read_bounded_line(&mut r, 1024).await.unwrap(), None);
3884 }
3885
3886 #[tokio::test]
3887 async fn read_bounded_line_errors_and_resyncs_on_an_oversized_line() {
3888 let mut data = vec![b'x'; 20];
3889 data.push(b'\n');
3890 data.extend_from_slice(b"next\n");
3891 let mut r = cursor(&data);
3892 let err = read_bounded_line(&mut r, 10).await.unwrap_err();
3893 assert!(err.to_string().contains("10 byte cap"));
3894 assert_eq!(
3897 read_bounded_line(&mut r, 1024).await.unwrap(),
3898 Some("next".to_string())
3899 );
3900 }
3901
3902 #[test]
3903 fn constant_time_eq_matches_equal_slices() {
3904 assert!(constant_time_eq(b"abc123", b"abc123"));
3905 }
3906
3907 #[test]
3908 fn constant_time_eq_rejects_different_length_or_content() {
3909 assert!(!constant_time_eq(b"abc123", b"abc1234"));
3910 assert!(!constant_time_eq(b"abc123", b"xbc123"));
3911 }
3912
3913 #[test]
3914 fn generate_token_is_64_hex_chars_and_varies() {
3915 let a = generate_token();
3916 let b = generate_token();
3917 assert_eq!(a.len(), 64);
3918 assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
3919 assert_ne!(a, b, "two calls must not mint the same token");
3920 }
3921
3922 #[cfg(feature = "adapter-api")]
3923 #[tokio::test]
3924 async fn credential_revocation_waits_for_registered_attachment_ack() {
3925 let revocation = Arc::new(RuntimeCredentialRevocation::new());
3926 let attachment = revocation.register();
3927 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
3928 let task = tokio::spawn({
3929 let revocation = revocation.clone();
3930 async move {
3931 let _ = started_tx.send(());
3932 revocation.revoke_and_wait().await;
3933 }
3934 });
3935
3936 started_rx.await.unwrap();
3937 tokio::task::yield_now().await;
3938 assert!(
3939 !task.is_finished(),
3940 "revoke must remain pending while the attachment is registered"
3941 );
3942
3943 drop(attachment);
3944 tokio::time::timeout(std::time::Duration::from_secs(1), task)
3945 .await
3946 .expect("attachment acknowledgement must release revoke")
3947 .unwrap();
3948 }
3949
3950 #[cfg(feature = "adapter-api")]
3951 #[tokio::test]
3952 async fn credential_revocation_has_no_check_to_wait_lost_wakeup() {
3953 for _ in 0..10_000 {
3954 let revocation = Arc::new(RuntimeCredentialRevocation::new());
3955 let attachment = revocation.register();
3956 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
3957 let task = tokio::spawn({
3958 let revocation = revocation.clone();
3959 async move {
3960 let _ = started_tx.send(());
3961 revocation.revoke_and_wait().await;
3962 }
3963 });
3964
3965 started_rx.await.unwrap();
3966 drop(attachment);
3967 tokio::time::timeout(std::time::Duration::from_secs(1), task)
3968 .await
3969 .expect("revoke lost its attachment-drained wakeup")
3970 .unwrap();
3971 }
3972 }
3973
3974 #[test]
3975 fn request_history_compaction_deduplicates_and_orders_resolutions() {
3976 let request = |sequence, id| {
3977 FrontendEvent::new(
3978 sequence,
3979 json!({"type": "request", "request": {"id": id, "kind": "approval", "payload": {}}}),
3980 )
3981 };
3982 let resolved = |sequence, id| {
3983 FrontendEvent::new(
3984 sequence,
3985 json!({"type": "request_resolved", "request_id": id, "response": {"kind": "approval", "request_id": id, "decision": "allow"}}),
3986 )
3987 };
3988 let replay = VecDeque::from([
3989 request(1, 2),
3990 resolved(2, 2),
3991 request(3, 1),
3992 resolved(4, 1),
3993 request(5, 2),
3994 resolved(6, 1),
3995 ]);
3996
3997 let compacted = compact_frontend_request_history(&replay);
3998 assert_eq!(compacted.len(), 4);
3999 assert_eq!(compacted[0]["request"]["id"], 1);
4000 assert_eq!(compacted[1]["request_id"], 1);
4001 assert_eq!(compacted[2]["request"]["id"], 2);
4002 assert_eq!(compacted[3]["request_id"], 2);
4003 }
4004}