返回 CodeWhale
thread_control.rs
根目录 / crates / app-server / src / thread_control.rs
1 //! Compatibility control projection onto the already held canonical owner.
2 //! SQLite supplies immutable migration input and committed alias receipts only.
3 use super::*;
4 use codewhale_protocol::{
5 CanonicalHistoryImportRequest, CanonicalHistoryOptions, CanonicalHistorySource,
6 CanonicalThreadMutation, CanonicalThreadMutationRequest, CanonicalThreadOperationKind,
7 CanonicalThreadOperationLookup, CanonicalThreadOperationRecovery,
8 CanonicalThreadOperationStatus, CanonicalThreadReceipt, CanonicalThreadSnapshot,
9 LegacyThreadHistory, MAX_CANONICAL_HISTORY_BYTES, MAX_CANONICAL_HISTORY_ENTRIES,
10 RuntimeOwnerReceipt, SessionSource, Thread, ThreadStatus,
11 };
12
13 #[derive(Debug, thiserror::Error)]
14 #[error("runtime API returned {status}: {detail}")]
15 struct HttpFailure {
16 status: StatusCode,
17 detail: String,
18 }
19
20 #[derive(Debug, thiserror::Error)]
21 #[error(
22 "canonical operation {operation} completed as thread {thread} / session {session}; inspect this result without replay: {source}"
23 )]
24 struct CommittedControlFailure {
25 operation: String,
26 thread: String,
27 session: String,
28 #[source]
29 source: anyhow::Error,
30 }
31
32 fn committed_failure(receipt: &CanonicalThreadReceipt, source: anyhow::Error) -> anyhow::Error {
33 CommittedControlFailure {
34 operation: receipt.operation_key.clone(),
35 thread: receipt.runtime_thread_id.clone(),
36 session: receipt.session_id.clone(),
37 source,
38 }
39 .into()
40 }
41
42 fn owner(state: &AppState) -> Result<RuntimeOwnerReceipt> {
43 state.captured_owner.clone().context(
44 "thread controls require the authenticated canonical owner; no standalone history writer",
45 )
46 }
47
48 fn endpoint(bridge: &RuntimeBridge, segments: &[&str]) -> Result<reqwest::Url> {
49 let mut url = reqwest::Url::parse(&bridge.base_url)?;
50 url.path_segments_mut()
51 .map_err(|_| anyhow!("invalid canonical owner URL"))?
52 .clear()
53 .extend(segments.iter().copied());
54 Ok(url)
55 }
56
57 /// Reused by every bridge JSON read, including full canonical history. The
58 /// declared length and actual streamed bytes must both fit; nothing is cut.
59 pub(super) async fn read_json_response(mut response: reqwest::Response) -> Result<Value> {
60 let status = response.status();
61 if response
62 .content_length()
63 .is_some_and(|len| len > MAX_CANONICAL_HISTORY_BYTES as u64)
64 {
65 bail!("canonical response exceeds complete-document bound");
66 }
67 let mut bytes = Vec::new();
68 while let Some(chunk) = response.chunk().await? {
69 if bytes
70 .len()
71 .checked_add(chunk.len())
72 .is_none_or(|len| len > MAX_CANONICAL_HISTORY_BYTES)
73 {
74 bail!("canonical response exceeds complete-document bound");
75 }
76 bytes.extend_from_slice(&chunk);
77 }
78 if !status.is_success() {
79 let detail = String::from_utf8_lossy(&bytes);
80 return Err(HttpFailure {
81 status,
82 detail: detail.trim().to_owned(),
83 }
84 .into());
85 }
86 if bytes.is_empty() && status == StatusCode::NO_CONTENT {
87 return Ok(Value::Null);
88 }
89 serde_json::from_slice(&bytes).context("invalid canonical Runtime JSON")
90 }
91
92 async fn request(
93 bridge: &RuntimeBridge,
94 method: Method,
95 segments: &[&str],
96 body: Option<&Value>,
97 ) -> Result<Value> {
98 let mut request = bridge.authed(bridge.client.request(method, endpoint(bridge, segments)?));
99 if let Some(body) = body {
100 anyhow::ensure!(
101 serde_json::to_vec(body)?.len() <= MAX_CANONICAL_HISTORY_BYTES,
102 "encoded canonical control exceeds bound"
103 );
104 request = request.json(body);
105 }
106 tokio::time::timeout(Duration::from_secs(30), bridge.request_json(request))
107 .await
108 .context(
109 "canonical control deadline expired; outcome may be committed, no automatic replay",
110 )?
111 }
112
113 #[cfg(any(unix, windows))]
114 async fn store_work<T, F>(work: F) -> Result<T>
115 where
116 T: Send + 'static,
117 F: FnOnce() -> Result<T> + Send + 'static,
118 {
119 daemon_socket::owner_work(work).await
120 }
121 #[cfg(not(any(unix, windows)))]
122 async fn store_work<T, F>(_work: F) -> Result<T>
123 where
124 T: Send + 'static,
125 F: FnOnce() -> Result<T> + Send + 'static,
126 {
127 bail!("canonical owner attachment is unsupported on this platform")
128 }
129
130 struct LegacySource {
131 metadata: codewhale_state::ThreadMetadata,
132 previous: Option<String>,
133 receipt: Option<CanonicalThreadReceipt>,
134 history: Option<LegacyThreadHistory>,
135 }
136
137 async fn legacy_source(
138 state: &AppState,
139 key: &str,
140 owner: &RuntimeOwnerReceipt,
141 ) -> Result<Option<(StateStore, LegacySource)>> {
142 let store = state.runtime.read().await.state_store().clone();
143 let read_store = store.clone();
144 let key = key.to_owned();
145 let owner = owner.clone();
146 let source = store_work(move || {
147 let Some(metadata) = read_store.get_thread(&key)? else {
148 return Ok(None);
149 };
150 let previous = read_store.get_runtime_thread_link(&key)?;
151 let receipt = read_store.get_canonical_runtime_link(&key, &owner)?;
152 let history = if receipt.is_none() {
153 Some(read_store.snapshot_legacy_thread_history(&key)?)
154 } else {
155 None
156 };
157 Ok(Some(LegacySource {
158 metadata,
159 previous,
160 receipt,
161 history,
162 }))
163 })
164 .await?;
165 Ok(source.map(|source| (store, source)))
166 }
167
168 fn migration_key(history: &LegacyThreadHistory) -> Result<String> {
169 // JSON tuple delimiters make this an exact source identity even when a
170 // legacy ID contains punctuation. Overlong identities visibly refuse.
171 let key = format!(
172 "legacy-import:{}",
173 serde_json::to_string(&(&history.state_store_id, &history.thread_id))?
174 );
175 anyhow::ensure!(
176 key.len() <= 128 && !key.chars().any(char::is_control),
177 "legacy operation identity exceeds the owner's bound; source retained for recovery"
178 );
179 Ok(key)
180 }
181
182 fn verify_receipt(
183 receipt: &CanonicalThreadReceipt,
184 owner: &RuntimeOwnerReceipt,
185 operation: &str,
186 ) -> Result<()> {
187 anyhow::ensure!(
188 receipt.version == 1
189 && receipt.data_dir == owner.data_dir
190 && receipt.execution_scope == owner.execution_scope
191 && receipt.operation_key == operation,
192 "canonical import receipt does not match the selected owner operation"
193 );
194 Ok(())
195 }
196
197 fn check_workspace(state: &AppState, workspace: &Path) -> Result<()> {
198 if let Some(selected) = &state.frontend_workspace {
199 anyhow::ensure!(
200 workspace == selected,
201 "thread belongs to another selected frontend workspace; attach its owner scope"
202 );
203 }
204 Ok(())
205 }
206
207 /// Holds the existing bridge serialization through import and exact source
208 /// CAS. A cancelled waiter does not cancel publication of a committed result.
209 pub(super) async fn resolve(
210 state: &AppState,
211 key: &str,
212 execution: bool,
213 ) -> Result<(String, PathBuf)> {
214 let state = state.clone();
215 let key = key.to_owned();
216 anyhow::ensure!(
217 !key.is_empty() && key.len() <= 1024 && !key.chars().any(char::is_control),
218 "invalid thread identity"
219 );
220 tokio::spawn(async move { resolve_owned(&state, &key, execution).await })
221 .await
222 .context("canonical resolution task failed")?
223 }
224
225 async fn resolve_owned(state: &AppState, key: &str, execution: bool) -> Result<(String, PathBuf)> {
226 let owner = owner(state)?;
227 let source = legacy_source(state, key, &owner).await?;
228 let bridge = acquire_live_runtime_bridge(state)
229 .await
230 .map_err(|e| anyhow!("{}", e.message))?;
231 let id = if let Some((store, source)) = source {
232 if execution || source.receipt.is_none() {
233 check_workspace(state, &source.metadata.cwd)?;
234 }
235 if let Some(receipt) = source.receipt {
236 receipt.runtime_thread_id
237 } else {
238 // Existing targets retain their canonical active branch; the
239 // owner admits complete legacy branches without replacing work.
240 let history = source.history.context("legacy source snapshot missing")?;
241 let operation = migration_key(&history)?;
242 let request_body = CanonicalHistoryImportRequest {
243 version: 1,
244 target_runtime_thread_id: source.previous.clone(),
245 operation_key: operation.clone(),
246 expected_data_dir: owner.data_dir.clone(),
247 expected_execution_scope: owner.execution_scope.clone(),
248 workspace: source.metadata.cwd,
249 model: None,
250 history: history.clone(),
251 };
252 let value = request(&bridge, Method::POST, &["v1", "thread-history", "import"], Some(&serde_json::to_value(request_body)?)).await
253 .with_context(|| format!("canonical import {operation} may have committed; inspect or retry the same operation, never remint"))?;
254 let receipt: CanonicalThreadReceipt = serde_json::from_value(value)?;
255 verify_receipt(&receipt, &owner, &operation)?;
256 let publication = receipt.clone();
257 let thread_key = key.to_owned();
258 let previous = source.previous;
259 store_work(move || store.publish_canonical_runtime_link(&thread_key, previous.as_deref(), &history, &owner, &publication)).await
260 .context("canonical result committed but compatibility publication failed; result and source retained for recovery")?;
261 receipt.runtime_thread_id
262 }
263 } else {
264 key.to_owned()
265 };
266 let value = request(&bridge, Method::GET, &["v1", "threads", &id], None)
267 .await
268 .context(
269 "canonical target is unavailable; existing identity retained, no replacement thread",
270 )?;
271 let record = value.get("thread").unwrap_or(&value);
272 anyhow::ensure!(
273 record.get("id").and_then(Value::as_str) == Some(id.as_str()),
274 "canonical resolution returned another target identity"
275 );
276 let thread = project_thread(record, key)?;
277 if execution {
278 check_workspace(state, &thread.cwd)?;
279 }
280 state
281 .runtime_thread_map
282 .lock()
283 .await
284 .insert(key.to_owned(), id.clone());
285 Ok((id, thread.cwd))
286 }
287
288 fn timestamp(record: &Value, key: &str) -> Result<i64> {
289 Ok(chrono::DateTime::parse_from_rfc3339(
290 record
291 .get(key)
292 .and_then(Value::as_str)
293 .context("canonical timestamp missing")?,
294 )?
295 .timestamp())
296 }
297
298 fn project_thread(record: &Value, public_id: &str) -> Result<Thread> {
299 let id = record
300 .get("id")
301 .and_then(Value::as_str)
302 .context("canonical thread identity missing")?;
303 anyhow::ensure!(!id.is_empty(), "canonical thread identity missing");
304 Ok(Thread {
305 id: public_id.to_owned(),
306 preview: String::new(),
307 ephemeral: false,
308 model_provider: record
309 .get("model_provider_id")
310 .or_else(|| record.get("model_provider"))
311 .and_then(Value::as_str)
312 .context("canonical provider identity missing")?
313 .to_owned(),
314 created_at: timestamp(record, "created_at")?,
315 updated_at: timestamp(record, "updated_at")?,
316 status: if record.get("archived").and_then(Value::as_bool) == Some(true) {
317 ThreadStatus::Archived
318 } else {
319 ThreadStatus::Idle
320 },
321 path: None,
322 cwd: serde_json::from_value(
323 record
324 .get("workspace")
325 .cloned()
326 .context("canonical workspace missing")?,
327 )?,
328 cli_version: env!("CARGO_PKG_VERSION").to_owned(),
329 source: SessionSource::Api,
330 name: record
331 .get("title")
332 .and_then(Value::as_str)
333 .map(str::to_owned),
334 })
335 }
336
337 fn response(id: String) -> ThreadResponse {
338 ThreadResponse {
339 thread_id: id,
340 status: "ok".into(),
341 thread: None,
342 threads: Vec::new(),
343 goal: None,
344 model: None,
345 model_provider: None,
346 cwd: None,
347 approval_policy: None,
348 sandbox: None,
349 events: Vec::new(),
350 data: json!({}),
351 }
352 }
353
354 async fn list(
355 state: &AppState,
356 params: codewhale_protocol::ThreadListParams,
357 ) -> Result<ThreadResponse> {
358 let owner = owner(state)?;
359 let store = state.runtime.read().await.state_store().clone();
360 let legacy = store_work(move || {
361 let rows = store.list_threads(codewhale_state::ThreadListFilters {
362 include_archived: true,
363 limit: Some(MAX_CANONICAL_HISTORY_ENTRIES + 1),
364 })?;
365 anyhow::ensure!(
366 rows.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
367 "legacy metadata list exceeds bound"
368 );
369 rows.into_iter()
370 .map(|row| {
371 let receipt = store.get_canonical_runtime_link(&row.id, &owner)?;
372 Ok((row, receipt))
373 })
374 .collect::<Result<Vec<_>>>()
375 })
376 .await?;
377 let limit = params.limit.unwrap_or(50);
378 anyhow::ensure!(
379 limit <= MAX_CANONICAL_HISTORY_ENTRIES,
380 "thread list exceeds bound"
381 );
382 let bridge = acquire_live_runtime_bridge(state)
383 .await
384 .map_err(|e| anyhow!("{}", e.message))?;
385 let mut url = endpoint(&bridge, &["v1", "threads"])?;
386 url.query_pairs_mut()
387 .append_pair("include_archived", "true")
388 .append_pair("limit", &(MAX_CANONICAL_HISTORY_ENTRIES + 1).to_string());
389 let rows = bridge
390 .request_json(bridge.authed(bridge.client.get(url)))
391 .await?;
392 let rows = rows
393 .as_array()
394 .context("canonical thread list is not an array")?;
395 anyhow::ensure!(
396 rows.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
397 "canonical thread list exceeds bound"
398 );
399 let mut canonical = HashMap::new();
400 let active = running_ids(&bridge).await?;
401 for row in rows {
402 let id = row
403 .get("id")
404 .and_then(Value::as_str)
405 .context("canonical list identity missing")?;
406 anyhow::ensure!(
407 canonical.insert(id.to_owned(), row).is_none(),
408 "duplicate canonical list identity"
409 );
410 }
411 let mut result = response("list".into());
412 let mut aliased = std::collections::HashSet::new();
413 for (metadata, receipt) in legacy {
414 if let Some(receipt) = receipt {
415 let row = canonical
416 .get(&receipt.runtime_thread_id)
417 .context("bound compatibility alias has no canonical record; recovery required")?;
418 let mut thread = project_thread(row, &metadata.id)?;
419 if active.contains(&receipt.runtime_thread_id)
420 && thread.status != ThreadStatus::Archived
421 {
422 thread.status = ThreadStatus::Running;
423 }
424 thread.preview = metadata.preview;
425 thread.source = serde_json::from_value(serde_json::to_value(metadata.source)?)?;
426 if params.include_archived || thread.status != ThreadStatus::Archived {
427 result.threads.push(thread);
428 }
429 aliased.insert(receipt.runtime_thread_id);
430 } else {
431 // Read projection only: listing never imports or mutates history.
432 if params.include_archived || !metadata.archived {
433 result
434 .threads
435 .push(serde_json::from_value(serde_json::to_value(metadata)?)?);
436 }
437 }
438 }
439 for (id, row) in canonical {
440 if !aliased.contains(&id) {
441 let mut thread = project_thread(row, &id)?;
442 if active.contains(&id) && thread.status != ThreadStatus::Archived {
443 thread.status = ThreadStatus::Running;
444 }
445 if params.include_archived || thread.status != ThreadStatus::Archived {
446 result.threads.push(thread);
447 }
448 }
449 }
450 anyhow::ensure!(
451 result.threads.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
452 "combined thread list exceeds bound"
453 );
454 result.threads.sort_by(|a, b| {
455 b.updated_at
456 .cmp(&a.updated_at)
457 .then_with(|| a.id.cmp(&b.id))
458 });
459 result.threads.truncate(limit); // Explicit caller list limit; history is never truncated.
460 Ok(result)
461 }
462
463 pub(super) async fn handle(
464 state: &AppState,
465 request_value: ThreadRequest,
466 ) -> std::result::Result<ThreadResponse, JsonRpcError> {
467 let key = match &request_value {
468 ThreadRequest::Read(p) => Some(p.thread_id.as_str()),
469 ThreadRequest::SetName(p) => Some(p.thread_id.as_str()),
470 ThreadRequest::Archive { thread_id } | ThreadRequest::Unarchive { thread_id } => {
471 Some(thread_id.as_str())
472 }
473 ThreadRequest::GoalGet(p) => Some(p.thread_id.as_str()),
474 ThreadRequest::GoalSet(p) => Some(p.thread_id.as_str()),
475 ThreadRequest::GoalClear(p) => Some(p.thread_id.as_str()),
476 ThreadRequest::Resume(p) => Some(p.thread_id.as_str()),
477 ThreadRequest::Fork(p) => Some(p.thread_id.as_str()),
478 ThreadRequest::GoalRecordProgress(p) => Some(p.thread_id.as_str()),
479 ThreadRequest::Message { thread_id, .. } => Some(thread_id.as_str()),
480 ThreadRequest::Create { .. } | ThreadRequest::Start(_) | ThreadRequest::List(_) => None,
481 }
482 .map(str::to_owned);
483 let result = handle_owned(state, request_value)
484 .await
485 .map_err(|error| rpc_error(error, key.as_deref()))?;
486 if serde_json::to_vec(&result)
487 .map_err(|e| JsonRpcError::internal(e.to_string()))?
488 .len()
489 > MAX_CANONICAL_HISTORY_BYTES
490 {
491 let message = "complete thread response exceeds transport bound; no truncation";
492 let error = result
493 .data
494 .get("receipt")
495 .and_then(|value| serde_json::from_value::<CanonicalThreadReceipt>(value.clone()).ok())
496 .map(|receipt| committed_failure(&receipt, anyhow!(message)))
497 .unwrap_or_else(|| anyhow!(message));
498 return Err(JsonRpcError::runtime_unavailable(format!("{error:#}")));
499 }
500 Ok(result)
501 }
502
503 pub(super) fn rpc_error(error: anyhow::Error, key: Option<&str>) -> JsonRpcError {
504 if error.downcast_ref::<CommittedControlFailure>().is_some() {
505 JsonRpcError::runtime_unavailable(format!("{error:#}"))
506 } else if error
507 .downcast_ref::<HttpFailure>()
508 .is_some_and(|error| error.status == StatusCode::NOT_FOUND)
509 {
510 JsonRpcError::thread_not_found(key.unwrap_or("unknown"))
511 } else {
512 JsonRpcError::runtime_unavailable(format!("{error:#}"))
513 }
514 }
515
516 async fn running_ids(bridge: &RuntimeBridge) -> Result<std::collections::HashSet<String>> {
517 let value = request(bridge, Method::GET, &["v1", "threads", "running"], None).await?;
518 let rows = value
519 .as_array()
520 .context("canonical running list is not an array")?;
521 anyhow::ensure!(
522 rows.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
523 "canonical running list exceeds bound"
524 );
525 rows.iter()
526 .map(|row| {
527 row.get("thread_id")
528 .and_then(Value::as_str)
529 .map(str::to_owned)
530 .context("canonical running identity missing")
531 })
532 .collect()
533 }
534
535 fn operation_key(key: Option<String>) -> Result<String> {
536 let key = key.context("canonical mutation requires a client-captured operation_key; retain it when retrying an uncertain outcome")?;
537 anyhow::ensure!(
538 !key.is_empty() && key.len() <= 128 && !key.chars().any(char::is_control),
539 "invalid canonical operation_key"
540 );
541 Ok(key)
542 }
543
544 fn selected_workspace(state: &AppState, explicit: Option<PathBuf>) -> Result<PathBuf> {
545 let workspace = explicit.or_else(|| state.frontend_workspace.clone()).context(
546 "canonical mutation needs the owner's acknowledged workspace or an explicit checked selection",
547 )?;
548 check_workspace(state, &workspace)?;
549 Ok(workspace)
550 }
551
552 async fn full_snapshot(
553 state: &AppState,
554 bridge: &RuntimeBridge,
555 id: &str,
556 ) -> Result<CanonicalThreadSnapshot> {
557 let snapshot: CanonicalThreadSnapshot = serde_json::from_value(
558 request(bridge, Method::GET, &["v1", "threads", id, "history"], None).await?,
559 )?;
560 let selected = owner(state)?;
561 anyhow::ensure!(
562 snapshot.version == 1
563 && snapshot.data_dir == selected.data_dir
564 && snapshot.execution_scope == selected.execution_scope
565 && snapshot.runtime_thread_id == id,
566 "canonical full history response has another owner/thread binding"
567 );
568 Ok(snapshot)
569 }
570
571 async fn mutate(
572 state: &AppState,
573 operation: String,
574 workspace: PathBuf,
575 mutation: CanonicalThreadMutation,
576 status: &str,
577 public_alias: Option<&str>,
578 ) -> Result<ThreadResponse> {
579 let owner = owner(state)?;
580 check_workspace(state, &workspace)?;
581 let bridge = acquire_live_runtime_bridge(state)
582 .await
583 .map_err(|error| anyhow!("{}", error.message))?;
584 let body = CanonicalThreadMutationRequest {
585 version: 1,
586 operation_key: operation.clone(),
587 expected_data_dir: owner.data_dir.clone(),
588 expected_execution_scope: owner.execution_scope.clone(),
589 workspace,
590 mutation,
591 };
592 let value = request(
593 &bridge,
594 Method::POST,
595 &["v1", "thread-history", "mutate"],
596 Some(&serde_json::to_value(body)?),
597 )
598 .await
599 .with_context(|| format!("canonical operation {operation} may have committed; retain this key, no automatic replay or replacement"))?;
600 let receipt: CanonicalThreadReceipt = serde_json::from_value(value).with_context(|| {
601 format!("canonical operation {operation} returned an invalid result; retain this key, no replacement")
602 })?;
603 verify_receipt(&receipt, &owner, &operation).with_context(|| {
604 format!("canonical operation {operation} returned an unbound receipt; retain this key without replay or replacement")
605 })?;
606 present_receipt(&bridge, receipt, status, public_alias).await
607 }
608
609 async fn present_receipt(
610 bridge: &RuntimeBridge,
611 receipt: CanonicalThreadReceipt,
612 status: &str,
613 public_alias: Option<&str>,
614 ) -> Result<ThreadResponse> {
615 let public_id = public_alias.unwrap_or(&receipt.runtime_thread_id);
616 let value = request(
617 bridge,
618 Method::GET,
619 &["v1", "threads", &receipt.runtime_thread_id],
620 None,
621 )
622 .await
623 .map_err(|error| committed_failure(&receipt, error.context("metadata is unavailable")))?;
624 let record = value.get("thread").unwrap_or(&value);
625 if record.get("id").and_then(Value::as_str) != Some(receipt.runtime_thread_id.as_str()) {
626 return Err(committed_failure(
627 &receipt,
628 anyhow!("another metadata identity was returned"),
629 ));
630 }
631 let mut result = response(public_id.to_owned());
632 result.status = status.to_owned();
633 result.thread = Some(project_thread(record, public_id).map_err(|error| {
634 committed_failure(&receipt, error.context("metadata projection failed"))
635 })?);
636 result.model = record
637 .get("model")
638 .and_then(Value::as_str)
639 .map(str::to_owned);
640 result.model_provider = result
641 .thread
642 .as_ref()
643 .map(|thread| thread.model_provider.clone());
644 result.cwd = result.thread.as_ref().map(|thread| thread.cwd.clone());
645 result.data = json!({"receipt":receipt,"thread":record});
646 Ok(result)
647 }
648
649 async fn source_identity(state: &AppState, key: &str) -> Result<String> {
650 let selected = owner(state)?;
651 let store = state.runtime.read().await.state_store().clone();
652 let key = key.to_owned();
653 store_work(move || {
654 if let Some(receipt) = store.get_canonical_runtime_link(&key, &selected)? {
655 Ok(receipt.runtime_thread_id)
656 } else {
657 Ok(store.get_runtime_thread_link(&key)?.unwrap_or(key))
658 }
659 })
660 .await
661 }
662
663 async fn recover_operation(
664 state: &AppState,
665 operation: &str,
666 workspace: &Path,
667 kind: CanonicalThreadOperationKind,
668 source: Option<&str>,
669 status: &str,
670 public_alias: Option<&str>,
671 ) -> Result<Option<ThreadResponse>> {
672 let selected = owner(state)?;
673 let bridge = acquire_live_runtime_bridge(state)
674 .await
675 .map_err(|error| anyhow!("{}", error.message))?;
676 let body = CanonicalThreadOperationLookup {
677 version: 1,
678 operation_key: operation.to_owned(),
679 expected_data_dir: selected.data_dir.clone(),
680 expected_execution_scope: selected.execution_scope.clone(),
681 workspace: workspace.to_owned(),
682 };
683 let outcome: CanonicalThreadOperationStatus = serde_json::from_value(request(
684 &bridge, Method::POST, &["v1","thread-history","operations","lookup"],
685 Some(&serde_json::to_value(&body)?),
686 ).await.with_context(|| format!("canonical key lookup {operation} is unavailable; retain key, no absence inference or replacement"))?)
687 .with_context(|| format!("canonical key lookup {operation} returned an invalid result; retain key, no absence inference or replacement"))?;
688 let (receipt, association, committed) = match outcome {
689 CanonicalThreadOperationStatus::Absent => return Ok(None),
690 CanonicalThreadOperationStatus::Pending {
691 receipt,
692 association,
693 } => (receipt, association, false),
694 CanonicalThreadOperationStatus::Committed {
695 receipt,
696 association,
697 } => (receipt, association, true),
698 };
699 verify_receipt(&receipt, &selected, operation)?;
700 if association.kind != kind || association.source_runtime_thread_id.as_deref() != source {
701 let error = anyhow!("retained key belongs to another action or source; no new operation");
702 return Err(if committed {
703 committed_failure(&receipt, error)
704 } else {
705 error.context(format!(
706 "canonical operation {operation} is pending; inspect reserved thread {}, no replay",
707 receipt.runtime_thread_id
708 ))
709 });
710 }
711 if !committed {
712 // Lookup is read-only. Explicit retained-key recovery can settle only
713 // the owner's already prepared target, never reconstruct this source.
714 let recovery = CanonicalThreadOperationRecovery {
715 operation: body,
716 association: association.clone(),
717 };
718 let recovered: CanonicalThreadOperationStatus = serde_json::from_value(
719 request(
720 &bridge,
721 Method::POST,
722 &["v1", "thread-history", "operations", "recover"],
723 Some(&serde_json::to_value(recovery)?),
724 )
725 .await
726 .with_context(|| format!(
727 "canonical operation {operation} remains uncertain as thread {} / session {}; retain key, no resume, replay or replacement",
728 receipt.runtime_thread_id, receipt.session_id
729 ))?,
730 ).with_context(|| format!(
731 "canonical operation {operation} returned an invalid recovery result for reserved thread {} / session {}; retain key, no replacement",
732 receipt.runtime_thread_id, receipt.session_id
733 ))?;
734 match recovered {
735 CanonicalThreadOperationStatus::Committed {
736 receipt: recovered_receipt,
737 association: recovered_association,
738 } => {
739 anyhow::ensure!(
740 recovered_receipt == receipt && recovered_association == association,
741 "canonical recovery changed retained operation {operation}; reserved thread {} / session {}, no replay",
742 receipt.runtime_thread_id,
743 receipt.session_id
744 );
745 }
746 CanonicalThreadOperationStatus::Pending {
747 receipt: pending_receipt,
748 association: pending_association,
749 } => {
750 anyhow::ensure!(
751 pending_receipt == receipt && pending_association == association,
752 "canonical recovery changed retained operation {operation}; no replay"
753 );
754 bail!(
755 "canonical operation {operation} is pending as thread {} / session {}; no resume, replay or replacement",
756 receipt.runtime_thread_id,
757 receipt.session_id
758 );
759 }
760 CanonicalThreadOperationStatus::Absent => bail!(
761 "canonical retained operation {operation} disappeared as thread {} / session {}; no resume, replay or replacement",
762 receipt.runtime_thread_id,
763 receipt.session_id
764 ),
765 }
766 }
767 present_receipt(&bridge, receipt, status, public_alias)
768 .await
769 .map(Some)
770 }
771
772 async fn resume_or_fork(
773 state: &AppState,
774 key: &str,
775 operation: String,
776 cwd: Option<PathBuf>,
777 fork: bool,
778 options: CanonicalHistoryOptions,
779 ) -> Result<ThreadResponse> {
780 let selected_workspace = selected_workspace(state, cwd.clone())?;
781 let expected_source = source_identity(state, key).await?;
782 if let Some(result) = recover_operation(
783 state,
784 &operation,
785 &selected_workspace,
786 if fork {
787 CanonicalThreadOperationKind::Fork
788 } else {
789 CanonicalThreadOperationKind::Resume
790 },
791 Some(&expected_source),
792 if fork { "forked" } else { "resumed" },
793 (!fork).then_some(key),
794 )
795 .await?
796 {
797 return Ok(result);
798 }
799 let (id, workspace) = resolve(state, key, !fork).await?;
800 if let Some(cwd) = cwd {
801 anyhow::ensure!(
802 fork || cwd == workspace,
803 "moving a history source to another workspace requires the canonical owner's explicit admission"
804 );
805 }
806 let bridge = acquire_live_runtime_bridge(state)
807 .await
808 .map_err(|error| anyhow!("{}", error.message))?;
809 let snapshot = full_snapshot(state, &bridge, &id).await?;
810 drop(bridge);
811 let source = CanonicalHistorySource::Thread {
812 runtime_thread_id: id,
813 expected_document_digest: snapshot.document_digest,
814 };
815 let mutation = if fork {
816 CanonicalThreadMutation::Fork {
817 source,
818 options,
819 selected_entry_id: None,
820 }
821 } else {
822 CanonicalThreadMutation::Resume { source, options }
823 };
824 // Only resume preserves the old public alias. Fork returns its new actual
825 // canonical identity and never creates another SQLite transcript writer.
826 mutate(
827 state,
828 operation,
829 if fork { selected_workspace } else { workspace },
830 mutation,
831 if fork { "forked" } else { "resumed" },
832 (!fork).then_some(key),
833 )
834 .await
835 }
836
837 fn history_parameters(
838 params: Value,
839 ) -> Result<(String, String, Option<PathBuf>, CanonicalHistoryOptions)> {
840 let mut fields = params
841 .as_object()
842 .cloned()
843 .context("invalid typed history control")?;
844 let key = serde_json::from_value(
845 fields
846 .remove("thread_id")
847 .context("thread identity missing")?,
848 )?;
849 let operation = operation_key(
850 fields
851 .remove("operation_key")
852 .map(serde_json::from_value)
853 .transpose()?,
854 )?;
855 let cwd = fields
856 .get("cwd")
857 .cloned()
858 .map(serde_json::from_value)
859 .transpose()?;
860 let source_path = fields
861 .remove("path")
862 .map(serde_json::from_value)
863 .transpose()?;
864 let offered_history = fields
865 .remove("history")
866 .map(serde_json::from_value)
867 .transpose()?
868 .unwrap_or_default();
869 // Complete canonical history is durable for both values of the old flag.
870 fields.remove("persist_extended_history");
871 Ok((
872 key,
873 operation,
874 cwd,
875 CanonicalHistoryOptions {
876 offered_history,
877 overrides: Value::Object(fields),
878 source_path,
879 expected_session_goal_digest: None,
880 },
881 ))
882 }
883
884 async fn creation(state: &AppState, request_value: ThreadRequest) -> Result<ThreadResponse> {
885 match request_value {
886 ThreadRequest::Create { metadata } => {
887 let mut config = metadata.as_object().cloned().context(
888 "thread/create metadata must carry an operation_key and existing create-thread fields",
889 )?;
890 let operation = operation_key(
891 config
892 .remove("operation_key")
893 .map(serde_json::from_value)
894 .transpose()?,
895 )?;
896 let explicit = config
897 .get("workspace")
898 .cloned()
899 .map(serde_json::from_value)
900 .transpose()?;
901 let workspace = selected_workspace(state, explicit)?;
902 if let Some(result) = recover_operation(
903 state,
904 &operation,
905 &workspace,
906 CanonicalThreadOperationKind::Create,
907 None,
908 "created",
909 None,
910 )
911 .await?
912 {
913 return Ok(result);
914 }
915 config.insert("workspace".into(), serde_json::to_value(&workspace)?);
916 mutate(
917 state,
918 operation,
919 workspace,
920 CanonicalThreadMutation::Create {
921 config: Value::Object(config),
922 },
923 "created",
924 None,
925 )
926 .await
927 }
928 ThreadRequest::Start(params) => {
929 let operation = operation_key(params.operation_key)?;
930 let workspace = selected_workspace(state, params.cwd)?;
931 if let Some(result) = recover_operation(
932 state,
933 &operation,
934 &workspace,
935 CanonicalThreadOperationKind::Create,
936 None,
937 "started",
938 None,
939 )
940 .await?
941 {
942 return Ok(result);
943 }
944 let mut config = json!({"workspace":workspace});
945 if let Some(model) = params.model {
946 config["model"] = json!(model);
947 }
948 if let Some(provider) = params.model_provider {
949 config["model_provider"] = json!(provider);
950 }
951 // Canonical history is always complete and durable; the legacy
952 // extended-history flag cannot turn off branches or receipts.
953 mutate(
954 state,
955 operation,
956 workspace,
957 CanonicalThreadMutation::Create { config },
958 "started",
959 None,
960 )
961 .await
962 }
963 ThreadRequest::Resume(params) => {
964 let (key, operation, cwd, options) = history_parameters(serde_json::to_value(params)?)?;
965 resume_or_fork(state, &key, operation, cwd, false, options).await
966 }
967 ThreadRequest::Fork(params) => {
968 let (key, operation, cwd, options) = history_parameters(serde_json::to_value(params)?)?;
969 resume_or_fork(state, &key, operation, cwd, true, options).await
970 }
971 _ => unreachable!("closed creation request dispatch"),
972 }
973 }
974
975 async fn handle_owned(state: &AppState, request_value: ThreadRequest) -> Result<ThreadResponse> {
976 if let ThreadRequest::List(params) = request_value {
977 return list(state, params).await;
978 }
979 if matches!(
980 &request_value,
981 ThreadRequest::Create { .. }
982 | ThreadRequest::Start(_)
983 | ThreadRequest::Resume(_)
984 | ThreadRequest::Fork(_)
985 ) {
986 return creation(state, request_value).await;
987 }
988 let (key, method, action, body) = match request_value {
989 ThreadRequest::Read(p) => (p.thread_id, Method::GET, None, None),
990 ThreadRequest::SetName(p) => (
991 p.thread_id,
992 Method::PATCH,
993 None,
994 Some(json!({"title":p.name})),
995 ),
996 ThreadRequest::Archive { thread_id } => (
997 thread_id,
998 Method::PATCH,
999 None,
1000 Some(json!({"archived":true})),
1001 ),
1002 ThreadRequest::Unarchive { thread_id } => (
1003 thread_id,
1004 Method::PATCH,
1005 None,
1006 Some(json!({"archived":false})),
1007 ),
1008 ThreadRequest::GoalGet(p) => (p.thread_id, Method::GET, Some("goal"), None),
1009 ThreadRequest::GoalSet(p) => (
1010 p.thread_id,
1011 Method::PUT,
1012 Some("goal"),
1013 Some(json!({"objective":p.objective,"token_budget":p.token_budget})),
1014 ),
1015 ThreadRequest::GoalClear(p) => (p.thread_id, Method::DELETE, Some("goal"), None),
1016 ThreadRequest::GoalRecordProgress(_) => bail!(
1017 "goal progress requires producing Engine usage/time receipts; client deltas are not accounting authority"
1018 ),
1019 ThreadRequest::Create { .. }
1020 | ThreadRequest::Start(_)
1021 | ThreadRequest::Resume(_)
1022 | ThreadRequest::Fork(_) => unreachable!("creation handled above"),
1023 ThreadRequest::Message { .. } => {
1024 bail!("thread messages use the existing canonical turn transport")
1025 }
1026 ThreadRequest::List(_) => unreachable!("list handled above"),
1027 };
1028 let (id, _) = resolve(state, &key, false).await?;
1029 let bridge = acquire_live_runtime_bridge(state)
1030 .await
1031 .map_err(|e| anyhow!("{}", e.message))?;
1032 let mut result = response(key.clone());
1033 let mut segments = vec!["v1", "threads", &id];
1034 if let Some(action) = action {
1035 segments.push(action)
1036 }
1037 let value = match request(&bridge, method.clone(), &segments, body.as_ref()).await {
1038 Ok(value) => value,
1039 Err(error)
1040 if action == Some("goal")
1041 && (method == Method::GET || method == Method::DELETE)
1042 && error
1043 .downcast_ref::<HttpFailure>()
1044 .is_some_and(|e| e.status == StatusCode::NOT_FOUND) =>
1045 {
1046 // Distinguish an absent goal from a concurrently removed thread.
1047 request(&bridge, Method::GET, &["v1", "threads", &id], None).await?;
1048 if method == Method::DELETE {
1049 result.status = "empty".into();
1050 }
1051 Value::Null
1052 }
1053 Err(error) => return Err(error),
1054 };
1055 if action == Some("goal") {
1056 if method == Method::DELETE && result.status != "empty" {
1057 result.status = "cleared".into();
1058 }
1059 if !value.is_null() {
1060 let mut goal: codewhale_protocol::ThreadGoal = serde_json::from_value(value.clone())?;
1061 anyhow::ensure!(
1062 goal.thread_id == id,
1063 "canonical goal response belongs to another thread"
1064 );
1065 goal.thread_id = key;
1066 result.goal = Some(goal);
1067 }
1068 result.data =
1069 json!({"goal":result.goal,"cleared":method==Method::DELETE && result.status!="empty"});
1070 } else {
1071 let record = value.get("thread").unwrap_or(&value);
1072 anyhow::ensure!(
1073 record.get("id").and_then(Value::as_str) == Some(id.as_str()),
1074 "canonical control returned another target identity"
1075 );
1076 let mut thread = project_thread(record, &key)?;
1077 if running_ids(&bridge).await?.contains(&id) && thread.status != ThreadStatus::Archived {
1078 thread.status = ThreadStatus::Running;
1079 }
1080 result.thread = Some(thread);
1081 result.model = record
1082 .get("model")
1083 .and_then(Value::as_str)
1084 .map(str::to_owned);
1085 result.model_provider = result.thread.as_ref().map(|t| t.model_provider.clone());
1086 result.cwd = result.thread.as_ref().map(|t| t.cwd.clone());
1087 result.data = value;
1088 if method == Method::GET {
1089 let snapshot = full_snapshot(state, &bridge, &id).await?;
1090 result.data["history"] = serde_json::to_value(snapshot)?;
1091 }
1092 }
1093 Ok(result)
1094 }
1095
1096 /// Explicit CLI startup facts. The helper fills scheduler facts only from
1097 /// the authenticated owner's routing, never from a guessed local Config.
1098 #[derive(Debug, Clone, PartialEq, Eq)]
1099 pub struct ThreadControlSelection {
1100 pub workspace: Option<PathBuf>,
1101 pub config_profile: Option<String>,
1102 pub config_source: Option<PathBuf>,
1103 }
1104
1105 #[cfg(any(unix, windows))]
1106 pub async fn request_thread_control(
1107 config_path: Option<PathBuf>,
1108 selected: Option<PathBuf>,
1109 selection: Option<ThreadControlSelection>,
1110 request: ThreadRequest,
1111 ) -> Result<ThreadResponse> {
1112 let mut client = daemon_client::connect(config_path.clone(), selected).await?;
1113 if let Some(selection) = selection {
1114 let routing = client
1115 .routing()
1116 .context("selected owner has no acknowledged control scope")?;
1117 let scope = RuntimeFrontendScope {
1118 workers: routing
1119 .workers
1120 .context("selected owner has no captured worker setting")?,
1121 workspace: selection
1122 .workspace
1123 .or_else(|| routing.workspace.clone())
1124 .context("selected owner has no acknowledged workspace")?,
1125 config_profile: selection.config_profile,
1126 config_source: selection.config_source,
1127 };
1128 scope.validate_bounds()?;
1129 let owner = client.receipt().clone();
1130 let socket = owner.socket_path.clone();
1131 drop(client);
1132 client = daemon_client::connect_scoped_control_if_published(
1133 config_path,
1134 Some(socket),
1135 scope,
1136 owner,
1137 )
1138 .await?
1139 .context("captured owner disappeared during scope admission; no replacement")?;
1140 }
1141 let request_id = json!(format!("thread-control-{}", Uuid::new_v4()));
1142 client
1143 .send(
1144 request_id.clone(),
1145 "thread/request",
1146 serde_json::to_value(request)?,
1147 )
1148 .await?;
1149 tokio::time::timeout(Duration::from_secs(30), async {
1150 let mut notifications = 0usize;
1151 loop {
1152 let frame = client
1153 .recv()
1154 .await?
1155 .context("selected owner closed; outcome uncertain, not replayed")?;
1156 if frame.get("id") == Some(&request_id) {
1157 if let Some(error) = frame.get("error") {
1158 bail!("canonical owner refused thread control: {error}")
1159 }
1160 return serde_json::from_value(
1161 frame
1162 .get("result")
1163 .cloned()
1164 .context("canonical control result missing")?,
1165 )
1166 .context("invalid canonical control response");
1167 }
1168 anyhow::ensure!(
1169 frame.get("id").is_none(),
1170 "canonical owner returned another request identity"
1171 );
1172 notifications += 1;
1173 anyhow::ensure!(
1174 notifications <= MAX_CANONICAL_HISTORY_ENTRIES,
1175 "canonical owner notification bound exceeded"
1176 );
1177 }
1178 })
1179 .await
1180 .context("canonical control deadline expired; outcome uncertain, not replayed")?
1181 }
1182 #[cfg(not(any(unix, windows)))]
1183 pub async fn request_thread_control(
1184 _config_path: Option<PathBuf>,
1185 _selected: Option<PathBuf>,
1186 _selection: Option<ThreadControlSelection>,
1187 _request: ThreadRequest,
1188 ) -> Result<ThreadResponse> {
1189 bail!("canonical owner attachment is unsupported on this platform")
1190 }
1191
1192 #[cfg(test)]
1193 mod tests;
1194
1195 #[cfg(test)]
1196 pub(super) async fn compatibility_fixture()
1197 -> (AppState, tempfile::TempDir, tokio::task::JoinHandle<()>) {
1198 tests::compatibility_fixture().await
1199 }
1200
1201 #[cfg(test)]
1202 pub(super) fn compatibility_router() -> (AppState, tempfile::TempDir, Router) {
1203 tests::compatibility_router()
1204 }
1205
1205 lines RUST