Skip to main content

pond/adapter/
pi_coding_agent.rs

1//! pi-coding-agent adapter (github.com/badlogic/pi-mono).
2//!
3//! Source path: `~/.pi/agent/sessions/<project-slug>/<ISO-ts>_<uuid>.jsonl`.
4//! One `.jsonl` file per session; each line is a typed record linked via a
5//! `parentId` -> `id` chain (pi-coding-agent's leaf-cursor DAG). Top-level types:
6//! `session` (consumed up front for Session), `model_change` /
7//! `thinking_level_change` / `compaction` (session-state carriers kept as
8//! System messages), and `message` (the per-turn model interaction, with
9//! roles user / assistant / toolResult and content nested under `.message`).
10//!
11//! v1 stores the linear log ordered by source line; pi-coding-agent's `parentId` fork
12//! graph (spec.md#deferred: multi-level fork lineage) is not collapsed into
13//! `parent_message_id` but preserved verbatim in `options.source.raw_record`
14//! for a future branching consumer.
15
16use std::path::{Path, PathBuf};
17
18use chrono::{DateTime, SecondsFormat, Utc};
19use serde_json::{Value, json};
20
21use crate::{
22    sessions::IngestEvent,
23    wire::{Message, Part, PartKind, Provenance, ProviderOptions, Session},
24};
25
26use super::{
27    Adapter, AdapterError, AdapterFactory, AdapterYieldStream, DiscoverFuture, Env,
28    RestoreFidelity, RestoredFile, SkipOracle, by_timestamp_then_id, compact_json, config_path,
29    empty_options,
30    extract::{Extracted, extract_compact_repr, extract_raw_record, extract_str},
31    extracted_text,
32    jsonl::{
33        BoundedRow, JsonlTree, jsonl_tree_discover, jsonl_tree_events, peek_last_line, source_line,
34    },
35    jsonl_bytes, part_id, part_ordinal, raw_record,
36};
37
38const NAME: &str = "pi-coding-agent";
39
40/// Stateless factory: opens [`PiCodingAgentAdapter`] instances and probes for the
41/// canonical install location under `~/.pi/agent/sessions`.
42pub struct PiCodingAgentFactory;
43
44impl AdapterFactory for PiCodingAgentFactory {
45    fn name(&self) -> &'static str {
46        NAME
47    }
48
49    fn open(&self, config: Value) -> Result<Box<dyn Adapter>, AdapterError> {
50        Ok(Box::new(PiCodingAgentAdapter::new(config_path(
51            NAME, config,
52        )?)))
53    }
54
55    fn probe_default(&self, env: &Env) -> Option<Value> {
56        let path = env.home.join(".pi").join("agent").join("sessions");
57        path.exists().then(|| json!({ "path": path }))
58    }
59
60    fn serialize(
61        &self,
62        session: &crate::sessions::SessionWithMessages,
63        fidelity: RestoreFidelity,
64    ) -> Result<Vec<RestoredFile>, AdapterError> {
65        serialize_session(session, fidelity)
66    }
67}
68
69fn serialize_session(
70    session: &crate::sessions::SessionWithMessages,
71    fidelity: RestoreFidelity,
72) -> Result<Vec<RestoredFile>, AdapterError> {
73    // Native replays the verbatim `options.source.raw_record` rows (the
74    // `session` line first, then one per message in source order); `pi_record`
75    // below is the foreign-only reconstruction. Replay echoes a frozen
76    // snapshot - safe only while canonical is append-only
77    // (spec.md#adapter-integrity-additive-sync).
78    //
79    // spec.md#adapter-native-restore-lossless: if Native is requested but the
80    // session has no stored `raw_record`, we downgrade to foreign and stamp
81    // `actual_fidelity` so the caller can signal the downgrade. Mirrors
82    // opencode's behavior - both adapters serve the best they can and tell
83    // the truth about what they served.
84    let session_raw = raw_record(&session.session.options);
85    let actual = match fidelity {
86        RestoreFidelity::Native if session_raw.is_some() => RestoreFidelity::Native,
87        _ => RestoreFidelity::Foreign,
88    };
89
90    let mut records = Vec::new();
91    if actual == RestoreFidelity::Native {
92        records.push(session_raw.unwrap_or_else(|| pi_session_record(session)));
93    } else {
94        records.push(pi_session_record(session));
95    }
96
97    // Sort message references rather than cloning the whole vec; restore is a
98    // hot path when users `pond restore` large sessions.
99    let mut messages: Vec<&crate::sessions::MessageWithParts> = session.messages.iter().collect();
100    if actual == RestoreFidelity::Native {
101        messages.sort_by(|left, right| {
102            source_line(left.message.options())
103                .cmp(&source_line(right.message.options()))
104                .then_with(|| by_timestamp_then_id(left, right))
105        });
106    } else {
107        messages.sort_by(|left, right| by_timestamp_then_id(left, right));
108    }
109
110    for message in &messages {
111        if actual == RestoreFidelity::Native
112            && let Some(raw) = raw_record(message.message.options())
113        {
114            records.push(raw);
115            continue;
116        }
117        // Foreign restore: a System carrier (model/thinking/compaction record)
118        // has no idiomatic home in another client's transcript - drop it; the
119        // content stays in canonical (spec.md#adapter-native-restore-lossless,
120        // foreign clause).
121        if matches!(message.message, Message::System { .. }) {
122            continue;
123        }
124        records.push(pi_message_record(message));
125    }
126
127    Ok(vec![RestoredFile::new(
128        pi_relative_path(session),
129        jsonl_bytes(NAME, &records)?,
130        actual,
131    )])
132}
133
134/// Reproduce the on-disk `sessions/<slug>/<file>.jsonl` path from the slug and
135/// file name captured at ingest. Falls back to a derived name when a foreign
136/// session never carried them.
137fn pi_relative_path(session: &crate::sessions::SessionWithMessages) -> PathBuf {
138    let source = session.session.options.get("source");
139    let slug = source
140        .and_then(|s| s.get("project_slug"))
141        .and_then(Value::as_str)
142        .map(ToOwned::to_owned)
143        .unwrap_or_else(|| encode_project(&session.session.project));
144    let file_name = source
145        .and_then(|s| s.get("file_name"))
146        .and_then(Value::as_str)
147        .map(ToOwned::to_owned)
148        .unwrap_or_else(|| {
149            let ts = session.session.created_at.format("%Y-%m-%dT%H-%M-%S-%3fZ");
150            format!("{ts}_{}.jsonl", session.session.id)
151        });
152    PathBuf::from("sessions").join(slug).join(file_name)
153}
154
155fn encode_project(project: &str) -> String {
156    project.replace(['/', '.'], "-")
157}
158
159fn pi_session_record(session: &crate::sessions::SessionWithMessages) -> Value {
160    json!({
161        "type": "session",
162        "version": 3,
163        "id": session.session.id,
164        "timestamp": session.session.created_at.to_rfc3339_opts(SecondsFormat::Millis, true),
165        "cwd": &*session.session.project,
166    })
167}
168
169fn pi_message_record(message: &crate::sessions::MessageWithParts) -> Value {
170    json!({
171        "type": "message",
172        "id": message.message.id(),
173        "parentId": message.message.options().get("source").and_then(|s| s.get("parent_id")),
174        "timestamp": message.message.timestamp().to_rfc3339_opts(SecondsFormat::Millis, true),
175        "message": pi_inner_message(message),
176    })
177}
178
179fn pi_inner_message(message: &crate::sessions::MessageWithParts) -> Value {
180    let epoch_ms = message.message.timestamp().timestamp_millis();
181    match &message.message {
182        Message::User { .. } => json!({
183            "role": "user",
184            "content": message.parts.iter().map(pi_content_item).collect::<Vec<_>>(),
185            "timestamp": epoch_ms,
186        }),
187        Message::Assistant { .. } => json!({
188            "role": "assistant",
189            "content": message.parts.iter().map(pi_content_item).collect::<Vec<_>>(),
190            "timestamp": epoch_ms,
191        }),
192        Message::Tool { .. } => {
193            // spec.md#adapter-native-restore-lossless (foreign clause): a
194            // canonical Tool message with no ToolResult part - or with parts
195            // that lack call_id/name - serializes with empty-string slots.
196            // That's lossy by design for foreign restore; the unaltered
197            // source still lives in canonical and in `raw_record`.
198            let part = message.parts.first();
199            let (call_id, name, is_error, result) = match part.map(|p| &p.kind) {
200                Some(PartKind::ToolResult {
201                    call_id,
202                    name,
203                    is_failure,
204                    result,
205                }) => (
206                    extracted_text(call_id).to_owned(),
207                    extracted_text(name).to_owned(),
208                    *is_failure,
209                    result.clone(),
210                ),
211                _ => (String::new(), String::new(), false, Value::Null),
212            };
213            json!({
214                "role": "toolResult",
215                "toolCallId": call_id,
216                "toolName": name,
217                "content": result,
218                "isError": is_error,
219                "timestamp": epoch_ms,
220            })
221        }
222        // serialize_session drops System carriers before reaching here in
223        // foreign mode, and native mode replays the source row verbatim - so
224        // this arm only fires if a caller invokes pi_message_record on a
225        // System message directly. Unreachable on every legitimate path.
226        Message::System { .. } => {
227            unreachable!("System messages are not serialized through pi_inner_message")
228        }
229    }
230}
231
232fn pi_content_item(part: &Part) -> Value {
233    match &part.kind {
234        PartKind::Text { text } => json!({"type": "text", "text": extracted_text(text)}),
235        PartKind::Reasoning { text } => json!({
236            "type": "thinking",
237            "thinking": extracted_text(text),
238            "thinkingSignature": part
239                .options
240                .get("pi")
241                .and_then(|p| p.get("thinking_signature")),
242        }),
243        PartKind::ToolCall {
244            call_id,
245            name,
246            params,
247            ..
248        } => json!({
249            "type": "toolCall",
250            "id": extracted_text(call_id),
251            "name": extracted_text(name),
252            "arguments": params,
253        }),
254        other => json!({
255            "type": "text",
256            "text": compact_json(&serde_json::to_value(other).unwrap_or(Value::Null)),
257        }),
258    }
259}
260
261/// Configured pi-coding-agent reader. Walks a tree of `*.jsonl` files under [`Self::root`]
262/// and yields canonical events in source order per session.
263#[derive(Debug, Clone)]
264pub struct PiCodingAgentAdapter {
265    root: PathBuf,
266}
267
268impl PiCodingAgentAdapter {
269    pub fn new(root: impl Into<PathBuf>) -> Self {
270        Self { root: root.into() }
271    }
272}
273
274impl Adapter for PiCodingAgentAdapter {
275    fn discover(&self) -> DiscoverFuture<'_> {
276        jsonl_tree_discover(self)
277    }
278
279    fn events_with<'a>(&'a self, oracle: &'a dyn SkipOracle) -> AdapterYieldStream<'a> {
280        jsonl_tree_events(self, oracle)
281    }
282
283    fn plan<'a>(&'a self, oracle: &'a dyn SkipOracle) -> crate::adapter::PlanFuture<'a> {
284        crate::adapter::jsonl::jsonl_tree_plan(self, oracle)
285    }
286}
287
288impl JsonlTree for PiCodingAgentAdapter {
289    // pi-coding-agent's `toolResult` records carry their own `toolName`, so unlike
290    // claude-code / codex-cli the adapter needs no per-file call_id -> name map.
291    type State = ();
292
293    fn name(&self) -> &'static str {
294        NAME
295    }
296
297    fn root(&self) -> &Path {
298        &self.root
299    }
300
301    fn peek_session_id(&self, _path: &Path, first_line: &str) -> Option<String> {
302        let row: Value = serde_json::from_str(first_line).ok()?;
303        if row.get("type").and_then(Value::as_str) == Some("session") {
304            row.get("id").and_then(Value::as_str).map(ToOwned::to_owned)
305        } else {
306            None
307        }
308    }
309
310    fn peek_watermark(&self, path: &Path) -> crate::adapter::SourceWatermark {
311        // The `session` header is eventless; every other row is a message
312        // carrying its own `timestamp`. The transcript is append-ordered, so the
313        // last line is the latest message; a `session`-only or timestamp-less
314        // tail is opaque and the file re-reads (safe).
315        let last_ts = || -> Option<i64> {
316            let row: Value = serde_json::from_str(&peek_last_line(path)?).ok()?;
317            if row.get("type").and_then(Value::as_str) == Some("session") {
318                return None;
319            }
320            let text = row.get("timestamp").and_then(Value::as_str)?;
321            Some(
322                DateTime::parse_from_rfc3339(text)
323                    .ok()?
324                    .with_timezone(&Utc)
325                    .timestamp_micros(),
326            )
327        };
328        match last_ts() {
329            Some(ts) => crate::adapter::SourceWatermark::At(ts),
330            None => crate::adapter::SourceWatermark::Opaque,
331        }
332    }
333
334    fn session(&self, path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
335        session_from_rows(path, rows)
336    }
337
338    fn events_from_row(
339        &self,
340        session: &Session,
341        row: &BoundedRow,
342        _state: &mut Self::State,
343    ) -> Result<Vec<IngestEvent>, String> {
344        events_from_row(&session.id, row.line, &row.value, session.created_at)
345    }
346}
347
348fn session_from_rows(path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
349    let path_display = path.display().to_string();
350    let first = rows
351        .first()
352        .ok_or_else(|| AdapterError::schema(NAME, path_display.clone(), "empty jsonl session"))?;
353    let row = &first.value;
354    let at_first = format!("{path_display}:{}", first.line);
355    if row.get("type").and_then(Value::as_str) != Some("session") {
356        return Err(AdapterError::schema(
357            NAME,
358            at_first,
359            "first row must be a `session` record",
360        ));
361    }
362    let id = row
363        .get("id")
364        .and_then(Value::as_str)
365        .ok_or_else(|| AdapterError::schema(NAME, at_first.clone(), "session record missing id"))?
366        .to_owned();
367    let created_at = row
368        .get("timestamp")
369        .and_then(Value::as_str)
370        .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
371        .map(|dt| dt.with_timezone(&Utc))
372        .ok_or_else(|| {
373            AdapterError::schema(
374                NAME,
375                at_first.clone(),
376                "session record has no parseable timestamp",
377            )
378        })?;
379    let project = extract_str(row, "cwd").ok_or_else(|| {
380        // spec.md#model-project-non-empty: pi always records `cwd` on the
381        // session line; its absence is a malformed session, not a default.
382        AdapterError::schema(NAME, at_first, "session record missing cwd")
383    })?;
384
385    // Capture the exact on-disk path components so native restore reproduces
386    // the source file byte-for-byte (the slug encoding and filename timestamp
387    // are not recomputable from canonical alone).
388    let project_slug = path
389        .parent()
390        .and_then(|p| p.file_name())
391        .and_then(|n| n.to_str())
392        .map(ToOwned::to_owned);
393    let file_name = path
394        .file_name()
395        .and_then(|n| n.to_str())
396        .map(ToOwned::to_owned);
397
398    let mut options = ProviderOptions::new();
399    options.insert(
400        "source".to_owned(),
401        json!({
402            "adapter": NAME,
403            "version": row.get("version"),
404            "project_slug": project_slug,
405            "file_name": file_name,
406            "raw_record": extract_raw_record(row),
407        }),
408    );
409
410    Ok(Session {
411        id,
412        parent_session_id: None,
413        parent_message_id: None,
414        source_agent: NAME.to_owned(),
415        created_at,
416        project,
417        options,
418    })
419}
420
421/// Map one pi JSONL record into zero-or-more `IngestEvent`s. `session` is
422/// consumed up front (eventless here); `model_change` / `thinking_level_change`
423/// / `compaction` become System carriers; `message` becomes a User / Assistant
424/// / Tool message plus its content Parts.
425fn events_from_row(
426    session_id: &str,
427    line: usize,
428    row: &Value,
429    default_timestamp: DateTime<Utc>,
430) -> Result<Vec<IngestEvent>, String> {
431    let kind = row.get("type").and_then(Value::as_str);
432    let timestamp = row
433        .get("timestamp")
434        .and_then(Value::as_str)
435        .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
436        .map(|dt| dt.with_timezone(&Utc))
437        .unwrap_or(default_timestamp);
438    let id = row
439        .get("id")
440        .and_then(Value::as_str)
441        .map_or_else(|| format!("{session_id}:{line}"), ToOwned::to_owned);
442
443    match kind {
444        // Consumed up front by `session_from_rows`.
445        Some("session") => Ok(Vec::new()),
446        Some("message") => {
447            let message_value = row
448                .get("message")
449                .ok_or_else(|| "message record missing `message` field".to_owned())?;
450            message_events(session_id, &id, timestamp, row, message_value, line)
451        }
452        // Session-state carriers: keep the human-meaningful field as content,
453        // the rest of the record in `options.source.raw_record`.
454        Some("compaction") => Ok(vec![carrier_event(
455            session_id,
456            &id,
457            timestamp,
458            row,
459            line,
460            extract_str(row, "summary"),
461        )]),
462        Some("model_change") | Some("thinking_level_change") => Ok(vec![carrier_event(
463            session_id,
464            &id,
465            timestamp,
466            row,
467            line,
468            extract_str(row, "type"),
469        )]),
470        // Unknown record type: preserve it as a System carrier rather than
471        // dropping (spec.md#adapter-integrity-no-silent-drops). The raw record
472        // survives in options; the type label is the content.
473        _ => Ok(vec![carrier_event(
474            session_id,
475            &id,
476            timestamp,
477            row,
478            line,
479            extract_str(row, "type"),
480        )]),
481    }
482}
483
484fn carrier_event(
485    session_id: &str,
486    id: &str,
487    timestamp: DateTime<Utc>,
488    row: &Value,
489    line: usize,
490    content: Option<Extracted<String>>,
491) -> IngestEvent {
492    IngestEvent::Message(Message::System {
493        id: id.to_owned(),
494        session_id: session_id.to_owned(),
495        timestamp,
496        content,
497        options: row_options(row, line),
498    })
499}
500
501fn message_events(
502    session_id: &str,
503    id: &str,
504    timestamp: DateTime<Utc>,
505    row: &Value,
506    message_value: &Value,
507    line: usize,
508) -> Result<Vec<IngestEvent>, String> {
509    let role = message_value
510        .get("role")
511        .and_then(Value::as_str)
512        .ok_or_else(|| "message missing role".to_owned())?;
513    let content = message_value
514        .get("content")
515        .and_then(Value::as_array)
516        .cloned()
517        .unwrap_or_default();
518
519    let mut parts = Vec::new();
520    let message = match role {
521        "user" => {
522            // spec.md#model-part-provenance: pi user messages are genuine human
523            // prompts; harness-injected context arrives as separate records
524            // (compaction, model_change), not inside a user turn.
525            for (ordinal, item) in content.iter().enumerate() {
526                parts.push(user_part(session_id, id, ordinal, item));
527            }
528            Message::User {
529                id: id.to_owned(),
530                session_id: session_id.to_owned(),
531                timestamp,
532                options: row_options(row, line),
533            }
534        }
535        "assistant" => {
536            for (ordinal, item) in content.iter().enumerate() {
537                parts.push(assistant_part(session_id, id, ordinal, item));
538            }
539            Message::Assistant {
540                id: id.to_owned(),
541                session_id: session_id.to_owned(),
542                timestamp,
543                options: assistant_options(row, message_value, line),
544            }
545        }
546        "toolResult" => {
547            parts.push(tool_result_part(session_id, id, message_value));
548            Message::Tool {
549                id: id.to_owned(),
550                session_id: session_id.to_owned(),
551                timestamp,
552                options: row_options(row, line),
553            }
554        }
555        // Unknown nested roles are still parseable source records. Preserve the
556        // row as a System carrier instead of turning it into a counted drop.
557        _ => Message::System {
558            id: id.to_owned(),
559            session_id: session_id.to_owned(),
560            timestamp,
561            content: extract_str(message_value, "role"),
562            options: row_options(row, line),
563        },
564    };
565
566    let mut events = Vec::with_capacity(parts.len() + 1);
567    events.push(IngestEvent::Message(message));
568    events.extend(parts.into_iter().map(IngestEvent::Part));
569    Ok(events)
570}
571
572fn user_part(session_id: &str, message_id: &str, ordinal: usize, item: &Value) -> Part {
573    let kind = match item.get("type").and_then(Value::as_str) {
574        Some("text") => PartKind::Text {
575            text: extract_str(item, "text"),
576        },
577        // Anything else (e.g. an `image` content item) is preserved losslessly
578        // as a compact-JSON Text Part rather than dropped.
579        _ => PartKind::Text {
580            text: Some(extract_compact_repr(item)),
581        },
582    };
583    Part {
584        session_id: session_id.to_owned(),
585        id: part_id(message_id, ordinal),
586        message_id: message_id.to_owned(),
587        ordinal: part_ordinal(ordinal),
588        // spec.md#model-part-provenance: a genuine human prompt is conversation.
589        provenance: Provenance::Conversational,
590        options: empty_options(),
591        kind,
592    }
593}
594
595fn assistant_part(session_id: &str, message_id: &str, ordinal: usize, item: &Value) -> Part {
596    // spec.md#model-part-provenance: assistant text, reasoning, and tool calls
597    // are model-authored, hence conversational.
598    let (kind, options) = match item.get("type").and_then(Value::as_str) {
599        Some("text") => (
600            PartKind::Text {
601                text: extract_str(item, "text"),
602            },
603            empty_options(),
604        ),
605        Some("thinking") => (
606            PartKind::Reasoning {
607                text: extract_str(item, "thinking"),
608            },
609            thinking_options(item),
610        ),
611        Some("toolCall") => (
612            PartKind::ToolCall {
613                call_id: extract_str(item, "id"),
614                name: extract_str(item, "name"),
615                params: item.get("arguments").cloned().unwrap_or(Value::Null),
616                provider_executed: false,
617            },
618            empty_options(),
619        ),
620        // Lossless fallback for an unrecognised assistant content shape.
621        _ => (
622            PartKind::Text {
623                text: Some(extract_compact_repr(item)),
624            },
625            empty_options(),
626        ),
627    };
628    Part {
629        session_id: session_id.to_owned(),
630        id: part_id(message_id, ordinal),
631        message_id: message_id.to_owned(),
632        ordinal: part_ordinal(ordinal),
633        provenance: Provenance::Conversational,
634        options,
635        kind,
636    }
637}
638
639fn tool_result_part(session_id: &str, message_id: &str, message_value: &Value) -> Part {
640    Part {
641        session_id: session_id.to_owned(),
642        id: part_id(message_id, 0),
643        message_id: message_id.to_owned(),
644        ordinal: 0,
645        // spec.md#model-part-provenance: tool output is runtime-produced.
646        provenance: Provenance::Injected,
647        options: empty_options(),
648        kind: PartKind::ToolResult {
649            call_id: extract_str(message_value, "toolCallId"),
650            name: extract_str(message_value, "toolName"),
651            is_failure: message_value
652                .get("isError")
653                .and_then(Value::as_bool)
654                .unwrap_or(false),
655            // The whole `content` array (text and/or image items) is the
656            // faithful result payload.
657            result: message_value.get("content").cloned().unwrap_or(Value::Null),
658        },
659    }
660}
661
662fn row_options(row: &Value, line: usize) -> ProviderOptions {
663    let mut options = ProviderOptions::new();
664    options.insert(
665        "source".to_owned(),
666        json!({
667            "line": line,
668            "parent_id": row.get("parentId"),
669            "raw_type": row.get("type"),
670            "raw_record": extract_raw_record(row),
671        }),
672    );
673    options
674}
675
676fn assistant_options(row: &Value, message_value: &Value, line: usize) -> ProviderOptions {
677    let mut options = row_options(row, line);
678    options.insert(
679        "pi".to_owned(),
680        json!({
681            "api": message_value.get("api"),
682            "provider": message_value.get("provider"),
683            "model": message_value.get("model"),
684            "usage": message_value.get("usage"),
685            "stop_reason": message_value.get("stopReason"),
686            "response_id": message_value.get("responseId"),
687        }),
688    );
689    options
690}
691
692fn thinking_options(item: &Value) -> ProviderOptions {
693    let mut options = ProviderOptions::new();
694    if let Some(signature) = item.get("thinkingSignature") {
695        options.insert("pi".to_owned(), json!({ "thinking_signature": signature }));
696    }
697    options
698}
699
700#[cfg(test)]
701mod tests {
702    //! End-to-end test for the pi-coding-agent adapter: ingest the committed fixture corpus
703    //! and assert pond's canonical Session/Message/Part shape comes out the
704    //! other side. The fixture lives under `tests/fixtures/adapter/pi-coding-agent/`.
705    #![allow(clippy::expect_used, clippy::unwrap_used)]
706
707    use super::*;
708    use crate::{handlers::ingest_adapter, sessions::Store, wire::PartKind};
709    use tempfile::TempDir;
710
711    // Manifest-dir anchored: unit tests must not depend on the process cwd
712    // (figment::Jail chdirs the whole test process while config tests run).
713    const FIXTURES: &str = concat!(
714        env!("CARGO_MANIFEST_DIR"),
715        "/tests/fixtures/adapter/pi-coding-agent/sessions"
716    );
717
718    #[test]
719    fn probe_default_finds_pi_sessions_under_home() -> anyhow::Result<()> {
720        crate::adapter::test_support::assert_probe_default(
721            &PiCodingAgentFactory,
722            &[".pi", "agent", "sessions"],
723        )
724    }
725
726    #[tokio::test(flavor = "multi_thread")]
727    async fn native_restore_is_value_equal_to_fixture_corpus() -> anyhow::Result<()> {
728        let adapter = PiCodingAgentAdapter::new(FIXTURES);
729        crate::adapter::test_support::assert_native_restore(
730            &PiCodingAgentFactory,
731            &adapter,
732            // pi-coding-agent relative paths embed the `sessions/` segment, so the corpus
733            // root is FIXTURES' parent, not FIXTURES itself.
734            std::path::Path::new(FIXTURES)
735                .parent()
736                .expect("FIXTURES is nested under a corpus root"),
737        )
738        .await
739    }
740
741    #[tokio::test(flavor = "multi_thread")]
742    async fn pi_coding_agent_adapter_ingests_fixture_corpus_into_canonical_shape()
743    -> anyhow::Result<()> {
744        let temp = TempDir::new()?;
745        let store = Store::open_local(temp.path()).await?;
746        let adapter = PiCodingAgentAdapter::new(FIXTURES);
747
748        let summary = ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
749        assert!(summary.accepted() > 0, "ingest must accept rows");
750        assert_eq!(summary.dropped_events, 0, "no per-event drops expected");
751        assert_eq!(
752            summary.dropped_sessions, 0,
753            "no session-level rejections expected"
754        );
755        assert_eq!(summary.skipped_files, 0, "no whole-file skips expected");
756
757        let (sessions, messages, parts) = store.row_counts().await?;
758        assert!(sessions > 0, "at least one pi-coding-agent session");
759        assert!(messages > 0, "at least one pi-coding-agent message");
760        assert!(parts > 0, "at least one pi-coding-agent Part");
761
762        let mut saw_tool_call = false;
763        let mut saw_tool_result = false;
764        let mut saw_reasoning = false;
765        for session_id in store.session_ids().await? {
766            let session = store
767                .get_session(&session_id)
768                .await?
769                .expect("session round-trips");
770            assert_eq!(session.session.source_agent, NAME);
771            assert!(
772                !(*session.session.project).is_empty(),
773                "spec.md#model-project-non-empty: project must be a real cwd",
774            );
775            for stored in &session.messages {
776                for part in &stored.parts {
777                    match &part.kind {
778                        PartKind::ToolCall { .. } => saw_tool_call = true,
779                        PartKind::ToolResult { .. } => saw_tool_result = true,
780                        PartKind::Reasoning { .. } => saw_reasoning = true,
781                        _ => {}
782                    }
783                }
784            }
785        }
786        assert!(saw_tool_call, "corpus has assistant tool calls");
787        assert!(saw_tool_result, "corpus has tool results");
788        assert!(saw_reasoning, "corpus has assistant reasoning");
789        Ok(())
790    }
791
792    #[test]
793    fn unknown_nested_message_role_becomes_system_carrier() -> anyhow::Result<()> {
794        let row = json!({
795            "type": "message",
796            "id": "mystery-message",
797            "message": {
798                "role": "mysteryRole",
799                "content": [{"type": "text", "text": "not yet understood"}]
800            }
801        });
802        let events = events_from_row(
803            "session-1",
804            42,
805            &row,
806            DateTime::parse_from_rfc3339("2026-04-28T18:47:32.280Z")?.with_timezone(&Utc),
807        )
808        .map_err(anyhow::Error::msg)?;
809
810        assert_eq!(events.len(), 1);
811        let IngestEvent::Message(Message::System {
812            id,
813            content,
814            options,
815            ..
816        }) = &events[0]
817        else {
818            panic!("unknown role must produce a System carrier");
819        };
820        assert_eq!(id, "mystery-message");
821        assert_eq!(content.as_deref().map(String::as_str), Some("mysteryRole"));
822        assert_eq!(
823            raw_record(options)
824                .and_then(|raw| raw.get("message").cloned())
825                .and_then(|message| message.get("role").cloned()),
826            Some(json!("mysteryRole")),
827        );
828        Ok(())
829    }
830
831    #[tokio::test(flavor = "multi_thread")]
832    async fn fork_parent_ids_and_compaction_summary_are_preserved() -> anyhow::Result<()> {
833        let temp = TempDir::new()?;
834        let root = temp.path().join("sessions");
835        let path = root
836            .join("project")
837            .join("2026-05-01T00-00-00-000Z_fork.jsonl");
838        write_jsonl_file(
839            &path,
840            &[
841                json!({
842                    "type": "session",
843                    "version": 3,
844                    "id": "pi-fork-session",
845                    "timestamp": "2026-05-01T00:00:00.000Z",
846                    "cwd": "/tmp/pi-fork",
847                }),
848                json!({
849                    "type": "message",
850                    "id": "parent-message",
851                    "timestamp": "2026-05-01T00:00:01.000Z",
852                    "message": {
853                        "role": "user",
854                        "content": [{"type": "text", "text": "parent"}],
855                    },
856                }),
857                json!({
858                    "type": "message",
859                    "id": "child-a",
860                    "parentId": "parent-message",
861                    "timestamp": "2026-05-01T00:00:02.000Z",
862                    "message": {
863                        "role": "assistant",
864                        "content": [{"type": "text", "text": "branch a"}],
865                    },
866                }),
867                json!({
868                    "type": "message",
869                    "id": "child-b",
870                    "parentId": "parent-message",
871                    "timestamp": "2026-05-01T00:00:03.000Z",
872                    "message": {
873                        "role": "assistant",
874                        "content": [{"type": "text", "text": "branch b"}],
875                    },
876                }),
877                json!({
878                    "type": "compaction",
879                    "id": "compact-1",
880                    "parentId": "child-b",
881                    "timestamp": "2026-05-01T00:00:04.000Z",
882                    "summary": "compact summary",
883                }),
884            ],
885        )?;
886
887        let store = Store::open_local(temp.path().join("store")).await?;
888        let summary = ingest_adapter(
889            &store,
890            &PiCodingAgentAdapter::new(&root),
891            &crate::adapter::NoopOracle,
892            |_| {},
893        )
894        .await?;
895        assert_eq!(summary.dropped_events, 0);
896
897        let session = store
898            .get_session("pi-fork-session")
899            .await?
900            .expect("fixture session lands");
901        let child_a = session
902            .messages
903            .iter()
904            .find(|stored| stored.message.id() == "child-a")
905            .expect("first fork child lands");
906        let child_b = session
907            .messages
908            .iter()
909            .find(|stored| stored.message.id() == "child-b")
910            .expect("second fork child lands");
911        for child in [child_a, child_b] {
912            assert_eq!(
913                child
914                    .message
915                    .options()
916                    .get("source")
917                    .and_then(|source| source.get("parent_id"))
918                    .and_then(Value::as_str),
919                Some("parent-message"),
920            );
921        }
922        assert!(source_line(child_a.message.options()) < source_line(child_b.message.options()));
923
924        let compact = session
925            .messages
926            .iter()
927            .find(|stored| stored.message.id() == "compact-1")
928            .expect("compaction carrier lands");
929        let Message::System { content, .. } = &compact.message else {
930            panic!("compaction is preserved as a System carrier");
931        };
932        assert_eq!(
933            content.as_deref().map(String::as_str),
934            Some("compact summary")
935        );
936        Ok(())
937    }
938
939    #[tokio::test(flavor = "multi_thread")]
940    async fn foreign_serialization_reparses_as_pi_coding_agent() -> anyhow::Result<()> {
941        let temp = TempDir::new()?;
942        let origin_store = Store::open_local(temp.path().join("origin-store")).await?;
943        let origin = crate::adapter::OpencodeAdapter::new(concat!(
944            env!("CARGO_MANIFEST_DIR"),
945            "/tests/fixtures/adapter/opencode/storage"
946        ));
947        ingest_adapter(&origin_store, &origin, &crate::adapter::NoopOracle, |_| {}).await?;
948        let session_id = origin_store
949            .session_ids()
950            .await?
951            .into_iter()
952            .next()
953            .expect("opencode fixture has sessions");
954        let session = origin_store
955            .get_session(&session_id)
956            .await?
957            .expect("fixture session is readable");
958
959        let restored_root = temp.path().join("pi-corpus");
960        crate::adapter::write_restored_files(
961            &restored_root,
962            PiCodingAgentFactory.serialize(&session, RestoreFidelity::Foreign)?,
963        )?;
964        let restored_store = Store::open_local(temp.path().join("restored-store")).await?;
965        let summary = ingest_adapter(
966            &restored_store,
967            &PiCodingAgentAdapter::new(restored_root.join("sessions")),
968            &crate::adapter::NoopOracle,
969            |_| {},
970        )
971        .await?;
972
973        assert!(summary.accepted() > 0);
974        assert_eq!(summary.dropped_events, 0);
975        Ok(())
976    }
977
978    /// spec.md#model-part-provenance: a tool result is harness-injected; an
979    /// assistant turn's text/reasoning/tool-call parts are conversation.
980    #[tokio::test(flavor = "multi_thread")]
981    async fn tool_results_are_injected_assistant_parts_are_conversational() -> anyhow::Result<()> {
982        let temp = TempDir::new()?;
983        let store = Store::open_local(temp.path()).await?;
984        let adapter = PiCodingAgentAdapter::new(FIXTURES);
985        ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
986
987        for session_id in store.session_ids().await? {
988            let session = store
989                .get_session(&session_id)
990                .await?
991                .expect("session round-trips");
992            for stored in &session.messages {
993                for part in &stored.parts {
994                    match &part.kind {
995                        PartKind::ToolResult { .. } => {
996                            assert_eq!(part.provenance, Provenance::Injected);
997                        }
998                        PartKind::ToolCall { .. } | PartKind::Reasoning { .. } => {
999                            assert_eq!(part.provenance, Provenance::Conversational);
1000                        }
1001                        _ => {}
1002                    }
1003                }
1004            }
1005        }
1006        Ok(())
1007    }
1008
1009    fn write_jsonl_file(path: &std::path::Path, records: &[Value]) -> anyhow::Result<()> {
1010        if let Some(parent) = path.parent() {
1011            std::fs::create_dir_all(parent)?;
1012        }
1013        std::fs::write(path, jsonl_bytes(NAME, records)?)?;
1014        Ok(())
1015    }
1016}