use std::path::{Path, PathBuf};
use chrono::{DateTime, SecondsFormat, Utc};
use rusqlite::Connection;
use serde::Deserialize;
use serde_json::{Value, json};
use tokio::sync::mpsc;
use tokio_stream::StreamExt;
use crate::{
sessions::IngestEvent,
wire::{Message, Part, PartKind, Provenance, ProviderOptions, Session},
};
use super::{
Adapter, AdapterError, AdapterFactory, AdapterYield, AdapterYieldStream, DiscoverFuture, Env,
RestoreFidelity, RestoredFile, SkipOracle, SkipReason, SourceWatermark, SyncPlan,
by_timestamp_then_id, compact_json, empty_options, expand_home,
extract::{Extracted, extract_compact_repr, extract_raw_record, extract_str},
extracted_text, is_session_fresh,
jsonl::{
BoundedRow, JsonlTree, jsonl_tree_discover, jsonl_tree_events, parse_bounded,
peek_last_line, source_line,
},
jsonl_bytes, part_id, part_ordinal, raw_record,
sqlite::{CHANNEL_CAP, ColKind, columns_sql, db_error, emit, join_error, open_db, row_to_json},
};
const NAME: &str = "pi-coding-agent";
const SUPPORTED_JSONL_VERSION: i64 = 4;
const V3_FORMAT: i64 = 3;
pub struct PiCodingAgentFactory;
impl AdapterFactory for PiCodingAgentFactory {
fn name(&self) -> &'static str {
NAME
}
fn open(&self, config: Value) -> Result<Box<dyn Adapter>, AdapterError> {
Ok(Box::new(PiCodingAgentAdapter::from_config(config)?))
}
fn probe_default(&self, env: &Env) -> Option<Value> {
let path = env.home.join(".pi").join("agent").join("sessions");
path.exists().then(|| json!({ "path": path }))
}
fn serialize(
&self,
session: &crate::sessions::SessionWithMessages,
fidelity: RestoreFidelity,
) -> Result<Vec<RestoredFile>, AdapterError> {
serialize_session(session, fidelity)
}
}
#[derive(Debug, Clone, Deserialize)]
struct PiCodingAgentConfig {
path: PathBuf,
#[serde(default)]
sqlite_path: Option<PathBuf>,
}
fn replay_source_rows(session: &crate::sessions::SessionWithMessages) -> Option<Vec<Value>> {
let mut records = vec![raw_record(&session.session.options)?];
let mut messages: Vec<&crate::sessions::MessageWithParts> = session.messages.iter().collect();
messages.sort_by(|left, right| {
source_line(left.message.options())
.cmp(&source_line(right.message.options()))
.then_with(|| by_timestamp_then_id(left, right))
});
for message in messages {
records.push(raw_record(message.message.options())?);
}
Some(records)
}
fn serialize_session(
session: &crate::sessions::SessionWithMessages,
fidelity: RestoreFidelity,
) -> Result<Vec<RestoredFile>, AdapterError> {
let native_replayable =
session_format(&session.session) == V3_FORMAT && fidelity == RestoreFidelity::Native;
if native_replayable
&& let Some(records) = replay_source_rows(session)
&& let Some(path) = captured_relative_path(session)
{
return Ok(vec![RestoredFile::new(
path,
jsonl_bytes(NAME, &records)?,
RestoreFidelity::Native,
)]);
}
let mut records = vec![pi_session_record(session)];
let mut messages: Vec<&crate::sessions::MessageWithParts> = session.messages.iter().collect();
messages.sort_by(|left, right| by_timestamp_then_id(left, right));
for message in messages {
if matches!(message.message, Message::System { .. }) {
continue;
}
records.push(pi_message_record(message));
}
Ok(vec![RestoredFile::new(
reconstructed_relative_path(session),
jsonl_bytes(NAME, &records)?,
RestoreFidelity::Foreign,
)])
}
fn captured_relative_path(session: &crate::sessions::SessionWithMessages) -> Option<PathBuf> {
let slug = captured_slug(session)?;
let file_name = session
.session
.options
.get("source")?
.get("file_name")
.and_then(Value::as_str)?;
Some(PathBuf::from("sessions").join(slug).join(file_name))
}
fn captured_slug(session: &crate::sessions::SessionWithMessages) -> Option<&str> {
session
.session
.options
.get("source")?
.get("project_slug")
.and_then(Value::as_str)
}
fn reconstructed_relative_path(session: &crate::sessions::SessionWithMessages) -> PathBuf {
let slug = captured_slug(session)
.map(ToOwned::to_owned)
.unwrap_or_else(|| encode_project(&session.session.project));
let ts = session.session.created_at.format("%Y-%m-%dT%H-%M-%S-%3fZ");
PathBuf::from("sessions")
.join(slug)
.join(format!("{ts}-pond-v3_{}.jsonl", session.session.id))
}
fn encode_project(project: &str) -> String {
let body: String = project
.strip_prefix(['/', '\\'])
.unwrap_or(project)
.chars()
.map(|c| {
if matches!(c, '/' | '\\' | ':') {
'-'
} else {
c
}
})
.collect();
format!("--{body}--")
}
fn pi_session_record(session: &crate::sessions::SessionWithMessages) -> Value {
json!({
"type": "session",
"version": V3_FORMAT,
"id": session.session.id,
"timestamp": session.session.created_at.to_rfc3339_opts(SecondsFormat::Millis, true),
"cwd": &*session.session.project,
})
}
fn pi_message_record(message: &crate::sessions::MessageWithParts) -> Value {
json!({
"type": "message",
"id": message.message.id(),
"parentId": message.message.options().get("source").and_then(|s| s.get("parent_id")),
"timestamp": message.message.timestamp().to_rfc3339_opts(SecondsFormat::Millis, true),
"message": pi_inner_message(message),
})
}
fn pi_inner_message(message: &crate::sessions::MessageWithParts) -> Value {
let epoch_ms = message.message.timestamp().timestamp_millis();
match &message.message {
Message::User { .. } => json!({
"role": "user",
"content": message.parts.iter().map(pi_content_item).collect::<Vec<_>>(),
"timestamp": epoch_ms,
}),
Message::Assistant { .. } => json!({
"role": "assistant",
"content": message.parts.iter().map(pi_content_item).collect::<Vec<_>>(),
"timestamp": epoch_ms,
}),
Message::Tool { .. } => {
let part = message.parts.first();
let (call_id, name, is_error, result) = match part.map(|p| &p.kind) {
Some(PartKind::ToolResult {
call_id,
name,
is_failure,
result,
}) => (
extracted_text(call_id).to_owned(),
extracted_text(name).to_owned(),
*is_failure,
result.clone(),
),
_ => (String::new(), String::new(), false, Value::Null),
};
json!({
"role": "toolResult",
"toolCallId": call_id,
"toolName": name,
"content": result,
"isError": is_error,
"timestamp": epoch_ms,
})
}
Message::System { .. } => {
unreachable!("System messages are not serialized through pi_inner_message")
}
}
}
fn pi_content_item(part: &Part) -> Value {
match &part.kind {
PartKind::Text { text } => json!({"type": "text", "text": extracted_text(text)}),
PartKind::Reasoning { text } => json!({
"type": "thinking",
"thinking": extracted_text(text),
"thinkingSignature": part
.options
.get("pi")
.and_then(|p| p.get("thinking_signature")),
}),
PartKind::ToolCall {
call_id,
name,
params,
..
} => json!({
"type": "toolCall",
"id": extracted_text(call_id),
"name": extracted_text(name),
"arguments": params,
}),
other => json!({
"type": "text",
"text": compact_json(&serde_json::to_value(other).unwrap_or(Value::Null)),
}),
}
}
#[derive(Debug, Clone)]
pub struct PiCodingAgentAdapter {
root: PathBuf,
sqlite_path: Option<PathBuf>,
}
impl PiCodingAgentAdapter {
pub fn new(root: impl Into<PathBuf>) -> Self {
Self {
root: root.into(),
sqlite_path: None,
}
}
pub fn with_sqlite(mut self, path: impl Into<PathBuf>) -> Self {
self.sqlite_path = Some(path.into());
self
}
fn from_config(config: Value) -> Result<Self, AdapterError> {
let cfg: PiCodingAgentConfig = serde_json::from_value(config)
.map_err(|err| AdapterError::config(NAME, format!("bad config blob: {err}")))?;
Ok(Self {
root: expand_home(cfg.path),
sqlite_path: cfg.sqlite_path.map(expand_home),
})
}
}
impl Adapter for PiCodingAgentAdapter {
fn discover(&self) -> DiscoverFuture<'_> {
let files = jsonl_tree_discover(self);
let Some(db_path) = self.sqlite_path.clone() else {
return files;
};
Box::pin(async move {
let db_sessions = tokio::task::spawn_blocking(move || {
let conn = open_db(NAME, &db_path)?;
conn.query_row("SELECT COUNT(*) FROM sessions", [], |row| {
row.get::<_, i64>(0)
})
.map(|count| usize::try_from(count).unwrap_or(0))
.map_err(|error| db_error(NAME, &db_path, "count sessions", &error))
});
let (files, db_sessions) = tokio::join!(files, db_sessions);
Ok(files? + db_sessions.map_err(|join| join_error(NAME, join))??)
})
}
fn events_with<'a>(&'a self, oracle: &'a dyn SkipOracle) -> AdapterYieldStream<'a> {
let files = jsonl_tree_events(self, oracle);
match self.sqlite_path.clone() {
None => files,
Some(db_path) => Box::pin(files.chain(sqlite_events(db_path, oracle))),
}
}
fn plan<'a>(&'a self, oracle: &'a dyn SkipOracle) -> crate::adapter::PlanFuture<'a> {
let files = crate::adapter::jsonl::jsonl_tree_plan(self, oracle);
let Some(db_path) = self.sqlite_path.clone() else {
return files;
};
Box::pin(async move {
let (files, db) = tokio::join!(files, sqlite_plan(db_path, oracle));
let db = db?;
Ok(files?.map(|files| SyncPlan {
sessions: files.sessions + db.sessions,
fresh: files.fresh + db.fresh,
pending: files.pending + db.pending,
}))
})
}
}
impl JsonlTree for PiCodingAgentAdapter {
type State = ();
fn name(&self) -> &'static str {
NAME
}
fn root(&self) -> &Path {
&self.root
}
fn peek_session_id(&self, _path: &Path, first_line: &str) -> Option<String> {
let row: Value = serde_json::from_str(first_line).ok()?;
if !is_session_head(&row) {
return None;
}
row.get("id").and_then(Value::as_str).map(ToOwned::to_owned)
}
fn peek_watermark(&self, path: &Path) -> SourceWatermark {
let last = || -> Option<i64> {
let row: Value = serde_json::from_str(&peek_last_line(path)?).ok()?;
if is_session_head(&row) {
return None;
}
row_timestamp(&row).map(|ts| ts.timestamp_micros())
};
match last() {
Some(ts) => SourceWatermark::At(ts),
None => SourceWatermark::Opaque,
}
}
fn unsupported_reason(&self, path: &Path, rows: &[BoundedRow]) -> Option<String> {
let row = &rows.first()?.value;
let version = row.get("version").and_then(Value::as_i64)?;
(row.get("kind").and_then(Value::as_str) == Some("header")
&& version != SUPPORTED_JSONL_VERSION)
.then(|| {
format!(
"{}: pi session format version {version} is newer than this pond build \
understands (supported: {SUPPORTED_JSONL_VERSION}); upgrade pond",
path.display(),
)
})
}
fn session(&self, path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
session_from_rows(path, rows)
}
fn events_from_row(
&self,
session: &Session,
row: &BoundedRow,
_state: &mut Self::State,
) -> Result<Vec<IngestEvent>, String> {
match session_format(session) {
SUPPORTED_JSONL_VERSION => v4_events_from_mutation(
&session.id,
row.line as i64,
&row.value,
session.created_at,
),
_ => v3_events_from_row(&session.id, row.line, &row.value, session.created_at),
}
}
}
fn is_session_head(row: &Value) -> bool {
row.get("type").and_then(Value::as_str) == Some("session")
|| row.get("kind").and_then(Value::as_str) == Some("header")
}
pub(crate) fn row_timestamp(row: &Value) -> Option<DateTime<Utc>> {
match row.get("timestamp") {
Some(Value::String(text)) => DateTime::parse_from_rfc3339(text)
.ok()
.map(|dt| dt.with_timezone(&Utc)),
Some(Value::Number(number)) => DateTime::from_timestamp_millis(number.as_i64()?),
_ => None,
}
}
fn session_format(session: &Session) -> i64 {
session
.options
.get("source")
.and_then(|source| source.get("format"))
.and_then(Value::as_i64)
.unwrap_or(V3_FORMAT)
}
fn session_from_rows(path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
let path_display = path.display().to_string();
let first = rows
.first()
.ok_or_else(|| AdapterError::schema(NAME, path_display.clone(), "empty jsonl session"))?;
let row = &first.value;
let at_first = format!("{path_display}:{}", first.line);
let placement = SourcePlacement::from_path(path);
if row.get("kind").and_then(Value::as_str) == Some("header") {
return v4_session_from_header(row, &at_first, &placement);
}
if row.get("type").and_then(Value::as_str) != Some("session") {
return Err(AdapterError::schema(
NAME,
at_first,
"first row must be a v3 `session` record or a v4 `header`",
));
}
v3_session_from_row(NAME, row, &at_first, &placement)
}
#[derive(Debug, Default, Clone)]
pub(crate) struct SourcePlacement {
pub(crate) project_slug: Option<String>,
pub(crate) file_name: Option<String>,
}
impl SourcePlacement {
pub(crate) fn from_path(path: &Path) -> Self {
Self {
project_slug: path
.parent()
.and_then(|parent| parent.file_name())
.and_then(|name| name.to_str())
.map(ToOwned::to_owned),
file_name: path
.file_name()
.and_then(|name| name.to_str())
.map(ToOwned::to_owned),
}
}
}
pub(crate) fn v3_session_from_row(
name: &'static str,
row: &Value,
at_first: &str,
placement: &SourcePlacement,
) -> Result<Session, AdapterError> {
let id = row
.get("id")
.and_then(Value::as_str)
.ok_or_else(|| {
AdapterError::schema(name, at_first.to_owned(), "session record missing id")
})?
.to_owned();
let created_at = row
.get("timestamp")
.and_then(Value::as_str)
.and_then(|text| DateTime::parse_from_rfc3339(text).ok())
.map(|dt| dt.with_timezone(&Utc))
.ok_or_else(|| {
AdapterError::schema(
name,
at_first.to_owned(),
"session record has no parseable timestamp",
)
})?;
let project = extract_str(row, "cwd").ok_or_else(|| {
AdapterError::schema(name, at_first.to_owned(), "session record missing cwd")
})?;
let mut options = ProviderOptions::new();
options.insert(
"source".to_owned(),
json!({
"adapter": name,
"format": V3_FORMAT,
"version": row.get("version"),
"project_slug": placement.project_slug,
"file_name": placement.file_name,
"raw_record": extract_raw_record(row),
}),
);
Ok(Session {
id,
parent_session_id: None,
parent_message_id: None,
source_agent: name.to_owned(),
created_at,
project,
options,
})
}
pub(crate) fn v3_events_from_row(
session_id: &str,
line: usize,
row: &Value,
default_timestamp: DateTime<Utc>,
) -> Result<Vec<IngestEvent>, String> {
let kind = row.get("type").and_then(Value::as_str);
let timestamp = row_timestamp(row).unwrap_or(default_timestamp);
let id = row
.get("id")
.and_then(Value::as_str)
.map_or_else(|| format!("{session_id}:{line}"), ToOwned::to_owned);
let order = line as i64;
match kind {
Some("session") => Ok(Vec::new()),
Some("message") => {
let message_value = row
.get("message")
.ok_or_else(|| "message record missing `message` field".to_owned())?;
message_events(session_id, &id, timestamp, row, message_value, order)
}
Some("compaction") => Ok(vec![carrier_event(
session_id,
&id,
timestamp,
row,
order,
extract_str(row, "summary"),
)]),
_ => Ok(vec![carrier_event(
session_id,
&id,
timestamp,
row,
order,
extract_str(row, "type"),
)]),
}
}
fn v4_session_from_header(
row: &Value,
at_first: &str,
placement: &SourcePlacement,
) -> Result<Session, AdapterError> {
let id = row
.get("id")
.and_then(Value::as_str)
.ok_or_else(|| AdapterError::schema(NAME, at_first.to_owned(), "v4 header missing id"))?
.to_owned();
let created_at = row
.get("createdAt")
.and_then(Value::as_i64)
.and_then(DateTime::from_timestamp_millis)
.ok_or_else(|| {
AdapterError::schema(
NAME,
at_first.to_owned(),
"v4 header has no parseable createdAt",
)
})?;
let project = extract_str(row, "cwd")
.ok_or_else(|| AdapterError::schema(NAME, at_first.to_owned(), "v4 header missing cwd"))?;
let mut options = ProviderOptions::new();
options.insert(
"source".to_owned(),
json!({
"adapter": NAME,
"format": SUPPORTED_JSONL_VERSION,
"version": row.get("version"),
"project_slug": placement.project_slug,
"file_name": placement.file_name,
"metadata": row.get("metadata"),
"legacy_parent_session_path": row.get("legacyParentSessionPath"),
"raw_record": extract_raw_record(row),
}),
);
Ok(Session {
id,
parent_session_id: row
.get("parentSessionId")
.and_then(Value::as_str)
.map(ToOwned::to_owned),
parent_message_id: None,
source_agent: NAME.to_owned(),
created_at,
project,
options,
})
}
fn v4_events_from_mutation(
session_id: &str,
order: i64,
row: &Value,
default_timestamp: DateTime<Utc>,
) -> Result<Vec<IngestEvent>, String> {
if is_session_head(row) {
return Ok(Vec::new());
}
let timestamp = row_timestamp(row).unwrap_or(default_timestamp);
let id = row.get("id").and_then(Value::as_str).map_or_else(
|| {
let key = row.get("seq").and_then(Value::as_i64).unwrap_or(order);
format!("{session_id}:{key}")
},
ToOwned::to_owned,
);
let carrier = |content| {
Ok(vec![carrier_event(
session_id, &id, timestamp, row, order, content,
)])
};
match row.get("kind").and_then(Value::as_str) {
Some("entry") => match row.get("type").and_then(Value::as_str) {
Some("message") => {
let message_value = row
.get("message")
.ok_or_else(|| "v4 message entry missing `message` field".to_owned())?;
message_events(session_id, &id, timestamp, row, message_value, order)
}
Some("compaction" | "branch_summary") => carrier(extract_str(row, "summary")),
Some("custom") => carrier(extract_str(row, "customType")),
_ => carrier(extract_str(row, "type")),
},
Some("record") => carrier(extract_str(row, "type")),
Some("lane") => carrier(extract_str(row, "lane")),
Some("fact") => carrier(
extract_str(row, "name")
.or_else(|| extract_str(row, "label"))
.or_else(|| extract_str(row, "fact")),
),
_ => carrier(extract_str(row, "kind")),
}
}
fn carrier_event(
session_id: &str,
id: &str,
timestamp: DateTime<Utc>,
row: &Value,
order: i64,
content: Option<Extracted<String>>,
) -> IngestEvent {
IngestEvent::Message(Message::System {
id: id.to_owned(),
session_id: session_id.to_owned(),
timestamp,
content,
options: row_options(row, order),
})
}
fn message_events(
session_id: &str,
id: &str,
timestamp: DateTime<Utc>,
row: &Value,
message_value: &Value,
order: i64,
) -> Result<Vec<IngestEvent>, String> {
let role = message_value
.get("role")
.and_then(Value::as_str)
.ok_or_else(|| "message missing role".to_owned())?;
let content = message_value
.get("content")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
let mut parts = Vec::new();
let message = match role {
"user" => {
for (ordinal, item) in content.iter().enumerate() {
parts.push(user_part(session_id, id, ordinal, item));
}
Message::User {
id: id.to_owned(),
session_id: session_id.to_owned(),
timestamp,
options: row_options(row, order),
}
}
"assistant" => {
for (ordinal, item) in content.iter().enumerate() {
parts.push(assistant_part(session_id, id, ordinal, item));
}
Message::Assistant {
id: id.to_owned(),
session_id: session_id.to_owned(),
timestamp,
options: assistant_options(row, message_value, order),
}
}
"toolResult" => {
parts.push(tool_result_part(session_id, id, message_value));
Message::Tool {
id: id.to_owned(),
session_id: session_id.to_owned(),
timestamp,
options: row_options(row, order),
}
}
_ => Message::System {
id: id.to_owned(),
session_id: session_id.to_owned(),
timestamp,
content: extract_str(message_value, "role"),
options: row_options(row, order),
},
};
let mut events = Vec::with_capacity(parts.len() + 1);
events.push(IngestEvent::Message(message));
events.extend(parts.into_iter().map(IngestEvent::Part));
Ok(events)
}
fn user_part(session_id: &str, message_id: &str, ordinal: usize, item: &Value) -> Part {
let kind = match item.get("type").and_then(Value::as_str) {
Some("text") => PartKind::Text {
text: extract_str(item, "text"),
},
_ => PartKind::Text {
text: Some(extract_compact_repr(item)),
},
};
Part {
session_id: session_id.to_owned(),
id: part_id(message_id, ordinal),
message_id: message_id.to_owned(),
ordinal: part_ordinal(ordinal),
provenance: Provenance::Conversational,
options: empty_options(),
kind,
}
}
fn assistant_part(session_id: &str, message_id: &str, ordinal: usize, item: &Value) -> Part {
let (kind, options) = match item.get("type").and_then(Value::as_str) {
Some("text") => (
PartKind::Text {
text: extract_str(item, "text"),
},
empty_options(),
),
Some("thinking") => (
PartKind::Reasoning {
text: extract_str(item, "thinking"),
},
thinking_options(item),
),
Some("toolCall") => (
PartKind::ToolCall {
call_id: extract_str(item, "id"),
name: extract_str(item, "name"),
params: item.get("arguments").cloned().unwrap_or(Value::Null),
provider_executed: false,
},
empty_options(),
),
_ => (
PartKind::Text {
text: Some(extract_compact_repr(item)),
},
empty_options(),
),
};
Part {
session_id: session_id.to_owned(),
id: part_id(message_id, ordinal),
message_id: message_id.to_owned(),
ordinal: part_ordinal(ordinal),
provenance: Provenance::Conversational,
options,
kind,
}
}
fn tool_result_part(session_id: &str, message_id: &str, message_value: &Value) -> Part {
Part {
session_id: session_id.to_owned(),
id: part_id(message_id, 0),
message_id: message_id.to_owned(),
ordinal: 0,
provenance: Provenance::Injected,
options: empty_options(),
kind: PartKind::ToolResult {
call_id: extract_str(message_value, "toolCallId"),
name: extract_str(message_value, "toolName"),
is_failure: message_value
.get("isError")
.and_then(Value::as_bool)
.unwrap_or(false),
result: message_value.get("content").cloned().unwrap_or(Value::Null),
},
}
}
fn row_options(row: &Value, order: i64) -> ProviderOptions {
let mut source = serde_json::Map::new();
source.insert("line".to_owned(), json!(order));
for key in ["seq", "lane"] {
if let Some(value) = row.get(key) {
source.insert(key.to_owned(), value.clone());
}
}
source.insert(
"parent_id".to_owned(),
row.get("parentId").cloned().unwrap_or(Value::Null),
);
source.insert(
"raw_type".to_owned(),
row.get("type").cloned().unwrap_or(Value::Null),
);
source.insert("raw_record".to_owned(), extract_raw_record(row));
let mut options = ProviderOptions::new();
options.insert("source".to_owned(), Value::Object(source));
options
}
fn assistant_options(row: &Value, message_value: &Value, order: i64) -> ProviderOptions {
let mut options = row_options(row, order);
options.insert(
"pi".to_owned(),
json!({
"api": message_value.get("api"),
"provider": message_value.get("provider"),
"model": message_value.get("model"),
"usage": message_value.get("usage"),
"stop_reason": message_value.get("stopReason"),
"response_id": message_value.get("responseId"),
}),
);
options
}
fn thinking_options(item: &Value) -> ProviderOptions {
let mut options = ProviderOptions::new();
if let Some(signature) = item.get("thinkingSignature") {
options.insert("pi".to_owned(), json!({ "thinking_signature": signature }));
}
options
}
const DB_SESSION_COLUMNS: &[(&str, ColKind)] = &[
("id", ColKind::Str),
("created_at", ColKind::Str),
("cwd", ColKind::Str),
("parent_session_id", ColKind::Str),
("metadata", ColKind::Str),
];
struct DbSession {
id: String,
row: Value,
}
fn list_db_sessions(conn: &Connection, db_path: &Path) -> Result<Vec<DbSession>, AdapterError> {
let sql = format!(
"SELECT {} FROM sessions ORDER BY id",
columns_sql(DB_SESSION_COLUMNS)
);
let mut stmt = conn
.prepare(&sql)
.map_err(|error| db_error(NAME, db_path, "prepare sessions", &error))?;
let rows = stmt
.query_map([], |row| row_to_json(row, DB_SESSION_COLUMNS))
.map_err(|error| db_error(NAME, db_path, "query sessions", &error))?;
let mut out = Vec::new();
for row in rows {
let row = row.map_err(|error| db_error(NAME, db_path, "read session row", &error))?;
let Some(id) = row.get("id").and_then(Value::as_str).map(ToOwned::to_owned) else {
continue;
};
out.push(DbSession { id, row });
}
Ok(out)
}
const WATERMARK_SQL: &[&str] = &[
"SELECT seq, timestamp FROM entries WHERE session_id = ?1 ORDER BY seq DESC LIMIT 1",
"SELECT seq, timestamp FROM records WHERE session_id = ?1 ORDER BY seq DESC LIMIT 1",
"SELECT seq, NULL FROM lane_moves WHERE session_id = ?1 ORDER BY seq DESC LIMIT 1",
"SELECT seq, NULL FROM facts WHERE session_id = ?1 ORDER BY seq DESC LIMIT 1",
];
fn db_session_watermark(conn: &Connection, session_id: &str) -> Option<i64> {
let mut best: Option<(i64, Option<String>)> = None;
for sql in WATERMARK_SQL {
let row = conn
.prepare_cached(sql)
.ok()?
.query_row([session_id], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?))
})
.ok();
if let Some((seq, timestamp)) = row
&& best.as_ref().is_none_or(|(best_seq, _)| seq > *best_seq)
{
best = Some((seq, timestamp));
}
}
parse_db_timestamp(&best?.1?).map(|dt| dt.timestamp_micros())
}
fn parse_db_timestamp(text: &str) -> Option<DateTime<Utc>> {
DateTime::parse_from_rfc3339(text)
.ok()
.map(|dt| dt.with_timezone(&Utc))
}
fn sqlite_plan<'a>(
db_path: PathBuf,
oracle: &'a dyn SkipOracle,
) -> impl std::future::Future<Output = Result<SyncPlan, AdapterError>> + Send + 'a {
let oracle_is_empty = oracle.is_empty();
async move {
let heads = tokio::task::spawn_blocking(move || db_heads(&db_path, !oracle_is_empty))
.await
.map_err(|join| join_error(NAME, join))??;
if oracle_is_empty {
return Ok(SyncPlan::all_pending(heads.len()));
}
Ok(SyncPlan::from_heads(
oracle,
heads.iter().map(|(session, ts)| {
(
Some(session.id.as_str()),
ts.map_or(SourceWatermark::Opaque, SourceWatermark::At),
)
}),
))
}
}
fn db_heads(db_path: &Path, peek: bool) -> Result<Vec<(DbSession, Option<i64>)>, AdapterError> {
let conn = open_db(NAME, db_path)?;
Ok(list_db_sessions(&conn, db_path)?
.into_iter()
.map(|session| {
let watermark = peek
.then(|| db_session_watermark(&conn, &session.id))
.flatten();
(session, watermark)
})
.collect())
}
fn sqlite_events<'a>(db_path: PathBuf, oracle: &'a dyn SkipOracle) -> AdapterYieldStream<'a> {
Box::pin(async_stream::stream! {
let peek = !oracle.is_empty();
let list_path = db_path.clone();
let listed = tokio::task::spawn_blocking(move || db_heads(&list_path, peek)).await;
let heads = match listed {
Ok(Ok(heads)) => heads,
Ok(Err(error)) => { yield Err(error); return; }
Err(join) => { yield Err(join_error(NAME, join)); return; }
};
let mut survivors = Vec::with_capacity(heads.len());
let mut fresh = 0usize;
for (session, watermark) in heads {
if is_session_fresh(oracle, &session.id, watermark) {
fresh += 1;
continue;
}
survivors.push(session.row);
}
if fresh > 0 {
yield Ok(AdapterYield::SkippedBatch { reason: SkipReason::Fresh, count: fresh });
}
let (tx, mut rx) = mpsc::channel(CHANNEL_CAP);
let handle = tokio::task::spawn_blocking(move || read_db_sessions(&db_path, &survivors, &tx));
while let Some(item) = rx.recv().await {
yield item;
}
if let Err(join) = handle.await {
yield Err(join_error(NAME, join));
}
})
}
fn read_db_sessions(
db_path: &Path,
sessions: &[Value],
tx: &mpsc::Sender<Result<AdapterYield, AdapterError>>,
) -> bool {
let conn = match open_db(NAME, db_path) {
Ok(conn) => conn,
Err(error) => {
emit!(tx, Err(error));
return true;
}
};
for row in sessions {
if !read_db_session(&conn, db_path, row, tx) {
return false;
}
}
true
}
fn read_db_session(
conn: &Connection,
db_path: &Path,
row: &Value,
tx: &mpsc::Sender<Result<AdapterYield, AdapterError>>,
) -> bool {
let at = |session_id: &str| format!("{}#{session_id}", db_path.display());
let header = db_header_value(row);
let session_id = row.get("id").and_then(Value::as_str).unwrap_or_default();
let session =
match v4_session_from_header(&header, &at(session_id), &SourcePlacement::default()) {
Ok(session) => session,
Err(error) => {
emit!(tx, Err(error));
return true;
}
};
emit!(
tx,
Ok(AdapterYield::Event(IngestEvent::Session(session.clone())))
);
let mutations = match db_mutations(conn, db_path, &session.id) {
Ok(mutations) => mutations,
Err(error) => {
emit!(tx, Err(error));
return true;
}
};
for (seq, mutation) in mutations {
let mutation = match mutation {
Ok(mutation) => mutation,
Err(error) => {
emit!(tx, Err(error));
continue;
}
};
match v4_events_from_mutation(&session.id, seq, &mutation, session.created_at) {
Ok(events) => {
for event in events {
emit!(tx, Ok(AdapterYield::Event(event)));
}
}
Err(message) => emit!(
tx,
Err(AdapterError::schema(
NAME,
format!("{}:{seq}", at(&session.id)),
message
))
),
}
}
true
}
fn db_header_value(row: &Value) -> Value {
let mut header = serde_json::Map::new();
header.insert("kind".to_owned(), json!("header"));
header.insert("version".to_owned(), json!(SUPPORTED_JSONL_VERSION));
if let Some(id) = row.get("id") {
header.insert("id".to_owned(), id.clone());
}
if let Some(created_at) = row
.get("created_at")
.and_then(Value::as_str)
.and_then(parse_db_timestamp)
{
header.insert("createdAt".to_owned(), json!(created_at.timestamp_millis()));
}
if let Some(cwd) = row.get("cwd") {
header.insert("cwd".to_owned(), cwd.clone());
}
if let Some(parent) = row.get("parent_session_id") {
header.insert("parentSessionId".to_owned(), parent.clone());
}
if let Some(text) = row.get("metadata").and_then(Value::as_str) {
header.insert(
"metadata".to_owned(),
serde_json::from_str(text).unwrap_or_else(|_| json!(text)),
);
}
Value::Object(header)
}
struct EntryRow {
seq: i64,
id: String,
parent_id: Option<String>,
kind: String,
timestamp: String,
payload: String,
}
type DbMutation = (i64, Result<Value, AdapterError>);
fn db_mutations(
conn: &Connection,
db_path: &Path,
session_id: &str,
) -> Result<Vec<DbMutation>, AdapterError> {
let at = |table: &str, seq: i64| format!("{}#{session_id}.{table}:{seq}", db_path.display());
let mut out = Vec::new();
for row in query_db_rows(
conn,
db_path,
session_id,
"SELECT seq, id, parent_id, type, timestamp, payload FROM entries WHERE session_id = ?1",
|row| {
Ok(EntryRow {
seq: row.get(0)?,
id: row.get(1)?,
parent_id: row.get(2)?,
kind: row.get(3)?,
timestamp: row.get(4)?,
payload: row.get(5)?,
})
},
)? {
let seq = row.seq;
let mutation = payload_object(&row.payload, || at("entries", seq)).map(|mut map| {
map.insert("kind".to_owned(), json!("entry"));
map.insert("id".to_owned(), json!(row.id));
map.insert(
"parentId".to_owned(),
row.parent_id.map_or(Value::Null, |parent| json!(parent)),
);
map.insert("type".to_owned(), json!(row.kind));
map.insert("seq".to_owned(), json!(seq));
insert_db_timestamp(&mut map, &row.timestamp);
Value::Object(map)
});
out.push((seq, mutation));
}
for (seq, timestamp, payload) in query_db_rows(
conn,
db_path,
session_id,
"SELECT seq, timestamp, payload FROM records WHERE session_id = ?1",
|row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
))
},
)? {
let mutation = payload_object(&payload, || at("records", seq)).map(|mut map| {
map.insert("kind".to_owned(), json!("record"));
map.insert("seq".to_owned(), json!(seq));
insert_db_timestamp(&mut map, ×tamp);
Value::Object(map)
});
out.push((seq, mutation));
}
for (seq, lane, leaf_id) in query_db_rows(
conn,
db_path,
session_id,
"SELECT seq, lane, leaf_id FROM lane_moves WHERE session_id = ?1",
|row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
))
},
)? {
out.push((
seq,
Ok(json!({ "kind": "lane", "seq": seq, "lane": lane, "leafId": leaf_id })),
));
}
for (seq, kind, key, value) in query_db_rows(
conn,
db_path,
session_id,
"SELECT seq, kind, key, value FROM facts WHERE session_id = ?1",
|row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, Option<String>>(3)?,
))
},
)? {
let mut map = serde_json::Map::new();
map.insert("kind".to_owned(), json!("fact"));
map.insert("seq".to_owned(), json!(seq));
map.insert("fact".to_owned(), json!(kind));
if let Some(target) = key {
map.insert("targetId".to_owned(), json!(target));
}
let mutation = match value {
None => Ok(Value::Object(map)),
Some(text) => {
parse_bounded(NAME, text.as_bytes(), || at("facts", seq)).map(|decoded| {
map.insert(
if kind == "name" { "name" } else { "label" }.to_owned(),
decoded,
);
Value::Object(map)
})
}
};
out.push((seq, mutation));
}
out.sort_by_key(|(seq, _)| *seq);
Ok(out)
}
fn query_db_rows<T>(
conn: &Connection,
db_path: &Path,
session_id: &str,
sql: &str,
map: impl Fn(&rusqlite::Row) -> rusqlite::Result<T>,
) -> Result<Vec<T>, AdapterError> {
let mut stmt = conn
.prepare_cached(sql)
.map_err(|error| db_error(NAME, db_path, "prepare mutations", &error))?;
let rows = stmt
.query_map([session_id], map)
.map_err(|error| db_error(NAME, db_path, "query mutations", &error))?;
rows.map(|row| row.map_err(|error| db_error(NAME, db_path, "read mutation row", &error)))
.collect()
}
fn payload_object(
text: &str,
at: impl Fn() -> String,
) -> Result<serde_json::Map<String, Value>, AdapterError> {
match parse_bounded(NAME, text.as_bytes(), &at)? {
Value::Object(map) => Ok(map),
other => Err(AdapterError::schema(
NAME,
at(),
format!(
"payload column is {}, not a JSON object",
match other {
Value::Null => "null",
Value::Bool(_) => "a boolean",
Value::Number(_) => "a number",
Value::String(_) => "a string",
Value::Array(_) => "an array",
Value::Object(_) => unreachable!("matched above"),
}
),
)),
}
}
fn insert_db_timestamp(map: &mut serde_json::Map<String, Value>, text: &str) {
if let Some(dt) = parse_db_timestamp(text) {
map.insert("timestamp".to_owned(), json!(dt.timestamp_millis()));
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::expect_used, clippy::unwrap_used)]
use super::*;
use crate::{handlers::ingest_adapter, sessions::Store, wire::PartKind};
use tempfile::TempDir;
const FIXTURES: &str = concat!(
env!("CARGO_MANIFEST_DIR"),
"/tests/fixtures/adapter/pi-coding-agent/sessions"
);
#[test]
fn probe_default_finds_pi_sessions_under_home() -> anyhow::Result<()> {
crate::adapter::test_support::assert_probe_default(
&PiCodingAgentFactory,
&[".pi", "agent", "sessions"],
)
}
#[test]
fn encode_project_matches_pi_session_directory_name() {
assert_eq!(encode_project("/Users/user/pond"), "--Users-user-pond--");
assert_eq!(
encode_project(r"C:\Users\user\pond"),
"--C--Users-user-pond--"
);
assert_eq!(encode_project(r"\\host\share\pond"), "---host-share-pond--");
}
#[tokio::test(flavor = "multi_thread")]
async fn codec_replays_every_fixture_back_to_its_source_bytes() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
let corpus = Path::new(FIXTURES)
.parent()
.expect("FIXTURES is nested under a corpus root");
let mut replayed = std::collections::BTreeSet::new();
for session_id in store.session_ids().await? {
let session = store
.get_session(&session_id)
.await?
.expect("session round-trips");
let records =
replay_source_rows(&session).expect("a fixture session carries its source rows");
let relative = captured_relative_path(&session)
.expect("a fixture session carries its file placement");
let expected = std::fs::read_to_string(corpus.join(&relative))?;
let expected: Vec<Value> = expected
.lines()
.filter(|line| !line.trim().is_empty())
.map(serde_json::from_str)
.collect::<Result<_, _>>()?;
assert_eq!(
records,
expected,
"replay mismatch at {}",
relative.display()
);
replayed.insert(relative);
}
let mut on_disk = std::collections::BTreeSet::new();
for dir in std::fs::read_dir(FIXTURES)? {
for file in std::fs::read_dir(dir?.path())? {
let path = file?.path();
if path.extension().and_then(|ext| ext.to_str()) == Some("jsonl") {
on_disk.insert(path.strip_prefix(corpus)?.to_path_buf());
}
}
}
assert_eq!(
replayed, on_disk,
"every fixture file must be covered by exactly one replayed session",
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn resume_emits_v3_for_every_origin_and_reports_the_downgrade() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
for (session_id, expected_fidelity) in [
(V4_SESSION, RestoreFidelity::Foreign),
(
"019dd55d-99a4-7344-aa11-d1d71d2c80fb",
RestoreFidelity::Native,
),
] {
let session = store
.get_session(session_id)
.await?
.unwrap_or_else(|| panic!("{session_id} lands"));
let files = PiCodingAgentFactory.serialize(&session, RestoreFidelity::Native)?;
let file = files.first().expect("one file per session");
assert_eq!(
file.actual_fidelity, expected_fidelity,
"{session_id} fidelity",
);
let head: Value =
serde_json::from_str(std::str::from_utf8(&file.bytes)?.lines().next().unwrap())?;
assert_eq!(head.get("type"), Some(&json!("session")), "{session_id}");
assert_eq!(head.get("version"), Some(&json!(3)), "{session_id}");
}
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn a_downgraded_restore_is_named_apart_from_its_source_file() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
let v4 = store
.get_session(V4_SESSION)
.await?
.expect("v4 fixture session lands");
let captured =
captured_relative_path(&v4).expect("a v4 fixture session carries its placement");
let restored = PiCodingAgentFactory.serialize(&v4, RestoreFidelity::Native)?;
let path = restored
.first()
.expect("one file per session")
.relative_path
.clone();
assert_ne!(
path, captured,
"the v3 reconstruction may not claim the v4 source file's name",
);
assert_eq!(
path.parent(),
captured.parent(),
"it still lands in the session's own project-slug directory",
);
assert!(
path.to_str()
.expect("utf-8 path")
.ends_with(&format!("_{V4_SESSION}.jsonl")),
"the pi-pond extension picks the file to open by this suffix: {}",
path.display(),
);
let again = PiCodingAgentFactory.serialize(&v4, RestoreFidelity::Native)?;
assert_eq!(
again.first().expect("one file per session").relative_path,
path,
"the name is deterministic, so a second resume refuses against it",
);
let v3 = store
.get_session("019dd55d-99a4-7344-aa11-d1d71d2c80fb")
.await?
.expect("v3 fixture session lands");
let v3_files = PiCodingAgentFactory.serialize(&v3, RestoreFidelity::Native)?;
assert_eq!(
v3_files
.first()
.expect("one file per session")
.relative_path,
captured_relative_path(&v3).expect("a v3 session carries its placement"),
"a v3-origin replay still restores to its own source file",
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn pi_coding_agent_adapter_ingests_fixture_corpus_into_canonical_shape()
-> anyhow::Result<()> {
let temp = TempDir::new()?;
let store = Store::open_local(temp.path()).await?;
let adapter = PiCodingAgentAdapter::new(FIXTURES);
let summary = ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
assert!(summary.accepted() > 0, "ingest must accept rows");
assert_eq!(summary.dropped_events, 0, "no per-event drops expected");
assert_eq!(
summary.dropped_sessions, 0,
"no session-level rejections expected"
);
assert_eq!(summary.skipped_files, 0, "no whole-file skips expected");
let (sessions, messages, parts) = store.row_counts().await?;
assert!(sessions > 0, "at least one pi-coding-agent session");
assert!(messages > 0, "at least one pi-coding-agent message");
assert!(parts > 0, "at least one pi-coding-agent Part");
let mut saw_tool_call = false;
let mut saw_tool_result = false;
let mut saw_reasoning = false;
for session_id in store.session_ids().await? {
let session = store
.get_session(&session_id)
.await?
.expect("session round-trips");
assert_eq!(session.session.source_agent, NAME);
assert!(
!(*session.session.project).is_empty(),
"spec.md#model-project-non-empty: project must be a real cwd",
);
for stored in &session.messages {
for part in &stored.parts {
match &part.kind {
PartKind::ToolCall { .. } => saw_tool_call = true,
PartKind::ToolResult { .. } => saw_tool_result = true,
PartKind::Reasoning { .. } => saw_reasoning = true,
_ => {}
}
}
}
}
assert!(saw_tool_call, "corpus has assistant tool calls");
assert!(saw_tool_result, "corpus has tool results");
assert!(saw_reasoning, "corpus has assistant reasoning");
Ok(())
}
#[test]
fn unknown_nested_message_role_becomes_system_carrier() -> anyhow::Result<()> {
let row = json!({
"type": "message",
"id": "mystery-message",
"message": {
"role": "mysteryRole",
"content": [{"type": "text", "text": "not yet understood"}]
}
});
let events = v3_events_from_row(
"session-1",
42,
&row,
DateTime::parse_from_rfc3339("2026-04-28T18:47:32.280Z")?.with_timezone(&Utc),
)
.map_err(anyhow::Error::msg)?;
assert_eq!(events.len(), 1);
let IngestEvent::Message(Message::System {
id,
content,
options,
..
}) = &events[0]
else {
panic!("unknown role must produce a System carrier");
};
assert_eq!(id, "mystery-message");
assert_eq!(content.as_deref().map(String::as_str), Some("mysteryRole"));
assert_eq!(
raw_record(options)
.and_then(|raw| raw.get("message").cloned())
.and_then(|message| message.get("role").cloned()),
Some(json!("mysteryRole")),
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn fork_parent_ids_and_compaction_summary_are_preserved() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let root = temp.path().join("sessions");
let path = root
.join("project")
.join("2026-05-01T00-00-00-000Z_fork.jsonl");
write_jsonl_file(
&path,
&[
json!({
"type": "session",
"version": 3,
"id": "pi-fork-session",
"timestamp": "2026-05-01T00:00:00.000Z",
"cwd": "/tmp/pi-fork",
}),
json!({
"type": "message",
"id": "parent-message",
"timestamp": "2026-05-01T00:00:01.000Z",
"message": {
"role": "user",
"content": [{"type": "text", "text": "parent"}],
},
}),
json!({
"type": "message",
"id": "child-a",
"parentId": "parent-message",
"timestamp": "2026-05-01T00:00:02.000Z",
"message": {
"role": "assistant",
"content": [{"type": "text", "text": "branch a"}],
},
}),
json!({
"type": "message",
"id": "child-b",
"parentId": "parent-message",
"timestamp": "2026-05-01T00:00:03.000Z",
"message": {
"role": "assistant",
"content": [{"type": "text", "text": "branch b"}],
},
}),
json!({
"type": "compaction",
"id": "compact-1",
"parentId": "child-b",
"timestamp": "2026-05-01T00:00:04.000Z",
"summary": "compact summary",
}),
],
)?;
let store = Store::open_local(temp.path().join("store")).await?;
let summary = ingest_adapter(
&store,
&PiCodingAgentAdapter::new(&root),
&crate::adapter::NoopOracle,
|_| {},
)
.await?;
assert_eq!(summary.dropped_events, 0);
let session = store
.get_session("pi-fork-session")
.await?
.expect("fixture session lands");
let child_a = session
.messages
.iter()
.find(|stored| stored.message.id() == "child-a")
.expect("first fork child lands");
let child_b = session
.messages
.iter()
.find(|stored| stored.message.id() == "child-b")
.expect("second fork child lands");
for child in [child_a, child_b] {
assert_eq!(
child
.message
.options()
.get("source")
.and_then(|source| source.get("parent_id"))
.and_then(Value::as_str),
Some("parent-message"),
);
}
assert!(source_line(child_a.message.options()) < source_line(child_b.message.options()));
let compact = session
.messages
.iter()
.find(|stored| stored.message.id() == "compact-1")
.expect("compaction carrier lands");
let Message::System { content, .. } = &compact.message else {
panic!("compaction is preserved as a System carrier");
};
assert_eq!(
content.as_deref().map(String::as_str),
Some("compact summary")
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn foreign_serialization_reparses_as_pi_coding_agent() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let origin_store = Store::open_local(temp.path().join("origin-store")).await?;
let origin = crate::adapter::OpencodeAdapter::new(concat!(
env!("CARGO_MANIFEST_DIR"),
"/tests/fixtures/adapter/opencode/storage"
));
ingest_adapter(&origin_store, &origin, &crate::adapter::NoopOracle, |_| {}).await?;
let session_id = origin_store
.session_ids()
.await?
.into_iter()
.next()
.expect("opencode fixture has sessions");
let session = origin_store
.get_session(&session_id)
.await?
.expect("fixture session is readable");
let restored_root = temp.path().join("pi-corpus");
crate::adapter::write_restored_files(
&restored_root,
PiCodingAgentFactory.serialize(&session, RestoreFidelity::Foreign)?,
)?;
let restored_store = Store::open_local(temp.path().join("restored-store")).await?;
let summary = ingest_adapter(
&restored_store,
&PiCodingAgentAdapter::new(restored_root.join("sessions")),
&crate::adapter::NoopOracle,
|_| {},
)
.await?;
assert!(summary.accepted() > 0);
assert_eq!(summary.dropped_events, 0);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn tool_results_are_injected_assistant_parts_are_conversational() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let store = Store::open_local(temp.path()).await?;
let adapter = PiCodingAgentAdapter::new(FIXTURES);
ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
for session_id in store.session_ids().await? {
let session = store
.get_session(&session_id)
.await?
.expect("session round-trips");
for stored in &session.messages {
for part in &stored.parts {
match &part.kind {
PartKind::ToolResult { .. } => {
assert_eq!(part.provenance, Provenance::Injected);
}
PartKind::ToolCall { .. } | PartKind::Reasoning { .. } => {
assert_eq!(part.provenance, Provenance::Conversational);
}
_ => {}
}
}
}
}
Ok(())
}
const V4_SESSION: &str = "v4-main-session";
const V4_FORK: &str = "v4-fork-session";
fn scratch_v4_corpus(temp: &TempDir) -> anyhow::Result<(PathBuf, PathBuf)> {
let source = Path::new(FIXTURES)
.join("--Users-user-Projects-harness-v2--")
.join("2026-08-06T00-00-01-000Z_v4-main-session.jsonl");
let root = temp.path().join("sessions");
let dir = root.join("--Users-user-Projects-harness-v2--");
std::fs::create_dir_all(&dir)?;
let dest = dir.join("2026-08-06T00-00-01-000Z_v4-main-session.jsonl");
std::fs::copy(&source, &dest)?;
Ok((root, dest))
}
async fn ingest_into_temp_store(
temp: &TempDir,
adapter: &PiCodingAgentAdapter,
) -> anyhow::Result<(Store, crate::sessions::IngestSummary)> {
let store = Store::open_local(temp.path().join("store")).await?;
let summary = ingest_adapter(&store, adapter, &crate::adapter::NoopOracle, |_| {}).await?;
Ok((store, summary))
}
#[tokio::test(flavor = "multi_thread")]
async fn v4_header_carries_lineage_project_and_metadata() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, summary) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
assert_eq!(summary.dropped_events, 0);
assert_eq!(summary.dropped_sessions, 0);
let main = store
.get_session(V4_SESSION)
.await?
.expect("v4 fixture session lands");
assert_eq!(main.session.source_agent, NAME);
assert_eq!(&*main.session.project, "/Users/user/Projects/harness-v2");
assert_eq!(main.session.parent_session_id, None);
let source = main
.session
.options
.get("source")
.expect("source options present");
assert_eq!(source.get("format"), Some(&json!(4)));
assert_eq!(
source.get("metadata"),
Some(&json!({"harness": "v2", "fixture": "pond"})),
"spec.md#model-lossless-projection: the header metadata bag is carried verbatim",
);
let fork = store
.get_session(V4_FORK)
.await?
.expect("v4 fork session lands");
assert_eq!(
fork.session.parent_session_id.as_deref(),
Some(V4_SESSION),
"adapter-lineage: a v4 fork names its parent session",
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn v4_orchestration_mutations_are_system_carriers() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
let session = store
.get_session(V4_SESSION)
.await?
.expect("v4 fixture session lands");
let find = |id: &str| {
session
.messages
.iter()
.find(|stored| stored.message.id() == id)
.unwrap_or_else(|| panic!("{id} lands"))
.message
.clone()
};
for (id, content) in [
("v4-run-1", "operation_started"),
("v4-tool-1", "tool_started"),
("v4-usage-1", "usage"),
("v4-entry-model", "model_change"),
(
"v4-entry-compaction",
"Compacted the storage-rewrite discussion.",
),
(
"v4-entry-branch-summary",
"Explored the SQLite backend, came back.",
),
("v4-entry-custom", "pond-fixture"),
] {
let Message::System {
content: got,
options,
..
} = find(id)
else {
panic!("{id} must be a System carrier, not conversation");
};
assert_eq!(got.as_deref().map(String::as_str), Some(content));
assert!(
raw_record(&options).is_some(),
"{id} must keep its whole source mutation",
);
}
let lane = find(&format!("{V4_SESSION}:19"));
let Message::System { content, .. } = &lane else {
panic!("a lane pointer move is a carrier");
};
assert_eq!(content.as_deref().map(String::as_str), Some("side"));
let name_fact = find(&format!("{V4_SESSION}:22"));
let Message::System { content, .. } = &name_fact else {
panic!("a name fact is a carrier");
};
assert_eq!(
content.as_deref().map(String::as_str),
Some("harness-v2 storage rewrite"),
);
assert!(
matches!(find("v4-entry-user"), Message::User { .. }),
"a v4 message entry is real conversation",
);
assert!(matches!(
find("v4-entry-assistant"),
Message::Assistant { .. }
));
assert!(matches!(find("v4-entry-tool-result"), Message::Tool { .. }));
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn v4_unknown_mutation_kind_degrades_to_a_carrier() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (root, path) = scratch_v4_corpus(&temp)?;
let mut text = std::fs::read_to_string(&path)?;
text.push_str(
r#"{"kind":"telepathy","seq":99,"id":"future-1","timestamp":1785974999000,"payload":{"from":"tomorrow"}}"#,
);
text.push('\n');
std::fs::write(&path, text)?;
let (store, summary) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(&root)).await?;
assert_eq!(
summary.dropped_events, 0,
"unknown-but-well-formed is preserved"
);
let session = store
.get_session(V4_SESSION)
.await?
.expect("session still lands");
let future = session
.messages
.iter()
.find(|stored| stored.message.id() == "future-1")
.expect("the future mutation is preserved");
let Message::System {
content, options, ..
} = &future.message
else {
panic!("an unknown mutation kind must become a carrier");
};
assert_eq!(content.as_deref().map(String::as_str), Some("telepathy"));
assert_eq!(
raw_record(options).and_then(|raw| raw.get("payload").cloned()),
Some(json!({"from": "tomorrow"})),
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn v4_torn_tail_keeps_the_whole_lines_and_counts_the_partial_one() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (root, path) = scratch_v4_corpus(&temp)?;
let mut text = std::fs::read_to_string(&path)?;
let whole_lines = text.lines().count();
text.push_str(r#"{"kind":"entry","lane":"main","id":"v4-torn-part"#);
std::fs::write(&path, text)?;
let (store, summary) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(&root)).await?;
assert_eq!(
summary.skipped_files, 1,
"the torn line surfaces as a named skip, not a silence",
);
let session = store
.get_session(V4_SESSION)
.await?
.expect("session still lands");
assert_eq!(
session.messages.len(),
whole_lines - 1,
"every complete mutation before the tear still ingests (the header is eventless)",
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn v4_unsupported_version_is_a_named_skip() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (root, path) = scratch_v4_corpus(&temp)?;
let text = std::fs::read_to_string(&path)?;
std::fs::write(&path, text.replacen(r#""version":4"#, r#""version":5"#, 1))?;
let (store, summary) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(&root)).await?;
assert_eq!(summary.skipped_files, 1, "the file is skipped visibly");
assert_eq!(summary.dropped_events, 0);
assert!(
store.session_ids().await?.is_empty(),
"nothing from an undecodable version reaches the store",
);
Ok(())
}
#[test]
fn unsupported_reason_reads_rows_not_the_file() {
let adapter = PiCodingAgentAdapter::new("/tmp/pond-test-root");
let absent = Path::new("/tmp/pond-test-root/does-not-exist.jsonl");
let head = |version: i64| {
vec![BoundedRow {
line: 1,
value: json!({"kind": "header", "version": version, "id": "s1"}),
}]
};
assert!(
adapter
.unsupported_reason(absent, &head(SUPPORTED_JSONL_VERSION + 1))
.is_some_and(|reason| reason.contains("upgrade pond")),
);
assert!(
adapter
.unsupported_reason(absent, &head(SUPPORTED_JSONL_VERSION))
.is_none(),
);
assert!(adapter.unsupported_reason(absent, &[]).is_none());
}
#[test]
fn v4_trailing_timestampless_mutation_forces_a_reread() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (root, path) = scratch_v4_corpus(&temp)?;
let adapter = PiCodingAgentAdapter::new(&root);
assert_eq!(
adapter.peek_watermark(&path),
SourceWatermark::Opaque,
"the fixture ends on a `fact` mutation, which has no timestamp",
);
let mut text = std::fs::read_to_string(&path)?;
text.push_str(
r#"{"kind":"entry","lane":"main","id":"later","type":"custom","customType":"x","parentId":null,"seq":24,"timestamp":1785974500000}"#,
);
text.push('\n');
std::fs::write(&path, text)?;
assert_eq!(
adapter.peek_watermark(&path),
SourceWatermark::At(1_785_974_500_000_000),
"a timestamped tail yields a real watermark",
);
Ok(())
}
const SQLITE_DB: &str = concat!(
env!("CARGO_MANIFEST_DIR"),
"/tests/fixtures/adapter/pi-coding-agent/sqlite/pi-sessions.sqlite"
);
fn sqlite_only_adapter(temp: &TempDir) -> anyhow::Result<PiCodingAgentAdapter> {
let root = temp.path().join("empty-sessions");
std::fs::create_dir_all(&root)?;
Ok(PiCodingAgentAdapter::new(root).with_sqlite(SQLITE_DB))
}
#[tokio::test(flavor = "multi_thread")]
async fn sqlite_backend_maps_through_the_same_v4_mapper() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let adapter = sqlite_only_adapter(&temp)?;
let (store, summary) = ingest_into_temp_store(&temp, &adapter).await?;
assert_eq!(summary.dropped_events, 0);
assert_eq!(summary.dropped_sessions, 0);
let mut ids = store.session_ids().await?;
ids.sort();
assert_eq!(ids, ["sqlite-child-session", "sqlite-main-session"]);
let main = store
.get_session("sqlite-main-session")
.await?
.expect("sqlite session lands");
assert_eq!(main.session.source_agent, NAME);
assert_eq!(&*main.session.project, "/Users/user/Projects/harness-v2");
assert_eq!(
main.session
.options
.get("source")
.and_then(|source| source.get("format")),
Some(&json!(4)),
"database rows are mapped as v4 mutations, not a fourth shape",
);
let child = store
.get_session("sqlite-child-session")
.await?
.expect("sqlite child lands");
assert_eq!(
child.session.parent_session_id.as_deref(),
Some("sqlite-main-session"),
);
let mut saw_tool_call = false;
let mut saw_tool_result = false;
for stored in &main.messages {
for part in &stored.parts {
match &part.kind {
PartKind::ToolCall { name, .. } => {
saw_tool_call = true;
assert_eq!(name.as_deref().map(String::as_str), Some("bash"));
}
PartKind::ToolResult { .. } => saw_tool_result = true,
_ => {}
}
}
}
assert!(saw_tool_call && saw_tool_result);
let carriers: Vec<&str> = main
.messages
.iter()
.filter_map(|stored| match &stored.message {
Message::System { content, .. } => content.as_deref().map(String::as_str),
_ => None,
})
.collect();
assert!(
carriers.contains(&"sqlite backend probe"),
"name fact lands"
);
assert!(carriers.contains(&"answer"), "label fact lands");
assert!(carriers.contains(&"side"), "lane move lands");
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn sqlite_a_corrupt_row_drops_only_itself() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let db_path = temp.path().join("pi-sessions.sqlite");
Connection::open(&db_path)?.execute_batch(
"CREATE TABLE sessions (
id TEXT PRIMARY KEY, created_at TEXT NOT NULL, cwd TEXT NOT NULL,
parent_session_id TEXT NULL, metadata TEXT NULL);
CREATE TABLE entries (
session_id TEXT NOT NULL, seq INTEGER NOT NULL, id TEXT NOT NULL,
parent_id TEXT NULL, type TEXT NOT NULL, timestamp TEXT NOT NULL,
payload TEXT NOT NULL, PRIMARY KEY (session_id, id));
CREATE TABLE records (
session_id TEXT NOT NULL, seq INTEGER NOT NULL, id TEXT NOT NULL,
lane TEXT NOT NULL, run_id TEXT NULL, type TEXT NOT NULL,
op_kind TEXT NULL, timestamp TEXT NOT NULL, payload TEXT NOT NULL,
PRIMARY KEY (session_id, id));
CREATE TABLE lane_moves (
session_id TEXT NOT NULL, seq INTEGER NOT NULL, lane TEXT NOT NULL,
leaf_id TEXT NULL, PRIMARY KEY (session_id, seq));
CREATE TABLE facts (
session_id TEXT NOT NULL, seq INTEGER NOT NULL, kind TEXT NOT NULL,
key TEXT NULL, value TEXT NULL, PRIMARY KEY (session_id, seq));
INSERT INTO sessions VALUES
('torn-session', '2026-08-06T00:00:00.000Z', '/Users/user/Projects/torn', NULL, NULL);
INSERT INTO entries (session_id, seq, id, parent_id, type, timestamp, payload) VALUES
('torn-session', 1, 'good-user', NULL, 'message', '2026-08-06T00:00:01.000Z',
'{\"message\":{\"role\":\"user\",\"content\":[{\"type\":\"text\",\"text\":\"kept\"}]}}'),
('torn-session', 2, 'corrupt-entry', NULL, 'message', '2026-08-06T00:00:02.000Z',
'{\"message\":{\"role\":\"user\",'),
('torn-session', 3, 'good-assistant', 'good-user', 'message', '2026-08-06T00:00:03.000Z',
'{\"message\":{\"role\":\"assistant\",\"content\":[{\"type\":\"text\",\"text\":\"answered\"}]}}');
INSERT INTO facts (session_id, seq, kind, key, value) VALUES
('torn-session', 4, 'name', NULL, '\"still named\"'),
('torn-session', 5, 'label', 'good-user', '{\"torn');",
)?;
let root = temp.path().join("empty-sessions");
std::fs::create_dir_all(&root)?;
let adapter = PiCodingAgentAdapter::new(root).with_sqlite(&db_path);
let (store, summary) = ingest_into_temp_store(&temp, &adapter).await?;
assert_eq!(
summary.dropped_events, 2,
"the two unparseable payloads are counted drops, nothing else",
);
assert_eq!(
summary.skipped_files, 0,
"a bad row is not a whole-source failure",
);
let session = store
.get_session("torn-session")
.await?
.expect("the session survives its bad rows");
let ids: Vec<&str> = session
.messages
.iter()
.map(|stored| stored.message.id())
.collect();
assert!(
ids.contains(&"good-user") && ids.contains(&"good-assistant"),
"the readable conversation lands: {ids:?}",
);
assert!(
ids.contains(&"torn-session:4"),
"the readable fact lands: {ids:?}",
);
assert!(
!ids.contains(&"corrupt-entry") && !ids.contains(&"torn-session:5"),
"only the unreadable rows are missing: {ids:?}",
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn sqlite_resume_emits_a_v3_jsonl_pi_can_open() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let adapter = sqlite_only_adapter(&temp)?;
let (store, _) = ingest_into_temp_store(&temp, &adapter).await?;
let session = store
.get_session("sqlite-main-session")
.await?
.expect("sqlite session lands");
let files = PiCodingAgentFactory.serialize(&session, RestoreFidelity::Native)?;
let file = files.first().expect("one file per session");
assert_eq!(
file.actual_fidelity,
RestoreFidelity::Foreign,
"a database-origin session cannot replay into a format pi reads",
);
assert_eq!(
file.relative_path,
Path::new("sessions")
.join("--Users-user-Projects-harness-v2--")
.join("2026-08-06T00-00-32-000Z-pond-v3_sqlite-main-session.jsonl"),
"a reconstruction lands in pi's own directory layout under pond's own name",
);
let head: Value =
serde_json::from_str(std::str::from_utf8(&file.bytes)?.lines().next().unwrap())?;
assert_eq!(head.get("type"), Some(&json!("session")));
assert_eq!(head.get("version"), Some(&json!(3)));
let restored_root = temp.path().join("pi-home");
crate::adapter::write_restored_files(&restored_root, files)?;
let reread = Store::open_local(temp.path().join("reread-store")).await?;
let summary = ingest_adapter(
&reread,
&PiCodingAgentAdapter::new(restored_root.join("sessions")),
&crate::adapter::NoopOracle,
|_| {},
)
.await?;
assert_eq!(summary.dropped_events, 0);
let round_tripped = reread
.get_session("sqlite-main-session")
.await?
.expect("the resumed file reads back as the same session");
assert_eq!(round_tripped.session.project, session.session.project);
assert_eq!(round_tripped.session.created_at, session.session.created_at);
let conversational = |session: &crate::sessions::SessionWithMessages| {
session
.messages
.iter()
.filter(|stored| !matches!(stored.message, Message::System { .. }))
.count()
};
assert_eq!(conversational(&round_tripped), conversational(&session));
assert!(conversational(&session) > 0, "the fixture has conversation");
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn restore_refuses_to_overwrite_an_existing_file() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
let session = store
.get_session(V4_SESSION)
.await?
.expect("v4 fixture session lands");
let files = PiCodingAgentFactory.serialize(&session, RestoreFidelity::Native)?;
let root = temp.path().join("pi-home");
crate::adapter::write_restored_files(&root, files.clone())?;
let sibling = root.join("keep-me.txt");
std::fs::write(&sibling, b"untouched")?;
let error = crate::adapter::write_restored_files(&root, files)
.expect_err("a second restore into the same root must be refused");
assert!(
error.to_string().contains("refusing to overwrite"),
"the refusal names what it found: {error}",
);
assert_eq!(
std::fs::read(&sibling)?,
b"untouched",
"a refused restore leaves the directory exactly as it was",
);
Ok(())
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn restore_refuses_a_destination_occupied_by_a_dangling_symlink() -> anyhow::Result<()> {
let temp = TempDir::new()?;
let (store, _) =
ingest_into_temp_store(&temp, &PiCodingAgentAdapter::new(FIXTURES)).await?;
let session = store
.get_session(V4_SESSION)
.await?
.expect("v4 fixture session lands");
let files = PiCodingAgentFactory.serialize(&session, RestoreFidelity::Native)?;
let root = temp.path().join("pi-home");
let outside = temp.path().join("outside-the-root.jsonl");
let dest = crate::adapter::restore_destinations(&root, &files)?
.into_iter()
.next()
.expect("one file per session");
std::fs::create_dir_all(dest.parent().expect("dest has a parent"))?;
std::os::unix::fs::symlink(&outside, &dest)?;
let error = crate::adapter::write_restored_files(&root, files)
.expect_err("a destination occupied by a symlink must be refused");
assert!(
error.to_string().contains("refusing to overwrite"),
"the refusal names the occupied path: {error}",
);
assert!(
!outside.exists(),
"nothing may be written through the link, outside the restore root",
);
assert!(
std::fs::symlink_metadata(&dest).is_ok(),
"the pre-existing link is not this restore's to remove",
);
Ok(())
}
fn write_jsonl_file(path: &std::path::Path, records: &[Value]) -> anyhow::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(path, jsonl_bytes(NAME, records)?)?;
Ok(())
}
}