1use 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
40pub 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 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 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 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
134fn 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 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 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#[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 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 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 AdapterError::schema(NAME, at_first, "session record missing cwd")
383 })?;
384
385 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
421fn 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 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 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 _ => 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 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 _ => 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 _ => 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 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 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 _ => (
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 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 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 #![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 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 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 #[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}