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