Skip to main content

cli/harness/
mod.rs

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