1use std::{
3 collections::{BTreeMap, BTreeSet},
4 fs::{self, OpenOptions},
5 io::{BufRead, BufReader, BufWriter, Write},
6 path::{Path, PathBuf},
7};
8
9use anyhow::{Result, anyhow};
10use base64::Engine as _;
11use chrono::Utc;
12use heddle_core::{
13 ExplicitAgentBind, SessionAttachFacts, SessionLookupFact, SessionPolicy, TokenSidFact,
14 WorktreeSessionFact, decide_session_attach, first_value_string, map_from_pairs,
15 merge_string_vec, opencode_tool_name, opencode_tool_status,
16 parse_relay_payload as core_parse_relay_payload,
17 should_rotate_segment as pure_should_rotate_segment, value_array_join, value_cost_micros,
18 value_cost_micros_u64, value_string, value_string_array, value_u64, value_u64_string,
19};
20use objects::{
21 fs_atomic::write_file_atomic,
22 object::{
23 ContentHash, DiffKind, NativeToolCallRefV1, Session, StateId, ThreadName, TimelineBranchId,
24 TimelineLabel, TimelineOperationBodyV1, TimelineOperationEnvelope, TimelineStepId,
25 TimelineToolCallStatus, TimelineToolPayloadMetadata, ToolCallFinishedV1, ToolCallStartedV1,
26 Tree,
27 },
28 store::{ActorPresence, ActorPresenceStore, ActorPresenceStatus, AgentUsageSummary, ObjectStore},
29};
30use oplog::OpLogRecorder;
31use refs::Head;
32use repo::{
33 Repository, SessionManager, Thread, ThreadFreshness, ThreadIntegrationPolicy, ThreadManager,
34 ThreadMode, ThreadState, TimelineStore, TimelineView,
35};
36use serde::{Deserialize, Serialize};
37use serde_json::Value;
38use wire::{
39 HarnessIdentity, ProgressCheckpoint, SessionDiffSummary, SessionReportEnvelope,
40 TranscriptAttachmentRef, UsageTotals, WorktreeChangeBaseline,
41};
42
43mod claude_hook;
44mod probe;
45
46use self::probe::{HarnessProbeInput, HarnessProbeResult, probe_harness_actor};
47use crate::{
48 cli::{
49 Cli,
50 commands::{
51 snapshot::{
52 SnapshotAgentOverrides, create_snapshot, summarize_confidence,
53 summarize_verification,
54 },
55 worktree_cmd::helpers::{prepare_worktree_target, write_isolated_checkout},
56 },
57 style, worktree_status_options,
58 },
59 config::{
60 HarnessMode, HarnessTranscriptMode, HarnessTransport, UserConfig, UserHarnessOverride,
61 UserHarnessRootThreadPolicy, UserHarnessSubagentThreadPolicy, UserThreadWorkspaceMode,
62 },
63};
64
65pub(crate) fn probe_current_process_harness(
66 repo: &Repository,
67 current_provider: Option<String>,
68 current_model: Option<String>,
69 current_policy: Option<String>,
70) -> Result<HarnessProbeResult> {
71 probe_harness_actor(&HarnessProbeInput {
72 argv: detected_harness_argv().or_else(|| Some(std::env::args().collect())),
73 env_hints: harness_env_hints(),
74 explicit_harness: None,
75 explicit_provider: None,
76 explicit_model: None,
77 explicit_thinking_level: None,
78 explicit_policy: None,
79 probe_metadata: BTreeMap::new(),
80 current_provider,
81 current_model,
82 current_policy,
83 repo_root: repo.root().display().to_string(),
84 })
85}
86
87fn harness_env_hints() -> BTreeMap<String, String> {
88 std::env::vars()
89 .filter(|(key, value)| {
90 !value.trim().is_empty()
91 && (key.starts_with("HEDDLE_AGENT_")
92 || key.starts_with("CODEX_")
93 || key.starts_with("CLAUDE")
94 || key.starts_with("ANTHROPIC_")
95 || key.starts_with("OPENAI_")
96 || key.starts_with("OPENCODE_")
97 || key.starts_with("AIDER_")
98 || matches!(
99 key.as_str(),
100 "MODEL" | "REASONING_EFFORT" | "THINKING_LEVEL"
101 ))
102 })
103 .collect()
104}
105
106fn detected_harness_argv() -> Option<Vec<String>> {
107 detected_harness_argv_impl()
108}
109
110#[cfg(target_os = "linux")]
111fn detected_harness_argv_impl() -> Option<Vec<String>> {
112 let mut pid = std::process::id();
113 for _ in 0..8 {
114 let stat = fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
115 let ppid = stat.split_whitespace().nth(3)?.parse::<u32>().ok()?;
116 if ppid == 0 || ppid == pid {
117 return None;
118 }
119 pid = ppid;
120 let raw = fs::read(format!("/proc/{pid}/cmdline")).ok()?;
121 let argv = raw
122 .split(|byte| *byte == 0)
123 .filter(|part| !part.is_empty())
124 .map(|part| String::from_utf8_lossy(part).to_string())
125 .collect::<Vec<_>>();
126 let program = argv.first().map(|arg| arg.to_ascii_lowercase())?;
127 if ["codex", "claude", "opencode", "aider"]
128 .iter()
129 .any(|needle| program.contains(needle))
130 {
131 return Some(argv);
132 }
133 }
134 None
135}
136
137#[cfg(not(target_os = "linux"))]
138fn detected_harness_argv_impl() -> Option<Vec<String>> {
139 None
140}
141
142pub fn cmd_harness_bridge(cli: &Cli) -> Result<()> {
143 let repo = cli.open_repo()?;
144 let mut runtime = init_harness_runtime(&repo)?;
145
146 let stdin = std::io::stdin();
147 let stdout = std::io::stdout();
148 let reader = BufReader::new(stdin.lock());
149 let mut writer = BufWriter::new(stdout.lock());
150
151 for line in reader.lines() {
152 let line = line?;
153 if line.trim().is_empty() {
154 continue;
155 }
156 let response = match serde_json::from_str::<BridgeRequest>(&line) {
157 Ok(request) => runtime.handle_request(request),
158 Err(err) => BridgeResponse::error(
159 None,
160 "invalid_request",
161 format!("failed to parse request: {err}"),
162 ),
163 };
164 serde_json::to_writer(&mut writer, &response)?;
165 writer.write_all(b"\n")?;
166 writer.flush()?;
167 }
168
169 Ok(())
170}
171
172pub(crate) fn relay_harness_event(
173 repo: &Repository,
174 harness: &str,
175 event: &str,
176 payload: &str,
177) -> Result<()> {
178 let mut runtime = init_harness_runtime(repo)?;
179 let (json, warning) = parse_relay_payload(payload);
180 if let Some(warning) = warning {
181 eprintln!("{}", style::warn(&warning));
182 }
183 match harness {
184 "codex" => relay_codex(&mut runtime, event, &json),
185 "claude-code" => relay_claude(&mut runtime, event, &json),
186 "opencode" => relay_opencode(&mut runtime, event, &json),
187 other => Err(anyhow!("unsupported harness relay: {other}")),
188 }
189}
190
191fn init_harness_runtime(repo: &Repository) -> Result<HarnessBridgeRuntime> {
192 let (user_config, warning) = load_harness_user_config(UserConfig::default_path());
193 if let Some(warning) = warning {
194 eprintln!("{}", style::warn(&warning));
195 }
196 Ok(HarnessBridgeRuntime::new(
197 Repository::open(repo.root())?,
198 user_config,
199 ))
200}
201
202fn load_harness_user_config(default_path: Option<PathBuf>) -> (UserConfig, Option<String>) {
203 let Some(path) = default_path else {
204 return (UserConfig::default(), None);
205 };
206 match UserConfig::load(&path) {
207 Ok(config) => (config, None),
208 Err(err) if is_not_found(&err) => (UserConfig::default(), None),
209 Err(err) => {
210 let warning = format!(
211 "warning: failed to load user config from {}: {err}; continuing with defaults",
212 path.display()
213 );
214 (UserConfig::default(), Some(warning))
215 }
216 }
217}
218
219fn is_not_found(err: &anyhow::Error) -> bool {
220 err.downcast_ref::<std::io::Error>()
221 .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound)
222}
223
224fn parse_relay_payload(payload: &str) -> (Value, Option<String>) {
225 core_parse_relay_payload(payload)
226}
227
228struct HarnessBridgeRuntime {
229 repo: Repository,
230 user_config: UserConfig,
231 reports: SessionReportStore,
232}
233
234struct RegistryEntryRequest<'a> {
235 heddle_session_id: &'a str,
236 thread_name: Option<&'a str>,
237 thread_id: Option<&'a str>,
238 identity: &'a ResolvedIdentity,
239 probe: &'a HarnessProbeResult,
240 attach: &'a ResolvedAttachment,
241 client_instance_id: Option<&'a str>,
242 requested_entry: Option<&'a ActorPresence>,
243}
244
245struct CanonicalActorSessionRequest<'a> {
246 tentative_session: Session,
247 tentative_owns_session: bool,
248 entry: &'a ActorPresence,
249 probe: &'a HarnessProbeResult,
250 attach: &'a mut ResolvedAttachment,
251}
252
253struct AttachmentResolutionInput<'a> {
254 requested_entry: Option<&'a ActorPresence>,
255 explicit_heddle_session_id: Option<&'a str>,
256 client_instance_id: Option<&'a str>,
257 probe: &'a HarnessProbeResult,
258 token_claims: Option<&'a TokenClaims>,
259}
260
261fn relay_codex(runtime: &mut HarnessBridgeRuntime, _event: &str, payload: &Value) -> Result<()> {
262 let metadata = map_from_pairs([
263 (
264 "client_name",
265 value_string(payload, &["client"]).or_else(|| value_string(payload, &["client_name"])),
266 ),
267 ("model", value_string(payload, &["model"])),
268 (
269 "model_provider",
270 value_string(payload, &["model_provider"])
271 .or_else(|| value_string(payload, &["provider"])),
272 ),
273 (
274 "model_reasoning_effort",
275 value_string(payload, &["reasoning_effort"]),
276 ),
277 ]);
278 let opened = runtime.open_session(OpenSessionParams {
279 harness: Some("codex".to_string()),
280 summary: value_string(payload, &["message"]),
281 probe_metadata: metadata,
282 ..OpenSessionParams::default()
283 })?;
284 runtime.update_progress(UpdateProgressParams {
285 heddle_session_id: opened.heddle_session_id,
286 summary: value_string(payload, &["message"]),
287 harness: Some("codex".to_string()),
288 ..UpdateProgressParams::default()
289 })?;
290 Ok(())
291}
292
293fn relay_claude(runtime: &mut HarnessBridgeRuntime, event: &str, payload: &Value) -> Result<()> {
294 let metadata = map_from_pairs([
295 ("session_id", value_string(payload, &["session_id"])),
296 ("agent_id", value_string(payload, &["agent_id"])),
297 ("session_name", value_string(payload, &["session_name"])),
298 (
299 "transcript_path",
300 value_string(payload, &["transcript_path"]),
301 ),
302 (
303 "model",
304 value_string(payload, &["model", "id"]).or_else(|| value_string(payload, &["model"])),
305 ),
306 (
307 "model_display_name",
308 value_string(payload, &["model", "display_name"]),
309 ),
310 ("effort", value_string(payload, &["effort"])),
311 ("hook_event", Some(event.to_string())),
312 (
313 "status_line",
314 (event == "StatusLine").then(|| "1".to_string()),
315 ),
316 (
317 "touched_paths",
318 value_array_join(payload, &["tool_response", "filePaths"])
319 .or_else(|| value_string(payload, &["file_path"])),
320 ),
321 (
322 "input_tokens",
323 value_u64_string(payload, &["context_window", "total_input_tokens"]),
324 ),
325 (
326 "output_tokens",
327 value_u64_string(payload, &["context_window", "total_output_tokens"]),
328 ),
329 (
330 "cost_micros_usd",
331 value_cost_micros(payload, &["cost", "total_cost_usd"]),
332 ),
333 ]);
334 let opened = runtime.open_session(OpenSessionParams {
335 harness: Some("claude-code".to_string()),
336 model: value_string(payload, &["model", "display_name"])
337 .or_else(|| value_string(payload, &["model", "id"]))
338 .or_else(|| value_string(payload, &["model"])),
339 summary: value_string(payload, &["message"]).or_else(|| value_string(payload, &["reason"])),
340 probe_metadata: metadata.clone(),
341 ..OpenSessionParams::default()
342 })?;
343 match event {
344 "SessionEnd" => {
345 runtime.close_session(CloseSessionParams {
346 heddle_session_id: opened.heddle_session_id,
347 summary: value_string(payload, &["reason"])
348 .or_else(|| value_string(payload, &["stop_hook_active"])),
349 outcome: Some("completed".to_string()),
350 ..CloseSessionParams::default()
351 })?;
352 }
353 "StatusLine" => {
354 runtime.update_progress(UpdateProgressParams {
355 heddle_session_id: opened.heddle_session_id.clone(),
356 harness: Some("claude-code".to_string()),
357 status: Some("StatusLine".to_string()),
358 message: value_string(payload, &["session_name"])
359 .or_else(|| value_string(payload, &["cwd"]))
360 .or_else(|| value_string(payload, &["workspace", "current_dir"])),
361 probe_metadata: metadata.clone(),
362 ..UpdateProgressParams::default()
363 })?;
364 runtime.record_usage(RecordUsageParams {
365 heddle_session_id: opened.heddle_session_id,
366 input_tokens: value_u64(payload, &["context_window", "total_input_tokens"]),
367 output_tokens: value_u64(payload, &["context_window", "total_output_tokens"]),
368 reasoning_tokens: value_u64(payload, &["context_window", "total_reasoning_tokens"]),
369 cache_creation_tokens: None,
370 cache_read_tokens: None,
371 tool_calls: None,
372 cost_micros_usd: value_cost_micros_u64(payload, &["cost", "total_cost_usd"]),
373 })?;
374 }
375 "Stop" => {
376 runtime.update_progress(UpdateProgressParams {
377 heddle_session_id: opened.heddle_session_id,
378 harness: Some("claude-code".to_string()),
379 status: Some("Stop".to_string()),
380 message: value_string(payload, &["message"])
381 .or_else(|| value_string(payload, &["result"]))
382 .or_else(|| value_string(payload, &["stop_reason"])),
383 probe_metadata: metadata,
384 ..UpdateProgressParams::default()
385 })?;
386 if let Err(err) = claude_hook::handle_stop_capture(
387 &runtime.repo,
388 &runtime.user_config,
389 payload,
390 "Claude Code turn",
391 ) {
392 tracing::warn!(?err, "heddle Stop hook capture failed");
393 }
394 }
395 "SubagentStop" => {
396 runtime.update_progress(UpdateProgressParams {
397 heddle_session_id: opened.heddle_session_id,
398 harness: Some("claude-code".to_string()),
399 status: Some("SubagentStop".to_string()),
400 touched_paths: csv_from_value(metadata.get("touched_paths")),
401 probe_metadata: metadata,
402 ..UpdateProgressParams::default()
403 })?;
404 if let Err(err) = claude_hook::handle_stop_capture(
405 &runtime.repo,
406 &runtime.user_config,
407 payload,
408 "Claude Code subagent turn",
409 ) {
410 tracing::warn!(?err, "heddle SubagentStop hook capture failed");
411 }
412 if let Err(err) = claude_hook::mark_subagent_complete(&runtime.repo, payload) {
413 tracing::debug!(?err, "heddle SubagentStop mark-complete failed");
414 }
415 }
416 "SubagentStart" => {
417 runtime.update_progress(UpdateProgressParams {
424 heddle_session_id: opened.heddle_session_id,
425 harness: Some("claude-code".to_string()),
426 status: Some("SubagentStart".to_string()),
427 touched_paths: csv_from_value(metadata.get("touched_paths")),
428 probe_metadata: metadata,
429 ..UpdateProgressParams::default()
430 })?;
431 }
432 "UserPromptSubmit" => {
433 runtime.update_progress(UpdateProgressParams {
434 heddle_session_id: opened.heddle_session_id.clone(),
435 harness: Some("claude-code".to_string()),
436 status: Some("UserPromptSubmit".to_string()),
437 touched_paths: csv_from_value(metadata.get("touched_paths")),
438 probe_metadata: metadata,
439 ..UpdateProgressParams::default()
440 })?;
441 if let Err(err) = claude_hook::handle_user_prompt_segment_rotate(
442 &runtime.repo,
443 &opened.heddle_session_id,
444 payload,
445 ) {
446 tracing::debug!(?err, "heddle UserPromptSubmit segment rotation failed");
447 }
448 }
449 "PreToolUse" => {
450 runtime.update_progress(UpdateProgressParams {
451 heddle_session_id: opened.heddle_session_id,
452 harness: Some("claude-code".to_string()),
453 status: Some("PreToolUse".to_string()),
454 touched_paths: csv_from_value(metadata.get("touched_paths")),
455 probe_metadata: metadata,
456 ..UpdateProgressParams::default()
457 })?;
458 if let Err(err) = claude_hook::handle_pre_tool_use(&runtime.repo, payload) {
459 tracing::debug!(?err, "heddle PreToolUse context inject skipped");
460 }
461 }
462 _ => {
463 runtime.update_progress(UpdateProgressParams {
464 heddle_session_id: opened.heddle_session_id,
465 harness: Some("claude-code".to_string()),
466 status: Some(event.to_string()),
467 touched_paths: csv_from_value(metadata.get("touched_paths")),
468 probe_metadata: metadata,
469 ..UpdateProgressParams::default()
470 })?;
471 }
472 }
473 Ok(())
474}
475
476fn relay_opencode(runtime: &mut HarnessBridgeRuntime, event: &str, payload: &Value) -> Result<()> {
477 let metadata = map_from_pairs([
478 (
479 "session_id",
480 value_string(payload, &["sessionID"])
481 .or_else(|| value_string(payload, &["session_id"])),
482 ),
483 (
484 "parent_id",
485 value_string(payload, &["parentID"]).or_else(|| value_string(payload, &["parent_id"])),
486 ),
487 (
488 "client_name",
489 value_string(payload, &["client"]).or_else(|| std::env::var("OPENCODE_CLIENT").ok()),
490 ),
491 ("model", value_string(payload, &["model"])),
492 ("provider", value_string(payload, &["provider"])),
493 ("hook_event", Some(event.to_string())),
494 (
495 "touched_paths",
496 value_string(payload, &["file", "path"]).or_else(|| value_string(payload, &["path"])),
497 ),
498 ]);
499 let opened = runtime.open_session(OpenSessionParams {
500 harness: Some("opencode".to_string()),
501 model: value_string(payload, &["model"]),
502 provider: value_string(payload, &["provider"]),
503 probe_metadata: metadata.clone(),
504 ..OpenSessionParams::default()
505 })?;
506 let session_id = opened.heddle_session_id.clone();
507 runtime.update_progress(UpdateProgressParams {
508 heddle_session_id: session_id.clone(),
509 harness: Some("opencode".to_string()),
510 status: Some(event.to_string()),
511 touched_paths: csv_from_value(metadata.get("touched_paths")),
512 probe_metadata: metadata,
513 ..UpdateProgressParams::default()
514 })?;
515 if let Err(err) = record_opencode_timeline_event(runtime, event, payload, &opened) {
516 tracing::debug!(?err, event, "heddle OpenCode timeline recording skipped");
517 }
518 Ok(())
519}
520
521#[derive(Clone, Copy, Debug, PartialEq, Eq)]
522enum TimelineToolEvent {
523 Started,
524 Finished,
525}
526
527trait HarnessTimelineExtractor {
528 fn timeline_event(&self, event: &str) -> Option<TimelineToolEvent>;
529 fn native_tool_call(&self, payload: &Value) -> Option<NativeToolCallRefV1>;
530 fn tool_name(&self, payload: &Value) -> String;
531 fn tool_status(&self, payload: &Value) -> TimelineToolCallStatus;
532 fn payload_metadata(&self, event: &str, payload: &Value)
533 -> Result<TimelineToolPayloadMetadata>;
534 fn touched_paths(&self, payload: &Value) -> Vec<String>;
535 fn capture_intent(&self, native: &NativeToolCallRefV1, payload: &Value) -> String;
536
537 fn timeline_thread(
538 &self,
539 runtime: &HarnessBridgeRuntime,
540 opened: &OpenSessionResult,
541 ) -> Result<String> {
542 if let Some(report) = runtime.reports.load(&opened.heddle_session_id)?
543 && let Some(thread) = report.thread
544 {
545 return Ok(thread);
546 }
547 match runtime.repo.head_ref()? {
548 Head::Attached { thread } => Ok(thread.to_string()),
549 Head::Detached { .. } => Ok("main".to_string()),
550 }
551 }
552
553 fn stable_step_id(&self, native: &NativeToolCallRefV1) -> TimelineStepId {
554 let key = format!(
555 "{}\0{}\0{}\0{}",
556 native.harness,
557 native.session_id.as_deref().unwrap_or(""),
558 native.message_id.as_deref().unwrap_or(""),
559 native.tool_call_id
560 );
561 let hash =
562 ContentHash::compute_typed("timeline-native-tool-call-v1", key.as_bytes()).to_hex();
563 TimelineStepId::new(format!("tls-{}", &hash[..24]))
564 }
565
566 fn started_labels(&self, _payload: &Value) -> Vec<TimelineLabel> {
567 vec![TimelineLabel::ExternalSideEffectsUnknown]
568 }
569
570 fn finished_labels(&self, changed: bool, _payload: &Value) -> Vec<TimelineLabel> {
571 if changed {
572 vec![
573 TimelineLabel::RepoReversible,
574 TimelineLabel::ExternalSideEffectsUnknown,
575 ]
576 } else {
577 vec![TimelineLabel::ExternalSideEffectsUnknown]
578 }
579 }
580}
581
582struct OpenCodeTimelineExtractor;
583
584impl HarnessTimelineExtractor for OpenCodeTimelineExtractor {
585 fn timeline_event(&self, event: &str) -> Option<TimelineToolEvent> {
586 match event {
587 "tool.execute.before" => Some(TimelineToolEvent::Started),
588 "tool.execute.after" => Some(TimelineToolEvent::Finished),
589 _ => None,
590 }
591 }
592
593 fn native_tool_call(&self, payload: &Value) -> Option<NativeToolCallRefV1> {
594 opencode_native_tool_call(payload)
595 }
596
597 fn tool_name(&self, payload: &Value) -> String {
598 opencode_tool_name(payload)
599 }
600
601 fn tool_status(&self, payload: &Value) -> TimelineToolCallStatus {
602 opencode_tool_status(payload)
603 }
604
605 fn payload_metadata(
606 &self,
607 event: &str,
608 payload: &Value,
609 ) -> Result<TimelineToolPayloadMetadata> {
610 opencode_payload_metadata(event, payload)
611 }
612
613 fn touched_paths(&self, payload: &Value) -> Vec<String> {
614 opencode_touched_paths(payload)
615 }
616
617 fn capture_intent(&self, native: &NativeToolCallRefV1, payload: &Value) -> String {
618 format!(
619 "OpenCode {} tool call {}",
620 self.tool_name(payload),
621 native.tool_call_id
622 )
623 }
624}
625
626fn record_opencode_timeline_event(
627 runtime: &mut HarnessBridgeRuntime,
628 event: &str,
629 payload: &Value,
630 opened: &OpenSessionResult,
631) -> Result<()> {
632 record_timeline_event(runtime, event, payload, opened, &OpenCodeTimelineExtractor)
633}
634
635fn record_timeline_event<E: HarnessTimelineExtractor>(
636 runtime: &mut HarnessBridgeRuntime,
637 event: &str,
638 payload: &Value,
639 opened: &OpenSessionResult,
640 extractor: &E,
641) -> Result<()> {
642 match extractor.timeline_event(event) {
643 Some(TimelineToolEvent::Started) => {
644 record_timeline_tool_started(runtime, event, payload, opened, extractor)
645 }
646 Some(TimelineToolEvent::Finished) => {
647 record_timeline_tool_finished(runtime, event, payload, opened, extractor)
648 }
649 None => Ok(()),
650 }
651}
652
653fn record_timeline_tool_started<E: HarnessTimelineExtractor>(
654 runtime: &mut HarnessBridgeRuntime,
655 event: &str,
656 payload: &Value,
657 opened: &OpenSessionResult,
658 extractor: &E,
659) -> Result<()> {
660 let Some(native) = extractor.native_tool_call(payload) else {
661 return Ok(());
662 };
663 let Some(before_state) = current_state_id(&runtime.repo)? else {
664 return Ok(());
665 };
666 let thread = extractor.timeline_thread(runtime, opened)?;
667 let store = TimelineStore::open(runtime.repo.heddle_dir())?;
668 let _record_guard = store.lock_recording(&thread)?;
669 let view = TimelineView::rebuild(&store)?;
670 let step_id = extractor.stable_step_id(&native);
671 let (branch_id, parent_step_id) = timeline_position_for_new_tool_step(&view, &thread, &step_id);
672 let envelope = TimelineOperationEnvelope::new(
673 TimelineOperationBodyV1::ToolCallStarted(ToolCallStartedV1 {
674 thread,
675 step_id,
676 branch_id,
677 parent_step_id,
678 native,
679 tool_name: extractor.tool_name(payload),
680 before_state,
681 payload: Some(extractor.payload_metadata(event, payload)?),
682 started_at_ms: Utc::now().timestamp_millis(),
683 }),
684 extractor.started_labels(payload),
685 );
686 store.write_operation(&envelope)?;
687 Ok(())
688}
689
690fn record_timeline_tool_finished<E: HarnessTimelineExtractor>(
691 runtime: &mut HarnessBridgeRuntime,
692 event: &str,
693 payload: &Value,
694 opened: &OpenSessionResult,
695 extractor: &E,
696) -> Result<()> {
697 let Some(native) = extractor.native_tool_call(payload) else {
698 return Ok(());
699 };
700 let Some(fallback_state) = current_state_id(&runtime.repo)? else {
701 return Ok(());
702 };
703 let thread = extractor.timeline_thread(runtime, opened)?;
704 let store = TimelineStore::open(runtime.repo.heddle_dir())?;
705 let _record_guard = store.lock_recording(&thread)?;
706 let before_view = TimelineView::rebuild(&store)?;
707 let step_id = extractor.stable_step_id(&native);
708 let (branch_id, _) = timeline_position_for_new_tool_step(&before_view, &thread, &step_id);
709 let before_state = before_view
710 .step(&thread, &step_id)
711 .and_then(|step| step.before_state)
712 .unwrap_or(fallback_state);
713 let has_worktree_changes_before_capture = !collect_worktree_changes(&runtime.repo)?.is_empty();
714 let mut capture_failed = false;
715 let capture_state = if !has_worktree_changes_before_capture {
716 None
717 } else {
718 let intent = extractor.capture_intent(&native, payload);
719 match create_snapshot(
720 &runtime.repo,
721 &runtime.user_config,
722 Some(intent),
723 None,
724 SnapshotAgentOverrides {
725 provider: opened.provider.clone(),
726 model: opened.model.clone(),
727 session: native.session_id.clone(),
728 segment: None,
729 policy: None,
730 no_policy: false,
731 no_agent: false,
732 },
733 ) {
734 Ok(_) => runtime.repo.head()?,
735 Err(err) => {
736 capture_failed = true;
737 tracing::warn!(?err, "heddle timeline tool capture failed");
738 None
739 }
740 }
741 };
742 let after_state = current_state_id(&runtime.repo)?.unwrap_or(fallback_state);
743 let mut touched_paths = extractor.touched_paths(payload);
744 merge_string_vec(
745 &mut touched_paths,
746 changed_paths_between_states(&runtime.repo, before_state, after_state)?,
747 );
748 let changed = before_state != after_state;
749 let mut labels = extractor.finished_labels(changed, payload);
750 if capture_failed {
751 merge_timeline_labels(&mut labels, vec![TimelineLabel::CaptureFailed]);
752 }
753 let envelope = TimelineOperationEnvelope::new(
754 TimelineOperationBodyV1::ToolCallFinished(ToolCallFinishedV1 {
755 thread,
756 step_id,
757 branch_id,
758 native,
759 status: extractor.tool_status(payload),
760 before_state,
761 after_state,
762 capture_state,
763 capture_oplog_batch_id: None,
764 changed,
765 touched_paths,
766 payload: Some(extractor.payload_metadata(event, payload)?),
767 finished_at_ms: Utc::now().timestamp_millis(),
768 }),
769 labels,
770 );
771 store.write_operation(&envelope)?;
772 Ok(())
773}
774
775fn opencode_native_tool_call(payload: &Value) -> Option<NativeToolCallRefV1> {
776 let tool_call_id = first_value_string(
777 payload,
778 &[
779 &["toolCallID"],
780 &["tool_call_id"],
781 &["toolCallId"],
782 &["callID"],
783 &["call_id"],
784 &["tool", "callID"],
785 &["tool", "call_id"],
786 &["tool", "id"],
787 &["toolCall", "id"],
788 &["tool_call", "id"],
789 &["id"],
790 ],
791 )?;
792 Some(NativeToolCallRefV1 {
793 harness: "opencode".to_string(),
794 session_id: value_string(payload, &["sessionID"])
795 .or_else(|| value_string(payload, &["session_id"])),
796 message_id: value_string(payload, &["messageID"])
797 .or_else(|| value_string(payload, &["message_id"]))
798 .or_else(|| value_string(payload, &["message", "id"])),
799 tool_call_id,
800 })
801}
802
803fn timeline_position_for_new_tool_step(
804 view: &TimelineView,
805 thread: &str,
806 step_id: &TimelineStepId,
807) -> (TimelineBranchId, Option<TimelineStepId>) {
808 let branch_id = view
809 .status(thread)
810 .and_then(|status| status.current_branch_id.clone())
811 .unwrap_or_else(|| TimelineBranchId::new("tlb-main"));
812 let parent_step_id = view
813 .status(thread)
814 .and_then(|status| status.current_step_id.clone())
815 .filter(|current| current != step_id);
816 (branch_id, parent_step_id)
817}
818
819fn current_state_id(repo: &Repository) -> Result<Option<StateId>> {
820 Ok(repo
821 .current_state()?
822 .map(|state| state.state_id)
823 .or(repo.head()?))
824}
825
826fn opencode_payload_metadata(event: &str, payload: &Value) -> Result<TimelineToolPayloadMetadata> {
829 let tool_name = opencode_tool_name(payload);
830 let tool_call_id = opencode_native_tool_call(payload)
831 .map(|native| native.tool_call_id)
832 .unwrap_or_default();
833 let raw = serde_json::to_vec(payload)?;
834 let hash = ContentHash::compute_typed("timeline-tool-payload", &raw);
835 let summary = if tool_call_id.is_empty() {
836 format!("OpenCode {event}: {tool_name}")
837 } else {
838 format!("OpenCode {event}: {tool_name} ({tool_call_id})")
839 };
840 Ok(TimelineToolPayloadMetadata {
841 summary: Some(summary),
842 hash: Some(hash),
843 })
844}
845
846fn opencode_touched_paths(payload: &Value) -> Vec<String> {
847 let mut paths = Vec::new();
848 for path in [
849 value_string(payload, &["file", "path"]),
850 value_string(payload, &["path"]),
851 value_string(payload, &["tool", "path"]),
852 value_string(payload, &["tool", "input", "file_path"]),
853 value_string(payload, &["input", "file_path"]),
854 ]
855 .into_iter()
856 .flatten()
857 {
858 if !path.trim().is_empty() && !paths.contains(&path) {
859 paths.push(path);
860 }
861 }
862 for value_path in [
863 &["paths"][..],
864 &["files"][..],
865 &["tool", "input", "paths"][..],
866 &["input", "paths"][..],
867 ] {
868 if let Some(items) = value_string_array(payload, value_path) {
869 merge_string_vec(&mut paths, items);
870 }
871 }
872 paths
873}
874
875fn merge_timeline_labels(target: &mut Vec<TimelineLabel>, incoming: Vec<TimelineLabel>) {
876 for label in incoming {
877 if !target.contains(&label) {
878 target.push(label);
879 }
880 }
881}
882
883fn csv_from_value(value: Option<&String>) -> Vec<String> {
884 value
885 .map(|value| {
886 value
887 .split(',')
888 .map(|item| item.trim().to_string())
889 .filter(|item| !item.is_empty())
890 .collect()
891 })
892 .unwrap_or_default()
893}
894
895impl HarnessBridgeRuntime {
896 fn new(repo: Repository, user_config: UserConfig) -> Self {
897 let reports = SessionReportStore::new(repo.root());
898 Self {
899 repo,
900 user_config,
901 reports,
902 }
903 }
904
905 fn handle_request(&mut self, request: BridgeRequest) -> BridgeResponse {
906 let response = match request.method.as_str() {
907 "open_session" => self
908 .decode_params::<OpenSessionParams>(request.params)
909 .and_then(|params| self.open_session(params))
910 .and_then(to_json_value),
911 "update_progress" => self
912 .decode_params::<UpdateProgressParams>(request.params)
913 .and_then(|params| self.update_progress(params))
914 .and_then(to_json_value),
915 "record_usage" => self
916 .decode_params::<RecordUsageParams>(request.params)
917 .and_then(|params| self.record_usage(params))
918 .and_then(to_json_value),
919 "record_touched_paths" => self
920 .decode_params::<RecordTouchedPathsParams>(request.params)
921 .and_then(|params| self.record_touched_paths(params))
922 .and_then(to_json_value),
923 "close_session" => self
924 .decode_params::<CloseSessionParams>(request.params)
925 .and_then(|params| self.close_session(params))
926 .and_then(to_json_value),
927 "flush_reports" => self
928 .decode_params::<FlushReportsParams>(request.params)
929 .and_then(|params| self.flush_reports(params))
930 .and_then(to_json_value),
931 other => Err(anyhow!("unknown method '{other}'")),
932 };
933
934 match response {
935 Ok(result) => BridgeResponse::ok(request.id, result),
936 Err(err) => BridgeResponse::error(request.id, "bridge_error", err.to_string()),
937 }
938 }
939
940 fn decode_params<T: for<'de> Deserialize<'de>>(&self, value: Value) -> Result<T> {
941 serde_json::from_value(value).map_err(|err| anyhow!(err))
942 }
943
944 fn open_session(&mut self, params: OpenSessionParams) -> Result<OpenSessionResult> {
945 if self.user_config.harness.mode == HarnessMode::Off {
946 return Err(anyhow!("harness integration is disabled in user config"));
947 }
948
949 let requested_transport = params
950 .transport
951 .unwrap_or(self.user_config.harness.transport);
952 let transcript_mode = params
953 .transcript_mode
954 .unwrap_or(self.user_config.harness.transcript);
955 let env_hints = merged_env_hints(¶ms.env_hints);
956 let token_claims = user_config_token_claims(&self.user_config);
957 let current_session = SessionManager::new(self.repo.root()).get_current_session()?;
958 let current_segment = current_session
959 .as_ref()
960 .and_then(|session| session.current_segment());
961 let probe = probe_harness_actor(&HarnessProbeInput {
962 argv: params.argv.clone(),
963 env_hints: env_hints.clone(),
964 explicit_harness: params.harness.clone(),
965 explicit_provider: params.provider.clone(),
966 explicit_model: params.model.clone(),
967 explicit_thinking_level: params.thinking_level.clone(),
968 explicit_policy: params.policy.clone(),
969 probe_metadata: params.probe_metadata.clone(),
970 current_provider: current_segment.map(|segment| segment.provider.clone()),
971 current_model: current_segment.map(|segment| segment.model.clone()),
972 current_policy: current_segment.and_then(|segment| segment.policy_id.clone()),
973 repo_root: self.repo.root().display().to_string(),
974 })?;
975 let identity = resolve_identity(
976 &self.repo,
977 &self.user_config,
978 IdentityHints {
979 harness: params.harness.clone(),
980 provider: params.provider.clone(),
981 model: params.model.clone(),
982 thinking_level: params.thinking_level.clone(),
983 policy: params.policy.clone(),
984 probe: probe.clone(),
985 },
986 )?;
987 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
988 let requested_entry = resolve_requested_registry_entry(
989 ®istry,
990 params.agent_session_id.as_deref(),
991 params.client_instance_id.as_deref(),
992 )?;
993
994 if self.user_config.harness.mode == HarnessMode::Required
995 && (identity.harness.is_none()
996 || identity.provider.is_none()
997 || identity.model.is_none())
998 {
999 return Err(anyhow!(
1000 "harness mode is 'required' but harness/provider/model could not be resolved"
1001 ));
1002 }
1003
1004 let mut sessions = SessionManager::new(self.repo.root());
1005 let principal = self.repo.get_principal()?;
1006 let mut attach = resolve_actor_attachment(
1007 ®istry,
1008 &self.repo,
1009 &mut sessions,
1010 AttachmentResolutionInput {
1011 requested_entry: requested_entry.as_ref(),
1012 explicit_heddle_session_id: params.heddle_session_id.as_deref(),
1013 client_instance_id: params.client_instance_id.as_deref(),
1014 probe: &probe,
1015 token_claims: token_claims.as_ref(),
1016 },
1017 )?;
1018 let (session, owns_session) = match &attach.target {
1019 AttachTarget::ExistingSession(session) => {
1020 let segment_id = session.current_segment_id.clone().unwrap_or_default();
1021 sessions.set_current_session(&session.id, &segment_id)?;
1022 (session.clone(), false)
1023 }
1024 AttachTarget::CreateNew {
1025 _because_claimed: _,
1026 } => {
1027 let session = sessions.start_session(
1028 principal,
1029 identity
1030 .provider
1031 .clone()
1032 .unwrap_or_else(|| "unknown".to_string()),
1033 identity
1034 .model
1035 .clone()
1036 .unwrap_or_else(|| "unknown".to_string()),
1037 identity.policy.clone(),
1038 )?;
1039 (session, true)
1040 }
1041 };
1042
1043 let (thread_name, thread_id) =
1044 self.resolve_harness_thread_binding(¶ms, &probe, &identity)?;
1045 let entry = self.ensure_registry_entry(RegistryEntryRequest {
1046 heddle_session_id: &session.id,
1047 thread_name: thread_name.as_deref(),
1048 thread_id: thread_id.as_deref(),
1049 identity: &identity,
1050 probe: &probe,
1051 attach: &attach,
1052 client_instance_id: params.client_instance_id.as_deref(),
1053 requested_entry: requested_entry.as_ref(),
1054 })?;
1055 let (session, owns_session) = self.reuse_canonical_actor_session(
1056 &mut sessions,
1057 CanonicalActorSessionRequest {
1058 tentative_session: session,
1059 tentative_owns_session: owns_session,
1060 entry: &entry,
1061 probe: &probe,
1062 attach: &mut attach,
1063 },
1064 )?;
1065
1066 let mut segment_id = session.current_segment_id.clone().unwrap_or_default();
1067 if should_rotate_segment(&session, &identity) {
1068 let segment = sessions.add_segment(
1069 &session.id,
1070 identity
1071 .provider
1072 .clone()
1073 .unwrap_or_else(|| "unknown".to_string()),
1074 identity
1075 .model
1076 .clone()
1077 .unwrap_or_else(|| "unknown".to_string()),
1078 identity.policy.clone(),
1079 )?;
1080 segment_id = segment.id;
1081 }
1082
1083 let base_state = self
1084 .repo
1085 .current_state()?
1086 .map(|state| state.state_id.to_string_full())
1087 .or_else(|| {
1088 self.repo
1089 .head()
1090 .ok()
1091 .flatten()
1092 .map(|id| id.to_string_full())
1093 });
1094 let worktree_changes_at_open = capture_worktree_change_snapshot(&self.repo)?;
1095 let opened_at = Utc::now().to_rfc3339();
1096 let mut report = SessionReportEnvelope {
1097 version: 1,
1098 heddle_session_id: session.id.clone(),
1099 heddle_segment_id: (!segment_id.is_empty()).then_some(segment_id.clone()),
1100 agent_session_id: Some(entry.session_id.clone()),
1101 client_instance_id: entry.client_instance_id.clone(),
1102 native_actor_key: entry.native_actor_key.clone(),
1103 native_parent_actor_key: entry.native_parent_actor_key.clone(),
1104 native_instance_key: entry.native_instance_key.clone(),
1105 repo_root: self.repo.root().display().to_string(),
1106 thread: thread_name.clone(),
1107 thread_id,
1108 task: params.task.clone(),
1109 summary: params.summary.clone(),
1110 opened_at,
1111 closed_at: None,
1112 base_state_at_open: base_state.clone(),
1113 worktree_changes_at_open,
1114 head_state_at_close: None,
1115 transport_mode: transport_mode_name(requested_transport).to_string(),
1116 transcript_mode: transcript_mode_name(transcript_mode).to_string(),
1117 outcome: None,
1118 harness: identity.to_transport_identity(),
1119 progress: Vec::new(),
1120 usage: UsageTotals::default(),
1121 touched_paths: Vec::new(),
1122 changed_paths: Vec::new(),
1123 diff_summary: None,
1124 transcript_refs: Vec::new(),
1125 last_progress_at: None,
1126 report_flush_state: Some("pending-local".to_string()),
1127 attach_reason: Some(attach.attach_reason.clone()),
1128 attach_precedence: attach.precedence.clone(),
1129 winning_attach_rule: Some(attach.winning_rule.clone()),
1130 probe_source: probe.probe_source.clone(),
1131 probe_confidence: probe.confidence,
1132 pending_flush: true,
1133 last_flushed_at: None,
1134 owns_session,
1135 };
1136 merge_unique_paths(&mut report.touched_paths, probe.touched_paths.clone());
1137 merge_usage(&mut report.usage, &probe.usage_totals);
1138 if transcript_mode != HarnessTranscriptMode::Off {
1139 report.transcript_refs = probe.transcript_refs.clone();
1140 }
1141 self.reports.save(&report)?;
1142 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1143 if matches!(requested_transport, HarnessTransport::Direct) {
1144 enqueue_report(&self.reports, &mut report)?;
1145 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1146 }
1147
1148 Ok(OpenSessionResult {
1149 heddle_session_id: report.heddle_session_id.clone(),
1150 heddle_segment_id: report.heddle_segment_id.clone(),
1151 agent_session_id: report.agent_session_id.clone(),
1152 created_session: owns_session,
1153 harness: report.harness.harness.clone(),
1154 provider: report.harness.provider.clone(),
1155 model: report.harness.model.clone(),
1156 thinking_level: report.harness.thinking_level.clone(),
1157 report_flush_state: report.report_flush_state.clone(),
1158 attach_reason: report.attach_reason.clone(),
1159 })
1160 }
1161
1162 fn update_progress(&mut self, params: UpdateProgressParams) -> Result<SessionMutationResult> {
1163 let mut report = self
1164 .reports
1165 .load(¶ms.heddle_session_id)?
1166 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1167 let current_session = SessionManager::new(self.repo.root()).get_current_session()?;
1168 let current_segment = current_session
1169 .as_ref()
1170 .and_then(|session| session.current_segment());
1171 let probe = probe_harness_actor(&HarnessProbeInput {
1172 argv: params.argv.clone(),
1173 env_hints: merged_env_hints(¶ms.env_hints),
1174 explicit_harness: params.harness.clone(),
1175 explicit_provider: params.provider.clone(),
1176 explicit_model: params.model.clone(),
1177 explicit_thinking_level: params.thinking_level.clone(),
1178 explicit_policy: params.policy.clone(),
1179 probe_metadata: params.probe_metadata.clone(),
1180 current_provider: current_segment.map(|segment| segment.provider.clone()),
1181 current_model: current_segment.map(|segment| segment.model.clone()),
1182 current_policy: current_segment.and_then(|segment| segment.policy_id.clone()),
1183 repo_root: self.repo.root().display().to_string(),
1184 })?;
1185 let identity = resolve_identity(
1186 &self.repo,
1187 &self.user_config,
1188 IdentityHints {
1189 harness: params.harness.clone(),
1190 provider: params.provider.clone(),
1191 model: params.model.clone(),
1192 thinking_level: params.thinking_level.clone(),
1193 policy: params.policy.clone(),
1194 probe: probe.clone(),
1195 },
1196 )?;
1197 self.ensure_segment_for_report(&mut report, &identity)?;
1198 if report.harness.harness.is_none() {
1199 report.harness.harness = identity.harness.clone();
1200 }
1201 if report.harness.provider.is_none() {
1202 report.harness.provider = identity.provider.clone();
1203 }
1204 if report.harness.model.is_none() {
1205 report.harness.model = identity.model.clone();
1206 }
1207 if report.harness.thinking_level.is_none() {
1208 report.harness.thinking_level = identity.thinking_level.clone();
1209 }
1210 if report.harness.policy.is_none() {
1211 report.harness.policy = identity.policy.clone();
1212 }
1213 if report.native_actor_key.is_none() {
1214 report.native_actor_key = probe.native_actor_key.clone();
1215 }
1216 if report.native_parent_actor_key.is_none() {
1217 report.native_parent_actor_key = probe.native_parent_actor_key.clone();
1218 }
1219 if report.native_instance_key.is_none() {
1220 report.native_instance_key = probe.native_instance_key.clone();
1221 }
1222 if report.probe_source.is_none() {
1223 report.probe_source = probe.probe_source.clone();
1224 }
1225 if report.probe_confidence.is_none() {
1226 report.probe_confidence = probe.confidence;
1227 }
1228
1229 let recorded_at = Utc::now().to_rfc3339();
1230 let checkpoint = ProgressCheckpoint {
1231 status: params.status.clone(),
1232 message: params.message.clone(),
1233 completed_steps: params.completed_steps,
1234 total_steps: params.total_steps,
1235 touched_paths: normalize_paths(
1236 params
1237 .touched_paths
1238 .into_iter()
1239 .chain(probe.touched_paths)
1240 .collect::<Vec<_>>(),
1241 ),
1242 recorded_at: recorded_at.clone(),
1243 };
1244 merge_unique_paths(
1245 &mut report.touched_paths,
1246 checkpoint.touched_paths.iter().cloned(),
1247 );
1248 merge_usage(&mut report.usage, &probe.usage_totals);
1249 if report.transcript_mode != "off" && report.transcript_refs.is_empty() {
1250 report.transcript_refs = probe.transcript_refs;
1251 }
1252 report.progress.push(checkpoint);
1253 if let Some(summary) = params.summary {
1254 report.summary = Some(summary);
1255 }
1256 report.last_progress_at = Some(recorded_at);
1257 mark_pending_flush(&mut report);
1258 self.persist_report(report)
1259 }
1260
1261 fn resolve_harness_thread_binding(
1262 &self,
1263 params: &OpenSessionParams,
1264 probe: &HarnessProbeResult,
1265 identity: &ResolvedIdentity,
1266 ) -> Result<(Option<String>, Option<String>)> {
1267 if let Some(thread) = params.thread.clone() {
1268 let thread_id = thread_id_for_name(&self.repo, Some(&thread))?;
1269 return Ok((Some(thread), thread_id));
1270 }
1271
1272 let current_attached = match self.repo.head_ref()? {
1273 Head::Attached { thread } => Some(thread.to_string()),
1274 Head::Detached { .. } => None,
1275 };
1276
1277 if !probe.attach_hints.root_actor
1278 && self.user_config.harness.threading.subagent
1279 == UserHarnessSubagentThreadPolicy::CreateChild
1280 && let Some(parent_thread) =
1281 resolve_parent_thread_for_subagent(&self.repo, probe, current_attached.as_deref())?
1282 && can_create_harness_thread(&self.repo, Some(&parent_thread), Some(&parent_thread))?
1283 {
1284 let name = allocate_thread_name(
1285 &self.repo,
1286 &format!(
1287 "{}/{}",
1288 parent_thread,
1289 sanitize_name(&preferred_thread_slug(params, probe, identity))
1290 ),
1291 )?;
1292 self.ensure_harness_thread(
1293 &name,
1294 Some(&parent_thread),
1295 Some(&parent_thread),
1296 params.task.clone(),
1297 )?;
1298 let thread_id = thread_id_for_name(&self.repo, Some(&name))?;
1299 return Ok((Some(name), thread_id));
1300 }
1301
1302 if probe.attach_hints.root_actor
1303 && self.user_config.harness.threading.root_actor
1304 == UserHarnessRootThreadPolicy::CreateNew
1305 && let Some(current) = current_attached.clone()
1306 && can_create_harness_thread(&self.repo, Some(¤t), None)?
1307 {
1308 let name = allocate_thread_name(
1309 &self.repo,
1310 &format!(
1311 "{}/{}",
1312 current,
1313 sanitize_name(&preferred_thread_slug(params, probe, identity))
1314 ),
1315 )?;
1316 self.ensure_harness_thread(&name, Some(¤t), None, params.task.clone())?;
1317 let thread_id = thread_id_for_name(&self.repo, Some(&name))?;
1318 return Ok((Some(name), thread_id));
1319 }
1320
1321 let thread_id = thread_id_for_name(&self.repo, current_attached.as_deref())?;
1322 Ok((current_attached, thread_id))
1323 }
1324
1325 fn ensure_harness_thread(
1326 &self,
1327 name: &str,
1328 target_thread: Option<&str>,
1329 parent_thread: Option<&str>,
1330 task: Option<String>,
1331 ) -> Result<()> {
1332 let manager = ThreadManager::new(self.repo.heddle_dir());
1333 if manager.load(name)?.is_some() {
1334 return Ok(());
1335 }
1336
1337 let base_state = self
1338 .resolve_harness_thread_base_state(target_thread, parent_thread)?
1339 .ok_or_else(|| anyhow!("No current state to start a thread from"))?;
1340 let tn = ThreadName::new(name);
1341 if self.repo.refs().get_thread(&tn)?.is_none() {
1342 self.repo
1343 .refs()
1344 .set_thread_cas(&tn, refs::RefExpectation::Missing, &base_state)?;
1345 self.repo.oplog().record_thread_create(
1350 &tn,
1351 &base_state,
1352 None,
1353 Some(&self.repo.op_scope()),
1354 )?;
1355 }
1356
1357 let workspace_mode = self
1358 .user_config
1359 .harness
1360 .threading
1361 .workspace_default
1362 .unwrap_or(UserThreadWorkspaceMode::Materialized);
1363 let thread_mode = match workspace_mode {
1364 UserThreadWorkspaceMode::Materialized | UserThreadWorkspaceMode::Auto => {
1365 ThreadMode::Materialized
1366 }
1367 UserThreadWorkspaceMode::Virtualized => ThreadMode::Virtualized,
1368 UserThreadWorkspaceMode::Solid => ThreadMode::Solid,
1369 };
1370 let path = match thread_mode {
1371 ThreadMode::Solid | ThreadMode::Materialized => {
1372 default_private_thread_path(&self.repo, name)
1373 }
1374 ThreadMode::Virtualized => default_private_thread_path(&self.repo, name),
1377 };
1378 let abs_path = prepare_worktree_target(&self.repo, &path, Some(name))?.path;
1379 write_isolated_checkout(&self.repo, &abs_path, &base_state, Some(name))?;
1380
1381 let base_state_obj = self
1382 .repo
1383 .store()
1384 .get_state(&base_state)?
1385 .ok_or_else(|| anyhow!("Base state '{}' not found", base_state.short()))?;
1386 let thread = Thread {
1387 id: name.to_string(),
1388 thread: name.to_string(),
1389 target_thread: target_thread.map(ToString::to_string),
1390 parent_thread: parent_thread.map(ToString::to_string),
1391 mode: thread_mode.clone(),
1392 state: ThreadState::Active,
1393 base_state: base_state.short(),
1394 base_root: base_state_obj.tree.short(),
1395 current_state: Some(base_state.short()),
1396 merged_state: None,
1397 task,
1398 execution_path: abs_path.clone(),
1399 materialized_path: match thread_mode {
1400 ThreadMode::Solid => Some(abs_path),
1401 ThreadMode::Materialized | ThreadMode::Virtualized => None,
1405 },
1406 changed_paths: vec![],
1407 impact_categories: vec![],
1408 heavy_impact_paths: vec![],
1409 promotion_suggested: false,
1410 freshness: if target_thread.is_some() {
1411 ThreadFreshness::Current
1412 } else {
1413 ThreadFreshness::Unknown
1414 },
1415 verification_summary: summarize_verification(base_state_obj.verification.as_ref()),
1416 confidence_summary: summarize_confidence(base_state_obj.confidence),
1417 integration_policy_result: ThreadIntegrationPolicy::default(),
1418 created_at: Utc::now(),
1419 updated_at: Utc::now(),
1420 ephemeral: None,
1421 auto: true,
1426 shared_target_dir: None,
1429 };
1430 manager.save(&thread)?;
1431 Ok(())
1432 }
1433
1434 fn resolve_harness_thread_base_state(
1435 &self,
1436 target_thread: Option<&str>,
1437 parent_thread: Option<&str>,
1438 ) -> Result<Option<objects::object::StateId>> {
1439 resolve_harness_thread_base_state(&self.repo, target_thread, parent_thread)
1440 }
1441
1442 fn record_usage(&mut self, params: RecordUsageParams) -> Result<SessionMutationResult> {
1443 let mut report = self
1444 .reports
1445 .load(¶ms.heddle_session_id)?
1446 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1447 if let Some(input) = params.input_tokens {
1448 report.usage.input_tokens = Some(max_u64(report.usage.input_tokens, input));
1449 }
1450 if let Some(output) = params.output_tokens {
1451 report.usage.output_tokens = Some(max_u64(report.usage.output_tokens, output));
1452 }
1453 if let Some(reasoning) = params.reasoning_tokens {
1454 report.usage.reasoning_tokens = Some(max_u64(report.usage.reasoning_tokens, reasoning));
1455 }
1456 if let Some(cache_creation) = params.cache_creation_tokens {
1457 report.usage.cache_creation_tokens =
1458 Some(max_u64(report.usage.cache_creation_tokens, cache_creation));
1459 }
1460 if let Some(cache_read) = params.cache_read_tokens {
1461 report.usage.cache_read_tokens =
1462 Some(max_u64(report.usage.cache_read_tokens, cache_read));
1463 }
1464 if let Some(tool_calls) = params.tool_calls {
1465 report.usage.tool_calls = Some(max_u32(report.usage.tool_calls, tool_calls));
1466 }
1467 if let Some(cost) = params.cost_micros_usd {
1468 report.usage.cost_micros_usd = Some(max_u64(report.usage.cost_micros_usd, cost));
1469 }
1470 mark_pending_flush(&mut report);
1471 self.persist_report(report)
1472 }
1473
1474 fn record_touched_paths(
1475 &mut self,
1476 params: RecordTouchedPathsParams,
1477 ) -> Result<SessionMutationResult> {
1478 let mut report = self
1479 .reports
1480 .load(¶ms.heddle_session_id)?
1481 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1482 merge_unique_paths(&mut report.touched_paths, normalize_paths(params.paths));
1483 mark_pending_flush(&mut report);
1484 self.persist_report(report)
1485 }
1486
1487 fn close_session(&mut self, params: CloseSessionParams) -> Result<CloseSessionResult> {
1488 let mut report = self
1489 .reports
1490 .load(¶ms.heddle_session_id)?
1491 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1492 report.closed_at = Some(Utc::now().to_rfc3339());
1493 report.outcome = params.outcome.clone();
1494 if let Some(summary) = params.summary {
1495 report.summary = Some(summary);
1496 }
1497 if let Some(transcript_refs) = params.transcript_refs {
1498 report.transcript_refs = transcript_refs;
1499 }
1500 let final_diff = compute_final_diff(
1501 &self.repo,
1502 report.base_state_at_open.as_deref(),
1503 &report.worktree_changes_at_open,
1504 )?;
1505 report.head_state_at_close = final_diff.head_state;
1506 report.changed_paths = final_diff.changed_paths;
1507 report.diff_summary = Some(final_diff.diff_summary);
1508 mark_pending_flush(&mut report);
1509 if report.owns_session {
1510 let mut sessions = SessionManager::new(self.repo.root());
1511 if let Ok(Some(session)) = sessions.get_session(&report.heddle_session_id)
1512 && session.is_active()
1513 {
1514 let _ = sessions.end_session(Some(&report.heddle_session_id));
1515 }
1516 }
1517
1518 let transport = params
1519 .transport
1520 .unwrap_or(self.user_config.harness.transport);
1521 if matches!(transport, HarnessTransport::Direct | HarnessTransport::End) {
1522 enqueue_report(&self.reports, &mut report)?;
1523 } else {
1524 self.reports.save(&report)?;
1525 }
1526 self.sync_registry_from_report(&report, ActorPresenceStatus::Complete)?;
1527 Ok(CloseSessionResult {
1528 heddle_session_id: report.heddle_session_id,
1529 changed_paths: report.changed_paths,
1530 diff_summary: report.diff_summary.unwrap_or_default(),
1531 report_flush_state: report.report_flush_state,
1532 })
1533 }
1534
1535 fn flush_reports(&mut self, params: FlushReportsParams) -> Result<FlushReportsResult> {
1536 let mut flushed = 0usize;
1537 let session_ids = match params.heddle_session_id {
1538 Some(session_id) => vec![session_id],
1539 None => self.reports.list_pending()?,
1540 };
1541 for session_id in session_ids {
1542 let Some(mut report) = self.reports.load(&session_id)? else {
1543 continue;
1544 };
1545 if !report.pending_flush {
1546 continue;
1547 }
1548 enqueue_report(&self.reports, &mut report)?;
1549 let status = if report.closed_at.is_some() {
1550 ActorPresenceStatus::Complete
1551 } else {
1552 ActorPresenceStatus::Active
1553 };
1554 self.sync_registry_from_report(&report, status)?;
1555 flushed += 1;
1556 }
1557 Ok(FlushReportsResult { flushed })
1558 }
1559
1560 fn persist_report(
1561 &mut self,
1562 mut report: SessionReportEnvelope,
1563 ) -> Result<SessionMutationResult> {
1564 let transport = transport_from_report(&report, self.user_config.harness.transport);
1565 match transport {
1566 HarnessTransport::Direct => {
1567 enqueue_report(&self.reports, &mut report)?;
1568 }
1569 HarnessTransport::Spool | HarnessTransport::End => {
1570 self.reports.save(&report)?;
1571 }
1572 }
1573 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1574 Ok(SessionMutationResult {
1575 heddle_session_id: report.heddle_session_id,
1576 heddle_segment_id: report.heddle_segment_id,
1577 report_flush_state: report.report_flush_state,
1578 })
1579 }
1580
1581 fn ensure_segment_for_report(
1582 &self,
1583 report: &mut SessionReportEnvelope,
1584 identity: &ResolvedIdentity,
1585 ) -> Result<()> {
1586 let mut sessions = SessionManager::new(self.repo.root());
1587 let Some(session) = sessions.get_session(&report.heddle_session_id)? else {
1588 return Ok(());
1589 };
1590 if !session.is_active() || !should_rotate_segment(&session, identity) {
1591 return Ok(());
1592 }
1593 let segment = sessions.add_segment(
1594 &report.heddle_session_id,
1595 identity
1596 .provider
1597 .clone()
1598 .unwrap_or_else(|| "unknown".to_string()),
1599 identity
1600 .model
1601 .clone()
1602 .unwrap_or_else(|| "unknown".to_string()),
1603 identity.policy.clone(),
1604 )?;
1605 report.heddle_segment_id = Some(segment.id);
1606 if identity.provider.is_some() {
1607 report.harness.provider = identity.provider.clone();
1608 }
1609 if identity.model.is_some() {
1610 report.harness.model = identity.model.clone();
1611 }
1612 if identity.policy.is_some() {
1613 report.harness.policy = identity.policy.clone();
1614 }
1615 if identity.thinking_level.is_some() {
1616 report.harness.thinking_level = identity.thinking_level.clone();
1617 }
1618 Ok(())
1619 }
1620
1621 fn ensure_registry_entry(&self, request: RegistryEntryRequest<'_>) -> Result<ActorPresence> {
1622 let RegistryEntryRequest {
1623 heddle_session_id,
1624 thread_name,
1625 thread_id,
1626 identity,
1627 probe,
1628 attach,
1629 client_instance_id,
1630 requested_entry,
1631 } = request;
1632 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
1633 let fallback_entry = if client_instance_id.is_some()
1634 || probe.native_actor_key.is_some()
1635 || probe.native_instance_key.is_some()
1636 {
1637 None
1638 } else {
1639 find_matching_registry_entry(®istry, &self.repo, heddle_session_id, thread_name)?
1640 };
1641 if let Some(entry) = requested_entry
1642 .cloned()
1643 .or_else(|| attach.matched_entry.clone())
1644 .or(fallback_entry)
1645 {
1646 return registry
1647 .update_entry(&entry.session_id, |existing| {
1648 if client_instance_id.is_some() {
1649 existing.client_instance_id = client_instance_id.map(ToString::to_string);
1650 }
1651 if probe.native_actor_key.is_some() {
1652 existing.native_actor_key = probe.native_actor_key.clone();
1653 }
1654 if probe.native_parent_actor_key.is_some() {
1655 existing.native_parent_actor_key = probe.native_parent_actor_key.clone();
1656 }
1657 if probe.native_instance_key.is_some() {
1658 existing.native_instance_key = probe.native_instance_key.clone();
1659 }
1660 existing.heddle_session_id = Some(heddle_session_id.to_string());
1661 existing.thread_id = thread_id.map(ToString::to_string);
1662 if let Some(thread_name) = thread_name {
1663 existing.thread = thread_name.to_string();
1664 }
1665 existing.path = Some(self.repo.root().to_path_buf());
1666 if identity.provider.is_some() {
1667 existing.provider = identity.provider.clone();
1668 }
1669 if identity.model.is_some() {
1670 existing.model = identity.model.clone();
1671 }
1672 if identity.harness.is_some() {
1673 existing.harness = identity.harness.clone();
1674 }
1675 if identity.thinking_level.is_some() {
1676 existing.thinking_level = identity.thinking_level.clone();
1677 }
1678 existing.attach_reason = Some(attach.attach_reason.clone());
1679 existing.attach_precedence = attach.precedence.clone();
1680 existing.winning_attach_rule = Some(attach.winning_rule.clone());
1681 existing.probe_source = probe.probe_source.clone();
1682 existing.probe_confidence = probe.confidence;
1683 existing.status = ActorPresenceStatus::Active;
1684 existing.completed_at = None;
1685 })?
1686 .ok_or_else(|| anyhow!("registry entry disappeared during update"));
1687 }
1688
1689 if client_instance_id.is_none() && probe.native_actor_key.is_some() {
1690 let (entry, _) = registry.find_or_create_active_entry(
1691 |entry| {
1692 claude_actor_compatible(entry, probe, self.repo.root())
1693 && entry.native_actor_key == probe.native_actor_key
1694 },
1695 |existing| {
1696 if client_instance_id.is_some() {
1697 existing.client_instance_id = client_instance_id.map(ToString::to_string);
1698 }
1699 if existing.heddle_session_id.is_none() {
1700 existing.heddle_session_id = Some(heddle_session_id.to_string());
1701 }
1702 existing.thread_id = thread_id.map(ToString::to_string);
1703 if let Some(thread_name) = thread_name {
1704 existing.thread = thread_name.to_string();
1705 }
1706 existing.path = Some(self.repo.root().to_path_buf());
1707 if identity.provider.is_some() {
1708 existing.provider = identity.provider.clone();
1709 }
1710 if identity.model.is_some() {
1711 existing.model = identity.model.clone();
1712 }
1713 if identity.harness.is_some() {
1714 existing.harness = identity.harness.clone();
1715 }
1716 if identity.thinking_level.is_some() {
1717 existing.thinking_level = identity.thinking_level.clone();
1718 }
1719 if probe.native_parent_actor_key.is_some() {
1720 existing.native_parent_actor_key = probe.native_parent_actor_key.clone();
1721 }
1722 if probe.native_instance_key.is_some() {
1723 existing.native_instance_key = probe.native_instance_key.clone();
1724 }
1725 existing.attach_reason = Some(attach.attach_reason.clone());
1726 existing.attach_precedence = attach.precedence.clone();
1727 existing.winning_attach_rule = Some(attach.winning_rule.clone());
1728 existing.probe_source = probe.probe_source.clone();
1729 existing.probe_confidence = probe.confidence;
1730 existing.status = ActorPresenceStatus::Active;
1731 existing.completed_at = None;
1732 },
1733 |session_id| {
1734 Ok(ActorPresence {
1735 session_id: session_id.to_string(),
1736 client_instance_id: client_instance_id.map(ToString::to_string),
1737 native_actor_key: probe.native_actor_key.clone(),
1738 native_parent_actor_key: probe.native_parent_actor_key.clone(),
1739 native_instance_key: probe.native_instance_key.clone(),
1740 heddle_session_id: Some(heddle_session_id.to_string()),
1741 thread_id: thread_id.map(ToString::to_string),
1742 thread: thread_name.unwrap_or("detached").to_string(),
1743 anchor_state: self.repo.head()?.map(|id| id.to_string_full()),
1744 anchor_root: None,
1745 path: Some(self.repo.root().to_path_buf()),
1746 base_state: self.repo.head()?.map(|id| id.short()).unwrap_or_default(),
1747 started_at: Utc::now(),
1748 provider: identity.provider.clone(),
1749 model: identity.model.clone(),
1750 harness: identity.harness.clone(),
1751 thinking_level: identity.thinking_level.clone(),
1752 usage_summary: AgentUsageSummary::default(),
1753 last_progress_at: None,
1754 report_flush_state: Some("pending-local".to_string()),
1755 attach_reason: Some(attach.attach_reason.clone()),
1756 task_assignment_id: None,
1757 attach_precedence: attach.precedence.clone(),
1758 winning_attach_rule: Some(attach.winning_rule.clone()),
1759 probe_source: probe.probe_source.clone(),
1760 probe_confidence: probe.confidence,
1761 status: ActorPresenceStatus::Active,
1762 completed_at: None,
1763 context_queries: vec![],
1764 })
1765 },
1766 )?;
1767 return Ok(entry);
1768 }
1769
1770 Ok(registry.create_generated_entry(|session_id| {
1771 Ok(ActorPresence {
1772 session_id: session_id.to_string(),
1773 client_instance_id: client_instance_id.map(ToString::to_string),
1774 native_actor_key: probe.native_actor_key.clone(),
1775 native_parent_actor_key: probe.native_parent_actor_key.clone(),
1776 native_instance_key: probe.native_instance_key.clone(),
1777 heddle_session_id: Some(heddle_session_id.to_string()),
1778 thread_id: thread_id.map(ToString::to_string),
1779 thread: thread_name.unwrap_or("detached").to_string(),
1780 anchor_state: self.repo.head()?.map(|id| id.to_string_full()),
1781 anchor_root: None,
1782 path: Some(self.repo.root().to_path_buf()),
1783 base_state: self.repo.head()?.map(|id| id.short()).unwrap_or_default(),
1784 started_at: Utc::now(),
1785 provider: identity.provider.clone(),
1786 model: identity.model.clone(),
1787 harness: identity.harness.clone(),
1788 thinking_level: identity.thinking_level.clone(),
1789 usage_summary: AgentUsageSummary::default(),
1790 last_progress_at: None,
1791 report_flush_state: Some("pending-local".to_string()),
1792 attach_reason: Some(attach.attach_reason.clone()),
1793 task_assignment_id: None,
1794 attach_precedence: attach.precedence.clone(),
1795 winning_attach_rule: Some(attach.winning_rule.clone()),
1796 probe_source: probe.probe_source.clone(),
1797 probe_confidence: probe.confidence,
1798 status: ActorPresenceStatus::Active,
1799 completed_at: None,
1800 context_queries: vec![],
1801 })
1802 })?)
1803 }
1804
1805 fn reuse_canonical_actor_session(
1806 &self,
1807 sessions: &mut SessionManager,
1808 request: CanonicalActorSessionRequest<'_>,
1809 ) -> Result<(Session, bool)> {
1810 let CanonicalActorSessionRequest {
1811 tentative_session,
1812 tentative_owns_session,
1813 entry,
1814 probe,
1815 attach,
1816 } = request;
1817 let Some(canonical_session_id) = entry.heddle_session_id.as_deref() else {
1818 return Ok((tentative_session, tentative_owns_session));
1819 };
1820 if canonical_session_id == tentative_session.id {
1821 return Ok((tentative_session, tentative_owns_session));
1822 }
1823
1824 if tentative_owns_session
1825 && let Ok(Some(session)) = sessions.get_session(&tentative_session.id)
1826 && session.is_active()
1827 {
1828 let _ = sessions.end_session(Some(&tentative_session.id));
1829 }
1830
1831 let canonical_session = sessions
1832 .get_session(canonical_session_id)?
1833 .ok_or_else(|| anyhow!("session not found: {canonical_session_id}"))?;
1834 let canonical_segment_id = canonical_session
1835 .current_segment_id
1836 .clone()
1837 .unwrap_or_default();
1838 sessions.set_current_session(canonical_session_id, &canonical_segment_id)?;
1839
1840 if let Some(native_actor_key) = probe
1841 .native_actor_key
1842 .as_deref()
1843 .or(entry.native_actor_key.as_deref())
1844 {
1845 attach.precedence.push(format!(
1846 "post-create-native-actor-key:{native_actor_key}:matched"
1847 ));
1848 attach.attach_reason = format!(
1849 "reused existing native actor {} on Heddle session {}",
1850 native_actor_key, canonical_session_id
1851 );
1852 attach.winning_rule = "native-actor-key-post-create".to_string();
1853 }
1854
1855 Ok((canonical_session, false))
1856 }
1857
1858 fn sync_registry_from_report(
1859 &self,
1860 report: &SessionReportEnvelope,
1861 status: ActorPresenceStatus,
1862 ) -> Result<()> {
1863 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
1864 let entry = if let Some(agent_session_id) = &report.agent_session_id {
1865 registry.update_entry(agent_session_id, |entry| {
1866 if report.client_instance_id.is_some() {
1867 entry.client_instance_id = report.client_instance_id.clone();
1868 }
1869 if report.native_actor_key.is_some() {
1870 entry.native_actor_key = report.native_actor_key.clone();
1871 }
1872 if report.native_parent_actor_key.is_some() {
1873 entry.native_parent_actor_key = report.native_parent_actor_key.clone();
1874 }
1875 if report.native_instance_key.is_some() {
1876 entry.native_instance_key = report.native_instance_key.clone();
1877 }
1878 entry.heddle_session_id = Some(report.heddle_session_id.clone());
1879 entry.path = Some(self.repo.root().to_path_buf());
1880 entry.harness = report.harness.harness.clone();
1881 entry.provider = report.harness.provider.clone();
1882 entry.model = report.harness.model.clone();
1883 entry.thinking_level = report.harness.thinking_level.clone();
1884 entry.usage_summary = usage_to_summary(&report.usage);
1885 entry.last_progress_at =
1886 report.last_progress_at.as_deref().and_then(parse_timestamp);
1887 entry.report_flush_state = report.report_flush_state.clone();
1888 entry.attach_reason = report.attach_reason.clone();
1889 entry.attach_precedence = report.attach_precedence.clone();
1890 entry.winning_attach_rule = report.winning_attach_rule.clone();
1891 entry.probe_source = report.probe_source.clone();
1892 entry.probe_confidence = report.probe_confidence;
1893 entry.status = status.clone();
1894 entry.completed_at = match status {
1895 ActorPresenceStatus::Active => None,
1896 ActorPresenceStatus::Abandoned | ActorPresenceStatus::Complete | ActorPresenceStatus::Merged => {
1897 Some(Utc::now())
1898 }
1899 };
1900 })?
1901 } else {
1902 None
1903 };
1904
1905 if entry.is_none() {
1906 let resolved = self.ensure_registry_entry(RegistryEntryRequest {
1907 heddle_session_id: &report.heddle_session_id,
1908 thread_name: report.thread.as_deref(),
1909 thread_id: report.thread_id.as_deref(),
1910 identity: &ResolvedIdentity {
1911 harness: report.harness.harness.clone(),
1912 provider: report.harness.provider.clone(),
1913 model: report.harness.model.clone(),
1914 thinking_level: report.harness.thinking_level.clone(),
1915 policy: report.harness.policy.clone(),
1916 },
1917 probe: &HarnessProbeResult {
1918 native_actor_key: report.native_actor_key.clone(),
1919 native_parent_actor_key: report.native_parent_actor_key.clone(),
1920 native_instance_key: report.native_instance_key.clone(),
1921 probe_source: report.probe_source.clone(),
1922 confidence: report.probe_confidence,
1923 ..HarnessProbeResult::default()
1924 },
1925 attach: &ResolvedAttachment {
1926 target: AttachTarget::CreateNew {
1927 _because_claimed: false,
1928 },
1929 matched_entry: None,
1930 attach_reason: report.attach_reason.clone().unwrap_or_else(|| {
1931 format!(
1932 "created actor for Heddle session {}",
1933 report.heddle_session_id
1934 )
1935 }),
1936 precedence: report.attach_precedence.clone(),
1937 winning_rule: report
1938 .winning_attach_rule
1939 .clone()
1940 .unwrap_or_else(|| "report-sync".to_string()),
1941 },
1942 client_instance_id: report.client_instance_id.as_deref(),
1943 requested_entry: None,
1944 })?;
1945 let mut report = report.clone();
1946 report.agent_session_id = Some(resolved.session_id);
1947 self.reports.save(&report)?;
1948 }
1949 Ok(())
1950 }
1951}
1952
1953#[derive(Debug, Clone, Default)]
1954struct ResolvedIdentity {
1955 harness: Option<String>,
1956 provider: Option<String>,
1957 model: Option<String>,
1958 thinking_level: Option<String>,
1959 policy: Option<String>,
1960}
1961
1962impl ResolvedIdentity {
1963 fn to_transport_identity(&self) -> HarnessIdentity {
1964 HarnessIdentity {
1965 harness: self.harness.clone(),
1966 provider: self.provider.clone(),
1967 model: self.model.clone(),
1968 thinking_level: self.thinking_level.clone(),
1969 policy: self.policy.clone(),
1970 }
1971 }
1972}
1973
1974struct IdentityHints {
1975 harness: Option<String>,
1976 provider: Option<String>,
1977 model: Option<String>,
1978 thinking_level: Option<String>,
1979 policy: Option<String>,
1980 probe: HarnessProbeResult,
1981}
1982
1983fn resolve_identity(
1984 repo: &Repository,
1985 user_config: &UserConfig,
1986 hints: IdentityHints,
1987) -> Result<ResolvedIdentity> {
1988 let current_session = SessionManager::new(repo.root()).get_current_session()?;
1989 let current_segment = current_session
1990 .as_ref()
1991 .and_then(|session| session.current_segment());
1992 let token_claims = if user_config.harness.auto_infer {
1993 user_config_token_claims(user_config)
1994 } else {
1995 None
1996 };
1997 let harness_override = resolved_harness_override(
1998 user_config,
1999 hints.harness.as_deref(),
2000 hints.probe.harness.as_deref(),
2001 );
2002
2003 Ok(ResolvedIdentity {
2004 harness: hints.harness.or(hints.probe.harness),
2005 provider: hints
2006 .provider
2007 .or(hints.probe.provider)
2008 .or_else(|| current_segment.map(|segment| segment.provider.clone()))
2009 .or_else(|| {
2010 token_claims
2011 .as_ref()
2012 .and_then(|claims| claims.agent_provider.clone())
2013 })
2014 .or_else(|| harness_override.and_then(|entry| entry.provider.clone()))
2015 .or_else(|| user_config.agent.provider.clone()),
2016 model: hints
2017 .model
2018 .or(hints.probe.model)
2019 .or_else(|| current_segment.map(|segment| segment.model.clone()))
2020 .or_else(|| {
2021 token_claims
2022 .as_ref()
2023 .and_then(|claims| claims.agent_model.clone())
2024 })
2025 .or_else(|| harness_override.and_then(|entry| entry.model.clone()))
2026 .or_else(|| user_config.agent.model.clone()),
2027 thinking_level: hints
2028 .thinking_level
2029 .or(hints.probe.thinking_level)
2030 .or_else(|| harness_override.and_then(|entry| entry.thinking_level.clone())),
2031 policy: hints
2032 .policy
2033 .or(hints.probe.policy)
2034 .or_else(|| current_segment.and_then(|segment| segment.policy_id.clone()))
2035 .or_else(|| harness_override.and_then(|entry| entry.policy.clone()))
2036 .or_else(|| user_config.agent.default_policy.clone()),
2037 })
2038}
2039
2040fn resolved_harness_override<'a>(
2041 user_config: &'a UserConfig,
2042 explicit: Option<&str>,
2043 fingerprint: Option<&str>,
2044) -> Option<&'a UserHarnessOverride> {
2045 explicit
2046 .and_then(|name| user_config.harness.harnesses.get(name))
2047 .or_else(|| fingerprint.and_then(|name| user_config.harness.harnesses.get(name)))
2048}
2049
2050enum AttachTarget {
2051 ExistingSession(objects::object::Session),
2052 CreateNew { _because_claimed: bool },
2053}
2054
2055struct ResolvedAttachment {
2056 target: AttachTarget,
2057 matched_entry: Option<ActorPresence>,
2058 attach_reason: String,
2059 precedence: Vec<String>,
2060 winning_rule: String,
2061}
2062
2063fn resolve_actor_attachment(
2064 registry: &ActorPresenceStore,
2065 repo: &Repository,
2066 sessions: &mut SessionManager,
2067 input: AttachmentResolutionInput<'_>,
2068) -> Result<ResolvedAttachment> {
2069 let AttachmentResolutionInput {
2070 requested_entry,
2071 explicit_heddle_session_id,
2072 client_instance_id,
2073 probe,
2074 token_claims,
2075 } = input;
2076
2077 let mut sessions_by_id: BTreeMap<String, Session> = BTreeMap::new();
2079 let mut matched_by_session: BTreeMap<String, ActorPresence> = BTreeMap::new();
2080 let mut facts = SessionAttachFacts {
2081 root_actor: probe.attach_hints.root_actor,
2082 ..SessionAttachFacts::default()
2083 };
2084
2085 if let Some(entry) = requested_entry
2086 && let Some(bound_session_id) = entry.heddle_session_id.as_deref()
2087 {
2088 let session = sessions
2089 .get_session(bound_session_id)?
2090 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2091 if !session.is_active() {
2092 return Err(anyhow!("session is not active: {bound_session_id}"));
2093 }
2094 matched_by_session.insert(session.id.clone(), entry.clone());
2095 sessions_by_id.insert(session.id.clone(), session);
2096 facts.explicit_agent = Some(ExplicitAgentBind {
2097 agent_session_id: entry.session_id.clone(),
2098 heddle_session_id: bound_session_id.to_string(),
2099 });
2100 }
2101
2102 if facts.explicit_agent.is_none()
2103 && let Some(session_id) = explicit_heddle_session_id
2104 {
2105 ensure_requested_entry_matches_session(requested_entry, session_id)?;
2106 let session = sessions
2107 .get_session(session_id)?
2108 .ok_or_else(|| anyhow!("session not found: {session_id}"))?;
2109 if !session.is_active() {
2110 return Err(anyhow!("session is not active: {session_id}"));
2111 }
2112 sessions_by_id.insert(session.id.clone(), session);
2113 facts.explicit_heddle_session_id = Some(session_id.to_string());
2114 }
2115
2116 if client_instance_id.is_none()
2117 && let Some(native_actor_key) = probe.native_actor_key.as_deref()
2118 {
2119 if let Some(entry) = registry.find_active_by_native_actor_key(native_actor_key)?
2120 && claude_actor_compatible(&entry, probe, repo.root())
2121 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2122 {
2123 let session = sessions
2124 .get_session(&bound_session_id)?
2125 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2126 if session.is_active() {
2127 matched_by_session.insert(session.id.clone(), entry);
2128 sessions_by_id.insert(session.id.clone(), session);
2129 facts.native_actor = SessionLookupFact::Hit {
2130 key: native_actor_key.to_string(),
2131 session_id: bound_session_id,
2132 };
2133 } else {
2134 facts.native_actor = SessionLookupFact::Miss {
2135 key: native_actor_key.to_string(),
2136 };
2137 }
2138 } else {
2139 facts.native_actor = SessionLookupFact::Miss {
2140 key: native_actor_key.to_string(),
2141 };
2142 }
2143 }
2144
2145 if let Some(client_instance_id) = client_instance_id {
2146 if let Some(entry) = registry.find_active_by_client_instance_id(client_instance_id)?
2147 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2148 {
2149 let session = sessions
2150 .get_session(&bound_session_id)?
2151 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2152 if session.is_active() {
2153 matched_by_session.insert(session.id.clone(), entry);
2154 sessions_by_id.insert(session.id.clone(), session);
2155 facts.client_instance = SessionLookupFact::Hit {
2156 key: client_instance_id.to_string(),
2157 session_id: bound_session_id,
2158 };
2159 } else {
2160 facts.client_instance = SessionLookupFact::Miss {
2161 key: client_instance_id.to_string(),
2162 };
2163 }
2164 } else {
2165 facts.client_instance = SessionLookupFact::Miss {
2166 key: client_instance_id.to_string(),
2167 };
2168 }
2169 }
2170
2171 if let Some(native_instance_key) = probe.native_instance_key.as_deref() {
2172 if let Some(entry) =
2173 registry.find_active_by_native_instance_key_at_path(native_instance_key, repo.root())?
2174 && claude_actor_compatible(&entry, probe, repo.root())
2175 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2176 {
2177 let session = sessions
2178 .get_session(&bound_session_id)?
2179 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2180 if session.is_active() {
2181 matched_by_session.insert(session.id.clone(), entry);
2182 sessions_by_id.insert(session.id.clone(), session);
2183 facts.native_instance = SessionLookupFact::Hit {
2184 key: native_instance_key.to_string(),
2185 session_id: bound_session_id,
2186 };
2187 } else {
2188 facts.native_instance = SessionLookupFact::Miss {
2189 key: native_instance_key.to_string(),
2190 };
2191 }
2192 } else {
2193 facts.native_instance = SessionLookupFact::Miss {
2194 key: native_instance_key.to_string(),
2195 };
2196 }
2197 }
2198
2199 if probe.attach_hints.root_actor
2200 && let Some(current) = sessions.get_current_session()?
2201 && current.is_active()
2202 {
2203 let claimed = session_claimed_by_other(
2204 registry,
2205 ¤t.id,
2206 requested_entry,
2207 client_instance_id,
2208 probe.native_actor_key.as_deref(),
2209 )?;
2210 sessions_by_id
2211 .entry(current.id.clone())
2212 .or_insert_with(|| current.clone());
2213 facts.current_worktree = if claimed {
2214 WorktreeSessionFact::Claimed {
2215 session_id: current.id.clone(),
2216 }
2217 } else {
2218 WorktreeSessionFact::Available {
2219 session_id: current.id.clone(),
2220 }
2221 };
2222 }
2223
2224 if let Some(claims) = token_claims
2225 && let Some(token_sid) = claims.sid.as_deref()
2226 && let Some(session) = sessions.get_session(token_sid)?
2227 && session.is_active()
2228 {
2229 let claimed = session_claimed_by_other(
2230 registry,
2231 &session.id,
2232 requested_entry,
2233 client_instance_id,
2234 probe.native_actor_key.as_deref(),
2235 )?;
2236 let session_id = session.id.clone();
2237 sessions_by_id.insert(session_id.clone(), session);
2238 facts.token_sid = if claimed {
2239 TokenSidFact::Claimed { session_id }
2240 } else {
2241 TokenSidFact::Available { session_id }
2242 };
2243 }
2244
2245 let decision = decide_session_attach(&facts);
2246 match decision.policy {
2247 SessionPolicy::AttachExisting { session_id, .. } => {
2248 let session = sessions_by_id
2249 .remove(&session_id)
2250 .or_else(|| sessions.get_session(&session_id).ok().flatten())
2251 .ok_or_else(|| anyhow!("session not found: {session_id}"))?;
2252 Ok(ResolvedAttachment {
2253 target: AttachTarget::ExistingSession(session),
2254 matched_entry: matched_by_session.remove(&session_id),
2255 attach_reason: decision.attach_reason,
2256 precedence: decision.precedence,
2257 winning_rule: decision.winning_rule.to_string(),
2258 })
2259 }
2260 SessionPolicy::CreateNew {
2261 because_claimed, ..
2262 } => Ok(ResolvedAttachment {
2263 target: AttachTarget::CreateNew {
2264 _because_claimed: because_claimed,
2265 },
2266 matched_entry: None,
2267 attach_reason: decision.attach_reason,
2268 precedence: decision.precedence,
2269 winning_rule: decision.winning_rule.to_string(),
2270 }),
2271 }
2272}
2273
2274fn claude_actor_compatible(
2275 entry: &ActorPresence,
2276 probe: &HarnessProbeResult,
2277 repo_root: &Path,
2278) -> bool {
2279 let Some(native_actor_key) = probe.native_actor_key.as_deref() else {
2280 return true;
2281 };
2282 if !native_actor_key.starts_with("claude-code:") {
2283 return true;
2284 }
2285 if native_actor_key.starts_with("claude-code:agent:") {
2286 return entry.native_actor_key.as_deref() == Some(native_actor_key);
2287 }
2288 if let Some(native_instance_key) = probe.native_instance_key.as_deref() {
2289 return entry.native_actor_key.as_deref() == Some(native_actor_key)
2290 && entry.native_instance_key.as_deref() == Some(native_instance_key);
2291 }
2292 let same_repo = entry
2293 .path
2294 .as_ref()
2295 .map(|path| path.canonicalize().unwrap_or_else(|_| path.clone()))
2296 .unwrap_or_default()
2297 == repo_root
2298 .canonicalize()
2299 .unwrap_or_else(|_| repo_root.to_path_buf());
2300 entry.native_actor_key.as_deref() == Some(native_actor_key)
2301 && same_repo
2302 && probe.confidence.unwrap_or_default() >= 0.9
2303}
2304
2305fn decode_token_claims(token: &str) -> Option<TokenClaims> {
2306 let payload = token.split('.').nth(1)?;
2307 let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
2308 .decode(payload.as_bytes())
2309 .ok()?;
2310 serde_json::from_slice(&decoded).ok()
2311}
2312
2313fn user_config_token_claims(user_config: &UserConfig) -> Option<TokenClaims> {
2314 user_config
2315 .remote_token()
2316 .ok()
2317 .flatten()
2318 .and_then(|token| decode_token_claims(&token.id))
2319}
2320
2321#[derive(Debug, Deserialize)]
2322struct TokenClaims {
2323 #[serde(default)]
2324 sid: Option<String>,
2325 #[serde(default)]
2326 agent_provider: Option<String>,
2327 #[serde(default)]
2328 agent_model: Option<String>,
2329}
2330
2331fn should_rotate_segment(session: &objects::object::Session, identity: &ResolvedIdentity) -> bool {
2332 let Some(segment) = session.current_segment() else {
2333 return false;
2334 };
2335 pure_should_rotate_segment(
2336 Some(segment.provider.as_str()),
2337 Some(segment.model.as_str()),
2338 identity.provider.as_deref(),
2339 identity.model.as_deref(),
2340 )
2341}
2342
2343fn thread_id_for_name(repo: &Repository, thread_name: Option<&str>) -> Result<Option<String>> {
2344 let Some(thread_name) = thread_name else {
2345 return Ok(None);
2346 };
2347 Ok(ThreadManager::new(repo.heddle_dir())
2348 .load(thread_name)?
2349 .map(|thread| thread.id))
2350}
2351
2352fn can_create_harness_thread(
2353 repo: &Repository,
2354 target_thread: Option<&str>,
2355 parent_thread: Option<&str>,
2356) -> Result<bool> {
2357 Ok(resolve_harness_thread_base_state(repo, target_thread, parent_thread)?.is_some())
2358}
2359
2360fn resolve_harness_thread_base_state(
2361 repo: &Repository,
2362 target_thread: Option<&str>,
2363 parent_thread: Option<&str>,
2364) -> Result<Option<objects::object::StateId>> {
2365 if let Some(head_state) = repo.head()? {
2366 return Ok(Some(head_state));
2367 }
2368
2369 for thread_name in [parent_thread, target_thread].into_iter().flatten() {
2370 if let Some(state) = resolve_named_thread_base_state(repo, thread_name)? {
2371 return Ok(Some(state));
2372 }
2373 }
2374
2375 Ok(None)
2376}
2377
2378fn resolve_named_thread_base_state(
2379 repo: &Repository,
2380 thread_name: &str,
2381) -> Result<Option<objects::object::StateId>> {
2382 if let Some(thread) = ThreadManager::new(repo.heddle_dir()).load(thread_name)?
2383 && let Some(state_spec) = thread
2384 .current_state
2385 .as_deref()
2386 .or(Some(thread.base_state.as_str()))
2387 && let Some(state_id) = repo
2388 .resolve_state(state_spec)?
2389 .or_else(|| objects::object::StateId::parse(state_spec).ok())
2390 {
2391 return Ok(Some(state_id));
2392 }
2393
2394 Ok(repo.refs().get_thread(&ThreadName::new(thread_name))?)
2395}
2396
2397fn resolve_parent_thread_for_subagent(
2398 repo: &Repository,
2399 probe: &HarnessProbeResult,
2400 current_attached: Option<&str>,
2401) -> Result<Option<String>> {
2402 if let Some(parent_key) = probe.native_parent_actor_key.as_deref() {
2403 let registry = ActorPresenceStore::new(repo.heddle_dir());
2404 if let Some(entry) = registry.find_active_by_native_actor_key(parent_key)? {
2405 return Ok(Some(entry.thread));
2406 }
2407 }
2408 Ok(current_attached.map(ToString::to_string))
2409}
2410
2411fn preferred_thread_slug(
2412 params: &OpenSessionParams,
2413 probe: &HarnessProbeResult,
2414 identity: &ResolvedIdentity,
2415) -> String {
2416 params
2417 .task
2418 .clone()
2419 .or_else(|| params.summary.clone())
2420 .or_else(|| probe.native_actor_key.as_deref().map(native_key_slug))
2421 .or_else(|| probe.native_instance_key.as_deref().map(native_key_slug))
2422 .or_else(|| identity.harness.clone())
2423 .unwrap_or_else(|| "work".to_string())
2424}
2425
2426fn native_key_slug(value: &str) -> String {
2427 value
2428 .rsplit(':')
2429 .next()
2430 .map(ToString::to_string)
2431 .unwrap_or_else(|| value.to_string())
2432}
2433
2434fn allocate_thread_name(repo: &Repository, base: &str) -> Result<String> {
2435 if ThreadManager::new(repo.heddle_dir()).load(base)?.is_none()
2436 && repo.refs().get_thread(&ThreadName::new(base))?.is_none()
2437 {
2438 return Ok(base.to_string());
2439 }
2440 for idx in 2..1000 {
2441 let candidate = format!("{base}-{idx}");
2442 if ThreadManager::new(repo.heddle_dir())
2443 .load(&candidate)?
2444 .is_none()
2445 && repo
2446 .refs()
2447 .get_thread(&ThreadName::new(&candidate))?
2448 .is_none()
2449 {
2450 return Ok(candidate);
2451 }
2452 }
2453 Err(anyhow!(
2454 "could not allocate a unique thread name from '{base}'"
2455 ))
2456}
2457
2458fn default_private_thread_path(repo: &Repository, name: &str) -> PathBuf {
2459 repo.managed_checkout_path(name)
2467}
2468
2469fn sanitize_name(name: &str) -> String {
2470 let mut out = String::new();
2471 let mut last_dash = false;
2472 for ch in name.chars() {
2473 if ch.is_ascii_alphanumeric() {
2474 out.push(ch.to_ascii_lowercase());
2475 last_dash = false;
2476 } else if !last_dash {
2477 out.push('-');
2478 last_dash = true;
2479 }
2480 }
2481 out.trim_matches('-').to_string()
2482}
2483
2484fn resolve_requested_registry_entry(
2485 registry: &ActorPresenceStore,
2486 agent_session_id: Option<&str>,
2487 client_instance_id: Option<&str>,
2488) -> Result<Option<ActorPresence>> {
2489 if let Some(agent_session_id) = agent_session_id {
2490 let entry = registry
2491 .load(agent_session_id)?
2492 .ok_or_else(|| anyhow!("agent session not found: {agent_session_id}"))?;
2493 if entry.status != ActorPresenceStatus::Active {
2494 return Err(anyhow!("agent session is not active: {agent_session_id}"));
2495 }
2496 return Ok(Some(entry));
2497 }
2498
2499 if let Some(client_instance_id) = client_instance_id {
2500 return Ok(registry.find_active_by_client_instance_id(client_instance_id)?);
2501 }
2502
2503 Ok(None)
2504}
2505
2506fn ensure_requested_entry_matches_session(
2507 requested_entry: Option<&ActorPresence>,
2508 heddle_session_id: &str,
2509) -> Result<()> {
2510 if let Some(entry) = requested_entry
2511 && let Some(bound_session_id) = entry.heddle_session_id.as_deref()
2512 && bound_session_id != heddle_session_id
2513 {
2514 return Err(anyhow!(
2515 "requested agent is already bound to a different heddle session: {}",
2516 entry.session_id
2517 ));
2518 }
2519 Ok(())
2520}
2521
2522fn session_claimed_by_other(
2523 registry: &ActorPresenceStore,
2524 heddle_session_id: &str,
2525 requested_entry: Option<&ActorPresence>,
2526 client_instance_id: Option<&str>,
2527 native_actor_key: Option<&str>,
2528) -> Result<bool> {
2529 if requested_entry.is_none() && client_instance_id.is_none() && native_actor_key.is_none() {
2530 return Ok(false);
2531 }
2532
2533 let Some(existing) = registry.find_active_by_heddle_session_id(heddle_session_id)? else {
2534 return Ok(false);
2535 };
2536 if let Some(requested) = requested_entry {
2537 return Ok(requested.session_id != existing.session_id);
2538 }
2539 if let Some(client_instance_id) = client_instance_id
2540 && existing.client_instance_id.as_deref() == Some(client_instance_id)
2541 {
2542 return Ok(false);
2543 }
2544 if let Some(native_actor_key) = native_actor_key
2545 && existing.native_actor_key.as_deref() == Some(native_actor_key)
2546 {
2547 return Ok(false);
2548 }
2549 Ok(true)
2550}
2551
2552fn find_matching_registry_entry(
2553 registry: &ActorPresenceStore,
2554 repo: &Repository,
2555 heddle_session_id: &str,
2556 thread_name: Option<&str>,
2557) -> Result<Option<ActorPresence>> {
2558 if let Some(entry) = registry.find_active_by_heddle_session_id(heddle_session_id)? {
2559 return Ok(Some(entry));
2560 }
2561 let canonical_root = repo
2562 .root()
2563 .canonicalize()
2564 .unwrap_or_else(|_| repo.root().to_path_buf());
2565 Ok(registry
2566 .list()?
2567 .into_iter()
2568 .filter(|entry| entry.status == ActorPresenceStatus::Active)
2569 .find(|entry| {
2570 entry
2571 .path
2572 .as_ref()
2573 .map(|path| path.canonicalize().unwrap_or_else(|_| path.clone()) == canonical_root)
2574 .unwrap_or(false)
2575 || thread_name.is_some_and(|thread| entry.thread == thread)
2576 }))
2577}
2578
2579fn merged_env_hints(extra: &BTreeMap<String, String>) -> BTreeMap<String, String> {
2580 let mut merged: BTreeMap<String, String> = std::env::vars()
2581 .filter(|(key, _)| inherited_harness_hint(key))
2582 .collect();
2583 for (key, value) in extra {
2584 merged.insert(key.clone(), value.clone());
2585 }
2586 merged
2587}
2588
2589fn inherited_harness_hint(key: &str) -> bool {
2590 if matches!(
2591 key,
2592 "OPENAI_MODEL"
2593 | "ANTHROPIC_MODEL"
2594 | "CLAUDE_MODEL"
2595 | "MODEL"
2596 | "OPENAI_REASONING_EFFORT"
2597 | "REASONING_EFFORT"
2598 | "THINKING_LEVEL"
2599 | "PROMPT_POLICY"
2600 ) {
2601 return false;
2602 }
2603
2604 key.starts_with("HEDDLE_")
2605 || key.starts_with("CODEX_")
2606 || key == "CLAUDECODE"
2607 || key.starts_with("OPENCODE_")
2608}
2609
2610fn to_json_value<T: Serialize>(value: T) -> Result<Value> {
2611 serde_json::to_value(value).map_err(|err| anyhow!(err))
2612}
2613
2614fn normalize_paths<I>(paths: I) -> Vec<String>
2615where
2616 I: IntoIterator<Item = String>,
2617{
2618 let mut ordered = BTreeSet::new();
2619 for path in paths {
2620 let normalized = path.trim().replace('\\', "/");
2621 if !normalized.is_empty() {
2622 ordered.insert(normalized);
2623 }
2624 }
2625 ordered.into_iter().collect()
2626}
2627
2628fn merge_unique_paths<I>(target: &mut Vec<String>, paths: I)
2629where
2630 I: IntoIterator<Item = String>,
2631{
2632 let mut merged: BTreeSet<String> = target.iter().cloned().collect();
2633 merged.extend(paths);
2634 *target = merged.into_iter().collect();
2635}
2636
2637fn max_u64(current: Option<u64>, candidate: u64) -> u64 {
2638 current
2639 .map(|value| value.max(candidate))
2640 .unwrap_or(candidate)
2641}
2642
2643fn max_u32(current: Option<u32>, candidate: u32) -> u32 {
2644 current
2645 .map(|value| value.max(candidate))
2646 .unwrap_or(candidate)
2647}
2648
2649fn merge_usage(target: &mut UsageTotals, incoming: &UsageTotals) {
2650 if let Some(input) = incoming.input_tokens {
2651 target.input_tokens = Some(max_u64(target.input_tokens, input));
2652 }
2653 if let Some(output) = incoming.output_tokens {
2654 target.output_tokens = Some(max_u64(target.output_tokens, output));
2655 }
2656 if let Some(reasoning) = incoming.reasoning_tokens {
2657 target.reasoning_tokens = Some(max_u64(target.reasoning_tokens, reasoning));
2658 }
2659 if let Some(cache_creation) = incoming.cache_creation_tokens {
2660 target.cache_creation_tokens = Some(max_u64(target.cache_creation_tokens, cache_creation));
2661 }
2662 if let Some(cache_read) = incoming.cache_read_tokens {
2663 target.cache_read_tokens = Some(max_u64(target.cache_read_tokens, cache_read));
2664 }
2665 if let Some(tool_calls) = incoming.tool_calls {
2666 target.tool_calls = Some(max_u32(target.tool_calls, tool_calls));
2667 }
2668 if let Some(cost) = incoming.cost_micros_usd {
2669 target.cost_micros_usd = Some(max_u64(target.cost_micros_usd, cost));
2670 }
2671}
2672
2673fn parse_timestamp(value: &str) -> Option<chrono::DateTime<Utc>> {
2674 chrono::DateTime::parse_from_rfc3339(value)
2675 .ok()
2676 .map(|dt| dt.with_timezone(&Utc))
2677}
2678
2679fn transport_from_report(
2680 report: &SessionReportEnvelope,
2681 fallback: HarnessTransport,
2682) -> HarnessTransport {
2683 match report.transport_mode.as_str() {
2684 "spool" => HarnessTransport::Spool,
2685 "direct" => HarnessTransport::Direct,
2686 "end" => HarnessTransport::End,
2687 _ => fallback,
2688 }
2689}
2690
2691fn mark_pending_flush(report: &mut SessionReportEnvelope) {
2692 report.pending_flush = true;
2693 report.report_flush_state = Some("pending-local".to_string());
2694}
2695
2696fn enqueue_report(store: &SessionReportStore, report: &mut SessionReportEnvelope) -> Result<()> {
2697 store.append_outbox(report)?;
2698 report.pending_flush = false;
2699 let flushed_at = Utc::now().to_rfc3339();
2700 report.last_flushed_at = Some(flushed_at);
2701 report.report_flush_state = Some("queued-local".to_string());
2702 store.save(report)?;
2703 Ok(())
2704}
2705
2706fn usage_to_summary(usage: &UsageTotals) -> AgentUsageSummary {
2707 AgentUsageSummary {
2708 input_tokens: usage.input_tokens,
2709 output_tokens: usage.output_tokens,
2710 reasoning_tokens: usage.reasoning_tokens,
2711 tool_calls: usage.tool_calls,
2712 cost_micros_usd: usage.cost_micros_usd,
2713 }
2714}
2715
2716fn transcript_mode_name(mode: HarnessTranscriptMode) -> &'static str {
2717 match mode {
2718 HarnessTranscriptMode::Off => "off",
2719 HarnessTranscriptMode::Summary => "summary",
2720 HarnessTranscriptMode::Full => "full",
2721 }
2722}
2723
2724fn transport_mode_name(mode: HarnessTransport) -> &'static str {
2725 match mode {
2726 HarnessTransport::Spool => "spool",
2727 HarnessTransport::Direct => "direct",
2728 HarnessTransport::End => "end",
2729 }
2730}
2731
2732struct FinalDiff {
2733 changed_paths: Vec<String>,
2734 diff_summary: SessionDiffSummary,
2735 head_state: Option<String>,
2736}
2737
2738fn compute_final_diff(
2739 repo: &Repository,
2740 base_state: Option<&str>,
2741 worktree_baseline: &[WorktreeChangeBaseline],
2742) -> Result<FinalDiff> {
2743 let mut changes: BTreeMap<String, DiffKind> = BTreeMap::new();
2744
2745 let head_state = repo.head()?;
2746 if let (Some(base_spec), Some(head_id)) = (base_state, head_state) {
2747 let base_id = repo
2748 .resolve_state(base_spec)?
2749 .or_else(|| objects::object::StateId::parse(base_spec).ok());
2750 if let Some(base_id) = base_id
2751 && base_id != head_id
2752 {
2753 let Some(base_state_obj) = repo.store().get_state(&base_id)? else {
2754 return Err(anyhow!("base state not found: {base_spec}"));
2755 };
2756 let Some(head_state_obj) = repo.store().get_state(&head_id)? else {
2757 return Err(anyhow!("head state not found: {}", head_id.short()));
2758 };
2759 for change in repo.diff_trees(&base_state_obj.tree, &head_state_obj.tree)? {
2760 changes.insert(change.path, change.kind);
2761 }
2762 }
2763 }
2764
2765 let baseline_paths: BTreeSet<(String, String)> = worktree_baseline
2766 .iter()
2767 .map(|change| (change.path.clone(), change.kind.clone()))
2768 .collect();
2769 for (path, kind) in collect_worktree_changes(repo)? {
2770 let kind_name = diff_kind_name(kind);
2771 if !baseline_paths.contains(&(path.clone(), kind_name.to_string())) {
2772 changes.insert(path, kind);
2773 }
2774 }
2775
2776 let diff_summary = SessionDiffSummary {
2777 changed_file_count: changes.len() as u32,
2778 added_files: changes
2779 .values()
2780 .filter(|kind| **kind == DiffKind::Added)
2781 .count() as u32,
2782 modified_files: changes
2783 .values()
2784 .filter(|kind| **kind == DiffKind::Modified)
2785 .count() as u32,
2786 deleted_files: changes
2787 .values()
2788 .filter(|kind| **kind == DiffKind::Deleted)
2789 .count() as u32,
2790 };
2791
2792 Ok(FinalDiff {
2793 changed_paths: changes.into_keys().collect(),
2794 diff_summary,
2795 head_state: head_state.map(|id| id.to_string_full()),
2796 })
2797}
2798
2799fn capture_worktree_change_snapshot(repo: &Repository) -> Result<Vec<WorktreeChangeBaseline>> {
2800 Ok(collect_worktree_changes(repo)?
2801 .into_iter()
2802 .map(|(path, kind)| WorktreeChangeBaseline {
2803 path,
2804 kind: diff_kind_name(kind).to_string(),
2805 })
2806 .collect())
2807}
2808
2809fn collect_worktree_changes(repo: &Repository) -> Result<BTreeMap<String, DiffKind>> {
2810 let status_options = worktree_status_options(Some(repo.config()));
2811 let worktree_tree = match repo.current_state()? {
2812 Some(state) => repo.require_tree(&state.tree)?,
2813 None => Tree::new(),
2814 };
2815 let status = repo.compare_worktree_cached_with_options(&worktree_tree, &status_options)?;
2816 let mut changes = BTreeMap::new();
2817 for path in status.added {
2818 changes.insert(path.display().to_string(), DiffKind::Added);
2819 }
2820 for path in status.modified {
2821 changes.insert(path.display().to_string(), DiffKind::Modified);
2822 }
2823 for path in status.deleted {
2824 changes.insert(path.display().to_string(), DiffKind::Deleted);
2825 }
2826 Ok(changes)
2827}
2828
2829fn changed_paths_between_states(
2830 repo: &Repository,
2831 before_state: StateId,
2832 after_state: StateId,
2833) -> Result<Vec<String>> {
2834 if before_state == after_state {
2835 return Ok(Vec::new());
2836 }
2837 let Some(before_state_obj) = repo.store().get_state(&before_state)? else {
2838 return Err(anyhow!(
2839 "timeline before state not found: {}",
2840 before_state.short()
2841 ));
2842 };
2843 let Some(after_state_obj) = repo.store().get_state(&after_state)? else {
2844 return Err(anyhow!(
2845 "timeline after state not found: {}",
2846 after_state.short()
2847 ));
2848 };
2849 let mut paths = BTreeSet::new();
2850 for change in repo.diff_trees(&before_state_obj.tree, &after_state_obj.tree)? {
2851 paths.insert(change.path);
2852 }
2853 Ok(paths.into_iter().collect())
2854}
2855
2856fn diff_kind_name(kind: DiffKind) -> &'static str {
2857 match kind {
2858 DiffKind::Added => "added",
2859 DiffKind::Modified => "modified",
2860 DiffKind::Deleted => "deleted",
2861 DiffKind::Unchanged => "unchanged",
2862 }
2863}
2864
2865struct SessionReportStore {
2866 dir: PathBuf,
2867}
2868
2869impl SessionReportStore {
2870 fn new(repo_root: &Path) -> Self {
2871 Self {
2872 dir: repo_root.join(".heddle/state").join("session-reports"),
2873 }
2874 }
2875
2876 fn session_path(&self, heddle_session_id: &str) -> PathBuf {
2877 self.dir.join(format!("{heddle_session_id}.json"))
2878 }
2879
2880 fn outbox_path(&self) -> PathBuf {
2881 self.dir.join("outbox.jsonl")
2882 }
2883
2884 fn load(&self, heddle_session_id: &str) -> Result<Option<SessionReportEnvelope>> {
2885 let path = self.session_path(heddle_session_id);
2886 if !path.exists() {
2887 return Ok(None);
2888 }
2889 let bytes = fs::read(path)?;
2890 Ok(Some(serde_json::from_slice(&bytes)?))
2891 }
2892
2893 fn save(&self, report: &SessionReportEnvelope) -> Result<()> {
2894 fs::create_dir_all(&self.dir)?;
2895 let path = self.session_path(&report.heddle_session_id);
2896 let bytes = serde_json::to_vec_pretty(report)?;
2897 write_file_atomic(&path, &bytes)?;
2898 Ok(())
2899 }
2900
2901 fn append_outbox(&self, report: &SessionReportEnvelope) -> Result<()> {
2902 fs::create_dir_all(&self.dir)?;
2903 let mut file = OpenOptions::new()
2904 .create(true)
2905 .append(true)
2906 .open(self.outbox_path())?;
2907 serde_json::to_writer(&mut file, report)?;
2908 file.write_all(b"\n")?;
2909 file.flush()?;
2910 Ok(())
2911 }
2912
2913 fn list_pending(&self) -> Result<Vec<String>> {
2914 if !self.dir.exists() {
2915 return Ok(Vec::new());
2916 }
2917 let mut ids = Vec::new();
2918 for entry in fs::read_dir(&self.dir)? {
2919 let entry = entry?;
2920 let path = entry.path();
2921 if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
2922 continue;
2923 }
2924 let bytes = fs::read(&path)?;
2925 let report: SessionReportEnvelope = serde_json::from_slice(&bytes)?;
2926 if report.pending_flush {
2927 ids.push(report.heddle_session_id);
2928 }
2929 }
2930 ids.sort();
2931 Ok(ids)
2932 }
2933}
2934
2935#[derive(Debug, Deserialize)]
2936struct BridgeRequest {
2937 #[serde(default)]
2938 id: Option<String>,
2939 method: String,
2940 #[serde(default)]
2941 params: Value,
2942}
2943
2944#[derive(Debug, Serialize)]
2945struct BridgeResponse {
2946 #[serde(default)]
2947 id: Option<String>,
2948 ok: bool,
2949 #[serde(skip_serializing_if = "Option::is_none")]
2950 result: Option<Value>,
2951 #[serde(skip_serializing_if = "Option::is_none")]
2952 error: Option<BridgeError>,
2953}
2954
2955impl BridgeResponse {
2956 fn ok(id: Option<String>, result: Value) -> Self {
2957 Self {
2958 id,
2959 ok: true,
2960 result: Some(result),
2961 error: None,
2962 }
2963 }
2964
2965 fn error(id: Option<String>, code: impl Into<String>, message: impl Into<String>) -> Self {
2966 Self {
2967 id,
2968 ok: false,
2969 result: None,
2970 error: Some(BridgeError {
2971 code: code.into(),
2972 message: message.into(),
2973 }),
2974 }
2975 }
2976}
2977
2978#[derive(Debug, Serialize)]
2979struct BridgeError {
2980 code: String,
2981 message: String,
2982}
2983
2984#[derive(Debug, Clone, Deserialize, Default)]
2985struct OpenSessionParams {
2986 #[serde(default)]
2987 heddle_session_id: Option<String>,
2988 #[serde(default)]
2989 agent_session_id: Option<String>,
2990 #[serde(default)]
2991 client_instance_id: Option<String>,
2992 #[serde(default)]
2993 thread: Option<String>,
2994 #[serde(default)]
2995 task: Option<String>,
2996 #[serde(default)]
2997 summary: Option<String>,
2998 #[serde(default)]
2999 harness: Option<String>,
3000 #[serde(default)]
3001 provider: Option<String>,
3002 #[serde(default)]
3003 model: Option<String>,
3004 #[serde(default)]
3005 thinking_level: Option<String>,
3006 #[serde(default)]
3007 policy: Option<String>,
3008 #[serde(default)]
3009 transport: Option<HarnessTransport>,
3010 #[serde(default)]
3011 transcript_mode: Option<HarnessTranscriptMode>,
3012 #[serde(default)]
3013 argv: Option<Vec<String>>,
3014 #[serde(default)]
3015 env_hints: BTreeMap<String, String>,
3016 #[serde(default)]
3017 probe_metadata: BTreeMap<String, String>,
3018}
3019
3020#[derive(Debug, Clone, Deserialize, Default)]
3021struct UpdateProgressParams {
3022 heddle_session_id: String,
3023 #[serde(default)]
3024 status: Option<String>,
3025 #[serde(default)]
3026 message: Option<String>,
3027 #[serde(default)]
3028 completed_steps: Option<u32>,
3029 #[serde(default)]
3030 total_steps: Option<u32>,
3031 #[serde(default)]
3032 touched_paths: Vec<String>,
3033 #[serde(default)]
3034 summary: Option<String>,
3035 #[serde(default)]
3036 harness: Option<String>,
3037 #[serde(default)]
3038 provider: Option<String>,
3039 #[serde(default)]
3040 model: Option<String>,
3041 #[serde(default)]
3042 thinking_level: Option<String>,
3043 #[serde(default)]
3044 policy: Option<String>,
3045 #[serde(default)]
3046 argv: Option<Vec<String>>,
3047 #[serde(default)]
3048 env_hints: BTreeMap<String, String>,
3049 #[serde(default)]
3050 probe_metadata: BTreeMap<String, String>,
3051}
3052
3053#[derive(Debug, Clone, Deserialize, Default)]
3054struct RecordUsageParams {
3055 heddle_session_id: String,
3056 #[serde(default)]
3057 input_tokens: Option<u64>,
3058 #[serde(default)]
3059 output_tokens: Option<u64>,
3060 #[serde(default)]
3061 reasoning_tokens: Option<u64>,
3062 #[serde(default)]
3063 cache_creation_tokens: Option<u64>,
3064 #[serde(default)]
3065 cache_read_tokens: Option<u64>,
3066 #[serde(default)]
3067 tool_calls: Option<u32>,
3068 #[serde(default)]
3069 cost_micros_usd: Option<u64>,
3070}
3071
3072#[derive(Debug, Clone, Deserialize, Default)]
3073struct RecordTouchedPathsParams {
3074 heddle_session_id: String,
3075 #[serde(default)]
3076 paths: Vec<String>,
3077}
3078
3079#[derive(Debug, Clone, Deserialize, Default)]
3080struct CloseSessionParams {
3081 heddle_session_id: String,
3082 #[serde(default)]
3083 outcome: Option<String>,
3084 #[serde(default)]
3085 summary: Option<String>,
3086 #[serde(default)]
3087 transcript_refs: Option<Vec<TranscriptAttachmentRef>>,
3088 #[serde(default)]
3089 transport: Option<HarnessTransport>,
3090}
3091
3092#[derive(Debug, Clone, Deserialize, Default)]
3093struct FlushReportsParams {
3094 #[serde(default)]
3095 heddle_session_id: Option<String>,
3096}
3097
3098#[derive(Debug, Serialize)]
3099struct OpenSessionResult {
3100 heddle_session_id: String,
3101 heddle_segment_id: Option<String>,
3102 agent_session_id: Option<String>,
3103 created_session: bool,
3104 harness: Option<String>,
3105 provider: Option<String>,
3106 model: Option<String>,
3107 thinking_level: Option<String>,
3108 report_flush_state: Option<String>,
3109 attach_reason: Option<String>,
3110}
3111
3112#[derive(Debug, Serialize)]
3113struct SessionMutationResult {
3114 heddle_session_id: String,
3115 heddle_segment_id: Option<String>,
3116 report_flush_state: Option<String>,
3117}
3118
3119#[derive(Debug, Serialize)]
3120struct CloseSessionResult {
3121 heddle_session_id: String,
3122 changed_paths: Vec<String>,
3123 diff_summary: SessionDiffSummary,
3124 report_flush_state: Option<String>,
3125}
3126
3127#[derive(Debug, Serialize)]
3128struct FlushReportsResult {
3129 flushed: usize,
3130}
3131
3132#[cfg(test)]
3133mod tests {
3134 #[cfg(unix)]
3135 use std::os::unix::fs::PermissionsExt;
3136
3137 use super::*;
3138
3139 fn init_repo() -> (tempfile::TempDir, Repository) {
3140 let temp = tempfile::TempDir::new().unwrap();
3141 let repo = Repository::init_default(temp.path()).unwrap();
3142 (temp, repo)
3143 }
3144
3145 #[test]
3146 fn harness_config_load_missing_path_defaults_without_warning() {
3147 let temp = tempfile::TempDir::new().unwrap();
3148 let missing = temp.path().join("missing-config.toml");
3149
3150 let (config, warning) = load_harness_user_config(Some(missing));
3151
3152 assert_eq!(config.harness.transport, HarnessTransport::Spool);
3153 assert!(warning.is_none());
3154 }
3155
3156 #[test]
3157 fn harness_config_load_malformed_path_warns_and_defaults() {
3158 let temp = tempfile::TempDir::new().unwrap();
3159 let path = temp.path().join("config.toml");
3160 std::fs::write(&path, "[harness\ntransport = \"direct\"\n").unwrap();
3161
3162 let (config, warning) = load_harness_user_config(Some(path.clone()));
3163
3164 assert_eq!(config.harness.transport, HarnessTransport::Spool);
3165 let warning = warning.expect("malformed config should produce a warning");
3166 assert!(warning.contains("failed to load user config"));
3167 assert!(warning.contains(&path.display().to_string()));
3168 assert!(warning.contains("continuing with defaults"));
3169 }
3170
3171 #[test]
3172 fn harness_config_load_valid_path_loads_without_warning() {
3173 let temp = tempfile::TempDir::new().unwrap();
3174 let path = temp.path().join("config.toml");
3175 std::fs::write(
3176 &path,
3177 "[harness]\ntransport = \"direct\"\ntranscript = \"summary\"\n",
3178 )
3179 .unwrap();
3180
3181 let (config, warning) = load_harness_user_config(Some(path));
3182
3183 assert_eq!(config.harness.transport, HarnessTransport::Direct);
3184 assert_eq!(config.harness.transcript, HarnessTranscriptMode::Summary);
3185 assert!(warning.is_none());
3186 }
3187
3188 #[test]
3189 fn relay_payload_parse_invalid_json_warns_and_uses_null() {
3190 let (value, warning) = parse_relay_payload("{not-json");
3191
3192 assert_eq!(value, Value::Null);
3193 let warning = warning.expect("invalid JSON should produce a warning");
3194 assert!(warning.contains("failed to parse harness relay payload as JSON"));
3195 assert!(warning.contains("continuing with null payload"));
3196 }
3197
3198 #[test]
3199 fn relay_payload_parse_empty_payload_uses_null_without_warning() {
3200 let (value, warning) = parse_relay_payload(" \n");
3201
3202 assert_eq!(value, Value::Null);
3203 assert!(warning.is_none());
3204 }
3205
3206 #[test]
3207 fn relay_payload_parse_valid_json_without_warning() {
3208 let (value, warning) = parse_relay_payload(r#"{"message":"hello"}"#);
3209
3210 assert_eq!(value["message"], "hello");
3211 assert!(warning.is_none());
3212 }
3213
3214 #[test]
3222 fn harness_default_path_matches_canonical_thread_dir() {
3223 let (_temp, repo) = init_repo();
3224 for id in ["foo", "parent/task", "feature/foo", "team@scope"] {
3225 let harness_path = default_private_thread_path(&repo, id);
3226 let canonical = repo.managed_checkout_path(id);
3227 assert_eq!(
3228 harness_path, canonical,
3229 "harness default must match the canonical thread_dir for {id:?}"
3230 );
3231 }
3232 }
3233
3234 #[test]
3235 fn inherited_harness_hints_exclude_ambient_model_identity() {
3236 assert!(!inherited_harness_hint("OPENAI_MODEL"));
3237 assert!(!inherited_harness_hint("ANTHROPIC_MODEL"));
3238 assert!(!inherited_harness_hint("CLAUDE_MODEL"));
3239 assert!(!inherited_harness_hint("MODEL"));
3240 assert!(!inherited_harness_hint("OPENAI_REASONING_EFFORT"));
3241 assert!(inherited_harness_hint("HEDDLE_AGENT_MODEL"));
3242 assert!(inherited_harness_hint("CODEX_SANDBOX"));
3243 assert!(inherited_harness_hint("CLAUDECODE"));
3244 }
3245
3246 #[test]
3247 fn open_session_creates_or_attaches() {
3248 let (_temp, repo) = init_repo();
3249 let user_config = UserConfig::default();
3250 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3251
3252 let created = runtime
3253 .open_session(OpenSessionParams {
3254 harness: Some("codex".to_string()),
3255 provider: Some("openai".to_string()),
3256 model: Some("gpt-5.4".to_string()),
3257 ..OpenSessionParams::default()
3258 })
3259 .unwrap();
3260 assert!(created.created_session);
3261
3262 let attached = runtime
3263 .open_session(OpenSessionParams {
3264 harness: Some("codex".to_string()),
3265 provider: Some("openai".to_string()),
3266 model: Some("gpt-5.4".to_string()),
3267 ..OpenSessionParams::default()
3268 })
3269 .unwrap();
3270 assert!(!attached.created_session);
3271 assert_eq!(created.heddle_session_id, attached.heddle_session_id);
3272 }
3273
3274 #[test]
3275 fn same_client_instance_reattaches_to_its_existing_session() {
3276 let (_temp, repo) = init_repo();
3277 let user_config = UserConfig::default();
3278 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3279
3280 let first = runtime
3281 .open_session(OpenSessionParams {
3282 client_instance_id: Some("client-a".to_string()),
3283 harness: Some("codex".to_string()),
3284 provider: Some("openai".to_string()),
3285 model: Some("gpt-5.4".to_string()),
3286 ..OpenSessionParams::default()
3287 })
3288 .unwrap();
3289 let second = runtime
3290 .open_session(OpenSessionParams {
3291 client_instance_id: Some("client-b".to_string()),
3292 harness: Some("codex".to_string()),
3293 provider: Some("openai".to_string()),
3294 model: Some("gpt-5.4".to_string()),
3295 ..OpenSessionParams::default()
3296 })
3297 .unwrap();
3298 let reopened = runtime
3299 .open_session(OpenSessionParams {
3300 client_instance_id: Some("client-a".to_string()),
3301 harness: Some("codex".to_string()),
3302 provider: Some("openai".to_string()),
3303 model: Some("gpt-5.4".to_string()),
3304 ..OpenSessionParams::default()
3305 })
3306 .unwrap();
3307
3308 assert_ne!(first.heddle_session_id, second.heddle_session_id);
3309 assert_eq!(first.heddle_session_id, reopened.heddle_session_id);
3310 assert_eq!(first.agent_session_id, reopened.agent_session_id);
3311 }
3312
3313 #[test]
3314 fn different_client_instances_do_not_share_the_current_session() {
3315 let (_temp, repo) = init_repo();
3316 let user_config = UserConfig::default();
3317 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3318
3319 let first = runtime
3320 .open_session(OpenSessionParams {
3321 client_instance_id: Some("client-a".to_string()),
3322 harness: Some("codex".to_string()),
3323 provider: Some("openai".to_string()),
3324 model: Some("gpt-5.4".to_string()),
3325 ..OpenSessionParams::default()
3326 })
3327 .unwrap();
3328 let second = runtime
3329 .open_session(OpenSessionParams {
3330 client_instance_id: Some("client-b".to_string()),
3331 harness: Some("codex".to_string()),
3332 provider: Some("openai".to_string()),
3333 model: Some("gpt-5.4".to_string()),
3334 ..OpenSessionParams::default()
3335 })
3336 .unwrap();
3337
3338 assert_ne!(first.heddle_session_id, second.heddle_session_id);
3339 assert_ne!(first.agent_session_id, second.agent_session_id);
3340 }
3341
3342 #[test]
3343 fn provider_model_change_creates_segment() {
3344 let (_temp, repo) = init_repo();
3345 let user_config = UserConfig::default();
3346 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3347
3348 let opened = runtime
3349 .open_session(OpenSessionParams {
3350 harness: Some("claude-code".to_string()),
3351 provider: Some("anthropic".to_string()),
3352 model: Some("claude-sonnet".to_string()),
3353 ..OpenSessionParams::default()
3354 })
3355 .unwrap();
3356 runtime
3357 .update_progress(UpdateProgressParams {
3358 heddle_session_id: opened.heddle_session_id.clone(),
3359 provider: Some("openai".to_string()),
3360 model: Some("gpt-5.4".to_string()),
3361 ..UpdateProgressParams::default()
3362 })
3363 .unwrap();
3364
3365 let report = runtime
3366 .reports
3367 .load(&opened.heddle_session_id)
3368 .unwrap()
3369 .unwrap();
3370 let expected_segment = format!("{}-seg-2", opened.heddle_session_id);
3371 assert_eq!(
3372 report.heddle_segment_id.as_deref(),
3373 Some(expected_segment.as_str())
3374 );
3375 }
3376
3377 #[test]
3378 fn blank_agent_model_hint_falls_through_to_detected_model_without_segment_rotation() {
3379 let (_temp, repo) = init_repo();
3380 let user_config = UserConfig::default();
3381 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3382 let blank_model_env = BTreeMap::from([
3383 ("HEDDLE_AGENT_PROVIDER".to_string(), "anthropic".to_string()),
3384 ("HEDDLE_AGENT_MODEL".to_string(), String::new()),
3385 ]);
3386
3387 let opened = runtime
3388 .open_session(OpenSessionParams {
3389 harness: Some("claude-code".to_string()),
3390 env_hints: blank_model_env.clone(),
3391 probe_metadata: BTreeMap::from([
3392 ("session_id".to_string(), "claude-sess-blank".to_string()),
3393 ("model".to_string(), "claude-opus-4-8[1m]".to_string()),
3394 ]),
3395 ..OpenSessionParams::default()
3396 })
3397 .unwrap();
3398 assert_eq!(opened.model.as_deref(), Some("claude-opus-4-8[1m]"));
3399
3400 let original_segment = opened.heddle_segment_id.clone();
3401 runtime
3402 .update_progress(UpdateProgressParams {
3403 heddle_session_id: opened.heddle_session_id.clone(),
3404 env_hints: blank_model_env,
3405 probe_metadata: BTreeMap::from([
3406 ("session_id".to_string(), "claude-sess-blank".to_string()),
3407 ("model".to_string(), "claude-opus-4-8[1m]".to_string()),
3408 ]),
3409 ..UpdateProgressParams::default()
3410 })
3411 .unwrap();
3412
3413 let report = runtime
3414 .reports
3415 .load(&opened.heddle_session_id)
3416 .unwrap()
3417 .unwrap();
3418 assert_eq!(report.harness.model.as_deref(), Some("claude-opus-4-8[1m]"));
3419 assert_eq!(report.heddle_segment_id, original_segment);
3420 }
3421
3422 #[test]
3423 fn close_session_captures_changed_paths_from_status_and_hints() {
3424 let (temp, repo) = init_repo();
3425 let config = UserConfig::default();
3426 let mut runtime = HarnessBridgeRuntime::new(repo, config);
3427
3428 let opened = runtime
3429 .open_session(OpenSessionParams {
3430 harness: Some("codex".to_string()),
3431 provider: Some("openai".to_string()),
3432 model: Some("gpt-5.4".to_string()),
3433 ..OpenSessionParams::default()
3434 })
3435 .unwrap();
3436 std::fs::write(temp.path().join("src.txt"), "hello\n").unwrap();
3437 runtime
3438 .record_touched_paths(RecordTouchedPathsParams {
3439 heddle_session_id: opened.heddle_session_id.clone(),
3440 paths: vec!["src.txt".to_string(), "notes.md".to_string()],
3441 })
3442 .unwrap();
3443 let closed = runtime
3444 .close_session(CloseSessionParams {
3445 heddle_session_id: opened.heddle_session_id.clone(),
3446 outcome: Some("completed".to_string()),
3447 ..CloseSessionParams::default()
3448 })
3449 .unwrap();
3450 let report = runtime
3451 .reports
3452 .load(&opened.heddle_session_id)
3453 .unwrap()
3454 .unwrap();
3455 assert!(closed.changed_paths.iter().any(|path| path == "src.txt"));
3456 assert!(!closed.changed_paths.iter().any(|path| path == "notes.md"));
3457 assert!(report.touched_paths.iter().any(|path| path == "src.txt"));
3458 assert!(report.touched_paths.iter().any(|path| path == "notes.md"));
3459 assert_eq!(
3460 closed.diff_summary.changed_file_count,
3461 closed.changed_paths.len() as u32
3462 );
3463 }
3464
3465 #[test]
3466 fn flush_reports_moves_pending_report_to_outbox() {
3467 let (_temp, repo) = init_repo();
3468 let user_config = UserConfig::default();
3469 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3470
3471 let opened = runtime
3472 .open_session(OpenSessionParams {
3473 harness: Some("codex".to_string()),
3474 provider: Some("openai".to_string()),
3475 model: Some("gpt-5.4".to_string()),
3476 ..OpenSessionParams::default()
3477 })
3478 .unwrap();
3479 let flushed = runtime
3480 .flush_reports(FlushReportsParams {
3481 heddle_session_id: Some(opened.heddle_session_id.clone()),
3482 })
3483 .unwrap();
3484 assert_eq!(flushed.flushed, 1);
3485 let report = runtime
3486 .reports
3487 .load(&opened.heddle_session_id)
3488 .unwrap()
3489 .unwrap();
3490 assert!(!report.pending_flush);
3491 assert_eq!(report.report_flush_state.as_deref(), Some("queued-local"));
3492 assert!(runtime.reports.outbox_path().exists());
3493 }
3494
3495 #[test]
3496 fn explicit_overrides_beat_fingerprint_and_user_defaults() {
3497 let (_temp, repo) = init_repo();
3498 let mut user_config = UserConfig::default();
3499 user_config.harness.harnesses.insert(
3500 "codex".to_string(),
3501 UserHarnessOverride {
3502 provider: Some("openai".to_string()),
3503 model: Some("gpt-default".to_string()),
3504 thinking_level: Some("medium".to_string()),
3505 policy: Some("default".to_string()),
3506 },
3507 );
3508 let identity = resolve_identity(
3509 &repo,
3510 &user_config,
3511 IdentityHints {
3512 harness: Some("codex".to_string()),
3513 provider: Some("openai".to_string()),
3514 model: Some("gpt-5.4".to_string()),
3515 thinking_level: Some("high".to_string()),
3516 policy: Some("custom".to_string()),
3517 probe: HarnessProbeResult::default(),
3518 },
3519 )
3520 .unwrap();
3521 assert_eq!(identity.model.as_deref(), Some("gpt-5.4"));
3522 assert_eq!(identity.thinking_level.as_deref(), Some("high"));
3523 assert_eq!(identity.policy.as_deref(), Some("custom"));
3524 }
3525
3526 #[test]
3527 fn transcript_mode_defaults_to_off_and_keeps_refs_empty() {
3528 let (_temp, repo) = init_repo();
3529 let user_config = UserConfig::default();
3530 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3531
3532 let opened = runtime
3533 .open_session(OpenSessionParams {
3534 harness: Some("codex".to_string()),
3535 provider: Some("openai".to_string()),
3536 model: Some("gpt-5.4".to_string()),
3537 ..OpenSessionParams::default()
3538 })
3539 .unwrap();
3540 let report = runtime
3541 .reports
3542 .load(&opened.heddle_session_id)
3543 .unwrap()
3544 .unwrap();
3545 assert_eq!(report.transcript_mode, "off");
3546 assert!(report.transcript_refs.is_empty());
3547 }
3548
3549 #[test]
3550 fn codex_thread_probe_reattaches_same_actor() {
3551 let (_temp, repo) = init_repo();
3552 let user_config = UserConfig::default();
3553 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3554
3555 let first = runtime
3556 .open_session(OpenSessionParams {
3557 harness: Some("codex".to_string()),
3558 probe_metadata: BTreeMap::from([
3559 ("thread_id".to_string(), "thr_123".to_string()),
3560 ("client_name".to_string(), "codex-tui".to_string()),
3561 ]),
3562 ..OpenSessionParams::default()
3563 })
3564 .unwrap();
3565 let second = runtime
3566 .open_session(OpenSessionParams {
3567 harness: Some("codex".to_string()),
3568 probe_metadata: BTreeMap::from([
3569 ("thread_id".to_string(), "thr_123".to_string()),
3570 ("client_name".to_string(), "codex-tui".to_string()),
3571 ]),
3572 ..OpenSessionParams::default()
3573 })
3574 .unwrap();
3575
3576 assert_eq!(first.agent_session_id, second.agent_session_id);
3577 assert_eq!(first.heddle_session_id, second.heddle_session_id);
3578 }
3579
3580 #[test]
3581 fn opencode_child_session_creates_distinct_actor_with_parent_key() {
3582 let (_temp, repo) = init_repo();
3583 let user_config = UserConfig::default();
3584 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3585
3586 let root = runtime
3587 .open_session(OpenSessionParams {
3588 harness: Some("opencode".to_string()),
3589 probe_metadata: BTreeMap::from([("session_id".to_string(), "root-1".to_string())]),
3590 ..OpenSessionParams::default()
3591 })
3592 .unwrap();
3593 let child = runtime
3594 .open_session(OpenSessionParams {
3595 harness: Some("opencode".to_string()),
3596 probe_metadata: BTreeMap::from([
3597 ("session_id".to_string(), "child-1".to_string()),
3598 ("parent_id".to_string(), "root-1".to_string()),
3599 ]),
3600 ..OpenSessionParams::default()
3601 })
3602 .unwrap();
3603
3604 assert_ne!(root.agent_session_id, child.agent_session_id);
3605 let report = runtime
3606 .reports
3607 .load(&child.heddle_session_id)
3608 .unwrap()
3609 .unwrap();
3610 assert_eq!(
3611 report.native_parent_actor_key.as_deref(),
3612 Some("opencode:session:root-1")
3613 );
3614 }
3615
3616 #[test]
3617 fn claude_resume_with_new_session_id_does_not_steal_existing_actor() {
3618 let (_temp, repo) = init_repo();
3619 let user_config = UserConfig::default();
3620 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3621
3622 let first = runtime
3623 .open_session(OpenSessionParams {
3624 harness: Some("claude-code".to_string()),
3625 probe_metadata: BTreeMap::from([
3626 ("session_id".to_string(), "sess-old".to_string()),
3627 (
3628 "transcript_path".to_string(),
3629 "/tmp/claude/session-a.jsonl".to_string(),
3630 ),
3631 ]),
3632 ..OpenSessionParams::default()
3633 })
3634 .unwrap();
3635 let resumed = runtime
3636 .open_session(OpenSessionParams {
3637 harness: Some("claude-code".to_string()),
3638 probe_metadata: BTreeMap::from([
3639 ("session_id".to_string(), "sess-new".to_string()),
3640 (
3641 "transcript_path".to_string(),
3642 "/tmp/claude/session-a.jsonl".to_string(),
3643 ),
3644 ]),
3645 ..OpenSessionParams::default()
3646 })
3647 .unwrap();
3648
3649 assert_ne!(first.agent_session_id, resumed.agent_session_id);
3650 assert_ne!(first.heddle_session_id, resumed.heddle_session_id);
3651 }
3652
3653 #[test]
3654 fn explicit_claude_harness_beats_generic_session_id_probe_match() {
3655 let (_temp, repo) = init_repo();
3656 let user_config = UserConfig::default();
3657 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3658
3659 let opened = runtime
3660 .open_session(OpenSessionParams {
3661 harness: Some("claude-code".to_string()),
3662 probe_metadata: BTreeMap::from([
3663 ("session_id".to_string(), "claude-sess-1".to_string()),
3664 ("hook_event".to_string(), "SubagentStop".to_string()),
3665 ]),
3666 ..OpenSessionParams::default()
3667 })
3668 .unwrap();
3669 let report = runtime
3670 .reports
3671 .load(&opened.heddle_session_id)
3672 .unwrap()
3673 .unwrap();
3674 assert_eq!(
3675 report.native_actor_key.as_deref(),
3676 Some("claude-code:session:claude-sess-1")
3677 );
3678 assert_eq!(report.harness.harness.as_deref(), Some("claude-code"));
3679 }
3680
3681 #[test]
3682 fn same_native_actor_key_reuses_existing_actor_after_tentative_session_creation() {
3683 let (_temp, repo) = init_repo();
3684 let user_config = UserConfig::default();
3685 let runtime = HarnessBridgeRuntime::new(repo, user_config);
3686 let principal = runtime.repo.get_principal().unwrap();
3687 let mut sessions = SessionManager::new(runtime.repo.root());
3688 let existing_session = sessions
3689 .start_session(
3690 principal.clone(),
3691 "anthropic".to_string(),
3692 "claude-opus-4-7[1m]".to_string(),
3693 None,
3694 )
3695 .unwrap();
3696 let tentative_session = sessions
3697 .start_session(
3698 principal,
3699 "anthropic".to_string(),
3700 "claude-opus-4-7[1m]".to_string(),
3701 None,
3702 )
3703 .unwrap();
3704
3705 let registry = ActorPresenceStore::new(runtime.repo.heddle_dir());
3706 let existing_entry = registry
3707 .create_generated_entry(|session_id| {
3708 Ok(ActorPresence {
3709 session_id: session_id.to_string(),
3710 client_instance_id: None,
3711 native_actor_key: Some(
3712 "claude-code:session:282396d3-554a-48aa-a9a8-8d1f0bd15fa5".to_string(),
3713 ),
3714 native_parent_actor_key: None,
3715 native_instance_key: Some(
3716 "claude-code:transcript:/tmp/claude/282396d3.jsonl".to_string(),
3717 ),
3718 heddle_session_id: Some(existing_session.id.clone()),
3719 thread_id: None,
3720 thread: "detached".to_string(),
3721 anchor_state: None,
3722 anchor_root: None,
3723 path: Some(runtime.repo.root().to_path_buf()),
3724 base_state: String::new(),
3725 started_at: Utc::now(),
3726 provider: Some("anthropic".to_string()),
3727 model: Some("claude-opus-4-7[1m]".to_string()),
3728 harness: Some("claude-code".to_string()),
3729 thinking_level: None,
3730 usage_summary: AgentUsageSummary::default(),
3731 last_progress_at: None,
3732 report_flush_state: Some("pending-local".to_string()),
3733 attach_reason: None,
3734 task_assignment_id: None,
3735 attach_precedence: vec![],
3736 winning_attach_rule: None,
3737 probe_source: Some("hook_payload".to_string()),
3738 probe_confidence: Some(1.0),
3739 status: ActorPresenceStatus::Active,
3740 completed_at: None,
3741 context_queries: vec![],
3742 })
3743 })
3744 .unwrap();
3745
3746 let probe = HarnessProbeResult {
3747 harness: Some("claude-code".to_string()),
3748 provider: Some("anthropic".to_string()),
3749 model: Some("claude-opus-4-7[1m]".to_string()),
3750 native_actor_key: Some(
3751 "claude-code:session:282396d3-554a-48aa-a9a8-8d1f0bd15fa5".to_string(),
3752 ),
3753 native_instance_key: Some(
3754 "claude-code:transcript:/tmp/claude/282396d3.jsonl".to_string(),
3755 ),
3756 probe_source: Some("hook_payload".to_string()),
3757 confidence: Some(1.0),
3758 ..HarnessProbeResult::default()
3759 };
3760 let identity = ResolvedIdentity {
3761 harness: Some("claude-code".to_string()),
3762 provider: Some("anthropic".to_string()),
3763 model: Some("claude-opus-4-7[1m]".to_string()),
3764 thinking_level: None,
3765 policy: None,
3766 };
3767 let mut attach = ResolvedAttachment {
3768 target: AttachTarget::CreateNew {
3769 _because_claimed: false,
3770 },
3771 matched_entry: None,
3772 attach_reason:
3773 "started new Heddle session because no compatible native actor match was found"
3774 .to_string(),
3775 precedence: vec!["native-actor-key:miss".to_string()],
3776 winning_rule: "create-new-session".to_string(),
3777 };
3778
3779 let resolved_entry = runtime
3780 .ensure_registry_entry(RegistryEntryRequest {
3781 heddle_session_id: &tentative_session.id,
3782 thread_name: None,
3783 thread_id: None,
3784 identity: &identity,
3785 probe: &probe,
3786 attach: &attach,
3787 client_instance_id: None,
3788 requested_entry: None,
3789 })
3790 .unwrap();
3791 assert_eq!(resolved_entry.session_id, existing_entry.session_id);
3792 assert_eq!(
3793 resolved_entry.heddle_session_id.as_deref(),
3794 Some(existing_session.id.as_str())
3795 );
3796
3797 let (canonical_session, owns_session) = runtime
3798 .reuse_canonical_actor_session(
3799 &mut sessions,
3800 CanonicalActorSessionRequest {
3801 tentative_session: tentative_session.clone(),
3802 tentative_owns_session: true,
3803 entry: &resolved_entry,
3804 probe: &probe,
3805 attach: &mut attach,
3806 },
3807 )
3808 .unwrap();
3809 assert_eq!(canonical_session.id, existing_session.id);
3810 assert!(!owns_session);
3811 assert!(
3812 attach
3813 .precedence
3814 .iter()
3815 .any(|step| step.starts_with("post-create-native-actor-key:"))
3816 );
3817 assert_eq!(attach.winning_rule, "native-actor-key-post-create");
3818 assert!(
3819 !sessions
3820 .get_session(&tentative_session.id)
3821 .unwrap()
3822 .unwrap()
3823 .is_active()
3824 );
3825 }
3826
3827 #[test]
3828 fn close_session_does_not_blame_preexisting_dirty_worktree() {
3829 let (temp, repo) = init_repo();
3830 std::fs::write(temp.path().join("preexisting.txt"), "already dirty\n").unwrap();
3831 let user_config = UserConfig::default();
3832 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
3833
3834 let opened = runtime
3835 .open_session(OpenSessionParams {
3836 harness: Some("claude-code".to_string()),
3837 provider: Some("anthropic".to_string()),
3838 model: Some("claude-opus-4-7[1m]".to_string()),
3839 ..OpenSessionParams::default()
3840 })
3841 .unwrap();
3842 let closed = runtime
3843 .close_session(CloseSessionParams {
3844 heddle_session_id: opened.heddle_session_id.clone(),
3845 outcome: Some("completed".to_string()),
3846 ..CloseSessionParams::default()
3847 })
3848 .unwrap();
3849 let report = runtime
3850 .reports
3851 .load(&opened.heddle_session_id)
3852 .unwrap()
3853 .unwrap();
3854
3855 assert!(
3856 report
3857 .worktree_changes_at_open
3858 .iter()
3859 .any(|change| change.path == "preexisting.txt")
3860 );
3861 assert!(
3862 !closed
3863 .changed_paths
3864 .iter()
3865 .any(|path| path == "preexisting.txt")
3866 );
3867 assert_eq!(closed.diff_summary.changed_file_count, 0);
3868 }
3869
3870 #[test]
3871 fn timeline_state_delta_paths_ignore_uncaptured_worktree_changes() {
3872 let (temp, repo) = init_repo();
3873 let repo_root = repo.root().to_path_buf();
3874 std::fs::write(repo_root.join("tracked.txt"), b"one\n").unwrap();
3875 let before = repo.snapshot(Some("seed".into()), None).unwrap();
3876 std::fs::write(repo_root.join("tracked.txt"), b"two\n").unwrap();
3877 let after = repo.snapshot(Some("advance".into()), None).unwrap();
3878 std::fs::write(temp.path().join("ambient.txt"), b"not in the state delta\n").unwrap();
3879
3880 assert_eq!(
3881 changed_paths_between_states(&repo, before.state_id, after.state_id).unwrap(),
3882 vec!["tracked.txt"]
3883 );
3884 }
3885
3886 #[test]
3887 fn relay_claude_stop_captures_state_with_agent_attribution() {
3888 let (temp, repo) = init_repo();
3889 let repo_root = repo.root().to_path_buf();
3890
3891 std::fs::write(repo_root.join("seed.txt"), b"hello").unwrap();
3893 let _ = repo.snapshot(Some("seed".into()), None).unwrap();
3894
3895 std::fs::write(repo_root.join("seed.txt"), b"hello, heddle").unwrap();
3897
3898 drop(repo);
3899
3900 let fresh_repo = Repository::open(temp.path()).unwrap();
3901 let user_config = UserConfig {
3902 principal: Some(crate::config::UserPrincipalConfig {
3903 name: "Ada Lovelace".to_string(),
3904 email: "ada@example.com".to_string(),
3905 }),
3906 ..UserConfig::default()
3907 };
3908 let mut runtime = HarnessBridgeRuntime::new(fresh_repo, user_config);
3909 let payload = serde_json::json!({
3910 "session_id": "claude-sess-123",
3911 "transcript_path": "/tmp/claude/x.jsonl",
3912 "model": {
3913 "id": "claude-opus-4-7",
3914 "display_name": "Claude Opus 4.7",
3915 },
3916 "message": "hook-driven capture test",
3917 "hook_event_name": "Stop",
3918 });
3919 relay_claude(&mut runtime, "Stop", &payload).unwrap();
3920 drop(runtime);
3921
3922 let verify = Repository::open(temp.path()).unwrap();
3923 let head_id = verify.head().unwrap().expect("HEAD after Stop capture");
3924 let state = verify
3925 .store()
3926 .get_state(&head_id)
3927 .unwrap()
3928 .expect("state for HEAD");
3929 let agent = state.attribution.agent.expect("agent attribution on state");
3930 assert_eq!(agent.provider, "anthropic");
3931 assert_eq!(agent.model, "Claude Opus 4.7");
3932 assert_eq!(
3933 state.intent.as_deref(),
3934 Some("hook-driven capture test"),
3935 "intent should be pulled from payload message",
3936 );
3937 }
3938
3939 #[test]
3940 fn relay_claude_stop_is_idempotent_when_clean() {
3941 let (temp, repo) = init_repo();
3942 let repo_root = repo.root().to_path_buf();
3943 std::fs::write(repo_root.join("seed.txt"), b"hello").unwrap();
3944 let seed = repo.snapshot(Some("seed".into()), None).unwrap();
3945 drop(repo);
3946
3947 let fresh_repo = Repository::open(temp.path()).unwrap();
3948 let mut runtime = HarnessBridgeRuntime::new(fresh_repo, UserConfig::default());
3949 let payload = serde_json::json!({
3950 "session_id": "claude-sess-clean",
3951 "model": {"id": "claude-sonnet-4-6"},
3952 });
3953 relay_claude(&mut runtime, "Stop", &payload).unwrap();
3954 drop(runtime);
3955
3956 let verify = Repository::open(temp.path()).unwrap();
3957 let head_id = verify.head().unwrap().expect("HEAD preserved");
3958 assert_eq!(
3959 head_id, seed.state_id,
3960 "no change expected when worktree is clean",
3961 );
3962 }
3963
3964 #[test]
3965 fn relay_claude_pre_tool_use_ignores_non_file_tool() {
3966 let (temp, repo) = init_repo();
3967 drop(repo);
3968 let fresh_repo = Repository::open(temp.path()).unwrap();
3969 let mut runtime = HarnessBridgeRuntime::new(fresh_repo, UserConfig::default());
3970 let payload = serde_json::json!({
3971 "session_id": "claude-sess-bash",
3972 "tool_name": "Bash",
3973 "tool_input": {"command": "ls"},
3974 });
3975 relay_claude(&mut runtime, "PreToolUse", &payload).unwrap();
3977 }
3978
3979 #[test]
3980 fn relay_opencode_tool_execute_before_records_timeline_step() {
3981 let (_temp, repo) = init_repo();
3982 let root = repo.root().to_path_buf();
3983 std::fs::write(root.join("seed.txt"), b"hello").unwrap();
3984 let seed = repo.snapshot(Some("seed".into()), None).unwrap();
3985 let mut runtime = HarnessBridgeRuntime::new(repo, UserConfig::default());
3986 let payload = opencode_tool_payload("call-1");
3987
3988 relay_opencode(&mut runtime, "tool.execute.before", &payload).unwrap();
3989
3990 let store = TimelineStore::open(runtime.repo.heddle_dir()).unwrap();
3991 let view = TimelineView::rebuild(&store).unwrap();
3992 let steps = view.steps_for_thread("main");
3993 assert_eq!(steps.len(), 1);
3994 let step = steps[0];
3995 assert_eq!(step.native.as_ref().unwrap().harness, "opencode");
3996 assert_eq!(step.native.as_ref().unwrap().tool_call_id, "call-1");
3997 assert_eq!(step.tool_name.as_deref(), Some("bash"));
3998 assert_eq!(step.before_state, Some(seed.state_id));
3999 assert!(step.status.is_none());
4000 assert!(step.payload_summary.as_deref().unwrap().contains("call-1"));
4001 assert!(step.payload_hash.is_some());
4002 assert!(
4003 step.labels
4004 .contains(&TimelineLabel::ExternalSideEffectsUnknown)
4005 );
4006 }
4007
4008 #[test]
4009 fn relay_opencode_tool_execute_after_captures_dirty_worktree() {
4010 let (_temp, repo) = init_repo();
4011 let root = repo.root().to_path_buf();
4012 std::fs::write(root.join("tracked.txt"), b"one\n").unwrap();
4013 let seed = repo.snapshot(Some("seed".into()), None).unwrap();
4014 let user_config = UserConfig {
4015 principal: Some(crate::config::UserPrincipalConfig {
4016 name: "Ada Lovelace".to_string(),
4017 email: "ada@example.com".to_string(),
4018 }),
4019 ..UserConfig::default()
4020 };
4021 let mut runtime = HarnessBridgeRuntime::new(repo, user_config);
4022 let payload = opencode_tool_payload("call-2");
4023
4024 relay_opencode(&mut runtime, "tool.execute.before", &payload).unwrap();
4025 std::fs::write(root.join("tracked.txt"), b"two\n").unwrap();
4026 relay_opencode(&mut runtime, "tool.execute.after", &payload).unwrap();
4027
4028 let head = runtime.repo.head().unwrap().expect("capture advanced HEAD");
4029 assert_ne!(head, seed.state_id);
4030 let store = TimelineStore::open(runtime.repo.heddle_dir()).unwrap();
4031 let view = TimelineView::rebuild(&store).unwrap();
4032 let steps = view.steps_for_thread("main");
4033 assert_eq!(steps.len(), 1, "before/after should merge by native id");
4034 let step = steps[0];
4035 assert_eq!(step.operation_ids.len(), 2);
4036 assert_eq!(step.status, Some(TimelineToolCallStatus::Succeeded));
4037 assert_eq!(step.before_state, Some(seed.state_id));
4038 assert_eq!(step.after_state, Some(head));
4039 assert_eq!(step.capture_state, Some(head));
4040 assert_eq!(step.changed, Some(true));
4041 assert!(step.touched_paths.contains(&"tracked.txt".to_string()));
4042 assert!(step.labels.contains(&TimelineLabel::RepoReversible));
4043 assert!(
4044 step.labels
4045 .contains(&TimelineLabel::ExternalSideEffectsUnknown)
4046 );
4047 assert!(!step.payload_summary.as_deref().unwrap().contains("SECRET"));
4048 assert!(step.payload_hash.is_some());
4049 }
4050
4051 #[cfg(unix)]
4052 #[test]
4053 fn relay_opencode_tool_execute_after_records_capture_failed_without_ambient_paths() {
4054 let (_temp, repo) = init_repo();
4055 let root = repo.root().to_path_buf();
4056 std::fs::write(root.join("seed.txt"), b"seed\n").unwrap();
4057 let seed = repo.snapshot(Some("seed".into()), None).unwrap();
4058 let mut runtime = HarnessBridgeRuntime::new(repo, UserConfig::default());
4059 let mut payload = opencode_tool_payload("call-capture-failed");
4060 payload["tool"]["input"]["file_path"] = serde_json::json!("hinted.txt");
4061 let hooks_dir = root.join(".heddle/hooks");
4062 std::fs::create_dir_all(&hooks_dir).unwrap();
4063 let hook_path = hooks_dir.join("pre-snapshot");
4064 std::fs::write(&hook_path, "#!/bin/sh\nexit 1\n").unwrap();
4065 let mut perms = std::fs::metadata(&hook_path).unwrap().permissions();
4066 perms.set_mode(0o755);
4067 std::fs::set_permissions(&hook_path, perms).unwrap();
4068
4069 relay_opencode(&mut runtime, "tool.execute.before", &payload).unwrap();
4070 std::fs::write(root.join("ambient.txt"), b"dirty but uncaptured\n").unwrap();
4071 relay_opencode(&mut runtime, "tool.execute.after", &payload).unwrap();
4072
4073 assert_eq!(
4074 runtime.repo.head().unwrap(),
4075 Some(seed.state_id),
4076 "capture failure must not advance HEAD"
4077 );
4078 let store = TimelineStore::open(runtime.repo.heddle_dir()).unwrap();
4079 let view = TimelineView::rebuild(&store).unwrap();
4080 let steps = view.steps_for_thread("main");
4081 assert_eq!(steps.len(), 1, "before/after should merge by native id");
4082 let step = steps[0];
4083 assert_eq!(step.operation_ids.len(), 2);
4084 assert_eq!(step.before_state, Some(seed.state_id));
4085 assert_eq!(step.after_state, Some(seed.state_id));
4086 assert_eq!(step.capture_state, None);
4087 assert_eq!(step.changed, Some(false));
4088 assert!(step.labels.contains(&TimelineLabel::CaptureFailed));
4089 assert!(
4090 !step.labels.contains(&TimelineLabel::RepoReversible),
4091 "failed captures are not repo-reversible"
4092 );
4093 assert_eq!(step.touched_paths, vec!["hinted.txt"]);
4094 }
4095
4096 #[test]
4097 fn relay_opencode_tool_execute_missing_tool_id_does_not_fail_or_record_timeline() {
4098 let (_temp, repo) = init_repo();
4099 let root = repo.root().to_path_buf();
4100 std::fs::write(root.join("seed.txt"), b"hello").unwrap();
4101 let _ = repo.snapshot(Some("seed".into()), None).unwrap();
4102 let mut runtime = HarnessBridgeRuntime::new(repo, UserConfig::default());
4103 let payload = serde_json::json!({
4104 "sessionID": "opencode-session",
4105 "model": "gpt-5.4",
4106 "provider": "openai",
4107 "tool": {"name": "bash"},
4108 });
4109
4110 relay_opencode(&mut runtime, "tool.execute.before", &payload).unwrap();
4111
4112 let store = TimelineStore::open(runtime.repo.heddle_dir()).unwrap();
4113 let view = TimelineView::rebuild(&store).unwrap();
4114 assert!(view.steps_for_thread("main").is_empty());
4115 let report_count = std::fs::read_dir(root.join(".heddle/state/session-reports"))
4116 .unwrap()
4117 .count();
4118 assert!(
4119 report_count > 0,
4120 "session progress should still be recorded"
4121 );
4122 }
4123
4124 fn opencode_tool_payload(call_id: &str) -> Value {
4125 serde_json::json!({
4126 "sessionID": "opencode-session",
4127 "messageID": "message-1",
4128 "toolCallID": call_id,
4129 "model": "gpt-5.4",
4130 "provider": "openai",
4131 "tool": {
4132 "name": "bash",
4133 "input": {
4134 "command": "echo SECRET",
4135 "file_path": "tracked.txt"
4136 }
4137 },
4138 "status": "success"
4139 })
4140 }
4141
4142 #[test]
4143 fn relay_claude_subagent_start_creates_child_entry_with_parent_key() {
4144 let (temp, repo) = init_repo();
4145 drop(repo);
4146 let fresh_repo = Repository::open(temp.path()).unwrap();
4147 let mut runtime = HarnessBridgeRuntime::new(fresh_repo, UserConfig::default());
4148 let payload = serde_json::json!({
4149 "session_id": "parent-claude-sess",
4150 "agent_id": "child-subagent-xyz",
4151 "model": {"id": "claude-sonnet-4-6"},
4152 });
4153 relay_claude(&mut runtime, "SubagentStart", &payload).unwrap();
4154 drop(runtime);
4155
4156 let verify = Repository::open(temp.path()).unwrap();
4157 let registry = ActorPresenceStore::new(verify.heddle_dir());
4158 let child = registry
4159 .find_active_by_native_actor_key("claude-code:agent:child-subagent-xyz")
4160 .unwrap()
4161 .expect("subagent ActorPresence should exist after SubagentStart");
4162 assert_eq!(
4163 child.native_parent_actor_key.as_deref(),
4164 Some("claude-code:session:parent-claude-sess"),
4165 "subagent must carry parent session linkage",
4166 );
4167 assert_eq!(child.status, ActorPresenceStatus::Active);
4168 }
4169
4170 #[test]
4171 fn relay_claude_subagent_stop_marks_child_entry_complete() {
4172 let (temp, repo) = init_repo();
4173 let repo_root = repo.root().to_path_buf();
4174 drop(repo);
4175
4176 let fresh = Repository::open(temp.path()).unwrap();
4178 let mut runtime = HarnessBridgeRuntime::new(fresh, UserConfig::default());
4179 let start_payload = serde_json::json!({
4180 "session_id": "parent-sess",
4181 "agent_id": "worker-1",
4182 "model": {"id": "claude-sonnet-4-6"},
4183 });
4184 relay_claude(&mut runtime, "SubagentStart", &start_payload).unwrap();
4185 drop(runtime);
4186
4187 std::fs::write(
4189 repo_root.join("child-output.txt"),
4190 b"subagent produced this",
4191 )
4192 .unwrap();
4193
4194 let fresh = Repository::open(temp.path()).unwrap();
4195 let mut runtime = HarnessBridgeRuntime::new(fresh, UserConfig::default());
4196 let stop_payload = serde_json::json!({
4197 "session_id": "parent-sess",
4198 "agent_id": "worker-1",
4199 "model": {
4200 "id": "claude-sonnet-4-6",
4201 "display_name": "Claude Sonnet 4.6",
4202 },
4203 });
4204 relay_claude(&mut runtime, "SubagentStop", &stop_payload).unwrap();
4205 drop(runtime);
4206
4207 let verify = Repository::open(temp.path()).unwrap();
4208 let registry = ActorPresenceStore::new(verify.heddle_dir());
4209 let child = registry
4210 .list()
4211 .unwrap()
4212 .into_iter()
4213 .find(|e| e.native_actor_key.as_deref() == Some("claude-code:agent:worker-1"))
4214 .expect("child entry should still exist");
4215 assert_eq!(
4216 child.status,
4217 ActorPresenceStatus::Complete,
4218 "SubagentStop should mark the child entry Complete",
4219 );
4220 }
4221
4222 #[test]
4223 fn relay_claude_user_prompt_submit_rotates_segment() {
4224 let (temp, repo) = init_repo();
4225 drop(repo);
4226
4227 let fresh = Repository::open(temp.path()).unwrap();
4228 let mut runtime = HarnessBridgeRuntime::new(fresh, UserConfig::default());
4229 let session_payload = serde_json::json!({
4231 "session_id": "claude-prompt-sess",
4232 "model": {"id": "claude-opus-4-7", "display_name": "Claude Opus 4.7"},
4233 });
4234 relay_claude(&mut runtime, "SessionStart", &session_payload).unwrap();
4235 let sessions_before = SessionManager::new(runtime.repo.root())
4236 .list_sessions(true)
4237 .unwrap();
4238 let initial_segments = sessions_before
4239 .iter()
4240 .find(|s| !s.segments.is_empty())
4241 .map(|s| s.segments.len())
4242 .unwrap_or(0);
4243
4244 let prompt_payload = serde_json::json!({
4246 "session_id": "claude-prompt-sess",
4247 "model": {"id": "claude-opus-4-7", "display_name": "Claude Opus 4.7"},
4248 "prompt": "write a new feature",
4249 });
4250 relay_claude(&mut runtime, "UserPromptSubmit", &prompt_payload).unwrap();
4251 drop(runtime);
4252
4253 let verify = Repository::open(temp.path()).unwrap();
4254 let sessions_after = SessionManager::new(verify.root())
4255 .list_sessions(true)
4256 .unwrap();
4257 let rotated = sessions_after
4258 .iter()
4259 .any(|s| s.segments.len() > initial_segments);
4260 assert!(
4261 rotated,
4262 "UserPromptSubmit must add at least one segment beyond the SessionStart baseline \
4263 (initial={initial_segments}, sessions_after={:?})",
4264 sessions_after
4265 .iter()
4266 .map(|s| (s.id.clone(), s.segments.len()))
4267 .collect::<Vec<_>>(),
4268 );
4269 }
4270}