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::{
29 ActorPresence, ActorPresenceStatus, ActorPresenceStore, AgentUsageSummary, ObjectStore,
30 },
31};
32use refs::Head;
33use repo::{
34 Repository, SessionManager, Thread, ThreadFreshness, ThreadIntegrationPolicy, ThreadManager,
35 ThreadMode, ThreadState, TimelineStore, TimelineView,
36};
37use serde::{Deserialize, Serialize};
38use serde_json::Value;
39use wire::{
40 HarnessIdentity, ProgressCheckpoint, SessionDiffSummary, SessionReportEnvelope,
41 TranscriptAttachmentRef, UsageTotals, WorktreeChangeBaseline,
42};
43
44mod claude_hook;
45mod probe;
46
47use self::probe::{HarnessProbeInput, HarnessProbeResult, probe_harness_actor};
48use crate::{
49 cli::{
50 Cli,
51 commands::{
52 snapshot::{
53 SnapshotAgentOverrides, create_snapshot, summarize_confidence,
54 summarize_verification,
55 },
56 worktree_cmd::helpers::{prepare_worktree_target, write_isolated_checkout},
57 },
58 style, worktree_status_options,
59 },
60 config::{
61 HarnessMode, HarnessTranscriptMode, HarnessTransport, UserConfig, UserHarnessOverride,
62 UserHarnessRootThreadPolicy, UserHarnessSubagentThreadPolicy, UserThreadWorkspaceMode,
63 },
64};
65
66pub(crate) fn probe_current_process_harness(
67 repo: &Repository,
68 current_provider: Option<String>,
69 current_model: Option<String>,
70 current_policy: Option<String>,
71) -> Result<HarnessProbeResult> {
72 probe_harness_actor(&HarnessProbeInput {
73 argv: detected_harness_argv().or_else(|| Some(std::env::args().collect())),
74 env_hints: harness_env_hints(),
75 explicit_harness: None,
76 explicit_provider: None,
77 explicit_model: None,
78 explicit_thinking_level: None,
79 explicit_policy: None,
80 probe_metadata: BTreeMap::new(),
81 current_provider,
82 current_model,
83 current_policy,
84 repo_root: repo.root().display().to_string(),
85 })
86}
87
88fn harness_env_hints() -> BTreeMap<String, String> {
89 std::env::vars()
90 .filter(|(key, value)| {
91 !value.trim().is_empty()
92 && (key.starts_with("HEDDLE_AGENT_")
93 || key.starts_with("CODEX_")
94 || key.starts_with("CLAUDE")
95 || key.starts_with("ANTHROPIC_")
96 || key.starts_with("OPENAI_")
97 || key.starts_with("OPENCODE_")
98 || key.starts_with("AIDER_")
99 || matches!(
100 key.as_str(),
101 "MODEL" | "REASONING_EFFORT" | "THINKING_LEVEL"
102 ))
103 })
104 .collect()
105}
106
107fn detected_harness_argv() -> Option<Vec<String>> {
108 detected_harness_argv_impl()
109}
110
111#[cfg(target_os = "linux")]
112fn detected_harness_argv_impl() -> Option<Vec<String>> {
113 let mut pid = std::process::id();
114 for _ in 0..8 {
115 let stat = fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
116 let ppid = stat.split_whitespace().nth(3)?.parse::<u32>().ok()?;
117 if ppid == 0 || ppid == pid {
118 return None;
119 }
120 pid = ppid;
121 let raw = fs::read(format!("/proc/{pid}/cmdline")).ok()?;
122 let argv = raw
123 .split(|byte| *byte == 0)
124 .filter(|part| !part.is_empty())
125 .map(|part| String::from_utf8_lossy(part).to_string())
126 .collect::<Vec<_>>();
127 let program = argv.first().map(|arg| arg.to_ascii_lowercase())?;
128 if ["codex", "claude", "opencode", "aider"]
129 .iter()
130 .any(|needle| program.contains(needle))
131 {
132 return Some(argv);
133 }
134 }
135 None
136}
137
138#[cfg(not(target_os = "linux"))]
139fn detected_harness_argv_impl() -> Option<Vec<String>> {
140 None
141}
142
143pub fn cmd_harness_bridge(cli: &Cli) -> Result<()> {
144 let repo = cli.open_repo()?;
145 let mut runtime = init_harness_runtime(&repo)?;
146
147 let stdin = std::io::stdin();
148 let stdout = std::io::stdout();
149 let reader = BufReader::new(stdin.lock());
150 let mut writer = BufWriter::new(stdout.lock());
151
152 for line in reader.lines() {
153 let line = line?;
154 if line.trim().is_empty() {
155 continue;
156 }
157 let response = match serde_json::from_str::<BridgeRequest>(&line) {
158 Ok(request) => runtime.handle_request(request),
159 Err(err) => BridgeResponse::error(
160 None,
161 "invalid_request",
162 format!("failed to parse request: {err}"),
163 ),
164 };
165 serde_json::to_writer(&mut writer, &response)?;
166 writer.write_all(b"\n")?;
167 writer.flush()?;
168 }
169
170 Ok(())
171}
172
173pub(crate) fn relay_harness_event(
174 repo: &Repository,
175 harness: &str,
176 event: &str,
177 payload: &str,
178) -> Result<()> {
179 let mut runtime = init_harness_runtime(repo)?;
180 let (json, warning) = parse_relay_payload(payload);
181 if let Some(warning) = warning {
182 eprintln!("{}", style::warn(&warning));
183 }
184 match harness {
185 "codex" => relay_codex(&mut runtime, event, &json),
186 "claude-code" => relay_claude(&mut runtime, event, &json),
187 "opencode" => relay_opencode(&mut runtime, event, &json),
188 other => Err(anyhow!("unsupported harness relay: {other}")),
189 }
190}
191
192fn init_harness_runtime(repo: &Repository) -> Result<HarnessBridgeRuntime> {
193 let (user_config, warning) = load_harness_user_config(UserConfig::default_path());
194 if let Some(warning) = warning {
195 eprintln!("{}", style::warn(&warning));
196 }
197 Ok(HarnessBridgeRuntime::new(
198 Repository::open(repo.root())?,
199 user_config,
200 ))
201}
202
203fn load_harness_user_config(default_path: Option<PathBuf>) -> (UserConfig, Option<String>) {
204 let Some(path) = default_path else {
205 return (UserConfig::default(), None);
206 };
207 match UserConfig::load(&path) {
208 Ok(config) => (config, None),
209 Err(err) if is_not_found(&err) => (UserConfig::default(), None),
210 Err(err) => {
211 let warning = format!(
212 "warning: failed to load user config from {}: {err}; continuing with defaults",
213 path.display()
214 );
215 (UserConfig::default(), Some(warning))
216 }
217 }
218}
219
220fn is_not_found(err: &anyhow::Error) -> bool {
221 err.downcast_ref::<std::io::Error>()
222 .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound)
223}
224
225fn parse_relay_payload(payload: &str) -> (Value, Option<String>) {
226 core_parse_relay_payload(payload)
227}
228
229struct HarnessBridgeRuntime {
230 repo: Repository,
231 user_config: UserConfig,
232 reports: SessionReportStore,
233}
234
235struct RegistryEntryRequest<'a> {
236 heddle_session_id: &'a str,
237 thread_name: Option<&'a str>,
238 thread_id: Option<&'a str>,
239 identity: &'a ResolvedIdentity,
240 probe: &'a HarnessProbeResult,
241 attach: &'a ResolvedAttachment,
242 client_instance_id: Option<&'a str>,
243 requested_entry: Option<&'a ActorPresence>,
244}
245
246struct CanonicalActorSessionRequest<'a> {
247 tentative_session: Session,
248 tentative_owns_session: bool,
249 entry: &'a ActorPresence,
250 probe: &'a HarnessProbeResult,
251 attach: &'a mut ResolvedAttachment,
252}
253
254struct AttachmentResolutionInput<'a> {
255 requested_entry: Option<&'a ActorPresence>,
256 explicit_heddle_session_id: Option<&'a str>,
257 client_instance_id: Option<&'a str>,
258 probe: &'a HarnessProbeResult,
259 token_claims: Option<&'a TokenClaims>,
260}
261
262fn relay_codex(runtime: &mut HarnessBridgeRuntime, _event: &str, payload: &Value) -> Result<()> {
263 let metadata = map_from_pairs([
264 (
265 "client_name",
266 value_string(payload, &["client"]).or_else(|| value_string(payload, &["client_name"])),
267 ),
268 ("model", value_string(payload, &["model"])),
269 (
270 "model_provider",
271 value_string(payload, &["model_provider"])
272 .or_else(|| value_string(payload, &["provider"])),
273 ),
274 (
275 "model_reasoning_effort",
276 value_string(payload, &["reasoning_effort"]),
277 ),
278 ]);
279 let opened = runtime.open_session(OpenSessionParams {
280 harness: Some("codex".to_string()),
281 summary: value_string(payload, &["message"]),
282 probe_metadata: metadata,
283 ..OpenSessionParams::default()
284 })?;
285 runtime.update_progress(UpdateProgressParams {
286 heddle_session_id: opened.heddle_session_id,
287 summary: value_string(payload, &["message"]),
288 harness: Some("codex".to_string()),
289 ..UpdateProgressParams::default()
290 })?;
291 Ok(())
292}
293
294fn relay_claude(runtime: &mut HarnessBridgeRuntime, event: &str, payload: &Value) -> Result<()> {
295 let metadata = map_from_pairs([
296 ("session_id", value_string(payload, &["session_id"])),
297 ("agent_id", value_string(payload, &["agent_id"])),
298 ("session_name", value_string(payload, &["session_name"])),
299 (
300 "transcript_path",
301 value_string(payload, &["transcript_path"]),
302 ),
303 (
304 "model",
305 value_string(payload, &["model", "id"]).or_else(|| value_string(payload, &["model"])),
306 ),
307 (
308 "model_display_name",
309 value_string(payload, &["model", "display_name"]),
310 ),
311 ("effort", value_string(payload, &["effort"])),
312 ("hook_event", Some(event.to_string())),
313 (
314 "status_line",
315 (event == "StatusLine").then(|| "1".to_string()),
316 ),
317 (
318 "touched_paths",
319 value_array_join(payload, &["tool_response", "filePaths"])
320 .or_else(|| value_string(payload, &["file_path"])),
321 ),
322 (
323 "input_tokens",
324 value_u64_string(payload, &["context_window", "total_input_tokens"]),
325 ),
326 (
327 "output_tokens",
328 value_u64_string(payload, &["context_window", "total_output_tokens"]),
329 ),
330 (
331 "cost_micros_usd",
332 value_cost_micros(payload, &["cost", "total_cost_usd"]),
333 ),
334 ]);
335 let opened = runtime.open_session(OpenSessionParams {
336 harness: Some("claude-code".to_string()),
337 model: value_string(payload, &["model", "display_name"])
338 .or_else(|| value_string(payload, &["model", "id"]))
339 .or_else(|| value_string(payload, &["model"])),
340 summary: value_string(payload, &["message"]).or_else(|| value_string(payload, &["reason"])),
341 probe_metadata: metadata.clone(),
342 ..OpenSessionParams::default()
343 })?;
344 match event {
345 "SessionEnd" => {
346 runtime.close_session(CloseSessionParams {
347 heddle_session_id: opened.heddle_session_id,
348 summary: value_string(payload, &["reason"])
349 .or_else(|| value_string(payload, &["stop_hook_active"])),
350 outcome: Some("completed".to_string()),
351 ..CloseSessionParams::default()
352 })?;
353 }
354 "StatusLine" => {
355 runtime.update_progress(UpdateProgressParams {
356 heddle_session_id: opened.heddle_session_id.clone(),
357 harness: Some("claude-code".to_string()),
358 status: Some("StatusLine".to_string()),
359 message: value_string(payload, &["session_name"])
360 .or_else(|| value_string(payload, &["cwd"]))
361 .or_else(|| value_string(payload, &["workspace", "current_dir"])),
362 probe_metadata: metadata.clone(),
363 ..UpdateProgressParams::default()
364 })?;
365 runtime.record_usage(RecordUsageParams {
366 heddle_session_id: opened.heddle_session_id,
367 input_tokens: value_u64(payload, &["context_window", "total_input_tokens"]),
368 output_tokens: value_u64(payload, &["context_window", "total_output_tokens"]),
369 reasoning_tokens: value_u64(payload, &["context_window", "total_reasoning_tokens"]),
370 cache_creation_tokens: None,
371 cache_read_tokens: None,
372 tool_calls: None,
373 cost_micros_usd: value_cost_micros_u64(payload, &["cost", "total_cost_usd"]),
374 })?;
375 }
376 "Stop" => {
377 runtime.update_progress(UpdateProgressParams {
378 heddle_session_id: opened.heddle_session_id,
379 harness: Some("claude-code".to_string()),
380 status: Some("Stop".to_string()),
381 message: value_string(payload, &["message"])
382 .or_else(|| value_string(payload, &["result"]))
383 .or_else(|| value_string(payload, &["stop_reason"])),
384 probe_metadata: metadata,
385 ..UpdateProgressParams::default()
386 })?;
387 if let Err(err) = claude_hook::handle_stop_capture(
388 &runtime.repo,
389 &runtime.user_config,
390 payload,
391 "Claude Code turn",
392 ) {
393 tracing::warn!(?err, "heddle Stop hook capture failed");
394 }
395 }
396 "SubagentStop" => {
397 runtime.update_progress(UpdateProgressParams {
398 heddle_session_id: opened.heddle_session_id,
399 harness: Some("claude-code".to_string()),
400 status: Some("SubagentStop".to_string()),
401 touched_paths: csv_from_value(metadata.get("touched_paths")),
402 probe_metadata: metadata,
403 ..UpdateProgressParams::default()
404 })?;
405 if let Err(err) = claude_hook::handle_stop_capture(
406 &runtime.repo,
407 &runtime.user_config,
408 payload,
409 "Claude Code subagent turn",
410 ) {
411 tracing::warn!(?err, "heddle SubagentStop hook capture failed");
412 }
413 if let Err(err) = claude_hook::mark_subagent_complete(&runtime.repo, payload) {
414 tracing::debug!(?err, "heddle SubagentStop mark-complete failed");
415 }
416 }
417 "SubagentStart" => {
418 runtime.update_progress(UpdateProgressParams {
425 heddle_session_id: opened.heddle_session_id,
426 harness: Some("claude-code".to_string()),
427 status: Some("SubagentStart".to_string()),
428 touched_paths: csv_from_value(metadata.get("touched_paths")),
429 probe_metadata: metadata,
430 ..UpdateProgressParams::default()
431 })?;
432 }
433 "UserPromptSubmit" => {
434 runtime.update_progress(UpdateProgressParams {
435 heddle_session_id: opened.heddle_session_id.clone(),
436 harness: Some("claude-code".to_string()),
437 status: Some("UserPromptSubmit".to_string()),
438 touched_paths: csv_from_value(metadata.get("touched_paths")),
439 probe_metadata: metadata,
440 ..UpdateProgressParams::default()
441 })?;
442 if let Err(err) = claude_hook::handle_user_prompt_segment_rotate(
443 &runtime.repo,
444 &opened.heddle_session_id,
445 payload,
446 ) {
447 tracing::debug!(?err, "heddle UserPromptSubmit segment rotation failed");
448 }
449 }
450 "PreToolUse" => {
451 runtime.update_progress(UpdateProgressParams {
452 heddle_session_id: opened.heddle_session_id,
453 harness: Some("claude-code".to_string()),
454 status: Some("PreToolUse".to_string()),
455 touched_paths: csv_from_value(metadata.get("touched_paths")),
456 probe_metadata: metadata,
457 ..UpdateProgressParams::default()
458 })?;
459 if let Err(err) = claude_hook::handle_pre_tool_use(&runtime.repo, payload) {
460 tracing::debug!(?err, "heddle PreToolUse context inject skipped");
461 }
462 }
463 _ => {
464 runtime.update_progress(UpdateProgressParams {
465 heddle_session_id: opened.heddle_session_id,
466 harness: Some("claude-code".to_string()),
467 status: Some(event.to_string()),
468 touched_paths: csv_from_value(metadata.get("touched_paths")),
469 probe_metadata: metadata,
470 ..UpdateProgressParams::default()
471 })?;
472 }
473 }
474 Ok(())
475}
476
477fn relay_opencode(runtime: &mut HarnessBridgeRuntime, event: &str, payload: &Value) -> Result<()> {
478 let metadata = map_from_pairs([
479 (
480 "session_id",
481 value_string(payload, &["sessionID"])
482 .or_else(|| value_string(payload, &["session_id"])),
483 ),
484 (
485 "parent_id",
486 value_string(payload, &["parentID"]).or_else(|| value_string(payload, &["parent_id"])),
487 ),
488 (
489 "client_name",
490 value_string(payload, &["client"]).or_else(|| std::env::var("OPENCODE_CLIENT").ok()),
491 ),
492 ("model", value_string(payload, &["model"])),
493 ("provider", value_string(payload, &["provider"])),
494 ("hook_event", Some(event.to_string())),
495 (
496 "touched_paths",
497 value_string(payload, &["file", "path"]).or_else(|| value_string(payload, &["path"])),
498 ),
499 ]);
500 let opened = runtime.open_session(OpenSessionParams {
501 harness: Some("opencode".to_string()),
502 model: value_string(payload, &["model"]),
503 provider: value_string(payload, &["provider"]),
504 probe_metadata: metadata.clone(),
505 ..OpenSessionParams::default()
506 })?;
507 let session_id = opened.heddle_session_id.clone();
508 runtime.update_progress(UpdateProgressParams {
509 heddle_session_id: session_id.clone(),
510 harness: Some("opencode".to_string()),
511 status: Some(event.to_string()),
512 touched_paths: csv_from_value(metadata.get("touched_paths")),
513 probe_metadata: metadata,
514 ..UpdateProgressParams::default()
515 })?;
516 if let Err(err) = record_opencode_timeline_event(runtime, event, payload, &opened) {
517 tracing::debug!(?err, event, "heddle OpenCode timeline recording skipped");
518 }
519 Ok(())
520}
521
522#[derive(Clone, Copy, Debug, PartialEq, Eq)]
523enum TimelineToolEvent {
524 Started,
525 Finished,
526}
527
528trait HarnessTimelineExtractor {
529 fn timeline_event(&self, event: &str) -> Option<TimelineToolEvent>;
530 fn native_tool_call(&self, payload: &Value) -> Option<NativeToolCallRefV1>;
531 fn tool_name(&self, payload: &Value) -> String;
532 fn tool_status(&self, payload: &Value) -> TimelineToolCallStatus;
533 fn payload_metadata(&self, event: &str, payload: &Value)
534 -> Result<TimelineToolPayloadMetadata>;
535 fn touched_paths(&self, payload: &Value) -> Vec<String>;
536 fn capture_intent(&self, native: &NativeToolCallRefV1, payload: &Value) -> String;
537
538 fn timeline_thread(
539 &self,
540 runtime: &HarnessBridgeRuntime,
541 opened: &OpenSessionResult,
542 ) -> Result<String> {
543 if let Some(report) = runtime.reports.load(&opened.heddle_session_id)?
544 && let Some(thread) = report.thread
545 {
546 return Ok(thread);
547 }
548 match runtime.repo.head_ref()? {
549 Head::Attached { thread } => Ok(thread.to_string()),
550 Head::Detached { .. } => Ok("main".to_string()),
551 }
552 }
553
554 fn stable_step_id(&self, native: &NativeToolCallRefV1) -> TimelineStepId {
555 let key = format!(
556 "{}\0{}\0{}\0{}",
557 native.harness,
558 native.session_id.as_deref().unwrap_or(""),
559 native.message_id.as_deref().unwrap_or(""),
560 native.tool_call_id
561 );
562 let hash =
563 ContentHash::compute_typed("timeline-native-tool-call-v1", key.as_bytes()).to_hex();
564 TimelineStepId::new(format!("tls-{}", &hash[..24]))
565 }
566
567 fn started_labels(&self, _payload: &Value) -> Vec<TimelineLabel> {
568 vec![TimelineLabel::ExternalSideEffectsUnknown]
569 }
570
571 fn finished_labels(&self, changed: bool, _payload: &Value) -> Vec<TimelineLabel> {
572 if changed {
573 vec![
574 TimelineLabel::RepoReversible,
575 TimelineLabel::ExternalSideEffectsUnknown,
576 ]
577 } else {
578 vec![TimelineLabel::ExternalSideEffectsUnknown]
579 }
580 }
581}
582
583struct OpenCodeTimelineExtractor;
584
585impl HarnessTimelineExtractor for OpenCodeTimelineExtractor {
586 fn timeline_event(&self, event: &str) -> Option<TimelineToolEvent> {
587 match event {
588 "tool.execute.before" => Some(TimelineToolEvent::Started),
589 "tool.execute.after" => Some(TimelineToolEvent::Finished),
590 _ => None,
591 }
592 }
593
594 fn native_tool_call(&self, payload: &Value) -> Option<NativeToolCallRefV1> {
595 opencode_native_tool_call(payload)
596 }
597
598 fn tool_name(&self, payload: &Value) -> String {
599 opencode_tool_name(payload)
600 }
601
602 fn tool_status(&self, payload: &Value) -> TimelineToolCallStatus {
603 opencode_tool_status(payload)
604 }
605
606 fn payload_metadata(
607 &self,
608 event: &str,
609 payload: &Value,
610 ) -> Result<TimelineToolPayloadMetadata> {
611 opencode_payload_metadata(event, payload)
612 }
613
614 fn touched_paths(&self, payload: &Value) -> Vec<String> {
615 opencode_touched_paths(payload)
616 }
617
618 fn capture_intent(&self, native: &NativeToolCallRefV1, payload: &Value) -> String {
619 format!(
620 "OpenCode {} tool call {}",
621 self.tool_name(payload),
622 native.tool_call_id
623 )
624 }
625}
626
627fn record_opencode_timeline_event(
628 runtime: &mut HarnessBridgeRuntime,
629 event: &str,
630 payload: &Value,
631 opened: &OpenSessionResult,
632) -> Result<()> {
633 record_timeline_event(runtime, event, payload, opened, &OpenCodeTimelineExtractor)
634}
635
636fn record_timeline_event<E: HarnessTimelineExtractor>(
637 runtime: &mut HarnessBridgeRuntime,
638 event: &str,
639 payload: &Value,
640 opened: &OpenSessionResult,
641 extractor: &E,
642) -> Result<()> {
643 match extractor.timeline_event(event) {
644 Some(TimelineToolEvent::Started) => {
645 record_timeline_tool_started(runtime, event, payload, opened, extractor)
646 }
647 Some(TimelineToolEvent::Finished) => {
648 record_timeline_tool_finished(runtime, event, payload, opened, extractor)
649 }
650 None => Ok(()),
651 }
652}
653
654fn record_timeline_tool_started<E: HarnessTimelineExtractor>(
655 runtime: &mut HarnessBridgeRuntime,
656 event: &str,
657 payload: &Value,
658 opened: &OpenSessionResult,
659 extractor: &E,
660) -> Result<()> {
661 let Some(native) = extractor.native_tool_call(payload) else {
662 return Ok(());
663 };
664 let Some(before_state) = current_state_id(&runtime.repo)? else {
665 return Ok(());
666 };
667 let thread = extractor.timeline_thread(runtime, opened)?;
668 let store = TimelineStore::open(runtime.repo.heddle_dir())?;
669 let _record_guard = store.lock_recording(&thread)?;
670 let view = TimelineView::rebuild(&store)?;
671 let step_id = extractor.stable_step_id(&native);
672 let (branch_id, parent_step_id) = timeline_position_for_new_tool_step(&view, &thread, &step_id);
673 let envelope = TimelineOperationEnvelope::new(
674 TimelineOperationBodyV1::ToolCallStarted(ToolCallStartedV1 {
675 thread,
676 step_id,
677 branch_id,
678 parent_step_id,
679 native,
680 tool_name: extractor.tool_name(payload),
681 before_state,
682 payload: Some(extractor.payload_metadata(event, payload)?),
683 started_at_ms: Utc::now().timestamp_millis(),
684 }),
685 extractor.started_labels(payload),
686 );
687 store.write_operation(&envelope)?;
688 Ok(())
689}
690
691fn record_timeline_tool_finished<E: HarnessTimelineExtractor>(
692 runtime: &mut HarnessBridgeRuntime,
693 event: &str,
694 payload: &Value,
695 opened: &OpenSessionResult,
696 extractor: &E,
697) -> Result<()> {
698 let Some(native) = extractor.native_tool_call(payload) else {
699 return Ok(());
700 };
701 let Some(fallback_state) = current_state_id(&runtime.repo)? else {
702 return Ok(());
703 };
704 let thread = extractor.timeline_thread(runtime, opened)?;
705 let store = TimelineStore::open(runtime.repo.heddle_dir())?;
706 let _record_guard = store.lock_recording(&thread)?;
707 let before_view = TimelineView::rebuild(&store)?;
708 let step_id = extractor.stable_step_id(&native);
709 let (branch_id, _) = timeline_position_for_new_tool_step(&before_view, &thread, &step_id);
710 let before_state = before_view
711 .step(&thread, &step_id)
712 .and_then(|step| step.before_state)
713 .unwrap_or(fallback_state);
714 let has_worktree_changes_before_capture = !collect_worktree_changes(&runtime.repo)?.is_empty();
715 let mut capture_failed = false;
716 let capture_state = if !has_worktree_changes_before_capture {
717 None
718 } else {
719 let intent = extractor.capture_intent(&native, payload);
720 match create_snapshot(
721 &runtime.repo,
722 &runtime.user_config,
723 Some(intent),
724 None,
725 SnapshotAgentOverrides {
726 provider: opened.provider.clone(),
727 model: opened.model.clone(),
728 session: native.session_id.clone(),
729 segment: None,
730 policy: None,
731 no_policy: false,
732 no_agent: false,
733 },
734 ) {
735 Ok(_) => runtime.repo.head()?,
736 Err(err) => {
737 capture_failed = true;
738 tracing::warn!(?err, "heddle timeline tool capture failed");
739 None
740 }
741 }
742 };
743 let after_state = current_state_id(&runtime.repo)?.unwrap_or(fallback_state);
744 let mut touched_paths = extractor.touched_paths(payload);
745 merge_string_vec(
746 &mut touched_paths,
747 changed_paths_between_states(&runtime.repo, before_state, after_state)?,
748 );
749 let changed = before_state != after_state;
750 let mut labels = extractor.finished_labels(changed, payload);
751 if capture_failed {
752 merge_timeline_labels(&mut labels, vec![TimelineLabel::CaptureFailed]);
753 }
754 let envelope = TimelineOperationEnvelope::new(
755 TimelineOperationBodyV1::ToolCallFinished(ToolCallFinishedV1 {
756 thread,
757 step_id,
758 branch_id,
759 native,
760 status: extractor.tool_status(payload),
761 before_state,
762 after_state,
763 capture_state,
764 capture_oplog_batch_id: None,
765 changed,
766 touched_paths,
767 payload: Some(extractor.payload_metadata(event, payload)?),
768 finished_at_ms: Utc::now().timestamp_millis(),
769 }),
770 labels,
771 );
772 store.write_operation(&envelope)?;
773 Ok(())
774}
775
776fn opencode_native_tool_call(payload: &Value) -> Option<NativeToolCallRefV1> {
777 let tool_call_id = first_value_string(
778 payload,
779 &[
780 &["toolCallID"],
781 &["tool_call_id"],
782 &["toolCallId"],
783 &["callID"],
784 &["call_id"],
785 &["tool", "callID"],
786 &["tool", "call_id"],
787 &["tool", "id"],
788 &["toolCall", "id"],
789 &["tool_call", "id"],
790 &["id"],
791 ],
792 )?;
793 Some(NativeToolCallRefV1 {
794 harness: "opencode".to_string(),
795 session_id: value_string(payload, &["sessionID"])
796 .or_else(|| value_string(payload, &["session_id"])),
797 message_id: value_string(payload, &["messageID"])
798 .or_else(|| value_string(payload, &["message_id"]))
799 .or_else(|| value_string(payload, &["message", "id"])),
800 tool_call_id,
801 })
802}
803
804fn timeline_position_for_new_tool_step(
805 view: &TimelineView,
806 thread: &str,
807 step_id: &TimelineStepId,
808) -> (TimelineBranchId, Option<TimelineStepId>) {
809 let branch_id = view
810 .status(thread)
811 .and_then(|status| status.current_branch_id.clone())
812 .unwrap_or_else(|| TimelineBranchId::new("tlb-main"));
813 let parent_step_id = view
814 .status(thread)
815 .and_then(|status| status.current_step_id.clone())
816 .filter(|current| current != step_id);
817 (branch_id, parent_step_id)
818}
819
820fn current_state_id(repo: &Repository) -> Result<Option<StateId>> {
821 Ok(repo
822 .current_state()?
823 .map(|state| state.state_id)
824 .or(repo.head()?))
825}
826
827fn opencode_payload_metadata(event: &str, payload: &Value) -> Result<TimelineToolPayloadMetadata> {
830 let tool_name = opencode_tool_name(payload);
831 let tool_call_id = opencode_native_tool_call(payload)
832 .map(|native| native.tool_call_id)
833 .unwrap_or_default();
834 let raw = serde_json::to_vec(payload)?;
835 let hash = ContentHash::compute_typed("timeline-tool-payload", &raw);
836 let summary = if tool_call_id.is_empty() {
837 format!("OpenCode {event}: {tool_name}")
838 } else {
839 format!("OpenCode {event}: {tool_name} ({tool_call_id})")
840 };
841 Ok(TimelineToolPayloadMetadata {
842 summary: Some(summary),
843 hash: Some(hash),
844 })
845}
846
847fn opencode_touched_paths(payload: &Value) -> Vec<String> {
848 let mut paths = Vec::new();
849 for path in [
850 value_string(payload, &["file", "path"]),
851 value_string(payload, &["path"]),
852 value_string(payload, &["tool", "path"]),
853 value_string(payload, &["tool", "input", "file_path"]),
854 value_string(payload, &["input", "file_path"]),
855 ]
856 .into_iter()
857 .flatten()
858 {
859 if !path.trim().is_empty() && !paths.contains(&path) {
860 paths.push(path);
861 }
862 }
863 for value_path in [
864 &["paths"][..],
865 &["files"][..],
866 &["tool", "input", "paths"][..],
867 &["input", "paths"][..],
868 ] {
869 if let Some(items) = value_string_array(payload, value_path) {
870 merge_string_vec(&mut paths, items);
871 }
872 }
873 paths
874}
875
876fn merge_timeline_labels(target: &mut Vec<TimelineLabel>, incoming: Vec<TimelineLabel>) {
877 for label in incoming {
878 if !target.contains(&label) {
879 target.push(label);
880 }
881 }
882}
883
884fn csv_from_value(value: Option<&String>) -> Vec<String> {
885 value
886 .map(|value| {
887 value
888 .split(',')
889 .map(|item| item.trim().to_string())
890 .filter(|item| !item.is_empty())
891 .collect()
892 })
893 .unwrap_or_default()
894}
895
896impl HarnessBridgeRuntime {
897 fn new(repo: Repository, user_config: UserConfig) -> Self {
898 let reports = SessionReportStore::new(repo.root());
899 Self {
900 repo,
901 user_config,
902 reports,
903 }
904 }
905
906 fn handle_request(&mut self, request: BridgeRequest) -> BridgeResponse {
907 let response = match request.method.as_str() {
908 "open_session" => self
909 .decode_params::<OpenSessionParams>(request.params)
910 .and_then(|params| self.open_session(params))
911 .and_then(to_json_value),
912 "update_progress" => self
913 .decode_params::<UpdateProgressParams>(request.params)
914 .and_then(|params| self.update_progress(params))
915 .and_then(to_json_value),
916 "record_usage" => self
917 .decode_params::<RecordUsageParams>(request.params)
918 .and_then(|params| self.record_usage(params))
919 .and_then(to_json_value),
920 "record_touched_paths" => self
921 .decode_params::<RecordTouchedPathsParams>(request.params)
922 .and_then(|params| self.record_touched_paths(params))
923 .and_then(to_json_value),
924 "close_session" => self
925 .decode_params::<CloseSessionParams>(request.params)
926 .and_then(|params| self.close_session(params))
927 .and_then(to_json_value),
928 "flush_reports" => self
929 .decode_params::<FlushReportsParams>(request.params)
930 .and_then(|params| self.flush_reports(params))
931 .and_then(to_json_value),
932 other => Err(anyhow!("unknown method '{other}'")),
933 };
934
935 match response {
936 Ok(result) => BridgeResponse::ok(request.id, result),
937 Err(err) => BridgeResponse::error(request.id, "bridge_error", err.to_string()),
938 }
939 }
940
941 fn decode_params<T: for<'de> Deserialize<'de>>(&self, value: Value) -> Result<T> {
942 serde_json::from_value(value).map_err(|err| anyhow!(err))
943 }
944
945 fn open_session(&mut self, params: OpenSessionParams) -> Result<OpenSessionResult> {
946 if self.user_config.harness.mode == HarnessMode::Off {
947 return Err(anyhow!("harness integration is disabled in user config"));
948 }
949
950 let requested_transport = params
951 .transport
952 .unwrap_or(self.user_config.harness.transport);
953 let transcript_mode = params
954 .transcript_mode
955 .unwrap_or(self.user_config.harness.transcript);
956 let env_hints = merged_env_hints(¶ms.env_hints);
957 let token_claims = active_token_claims();
958 let current_session = SessionManager::new(self.repo.root()).get_current_session()?;
959 let current_segment = current_session
960 .as_ref()
961 .and_then(|session| session.current_segment());
962 let probe = probe_harness_actor(&HarnessProbeInput {
963 argv: params.argv.clone(),
964 env_hints: env_hints.clone(),
965 explicit_harness: params.harness.clone(),
966 explicit_provider: params.provider.clone(),
967 explicit_model: params.model.clone(),
968 explicit_thinking_level: params.thinking_level.clone(),
969 explicit_policy: params.policy.clone(),
970 probe_metadata: params.probe_metadata.clone(),
971 current_provider: current_segment.map(|segment| segment.provider.clone()),
972 current_model: current_segment.map(|segment| segment.model.clone()),
973 current_policy: current_segment.and_then(|segment| segment.policy_id.clone()),
974 repo_root: self.repo.root().display().to_string(),
975 })?;
976 let identity = resolve_identity(
977 &self.repo,
978 &self.user_config,
979 IdentityHints {
980 harness: params.harness.clone(),
981 provider: params.provider.clone(),
982 model: params.model.clone(),
983 thinking_level: params.thinking_level.clone(),
984 policy: params.policy.clone(),
985 probe: probe.clone(),
986 },
987 )?;
988 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
989 let requested_entry = resolve_requested_registry_entry(
990 ®istry,
991 params.agent_session_id.as_deref(),
992 params.client_instance_id.as_deref(),
993 )?;
994
995 if self.user_config.harness.mode == HarnessMode::Required
996 && (identity.harness.is_none()
997 || identity.provider.is_none()
998 || identity.model.is_none())
999 {
1000 return Err(anyhow!(
1001 "harness mode is 'required' but harness/provider/model could not be resolved"
1002 ));
1003 }
1004
1005 let mut sessions = SessionManager::new(self.repo.root());
1006 let principal = self.repo.get_principal()?;
1007 let mut attach = resolve_actor_attachment(
1008 ®istry,
1009 &self.repo,
1010 &mut sessions,
1011 AttachmentResolutionInput {
1012 requested_entry: requested_entry.as_ref(),
1013 explicit_heddle_session_id: params.heddle_session_id.as_deref(),
1014 client_instance_id: params.client_instance_id.as_deref(),
1015 probe: &probe,
1016 token_claims: token_claims.as_ref(),
1017 },
1018 )?;
1019 let (session, owns_session) = match &attach.target {
1020 AttachTarget::ExistingSession(session) => {
1021 let segment_id = session.current_segment_id.clone().unwrap_or_default();
1022 sessions.set_current_session(&session.id, &segment_id)?;
1023 (session.clone(), false)
1024 }
1025 AttachTarget::CreateNew {
1026 _because_claimed: _,
1027 } => {
1028 let session = sessions.start_session(
1029 principal,
1030 identity
1031 .provider
1032 .clone()
1033 .unwrap_or_else(|| "unknown".to_string()),
1034 identity
1035 .model
1036 .clone()
1037 .unwrap_or_else(|| "unknown".to_string()),
1038 identity.policy.clone(),
1039 )?;
1040 (session, true)
1041 }
1042 };
1043
1044 let (thread_name, thread_id) =
1045 self.resolve_harness_thread_binding(¶ms, &probe, &identity)?;
1046 let entry = self.ensure_registry_entry(RegistryEntryRequest {
1047 heddle_session_id: &session.id,
1048 thread_name: thread_name.as_deref(),
1049 thread_id: thread_id.as_deref(),
1050 identity: &identity,
1051 probe: &probe,
1052 attach: &attach,
1053 client_instance_id: params.client_instance_id.as_deref(),
1054 requested_entry: requested_entry.as_ref(),
1055 })?;
1056 let (session, owns_session) = self.reuse_canonical_actor_session(
1057 &mut sessions,
1058 CanonicalActorSessionRequest {
1059 tentative_session: session,
1060 tentative_owns_session: owns_session,
1061 entry: &entry,
1062 probe: &probe,
1063 attach: &mut attach,
1064 },
1065 )?;
1066
1067 let mut segment_id = session.current_segment_id.clone().unwrap_or_default();
1068 if should_rotate_segment(&session, &identity) {
1069 let segment = sessions.add_segment(
1070 &session.id,
1071 identity
1072 .provider
1073 .clone()
1074 .unwrap_or_else(|| "unknown".to_string()),
1075 identity
1076 .model
1077 .clone()
1078 .unwrap_or_else(|| "unknown".to_string()),
1079 identity.policy.clone(),
1080 )?;
1081 segment_id = segment.id;
1082 }
1083
1084 let base_state = self
1085 .repo
1086 .current_state()?
1087 .map(|state| state.state_id.to_string_full())
1088 .or_else(|| {
1089 self.repo
1090 .head()
1091 .ok()
1092 .flatten()
1093 .map(|id| id.to_string_full())
1094 });
1095 let worktree_changes_at_open = capture_worktree_change_snapshot(&self.repo)?;
1096 let opened_at = Utc::now().to_rfc3339();
1097 let mut report = SessionReportEnvelope {
1098 version: 1,
1099 heddle_session_id: session.id.clone(),
1100 heddle_segment_id: (!segment_id.is_empty()).then_some(segment_id.clone()),
1101 agent_session_id: Some(entry.session_id.clone()),
1102 client_instance_id: entry.client_instance_id.clone(),
1103 native_actor_key: entry.native_actor_key.clone(),
1104 native_parent_actor_key: entry.native_parent_actor_key.clone(),
1105 native_instance_key: entry.native_instance_key.clone(),
1106 repo_root: self.repo.root().display().to_string(),
1107 thread: thread_name.clone(),
1108 thread_id,
1109 task: params.task.clone(),
1110 summary: params.summary.clone(),
1111 opened_at,
1112 closed_at: None,
1113 base_state_at_open: base_state.clone(),
1114 worktree_changes_at_open,
1115 head_state_at_close: None,
1116 transport_mode: transport_mode_name(requested_transport).to_string(),
1117 transcript_mode: transcript_mode_name(transcript_mode).to_string(),
1118 outcome: None,
1119 harness: identity.to_transport_identity(),
1120 progress: Vec::new(),
1121 usage: UsageTotals::default(),
1122 touched_paths: Vec::new(),
1123 changed_paths: Vec::new(),
1124 diff_summary: None,
1125 transcript_refs: Vec::new(),
1126 last_progress_at: None,
1127 report_flush_state: Some("pending-local".to_string()),
1128 attach_reason: Some(attach.attach_reason.clone()),
1129 attach_precedence: attach.precedence.clone(),
1130 winning_attach_rule: Some(attach.winning_rule.clone()),
1131 probe_source: probe.probe_source.clone(),
1132 probe_confidence: probe.confidence,
1133 pending_flush: true,
1134 last_flushed_at: None,
1135 owns_session,
1136 };
1137 merge_unique_paths(&mut report.touched_paths, probe.touched_paths.clone());
1138 merge_usage(&mut report.usage, &probe.usage_totals);
1139 if transcript_mode != HarnessTranscriptMode::Off {
1140 report.transcript_refs = probe.transcript_refs.clone();
1141 }
1142 self.reports.save(&report)?;
1143 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1144 if matches!(requested_transport, HarnessTransport::Direct) {
1145 enqueue_report(&self.reports, &mut report)?;
1146 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1147 }
1148
1149 Ok(OpenSessionResult {
1150 heddle_session_id: report.heddle_session_id.clone(),
1151 heddle_segment_id: report.heddle_segment_id.clone(),
1152 agent_session_id: report.agent_session_id.clone(),
1153 created_session: owns_session,
1154 harness: report.harness.harness.clone(),
1155 provider: report.harness.provider.clone(),
1156 model: report.harness.model.clone(),
1157 thinking_level: report.harness.thinking_level.clone(),
1158 report_flush_state: report.report_flush_state.clone(),
1159 attach_reason: report.attach_reason.clone(),
1160 })
1161 }
1162
1163 fn update_progress(&mut self, params: UpdateProgressParams) -> Result<SessionMutationResult> {
1164 let mut report = self
1165 .reports
1166 .load(¶ms.heddle_session_id)?
1167 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1168 let current_session = SessionManager::new(self.repo.root()).get_current_session()?;
1169 let current_segment = current_session
1170 .as_ref()
1171 .and_then(|session| session.current_segment());
1172 let probe = probe_harness_actor(&HarnessProbeInput {
1173 argv: params.argv.clone(),
1174 env_hints: merged_env_hints(¶ms.env_hints),
1175 explicit_harness: params.harness.clone(),
1176 explicit_provider: params.provider.clone(),
1177 explicit_model: params.model.clone(),
1178 explicit_thinking_level: params.thinking_level.clone(),
1179 explicit_policy: params.policy.clone(),
1180 probe_metadata: params.probe_metadata.clone(),
1181 current_provider: current_segment.map(|segment| segment.provider.clone()),
1182 current_model: current_segment.map(|segment| segment.model.clone()),
1183 current_policy: current_segment.and_then(|segment| segment.policy_id.clone()),
1184 repo_root: self.repo.root().display().to_string(),
1185 })?;
1186 let identity = resolve_identity(
1187 &self.repo,
1188 &self.user_config,
1189 IdentityHints {
1190 harness: params.harness.clone(),
1191 provider: params.provider.clone(),
1192 model: params.model.clone(),
1193 thinking_level: params.thinking_level.clone(),
1194 policy: params.policy.clone(),
1195 probe: probe.clone(),
1196 },
1197 )?;
1198 self.ensure_segment_for_report(&mut report, &identity)?;
1199 if report.harness.harness.is_none() {
1200 report.harness.harness = identity.harness.clone();
1201 }
1202 if report.harness.provider.is_none() {
1203 report.harness.provider = identity.provider.clone();
1204 }
1205 if report.harness.model.is_none() {
1206 report.harness.model = identity.model.clone();
1207 }
1208 if report.harness.thinking_level.is_none() {
1209 report.harness.thinking_level = identity.thinking_level.clone();
1210 }
1211 if report.harness.policy.is_none() {
1212 report.harness.policy = identity.policy.clone();
1213 }
1214 if report.native_actor_key.is_none() {
1215 report.native_actor_key = probe.native_actor_key.clone();
1216 }
1217 if report.native_parent_actor_key.is_none() {
1218 report.native_parent_actor_key = probe.native_parent_actor_key.clone();
1219 }
1220 if report.native_instance_key.is_none() {
1221 report.native_instance_key = probe.native_instance_key.clone();
1222 }
1223 if report.probe_source.is_none() {
1224 report.probe_source = probe.probe_source.clone();
1225 }
1226 if report.probe_confidence.is_none() {
1227 report.probe_confidence = probe.confidence;
1228 }
1229
1230 let recorded_at = Utc::now().to_rfc3339();
1231 let checkpoint = ProgressCheckpoint {
1232 status: params.status.clone(),
1233 message: params.message.clone(),
1234 completed_steps: params.completed_steps,
1235 total_steps: params.total_steps,
1236 touched_paths: normalize_paths(
1237 params
1238 .touched_paths
1239 .into_iter()
1240 .chain(probe.touched_paths)
1241 .collect::<Vec<_>>(),
1242 ),
1243 recorded_at: recorded_at.clone(),
1244 };
1245 merge_unique_paths(
1246 &mut report.touched_paths,
1247 checkpoint.touched_paths.iter().cloned(),
1248 );
1249 merge_usage(&mut report.usage, &probe.usage_totals);
1250 if report.transcript_mode != "off" && report.transcript_refs.is_empty() {
1251 report.transcript_refs = probe.transcript_refs;
1252 }
1253 report.progress.push(checkpoint);
1254 if let Some(summary) = params.summary {
1255 report.summary = Some(summary);
1256 }
1257 report.last_progress_at = Some(recorded_at);
1258 mark_pending_flush(&mut report);
1259 self.persist_report(report)
1260 }
1261
1262 fn resolve_harness_thread_binding(
1263 &self,
1264 params: &OpenSessionParams,
1265 probe: &HarnessProbeResult,
1266 identity: &ResolvedIdentity,
1267 ) -> Result<(Option<String>, Option<String>)> {
1268 if let Some(thread) = params.thread.clone() {
1269 let thread_id = thread_id_for_name(&self.repo, Some(&thread))?;
1270 return Ok((Some(thread), thread_id));
1271 }
1272
1273 let current_attached = match self.repo.head_ref()? {
1274 Head::Attached { thread } => Some(thread.to_string()),
1275 Head::Detached { .. } => None,
1276 };
1277
1278 if !probe.attach_hints.root_actor
1279 && self.user_config.harness.threading.subagent
1280 == UserHarnessSubagentThreadPolicy::CreateChild
1281 && let Some(parent_thread) =
1282 resolve_parent_thread_for_subagent(&self.repo, probe, current_attached.as_deref())?
1283 && can_create_harness_thread(&self.repo, Some(&parent_thread), Some(&parent_thread))?
1284 {
1285 let name = allocate_thread_name(
1286 &self.repo,
1287 &format!(
1288 "{}/{}",
1289 parent_thread,
1290 sanitize_name(&preferred_thread_slug(params, probe, identity))
1291 ),
1292 )?;
1293 self.ensure_harness_thread(
1294 &name,
1295 Some(&parent_thread),
1296 Some(&parent_thread),
1297 params.task.clone(),
1298 )?;
1299 let thread_id = thread_id_for_name(&self.repo, Some(&name))?;
1300 return Ok((Some(name), thread_id));
1301 }
1302
1303 if probe.attach_hints.root_actor
1304 && self.user_config.harness.threading.root_actor
1305 == UserHarnessRootThreadPolicy::CreateNew
1306 && let Some(current) = current_attached.clone()
1307 && can_create_harness_thread(&self.repo, Some(¤t), None)?
1308 {
1309 let name = allocate_thread_name(
1310 &self.repo,
1311 &format!(
1312 "{}/{}",
1313 current,
1314 sanitize_name(&preferred_thread_slug(params, probe, identity))
1315 ),
1316 )?;
1317 self.ensure_harness_thread(&name, Some(¤t), None, params.task.clone())?;
1318 let thread_id = thread_id_for_name(&self.repo, Some(&name))?;
1319 return Ok((Some(name), thread_id));
1320 }
1321
1322 let thread_id = thread_id_for_name(&self.repo, current_attached.as_deref())?;
1323 Ok((current_attached, thread_id))
1324 }
1325
1326 fn ensure_harness_thread(
1327 &self,
1328 name: &str,
1329 target_thread: Option<&str>,
1330 parent_thread: Option<&str>,
1331 task: Option<String>,
1332 ) -> Result<()> {
1333 let manager = ThreadManager::new(self.repo.heddle_dir());
1334 if manager.load(name)?.is_some() {
1335 return Ok(());
1336 }
1337
1338 let base_state = self
1339 .resolve_harness_thread_base_state(target_thread, parent_thread)?
1340 .ok_or_else(|| anyhow!("No current state to start a thread from"))?;
1341 let tn = ThreadName::new(name);
1342 if self.repo.refs().get_thread(&tn)?.is_none() {
1343 self.repo
1344 .set_thread_recorded_cas(&tn, refs::RefExpectation::Missing, &base_state)?;
1345 }
1346
1347 let workspace_mode = self
1348 .user_config
1349 .harness
1350 .threading
1351 .workspace_default
1352 .unwrap_or(UserThreadWorkspaceMode::Materialized);
1353 let thread_mode = match workspace_mode {
1354 UserThreadWorkspaceMode::Materialized | UserThreadWorkspaceMode::Auto => {
1355 ThreadMode::Materialized
1356 }
1357 UserThreadWorkspaceMode::Virtualized => ThreadMode::Virtualized,
1358 UserThreadWorkspaceMode::Solid => ThreadMode::Solid,
1359 };
1360 let path = match thread_mode {
1361 ThreadMode::Solid | ThreadMode::Materialized => {
1362 default_private_thread_path(&self.repo, name)
1363 }
1364 ThreadMode::Virtualized => default_private_thread_path(&self.repo, name),
1367 };
1368 let abs_path = prepare_worktree_target(&self.repo, &path, Some(name))?.path;
1369 write_isolated_checkout(&self.repo, &abs_path, &base_state, Some(name))?;
1370
1371 let base_state_obj = self
1372 .repo
1373 .store()
1374 .get_state(&base_state)?
1375 .ok_or_else(|| anyhow!("Base state '{}' not found", base_state.short()))?;
1376 let thread = Thread {
1377 id: name.to_string(),
1378 thread: name.to_string(),
1379 target_thread: target_thread.map(ToString::to_string),
1380 parent_thread: parent_thread.map(ToString::to_string),
1381 mode: thread_mode.clone(),
1382 state: ThreadState::Active,
1383 base_state: base_state.short(),
1384 base_root: base_state_obj.tree.short(),
1385 current_state: Some(base_state.short()),
1386 merged_state: None,
1387 task,
1388 execution_path: abs_path.clone(),
1389 materialized_path: match thread_mode {
1390 ThreadMode::Solid => Some(abs_path),
1391 ThreadMode::Materialized | ThreadMode::Virtualized => None,
1395 },
1396 changed_paths: vec![],
1397 impact_categories: vec![],
1398 heavy_impact_paths: vec![],
1399 promotion_suggested: false,
1400 freshness: if target_thread.is_some() {
1401 ThreadFreshness::Current
1402 } else {
1403 ThreadFreshness::Unknown
1404 },
1405 verification_summary: summarize_verification(base_state_obj.verification.as_ref()),
1406 confidence_summary: summarize_confidence(base_state_obj.confidence),
1407 integration_policy_result: ThreadIntegrationPolicy::default(),
1408 created_at: Utc::now(),
1409 updated_at: Utc::now(),
1410 ephemeral: None,
1411 auto: true,
1416 shared_target_dir: None,
1419 };
1420 manager.save(&thread)?;
1421 Ok(())
1422 }
1423
1424 fn resolve_harness_thread_base_state(
1425 &self,
1426 target_thread: Option<&str>,
1427 parent_thread: Option<&str>,
1428 ) -> Result<Option<objects::object::StateId>> {
1429 resolve_harness_thread_base_state(&self.repo, target_thread, parent_thread)
1430 }
1431
1432 fn record_usage(&mut self, params: RecordUsageParams) -> Result<SessionMutationResult> {
1433 let mut report = self
1434 .reports
1435 .load(¶ms.heddle_session_id)?
1436 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1437 if let Some(input) = params.input_tokens {
1438 report.usage.input_tokens = Some(max_u64(report.usage.input_tokens, input));
1439 }
1440 if let Some(output) = params.output_tokens {
1441 report.usage.output_tokens = Some(max_u64(report.usage.output_tokens, output));
1442 }
1443 if let Some(reasoning) = params.reasoning_tokens {
1444 report.usage.reasoning_tokens = Some(max_u64(report.usage.reasoning_tokens, reasoning));
1445 }
1446 if let Some(cache_creation) = params.cache_creation_tokens {
1447 report.usage.cache_creation_tokens =
1448 Some(max_u64(report.usage.cache_creation_tokens, cache_creation));
1449 }
1450 if let Some(cache_read) = params.cache_read_tokens {
1451 report.usage.cache_read_tokens =
1452 Some(max_u64(report.usage.cache_read_tokens, cache_read));
1453 }
1454 if let Some(tool_calls) = params.tool_calls {
1455 report.usage.tool_calls = Some(max_u32(report.usage.tool_calls, tool_calls));
1456 }
1457 if let Some(cost) = params.cost_micros_usd {
1458 report.usage.cost_micros_usd = Some(max_u64(report.usage.cost_micros_usd, cost));
1459 }
1460 mark_pending_flush(&mut report);
1461 self.persist_report(report)
1462 }
1463
1464 fn record_touched_paths(
1465 &mut self,
1466 params: RecordTouchedPathsParams,
1467 ) -> Result<SessionMutationResult> {
1468 let mut report = self
1469 .reports
1470 .load(¶ms.heddle_session_id)?
1471 .ok_or_else(|| anyhow!("session report not found for {}", params.heddle_session_id))?;
1472 merge_unique_paths(&mut report.touched_paths, normalize_paths(params.paths));
1473 mark_pending_flush(&mut report);
1474 self.persist_report(report)
1475 }
1476
1477 fn close_session(&mut self, params: CloseSessionParams) -> Result<CloseSessionResult> {
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 report.closed_at = Some(Utc::now().to_rfc3339());
1483 report.outcome = params.outcome.clone();
1484 if let Some(summary) = params.summary {
1485 report.summary = Some(summary);
1486 }
1487 if let Some(transcript_refs) = params.transcript_refs {
1488 report.transcript_refs = transcript_refs;
1489 }
1490 let final_diff = compute_final_diff(
1491 &self.repo,
1492 report.base_state_at_open.as_deref(),
1493 &report.worktree_changes_at_open,
1494 )?;
1495 report.head_state_at_close = final_diff.head_state;
1496 report.changed_paths = final_diff.changed_paths;
1497 report.diff_summary = Some(final_diff.diff_summary);
1498 mark_pending_flush(&mut report);
1499 if report.owns_session {
1500 let mut sessions = SessionManager::new(self.repo.root());
1501 if let Ok(Some(session)) = sessions.get_session(&report.heddle_session_id)
1502 && session.is_active()
1503 {
1504 let _ = sessions.end_session(Some(&report.heddle_session_id));
1505 }
1506 }
1507
1508 let transport = params
1509 .transport
1510 .unwrap_or(self.user_config.harness.transport);
1511 if matches!(transport, HarnessTransport::Direct | HarnessTransport::End) {
1512 enqueue_report(&self.reports, &mut report)?;
1513 } else {
1514 self.reports.save(&report)?;
1515 }
1516 self.sync_registry_from_report(&report, ActorPresenceStatus::Complete)?;
1517 Ok(CloseSessionResult {
1518 heddle_session_id: report.heddle_session_id,
1519 changed_paths: report.changed_paths,
1520 diff_summary: report.diff_summary.unwrap_or_default(),
1521 report_flush_state: report.report_flush_state,
1522 })
1523 }
1524
1525 fn flush_reports(&mut self, params: FlushReportsParams) -> Result<FlushReportsResult> {
1526 let mut flushed = 0usize;
1527 let session_ids = match params.heddle_session_id {
1528 Some(session_id) => vec![session_id],
1529 None => self.reports.list_pending()?,
1530 };
1531 for session_id in session_ids {
1532 let Some(mut report) = self.reports.load(&session_id)? else {
1533 continue;
1534 };
1535 if !report.pending_flush {
1536 continue;
1537 }
1538 enqueue_report(&self.reports, &mut report)?;
1539 let status = if report.closed_at.is_some() {
1540 ActorPresenceStatus::Complete
1541 } else {
1542 ActorPresenceStatus::Active
1543 };
1544 self.sync_registry_from_report(&report, status)?;
1545 flushed += 1;
1546 }
1547 Ok(FlushReportsResult { flushed })
1548 }
1549
1550 fn persist_report(
1551 &mut self,
1552 mut report: SessionReportEnvelope,
1553 ) -> Result<SessionMutationResult> {
1554 let transport = transport_from_report(&report, self.user_config.harness.transport);
1555 match transport {
1556 HarnessTransport::Direct => {
1557 enqueue_report(&self.reports, &mut report)?;
1558 }
1559 HarnessTransport::Spool | HarnessTransport::End => {
1560 self.reports.save(&report)?;
1561 }
1562 }
1563 self.sync_registry_from_report(&report, ActorPresenceStatus::Active)?;
1564 Ok(SessionMutationResult {
1565 heddle_session_id: report.heddle_session_id,
1566 heddle_segment_id: report.heddle_segment_id,
1567 report_flush_state: report.report_flush_state,
1568 })
1569 }
1570
1571 fn ensure_segment_for_report(
1572 &self,
1573 report: &mut SessionReportEnvelope,
1574 identity: &ResolvedIdentity,
1575 ) -> Result<()> {
1576 let mut sessions = SessionManager::new(self.repo.root());
1577 let Some(session) = sessions.get_session(&report.heddle_session_id)? else {
1578 return Ok(());
1579 };
1580 if !session.is_active() || !should_rotate_segment(&session, identity) {
1581 return Ok(());
1582 }
1583 let segment = sessions.add_segment(
1584 &report.heddle_session_id,
1585 identity
1586 .provider
1587 .clone()
1588 .unwrap_or_else(|| "unknown".to_string()),
1589 identity
1590 .model
1591 .clone()
1592 .unwrap_or_else(|| "unknown".to_string()),
1593 identity.policy.clone(),
1594 )?;
1595 report.heddle_segment_id = Some(segment.id);
1596 if identity.provider.is_some() {
1597 report.harness.provider = identity.provider.clone();
1598 }
1599 if identity.model.is_some() {
1600 report.harness.model = identity.model.clone();
1601 }
1602 if identity.policy.is_some() {
1603 report.harness.policy = identity.policy.clone();
1604 }
1605 if identity.thinking_level.is_some() {
1606 report.harness.thinking_level = identity.thinking_level.clone();
1607 }
1608 Ok(())
1609 }
1610
1611 fn ensure_registry_entry(&self, request: RegistryEntryRequest<'_>) -> Result<ActorPresence> {
1612 let RegistryEntryRequest {
1613 heddle_session_id,
1614 thread_name,
1615 thread_id,
1616 identity,
1617 probe,
1618 attach,
1619 client_instance_id,
1620 requested_entry,
1621 } = request;
1622 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
1623 let fallback_entry = if client_instance_id.is_some()
1624 || probe.native_actor_key.is_some()
1625 || probe.native_instance_key.is_some()
1626 {
1627 None
1628 } else {
1629 find_matching_registry_entry(®istry, &self.repo, heddle_session_id, thread_name)?
1630 };
1631 if let Some(entry) = requested_entry
1632 .cloned()
1633 .or_else(|| attach.matched_entry.clone())
1634 .or(fallback_entry)
1635 {
1636 return registry
1637 .update_entry(&entry.session_id, |existing| {
1638 if client_instance_id.is_some() {
1639 existing.client_instance_id = client_instance_id.map(ToString::to_string);
1640 }
1641 if probe.native_actor_key.is_some() {
1642 existing.native_actor_key = probe.native_actor_key.clone();
1643 }
1644 if probe.native_parent_actor_key.is_some() {
1645 existing.native_parent_actor_key = probe.native_parent_actor_key.clone();
1646 }
1647 if probe.native_instance_key.is_some() {
1648 existing.native_instance_key = probe.native_instance_key.clone();
1649 }
1650 existing.heddle_session_id = Some(heddle_session_id.to_string());
1651 existing.thread_id = thread_id.map(ToString::to_string);
1652 if let Some(thread_name) = thread_name {
1653 existing.thread = thread_name.to_string();
1654 }
1655 existing.path = Some(self.repo.root().to_path_buf());
1656 if identity.provider.is_some() {
1657 existing.provider = identity.provider.clone();
1658 }
1659 if identity.model.is_some() {
1660 existing.model = identity.model.clone();
1661 }
1662 if identity.harness.is_some() {
1663 existing.harness = identity.harness.clone();
1664 }
1665 if identity.thinking_level.is_some() {
1666 existing.thinking_level = identity.thinking_level.clone();
1667 }
1668 existing.attach_reason = Some(attach.attach_reason.clone());
1669 existing.attach_precedence = attach.precedence.clone();
1670 existing.winning_attach_rule = Some(attach.winning_rule.clone());
1671 existing.probe_source = probe.probe_source.clone();
1672 existing.probe_confidence = probe.confidence;
1673 existing.status = ActorPresenceStatus::Active;
1674 existing.completed_at = None;
1675 })?
1676 .ok_or_else(|| anyhow!("registry entry disappeared during update"));
1677 }
1678
1679 if client_instance_id.is_none() && probe.native_actor_key.is_some() {
1680 let (entry, _) = registry.find_or_create_active_entry(
1681 |entry| {
1682 claude_actor_compatible(entry, probe, self.repo.root())
1683 && entry.native_actor_key == probe.native_actor_key
1684 },
1685 |existing| {
1686 if client_instance_id.is_some() {
1687 existing.client_instance_id = client_instance_id.map(ToString::to_string);
1688 }
1689 if existing.heddle_session_id.is_none() {
1690 existing.heddle_session_id = Some(heddle_session_id.to_string());
1691 }
1692 existing.thread_id = thread_id.map(ToString::to_string);
1693 if let Some(thread_name) = thread_name {
1694 existing.thread = thread_name.to_string();
1695 }
1696 existing.path = Some(self.repo.root().to_path_buf());
1697 if identity.provider.is_some() {
1698 existing.provider = identity.provider.clone();
1699 }
1700 if identity.model.is_some() {
1701 existing.model = identity.model.clone();
1702 }
1703 if identity.harness.is_some() {
1704 existing.harness = identity.harness.clone();
1705 }
1706 if identity.thinking_level.is_some() {
1707 existing.thinking_level = identity.thinking_level.clone();
1708 }
1709 if probe.native_parent_actor_key.is_some() {
1710 existing.native_parent_actor_key = probe.native_parent_actor_key.clone();
1711 }
1712 if probe.native_instance_key.is_some() {
1713 existing.native_instance_key = probe.native_instance_key.clone();
1714 }
1715 existing.attach_reason = Some(attach.attach_reason.clone());
1716 existing.attach_precedence = attach.precedence.clone();
1717 existing.winning_attach_rule = Some(attach.winning_rule.clone());
1718 existing.probe_source = probe.probe_source.clone();
1719 existing.probe_confidence = probe.confidence;
1720 existing.status = ActorPresenceStatus::Active;
1721 existing.completed_at = None;
1722 },
1723 |session_id| {
1724 Ok(ActorPresence {
1725 session_id: session_id.to_string(),
1726 client_instance_id: client_instance_id.map(ToString::to_string),
1727 native_actor_key: probe.native_actor_key.clone(),
1728 native_parent_actor_key: probe.native_parent_actor_key.clone(),
1729 native_instance_key: probe.native_instance_key.clone(),
1730 heddle_session_id: Some(heddle_session_id.to_string()),
1731 thread_id: thread_id.map(ToString::to_string),
1732 thread: thread_name.unwrap_or("detached").to_string(),
1733 anchor_state: self.repo.head()?.map(|id| id.to_string_full()),
1734 anchor_root: None,
1735 path: Some(self.repo.root().to_path_buf()),
1736 base_state: self.repo.head()?.map(|id| id.short()).unwrap_or_default(),
1737 started_at: Utc::now(),
1738 provider: identity.provider.clone(),
1739 model: identity.model.clone(),
1740 harness: identity.harness.clone(),
1741 thinking_level: identity.thinking_level.clone(),
1742 usage_summary: AgentUsageSummary::default(),
1743 last_progress_at: None,
1744 report_flush_state: Some("pending-local".to_string()),
1745 attach_reason: Some(attach.attach_reason.clone()),
1746 task_assignment_id: None,
1747 attach_precedence: attach.precedence.clone(),
1748 winning_attach_rule: Some(attach.winning_rule.clone()),
1749 probe_source: probe.probe_source.clone(),
1750 probe_confidence: probe.confidence,
1751 status: ActorPresenceStatus::Active,
1752 completed_at: None,
1753 context_queries: vec![],
1754 })
1755 },
1756 )?;
1757 return Ok(entry);
1758 }
1759
1760 Ok(registry.create_generated_entry(|session_id| {
1761 Ok(ActorPresence {
1762 session_id: session_id.to_string(),
1763 client_instance_id: client_instance_id.map(ToString::to_string),
1764 native_actor_key: probe.native_actor_key.clone(),
1765 native_parent_actor_key: probe.native_parent_actor_key.clone(),
1766 native_instance_key: probe.native_instance_key.clone(),
1767 heddle_session_id: Some(heddle_session_id.to_string()),
1768 thread_id: thread_id.map(ToString::to_string),
1769 thread: thread_name.unwrap_or("detached").to_string(),
1770 anchor_state: self.repo.head()?.map(|id| id.to_string_full()),
1771 anchor_root: None,
1772 path: Some(self.repo.root().to_path_buf()),
1773 base_state: self.repo.head()?.map(|id| id.short()).unwrap_or_default(),
1774 started_at: Utc::now(),
1775 provider: identity.provider.clone(),
1776 model: identity.model.clone(),
1777 harness: identity.harness.clone(),
1778 thinking_level: identity.thinking_level.clone(),
1779 usage_summary: AgentUsageSummary::default(),
1780 last_progress_at: None,
1781 report_flush_state: Some("pending-local".to_string()),
1782 attach_reason: Some(attach.attach_reason.clone()),
1783 task_assignment_id: None,
1784 attach_precedence: attach.precedence.clone(),
1785 winning_attach_rule: Some(attach.winning_rule.clone()),
1786 probe_source: probe.probe_source.clone(),
1787 probe_confidence: probe.confidence,
1788 status: ActorPresenceStatus::Active,
1789 completed_at: None,
1790 context_queries: vec![],
1791 })
1792 })?)
1793 }
1794
1795 fn reuse_canonical_actor_session(
1796 &self,
1797 sessions: &mut SessionManager,
1798 request: CanonicalActorSessionRequest<'_>,
1799 ) -> Result<(Session, bool)> {
1800 let CanonicalActorSessionRequest {
1801 tentative_session,
1802 tentative_owns_session,
1803 entry,
1804 probe,
1805 attach,
1806 } = request;
1807 let Some(canonical_session_id) = entry.heddle_session_id.as_deref() else {
1808 return Ok((tentative_session, tentative_owns_session));
1809 };
1810 if canonical_session_id == tentative_session.id {
1811 return Ok((tentative_session, tentative_owns_session));
1812 }
1813
1814 if tentative_owns_session
1815 && let Ok(Some(session)) = sessions.get_session(&tentative_session.id)
1816 && session.is_active()
1817 {
1818 let _ = sessions.end_session(Some(&tentative_session.id));
1819 }
1820
1821 let canonical_session = sessions
1822 .get_session(canonical_session_id)?
1823 .ok_or_else(|| anyhow!("session not found: {canonical_session_id}"))?;
1824 let canonical_segment_id = canonical_session
1825 .current_segment_id
1826 .clone()
1827 .unwrap_or_default();
1828 sessions.set_current_session(canonical_session_id, &canonical_segment_id)?;
1829
1830 if let Some(native_actor_key) = probe
1831 .native_actor_key
1832 .as_deref()
1833 .or(entry.native_actor_key.as_deref())
1834 {
1835 attach.precedence.push(format!(
1836 "post-create-native-actor-key:{native_actor_key}:matched"
1837 ));
1838 attach.attach_reason = format!(
1839 "reused existing native actor {} on Heddle session {}",
1840 native_actor_key, canonical_session_id
1841 );
1842 attach.winning_rule = "native-actor-key-post-create".to_string();
1843 }
1844
1845 Ok((canonical_session, false))
1846 }
1847
1848 fn sync_registry_from_report(
1849 &self,
1850 report: &SessionReportEnvelope,
1851 status: ActorPresenceStatus,
1852 ) -> Result<()> {
1853 let registry = ActorPresenceStore::new(self.repo.heddle_dir());
1854 let entry = if let Some(agent_session_id) = &report.agent_session_id {
1855 registry.update_entry(agent_session_id, |entry| {
1856 if report.client_instance_id.is_some() {
1857 entry.client_instance_id = report.client_instance_id.clone();
1858 }
1859 if report.native_actor_key.is_some() {
1860 entry.native_actor_key = report.native_actor_key.clone();
1861 }
1862 if report.native_parent_actor_key.is_some() {
1863 entry.native_parent_actor_key = report.native_parent_actor_key.clone();
1864 }
1865 if report.native_instance_key.is_some() {
1866 entry.native_instance_key = report.native_instance_key.clone();
1867 }
1868 entry.heddle_session_id = Some(report.heddle_session_id.clone());
1869 entry.path = Some(self.repo.root().to_path_buf());
1870 entry.harness = report.harness.harness.clone();
1871 entry.provider = report.harness.provider.clone();
1872 entry.model = report.harness.model.clone();
1873 entry.thinking_level = report.harness.thinking_level.clone();
1874 entry.usage_summary = usage_to_summary(&report.usage);
1875 entry.last_progress_at =
1876 report.last_progress_at.as_deref().and_then(parse_timestamp);
1877 entry.report_flush_state = report.report_flush_state.clone();
1878 entry.attach_reason = report.attach_reason.clone();
1879 entry.attach_precedence = report.attach_precedence.clone();
1880 entry.winning_attach_rule = report.winning_attach_rule.clone();
1881 entry.probe_source = report.probe_source.clone();
1882 entry.probe_confidence = report.probe_confidence;
1883 entry.status = status.clone();
1884 entry.completed_at = match status {
1885 ActorPresenceStatus::Active => None,
1886 ActorPresenceStatus::Abandoned
1887 | ActorPresenceStatus::Complete
1888 | ActorPresenceStatus::Merged => Some(Utc::now()),
1889 };
1890 })?
1891 } else {
1892 None
1893 };
1894
1895 if entry.is_none() {
1896 let resolved = self.ensure_registry_entry(RegistryEntryRequest {
1897 heddle_session_id: &report.heddle_session_id,
1898 thread_name: report.thread.as_deref(),
1899 thread_id: report.thread_id.as_deref(),
1900 identity: &ResolvedIdentity {
1901 harness: report.harness.harness.clone(),
1902 provider: report.harness.provider.clone(),
1903 model: report.harness.model.clone(),
1904 thinking_level: report.harness.thinking_level.clone(),
1905 policy: report.harness.policy.clone(),
1906 },
1907 probe: &HarnessProbeResult {
1908 native_actor_key: report.native_actor_key.clone(),
1909 native_parent_actor_key: report.native_parent_actor_key.clone(),
1910 native_instance_key: report.native_instance_key.clone(),
1911 probe_source: report.probe_source.clone(),
1912 confidence: report.probe_confidence,
1913 ..HarnessProbeResult::default()
1914 },
1915 attach: &ResolvedAttachment {
1916 target: AttachTarget::CreateNew {
1917 _because_claimed: false,
1918 },
1919 matched_entry: None,
1920 attach_reason: report.attach_reason.clone().unwrap_or_else(|| {
1921 format!(
1922 "created actor for Heddle session {}",
1923 report.heddle_session_id
1924 )
1925 }),
1926 precedence: report.attach_precedence.clone(),
1927 winning_rule: report
1928 .winning_attach_rule
1929 .clone()
1930 .unwrap_or_else(|| "report-sync".to_string()),
1931 },
1932 client_instance_id: report.client_instance_id.as_deref(),
1933 requested_entry: None,
1934 })?;
1935 let mut report = report.clone();
1936 report.agent_session_id = Some(resolved.session_id);
1937 self.reports.save(&report)?;
1938 }
1939 Ok(())
1940 }
1941}
1942
1943#[derive(Debug, Clone, Default)]
1944struct ResolvedIdentity {
1945 harness: Option<String>,
1946 provider: Option<String>,
1947 model: Option<String>,
1948 thinking_level: Option<String>,
1949 policy: Option<String>,
1950}
1951
1952impl ResolvedIdentity {
1953 fn to_transport_identity(&self) -> HarnessIdentity {
1954 HarnessIdentity {
1955 harness: self.harness.clone(),
1956 provider: self.provider.clone(),
1957 model: self.model.clone(),
1958 thinking_level: self.thinking_level.clone(),
1959 policy: self.policy.clone(),
1960 }
1961 }
1962}
1963
1964struct IdentityHints {
1965 harness: Option<String>,
1966 provider: Option<String>,
1967 model: Option<String>,
1968 thinking_level: Option<String>,
1969 policy: Option<String>,
1970 probe: HarnessProbeResult,
1971}
1972
1973fn resolve_identity(
1974 repo: &Repository,
1975 user_config: &UserConfig,
1976 hints: IdentityHints,
1977) -> Result<ResolvedIdentity> {
1978 let current_session = SessionManager::new(repo.root()).get_current_session()?;
1979 let current_segment = current_session
1980 .as_ref()
1981 .and_then(|session| session.current_segment());
1982 let token_claims = if user_config.harness.auto_infer {
1983 active_token_claims()
1984 } else {
1985 None
1986 };
1987 let harness_override = resolved_harness_override(
1988 user_config,
1989 hints.harness.as_deref(),
1990 hints.probe.harness.as_deref(),
1991 );
1992
1993 Ok(ResolvedIdentity {
1994 harness: hints.harness.or(hints.probe.harness),
1995 provider: hints
1996 .provider
1997 .or(hints.probe.provider)
1998 .or_else(|| current_segment.map(|segment| segment.provider.clone()))
1999 .or_else(|| {
2000 token_claims
2001 .as_ref()
2002 .and_then(|claims| claims.agent_provider.clone())
2003 })
2004 .or_else(|| harness_override.and_then(|entry| entry.provider.clone()))
2005 .or_else(|| user_config.agent.provider.clone()),
2006 model: hints
2007 .model
2008 .or(hints.probe.model)
2009 .or_else(|| current_segment.map(|segment| segment.model.clone()))
2010 .or_else(|| {
2011 token_claims
2012 .as_ref()
2013 .and_then(|claims| claims.agent_model.clone())
2014 })
2015 .or_else(|| harness_override.and_then(|entry| entry.model.clone()))
2016 .or_else(|| user_config.agent.model.clone()),
2017 thinking_level: hints
2018 .thinking_level
2019 .or(hints.probe.thinking_level)
2020 .or_else(|| harness_override.and_then(|entry| entry.thinking_level.clone())),
2021 policy: hints
2022 .policy
2023 .or(hints.probe.policy)
2024 .or_else(|| current_segment.and_then(|segment| segment.policy_id.clone()))
2025 .or_else(|| harness_override.and_then(|entry| entry.policy.clone()))
2026 .or_else(|| user_config.agent.default_policy.clone()),
2027 })
2028}
2029
2030fn resolved_harness_override<'a>(
2031 user_config: &'a UserConfig,
2032 explicit: Option<&str>,
2033 fingerprint: Option<&str>,
2034) -> Option<&'a UserHarnessOverride> {
2035 explicit
2036 .and_then(|name| user_config.harness.harnesses.get(name))
2037 .or_else(|| fingerprint.and_then(|name| user_config.harness.harnesses.get(name)))
2038}
2039
2040enum AttachTarget {
2041 ExistingSession(objects::object::Session),
2042 CreateNew { _because_claimed: bool },
2043}
2044
2045struct ResolvedAttachment {
2046 target: AttachTarget,
2047 matched_entry: Option<ActorPresence>,
2048 attach_reason: String,
2049 precedence: Vec<String>,
2050 winning_rule: String,
2051}
2052
2053fn resolve_actor_attachment(
2054 registry: &ActorPresenceStore,
2055 repo: &Repository,
2056 sessions: &mut SessionManager,
2057 input: AttachmentResolutionInput<'_>,
2058) -> Result<ResolvedAttachment> {
2059 let AttachmentResolutionInput {
2060 requested_entry,
2061 explicit_heddle_session_id,
2062 client_instance_id,
2063 probe,
2064 token_claims,
2065 } = input;
2066
2067 let mut sessions_by_id: BTreeMap<String, Session> = BTreeMap::new();
2069 let mut matched_by_session: BTreeMap<String, ActorPresence> = BTreeMap::new();
2070 let mut facts = SessionAttachFacts {
2071 root_actor: probe.attach_hints.root_actor,
2072 ..SessionAttachFacts::default()
2073 };
2074
2075 if let Some(entry) = requested_entry
2076 && let Some(bound_session_id) = entry.heddle_session_id.as_deref()
2077 {
2078 let session = sessions
2079 .get_session(bound_session_id)?
2080 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2081 if !session.is_active() {
2082 return Err(anyhow!("session is not active: {bound_session_id}"));
2083 }
2084 matched_by_session.insert(session.id.clone(), entry.clone());
2085 sessions_by_id.insert(session.id.clone(), session);
2086 facts.explicit_agent = Some(ExplicitAgentBind {
2087 agent_session_id: entry.session_id.clone(),
2088 heddle_session_id: bound_session_id.to_string(),
2089 });
2090 }
2091
2092 if facts.explicit_agent.is_none()
2093 && let Some(session_id) = explicit_heddle_session_id
2094 {
2095 ensure_requested_entry_matches_session(requested_entry, session_id)?;
2096 let session = sessions
2097 .get_session(session_id)?
2098 .ok_or_else(|| anyhow!("session not found: {session_id}"))?;
2099 if !session.is_active() {
2100 return Err(anyhow!("session is not active: {session_id}"));
2101 }
2102 sessions_by_id.insert(session.id.clone(), session);
2103 facts.explicit_heddle_session_id = Some(session_id.to_string());
2104 }
2105
2106 if client_instance_id.is_none()
2107 && let Some(native_actor_key) = probe.native_actor_key.as_deref()
2108 {
2109 if let Some(entry) = registry.find_active_by_native_actor_key(native_actor_key)?
2110 && claude_actor_compatible(&entry, probe, repo.root())
2111 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2112 {
2113 let session = sessions
2114 .get_session(&bound_session_id)?
2115 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2116 if session.is_active() {
2117 matched_by_session.insert(session.id.clone(), entry);
2118 sessions_by_id.insert(session.id.clone(), session);
2119 facts.native_actor = SessionLookupFact::Hit {
2120 key: native_actor_key.to_string(),
2121 session_id: bound_session_id,
2122 };
2123 } else {
2124 facts.native_actor = SessionLookupFact::Miss {
2125 key: native_actor_key.to_string(),
2126 };
2127 }
2128 } else {
2129 facts.native_actor = SessionLookupFact::Miss {
2130 key: native_actor_key.to_string(),
2131 };
2132 }
2133 }
2134
2135 if let Some(client_instance_id) = client_instance_id {
2136 if let Some(entry) = registry.find_active_by_client_instance_id(client_instance_id)?
2137 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2138 {
2139 let session = sessions
2140 .get_session(&bound_session_id)?
2141 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2142 if session.is_active() {
2143 matched_by_session.insert(session.id.clone(), entry);
2144 sessions_by_id.insert(session.id.clone(), session);
2145 facts.client_instance = SessionLookupFact::Hit {
2146 key: client_instance_id.to_string(),
2147 session_id: bound_session_id,
2148 };
2149 } else {
2150 facts.client_instance = SessionLookupFact::Miss {
2151 key: client_instance_id.to_string(),
2152 };
2153 }
2154 } else {
2155 facts.client_instance = SessionLookupFact::Miss {
2156 key: client_instance_id.to_string(),
2157 };
2158 }
2159 }
2160
2161 if let Some(native_instance_key) = probe.native_instance_key.as_deref() {
2162 if let Some(entry) =
2163 registry.find_active_by_native_instance_key_at_path(native_instance_key, repo.root())?
2164 && claude_actor_compatible(&entry, probe, repo.root())
2165 && let Some(bound_session_id) = entry.heddle_session_id.clone()
2166 {
2167 let session = sessions
2168 .get_session(&bound_session_id)?
2169 .ok_or_else(|| anyhow!("session not found: {bound_session_id}"))?;
2170 if session.is_active() {
2171 matched_by_session.insert(session.id.clone(), entry);
2172 sessions_by_id.insert(session.id.clone(), session);
2173 facts.native_instance = SessionLookupFact::Hit {
2174 key: native_instance_key.to_string(),
2175 session_id: bound_session_id,
2176 };
2177 } else {
2178 facts.native_instance = SessionLookupFact::Miss {
2179 key: native_instance_key.to_string(),
2180 };
2181 }
2182 } else {
2183 facts.native_instance = SessionLookupFact::Miss {
2184 key: native_instance_key.to_string(),
2185 };
2186 }
2187 }
2188
2189 if probe.attach_hints.root_actor
2190 && let Some(current) = sessions.get_current_session()?
2191 && current.is_active()
2192 {
2193 let claimed = session_claimed_by_other(
2194 registry,
2195 ¤t.id,
2196 requested_entry,
2197 client_instance_id,
2198 probe.native_actor_key.as_deref(),
2199 )?;
2200 sessions_by_id
2201 .entry(current.id.clone())
2202 .or_insert_with(|| current.clone());
2203 facts.current_worktree = if claimed {
2204 WorktreeSessionFact::Claimed {
2205 session_id: current.id.clone(),
2206 }
2207 } else {
2208 WorktreeSessionFact::Available {
2209 session_id: current.id.clone(),
2210 }
2211 };
2212 }
2213
2214 if let Some(claims) = token_claims
2215 && let Some(token_sid) = claims.sid.as_deref()
2216 && let Some(session) = sessions.get_session(token_sid)?
2217 && session.is_active()
2218 {
2219 let claimed = session_claimed_by_other(
2220 registry,
2221 &session.id,
2222 requested_entry,
2223 client_instance_id,
2224 probe.native_actor_key.as_deref(),
2225 )?;
2226 let session_id = session.id.clone();
2227 sessions_by_id.insert(session_id.clone(), session);
2228 facts.token_sid = if claimed {
2229 TokenSidFact::Claimed { session_id }
2230 } else {
2231 TokenSidFact::Available { session_id }
2232 };
2233 }
2234
2235 let decision = decide_session_attach(&facts);
2236 match decision.policy {
2237 SessionPolicy::AttachExisting { session_id, .. } => {
2238 let session = sessions_by_id
2239 .remove(&session_id)
2240 .or_else(|| sessions.get_session(&session_id).ok().flatten())
2241 .ok_or_else(|| anyhow!("session not found: {session_id}"))?;
2242 Ok(ResolvedAttachment {
2243 target: AttachTarget::ExistingSession(session),
2244 matched_entry: matched_by_session.remove(&session_id),
2245 attach_reason: decision.attach_reason,
2246 precedence: decision.precedence,
2247 winning_rule: decision.winning_rule.to_string(),
2248 })
2249 }
2250 SessionPolicy::CreateNew {
2251 because_claimed, ..
2252 } => Ok(ResolvedAttachment {
2253 target: AttachTarget::CreateNew {
2254 _because_claimed: because_claimed,
2255 },
2256 matched_entry: None,
2257 attach_reason: decision.attach_reason,
2258 precedence: decision.precedence,
2259 winning_rule: decision.winning_rule.to_string(),
2260 }),
2261 }
2262}
2263
2264fn claude_actor_compatible(
2265 entry: &ActorPresence,
2266 probe: &HarnessProbeResult,
2267 repo_root: &Path,
2268) -> bool {
2269 let Some(native_actor_key) = probe.native_actor_key.as_deref() else {
2270 return true;
2271 };
2272 if !native_actor_key.starts_with("claude-code:") {
2273 return true;
2274 }
2275 if native_actor_key.starts_with("claude-code:agent:") {
2276 return entry.native_actor_key.as_deref() == Some(native_actor_key);
2277 }
2278 if let Some(native_instance_key) = probe.native_instance_key.as_deref() {
2279 return entry.native_actor_key.as_deref() == Some(native_actor_key)
2280 && entry.native_instance_key.as_deref() == Some(native_instance_key);
2281 }
2282 let same_repo = entry
2283 .path
2284 .as_ref()
2285 .map(|path| path.canonicalize().unwrap_or_else(|_| path.clone()))
2286 .unwrap_or_default()
2287 == repo_root
2288 .canonicalize()
2289 .unwrap_or_else(|_| repo_root.to_path_buf());
2290 entry.native_actor_key.as_deref() == Some(native_actor_key)
2291 && same_repo
2292 && probe.confidence.unwrap_or_default() >= 0.9
2293}
2294
2295fn decode_token_claims(token: &str) -> Option<TokenClaims> {
2296 let payload = token.split('.').nth(1)?;
2297 let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
2298 .decode(payload.as_bytes())
2299 .ok()?;
2300 serde_json::from_slice(&decoded).ok()
2301}
2302
2303fn active_token_claims() -> Option<TokenClaims> {
2308 #[cfg(feature = "client")]
2309 {
2310 crate::client::resolve_active_bearer()
2311 .ok()
2312 .flatten()
2313 .and_then(|token| decode_token_claims(&token.id))
2314 }
2315 #[cfg(not(feature = "client"))]
2316 {
2317 None
2318 }
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}