| 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 |