返回 CodeWhale
runtime_chat_relay.rs
根目录 / crates / tui / src / runtime_chat_relay.rs
1 //! Isolated native Runtime execution for account-owned Chat relay commands.
2 //!
3 //! The managed control plane supplies only opaque tenant/thread/turn bindings
4 //! and an exact non-secret provider route. Provider credentials and local
5 //! paths stay inside the Runtime. Every Chat thread is a dedicated
6 //! [`RuntimeThreadManager`] thread with an empty model-visible tool allowlist;
7 //! the active interactive TUI thread is never reused.
8
9 use std::{
10 collections::{BTreeMap, HashSet},
11 fs::{self, File},
12 path::{Path, PathBuf},
13 sync::Arc,
14 time::{Duration, Instant},
15 };
16
17 use anyhow::{Context, Result, bail};
18 use parking_lot::Mutex;
19 use serde::{Deserialize, Serialize};
20 use serde_json::{Value, json};
21 use sha2::{Digest, Sha256};
22
23 use crate::{
24 config::{Config, MemoryBackend, MemoryConfig, SkillsConfig},
25 core::engine::ISOLATED_CHAT_SYSTEM_PROMPT,
26 plugins::PluginRegistry,
27 runtime_threads::{
28 CreateThreadRequest, RuntimeEventRecord, RuntimeThreadManager, RuntimeThreadManagerConfig,
29 RuntimeTurnStatus, StartTurnRequest,
30 },
31 };
32
33 #[cfg(test)]
34 use crate::config::ContextConfig;
35
36 const STATE_SCHEMA_VERSION: u32 = 2;
37 const MAX_RELAY_ID_BYTES: usize = 240;
38 const MAX_OPERATION_KEY_BYTES: usize = 128;
39 const STATE_FILE: &str = "runtime-chat-bindings.json";
40 const SCOPE_LOCK_FILE: &str = "runtime-chat.owner.lock";
41
42 #[cfg(test)]
43 static TEST_STATE_PERSIST_FAILURES: std::sync::Mutex<Vec<(PathBuf, usize)>> =
44 std::sync::Mutex::new(Vec::new());
45
46 #[cfg(test)]
47 fn inject_state_persist_failures(path: &Path, count: usize) {
48 assert!(count > 0);
49 TEST_STATE_PERSIST_FAILURES
50 .lock()
51 .unwrap_or_else(std::sync::PoisonError::into_inner)
52 .push((path.to_path_buf(), count));
53 }
54
55 #[cfg(test)]
56 fn take_state_persist_failure(path: &Path) -> bool {
57 let mut failures = TEST_STATE_PERSIST_FAILURES
58 .lock()
59 .unwrap_or_else(std::sync::PoisonError::into_inner);
60 let Some(index) = failures.iter().position(|(target, _)| target == path) else {
61 return false;
62 };
63 if failures[index].1 > 1 {
64 failures[index].1 -= 1;
65 } else {
66 failures.remove(index);
67 }
68 true
69 }
70
71 #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)]
72 #[serde(rename_all = "camelCase", deny_unknown_fields)]
73 pub(crate) struct RuntimeChatPrompt {
74 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub max_output_tokens: Option<std::num::NonZeroU32>,
76 #[serde(rename = "type")]
77 pub command_type: String,
78 pub run_id: String,
79 pub turn_id: String,
80 pub operation_key: String,
81 pub runtime_binding_id: String,
82 pub runtime_thread_id: String,
83 pub prompt: String,
84 #[serde(default, skip_serializing_if = "Vec::is_empty")]
85 pub images: Vec<codewhale_protocol::runtime::RuntimeImageInput>,
86 #[serde(default, skip_serializing_if = "Option::is_none")]
87 pub system_prompt: Option<String>,
88 pub model: String,
89 pub model_provider: String,
90 pub model_provider_id: String,
91 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub reasoning_effort: Option<String>,
93 pub allowed_tools: Vec<String>,
94 pub mode: String,
95 pub requested_mode: String,
96 pub workspace: RuntimeChatWorkspace,
97 }
98
99 #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)]
100 #[serde(rename_all = "camelCase", deny_unknown_fields)]
101 pub(crate) struct RuntimeChatWorkspace {
102 pub id: String,
103 pub target_ref: String,
104 }
105
106 #[derive(Debug, Clone, PartialEq, Eq)]
107 pub(crate) struct RuntimeChatControlScope {
108 pub(crate) runtime_binding_id: String,
109 pub(crate) runtime_thread_id: String,
110 }
111
112 #[derive(Debug, Clone)]
113 pub(crate) struct RuntimeChatProjection {
114 pub run_id: String,
115 pub native_thread_id: String,
116 pub native_seq: u64,
117 pub source_event_id: String,
118 pub virtual_thread_id: String,
119 pub virtual_turn_id: String,
120 pub event: &'static str,
121 pub timestamp: String,
122 pub payload: Value,
123 }
124
125 #[derive(Clone)]
126 pub(crate) struct RuntimeChatRelayHost {
127 manager: Arc<RuntimeThreadManager>,
128 config: Arc<Config>,
129 state: Arc<Mutex<RelayState>>,
130 state_path: Arc<PathBuf>,
131 target_ref: Arc<String>,
132 session_id: Arc<String>,
133 _scope_lock: Arc<RelayScopeLock>,
134 apply_lock: Arc<tokio::sync::Mutex<()>>,
135 inference_ownership: Arc<Mutex<Option<crate::client::RuntimeChatInferenceOwnership>>>,
136 claimed_projections: Arc<Mutex<HashSet<(String, u64)>>>,
137 authorized_run_id: Arc<Mutex<Option<String>>>,
138 }
139
140 #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)]
141 #[serde(rename_all = "camelCase", deny_unknown_fields)]
142 struct RelayState {
143 schema_version: u32,
144 #[serde(default, skip_serializing_if = "Option::is_none")]
145 owner_scope_fingerprint: Option<String>,
146 #[serde(default)]
147 bindings: Vec<RelayThreadBinding>,
148 }
149
150 impl Default for RelayState {
151 fn default() -> Self {
152 Self {
153 schema_version: STATE_SCHEMA_VERSION,
154 owner_scope_fingerprint: None,
155 bindings: Vec::new(),
156 }
157 }
158 }
159
160 #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)]
161 #[serde(rename_all = "camelCase", deny_unknown_fields)]
162 struct RelayThreadBinding {
163 run_id: String,
164 runtime_binding_id: String,
165 virtual_thread_id: String,
166 native_thread_id: String,
167 model: String,
168 model_provider: String,
169 model_provider_id: String,
170 first_operation_fingerprint: String,
171 system_prompt_fingerprint: Option<String>,
172 #[serde(default)]
173 turns: BTreeMap<String, RelayTurnBinding>,
174 #[serde(default)]
175 projected_native_seq: u64,
176 }
177
178 #[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq)]
179 #[serde(rename_all = "camelCase", deny_unknown_fields)]
180 struct RelayTurnBinding {
181 native_turn_id: String,
182 operation_fingerprint: String,
183 request_fingerprint: String,
184 #[serde(default)]
185 terminal_projected: bool,
186 /// True only when the deterministic reservation was durably written but
187 /// native start was proven to have rejected before accepting provider work.
188 /// This distinguishes a retryable reservation from an already-projected
189 /// terminal turn, whose exact replay must remain settled.
190 #[serde(default)]
191 start_rejected: bool,
192 }
193
194 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
195 enum TurnReservationDisposition {
196 New,
197 Reopened,
198 ExistingUnsettled,
199 ExistingTerminal,
200 }
201
202 impl RelayState {
203 fn validate(&self) -> Result<()> {
204 if self.schema_version != STATE_SCHEMA_VERSION {
205 bail!("Runtime Chat binding state uses an unsupported schema");
206 }
207 if let Some(owner) = self.owner_scope_fingerprint.as_deref() {
208 validate_fingerprint(owner)?;
209 } else if !self.bindings.is_empty() {
210 bail!("Runtime Chat binding state has no account owner");
211 }
212 let mut binding_ids = HashSet::new();
213 let mut virtual_threads = HashSet::new();
214 let mut native_threads = HashSet::new();
215 for binding in &self.bindings {
216 validate_relay_id(&binding.run_id, "run id")?;
217 validate_relay_id(&binding.runtime_binding_id, "binding id")?;
218 validate_virtual_thread_id(&binding.virtual_thread_id)?;
219 validate_native_record_id(&binding.native_thread_id, "native thread id")?;
220 validate_route_id(&binding.model_provider, "provider id")?;
221 validate_route_id(&binding.model_provider_id, "model-provider id")?;
222 validate_model_id(&binding.model)?;
223 validate_fingerprint(&binding.first_operation_fingerprint)?;
224 if let Some(fingerprint) = binding.system_prompt_fingerprint.as_deref() {
225 validate_fingerprint(fingerprint)?;
226 }
227 if !binding_ids.insert(binding.runtime_binding_id.clone())
228 || !virtual_threads.insert(binding.virtual_thread_id.clone())
229 || !native_threads.insert(binding.native_thread_id.clone())
230 {
231 bail!("Runtime Chat binding state contains duplicate thread authority");
232 }
233 for (virtual_turn_id, turn) in &binding.turns {
234 validate_virtual_turn_id(virtual_turn_id)?;
235 validate_native_record_id(&turn.native_turn_id, "native turn id")?;
236 validate_fingerprint(&turn.operation_fingerprint)?;
237 validate_fingerprint(&turn.request_fingerprint)?;
238 }
239 }
240 Ok(())
241 }
242 }
243
244 impl RuntimeChatRelayHost {
245 pub(crate) fn open(
246 config: Config,
247 _plugin_registry: Arc<PluginRegistry>,
248 private_root: PathBuf,
249 target_ref: String,
250 session_id: String,
251 ) -> Result<Self, String> {
252 validate_owner_component(&target_ref, "target")?;
253 validate_owner_component(&session_id, "session")?;
254 let private_dir = scoped_private_dir(&private_root, &target_ref, &session_id);
255 fs::create_dir_all(&private_dir)
256 .map_err(|_| "Runtime Chat could not prepare private local state.".to_string())?;
257 #[cfg(unix)]
258 {
259 use std::os::unix::fs::PermissionsExt;
260 fs::set_permissions(&private_dir, fs::Permissions::from_mode(0o700))
261 .map_err(|_| "Runtime Chat could not protect private local state.".to_string())?;
262 }
263 let scope_lock =
264 RelayScopeLock::acquire(&private_dir.join(SCOPE_LOCK_FILE)).map_err(|error| {
265 // Only WouldBlock is genuine contention. Any other lock
266 // failure is a local IO fault, and misreporting it as
267 // ownership hides the cause (#5735's flake evidence).
268 let contention = error
269 .downcast_ref::<std::io::Error>()
270 .is_some_and(|io| io.kind() == std::io::ErrorKind::WouldBlock);
271 if contention {
272 "Another Codewhale process already owns this Runtime Chat account session."
273 .to_string()
274 } else {
275 format!("Runtime Chat could not take its owner lock: {error:#}")
276 }
277 })?;
278 let state_path = private_dir.join(STATE_FILE);
279 let state = load_state(&state_path).map_err(|_| {
280 "The saved Runtime Chat binding state could not be trusted.".to_string()
281 })?;
282 state.validate().map_err(|_| {
283 "The saved Runtime Chat binding state could not be trusted.".to_string()
284 })?;
285
286 // Reuse the native Runtime thread engine and durable records, but keep
287 // account Chat in its own private store and empty workspace. The
288 // dedicated system-prompt override below is the model-visible boundary;
289 // this workspace/config hardening also prevents local project, memory,
290 // instruction, and skill sources from becoming fallback context.
291 let (execution_config, chat_workspace) =
292 isolated_chat_execution_config(&config, &private_dir)?;
293 let relay_plugin_registry = Arc::new(PluginRegistry::empty(&chat_workspace));
294 let task_data_dir = private_dir.join("tasks");
295 let mut manager_cfg = RuntimeThreadManagerConfig::from_task_data_dir(task_data_dir);
296 manager_cfg.data_dir = private_dir.join("runtime");
297 let manager = RuntimeThreadManager::open_with_plugin_registry(
298 execution_config,
299 chat_workspace,
300 manager_cfg,
301 relay_plugin_registry,
302 )
303 .map_err(|_| "Runtime Chat could not open its isolated native thread store.".to_string())?;
304
305 Ok(Self {
306 manager: Arc::new(manager),
307 config: Arc::new(config),
308 state: Arc::new(Mutex::new(state)),
309 state_path: Arc::new(state_path),
310 target_ref: Arc::new(target_ref),
311 session_id: Arc::new(session_id),
312 _scope_lock: Arc::new(scope_lock),
313 apply_lock: Arc::new(tokio::sync::Mutex::new(())),
314 inference_ownership: Arc::new(Mutex::new(None)),
315 claimed_projections: Arc::new(Mutex::new(HashSet::new())),
316 authorized_run_id: Arc::new(Mutex::new(None)),
317 })
318 }
319
320 pub(crate) fn catalog(&self, challenge: &str) -> Result<Value, String> {
321 self.ensure_account_bound()?;
322 crate::runtime_api::runtime_chat_relay_catalog(&self.config, challenge)
323 }
324
325 pub(crate) fn catalog_payload_fingerprint(payload: &Value) -> Result<String, String> {
326 serde_json::to_vec(&canonical_json_value(payload))
327 .map(|bytes| hex_digest(Sha256::digest(bytes)))
328 .map_err(|_| "Runtime Chat could not fingerprint its safe catalog.".to_string())
329 }
330
331 pub(crate) fn bind_account(&self, account_ref: &str, target_ref: &str) -> Result<(), String> {
332 validate_owner_component(account_ref, "account")?;
333 if target_ref != self.target_ref.as_str() {
334 return Err("The Runtime Chat account owner does not match this target.".to_string());
335 }
336 let owner = owner_scope_fingerprint(account_ref, target_ref, &self.session_id);
337 self.persist_state_update(
338 "Runtime Chat could not persist its account ownership.",
339 |state| bind_owner_scope(state, &owner),
340 )
341 }
342
343 pub(crate) fn authorize_run(&self, run_id: &str) -> Result<(), String> {
344 validate_relay_id(run_id, "run id")
345 .map_err(|_| "The Runtime Chat attachment has an invalid run identity.".to_string())?;
346 self.ensure_account_bound()?;
347 let durable_other_run_is_unsettled = self.state.lock().bindings.iter().any(|binding| {
348 binding.run_id != run_id && binding.turns.values().any(|turn| !turn.terminal_projected)
349 });
350 if durable_other_run_is_unsettled {
351 return Err(
352 "Finish or interrupt the active Runtime Chat turn before attaching another run."
353 .to_string(),
354 );
355 }
356 if self.has_any_unsettled_turns() && self.inference_ownership.lock().is_none() {
357 let ownership = crate::client::try_acquire_runtime_chat_inference_ownership()
358 .ok_or_else(|| {
359 "Finish the active local turn before recovering Runtime Chat.".to_string()
360 })?;
361 *self.inference_ownership.lock() = Some(ownership);
362 }
363 *self.authorized_run_id.lock() = Some(run_id.to_string());
364 Ok(())
365 }
366
367 #[cfg(test)]
368 pub(crate) fn has_unsettled_authorized_turns(&self) -> bool {
369 self.authorized_run_id
370 .lock()
371 .as_deref()
372 .is_some_and(|run_id| self.has_unsettled_turns_for_run(run_id))
373 }
374
375 pub(crate) fn has_any_unsettled_turns(&self) -> bool {
376 self.state
377 .lock()
378 .bindings
379 .iter()
380 .any(|binding| binding.turns.values().any(|turn| !turn.terminal_projected))
381 }
382
383 #[cfg(test)]
384 async fn ensure_inference_ownership(&self) {
385 if self.inference_ownership.lock().is_some() {
386 return;
387 }
388 let ownership = crate::client::acquire_runtime_chat_inference_ownership().await;
389 let mut current = self.inference_ownership.lock();
390 if current.is_none() {
391 *current = Some(ownership);
392 }
393 }
394
395 fn try_ensure_inference_ownership(&self) -> Result<(), String> {
396 if self.inference_ownership.lock().is_some() {
397 return Ok(());
398 }
399 let ownership =
400 crate::client::try_acquire_runtime_chat_inference_ownership().ok_or_else(|| {
401 "Finish the active local provider work before starting Runtime Chat.".to_string()
402 })?;
403 let mut current = self.inference_ownership.lock();
404 if current.is_none() {
405 *current = Some(ownership);
406 }
407 Ok(())
408 }
409
410 pub(crate) fn recover_inference_ownership_for_pending_delivery(&self) -> Result<(), String> {
411 if self.inference_ownership.lock().is_some() {
412 return Ok(());
413 }
414 let ownership =
415 crate::client::try_acquire_runtime_chat_inference_ownership().ok_or_else(|| {
416 "Finish the active local turn before recovering Runtime Chat delivery.".to_string()
417 })?;
418 *self.inference_ownership.lock() = Some(ownership);
419 Ok(())
420 }
421
422 pub(crate) fn release_inference_ownership_if_settled(&self) {
423 if !self.has_any_unsettled_turns() {
424 self.inference_ownership.lock().take();
425 }
426 }
427
428 pub(crate) fn is_exact_prompt_replay(
429 &self,
430 command: &RuntimeChatPrompt,
431 ) -> Result<bool, String> {
432 command.validate_shape()?;
433 let operation_fingerprint = fingerprint(&command.operation_key);
434 let request_fingerprint = runtime_chat_request_fingerprint(command)?;
435 let state = self.state.lock();
436 let by_binding = state
437 .bindings
438 .iter()
439 .position(|binding| binding.runtime_binding_id == command.runtime_binding_id);
440 let by_thread = state
441 .bindings
442 .iter()
443 .position(|binding| binding.virtual_thread_id == command.runtime_thread_id);
444 let binding = match (by_binding, by_thread) {
445 (None, None) => return Ok(false),
446 (Some(binding_index), Some(thread_index)) if binding_index == thread_index => {
447 &state.bindings[binding_index]
448 }
449 _ => {
450 return Err(
451 "The Runtime Chat thread authority does not match its binding.".to_string(),
452 );
453 }
454 };
455 if binding.run_id != command.run_id {
456 return Err("The Runtime Chat replay belongs to another run.".to_string());
457 }
458 let Some(turn) = binding.turns.get(&command.turn_id) else {
459 return Ok(false);
460 };
461 if turn.operation_fingerprint == operation_fingerprint
462 && turn.request_fingerprint == request_fingerprint
463 {
464 return Ok(true);
465 }
466 Err("The Runtime Chat turn binding does not match its replay.".to_string())
467 }
468
469 pub(crate) fn scope_matches(&self, target_ref: &str, session_id: &str) -> bool {
470 self.target_ref.as_str() == target_ref && self.session_id.as_str() == session_id
471 }
472
473 #[cfg(test)]
474 pub(crate) fn configured_default_model_for_tests(&self) -> String {
475 self.config.default_model()
476 }
477
478 #[cfg(test)]
479 pub(crate) fn inference_ownership_is_held_for_tests(&self) -> bool {
480 self.inference_ownership.lock().is_some()
481 }
482
483 #[cfg(test)]
484 pub(crate) async fn acquire_inference_ownership_for_tests(&self) {
485 self.ensure_inference_ownership().await;
486 }
487
488 #[cfg(test)]
489 pub(crate) fn install_unsettled_turn_for_tests(
490 &self,
491 run_id: &str,
492 native_thread_id: &str,
493 virtual_thread_id: &str,
494 virtual_turn_id: &str,
495 ) -> Result<(), String> {
496 self.insert_binding(RelayThreadBinding {
497 run_id: run_id.to_string(),
498 runtime_binding_id: "binding_test_fixture".to_string(),
499 virtual_thread_id: virtual_thread_id.to_string(),
500 native_thread_id: native_thread_id.to_string(),
501 model: "model-1".to_string(),
502 model_provider: "custom".to_string(),
503 model_provider_id: "local-provider".to_string(),
504 first_operation_fingerprint: fingerprint("operation-test-fixture"),
505 system_prompt_fingerprint: None,
506 turns: BTreeMap::from([(
507 virtual_turn_id.to_string(),
508 RelayTurnBinding {
509 native_turn_id: "turn_test_fixture".to_string(),
510 operation_fingerprint: fingerprint("operation-test-fixture"),
511 request_fingerprint: fingerprint("request-test-fixture"),
512 terminal_projected: false,
513 start_rejected: false,
514 },
515 )]),
516 projected_native_seq: 0,
517 })
518 }
519
520 #[cfg(test)]
521 pub(crate) async fn install_prompt_replay_for_tests(
522 &self,
523 command: &RuntimeChatPrompt,
524 terminal_projected: bool,
525 ) -> Result<(), String> {
526 command.validate_shape()?;
527 let operation_fingerprint = fingerprint(&command.operation_key);
528 let request_fingerprint = runtime_chat_request_fingerprint(command)?;
529 let native_thread_id = self
530 .manager
531 .create_thread(CreateThreadRequest {
532 model: Some(command.model.clone()),
533 model_provider: Some(command.model_provider.clone()),
534 model_provider_id: Some(command.model_provider_id.clone()),
535 reasoning_effort: command.reasoning_effort.clone(),
536 allowed_tools: Some(Vec::new()),
537 workspace: None,
538 mode: Some("agent".to_string()),
539 permission_posture: Some("ask".to_string()),
540 allow_shell: Some(false),
541 trust_mode: Some(false),
542 auto_approve: Some(false),
543 archived: false,
544 system_prompt: Some(dedicated_chat_system_prompt(
545 command.system_prompt.as_deref(),
546 )),
547 task_id: None,
548 dynamic_tools: Vec::new(),
549 environments: Vec::new(),
550 })
551 .await
552 .map_err(|_| "Runtime Chat could not create its replay fixture.".to_string())?
553 .id;
554 let native_turn_id = reserved_native_turn_id(
555 &native_thread_id,
556 &command.runtime_binding_id,
557 &command.turn_id,
558 &operation_fingerprint,
559 );
560 self.insert_binding(RelayThreadBinding {
561 run_id: command.run_id.clone(),
562 runtime_binding_id: command.runtime_binding_id.clone(),
563 virtual_thread_id: command.runtime_thread_id.clone(),
564 native_thread_id,
565 model: command.model.clone(),
566 model_provider: command.model_provider.clone(),
567 model_provider_id: command.model_provider_id.clone(),
568 first_operation_fingerprint: operation_fingerprint.clone(),
569 system_prompt_fingerprint: command.system_prompt.as_deref().map(fingerprint),
570 turns: BTreeMap::from([(
571 command.turn_id.clone(),
572 RelayTurnBinding {
573 native_turn_id,
574 operation_fingerprint,
575 request_fingerprint,
576 terminal_projected,
577 start_rejected: false,
578 },
579 )]),
580 projected_native_seq: if terminal_projected { 1 } else { 0 },
581 })
582 }
583
584 fn has_unsettled_turns_for_run(&self, run_id: &str) -> bool {
585 self.state.lock().bindings.iter().any(|binding| {
586 binding.run_id == run_id && binding.turns.values().any(|turn| !turn.terminal_projected)
587 })
588 }
589
590 pub(crate) async fn apply_prompt(&self, command: &RuntimeChatPrompt) -> Result<(), String> {
591 self.ensure_account_bound()?;
592 if self.authorized_run_id.lock().as_deref() != Some(command.run_id.as_str()) {
593 return Err("The Runtime Chat command is not authorized for this run.".to_string());
594 }
595 let _apply = self.apply_lock.lock().await;
596 command.validate_shape()?;
597 let operation_fingerprint = fingerprint(&command.operation_key);
598 let request_fingerprint = runtime_chat_request_fingerprint(command)?;
599 let system_prompt_fingerprint = command.system_prompt.as_deref().map(fingerprint);
600
601 let existing = {
602 let state = self.state.lock();
603 let by_binding = state
604 .bindings
605 .iter()
606 .position(|binding| binding.runtime_binding_id == command.runtime_binding_id);
607 let by_thread = state
608 .bindings
609 .iter()
610 .position(|binding| binding.virtual_thread_id == command.runtime_thread_id);
611 match (by_binding, by_thread) {
612 (None, None) => None,
613 (Some(binding_index), Some(thread_index)) if binding_index == thread_index => {
614 Some(state.bindings[binding_index].clone())
615 }
616 _ => {
617 return Err(
618 "The Runtime Chat thread authority does not match its binding.".to_string(),
619 );
620 }
621 }
622 };
623 let exact_operation_replay = existing.as_ref().is_some_and(|binding| {
624 binding.turns.get(&command.turn_id).is_some_and(|turn| {
625 turn.operation_fingerprint == operation_fingerprint
626 && turn.request_fingerprint == request_fingerprint
627 })
628 });
629 if !exact_operation_replay {
630 self.validate_route(command)?;
631 }
632 if !exact_operation_replay && self.has_unsettled_turns_for_run(&command.run_id) {
633 return Err(
634 "Finish or interrupt the active Runtime Chat turn before starting another."
635 .to_string(),
636 );
637 }
638
639 let binding = if let Some(binding) = existing {
640 validate_existing_binding(
641 &binding,
642 command,
643 &operation_fingerprint,
644 system_prompt_fingerprint.as_deref(),
645 )?;
646 self.manager
647 .get_thread(&binding.native_thread_id)
648 .await
649 .map_err(|_| "The isolated Runtime Chat thread is unavailable.".to_string())?;
650 binding
651 } else {
652 let thread = self
653 .manager
654 .create_thread(CreateThreadRequest {
655 model: Some(command.model.clone()),
656 model_provider: Some(command.model_provider.clone()),
657 model_provider_id: Some(command.model_provider_id.clone()),
658 reasoning_effort: command.reasoning_effort.clone(),
659 allowed_tools: Some(Vec::new()),
660 workspace: None,
661 // The native engine's existing Act loop executes a pure
662 // chat turn once its model-visible tool catalog is empty.
663 // `chat` remains the account wire mode, not a second loop.
664 mode: Some("agent".to_string()),
665 permission_posture: Some("ask".to_string()),
666 allow_shell: Some(false),
667 trust_mode: Some(false),
668 auto_approve: Some(false),
669 archived: false,
670 system_prompt: Some(dedicated_chat_system_prompt(
671 command.system_prompt.as_deref(),
672 )),
673 task_id: None,
674 dynamic_tools: Vec::new(),
675 environments: Vec::new(),
676 })
677 .await
678 .map_err(|_| {
679 "Runtime Chat could not create an isolated native thread.".to_string()
680 })?;
681 let binding = RelayThreadBinding {
682 runtime_binding_id: command.runtime_binding_id.clone(),
683 run_id: command.run_id.clone(),
684 virtual_thread_id: command.runtime_thread_id.clone(),
685 native_thread_id: thread.id,
686 model: command.model.clone(),
687 model_provider: command.model_provider.clone(),
688 model_provider_id: command.model_provider_id.clone(),
689 first_operation_fingerprint: operation_fingerprint.clone(),
690 system_prompt_fingerprint,
691 turns: BTreeMap::new(),
692 projected_native_seq: 0,
693 };
694 if let Err(error) = self.insert_binding(binding.clone()) {
695 let _ = self
696 .manager
697 .discard_empty_thread(&binding.native_thread_id)
698 .await;
699 return Err(error);
700 }
701 binding
702 };
703
704 // Reject virtual-turn and operation-key drift before asking the native
705 // manager to start anything. RuntimeThreadManager independently
706 // enforces the same operation-key idempotency for crash replay.
707 if let Some(existing_turn) = binding.turns.get(&command.turn_id)
708 && (existing_turn.operation_fingerprint != operation_fingerprint
709 || existing_turn.request_fingerprint != request_fingerprint)
710 {
711 return Err("The Runtime Chat turn binding does not match its replay.".to_string());
712 }
713 if binding.turns.iter().any(|(virtual_turn_id, turn)| {
714 virtual_turn_id != &command.turn_id
715 && turn.operation_fingerprint == operation_fingerprint
716 }) {
717 return Err("The Runtime Chat operation key belongs to another turn.".to_string());
718 }
719
720 // Preallocate and persist the exact native id before the manager may
721 // submit anything to a provider. RuntimeThreadManager binds the same id
722 // inside its operation-key transaction, eliminating the crash window
723 // between native acceptance and relay correlation.
724 let reserved_native_turn_id = reserved_native_turn_id(
725 &binding.native_thread_id,
726 &command.runtime_binding_id,
727 &command.turn_id,
728 &operation_fingerprint,
729 );
730 let reservation = self.reserve_turn(
731 &command.runtime_binding_id,
732 &command.runtime_thread_id,
733 &command.turn_id,
734 &reserved_native_turn_id,
735 &operation_fingerprint,
736 &request_fingerprint,
737 )?;
738 if matches!(reservation, TurnReservationDisposition::ExistingTerminal) {
739 return Ok(());
740 }
741
742 // The exclusive provider-request lease is acquired only after all
743 // command/route/binding validation and deterministic reservation are
744 // durable, but before native Runtime can resolve Auto or dispatch any
745 // provider request. At the common client seam it also blocks advisor,
746 // detached-subagent, compaction, purge, and interactive requests that
747 // could otherwise feed this attached CWC run.
748 if let Err(error) = self.try_ensure_inference_ownership() {
749 if !matches!(reservation, TurnReservationDisposition::ExistingUnsettled) {
750 self.finish_turn_reservation(
751 &command.runtime_binding_id,
752 &command.runtime_thread_id,
753 &command.turn_id,
754 )?;
755 }
756 self.release_inference_ownership_if_settled();
757 return Err(error);
758 }
759
760 let turn = match self
761 .manager
762 .start_turn_with_reserved_id(
763 &binding.native_thread_id,
764 StartTurnRequest {
765 expected_workspace: None,
766 max_output_tokens: command.max_output_tokens,
767 prompt: command.prompt.clone(),
768 images: command.images.clone(),
769 operation_key: Some(command.operation_key.clone()),
770 input_summary: None,
771 model: Some(command.model.clone()),
772 reasoning_effort: command.reasoning_effort.clone(),
773 allowed_tools: Some(Vec::new()),
774 mode: Some("agent".to_string()),
775 permission_posture: Some("ask".to_string()),
776 allow_shell: Some(false),
777 trust_mode: Some(false),
778 auto_approve: Some(false),
779 dynamic_tools: Vec::new(),
780 environment_id: None,
781 model_provider: None,
782 model_provider_id: None,
783 },
784 &reserved_native_turn_id,
785 )
786 .await
787 {
788 Ok(turn) => turn,
789 Err(error) => {
790 // A process may crash after the relay reservation is durable
791 // but before Runtime persists its operation binding/turn. On
792 // exact replay that appears as ExistingUnsettled. Settle it
793 // only when the native thread snapshot positively proves the
794 // deterministic turn id was never accepted; any unreadable or
795 // present native record stays fail-closed and projectable.
796 let native_turn_is_durable = self
797 .manager
798 .get_thread_detail(&binding.native_thread_id)
799 .await
800 .map(|detail| {
801 detail
802 .turns
803 .iter()
804 .any(|turn| turn.id == reserved_native_turn_id)
805 })
806 .unwrap_or(true);
807 if !native_turn_is_durable {
808 self.finish_turn_reservation(
809 &command.runtime_binding_id,
810 &command.runtime_thread_id,
811 &command.turn_id,
812 )?;
813 }
814 self.release_inference_ownership_if_settled();
815 return Err(sanitized_runtime_error("start", &error));
816 }
817 };
818 debug_assert_eq!(turn.id, reserved_native_turn_id);
819 Ok(())
820 }
821
822 pub(crate) async fn interrupt(
823 &self,
824 run_id: &str,
825 scope: &RuntimeChatControlScope,
826 virtual_turn_id: &str,
827 ) -> Result<(), String> {
828 self.ensure_account_bound()?;
829 validate_relay_id(run_id, "run id")
830 .map_err(|_| "The Runtime Chat interrupt run is invalid.".to_string())?;
831 validate_relay_id(&scope.runtime_binding_id, "binding id")
832 .map_err(|_| "The Runtime Chat interrupt binding is invalid.".to_string())?;
833 validate_virtual_thread_id(&scope.runtime_thread_id)
834 .map_err(|_| "The Runtime Chat interrupt thread is invalid.".to_string())?;
835 validate_virtual_turn_id(virtual_turn_id)
836 .map_err(|_| "The Runtime Chat interrupt turn is invalid.".to_string())?;
837 let (native_thread_id, native_turn_id) = {
838 let state = self.state.lock();
839 resolve_interrupt_target(&state, run_id, scope, virtual_turn_id)?
840 };
841
842 let detail = self
843 .manager
844 .get_thread_detail(&native_thread_id)
845 .await
846 .map_err(|_| "The isolated Runtime Chat thread is unavailable.".to_string())?;
847 let turn = detail
848 .turns
849 .iter()
850 .find(|turn| turn.id == native_turn_id)
851 .ok_or_else(|| "The isolated Runtime Chat turn is unavailable.".to_string())?;
852 if !matches!(
853 turn.status,
854 RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress
855 ) {
856 // Exact replay after the terminal boundary is already applied.
857 return Ok(());
858 }
859 self.manager
860 .interrupt_turn(&native_thread_id, &native_turn_id)
861 .await
862 .map(|_| ())
863 .map_err(|error| sanitized_runtime_error("interrupt", &error))
864 }
865
866 /// Return at most one not-yet-journaled projection per poll. The caller
867 /// journals it before advancing `projected_native_seq`, so a crash can
868 /// cause a safe replay but never silently lose an accepted native event.
869 pub(crate) async fn pending_projections(&self) -> Result<Vec<RuntimeChatProjection>, String> {
870 let Some(authorized_run_id) = self.authorized_run_id.lock().clone() else {
871 return Ok(Vec::new());
872 };
873 let bindings = {
874 let state = self.state.lock();
875 if state.owner_scope_fingerprint.is_none() {
876 return Ok(Vec::new());
877 }
878 state
879 .bindings
880 .iter()
881 .filter(|binding| binding.run_id == authorized_run_id)
882 .cloned()
883 .collect::<Vec<_>>()
884 };
885 for binding in bindings {
886 let events = self
887 .manager
888 .events_since_async(
889 &binding.native_thread_id,
890 Some(binding.projected_native_seq),
891 )
892 .await
893 .map_err(|_| "Runtime Chat could not read its native event ledger.".to_string())?;
894 for event in events {
895 let key = (binding.native_thread_id.clone(), event.seq);
896 if self.claimed_projections.lock().contains(&key) {
897 break;
898 }
899 match project_native_event(&binding, &event) {
900 ProjectionDecision::Project(projection) => {
901 self.claimed_projections.lock().insert(key);
902 // Return immediately after acquiring the claim. No
903 // later binding read can fail and accidentally drop a
904 // previously acquired in-process claim.
905 return Ok(vec![*projection]);
906 }
907 ProjectionDecision::Skip => {
908 self.advance_projection_cursor(&binding.native_thread_id, event.seq, None)?;
909 }
910 ProjectionDecision::WaitForTurnBinding => break,
911 }
912 }
913 }
914 Ok(Vec::new())
915 }
916
917 pub(crate) fn mark_projected(
918 &self,
919 native_thread_id: &str,
920 native_seq: u64,
921 virtual_turn_id: &str,
922 event: &str,
923 ) -> Result<(), String> {
924 self.claimed_projections
925 .lock()
926 .remove(&(native_thread_id.to_string(), native_seq));
927 self.advance_projection_cursor(
928 native_thread_id,
929 native_seq,
930 (event == "turn.completed").then_some(virtual_turn_id),
931 )
932 }
933
934 pub(crate) fn release_projection(&self, native_thread_id: &str, native_seq: u64) {
935 self.claimed_projections
936 .lock()
937 .remove(&(native_thread_id.to_string(), native_seq));
938 }
939
940 pub(crate) fn release_all_projection_claims(&self) {
941 self.claimed_projections.lock().clear();
942 }
943
944 #[cfg(test)]
945 pub(crate) fn install_projection_claim_for_tests(
946 &self,
947 native_thread_id: &str,
948 native_seq: u64,
949 ) {
950 self.claimed_projections
951 .lock()
952 .insert((native_thread_id.to_string(), native_seq));
953 }
954
955 #[cfg(test)]
956 pub(crate) fn projection_is_claimed_for_tests(
957 &self,
958 native_thread_id: &str,
959 native_seq: u64,
960 ) -> bool {
961 self.claimed_projections
962 .lock()
963 .contains(&(native_thread_id.to_string(), native_seq))
964 }
965
966 #[cfg(test)]
967 async fn native_turn_count_for_binding_for_tests(&self, runtime_binding_id: &str) -> usize {
968 let native_thread_id = self
969 .state
970 .lock()
971 .bindings
972 .iter()
973 .find(|binding| binding.runtime_binding_id == runtime_binding_id)
974 .map(|binding| binding.native_thread_id.clone())
975 .expect("test binding exists");
976 self.manager
977 .get_thread_detail(&native_thread_id)
978 .await
979 .expect("test native thread detail")
980 .turns
981 .len()
982 }
983
984 fn validate_route(&self, command: &RuntimeChatPrompt) -> Result<(), String> {
985 let catalog = self.catalog("a2345678901234567890123456789012")?;
986 let provider = catalog
987 .get("providers")
988 .and_then(Value::as_array)
989 .and_then(|providers| providers.first())
990 .ok_or_else(|| "The active Runtime Chat route is unavailable.".to_string())?;
991 let route_matches = provider.get("id").and_then(Value::as_str)
992 == Some(command.model_provider.as_str())
993 && provider.get("modelProviderId").and_then(Value::as_str)
994 == Some(command.model_provider_id.as_str())
995 && provider
996 .get("models")
997 .and_then(Value::as_array)
998 .is_some_and(|models| {
999 models.iter().any(|model| {
1000 model.get("id").and_then(Value::as_str) == Some(command.model.as_str())
1001 && (command.images.is_empty()
1002 || model.get("imageInput").and_then(Value::as_str)
1003 == Some("supported"))
1004 })
1005 });
1006 if command.max_output_tokens.is_some()
1007 && !provider
1008 .get("models")
1009 .and_then(Value::as_array)
1010 .is_some_and(|models| {
1011 models.iter().any(|model| {
1012 model.get("id").and_then(Value::as_str) == Some(command.model.as_str())
1013 && model.get("outputTokenLimit").and_then(Value::as_str)
1014 == Some("supported")
1015 })
1016 })
1017 {
1018 return Err(
1019 "The selected Runtime Chat route does not support maxOutputTokens.".to_string(),
1020 );
1021 }
1022 if !route_matches {
1023 return Err(
1024 "The requested Runtime Chat route is not the active ready route.".to_string(),
1025 );
1026 }
1027 Ok(())
1028 }
1029
1030 fn insert_binding(&self, binding: RelayThreadBinding) -> Result<(), String> {
1031 self.persist_state_update(
1032 "Runtime Chat could not persist its private thread binding.",
1033 |state| {
1034 if state.owner_scope_fingerprint.is_none() {
1035 return Err("The Runtime Chat account owner is not established.".to_string());
1036 }
1037 if state.bindings.iter().any(|existing| {
1038 existing.runtime_binding_id == binding.runtime_binding_id
1039 || existing.virtual_thread_id == binding.virtual_thread_id
1040 }) {
1041 return Err(
1042 "The Runtime Chat thread binding changed while it was being created."
1043 .to_string(),
1044 );
1045 }
1046 state.bindings.push(binding);
1047 Ok(())
1048 },
1049 )
1050 }
1051
1052 fn reserve_turn(
1053 &self,
1054 runtime_binding_id: &str,
1055 virtual_thread_id: &str,
1056 virtual_turn_id: &str,
1057 native_turn_id: &str,
1058 operation_fingerprint: &str,
1059 request_fingerprint: &str,
1060 ) -> Result<TurnReservationDisposition, String> {
1061 validate_native_record_id(native_turn_id, "native turn id")
1062 .map_err(public_validation_error)?;
1063 validate_fingerprint(request_fingerprint).map_err(public_validation_error)?;
1064 self.persist_state_update(
1065 "Runtime Chat could not persist its private turn reservation.",
1066 |state| {
1067 let binding = state
1068 .bindings
1069 .iter_mut()
1070 .find(|binding| {
1071 binding.runtime_binding_id == runtime_binding_id
1072 && binding.virtual_thread_id == virtual_thread_id
1073 })
1074 .ok_or_else(|| "The Runtime Chat thread binding disappeared.".to_string())?;
1075 if let Some(existing) = binding.turns.get_mut(virtual_turn_id) {
1076 if existing.operation_fingerprint == operation_fingerprint
1077 && existing.request_fingerprint == request_fingerprint
1078 && existing.native_turn_id == native_turn_id
1079 {
1080 // A retry after a proven pre-dispatch failure reopens
1081 // the same deterministic reservation. Any accepted or
1082 // replayed native turn is therefore unsettled again
1083 // until its terminal event is durably projected.
1084 let disposition = if existing.terminal_projected && existing.start_rejected
1085 {
1086 TurnReservationDisposition::Reopened
1087 } else if existing.terminal_projected {
1088 TurnReservationDisposition::ExistingTerminal
1089 } else {
1090 TurnReservationDisposition::ExistingUnsettled
1091 };
1092 if matches!(disposition, TurnReservationDisposition::Reopened) {
1093 existing.terminal_projected = false;
1094 existing.start_rejected = false;
1095 }
1096 return Ok(disposition);
1097 }
1098 return Err(
1099 "The Runtime Chat turn binding does not match its replay.".to_string()
1100 );
1101 }
1102 if binding
1103 .turns
1104 .values()
1105 .any(|turn| turn.operation_fingerprint == operation_fingerprint)
1106 {
1107 return Err(
1108 "The Runtime Chat operation key belongs to another turn.".to_string()
1109 );
1110 }
1111 binding.turns.insert(
1112 virtual_turn_id.to_string(),
1113 RelayTurnBinding {
1114 native_turn_id: native_turn_id.to_string(),
1115 operation_fingerprint: operation_fingerprint.to_string(),
1116 request_fingerprint: request_fingerprint.to_string(),
1117 terminal_projected: false,
1118 start_rejected: false,
1119 },
1120 );
1121 Ok(TurnReservationDisposition::New)
1122 },
1123 )
1124 }
1125
1126 fn finish_turn_reservation(
1127 &self,
1128 runtime_binding_id: &str,
1129 virtual_thread_id: &str,
1130 virtual_turn_id: &str,
1131 ) -> Result<(), String> {
1132 self.persist_state_update(
1133 "Runtime Chat could not settle its rejected turn reservation.",
1134 |state| {
1135 let turn = state
1136 .bindings
1137 .iter_mut()
1138 .find(|binding| {
1139 binding.runtime_binding_id == runtime_binding_id
1140 && binding.virtual_thread_id == virtual_thread_id
1141 })
1142 .and_then(|binding| binding.turns.get_mut(virtual_turn_id))
1143 .ok_or_else(|| "The Runtime Chat turn reservation disappeared.".to_string())?;
1144 turn.terminal_projected = true;
1145 turn.start_rejected = true;
1146 Ok(())
1147 },
1148 )
1149 }
1150
1151 fn advance_projection_cursor(
1152 &self,
1153 native_thread_id: &str,
1154 native_seq: u64,
1155 terminal_virtual_turn_id: Option<&str>,
1156 ) -> Result<(), String> {
1157 self.persist_state_update(
1158 "Runtime Chat could not persist its event cursor.",
1159 |state| {
1160 let binding = state
1161 .bindings
1162 .iter_mut()
1163 .find(|binding| binding.native_thread_id == native_thread_id)
1164 .ok_or_else(|| "The Runtime Chat projection binding is unknown.".to_string())?;
1165 binding.projected_native_seq = binding.projected_native_seq.max(native_seq);
1166 if let Some(virtual_turn_id) = terminal_virtual_turn_id {
1167 let turn = binding.turns.get_mut(virtual_turn_id).ok_or_else(|| {
1168 "The Runtime Chat terminal projection has no turn binding.".to_string()
1169 })?;
1170 turn.terminal_projected = true;
1171 turn.start_rejected = false;
1172 }
1173 Ok(())
1174 },
1175 )
1176 }
1177
1178 fn persist_state_update<T>(
1179 &self,
1180 persistence_error: &'static str,
1181 update: impl FnOnce(&mut RelayState) -> Result<T, String>,
1182 ) -> Result<T, String> {
1183 let mut current = self.state.lock();
1184 let mut candidate = current.clone();
1185 let output = update(&mut candidate)?;
1186 persist_state(&self.state_path, &candidate).map_err(|_| persistence_error.to_string())?;
1187 *current = candidate;
1188 Ok(output)
1189 }
1190
1191 fn ensure_account_bound(&self) -> Result<(), String> {
1192 if self.state.lock().owner_scope_fingerprint.is_none() {
1193 return Err("The Runtime Chat account owner is not established.".to_string());
1194 }
1195 Ok(())
1196 }
1197 }
1198
1199 impl RuntimeChatPrompt {
1200 pub(crate) fn validate_shape(&self) -> Result<(), String> {
1201 crate::image_attach::prepare_runtime_images(&self.images)
1202 .map_err(|error| error.to_string())?;
1203 if !self.images.is_empty() && self.model.trim().eq_ignore_ascii_case("auto") {
1204 return Err(
1205 "Image inputs require an exact named model; Auto is unavailable for images."
1206 .to_string(),
1207 );
1208 }
1209 if self.command_type != "prompt.request" {
1210 return Err("Codewhale sent an unsupported Runtime Chat command.".to_string());
1211 }
1212 validate_relay_id(&self.run_id, "run id").map_err(public_validation_error)?;
1213 validate_virtual_turn_id(&self.turn_id).map_err(public_validation_error)?;
1214 validate_operation_key(&self.operation_key).map_err(public_validation_error)?;
1215 validate_relay_id(&self.runtime_binding_id, "binding id")
1216 .map_err(public_validation_error)?;
1217 validate_virtual_thread_id(&self.runtime_thread_id).map_err(public_validation_error)?;
1218 validate_model_id(&self.model).map_err(public_validation_error)?;
1219 validate_route_id(&self.model_provider, "provider id").map_err(public_validation_error)?;
1220 validate_route_id(&self.model_provider_id, "model-provider id")
1221 .map_err(public_validation_error)?;
1222 if self.prompt.trim().is_empty()
1223 || self.prompt.len() > 128 * 1024
1224 || self.prompt.contains('\0')
1225 {
1226 return Err("The Runtime Chat prompt is empty or oversized.".to_string());
1227 }
1228 if let Some(system_prompt) = self.system_prompt.as_deref()
1229 && (system_prompt.trim().is_empty()
1230 || system_prompt.len() > 64_000
1231 || system_prompt.contains('\0'))
1232 {
1233 return Err("The Runtime Chat system instructions are invalid.".to_string());
1234 }
1235 if !self.allowed_tools.is_empty() {
1236 return Err("Runtime-backed Chat does not grant work-execution tools.".to_string());
1237 }
1238 if self.mode != "chat" || self.requested_mode != "chat" {
1239 return Err("Runtime relay turns must use Chat mode.".to_string());
1240 }
1241 if let Some(reasoning) = self.reasoning_effort.as_deref()
1242 && crate::reasoning_preference::ReasoningEffort::parse_strict(reasoning).is_err()
1243 {
1244 return Err("The Runtime Chat reasoning effort is invalid.".to_string());
1245 }
1246 validate_relay_id(&self.workspace.id, "workspace id").map_err(public_validation_error)?;
1247 validate_relay_id(&self.workspace.target_ref, "target reference")
1248 .map_err(public_validation_error)?;
1249 Ok(())
1250 }
1251 }
1252
1253 impl RuntimeChatControlScope {
1254 pub(crate) fn validate_for_turn(&self, virtual_turn_id: &str) -> Result<(), String> {
1255 validate_relay_id(&self.runtime_binding_id, "binding id")
1256 .map_err(public_validation_error)?;
1257 validate_virtual_thread_id(&self.runtime_thread_id).map_err(public_validation_error)?;
1258 validate_virtual_turn_id(virtual_turn_id).map_err(public_validation_error)
1259 }
1260 }
1261
1262 fn validate_existing_binding(
1263 binding: &RelayThreadBinding,
1264 command: &RuntimeChatPrompt,
1265 operation_fingerprint: &str,
1266 system_prompt_fingerprint: Option<&str>,
1267 ) -> Result<(), String> {
1268 if binding.model != command.model
1269 || binding.run_id != command.run_id
1270 || binding.model_provider != command.model_provider
1271 || binding.model_provider_id != command.model_provider_id
1272 {
1273 return Err("The Runtime Chat thread is pinned to a different provider route.".to_string());
1274 }
1275 if command.system_prompt.is_some()
1276 && (binding.first_operation_fingerprint != operation_fingerprint
1277 || binding.system_prompt_fingerprint.as_deref() != system_prompt_fingerprint)
1278 {
1279 return Err(
1280 "Runtime Chat system instructions are allowed only on the first turn.".to_string(),
1281 );
1282 }
1283 Ok(())
1284 }
1285
1286 fn resolve_interrupt_target(
1287 state: &RelayState,
1288 run_id: &str,
1289 scope: &RuntimeChatControlScope,
1290 virtual_turn_id: &str,
1291 ) -> Result<(String, String), String> {
1292 let binding = state
1293 .bindings
1294 .iter()
1295 .find(|binding| {
1296 binding.runtime_binding_id == scope.runtime_binding_id
1297 && binding.virtual_thread_id == scope.runtime_thread_id
1298 })
1299 .ok_or_else(|| "The Runtime Chat interrupt binding is unknown.".to_string())?;
1300 if binding.run_id != run_id {
1301 return Err("The Runtime Chat interrupt binding belongs to another run.".to_string());
1302 }
1303 let turn = binding
1304 .turns
1305 .get(virtual_turn_id)
1306 .ok_or_else(|| "The Runtime Chat interrupt turn is unknown.".to_string())?;
1307 Ok((
1308 binding.native_thread_id.clone(),
1309 turn.native_turn_id.clone(),
1310 ))
1311 }
1312
1313 pub(crate) fn dedicated_chat_system_prompt(account_instructions: Option<&str>) -> String {
1314 match account_instructions
1315 .map(str::trim)
1316 .filter(|value| !value.is_empty())
1317 {
1318 Some(instructions) => format!(
1319 "{ISOLATED_CHAT_SYSTEM_PROMPT}\n\n<account_chat_instructions>\n{instructions}\n</account_chat_instructions>"
1320 ),
1321 None => ISOLATED_CHAT_SYSTEM_PROMPT.to_string(),
1322 }
1323 }
1324
1325 enum ProjectionDecision {
1326 Project(Box<RuntimeChatProjection>),
1327 Skip,
1328 WaitForTurnBinding,
1329 }
1330
1331 fn project_native_event(
1332 binding: &RelayThreadBinding,
1333 event: &RuntimeEventRecord,
1334 ) -> ProjectionDecision {
1335 let Some(native_turn_id) = event.turn_id.as_deref() else {
1336 return ProjectionDecision::Skip;
1337 };
1338 let Some((virtual_turn_id, _)) = binding
1339 .turns
1340 .iter()
1341 .find(|(_, turn)| turn.native_turn_id == native_turn_id)
1342 else {
1343 return ProjectionDecision::WaitForTurnBinding;
1344 };
1345 let base = || RuntimeChatProjection {
1346 run_id: binding.run_id.clone(),
1347 native_thread_id: binding.native_thread_id.clone(),
1348 native_seq: event.seq,
1349 source_event_id: source_event_id(&binding.native_thread_id, event.seq),
1350 virtual_thread_id: binding.virtual_thread_id.clone(),
1351 virtual_turn_id: virtual_turn_id.clone(),
1352 event: "",
1353 timestamp: event.timestamp.to_rfc3339(),
1354 payload: Value::Null,
1355 };
1356 match event.event.as_str() {
1357 "turn.started" => {
1358 let mut projection = base();
1359 projection.event = "turn.started";
1360 projection.payload = json!({
1361 "turn": {
1362 "model": event.payload.pointer("/turn/effective_model")
1363 .and_then(Value::as_str)
1364 .unwrap_or(binding.model.as_str()),
1365 "mode": "chat",
1366 },
1367 });
1368 ProjectionDecision::Project(Box::new(projection))
1369 }
1370 "item.delta"
1371 if event.payload.get("kind").and_then(Value::as_str) == Some("agent_message") =>
1372 {
1373 let Some(delta) = event.payload.get("delta").and_then(Value::as_str) else {
1374 return ProjectionDecision::Skip;
1375 };
1376 let mut projection = base();
1377 projection.event = "item.delta";
1378 projection.payload = json!({ "kind": "agent_message", "delta": delta });
1379 ProjectionDecision::Project(Box::new(projection))
1380 }
1381 "turn.completed" => {
1382 let turn = event
1383 .payload
1384 .get("turn")
1385 .cloned()
1386 .unwrap_or_else(|| json!({}));
1387 let status = turn
1388 .get("status")
1389 .and_then(Value::as_str)
1390 .unwrap_or("failed");
1391 let mut projection = base();
1392 projection.event = "turn.completed";
1393 projection.payload = json!({
1394 "turn": {
1395 "status": status,
1396 "usage": turn.get("usage").cloned().unwrap_or_else(|| json!({})),
1397 "effective_model": turn.get("effective_model").and_then(Value::as_str)
1398 .unwrap_or(binding.model.as_str()),
1399 "effective_provider": turn.get("effective_provider").and_then(Value::as_str)
1400 .unwrap_or(binding.model_provider.as_str()),
1401 "effective_billing_surface": turn.get("effective_billing_surface")
1402 .and_then(Value::as_str).unwrap_or("provider_byok"),
1403 "effective_billing_mode": turn.get("effective_billing_mode")
1404 .and_then(Value::as_str).unwrap_or("local"),
1405 },
1406 });
1407 ProjectionDecision::Project(Box::new(projection))
1408 }
1409 _ => ProjectionDecision::Skip,
1410 }
1411 }
1412
1413 fn validate_owner_component(value: &str, label: &str) -> Result<(), String> {
1414 if value.trim() != value || value.is_empty() || value.len() > 512 || value.contains('\0') {
1415 return Err(format!("The Runtime Chat {label} owner is invalid."));
1416 }
1417 Ok(())
1418 }
1419
1420 fn scoped_private_dir(root: &Path, target_ref: &str, session_id: &str) -> PathBuf {
1421 let mut hasher = Sha256::new();
1422 hasher.update(b"codewhale.runtime-chat-scope.v1\0");
1423 hasher.update(target_ref.as_bytes());
1424 hasher.update(b"\0");
1425 hasher.update(session_id.as_bytes());
1426 root.join(format!("scope-{}", hex_digest(hasher.finalize())))
1427 }
1428
1429 fn owner_scope_fingerprint(account_ref: &str, target_ref: &str, session_id: &str) -> String {
1430 let mut hasher = Sha256::new();
1431 hasher.update(b"codewhale.runtime-chat-owner.v1\0");
1432 hasher.update(account_ref.as_bytes());
1433 hasher.update(b"\0");
1434 hasher.update(target_ref.as_bytes());
1435 hasher.update(b"\0");
1436 hasher.update(session_id.as_bytes());
1437 hex_digest(hasher.finalize())
1438 }
1439
1440 fn bind_owner_scope(state: &mut RelayState, owner: &str) -> Result<(), String> {
1441 validate_fingerprint(owner)
1442 .map_err(|_| "The Runtime Chat account owner is invalid.".to_string())?;
1443 match state.owner_scope_fingerprint.as_deref() {
1444 Some(existing) if existing == owner => Ok(()),
1445 Some(_) => Err("This Runtime Chat session belongs to another account.".to_string()),
1446 None if state.bindings.is_empty() => {
1447 state.owner_scope_fingerprint = Some(owner.to_string());
1448 Ok(())
1449 }
1450 None => Err("The saved Runtime Chat state has no trusted account owner.".to_string()),
1451 }
1452 }
1453
1454 fn source_event_id(native_thread_id: &str, native_seq: u64) -> String {
1455 let mut hasher = Sha256::new();
1456 hasher.update(b"codewhale.runtime-chat-source-event.v1\0");
1457 hasher.update(native_thread_id.as_bytes());
1458 hasher.update(b"\0");
1459 hasher.update(native_seq.to_be_bytes());
1460 format!("native_event_{}", hex_digest(hasher.finalize()))
1461 }
1462
1463 fn reserved_native_turn_id(
1464 native_thread_id: &str,
1465 runtime_binding_id: &str,
1466 virtual_turn_id: &str,
1467 operation_fingerprint: &str,
1468 ) -> String {
1469 let mut hasher = Sha256::new();
1470 hasher.update(b"codewhale.runtime-chat-native-turn.v1\0");
1471 hasher.update(native_thread_id.as_bytes());
1472 hasher.update(b"\0");
1473 hasher.update(runtime_binding_id.as_bytes());
1474 hasher.update(b"\0");
1475 hasher.update(virtual_turn_id.as_bytes());
1476 hasher.update(b"\0");
1477 hasher.update(operation_fingerprint.as_bytes());
1478 format!("turn_{}", hex_digest(hasher.finalize()))
1479 }
1480
1481 fn hex_digest(digest: impl AsRef<[u8]>) -> String {
1482 digest
1483 .as_ref()
1484 .iter()
1485 .map(|byte| format!("{byte:02x}"))
1486 .collect()
1487 }
1488
1489 fn canonical_json_value(value: &Value) -> Value {
1490 match value {
1491 Value::Array(values) => Value::Array(values.iter().map(canonical_json_value).collect()),
1492 Value::Object(values) => {
1493 let mut keys = values.keys().collect::<Vec<_>>();
1494 keys.sort_unstable();
1495 let mut canonical = serde_json::Map::new();
1496 for key in keys {
1497 canonical.insert(key.clone(), canonical_json_value(&values[key]));
1498 }
1499 Value::Object(canonical)
1500 }
1501 _ => value.clone(),
1502 }
1503 }
1504
1505 fn isolated_chat_execution_config(
1506 config: &Config,
1507 private_dir: &Path,
1508 ) -> Result<(Config, PathBuf), String> {
1509 let workspace = private_dir.join("chat-workspace");
1510 let context_dir = private_dir.join("prompt-context");
1511 let skills_dir = context_dir.join("skills-empty");
1512 fs::create_dir_all(&workspace)
1513 .and_then(|()| fs::create_dir_all(&skills_dir))
1514 .map_err(|_| "Runtime Chat could not prepare its isolated Chat context.".to_string())?;
1515 #[cfg(unix)]
1516 {
1517 use std::os::unix::fs::PermissionsExt;
1518 for path in [&workspace, &context_dir, &skills_dir] {
1519 fs::set_permissions(path, fs::Permissions::from_mode(0o700)).map_err(|_| {
1520 "Runtime Chat could not protect its isolated Chat context.".to_string()
1521 })?;
1522 }
1523 }
1524
1525 let mut execution = config.clone();
1526 execution.runtime_chat_isolated = true;
1527 execution.skills_dir = Some(skills_dir.to_string_lossy().into_owned());
1528 execution.instructions = None;
1529 execution.notes_path = Some(
1530 context_dir
1531 .join("notes-disabled")
1532 .to_string_lossy()
1533 .into_owned(),
1534 );
1535 execution.mcp_config_path = Some(
1536 context_dir
1537 .join("mcp-disabled.json")
1538 .to_string_lossy()
1539 .into_owned(),
1540 );
1541 execution.memory = Some(MemoryConfig {
1542 enabled: Some(false),
1543 backend: Some(MemoryBackend::Off),
1544 });
1545 execution.memory_path = Some(
1546 context_dir
1547 .join("memory-disabled.md")
1548 .to_string_lossy()
1549 .into_owned(),
1550 );
1551 execution.context.project_pack = Some(false);
1552 execution
1553 .skills
1554 .get_or_insert_with(SkillsConfig::default)
1555 .scan_codewhale_only = Some(true);
1556 Ok((execution, workspace))
1557 }
1558
1559 #[derive(Debug)]
1560 struct RelayScopeLock {
1561 _file: File,
1562 }
1563
1564 impl RelayScopeLock {
1565 fn acquire(path: &Path) -> Result<Self> {
1566 if fs::symlink_metadata(path).is_ok_and(|metadata| metadata.file_type().is_symlink()) {
1567 bail!("refusing a symlinked Runtime Chat owner lock");
1568 }
1569 let mut options = fs::OpenOptions::new();
1570 options.read(true).write(true).create(true);
1571 #[cfg(unix)]
1572 {
1573 use std::os::unix::fs::OpenOptionsExt as _;
1574 options.mode(0o600).custom_flags(libc::O_NOFOLLOW);
1575 }
1576 #[cfg(windows)]
1577 {
1578 use std::os::windows::fs::OpenOptionsExt as _;
1579 use windows_sys::Win32::Storage::FileSystem::FILE_FLAG_OPEN_REPARSE_POINT;
1580 options.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT);
1581 }
1582 let file = options.open(path).context("open Runtime Chat owner lock")?;
1583 if !file
1584 .metadata()
1585 .context("inspect Runtime Chat owner lock")?
1586 .file_type()
1587 .is_file()
1588 {
1589 bail!("Runtime Chat owner lock is not a regular file");
1590 }
1591 #[cfg(unix)]
1592 {
1593 use std::os::unix::fs::PermissionsExt as _;
1594 file.set_permissions(fs::Permissions::from_mode(0o600))
1595 .context("protect Runtime Chat owner lock")?;
1596 }
1597 // Same-process drop-then-reopen can observe WouldBlock for a brief
1598 // window while the previous fd is still closing (#5735). Retry only
1599 // that contention; a lock that stays held is still ownership.
1600 let deadline = Instant::now() + Duration::from_millis(25);
1601 loop {
1602 match Self::try_lock_exclusive(&file) {
1603 Ok(()) => break,
1604 Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
1605 if Instant::now() >= deadline {
1606 return Err(error).context("acquire Runtime Chat owner lock");
1607 }
1608 std::thread::yield_now();
1609 std::thread::sleep(Duration::from_millis(1));
1610 }
1611 Err(error) => {
1612 return Err(error).context("acquire Runtime Chat owner lock");
1613 }
1614 }
1615 }
1616 Ok(Self { _file: file })
1617 }
1618
1619 fn try_lock_exclusive(file: &File) -> std::io::Result<()> {
1620 #[cfg(unix)]
1621 {
1622 use std::os::fd::AsRawFd as _;
1623 // SAFETY: `file` owns a valid descriptor for the duration of the
1624 // call and remains alive in `Self` for the full lock lifetime.
1625 if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } != 0 {
1626 return Err(std::io::Error::last_os_error());
1627 }
1628 Ok(())
1629 }
1630 #[cfg(windows)]
1631 {
1632 use std::os::windows::io::AsRawHandle as _;
1633 use windows_sys::Win32::Storage::FileSystem::LockFile;
1634 // SAFETY: `file` owns a valid handle that remains alive in `Self`.
1635 if unsafe { LockFile(file.as_raw_handle() as _, 0, 0, u32::MAX, u32::MAX) } == 0 {
1636 return Err(std::io::Error::last_os_error());
1637 }
1638 Ok(())
1639 }
1640 #[cfg(not(any(unix, windows)))]
1641 {
1642 let _ = file;
1643 Ok(())
1644 }
1645 }
1646 }
1647
1648 impl Drop for RelayScopeLock {
1649 fn drop(&mut self) {
1650 // close() also releases, but unlocking first lets a same-process
1651 // reopen proceed without racing the previous fd's teardown (#5735).
1652 #[cfg(unix)]
1653 {
1654 use std::os::fd::AsRawFd as _;
1655 // SAFETY: Drop runs only while `_file` still owns this descriptor.
1656 unsafe {
1657 libc::flock(self._file.as_raw_fd(), libc::LOCK_UN);
1658 }
1659 }
1660 #[cfg(windows)]
1661 {
1662 use std::os::windows::io::AsRawHandle as _;
1663 use windows_sys::Win32::Storage::FileSystem::UnlockFile;
1664 // SAFETY: Drop runs only while `_file` still owns this handle.
1665 unsafe {
1666 UnlockFile(self._file.as_raw_handle() as _, 0, 0, u32::MAX, u32::MAX);
1667 }
1668 }
1669 }
1670 }
1671
1672 fn load_state(path: &Path) -> Result<RelayState> {
1673 match fs::read(path) {
1674 Ok(bytes) => serde_json::from_slice(&bytes).context("decode Runtime Chat binding state"),
1675 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(RelayState::default()),
1676 Err(error) => Err(error).context("read Runtime Chat binding state"),
1677 }
1678 }
1679
1680 fn persist_state(path: &Path, state: &RelayState) -> Result<()> {
1681 state.validate()?;
1682 #[cfg(test)]
1683 if take_state_persist_failure(path) {
1684 bail!("injected Runtime Chat state persistence failure");
1685 }
1686 let body = serde_json::to_vec(state).context("encode Runtime Chat binding state")?;
1687 crate::utils::write_atomic(path, &body).context("persist Runtime Chat binding state")
1688 }
1689
1690 fn validate_relay_id(value: &str, label: &str) -> Result<()> {
1691 if value.is_empty()
1692 || value.len() > MAX_RELAY_ID_BYTES
1693 || value.contains("..")
1694 || value.contains("://")
1695 || !value.bytes().all(|byte| {
1696 byte.is_ascii_alphanumeric()
1697 || matches!(byte, b'.' | b'_' | b':' | b'@' | b'/' | b'+' | b'~' | b'-')
1698 })
1699 {
1700 bail!("invalid Runtime Chat {label}");
1701 }
1702 Ok(())
1703 }
1704
1705 fn validate_operation_key(value: &str) -> Result<()> {
1706 validate_relay_id(value, "operation key")?;
1707 if value.len() > MAX_OPERATION_KEY_BYTES {
1708 bail!("Runtime Chat operation key is too long");
1709 }
1710 Ok(())
1711 }
1712
1713 fn validate_virtual_thread_id(value: &str) -> Result<()> {
1714 if value.len() != 37
1715 || !value.starts_with("local_thread_")
1716 || !value[13..]
1717 .bytes()
1718 .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f'))
1719 {
1720 bail!("invalid Runtime Chat virtual thread id");
1721 }
1722 Ok(())
1723 }
1724
1725 fn validate_virtual_turn_id(value: &str) -> Result<()> {
1726 if value.len() != 35
1727 || !value.starts_with("local_turn_")
1728 || !value[11..]
1729 .bytes()
1730 .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f'))
1731 {
1732 bail!("invalid Runtime Chat virtual turn id");
1733 }
1734 Ok(())
1735 }
1736
1737 fn validate_native_record_id(value: &str, label: &str) -> Result<()> {
1738 if value.len() < 5
1739 || value.len() > 80
1740 || !value
1741 .bytes()
1742 .all(|byte| byte.is_ascii_alphanumeric() || byte == b'_')
1743 {
1744 bail!("invalid {label}");
1745 }
1746 Ok(())
1747 }
1748
1749 fn validate_route_id(value: &str, label: &str) -> Result<()> {
1750 if !crate::runtime_api::runtime_chat_route_id_is_safe(value) {
1751 bail!("invalid Runtime Chat {label}");
1752 }
1753 Ok(())
1754 }
1755
1756 fn validate_model_id(value: &str) -> Result<()> {
1757 if !crate::runtime_api::runtime_chat_model_id_is_safe(value) {
1758 bail!("invalid Runtime Chat model id");
1759 }
1760 Ok(())
1761 }
1762
1763 fn validate_fingerprint(value: &str) -> Result<()> {
1764 if value.len() != 64 || !value.bytes().all(|byte| byte.is_ascii_hexdigit()) {
1765 bail!("invalid Runtime Chat fingerprint");
1766 }
1767 Ok(())
1768 }
1769
1770 fn fingerprint(value: &str) -> String {
1771 hex_digest(Sha256::digest(value.as_bytes()))
1772 }
1773
1774 fn runtime_chat_request_fingerprint(command: &RuntimeChatPrompt) -> Result<String, String> {
1775 serde_json::to_value(command)
1776 .map(|value| canonical_json_value(&value))
1777 .and_then(|value| serde_json::to_vec(&value))
1778 .map(|bytes| hex_digest(Sha256::digest(bytes)))
1779 .map_err(|_| "Runtime Chat could not fingerprint its validated request.".to_string())
1780 }
1781
1782 fn public_validation_error(_error: anyhow::Error) -> String {
1783 "The Runtime Chat command contains an invalid opaque identity.".to_string()
1784 }
1785
1786 fn sanitized_runtime_error(action: &str, error: &anyhow::Error) -> String {
1787 let text = error.to_string().to_ascii_lowercase();
1788 if text.contains("operation_key") || text.contains("operation key") {
1789 return "The Runtime Chat operation key does not match its original turn.".to_string();
1790 }
1791 format!("The isolated Runtime Chat turn could not {action}.")
1792 }
1793
1794 #[cfg(test)]
1795 mod tests {
1796 use super::*;
1797 use chrono::Utc;
1798 use std::time::Duration;
1799
1800 fn open_host(root: &Path) -> RuntimeChatRelayHost {
1801 RuntimeChatRelayHost::open(
1802 Config::default(),
1803 Arc::new(PluginRegistry::empty(root)),
1804 root.to_path_buf(),
1805 "target_fixture".to_string(),
1806 "session_fixture".to_string(),
1807 )
1808 .expect("open Runtime Chat host")
1809 }
1810
1811 fn binding() -> RelayThreadBinding {
1812 RelayThreadBinding {
1813 run_id: "run_fixture".to_string(),
1814 runtime_binding_id: "binding_fixture".to_string(),
1815 virtual_thread_id: format!("local_thread_{}", "a".repeat(24)),
1816 native_thread_id: "thr_fixture".to_string(),
1817 model: "model-1".to_string(),
1818 model_provider: "custom".to_string(),
1819 model_provider_id: "local-provider".to_string(),
1820 first_operation_fingerprint: fingerprint("operation-1"),
1821 system_prompt_fingerprint: None,
1822 turns: BTreeMap::from([(
1823 format!("local_turn_{}", "b".repeat(24)),
1824 RelayTurnBinding {
1825 native_turn_id: "turn_native".to_string(),
1826 operation_fingerprint: fingerprint("operation-1"),
1827 request_fingerprint: fingerprint("request-1"),
1828 terminal_projected: false,
1829 start_rejected: false,
1830 },
1831 )]),
1832 projected_native_seq: 0,
1833 }
1834 }
1835
1836 #[test]
1837 fn persisted_binding_state_rejects_duplicate_or_nonopaque_authority() {
1838 let one = binding();
1839 let state = RelayState {
1840 schema_version: STATE_SCHEMA_VERSION,
1841 owner_scope_fingerprint: Some(fingerprint("owner")),
1842 bindings: vec![one.clone()],
1843 };
1844 state.validate().unwrap();
1845 let duplicate = RelayState {
1846 schema_version: STATE_SCHEMA_VERSION,
1847 owner_scope_fingerprint: Some(fingerprint("owner")),
1848 bindings: vec![one.clone(), one],
1849 };
1850 assert!(duplicate.validate().is_err());
1851 }
1852
1853 #[test]
1854 fn projection_rewrites_native_ids_to_virtual_chat_ids_and_keeps_route_receipt() {
1855 let binding = binding();
1856 let event = RuntimeEventRecord {
1857 schema_version: 2,
1858 seq: 9,
1859 timestamp: Utc::now(),
1860 thread_id: binding.native_thread_id.clone(),
1861 turn_id: Some("turn_native".to_string()),
1862 item_id: None,
1863 event: "turn.completed".to_string(),
1864 payload: json!({
1865 "turn": {
1866 "status": "completed",
1867 "usage": { "input_tokens": 2, "output_tokens": 3 },
1868 "effective_model": "model-1",
1869 "effective_provider": "custom",
1870 "effective_billing_surface": "provider_byok",
1871 "effective_billing_mode": "local",
1872 }
1873 }),
1874 };
1875 let ProjectionDecision::Project(projected) = project_native_event(&binding, &event) else {
1876 panic!("terminal event must project");
1877 };
1878 assert_eq!(projected.virtual_thread_id, binding.virtual_thread_id);
1879 assert!(projected.virtual_turn_id.starts_with("local_turn_"));
1880 assert_eq!(projected.event, "turn.completed");
1881 assert_eq!(
1882 projected.source_event_id,
1883 source_event_id(&binding.native_thread_id, event.seq)
1884 );
1885 assert_eq!(projected.source_event_id.len(), 77);
1886 assert_eq!(projected.payload["turn"]["effective_billing_mode"], "local");
1887 assert!(!projected.payload.to_string().contains("thr_fixture"));
1888 }
1889
1890 #[test]
1891 fn chat_command_shape_requires_empty_tools_and_exact_chat_modes() {
1892 let mut prompt = RuntimeChatPrompt {
1893 images: Vec::new(),
1894 max_output_tokens: None,
1895 command_type: "prompt.request".to_string(),
1896 run_id: "run_fixture".to_string(),
1897 turn_id: format!("local_turn_{}", "b".repeat(24)),
1898 operation_key: "operation-1".to_string(),
1899 runtime_binding_id: "binding_fixture".to_string(),
1900 runtime_thread_id: format!("local_thread_{}", "a".repeat(24)),
1901 prompt: "hello".to_string(),
1902 system_prompt: None,
1903 model: "model-1".to_string(),
1904 model_provider: "custom".to_string(),
1905 model_provider_id: "local-provider".to_string(),
1906 reasoning_effort: Some("high".to_string()),
1907 allowed_tools: Vec::new(),
1908 mode: "chat".to_string(),
1909 requested_mode: "chat".to_string(),
1910 workspace: RuntimeChatWorkspace {
1911 id: "workspace_fixture".to_string(),
1912 target_ref: "target_fixture".to_string(),
1913 },
1914 };
1915 prompt.validate_shape().unwrap();
1916 for reasoning in ["minimal", "ultra"] {
1917 prompt.reasoning_effort = Some(reasoning.into());
1918 prompt.validate_shape().unwrap();
1919 }
1920 prompt.reasoning_effort = Some("invented-effort".into());
1921 assert!(prompt.validate_shape().is_err());
1922 prompt.reasoning_effort = None;
1923 prompt.allowed_tools.push("bash".to_string());
1924 assert!(prompt.validate_shape().is_err());
1925 prompt.allowed_tools.clear();
1926 prompt.requested_mode = "work".to_string();
1927 assert!(prompt.validate_shape().is_err());
1928 }
1929
1930 #[tokio::test]
1931 async fn active_participant_rejects_new_and_exact_replay_without_blocking_or_starting() {
1932 let root = tempfile::tempdir().unwrap();
1933 let config = Config {
1934 provider: Some("ollama".to_string()),
1935 default_text_model: Some("relay-local:fixture".to_string()),
1936 ..Config::default()
1937 };
1938 let host = RuntimeChatRelayHost::open(
1939 config,
1940 Arc::new(PluginRegistry::empty(root.path())),
1941 root.path().to_path_buf(),
1942 "target_fixture".to_string(),
1943 "session_fixture".to_string(),
1944 )
1945 .unwrap();
1946 host.bind_account("account_fixture", "target_fixture")
1947 .unwrap();
1948 host.authorize_run("run_fixture").unwrap();
1949 let prompt = RuntimeChatPrompt {
1950 images: Vec::new(),
1951 max_output_tokens: None,
1952 command_type: "prompt.request".to_string(),
1953 run_id: "run_fixture".to_string(),
1954 turn_id: format!("local_turn_{}", "e".repeat(24)),
1955 operation_key: "operation-gate-fixture".to_string(),
1956 runtime_binding_id: "binding-gate-fixture".to_string(),
1957 runtime_thread_id: format!("local_thread_{}", "f".repeat(24)),
1958 prompt: "hello".to_string(),
1959 system_prompt: None,
1960 model: "relay-local:fixture".to_string(),
1961 model_provider: "ollama".to_string(),
1962 model_provider_id: "ollama".to_string(),
1963 reasoning_effort: None,
1964 allowed_tools: Vec::new(),
1965 mode: "chat".to_string(),
1966 requested_mode: "chat".to_string(),
1967 workspace: RuntimeChatWorkspace {
1968 id: "workspace_fixture".to_string(),
1969 target_ref: "target_fixture".to_string(),
1970 },
1971 };
1972
1973 let participant = crate::client::acquire_remote_control_inference_participant().await;
1974 let direct_error = host.try_ensure_inference_ownership().unwrap_err();
1975 assert!(
1976 direct_error.contains("active local provider work"),
1977 "{direct_error}"
1978 );
1979 for attempt in 0..2 {
1980 let error = tokio::time::timeout(Duration::from_secs(1), host.apply_prompt(&prompt))
1981 .await
1982 .expect("Runtime Chat admission must not deadlock behind the UI-owned writer")
1983 .unwrap_err();
1984 assert!(error.contains("active local provider work"), "{error}");
1985 assert!(
1986 !host.has_any_unsettled_turns(),
1987 "rejected attempt {attempt} must not leave a phantom reservation"
1988 );
1989 }
1990 assert_eq!(
1991 host.native_turn_count_for_binding_for_tests(&prompt.runtime_binding_id)
1992 .await,
1993 0,
1994 "neither the new command nor its exact replay may reach provider-backed native start"
1995 );
1996 drop(participant);
1997 host.try_ensure_inference_ownership().unwrap();
1998 host.release_inference_ownership_if_settled();
1999 }
2000
2001 #[tokio::test]
2002 async fn recovered_reservation_without_native_turn_settles_after_definitive_start_rejection() {
2003 let root = tempfile::tempdir().unwrap();
2004 let host = open_host(root.path());
2005 host.bind_account("account_fixture", "target_fixture")
2006 .unwrap();
2007 host.authorize_run("run_fixture").unwrap();
2008 let identity = host.config.active_provider_identity().unwrap();
2009 let provider = identity.provider;
2010 let model_provider_id = host
2011 .config
2012 .active_provider_identity()
2013 .unwrap()
2014 .persisted_id()
2015 .unwrap_or_else(|| provider.as_str())
2016 .to_string();
2017 let prompt = RuntimeChatPrompt {
2018 images: Vec::new(),
2019 max_output_tokens: None,
2020 command_type: "prompt.request".to_string(),
2021 run_id: "run_fixture".to_string(),
2022 turn_id: format!("local_turn_{}", "8".repeat(24)),
2023 operation_key: "operation-crash-window".to_string(),
2024 runtime_binding_id: "binding-crash-window".to_string(),
2025 runtime_thread_id: format!("local_thread_{}", "9".repeat(24)),
2026 prompt: "hello".to_string(),
2027 system_prompt: None,
2028 model: host.config.default_model(),
2029 model_provider: provider.as_str().to_string(),
2030 model_provider_id,
2031 reasoning_effort: None,
2032 allowed_tools: Vec::new(),
2033 mode: "chat".to_string(),
2034 requested_mode: "chat".to_string(),
2035 workspace: RuntimeChatWorkspace {
2036 id: "workspace_fixture".to_string(),
2037 target_ref: "target_fixture".to_string(),
2038 },
2039 };
2040 host.install_prompt_replay_for_tests(&prompt, false)
2041 .await
2042 .unwrap();
2043 assert!(host.has_any_unsettled_turns());
2044 assert_eq!(
2045 host.native_turn_count_for_binding_for_tests(&prompt.runtime_binding_id)
2046 .await,
2047 0
2048 );
2049
2050 let error = host.apply_prompt(&prompt).await.unwrap_err();
2051 assert!(error.contains("could not start"), "{error}");
2052 assert!(
2053 !host.has_any_unsettled_turns(),
2054 "a proven pre-native crash reservation must not hold ownership forever"
2055 );
2056 assert!(!host.inference_ownership_is_held_for_tests());
2057 assert_eq!(
2058 host.native_turn_count_for_binding_for_tests(&prompt.runtime_binding_id)
2059 .await,
2060 0
2061 );
2062 }
2063
2064 #[test]
2065 fn catalog_is_challenge_bound_active_ready_and_secret_free() {
2066 let mut config = Config {
2067 provider: Some("ollama".to_string()),
2068 default_text_model: Some(crate::config::DEFAULT_OLLAMA_MODEL.to_string()),
2069 ..Config::default()
2070 };
2071 config.set_legacy_root(Some("must-not-cross".to_string()), None);
2072 config.set_legacy_root(None, Some("http://127.0.0.1:11434/v1".to_string()));
2073 let challenge = "c".repeat(32);
2074 assert!(crate::runtime_api::runtime_chat_relay_catalog(&config, &challenge).is_err());
2075 config.default_text_model = Some("relay-local:fixture".to_string());
2076 let catalog = crate::runtime_api::runtime_chat_relay_catalog(&config, &challenge).unwrap();
2077 assert_eq!(catalog["protocol"], "codewhale.runtime-chat-relay.v1");
2078 assert_eq!(catalog["challenge"], challenge);
2079 assert_eq!(catalog["runtime"]["service"], "codewhale-runtime-api");
2080 assert_eq!(catalog["runtime"]["apiVersion"], "1.0");
2081 assert_eq!(catalog["runtime"]["capabilities"]["stable_event_ids"], true);
2082 assert_eq!(catalog["providers"].as_array().unwrap().len(), 1);
2083 assert_eq!(
2084 catalog["providers"][0]["models"][0]["imageInput"],
2085 "unknown"
2086 );
2087 assert_eq!(
2088 catalog["runtime"]["capabilities"]["turn_image_inputs"],
2089 true
2090 );
2091 let serialized = catalog.to_string();
2092 assert!(!serialized.contains("must-not-cross"));
2093 assert!(!serialized.contains("127.0.0.1"));
2094 assert!(!serialized.contains("baseUrl"));
2095 assert!(!serialized.contains("endpoint"));
2096 }
2097
2098 #[test]
2099 fn interrupt_resolution_is_bound_to_the_current_run() {
2100 let binding = binding();
2101 let turn_id = binding.turns.keys().next().unwrap().clone();
2102 let scope = RuntimeChatControlScope {
2103 runtime_binding_id: binding.runtime_binding_id.clone(),
2104 runtime_thread_id: binding.virtual_thread_id.clone(),
2105 };
2106 let state = RelayState {
2107 schema_version: STATE_SCHEMA_VERSION,
2108 owner_scope_fingerprint: Some(fingerprint("owner")),
2109 bindings: vec![binding],
2110 };
2111 assert!(resolve_interrupt_target(&state, "run_fixture", &scope, &turn_id).is_ok());
2112 assert!(resolve_interrupt_target(&state, "run_other", &scope, &turn_id).is_err());
2113 }
2114
2115 #[test]
2116 fn scoped_state_is_exclusive_account_bound_and_restart_stable() {
2117 let root = tempfile::tempdir().unwrap();
2118 let first_path = scoped_private_dir(root.path(), "target_fixture", "session_fixture");
2119 let same_path = scoped_private_dir(root.path(), "target_fixture", "session_fixture");
2120 let other_path = scoped_private_dir(root.path(), "target_fixture", "session_other");
2121 assert_eq!(first_path, same_path);
2122 assert_ne!(first_path, other_path);
2123 assert!(!first_path.to_string_lossy().contains("target_fixture"));
2124
2125 fs::create_dir_all(&first_path).unwrap();
2126 let lock_path = first_path.join(SCOPE_LOCK_FILE);
2127 let first_lock = RelayScopeLock::acquire(&lock_path).unwrap();
2128 assert!(RelayScopeLock::acquire(&lock_path).is_err());
2129
2130 let owner = owner_scope_fingerprint("account_one", "target_fixture", "session_fixture");
2131 let mut state = RelayState::default();
2132 bind_owner_scope(&mut state, &owner).unwrap();
2133 bind_owner_scope(&mut state, &owner).unwrap();
2134 assert!(
2135 bind_owner_scope(
2136 &mut state,
2137 &owner_scope_fingerprint("account_other", "target_fixture", "session_fixture")
2138 )
2139 .is_err()
2140 );
2141 state.bindings.push(binding());
2142 let state_path = first_path.join(STATE_FILE);
2143 persist_state(&state_path, &state).unwrap();
2144 assert_eq!(load_state(&state_path).unwrap(), state);
2145 assert!(
2146 fs::read_dir(&first_path)
2147 .unwrap()
2148 .filter_map(Result::ok)
2149 .all(|entry| !entry.file_name().to_string_lossy().ends_with(".tmp"))
2150 );
2151
2152 drop(first_lock);
2153 RelayScopeLock::acquire(&lock_path).unwrap();
2154 }
2155
2156 #[test]
2157 fn isolated_chat_prompt_drops_local_project_memory_and_skill_context() {
2158 let root = tempfile::tempdir().unwrap();
2159 let config = Config {
2160 skills_dir: Some("/Users/alice/CANARY_SKILLS".to_string()),
2161 instructions: Some(vec!["/Users/alice/CANARY_AGENTS.md".to_string()]),
2162 memory_path: Some("/Users/alice/CANARY_MEMORY.md".to_string()),
2163 memory: Some(MemoryConfig {
2164 enabled: Some(true),
2165 backend: Some(MemoryBackend::Native),
2166 }),
2167 context: ContextConfig {
2168 project_pack: Some(true),
2169 },
2170 ..Config::default()
2171 };
2172
2173 let (execution, workspace) = isolated_chat_execution_config(&config, root.path()).unwrap();
2174 assert!(workspace.starts_with(root.path()));
2175 assert_ne!(workspace, PathBuf::from("/Users/alice"));
2176 assert!(!execution.memory_enabled());
2177 assert!(execution.instructions_paths().is_empty());
2178 assert!(!execution.project_context_pack_enabled());
2179 assert!(execution.skills_dir().starts_with(root.path()));
2180 assert!(execution.memory_path().starts_with(root.path()));
2181 assert!(execution.mcp_config_path().starts_with(root.path()));
2182 assert!(execution.notes_path().starts_with(root.path()));
2183 assert!(execution.skills_config().scan_codewhale_only());
2184
2185 let prompt = dedicated_chat_system_prompt(None);
2186 assert_eq!(prompt, ISOLATED_CHAT_SYSTEM_PROMPT);
2187 for canary in [
2188 "CANARY_SKILLS",
2189 "CANARY_AGENTS",
2190 "CANARY_MEMORY",
2191 "/Users/alice",
2192 ] {
2193 assert!(!prompt.contains(canary));
2194 }
2195 let account_prompt = dedicated_chat_system_prompt(Some("Reply in short paragraphs."));
2196 assert!(account_prompt.starts_with(ISOLATED_CHAT_SYSTEM_PROMPT));
2197 assert!(account_prompt.contains("<account_chat_instructions>"));
2198 }
2199
2200 #[test]
2201 fn unsafe_model_ids_are_rejected_and_never_cross_the_catalog() {
2202 for model in [
2203 "/Users/alice/private-model",
2204 "~/.ssh/provider-model",
2205 "sk-live-secret",
2206 "../secrets/model",
2207 "redacted-local-path",
2208 "alice:hunter2@internal.example:443/model",
2209 "internal.example:443/model",
2210 "internal:443/model",
2211 ] {
2212 assert!(validate_model_id(model).is_err(), "accepted {model:?}");
2213 let config = Config {
2214 provider: Some("ollama".to_string()),
2215 default_text_model: Some(model.to_string()),
2216 ..Config::default()
2217 };
2218 if let Ok(catalog) =
2219 crate::runtime_api::runtime_chat_relay_catalog(&config, &"c".repeat(32))
2220 {
2221 let ids = catalog["providers"][0]["models"]
2222 .as_array()
2223 .unwrap()
2224 .iter()
2225 .filter_map(|entry| entry["id"].as_str())
2226 .collect::<Vec<_>>();
2227 assert!(!ids.contains(&model), "catalog leaked {model:?}");
2228 }
2229 }
2230 validate_model_id("anthropic/claude-sonnet-5").unwrap();
2231
2232 for provider_id in ["ghp_secret-provider", "hf_secret-provider", "glpat-secret"] {
2233 assert!(!crate::runtime_api::runtime_chat_route_id_is_safe(
2234 provider_id
2235 ));
2236 let config = Config {
2237 provider: Some(provider_id.to_string()),
2238 providers: Some(crate::config::ProvidersConfig {
2239 custom: std::collections::HashMap::from([(
2240 provider_id.to_string(),
2241 crate::config::ProviderConfig {
2242 kind: Some("openai-compatible".to_string()),
2243 api_key: Some("fixture-key".to_string()),
2244 base_url: Some("https://example.test/v1".to_string()),
2245 model: Some("safe-model".to_string()),
2246 ..Default::default()
2247 },
2248 )]),
2249 ..Default::default()
2250 }),
2251 ..Config::default()
2252 };
2253 let error = crate::runtime_api::runtime_chat_relay_catalog(&config, &"c".repeat(32))
2254 .unwrap_err();
2255 assert!(!error.contains(provider_id));
2256 }
2257 }
2258
2259 #[test]
2260 fn native_source_event_id_is_stable_across_transport_replay() {
2261 let first = source_event_id("thr_fixture", 41);
2262 assert_eq!(first, source_event_id("thr_fixture", 41));
2263 assert_ne!(first, source_event_id("thr_fixture", 42));
2264 assert_ne!(first, source_event_id("thr_other", 41));
2265 assert!(first.starts_with("native_event_"));
2266 assert!(
2267 first
2268 .bytes()
2269 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
2270 );
2271 }
2272
2273 #[test]
2274 fn failed_state_writes_never_become_in_memory_authority_and_exact_retry_reopens() {
2275 let root = tempfile::tempdir().unwrap();
2276 let host = open_host(root.path());
2277
2278 inject_state_persist_failures(&host.state_path, 1);
2279 assert!(
2280 host.bind_account("account_fixture", "target_fixture")
2281 .is_err()
2282 );
2283 assert!(host.state.lock().owner_scope_fingerprint.is_none());
2284 host.bind_account("account_fixture", "target_fixture")
2285 .unwrap();
2286
2287 let mut relay_binding = binding();
2288 relay_binding.turns.clear();
2289 inject_state_persist_failures(&host.state_path, 1);
2290 assert!(host.insert_binding(relay_binding.clone()).is_err());
2291 assert!(host.state.lock().bindings.is_empty());
2292 host.insert_binding(relay_binding.clone()).unwrap();
2293
2294 let virtual_turn_id = format!("local_turn_{}", "c".repeat(24));
2295 let operation_fingerprint = fingerprint("operation-state-fault");
2296 let request_fingerprint = fingerprint("request-state-fault");
2297 let native_turn_id = reserved_native_turn_id(
2298 &relay_binding.native_thread_id,
2299 &relay_binding.runtime_binding_id,
2300 &virtual_turn_id,
2301 &operation_fingerprint,
2302 );
2303 inject_state_persist_failures(&host.state_path, 1);
2304 assert!(
2305 host.reserve_turn(
2306 &relay_binding.runtime_binding_id,
2307 &relay_binding.virtual_thread_id,
2308 &virtual_turn_id,
2309 &native_turn_id,
2310 &operation_fingerprint,
2311 &request_fingerprint,
2312 )
2313 .is_err()
2314 );
2315 assert!(
2316 !host.state.lock().bindings[0]
2317 .turns
2318 .contains_key(&virtual_turn_id)
2319 );
2320 host.reserve_turn(
2321 &relay_binding.runtime_binding_id,
2322 &relay_binding.virtual_thread_id,
2323 &virtual_turn_id,
2324 &native_turn_id,
2325 &operation_fingerprint,
2326 &request_fingerprint,
2327 )
2328 .unwrap();
2329
2330 inject_state_persist_failures(&host.state_path, 1);
2331 assert!(
2332 host.advance_projection_cursor(&relay_binding.native_thread_id, 7, None)
2333 .is_err()
2334 );
2335 assert_eq!(host.state.lock().bindings[0].projected_native_seq, 0);
2336 host.advance_projection_cursor(&relay_binding.native_thread_id, 7, None)
2337 .unwrap();
2338
2339 drop(host);
2340 let reopened = open_host(root.path());
2341 let state = reopened.state.lock();
2342 assert!(state.owner_scope_fingerprint.is_some());
2343 assert_eq!(state.bindings.len(), 1);
2344 assert_eq!(state.bindings[0].projected_native_seq, 7);
2345 assert_eq!(
2346 state.bindings[0].turns[&virtual_turn_id].native_turn_id,
2347 native_turn_id
2348 );
2349 }
2350
2351 #[test]
2352 fn conflicting_replay_never_settles_a_live_turn_reservation() {
2353 let root = tempfile::tempdir().unwrap();
2354 let host = open_host(root.path());
2355 host.bind_account("account_fixture", "target_fixture")
2356 .unwrap();
2357 let mut relay_binding = binding();
2358 relay_binding.turns.clear();
2359 host.insert_binding(relay_binding.clone()).unwrap();
2360
2361 let virtual_turn_id = format!("local_turn_{}", "d".repeat(24));
2362 let operation_fingerprint = fingerprint("same-operation");
2363 let original_request = fingerprint("original-request");
2364 let changed_request = fingerprint("changed-request");
2365 let native_turn_id = reserved_native_turn_id(
2366 &relay_binding.native_thread_id,
2367 &relay_binding.runtime_binding_id,
2368 &virtual_turn_id,
2369 &operation_fingerprint,
2370 );
2371 assert_eq!(
2372 host.reserve_turn(
2373 &relay_binding.runtime_binding_id,
2374 &relay_binding.virtual_thread_id,
2375 &virtual_turn_id,
2376 &native_turn_id,
2377 &operation_fingerprint,
2378 &original_request,
2379 )
2380 .unwrap(),
2381 TurnReservationDisposition::New
2382 );
2383 assert!(
2384 host.reserve_turn(
2385 &relay_binding.runtime_binding_id,
2386 &relay_binding.virtual_thread_id,
2387 &virtual_turn_id,
2388 &native_turn_id,
2389 &operation_fingerprint,
2390 &changed_request,
2391 )
2392 .is_err()
2393 );
2394 assert!(host.has_any_unsettled_turns());
2395 assert!(!host.state.lock().bindings[0].turns[&virtual_turn_id].terminal_projected);
2396 }
2397
2398 #[test]
2399 fn rejected_start_retry_reopens_the_same_deterministic_reservation() {
2400 let root = tempfile::tempdir().unwrap();
2401 let host = open_host(root.path());
2402 host.bind_account("account_fixture", "target_fixture")
2403 .unwrap();
2404 let mut relay_binding = binding();
2405 relay_binding.turns.clear();
2406 host.insert_binding(relay_binding.clone()).unwrap();
2407
2408 let virtual_turn_id = format!("local_turn_{}", "e".repeat(24));
2409 let operation_fingerprint = fingerprint("retry-operation");
2410 let request_fingerprint = fingerprint("retry-request");
2411 let native_turn_id = reserved_native_turn_id(
2412 &relay_binding.native_thread_id,
2413 &relay_binding.runtime_binding_id,
2414 &virtual_turn_id,
2415 &operation_fingerprint,
2416 );
2417 assert_eq!(
2418 host.reserve_turn(
2419 &relay_binding.runtime_binding_id,
2420 &relay_binding.virtual_thread_id,
2421 &virtual_turn_id,
2422 &native_turn_id,
2423 &operation_fingerprint,
2424 &request_fingerprint,
2425 )
2426 .unwrap(),
2427 TurnReservationDisposition::New
2428 );
2429 host.finish_turn_reservation(
2430 &relay_binding.runtime_binding_id,
2431 &relay_binding.virtual_thread_id,
2432 &virtual_turn_id,
2433 )
2434 .unwrap();
2435 assert!(!host.has_any_unsettled_turns());
2436 assert_eq!(
2437 host.reserve_turn(
2438 &relay_binding.runtime_binding_id,
2439 &relay_binding.virtual_thread_id,
2440 &virtual_turn_id,
2441 &native_turn_id,
2442 &operation_fingerprint,
2443 &request_fingerprint,
2444 )
2445 .unwrap(),
2446 TurnReservationDisposition::Reopened
2447 );
2448 assert!(host.has_any_unsettled_turns());
2449 }
2450
2451 #[test]
2452 fn exact_replay_of_a_projected_terminal_turn_stays_settled() {
2453 let root = tempfile::tempdir().unwrap();
2454 let host = open_host(root.path());
2455 host.bind_account("account_fixture", "target_fixture")
2456 .unwrap();
2457 let mut relay_binding = binding();
2458 relay_binding.turns.clear();
2459 host.insert_binding(relay_binding.clone()).unwrap();
2460
2461 let virtual_turn_id = format!("local_turn_{}", "f".repeat(24));
2462 let operation_fingerprint = fingerprint("terminal-replay-operation");
2463 let request_fingerprint = fingerprint("terminal-replay-request");
2464 let native_turn_id = reserved_native_turn_id(
2465 &relay_binding.native_thread_id,
2466 &relay_binding.runtime_binding_id,
2467 &virtual_turn_id,
2468 &operation_fingerprint,
2469 );
2470 assert_eq!(
2471 host.reserve_turn(
2472 &relay_binding.runtime_binding_id,
2473 &relay_binding.virtual_thread_id,
2474 &virtual_turn_id,
2475 &native_turn_id,
2476 &operation_fingerprint,
2477 &request_fingerprint,
2478 )
2479 .unwrap(),
2480 TurnReservationDisposition::New
2481 );
2482 host.advance_projection_cursor(&relay_binding.native_thread_id, 11, Some(&virtual_turn_id))
2483 .unwrap();
2484 assert!(!host.has_any_unsettled_turns());
2485
2486 assert_eq!(
2487 host.reserve_turn(
2488 &relay_binding.runtime_binding_id,
2489 &relay_binding.virtual_thread_id,
2490 &virtual_turn_id,
2491 &native_turn_id,
2492 &operation_fingerprint,
2493 &request_fingerprint,
2494 )
2495 .unwrap(),
2496 TurnReservationDisposition::ExistingTerminal
2497 );
2498 assert!(
2499 !host.has_any_unsettled_turns(),
2500 "an idempotent replay of already-projected output cannot reopen provider work"
2501 );
2502 }
2503
2504 #[tokio::test]
2505 async fn restart_refuses_a_new_run_until_the_durable_old_turn_is_projected() {
2506 let root = tempfile::tempdir().unwrap();
2507 {
2508 let host = open_host(root.path());
2509 host.bind_account("account_fixture", "target_fixture")
2510 .unwrap();
2511 host.insert_binding(binding()).unwrap();
2512 host.authorize_run("run_fixture").unwrap();
2513 assert!(host.has_unsettled_authorized_turns());
2514 }
2515
2516 let reopened = open_host(root.path());
2517 reopened
2518 .bind_account("account_fixture", "target_fixture")
2519 .unwrap();
2520 assert!(reopened.authorize_run("run_other").is_err());
2521 reopened.authorize_run("run_fixture").unwrap();
2522 let relay_binding = reopened.state.lock().bindings[0].clone();
2523 let virtual_turn_id = relay_binding.turns.keys().next().unwrap().clone();
2524 reopened
2525 .mark_projected(
2526 &relay_binding.native_thread_id,
2527 9,
2528 &virtual_turn_id,
2529 "turn.completed",
2530 )
2531 .unwrap();
2532 assert!(!reopened.has_unsettled_authorized_turns());
2533 reopened.authorize_run("run_other").unwrap();
2534 }
2535 #[test]
2536 fn runtime_image_relay_hash_preserves_text_and_binds_order() {
2537 let legacy = json!({"type":"prompt.request","runId":"run_fixture","turnId":format!("local_turn_{}", "b".repeat(24)),"operationKey":"operation-1","runtimeBindingId":"binding_fixture","runtimeThreadId":format!("local_thread_{}", "a".repeat(24)),"prompt":"look","model":"deepseek-v4-flash-vision-exp","modelProvider":"deepseek","modelProviderId":"deepseek","allowedTools":[],"mode":"chat","requestedMode":"chat","workspace":{"id":"workspace_fixture","targetRef":"target_fixture"}});
2538 let mut command: RuntimeChatPrompt = serde_json::from_value(legacy.clone()).unwrap();
2539 command.validate_shape().unwrap();
2540 let expected = hex_digest(Sha256::digest(
2541 serde_json::to_vec(&canonical_json_value(&legacy)).unwrap(),
2542 ));
2543 assert_eq!(
2544 runtime_chat_request_fingerprint(&command).unwrap(),
2545 expected
2546 );
2547 let mut with_empty = legacy;
2548 with_empty["images"] = json!([]);
2549 let empty: RuntimeChatPrompt = serde_json::from_value(with_empty).unwrap();
2550 assert_eq!(runtime_chat_request_fingerprint(&empty).unwrap(), expected);
2551 command.images = vec![
2552 crate::image_attach::tests::runtime_image_fixture(1),
2553 crate::image_attach::tests::runtime_image_fixture(2),
2554 ];
2555 command.validate_shape().unwrap();
2556 let first = runtime_chat_request_fingerprint(&command).unwrap();
2557 command.images.reverse();
2558 assert_ne!(runtime_chat_request_fingerprint(&command).unwrap(), first);
2559 command.model = "auto".into();
2560 assert!(command.validate_shape().is_err());
2561 }
2562 }
2563
2563 lines RUST