feat(coding-agent): enabled importing foreign sessions from claude and codex

- Implemented session stores and metadata converters to import Claude and Codex sessions into OMP.
- Added `--from-claude` and `--from-codex` CLI flags and `/resume` command arguments for foreign session resolution.
- Updated session selector components and controllers to support listing and picking external agent sessions.
- Added comprehensive unit tests and documentation covering foreign session import functionality.
This commit is contained in:
can1357
2026-07-30 00:28:40 +02:00
parent 47d3317d32
commit f11641d5a8
20 changed files with 1756 additions and 80 deletions
@@ -0,0 +1,426 @@
import type * as fsTypes from "node:fs";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import type {
AssistantMessage,
ImageContent,
TextContent,
ThinkingContent,
ToolCall,
ToolResultMessage,
Usage,
UserMessage,
} from "@oh-my-pi/pi-ai";
import { isRecord } from "@oh-my-pi/pi-utils";
import { collectForeignJsonRecords, type ForeignJsonRecord, readForeignJsonRecords } from "./foreign-session-jsonl";
import type { ForeignSessionInfo, ForeignSessionStore } from "./foreign-session-store";
import type { ModelChangeEntry, SessionMessageEntry } from "./session-entries";
import { SessionManager } from "./session-manager";
interface ClaudeHistoryMetadata {
firstMessage?: string;
cwd?: string;
created: number;
modified: number;
messageCount: number;
}
interface ConvertedMessage {
readonly message: UserMessage | AssistantMessage | ToolResultMessage;
readonly suffix: string;
}
const EMPTY_USAGE: Usage = {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
function stringField(record: Record<string, unknown>, key: string): string | undefined {
const value = record[key];
return typeof value === "string" && value.length > 0 ? value : undefined;
}
function numberField(record: Record<string, unknown>, key: string): number {
const value = record[key];
return typeof value === "number" && Number.isFinite(value) ? value : 0;
}
function timestampMs(value: unknown, fallback: number): number {
if (typeof value === "number" && Number.isFinite(value)) return value;
if (typeof value === "string") {
const parsed = Date.parse(value);
if (Number.isFinite(parsed)) return parsed;
}
return fallback;
}
function isoTimestamp(value: unknown, fallback: number): string {
return new Date(timestampMs(value, fallback)).toISOString();
}
function cleanPreview(value: string): string | undefined {
const cleaned = value.replace(/\s+/g, " ").trim();
return cleaned.length > 0 ? cleaned : undefined;
}
async function readHistoryIndex(file: string): Promise<Map<string, ClaudeHistoryMetadata>> {
const metadata = new Map<string, ClaudeHistoryMetadata>();
try {
for await (const { value } of readForeignJsonRecords(file)) {
const id = stringField(value, "sessionId") ?? stringField(value, "session_id");
const timestamp = timestampMs(value.timestamp ?? value.ts, 0);
if (!id || timestamp <= 0) continue;
const previous = metadata.get(id);
const rawText = stringField(value, "display") ?? stringField(value, "text");
const text = rawText ? cleanPreview(rawText) : undefined;
const cwd = stringField(value, "project");
if (!previous) {
metadata.set(id, { created: timestamp, modified: timestamp, firstMessage: text, cwd, messageCount: 1 });
continue;
}
if (timestamp < previous.created) {
previous.created = timestamp;
if (text) previous.firstMessage = text;
}
previous.modified = Math.max(previous.modified, timestamp);
if (!previous.firstMessage && text) previous.firstMessage = text;
previous.messageCount += 1;
if (!previous.cwd && cwd) previous.cwd = cwd;
}
} catch {
// History is an optional index; project files remain independently discoverable.
}
return metadata;
}
async function readRegisteredProjects(root: string): Promise<string[]> {
const config = path.join(path.dirname(root), ".claude.json");
try {
const parsed: unknown = await Bun.file(config).json();
if (!isRecord(parsed) || !isRecord(parsed.projects)) return [];
const projects: string[] = [];
for (const project in parsed.projects) {
if (path.isAbsolute(project)) projects.push(project);
}
return projects;
} catch {
return [];
}
}
function projectCwd(encoded: string, registered: readonly string[]): string {
const exact = registered.find(project => project.replaceAll(path.sep, "-") === encoded);
if (exact) return exact;
if (!encoded.startsWith("-")) return encoded;
return encoded.replaceAll("-", path.sep);
}
async function projectFiles(root: string): Promise<Array<{ file: string; cwd: string }>> {
const registered = await readRegisteredProjects(root);
const found: Array<{ file: string; cwd: string }> = [];
for (const containerName of ["projects", ".projects"]) {
const container = path.join(root, containerName);
const projects = await fs.readdir(container, { withFileTypes: true }).catch(() => []);
for (const project of projects) {
if (!project.isDirectory()) continue;
const directory = path.join(container, project.name);
const entries = await fs.readdir(directory, { withFileTypes: true }).catch(() => []);
const cwd = projectCwd(project.name, registered);
for (const entry of entries) {
if (entry.isFile() && entry.name.endsWith(".jsonl")) {
found.push({ file: path.join(directory, entry.name), cwd });
}
}
}
}
return found;
}
function imageContent(value: unknown): ImageContent | undefined {
if (!isRecord(value) || value.type !== "image" || !isRecord(value.source)) return undefined;
const data = stringField(value.source, "data");
const mimeType = stringField(value.source, "media_type");
return data && mimeType ? { type: "image", data, mimeType } : undefined;
}
function userContent(value: unknown): string | (TextContent | ImageContent)[] | undefined {
if (typeof value === "string") return value;
if (!Array.isArray(value)) return undefined;
const content: (TextContent | ImageContent)[] = [];
for (const block of value) {
if (!isRecord(block)) continue;
if (block.type === "text" && typeof block.text === "string") content.push({ type: "text", text: block.text });
else {
const image = imageContent(block);
if (image) content.push(image);
}
}
return content.length > 0 ? content : undefined;
}
function toolResultContent(value: unknown): (TextContent | ImageContent)[] {
if (typeof value === "string") return [{ type: "text", text: value }];
if (!Array.isArray(value)) return [];
const content: (TextContent | ImageContent)[] = [];
for (const block of value) {
if (!isRecord(block)) continue;
if (block.type === "text" && typeof block.text === "string") content.push({ type: "text", text: block.text });
else {
const image = imageContent(block);
if (image) content.push(image);
}
}
return content;
}
function assistantContent(
value: unknown,
toolNames: Map<string, string>,
): (TextContent | ThinkingContent | ToolCall)[] {
if (!Array.isArray(value)) return [];
const content: (TextContent | ThinkingContent | ToolCall)[] = [];
for (const block of value) {
if (!isRecord(block)) continue;
if (block.type === "text" && typeof block.text === "string") {
content.push({ type: "text", text: block.text });
} else if (block.type === "thinking" && typeof block.thinking === "string") {
const thinking: ThinkingContent = { type: "thinking", thinking: block.thinking };
if (typeof block.signature === "string") thinking.thinkingSignature = block.signature;
content.push(thinking);
} else if (block.type === "tool_use") {
const id = stringField(block, "id");
const name = stringField(block, "name");
if (!id || !name) continue;
const argumentsValue = isRecord(block.input) ? block.input : {};
content.push({ type: "toolCall", id, name, arguments: argumentsValue });
toolNames.set(id, name);
}
}
return content;
}
function claudeUsage(value: unknown): Usage {
if (!isRecord(value)) return EMPTY_USAGE;
const input = numberField(value, "input_tokens");
const output = numberField(value, "output_tokens");
const cacheRead = numberField(value, "cache_read_input_tokens");
const cacheWrite = numberField(value, "cache_creation_input_tokens");
return {
input,
output,
cacheRead,
cacheWrite,
totalTokens: input + output + cacheRead + cacheWrite,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
}
function stopReason(value: unknown): AssistantMessage["stopReason"] {
if (value === "tool_use") return "toolUse";
if (value === "max_tokens") return "length";
if (value === "end_turn" || value === "stop_sequence" || value === "pause_turn") return "stop";
return value == null ? "stop" : "error";
}
function convertRecord(
record: Record<string, unknown>,
fallbackTimestamp: number,
toolNames: Map<string, string>,
): ConvertedMessage[] {
if (
(record.type !== "user" && record.type !== "assistant") ||
record.isSidechain === true ||
record.isMeta === true
) {
return [];
}
if (!isRecord(record.message)) return [];
const timestamp = timestampMs(record.timestamp, fallbackTimestamp);
if (record.type === "assistant") {
const content = assistantContent(record.message.content, toolNames);
const errorMessage = stringField(record, "error");
if (content.length === 0 && !errorMessage) return [];
const model = stringField(record.message, "model") ?? "unknown";
const message: AssistantMessage = {
role: "assistant",
content,
api: "anthropic-messages",
provider: "anthropic",
model,
usage: claudeUsage(record.message.usage),
stopReason: stopReason(record.message.stop_reason),
timestamp,
};
const responseId = stringField(record.message, "id");
if (responseId) message.responseId = responseId;
if (errorMessage) message.errorMessage = errorMessage;
if (typeof record.apiErrorStatus === "number" && Number.isFinite(record.apiErrorStatus)) {
message.errorStatus = record.apiErrorStatus;
}
return [{ message, suffix: "message" }];
}
const rawContent = record.message.content;
if (Array.isArray(rawContent)) {
const results: ConvertedMessage[] = [];
for (let index = 0; index < rawContent.length; index += 1) {
const block = rawContent[index];
if (!isRecord(block) || block.type !== "tool_result") continue;
const toolCallId = stringField(block, "tool_use_id");
if (!toolCallId) continue;
const message: ToolResultMessage = {
role: "toolResult",
toolCallId,
toolName: toolNames.get(toolCallId) ?? "unknown",
content: toolResultContent(block.content),
isError: block.is_error === true,
timestamp,
};
results.push({ message, suffix: `tool-${index}` });
}
if (results.length > 0) return results;
}
const content = userContent(rawContent);
if (content === undefined || (typeof content === "string" && content.length === 0)) return [];
return [{ message: { role: "user", content, timestamp }, suffix: "message" }];
}
function uniqueEntryId(base: string, used: Set<string>): string {
let id = base;
let suffix = 2;
while (used.has(id)) {
id = `${base}-${suffix}`;
suffix += 1;
}
used.add(id);
return id;
}
/** Imports Claude Code JSONL sessions into non-persistent OMP session managers. */
export class ClaudeSessionStore implements ForeignSessionStore {
readonly source = "claude";
readonly #root: string;
/** Creates a store rooted at Claude's data directory, or at a fixture root when supplied. */
constructor(root: string = path.join(os.homedir(), ".claude")) {
this.#root = path.resolve(root);
}
/** Lists indexed Claude sessions without reading transcript bodies. */
async list(): Promise<ForeignSessionInfo[]> {
const [history, files] = await Promise.all([
readHistoryIndex(path.join(this.#root, "history.jsonl")),
projectFiles(this.#root),
]);
const sessions: ForeignSessionInfo[] = [];
for (const item of files) {
try {
const stats = await fs.stat(item.file);
const id = path.basename(item.file, ".jsonl");
const indexed = history.get(id);
const createdMs = indexed?.created ?? (stats.birthtimeMs || stats.ctimeMs || stats.mtimeMs);
const modifiedMs = Math.max(indexed?.modified ?? 0, stats.mtimeMs);
sessions.push({
source: this.source,
id,
path: item.file,
cwd: indexed?.cwd ?? item.cwd,
created: new Date(createdMs),
modified: new Date(modifiedMs),
firstMessage: indexed?.firstMessage,
messageCount: indexed?.messageCount,
});
} catch {
// Files may disappear while Claude rotates its session store.
}
}
return sessions.sort(
(left, right) => right.modified.getTime() - left.modified.getTime() || left.path.localeCompare(right.path),
);
}
/** Loads and converts a Claude transcript while preserving its source tree and timestamps. */
async load(info: ForeignSessionInfo): Promise<SessionManager> {
if (info.source !== this.source) throw new Error(`Cannot load ${info.source} session with ClaudeSessionStore`);
let records: ForeignJsonRecord[];
let stats: fsTypes.Stats;
try {
[records, stats] = await Promise.all([collectForeignJsonRecords(info.path), fs.stat(info.path)]);
} catch (error) {
const detail = error instanceof Error ? error.message : String(error);
throw new Error(`Unable to read Claude session ${info.id}: ${detail}`);
}
if (records.length === 0 && stats.size > 0)
throw new Error(`Claude session ${info.id} contains no readable records`);
const sourceParents = new Map<string, string | null>();
let sourceCwd: string | undefined;
let sourceTitle: string | undefined;
let aiTitle: string | undefined;
for (const { value } of records) {
const uuid = stringField(value, "uuid");
if (uuid) sourceParents.set(uuid, stringField(value, "parentUuid") ?? null);
if (!sourceCwd) sourceCwd = stringField(value, "cwd");
if (value.type === "custom-title") sourceTitle = stringField(value, "customTitle") ?? sourceTitle;
if (value.type === "ai-title") aiTitle = stringField(value, "aiTitle") ?? aiTitle;
}
const manager = SessionManager.inMemory(sourceCwd ?? info.cwd);
const sourceTails = new Map<string, string>();
const usedIds = new Set<string>();
const toolNames = new Map<string, string>();
let lastModel: string | undefined;
let synthetic = 0;
const resolveParent = (sourceId: string | undefined): string | null => {
const seen = new Set<string>();
let cursor = sourceId;
while (cursor && !seen.has(cursor)) {
seen.add(cursor);
const retained = sourceTails.get(cursor);
if (retained) return retained;
cursor = sourceParents.get(cursor) ?? undefined;
}
return null;
};
for (const { value, line } of records) {
const converted = convertRecord(value, stats.mtimeMs, toolNames);
if (converted.length === 0) continue;
const sourceUuid = stringField(value, "uuid") ?? `line-${line}`;
let parentId = resolveParent(stringField(value, "parentUuid"));
const timestamp = isoTimestamp(value.timestamp, stats.mtimeMs);
if (value.type === "assistant" && isRecord(value.message)) {
const model = stringField(value.message, "model");
if (model && model !== lastModel) {
const id = uniqueEntryId(`claude-${sourceUuid}-model`, usedIds);
const entry: ModelChangeEntry = {
type: "model_change",
id,
parentId,
timestamp,
model: `anthropic/${model}`,
};
manager.ingestReplicatedEntry(entry);
parentId = id;
lastModel = model;
}
}
for (const item of converted) {
synthetic += 1;
const id = uniqueEntryId(`claude-${sourceUuid}-${item.suffix}-${synthetic}`, usedIds);
const entry: SessionMessageEntry = { type: "message", id, parentId, timestamp, message: item.message };
manager.ingestReplicatedEntry(entry);
parentId = id;
}
if (parentId) sourceTails.set(sourceUuid, parentId);
}
const title = sourceTitle ?? aiTitle ?? info.title;
if (title) await manager.setSessionName(title, sourceTitle ? "user" : "auto");
return manager;
}
}
@@ -0,0 +1,620 @@
import { Database } from "bun:sqlite";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import type {
AssistantMessage,
ImageContent,
TextContent,
ThinkingContent,
ToolCall,
ToolResultMessage,
UserMessage,
} from "@oh-my-pi/pi-ai";
import { isRecord } from "@oh-my-pi/pi-utils";
import { readForeignJsonRecords } from "./foreign-session-jsonl";
import type { ForeignSessionInfo, ForeignSessionStore } from "./foreign-session-store";
import type { ModelChangeEntry, SessionEntry, SessionMessageEntry } from "./session-entries";
import { SessionManager } from "./session-manager";
interface CodexThreadRow {
id: string;
rollout_path: string;
created_at: number | null;
updated_at: number | null;
cwd: string;
title: string | null;
first_user_message: string | null;
}
interface CodexIndexRow {
id: string;
thread_name: string;
updated_at: string;
}
interface ConvertedRecord {
message?: UserMessage | AssistantMessage | ToolResultMessage;
followingMessage?: ToolResultMessage;
model?: string;
rollbackTurns?: number;
title?: string;
timestamp?: number;
}
const EMPTY_USAGE: AssistantMessage["usage"] = {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
function stringField(record: Record<string, unknown>, key: string): string | undefined {
const value = record[key];
return typeof value === "string" ? value : undefined;
}
function numberField(record: Record<string, unknown>, key: string): number | undefined {
const value = record[key];
return typeof value === "number" && Number.isFinite(value) ? value : undefined;
}
function timestampMillis(value: unknown, fallback: number): number {
if (typeof value === "number" && Number.isFinite(value)) return value < 10_000_000_000 ? value * 1000 : value;
if (typeof value === "string") {
const parsed = Date.parse(value);
if (Number.isFinite(parsed)) return parsed;
}
return fallback;
}
function dateFromEpoch(value: number | null, fallback: Date): Date {
if (value === null || !Number.isFinite(value)) return fallback;
return new Date(value < 10_000_000_000 ? value * 1000 : value);
}
function imageFromUrl(value: unknown, detail: unknown): ImageContent | undefined {
if (typeof value !== "string") return undefined;
const match = /^data:([^;,]+);base64,(.+)$/s.exec(value);
if (!match) return undefined;
const resolution =
detail === "auto" || detail === "low" || detail === "high" || detail === "original" ? detail : undefined;
return { type: "image", mimeType: match[1], data: match[2], detail: resolution };
}
function responseContent(value: unknown): Array<TextContent | ImageContent> {
if (!Array.isArray(value)) return [];
const content: Array<TextContent | ImageContent> = [];
for (const item of value) {
if (!isRecord(item)) continue;
const type = stringField(item, "type");
const text = stringField(item, "text");
if ((type === "input_text" || type === "output_text" || type === "text") && text !== undefined) {
content.push({ type: "text", text });
continue;
}
if (type === "input_image") {
const image = imageFromUrl(item.image_url, item.detail);
if (image) content.push(image);
}
}
return content;
}
function textFromContent(value: unknown): string {
return responseContent(value)
.filter((part): part is TextContent => part.type === "text")
.map(part => part.text)
.join("");
}
function toolArguments(value: unknown): Record<string, unknown> {
if (isRecord(value)) return value;
if (typeof value !== "string") return {};
try {
const parsed: unknown = JSON.parse(value);
if (isRecord(parsed)) return parsed;
} catch {
// Custom tools intentionally carry non-JSON input.
}
return { input: value };
}
function toolOutputContent(value: unknown): Array<TextContent | ImageContent> {
const structured = responseContent(value);
if (structured.length > 0) return structured;
if (typeof value === "string") return [{ type: "text", text: value }];
if (value === undefined) return [];
try {
return [{ type: "text", text: JSON.stringify(value) }];
} catch {
return [{ type: "text", text: String(value) }];
}
}
function reasoningContent(payload: Record<string, unknown>): ThinkingContent[] {
const parts: ThinkingContent[] = [];
for (const key of ["summary", "content"]) {
const value = payload[key];
if (!Array.isArray(value)) continue;
for (const item of value) {
if (!isRecord(item)) continue;
const text = stringField(item, "text");
if (text) parts.push({ type: "thinking", thinking: text });
}
}
return parts;
}
function assistantMessage(
content: AssistantMessage["content"],
timestamp: number,
model: string,
stopReason: AssistantMessage["stopReason"],
): AssistantMessage {
return {
role: "assistant",
content,
api: "openai-codex-responses",
provider: "openai-codex",
model,
usage: EMPTY_USAGE,
stopReason,
timestamp,
};
}
async function firstJsonRecord(filePath: string): Promise<Record<string, unknown> | undefined> {
for await (const { value } of readForeignJsonRecords(filePath)) return value;
return undefined;
}
async function readJsonLines(filePath: string): Promise<Record<string, unknown>[]> {
const records: Record<string, unknown>[] = [];
for await (const { value } of readForeignJsonRecords(filePath)) records.push(value);
return records;
}
async function rolloutFiles(directory: string): Promise<string[]> {
let entries: fs.Dirent[];
try {
entries = await fs.promises.readdir(directory, { withFileTypes: true });
} catch {
return [];
}
const files: string[] = [];
for (const entry of entries) {
const child = path.join(directory, entry.name);
if (entry.isDirectory()) files.push(...(await rolloutFiles(child)));
else if (entry.isFile() && entry.name.endsWith(".jsonl")) files.push(child);
}
return files;
}
function rolloutId(filePath: string): string {
const match = /([0-9a-f]{8}-[0-9a-f-]{27})\.jsonl$/i.exec(filePath);
return match?.[1] ?? path.basename(filePath, ".jsonl");
}
async function stateDatabasePath(root: string): Promise<string | undefined> {
const names = await fs.promises.readdir(root).catch(() => []);
return names
.map(name => ({ name, version: /^state_(\d+)\.sqlite$/.exec(name) }))
.filter(item => item.version !== null)
.sort((left, right) => Number(right.version?.[1]) - Number(left.version?.[1]))
.map(item => path.join(root, item.name))
.at(0);
}
async function loadIndex(root: string): Promise<Map<string, CodexIndexRow>> {
const index = new Map<string, CodexIndexRow>();
const filePath = path.join(root, "session_index.jsonl");
let records: Record<string, unknown>[];
try {
records = await readJsonLines(filePath);
} catch {
return index;
}
for (const record of records) {
const id = stringField(record, "id");
const threadName = stringField(record, "thread_name");
const updatedAt = stringField(record, "updated_at");
if (id && threadName && updatedAt) index.set(id, { id, thread_name: threadName, updated_at: updatedAt });
}
return index;
}
function convertedResponseItem(
payload: Record<string, unknown>,
timestamp: number,
model: string,
toolNames: Map<string, string>,
): ConvertedRecord | undefined {
const type = stringField(payload, "type");
if (type === "message") {
const role = stringField(payload, "role");
const content = responseContent(payload.content);
if (content.length === 0) return undefined;
if (role === "user") return { message: { role: "user", content, timestamp } };
if (role === "assistant") return { message: assistantMessage(content, timestamp, model, "stop") };
return undefined;
}
if (type === "reasoning") {
const content = reasoningContent(payload);
return content.length > 0 ? { message: assistantMessage(content, timestamp, model, "stop") } : undefined;
}
if (type === "function_call" || type === "custom_tool_call") {
const callId = stringField(payload, "call_id") ?? stringField(payload, "id");
const name = stringField(payload, "name");
if (!callId || !name) return undefined;
toolNames.set(callId, name);
const call: ToolCall = {
type: "toolCall",
id: callId,
name,
arguments: toolArguments(type === "custom_tool_call" ? payload.input : payload.arguments),
customWireName: type === "custom_tool_call" ? name : undefined,
};
return { message: assistantMessage([call], timestamp, model, "toolUse") };
}
if (type === "function_call_output" || type === "custom_tool_call_output") {
const callId = stringField(payload, "call_id");
if (!callId) return undefined;
return {
message: {
role: "toolResult",
toolCallId: callId,
toolName: toolNames.get(callId) ?? "unknown",
content: toolOutputContent(payload.output),
isError: false,
timestamp,
},
};
}
if (type === "web_search_call" || type === "tool_search_call") {
const callId = stringField(payload, "call_id") ?? stringField(payload, "id");
if (!callId) return undefined;
const name = type === "web_search_call" ? "web_search" : "tool_search";
toolNames.set(callId, name);
const input = type === "web_search_call" ? payload.action : payload.arguments;
const call: ToolCall = { type: "toolCall", id: callId, name, arguments: toolArguments(input) };
return { message: assistantMessage([call], timestamp, model, "toolUse") };
}
if (type === "tool_search_output") {
const callId = stringField(payload, "call_id");
if (!callId) return undefined;
return {
message: {
role: "toolResult",
toolCallId: callId,
toolName: toolNames.get(callId) ?? "tool_search",
content: toolOutputContent(payload.tools),
isError: payload.status === "failed",
timestamp,
},
};
}
return undefined;
}
function convertedEvent(
payload: Record<string, unknown>,
timestamp: number,
model: string,
canonicalUserText: Set<string>,
canonicalAssistantText: Set<string>,
canonicalToolCalls: Set<string>,
toolNames: Map<string, string>,
): ConvertedRecord | undefined {
const type = stringField(payload, "type");
if (type === "user_message") {
const text = stringField(payload, "message");
if (!text || canonicalUserText.has(text)) return undefined;
return { message: { role: "user", content: text, timestamp } };
}
if (type === "agent_message") {
const text = stringField(payload, "message");
if (!text || canonicalAssistantText.has(text)) return undefined;
return { message: assistantMessage([{ type: "text", text }], timestamp, model, "stop") };
}
if (type === "agent_reasoning") {
const text = stringField(payload, "text");
if (!text || canonicalAssistantText.has(text)) return undefined;
return { message: assistantMessage([{ type: "thinking", thinking: text }], timestamp, model, "stop") };
}
if (type === "dynamic_tool_call_request") {
const callId = stringField(payload, "callId") ?? stringField(payload, "call_id");
const name = stringField(payload, "tool");
if (!callId || !name || canonicalToolCalls.has(callId)) return undefined;
toolNames.set(callId, name);
const call: ToolCall = { type: "toolCall", id: callId, name, arguments: toolArguments(payload.arguments) };
return { message: assistantMessage([call], timestamp, model, "toolUse") };
}
if (type === "dynamic_tool_call_response") {
const callId = stringField(payload, "call_id") ?? stringField(payload, "callId");
if (!callId || canonicalToolCalls.has(callId)) return undefined;
const error = stringField(payload, "error");
return {
message: {
role: "toolResult",
toolCallId: callId,
toolName: toolNames.get(callId) ?? stringField(payload, "tool") ?? "unknown",
content: error ? [{ type: "text", text: error }] : toolOutputContent(payload.content_items),
isError: error !== undefined || payload.success === false,
timestamp,
},
};
}
if (type === "web_search_end") {
const callId = stringField(payload, "call_id");
if (!callId) return undefined;
const name = toolNames.get(callId) ?? "web_search";
const result: ToolResultMessage = {
role: "toolResult",
toolCallId: callId,
toolName: name,
content: toolOutputContent(payload.results ?? payload.query),
isError: false,
timestamp,
};
if (canonicalToolCalls.has(callId) || toolNames.has(callId)) return { message: result };
toolNames.set(callId, name);
const call: ToolCall = {
type: "toolCall",
id: callId,
name,
arguments: toolArguments(payload.action ?? payload.query),
};
return {
message: assistantMessage([call], timestamp, model, "toolUse"),
followingMessage: result,
};
}
if (type === "mcp_tool_call_end") {
const callId = stringField(payload, "call_id");
if (!callId || canonicalToolCalls.has(callId) || !isRecord(payload.invocation)) return undefined;
const server = stringField(payload.invocation, "server");
const tool = stringField(payload.invocation, "tool");
if (!server || !tool) return undefined;
const name = `${server}/${tool}`;
toolNames.set(callId, name);
const call: ToolCall = {
type: "toolCall",
id: callId,
name,
arguments: toolArguments(payload.invocation.arguments),
};
const result = isRecord(payload.result) && isRecord(payload.result.Ok) ? payload.result.Ok : payload.result;
const error = isRecord(payload.result) && typeof payload.result.Err === "string" ? payload.result.Err : undefined;
return {
message: assistantMessage([call], timestamp, model, "toolUse"),
followingMessage: {
role: "toolResult",
toolCallId: callId,
toolName: name,
content: error
? [{ type: "text", text: error }]
: toolOutputContent(isRecord(result) ? result.content : result),
isError: error !== undefined || (isRecord(result) && result.isError === true),
timestamp,
},
};
}
if (type === "thread_name_updated") return { title: stringField(payload, "thread_name") };
if (type === "thread_rolled_back") return { rollbackTurns: numberField(payload, "num_turns") ?? 0 };
return undefined;
}
function canonicalTexts(records: Record<string, unknown>[]): {
users: Set<string>;
assistants: Set<string>;
toolCalls: Set<string>;
} {
const users = new Set<string>();
const assistants = new Set<string>();
const toolCalls = new Set<string>();
for (const record of records) {
if (record.type !== "response_item" || !isRecord(record.payload)) continue;
const callId = stringField(record.payload, "call_id") ?? stringField(record.payload, "id");
if (callId && typeof record.payload.type === "string" && record.payload.type.includes("call"))
toolCalls.add(callId);
if (record.payload.type === "message") {
const text = textFromContent(record.payload.content);
if (!text) continue;
if (record.payload.role === "user") users.add(text);
else if (record.payload.role === "assistant") assistants.add(text);
} else if (record.payload.type === "reasoning") {
for (const part of reasoningContent(record.payload)) assistants.add(part.thinking);
}
}
return { users, assistants, toolCalls };
}
function rollback(records: ConvertedRecord[], turns: number): void {
for (let remaining = turns; remaining > 0; remaining--) {
let userIndex = -1;
for (let index = records.length - 1; index >= 0; index--) {
if (records[index].message?.role === "user") {
userIndex = index;
break;
}
}
if (userIndex < 0) return;
records.splice(userIndex);
}
}
/** Imports locally stored OpenAI Codex sessions into OMP's in-memory session format. */
export class CodexSessionStore implements ForeignSessionStore {
/** Foreign-session source discriminator. */
readonly source = "codex";
readonly #root: string;
/** Uses the supplied Codex data root, or ~/.codex by default. */
constructor(rootDirectory: string = path.join(os.homedir(), ".codex")) {
this.#root = path.resolve(rootDirectory);
}
/** Lists Codex sessions from its state index without reading transcript bodies. */
async list(): Promise<ForeignSessionInfo[]> {
const databasePath = await stateDatabasePath(this.#root);
if (databasePath) {
try {
const database = new Database(databasePath, { readonly: true });
try {
const rows = database
.query<CodexThreadRow, []>(
"SELECT id, rollout_path, created_at, updated_at, cwd, title, first_user_message FROM threads",
)
.all();
const sessions: ForeignSessionInfo[] = [];
for (const row of rows) {
if (!row.id || !row.rollout_path || !row.cwd) continue;
const rolloutPath = path.isAbsolute(row.rollout_path)
? row.rollout_path
: path.join(this.#root, row.rollout_path);
const modified = dateFromEpoch(row.updated_at, new Date(0));
const created = dateFromEpoch(row.created_at, modified);
sessions.push({
source: "codex",
id: row.id,
path: rolloutPath,
cwd: row.cwd,
title: row.title ?? undefined,
created,
modified,
firstMessage: row.first_user_message ?? undefined,
});
}
sessions.sort(
(left, right) =>
right.modified.getTime() - left.modified.getTime() || left.id.localeCompare(right.id),
);
if (sessions.length > 0) return sessions;
} finally {
database.close();
}
} catch {
// Older Codex state databases fall back to the rollout metadata path.
}
}
const index = await loadIndex(this.#root);
const roots = ["sessions", ".sessions", "archived_sessions"].map(name => path.join(this.#root, name));
const files = (await Promise.all(roots.map(rolloutFiles))).flat();
const sessions: ForeignSessionInfo[] = [];
for (const filePath of files) {
const first = await firstJsonRecord(filePath);
if (first?.type !== "session_meta" || !isRecord(first.payload)) continue;
const id = stringField(first.payload, "id") ?? rolloutId(filePath);
const cwd = stringField(first.payload, "cwd");
if (!cwd) continue;
const stat = await fs.promises.stat(filePath);
const indexed = index.get(id);
const sourceCreated = stringField(first.payload, "timestamp");
const created = new Date(timestampMillis(sourceCreated, stat.birthtimeMs));
const modified = indexed ? new Date(timestampMillis(indexed.updated_at, stat.mtimeMs)) : stat.mtime;
sessions.push({
source: "codex",
id,
path: filePath,
cwd,
title: indexed?.thread_name,
created,
modified,
});
}
sessions.sort(
(left, right) => right.modified.getTime() - left.modified.getTime() || left.id.localeCompare(right.id),
);
return sessions;
}
/** Converts one Codex rollout into a non-persistent OMP session. */
async load(info: ForeignSessionInfo): Promise<SessionManager> {
if (info.source !== "codex") throw new Error(`Cannot load ${info.source} session with CodexSessionStore`);
let records: Record<string, unknown>[];
try {
records = await readJsonLines(info.path);
} catch (error) {
throw new Error(`Unable to read Codex session ${info.id} at ${info.path}`, { cause: error });
}
if (records.length === 0) throw new Error(`Codex session ${info.id} at ${info.path} is empty or malformed`);
const metadata = records.find(record => record.type === "session_meta" && isRecord(record.payload));
const cwd =
metadata && isRecord(metadata.payload) ? (stringField(metadata.payload, "cwd") ?? info.cwd) : info.cwd;
const manager = SessionManager.inMemory(cwd);
const canonical = canonicalTexts(records);
const converted: ConvertedRecord[] = [];
const toolNames = new Map<string, string>();
let model = "codex";
let fallbackTimestamp = info.created.getTime();
let title = info.title;
for (const record of records) {
const timestamp = timestampMillis(record.timestamp, fallbackTimestamp);
fallbackTimestamp = Math.max(fallbackTimestamp + 1, timestamp);
if (!isRecord(record.payload)) continue;
let item: ConvertedRecord | undefined;
if (record.type === "turn_context") {
const nextModel = stringField(record.payload, "model");
if (nextModel && nextModel !== model) {
model = nextModel;
item = { model, timestamp };
}
} else if (record.type === "response_item") {
item = convertedResponseItem(record.payload, timestamp, model, toolNames);
} else if (record.type === "event_msg") {
item = convertedEvent(
record.payload,
timestamp,
model,
canonical.users,
canonical.assistants,
canonical.toolCalls,
toolNames,
);
}
if (!item) continue;
if (item.rollbackTurns) rollback(converted, item.rollbackTurns);
else converted.push(item);
if (item.followingMessage) converted.push({ message: item.followingMessage });
if (item.title) title = item.title;
}
let parentId: string | null = null;
let ordinal = 0;
for (const item of converted) {
const id = `codex-${(++ordinal).toString(36)}`;
const timestamp = new Date(item.message?.timestamp ?? item.timestamp ?? fallbackTimestamp).toISOString();
let entry: SessionEntry | undefined;
if (item.message) {
const messageEntry: SessionMessageEntry = {
type: "message",
id,
parentId,
timestamp,
message: item.message,
};
entry = messageEntry;
} else if (item.model) {
const modelEntry: ModelChangeEntry = {
type: "model_change",
id,
parentId,
timestamp,
model: `openai-codex/${item.model}`,
};
entry = modelEntry;
}
if (!entry) continue;
manager.ingestReplicatedEntry(entry);
parentId = id;
}
if (title) await manager.setSessionName(title, "auto", "codex-import");
return manager;
}
}
@@ -0,0 +1,52 @@
import { directoryExists } from "@oh-my-pi/pi-utils";
import { ClaudeSessionStore } from "./claude-session-store";
import { CodexSessionStore } from "./codex-session-store";
import type { ForeignSessionInfo, ForeignSessionSource, ForeignSessionStore } from "./foreign-session-store";
import type { SessionInfo } from "./session-listing";
import type { SessionManager } from "./session-manager";
/** Construct the importer for a supported foreign session source. */
export function createForeignSessionStore(source: ForeignSessionSource): ForeignSessionStore {
return source === "claude" ? new ClaudeSessionStore() : new CodexSessionStore();
}
/** Display name for a supported foreign session source. */
export function foreignSessionSourceName(source: ForeignSessionSource): string {
return source === "claude" ? "Claude" : "Codex";
}
/** Convert lightweight foreign metadata for the existing session picker. */
export function foreignSessionInfoToSessionInfo(info: ForeignSessionInfo): SessionInfo {
const firstMessage = info.firstMessage ?? "(no messages)";
return {
path: info.path,
id: info.id,
cwd: info.cwd,
title: info.title,
created: info.created,
modified: info.modified,
messageCount: info.messageCount ?? 0,
size: 0,
firstMessage,
allMessagesText: firstMessage,
};
}
/** Import and persist one foreign session under a fresh OMP session identity. */
export async function persistForeignSession(
store: ForeignSessionStore,
info: ForeignSessionInfo,
options?: { fallbackCwd?: string; sessionDir?: string; suppressBreadcrumb?: boolean },
): Promise<SessionManager> {
const imported = await store.load(info);
imported.appendCustomEntry("foreign_session_import", {
source: info.source,
sourceId: info.id,
sourcePath: info.path,
sourceCwd: info.cwd,
});
if (options?.fallbackCwd && !(await directoryExists(imported.getCwd()))) {
await imported.moveTo(options.fallbackCwd);
}
return await imported.persistCopy(options);
}
@@ -0,0 +1,29 @@
import { isRecord, readLines } from "@oh-my-pi/pi-utils";
/** One readable object record from a foreign JSONL transcript. */
export interface ForeignJsonRecord {
readonly value: Record<string, unknown>;
readonly line: number;
}
/** Stream valid object records while tolerating malformed or truncated lines. */
export async function* readForeignJsonRecords(filePath: string): AsyncGenerator<ForeignJsonRecord> {
const decoder = new TextDecoder();
let line = 0;
for await (const bytes of readLines(Bun.file(filePath).stream())) {
line += 1;
try {
const value: unknown = JSON.parse(decoder.decode(bytes));
if (isRecord(value)) yield { value, line };
} catch {
// A partially-written line must not hide the readable transcript prefix.
}
}
}
/** Read every valid object record from a foreign JSONL transcript. */
export async function collectForeignJsonRecords(filePath: string): Promise<ForeignJsonRecord[]> {
const records: ForeignJsonRecord[] = [];
for await (const record of readForeignJsonRecords(filePath)) records.push(record);
return records;
}
@@ -0,0 +1,26 @@
import type { SessionManager } from "./session-manager";
/** External coding-agent session source supported by OMP imports. */
export type ForeignSessionSource = "claude" | "codex";
/** Lightweight source metadata used to choose a foreign session before loading its transcript. */
export interface ForeignSessionInfo {
readonly source: ForeignSessionSource;
readonly id: string;
readonly path: string;
readonly cwd: string;
readonly title?: string;
readonly created: Date;
readonly modified: Date;
readonly messageCount?: number;
readonly firstMessage?: string;
}
/** Lists and converts sessions owned by another coding agent. */
export interface ForeignSessionStore {
readonly source: ForeignSessionSource;
/** Lists source sessions without parsing complete transcripts. */
list(): Promise<ForeignSessionInfo[]>;
/** Converts one source session into a non-persistent OMP session. */
load(session: ForeignSessionInfo): Promise<SessionManager>;
}
@@ -1420,6 +1420,30 @@ export class SessionManager {
await this.#rewriteAtomically();
}
/** Persist this session's transcript as a newly identified OMP session. */
async persistCopy(
options?: { sessionDir?: string; suppressBreadcrumb?: boolean },
storage: SessionStorage = new FileSessionStorage(),
): Promise<SessionManager> {
const sessionDir = options?.sessionDir ?? SessionManager.getDefaultSessionDir(this.#cwd, undefined, storage);
const manager = new SessionManager(this.#cwd, sessionDir, true, storage);
manager.#suppressBreadcrumb = options?.suppressBreadcrumb === true;
manager.#resetToNewSession();
manager.#sessionName = this.#sessionName;
manager.#titleSource = this.#titleSource;
manager.#titleUpdatedAt = this.#titleUpdatedAt;
manager.#header.title = this.#sessionName;
manager.#header.titleSource = this.#titleSource;
manager.#additionalDirectories = [...this.#additionalDirectories];
manager.#header.additionalDirectories =
manager.#additionalDirectories.length > 0 ? [...manager.#additionalDirectories] : undefined;
manager.#entries = structuredClone(this.#entries);
manager.#index.rebuild(manager.#entries);
manager.#forceFileCreation = true;
await manager.#rewriteAtomically();
return manager;
}
/**
* Stage a synchronous group of entry appends and publish the resulting full
* journal with one atomic replace. A failed publish removes only the staged