Files
platform/src/servers/api/pi/storage.ts
T
pastilhasandClaude Opus 4.6 c443fe0fe2 pi session context persistence, fix new session button, add clear all sessions
- handle pi process death by nulling piProcess ref so next message respawns
- replay conversation history on respawn via --session flag
- fix user message JSONL format to array for pi compatibility
- fix container sessions mount to match storage path
- fix ChatPanelWrapper: key on inner component so usePiChat resets on new session
- add bulk delete sessions endpoint and clear all button in chat header

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-02-27 00:04:14 +00:00

830 lines
24 KiB
TypeScript

import * as fs from 'fs/promises';
import * as path from 'path';
import { randomUUID } from 'crypto';
import type {
SessionMeta,
Message,
GroupMeta,
MessageCost,
JnlSessionHeader,
JnlEntry,
JnlMessageEntry,
JnlSessionInfoEntry,
JnlTextContent,
JnlToolCall,
SessionIndex,
SessionIndexEntry,
} from './types';
// ── Path helpers ───────────────────────────────────────────────────────
const SESSIONS_ROOT = '.pi/agent/sessions';
const OFFICER_DIR = '.officer';
const INDEX_FILE = 'index.json';
const GROUPS_DIR = 'groups';
function getSessionsDir(baseCwd: string): string {
return path.join(baseCwd, SESSIONS_ROOT);
}
function getOfficerDir(baseCwd: string): string {
return path.join(baseCwd, SESSIONS_ROOT, OFFICER_DIR);
}
function getIndexPath(baseCwd: string): string {
return path.join(getOfficerDir(baseCwd), INDEX_FILE);
}
function getGroupsDir(baseCwd: string): string {
return path.join(getOfficerDir(baseCwd), GROUPS_DIR);
}
function getGroupPath(baseCwd: string, groupSlug: string): string {
return path.join(getGroupsDir(baseCwd), `${groupSlug}.json`);
}
export function encodeCwdDir(cwd: string): string {
return cwd.replace(/\//g, '-');
}
function buildJnlFilename(createdAt: number, sessionId: string): string {
return `${createdAt}_${sessionId}.jsonl`;
}
function buildRelativePath(cwd: string, createdAt: number, sessionId: string): string {
return path.join(encodeCwdDir(cwd), buildJnlFilename(createdAt, sessionId));
}
// ── Hex ID generator ───────────────────────────────────────────────────
export function generateHexId(): string {
const bytes = new Uint8Array(4);
crypto.getRandomValues(bytes);
return Array.from(bytes)
.map((b) => b.toString(16).padStart(2, '0'))
.join('');
}
// ── JSONL serialization ────────────────────────────────────────────────
function serializeJnlFile(header: JnlSessionHeader, entries: JnlEntry[]): string {
const lines = [JSON.stringify(header)];
for (const entry of entries) {
lines.push(JSON.stringify(entry));
}
return lines.join('\n') + '\n';
}
function parseJnlFile(content: string): { header: JnlSessionHeader; entries: JnlEntry[] } {
const lines = content.trim().split('\n');
if (lines.length === 0) {
throw new Error('Empty JSONL file');
}
const header = JSON.parse(lines[0]!) as JnlSessionHeader;
const entries: JnlEntry[] = [];
for (let i = 1; i < lines.length; i++) {
const line = lines[i]!.trim();
if (!line) continue;
entries.push(JSON.parse(line) as JnlEntry);
}
return { header, entries };
}
// ── Index management ───────────────────────────────────────────────────
async function loadIndex(baseCwd: string): Promise<SessionIndex> {
try {
const indexPath = getIndexPath(baseCwd);
const content = await fs.readFile(indexPath, 'utf-8');
return JSON.parse(content) as SessionIndex;
} catch {
return {};
}
}
async function saveIndex(baseCwd: string, index: SessionIndex): Promise<void> {
const officerDir = getOfficerDir(baseCwd);
await fs.mkdir(officerDir, { recursive: true });
const indexPath = getIndexPath(baseCwd);
await fs.writeFile(indexPath, JSON.stringify(index, null, 2));
}
// ── Message ↔ JSONL conversion ─────────────────────────────────────────
export function messagesToJnlEntries(messages: Message[], meta: SessionMeta): JnlEntry[] {
const entries: JnlEntry[] = [];
let prevId: string | undefined;
let i = 0;
while (i < messages.length) {
const msg = messages[i]!;
if (msg.role === 'user') {
const id = generateHexId();
const entry: JnlMessageEntry = {
type: 'message',
id,
parentId: prevId,
timestamp: new Date(msg.timestamp).toISOString(),
message: { role: 'user', content: [{ type: 'text' as const, text: msg.text ?? '' }] },
};
entries.push(entry);
prevId = id;
i++;
} else if (msg.role === 'assistant') {
const content: Array<JnlTextContent | JnlToolCall> = [];
if (msg.text) {
content.push({ type: 'text', text: msg.text });
}
// Collect following tool messages
const toolMessages: Message[] = [];
let j = i + 1;
while (j < messages.length && messages[j]!.role === 'tool') {
const toolMsg = messages[j]!;
content.push({
type: 'tool_use',
id: toolMsg.toolCallId ?? generateHexId(),
name: toolMsg.toolName ?? 'unknown',
input: toolMsg.toolInput ?? {},
});
toolMessages.push(toolMsg);
j++;
}
const assistId = generateHexId();
const assistEntry: JnlMessageEntry = {
type: 'message',
id: assistId,
parentId: prevId,
timestamp: new Date(msg.timestamp).toISOString(),
message: { role: 'assistant', content },
};
entries.push(assistEntry);
prevId = assistId;
// Emit toolResult entries for tools with output
for (const toolMsg of toolMessages) {
if (toolMsg.output !== undefined) {
const trId = generateHexId();
const trEntry: JnlMessageEntry = {
type: 'message',
id: trId,
parentId: prevId,
timestamp: new Date(toolMsg.timestamp).toISOString(),
message: {
role: 'toolResult',
toolCallId: toolMsg.toolCallId ?? '',
toolName: toolMsg.toolName ?? 'unknown',
content: [{ type: 'text', text: toolMsg.output }],
isError: toolMsg.isError,
},
};
entries.push(trEntry);
prevId = trId;
}
}
i = j;
} else if (msg.role === 'tool') {
// Standalone tool without preceding assistant (edge case)
const content: Array<JnlTextContent | JnlToolCall> = [
{
type: 'tool_use',
id: msg.toolCallId ?? generateHexId(),
name: msg.toolName ?? 'unknown',
input: msg.toolInput ?? {},
},
];
const assistId = generateHexId();
const assistEntry: JnlMessageEntry = {
type: 'message',
id: assistId,
parentId: prevId,
timestamp: new Date(msg.timestamp).toISOString(),
message: { role: 'assistant', content },
};
entries.push(assistEntry);
prevId = assistId;
if (msg.output !== undefined) {
const trId = generateHexId();
const trEntry: JnlMessageEntry = {
type: 'message',
id: trId,
parentId: prevId,
timestamp: new Date(msg.timestamp).toISOString(),
message: {
role: 'toolResult',
toolCallId: msg.toolCallId ?? '',
toolName: msg.toolName ?? 'unknown',
content: [{ type: 'text', text: msg.output }],
isError: msg.isError,
},
};
entries.push(trEntry);
prevId = trId;
}
i++;
} else {
i++;
}
}
// Append session_info entry
const infoId = generateHexId();
const infoEntry: JnlSessionInfoEntry = {
type: 'session_info',
id: infoId,
parentId: prevId,
timestamp: new Date(meta.updatedAt).toISOString(),
name: meta.title,
officer: {
cost: meta.cost,
model: meta.model,
groupSlug: meta.groupSlug ?? null,
messageCount: meta.messageCount,
createdAt: meta.createdAt,
updatedAt: meta.updatedAt,
context: meta.context,
contextId: meta.contextId,
},
};
entries.push(infoEntry);
return entries;
}
type ParsedSession = {
messages: Message[];
sessionInfo: JnlSessionInfoEntry | null;
};
export function jnlEntriesToMessages(entries: JnlEntry[]): ParsedSession {
const messages: Message[] = [];
let sessionInfo: JnlSessionInfoEntry | null = null;
for (const entry of entries) {
if (entry.type === 'session_info') {
sessionInfo = entry as JnlSessionInfoEntry;
continue;
}
if (entry.type !== 'message') continue;
const msgEntry = entry as JnlMessageEntry;
const ts = new Date(entry.timestamp).getTime();
if (msgEntry.message.role === 'user') {
const rawContent = msgEntry.message.content;
const text =
typeof rawContent === 'string'
? rawContent
: (rawContent as Array<JnlTextContent>).map((c) => c.text).join('\n');
messages.push({
id: entry.id,
timestamp: ts,
role: 'user',
text,
});
} else if (msgEntry.message.role === 'assistant') {
const contentBlocks = msgEntry.message.content as Array<JnlTextContent | JnlToolCall>;
let text = '';
const toolCalls: JnlToolCall[] = [];
for (const block of contentBlocks) {
if (block.type === 'text') {
text += (text ? '\n' : '') + block.text;
} else if (block.type === 'tool_use') {
toolCalls.push(block);
}
}
if (text) {
messages.push({
id: entry.id,
timestamp: ts,
role: 'assistant',
text,
});
}
// Create tool messages from tool_use blocks
for (const tc of toolCalls) {
messages.push({
id: randomUUID(),
timestamp: ts,
role: 'tool',
toolCallId: tc.id,
toolName: tc.name,
toolInput: tc.input,
});
}
} else if (msgEntry.message.role === 'toolResult') {
const trMsg = msgEntry.message as {
role: 'toolResult';
toolCallId: string;
toolName: string;
content: Array<JnlTextContent>;
isError?: boolean;
};
const outputText = trMsg.content.map((c) => c.text).join('\n');
// Find matching tool message and update it with output
for (let k = messages.length - 1; k >= 0; k--) {
const m = messages[k]!;
if (m.role === 'tool' && m.toolCallId === trMsg.toolCallId) {
m.output = outputText;
m.isError = trMsg.isError;
break;
}
}
}
}
return { messages, sessionInfo };
}
// ── Helper: build SessionMeta from index entry ─────────────────────────
function indexEntryToMeta(sessionId: string, entry: SessionIndexEntry): SessionMeta {
return {
id: sessionId,
title: entry.title,
model: entry.model,
cwd: entry.cwd,
createdAt: entry.createdAt,
updatedAt: entry.updatedAt,
messageCount: entry.messageCount,
cost: entry.cost,
groupSlug: entry.groupSlug ?? null,
context: entry.context,
contextId: entry.contextId,
};
}
function metaToIndexEntry(meta: SessionMeta, relFile: string): SessionIndexEntry {
return {
file: relFile,
title: meta.title,
model: meta.model,
cwd: meta.cwd,
createdAt: meta.createdAt,
updatedAt: meta.updatedAt,
messageCount: meta.messageCount,
cost: meta.cost,
groupSlug: meta.groupSlug ?? null,
context: meta.context,
contextId: meta.contextId,
};
}
// ── Path resolution ─────────────────────────────────────────────────────
export async function getSessionFilePath(baseCwd: string, sessionId: string): Promise<string | null> {
const index = await loadIndex(baseCwd);
const entry = index[sessionId];
if (!entry) return null;
return path.join(getSessionsDir(baseCwd), entry.file);
}
// ── Session CRUD ───────────────────────────────────────────────────────
export async function saveSession(
baseCwd: string,
sessionId: string,
meta: SessionMeta,
messages: Message[],
): Promise<void> {
const index = await loadIndex(baseCwd);
// Reuse existing file path or create new one
const existing = index[sessionId];
const relFile = existing?.file ?? buildRelativePath(meta.cwd, meta.createdAt, sessionId);
const absFile = path.join(getSessionsDir(baseCwd), relFile);
await fs.mkdir(path.dirname(absFile), { recursive: true });
const header: JnlSessionHeader = {
type: 'session',
version: 3,
id: sessionId,
timestamp: new Date(meta.createdAt).toISOString(),
cwd: meta.cwd,
};
const entries = messagesToJnlEntries(messages, meta);
const content = serializeJnlFile(header, entries);
await fs.writeFile(absFile, content);
// Update index
index[sessionId] = metaToIndexEntry(meta, relFile);
await saveIndex(baseCwd, index);
}
export async function loadSession(
baseCwd: string,
sessionId: string,
_groupSlug?: string | null,
): Promise<{ meta: SessionMeta; messages: Message[] }> {
const index = await loadIndex(baseCwd);
const entry = index[sessionId];
if (!entry) {
throw new Error(`Session not found: ${sessionId}`);
}
const absFile = path.join(getSessionsDir(baseCwd), entry.file);
const raw = await fs.readFile(absFile, 'utf-8');
const { entries } = parseJnlFile(raw);
const { messages, sessionInfo } = jnlEntriesToMessages(entries);
// Reconstruct meta from index entry (authoritative) enriched by session_info
const meta: SessionMeta = indexEntryToMeta(sessionId, entry);
// If session_info has officer data, prefer those for fields that may differ
if (sessionInfo?.officer) {
meta.cost = sessionInfo.officer.cost;
meta.messageCount = sessionInfo.officer.messageCount;
}
if (sessionInfo?.name) {
meta.title = sessionInfo.name;
}
return { meta, messages };
}
export async function sessionExists(
baseCwd: string,
sessionId: string,
_groupSlug?: string | null,
): Promise<boolean> {
const index = await loadIndex(baseCwd);
return sessionId in index;
}
export async function updateSessionMeta(
baseCwd: string,
sessionId: string,
updates: Partial<SessionMeta>,
_groupSlug?: string | null,
): Promise<SessionMeta> {
const { meta, messages } = await loadSession(baseCwd, sessionId);
const updatedMeta: SessionMeta = {
...meta,
...updates,
updatedAt: Date.now(),
};
await saveSession(baseCwd, sessionId, updatedMeta, messages);
return updatedMeta;
}
export async function deleteSession(
baseCwd: string,
sessionId: string,
_groupSlug?: string | null,
): Promise<void> {
const index = await loadIndex(baseCwd);
const entry = index[sessionId];
if (entry) {
const absFile = path.join(getSessionsDir(baseCwd), entry.file);
await fs.rm(absFile, { force: true });
// Clean up empty cwd directory
try {
const cwdDir = path.dirname(absFile);
const remaining = await fs.readdir(cwdDir);
if (remaining.length === 0) {
await fs.rmdir(cwdDir);
}
} catch {
// Ignore cleanup errors
}
}
delete index[sessionId];
await saveIndex(baseCwd, index);
}
type SessionFilter = {
context?: string;
contextId?: string;
};
export async function listUserSessions(baseCwd: string, filter?: SessionFilter): Promise<SessionMeta[]> {
let index = await loadIndex(baseCwd);
if (Object.keys(index).length === 0) {
index = await rebuildIndex(baseCwd);
}
let sessions: SessionMeta[] = Object.entries(index).map(([id, entry]) => indexEntryToMeta(id, entry));
if (filter?.context) {
if (filter.context === 'chat') {
// 'chat' matches sessions with no context or context='chat'
sessions = sessions.filter((s) => !s.context || s.context === 'chat');
} else {
sessions = sessions.filter((s) => s.context === filter.context && s.contextId === filter.contextId);
}
}
sessions.sort((a, b) => b.updatedAt - a.updatedAt);
return sessions;
}
export async function searchSessions(
baseCwd: string,
query: string,
): Promise<Array<SessionMeta & { preview?: string; relevance?: number }>> {
const index = await loadIndex(baseCwd);
const lowerQuery = query.toLowerCase();
const results: Array<SessionMeta & { preview?: string; relevance?: number }> = [];
for (const [sessionId, entry] of Object.entries(index)) {
let relevance = 0;
let preview = '';
// Check title (from index — fast)
if (entry.title.toLowerCase().includes(lowerQuery)) {
relevance += 1.0;
preview = entry.title;
}
// Check group membership — boost relevance if group matches
if (entry.groupSlug) {
try {
const group = await loadGroup(baseCwd, entry.groupSlug);
if (group.name.toLowerCase().includes(lowerQuery)) {
relevance += 2.0;
}
if (group.description?.toLowerCase().includes(lowerQuery)) {
relevance += 1.5;
}
} catch {
// Group file missing — skip boost
}
}
// Lazy content search — only read JSONL if title didn't match
if (relevance === 0 || !preview) {
try {
const absFile = path.join(getSessionsDir(baseCwd), entry.file);
const raw = await fs.readFile(absFile, 'utf-8');
const { entries } = parseJnlFile(raw);
const { messages } = jnlEntriesToMessages(entries);
for (const msg of messages) {
if (msg.text && msg.text.toLowerCase().includes(lowerQuery)) {
relevance += 0.5;
if (!preview) {
const idx = msg.text.toLowerCase().indexOf(lowerQuery);
const start = Math.max(0, idx - 50);
const end = Math.min(msg.text.length, idx + query.length + 50);
preview = '...' + msg.text.slice(start, end) + '...';
}
}
}
} catch {
// File unreadable — skip
}
}
if (relevance > 0) {
const meta = indexEntryToMeta(sessionId, entry);
results.push({ ...meta, preview, relevance });
}
}
results.sort((a, b) => (b.relevance ?? 0) - (a.relevance ?? 0));
return results;
}
// ── Group management ───────────────────────────────────────────────────
export async function saveGroup(baseCwd: string, groupMeta: GroupMeta): Promise<void> {
const groupsDir = getGroupsDir(baseCwd);
await fs.mkdir(groupsDir, { recursive: true });
const groupPath = getGroupPath(baseCwd, groupMeta.slug);
await fs.writeFile(groupPath, JSON.stringify(groupMeta, null, 2));
}
export async function loadGroup(baseCwd: string, groupSlug: string): Promise<GroupMeta> {
const groupPath = getGroupPath(baseCwd, groupSlug);
const content = await fs.readFile(groupPath, 'utf-8');
return JSON.parse(content) as GroupMeta;
}
export async function groupExists(baseCwd: string, groupSlug: string): Promise<boolean> {
try {
const groupPath = getGroupPath(baseCwd, groupSlug);
await fs.access(groupPath);
return true;
} catch {
return false;
}
}
export async function listGroups(baseCwd: string): Promise<GroupMeta[]> {
const groupsDir = getGroupsDir(baseCwd);
try {
await fs.access(groupsDir);
} catch {
return [];
}
const files = await fs.readdir(groupsDir);
const groups: GroupMeta[] = [];
for (const file of files) {
if (!file.endsWith('.json')) continue;
try {
const filePath = path.join(groupsDir, file);
const content = await fs.readFile(filePath, 'utf-8');
groups.push(JSON.parse(content) as GroupMeta);
} catch {
continue;
}
}
groups.sort((a, b) => b.updatedAt - a.updatedAt);
return groups;
}
export async function updateGroupMeta(
baseCwd: string,
groupSlug: string,
updates: Partial<GroupMeta>,
): Promise<GroupMeta> {
const groupMeta = await loadGroup(baseCwd, groupSlug);
const updated: GroupMeta = {
...groupMeta,
...updates,
slug: groupMeta.slug,
updatedAt: Date.now(),
};
await saveGroup(baseCwd, updated);
return updated;
}
export async function deleteGroup(baseCwd: string, groupSlug: string): Promise<void> {
const index = await loadIndex(baseCwd);
// Remove groupSlug from all sessions in this group
let changed = false;
for (const entry of Object.values(index)) {
if (entry.groupSlug === groupSlug) {
entry.groupSlug = null;
changed = true;
}
}
if (changed) {
await saveIndex(baseCwd, index);
}
// Delete group file
const groupPath = getGroupPath(baseCwd, groupSlug);
await fs.rm(groupPath, { force: true });
}
export async function moveSession(
baseCwd: string,
sessionId: string,
fromGroupSlug: string | null,
toGroupSlug: string | null,
): Promise<SessionMeta> {
const index = await loadIndex(baseCwd);
const entry = index[sessionId];
if (!entry) {
throw new Error(`Session not found: ${sessionId}`);
}
entry.groupSlug = toGroupSlug;
entry.updatedAt = Date.now();
await saveIndex(baseCwd, index);
// Update group session counts
if (fromGroupSlug) {
try {
const fromGroup = await loadGroup(baseCwd, fromGroupSlug);
fromGroup.sessionCount = Math.max(0, fromGroup.sessionCount - 1);
fromGroup.updatedAt = Date.now();
await saveGroup(baseCwd, fromGroup);
} catch {
// Group might not exist
}
}
if (toGroupSlug) {
try {
const toGroup = await loadGroup(baseCwd, toGroupSlug);
toGroup.sessionCount += 1;
toGroup.updatedAt = Date.now();
await saveGroup(baseCwd, toGroup);
} catch {
// Group might not exist
}
}
return indexEntryToMeta(sessionId, entry);
}
// ── Index recovery ─────────────────────────────────────────────────────
export async function rebuildIndex(baseCwd: string): Promise<SessionIndex> {
const sessionsDir = getSessionsDir(baseCwd);
const index: SessionIndex = {};
try {
await fs.access(sessionsDir);
} catch {
return index;
}
const cwdDirs = await fs.readdir(sessionsDir);
for (const dir of cwdDirs) {
if (dir === OFFICER_DIR) continue;
const dirPath = path.join(sessionsDir, dir);
const stat = await fs.stat(dirPath);
if (!stat.isDirectory()) continue;
const files = await fs.readdir(dirPath);
for (const file of files) {
if (!file.endsWith('.jsonl')) continue;
try {
const filePath = path.join(dirPath, file);
const raw = await fs.readFile(filePath, 'utf-8');
const { header, entries } = parseJnlFile(raw);
const relFile = path.join(dir, file);
let title = '';
let cost: MessageCost = { inputTokens: 0, outputTokens: 0, totalUSD: 0 };
let model = '';
let messageCount = 0;
let createdAt = new Date(header.timestamp).getTime();
let updatedAt = createdAt;
let groupSlug: string | null = null;
let context: string | undefined;
let contextId: string | undefined;
// Count message entries and find session_info
for (const entry of entries) {
if (entry.type === 'message') {
messageCount++;
} else if (entry.type === 'session_info') {
const info = entry as JnlSessionInfoEntry;
title = info.name;
if (info.officer) {
cost = info.officer.cost;
model = info.officer.model;
messageCount = info.officer.messageCount;
createdAt = info.officer.createdAt;
updatedAt = info.officer.updatedAt;
groupSlug = info.officer.groupSlug ?? null;
context = info.officer.context;
contextId = info.officer.contextId;
}
}
}
index[header.id] = {
file: relFile,
title,
model,
cwd: header.cwd,
createdAt,
updatedAt,
messageCount,
cost,
groupSlug,
context,
contextId,
};
} catch {
continue;
}
}
}
await saveIndex(baseCwd, index);
return index;
}