feat: add Codex transcript indexing and dual-source retrieval
Obelisk now indexes ~/.codex/ sessions alongside Claude Code, storing both under a unified schema with a column for provenance. Includes DB migration to ~/.obelisk/, guardian thread filtering, and source-aware query helpers.
This commit is contained in:
+2
-1
@@ -5,4 +5,5 @@ tests/
|
||||
node_modules/
|
||||
dist-renderer/
|
||||
release/
|
||||
|
||||
docs/
|
||||
.claude/
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
---
|
||||
name: obelisk
|
||||
description: >
|
||||
Search and query past Claude Code session history.
|
||||
Search and query past Claude Code and Codex session history.
|
||||
Reactive: when the user asks "how did I fix X", "what did we do last time", "find the session where", "上次怎么修的", "之前的session", "历史记录".
|
||||
Proactive: when the user references past work you lack context for, when you're about to modify a file with complex edit history, when the user says "继续之前的" or "continue where we left off", or when understanding prior decisions would improve your current response.
|
||||
Memory: when the user says "记住这个", "remember this", "写入记忆", "save this conclusion", or when you determine a retrieval result contains a conclusion worth persisting.
|
||||
@@ -13,11 +13,19 @@ allowed-tools:
|
||||
|
||||
# obelisk
|
||||
|
||||
Search and query Claude Code session history stored in `~/.claude/`.
|
||||
Search and query Claude Code and Codex session history stored in `~/.claude/`
|
||||
and `~/.codex/`.
|
||||
Obelisk indexes sessions, messages, tool calls, tool results, summaries,
|
||||
subagents, workflows, workflow agents, parent chains, and raw JSONL lines into
|
||||
SQLite + FTS5.
|
||||
|
||||
Obelisk has two transcript sources. Treat both as ordinary sessions by default:
|
||||
Claude rows use `source='claude'`; Codex rows use `source='codex'` and IDs
|
||||
prefixed with `codex:`. Use `source` only when provenance matters or the user
|
||||
asks to scope to one provider. Codex subagent child threads are mapped to the
|
||||
same `subagents` table; Codex workflow rows may be absent because Codex does not
|
||||
emit Claude-style workflow metadata.
|
||||
|
||||
Obelisk is a CodeAct memory layer: write a small JS query, run it locally, read
|
||||
the JSON, then answer. Do not turn history into a flat document or browse entire
|
||||
sessions by default.
|
||||
@@ -123,8 +131,8 @@ messages.
|
||||
Returns:
|
||||
|
||||
```js
|
||||
[{ message: { uuid, text, content_type, is_meta, role, timestamp, model, cwd },
|
||||
session: { id, title, project, started_at },
|
||||
[{ message: { uuid, text, content_type, is_meta, role, timestamp, model, cwd, source },
|
||||
session: { id, title, project, started_at, source },
|
||||
rank,
|
||||
context }]
|
||||
```
|
||||
@@ -148,13 +156,16 @@ be treated as the user's request by default. `search()` and `thread()` omit meta
|
||||
messages unless `includeMeta: true` is passed; `context()` and `trace()` preserve
|
||||
the original chain and expose `is_meta` on rows.
|
||||
|
||||
Opts: `{ limit, sessionId, project, after, before, cwd, includeMeta }`.
|
||||
Opts: `{ limit, sessionId, project, after, before, cwd, source, includeMeta }`.
|
||||
|
||||
`project` is a SQL `LIKE` filter over `sessions.project`, not an exact project
|
||||
identity. Results are already ordered by FTS5 rank; lower rank sorts earlier.
|
||||
Prefer returned order over manually interpreting numeric rank unless you are
|
||||
deliberately using FTS5 semantics.
|
||||
|
||||
`source` can be `'claude'`, `'codex'`, or omitted. Omitted means search all
|
||||
indexed sources.
|
||||
|
||||
### `context(uuid)`
|
||||
|
||||
Returns the full story around one indexed message:
|
||||
@@ -191,13 +202,13 @@ not replace `sql()`, but they are the default first-pass surface. Use `sql()`
|
||||
when you need an exact aggregation or a join the helper does not expose.
|
||||
|
||||
All list helpers accept a bounded `limit`. Many also accept:
|
||||
`{ project, after, before, sessionId, sessions, branch }`. Check the schema or a
|
||||
tiny sample before relying on less common filters.
|
||||
`{ project, after, before, sessionId, sessions, branch, source }`. Check the
|
||||
schema or a tiny sample before relying on less common filters.
|
||||
|
||||
- `overview(opts?)` -- compact orientation map. Returns current cwd/project if knowable, global project counts, and current-project recent sessions plus memory records. It is a map, not evidence.
|
||||
- `overview(opts?)` -- compact orientation map. Returns current cwd/project if knowable, global project/source counts, and current-project recent sessions plus memory records. It is a map, not evidence.
|
||||
- `sessions(opts?)` -- session rows, newest first. `project` is a SQL `LIKE` pattern.
|
||||
- `recent(n?)` -- shorthand for recent sessions.
|
||||
- `summaries(opts?)` -- summary rows, newest first: `{ id, session_id, timestamp, source, content, session_title, project }`.
|
||||
- `summaries(opts?)` -- summary rows, newest first: `{ id, session_id, timestamp, source, content, session_title, project }`; here `source` is the summary kind, not the transcript provider.
|
||||
- `subagents(opts?)` -- subagent metadata plus `messageCount`.
|
||||
- `workflows(opts?)` -- workflow runs, newest first.
|
||||
- `workflowTree(runId)` -- workflow row plus parsed `result` and `agents`; may include bulky `script` and `result_json`, so project compact fields.
|
||||
@@ -372,6 +383,6 @@ See `references/query-patterns.md` for longer recipes.
|
||||
## Notes
|
||||
|
||||
- First run builds the index. Later runs update incrementally.
|
||||
- DB location: `~/.claude/obelisk.sqlite`.
|
||||
- DB location: `~/.obelisk/obelisk.sqlite`; old `~/.claude/obelisk.sqlite` is copied forward if needed.
|
||||
- Query scripts run in a sandboxed VM with no filesystem or network access from inside the script.
|
||||
- Indexed text and stored tool inputs/results are truncated to 10k chars. Use `raw(uuid, { offset, limit })` for specific JSONL windows.
|
||||
|
||||
@@ -33,11 +33,13 @@ unless scoped evidence is insufficient and `query_plan` says why.
|
||||
|
||||
Project-like fields are distinct:
|
||||
|
||||
- `sessions.project`: stored Claude Code project slug.
|
||||
- `sessions.project`: provider-normalized project slug.
|
||||
- `memories.project`: stored project slug copied onto registered memory records.
|
||||
- `sessions.project_path`: absolute session path derived from message `cwd` when available; slug decoding is only a fallback.
|
||||
- `messages.cwd`: working directory at message time.
|
||||
- `sessions.source` / `messages.source`: transcript provider, currently `claude` or `codex`.
|
||||
- helper `project`: SQL `LIKE` over `sessions.project`, not exact membership.
|
||||
- helper `source`: optional provider filter. Omit it unless provenance matters.
|
||||
|
||||
For exact project membership, prefer helper filters or a scoped first pass when
|
||||
they are expressive enough; use `sql()` with `s.project = ?` or
|
||||
|
||||
+57
-27
@@ -10,24 +10,35 @@ contract for agents and humans; it is not the runtime source of truth.
|
||||
|
||||
## 1. Database Schema
|
||||
|
||||
Database location: `~/.claude/obelisk.sqlite`
|
||||
Database location: `~/.obelisk/obelisk.sqlite`. Older
|
||||
`~/.claude/obelisk.sqlite` databases are copied forward on first open when the
|
||||
new database does not exist.
|
||||
|
||||
Obelisk indexes two transcript sources into the same schema: Claude Code rows
|
||||
use `source='claude'`; Codex rows use `source='codex'` and synthetic IDs
|
||||
prefixed with `codex:`. Query helpers search all sources by default. Use helper
|
||||
`source` filters or raw SQL on `sessions.source` / `messages.source` only when
|
||||
provider provenance matters.
|
||||
|
||||
### sessions
|
||||
|
||||
One row per Claude Code session.
|
||||
One row per root session. Claude session IDs match JSONL filenames. Codex root
|
||||
session IDs are prefixed with `codex:`; Codex child threads are attached through
|
||||
`subagents` instead of becoming separate `sessions` rows.
|
||||
|
||||
```sql
|
||||
CREATE TABLE sessions (
|
||||
id TEXT PRIMARY KEY, -- session UUID (matches JSONL filename)
|
||||
title TEXT, -- AI-generated session title (may be NULL)
|
||||
project TEXT, -- Claude project slug (e.g. "-Users-tomiya-Code-quiet-zero")
|
||||
id TEXT PRIMARY KEY, -- Claude UUID or "codex:<thread-id>"
|
||||
title TEXT, -- AI-generated/session title (may be NULL)
|
||||
project TEXT, -- provider-normalized project slug (e.g. "-Users-tomiya-Code-quiet-zero")
|
||||
project_path TEXT, -- absolute session cwd-derived path, with slug fallback (e.g. "/Users/tomiya/Code/quiet-zero")
|
||||
started_at TEXT, -- ISO 8601 timestamp of first message
|
||||
ended_at TEXT, -- ISO 8601 timestamp of last message
|
||||
git_branch TEXT, -- git branch active during session (if any)
|
||||
version TEXT, -- Claude Code version string
|
||||
version TEXT, -- provider CLI/app version string
|
||||
message_count INTEGER DEFAULT 0, -- total user + assistant messages
|
||||
jsonl_path TEXT -- absolute path to source JSONL file
|
||||
jsonl_path TEXT, -- absolute path to source JSONL file
|
||||
source TEXT DEFAULT 'claude' -- "claude" or "codex"
|
||||
);
|
||||
```
|
||||
|
||||
@@ -53,18 +64,19 @@ CREATE TABLE messages (
|
||||
output_tokens INTEGER, -- token usage (assistant messages only)
|
||||
cwd TEXT, -- working directory at message time (may differ from session project_path)
|
||||
skill TEXT, -- skill that generated this response (e.g. "obelisk"), NULL if none
|
||||
turn_duration_ms INTEGER -- wall-clock duration of the turn ending at this message (from system turn_duration event)
|
||||
turn_duration_ms INTEGER, -- wall-clock duration of the turn ending at this message
|
||||
source TEXT DEFAULT 'claude' -- "claude" or "codex"
|
||||
);
|
||||
```
|
||||
|
||||
Indexes: `idx_messages_session(session_id)`, `idx_messages_agent(agent_id)`, `idx_messages_ts(session_id, timestamp)`.
|
||||
|
||||
`content_type` preserves the top-level Claude Code content block shape for the
|
||||
message row. Treat `text` as user/assistant visible language, `thinking` as
|
||||
trace/debug material, and `tool_use` as a marker that the assistant message
|
||||
contains tool calls. `tool_result` marks a tool-result message, but the
|
||||
structured payload remains in `tool_results`. Tool-call details remain in
|
||||
`tool_calls`. Messages whose top-level content is not one of these four raw
|
||||
`content_type` preserves the normalized transcript surface. Treat `text` as
|
||||
user/assistant visible language, `thinking` as trace/debug material, and
|
||||
`tool_use` as a marker that the assistant message contains tool calls.
|
||||
`tool_result` marks a tool-result message when the provider emits one as a
|
||||
message; structured payloads remain in `tool_results`. Tool-call details remain
|
||||
in `tool_calls`. Messages whose top-level content is not one of these raw
|
||||
message surfaces are `unknown`. Real user input is represented by `type='user'`
|
||||
and `content_type='text'`, not by a separate `user_message` content type.
|
||||
|
||||
@@ -298,27 +310,30 @@ Full-text search across all message text using FTS5.
|
||||
| `opts.after` | `string` | ISO 8601 lower bound on timestamp |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound on timestamp |
|
||||
| `opts.cwd` | `string` | Filter by working directory (supports LIKE) |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.includeMeta` | `boolean` | Include injected/control-plane messages (default `false`) |
|
||||
|
||||
**Scope note:** `sessions.project` is the stored Claude Code project slug,
|
||||
**Scope note:** `sessions.project` is the provider-normalized project slug,
|
||||
`sessions.project_path` is the absolute session path derived from message `cwd`
|
||||
when available, and `messages.cwd` is the working directory at message time.
|
||||
Helper `project` filters are fuzzy `LIKE` filters over `sessions.project`. For
|
||||
exact project membership, use `sql()` with `s.project = ?` or
|
||||
`s.project_path = ?`.
|
||||
when available, `messages.cwd` is the working directory at message time, and
|
||||
`source` is the provider. Helper `project` filters are fuzzy `LIKE` filters over
|
||||
`sessions.project`. For exact project membership, use `sql()` with
|
||||
`s.project = ?` or `s.project_path = ?`.
|
||||
|
||||
**Returns:** `Array<{ message, session, rank, context }>` where `message`
|
||||
includes `{ uuid, text, content_type, is_meta, role, timestamp, model, cwd }`
|
||||
and `context` is the 6 nearest non-meta messages by timestamp in the same
|
||||
session unless `includeMeta: true` is passed. It is temporal neighbor context,
|
||||
not a parent chain. `rank` is the FTS5 relevance score used by `ORDER BY rank`;
|
||||
lower values sort earlier, so treat the returned order as the relevance order
|
||||
unless you are deliberately using FTS5 ranking details.
|
||||
includes `{ uuid, text, content_type, is_meta, role, timestamp, model, cwd,
|
||||
source }`, `session` includes `source`, and `context` is the 6 nearest non-meta
|
||||
messages by timestamp in the same session unless `includeMeta: true` is passed.
|
||||
It is temporal neighbor context, not a parent chain. `rank` is the FTS5
|
||||
relevance score used by `ORDER BY rank`; lower values sort earlier, so treat the
|
||||
returned order as the relevance order unless you are deliberately using FTS5
|
||||
ranking details.
|
||||
|
||||
```js
|
||||
const hits = search('MCTS exploration');
|
||||
return hits.map(h => ({
|
||||
title: h.session.title,
|
||||
source: h.session.source,
|
||||
content_type: h.message.content_type,
|
||||
is_meta: h.message.is_meta,
|
||||
text: h.message.text?.slice(0, 200),
|
||||
@@ -418,6 +433,7 @@ All subagent spawns, with message counts. For backward compatibility, passing a
|
||||
|-------|------|-------------|
|
||||
| `opts.sessionId` | `string` | Restrict to one session |
|
||||
| `opts.project` | `string` | SQL `LIKE` pattern over `sessions.project` |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 100) |
|
||||
|
||||
**Returns:** `Array<{ ...subagent_row, messageCount }>`.
|
||||
@@ -437,6 +453,7 @@ Workflow executions. For backward compatibility, passing a string is treated as
|
||||
| `opts.project` | `string` | SQL `LIKE` pattern over `sessions.project` |
|
||||
| `opts.after` | `string` | ISO 8601 lower bound on timestamp |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound on timestamp |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 100) |
|
||||
|
||||
**Returns:** `Array<workflow_row>`.
|
||||
@@ -466,6 +483,7 @@ All tool calls that touched a specific file, across every session.
|
||||
| `filePath` | `string` | Absolute file path (required) |
|
||||
| `opts.after` | `string` | ISO 8601 lower bound |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 200) |
|
||||
|
||||
**Returns:** `Array<{ toolCall, session, timestamp }>`.
|
||||
@@ -488,6 +506,7 @@ Tool calls whose results contain error patterns (`Error`, `ENOENT`, `failed`, `p
|
||||
| `opts.project` | `string` | SQL `LIKE` pattern over `sessions.project` |
|
||||
| `opts.after` | `string` | ISO 8601 lower bound |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 50) |
|
||||
|
||||
**Returns:** `Array<{ toolCall, result, session, nextMessages }>`.
|
||||
@@ -523,11 +542,15 @@ string is treated as `sessionId`, and passing a number is treated as `limit`.
|
||||
| `opts.after` | `string` | ISO 8601 lower bound on summary timestamp |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound on summary timestamp |
|
||||
| `opts.branch` | `string` | Filter by source session git branch (exact match) |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 100) |
|
||||
|
||||
**Returns:** `Array<summary_row & { session_title, project }>` ordered by
|
||||
`timestamp` descending.
|
||||
|
||||
Note: `summaries.source` is the summary kind such as `away_summary`; provider
|
||||
filtering uses the joined session's `source`.
|
||||
|
||||
```js
|
||||
const rows = summaries({ project: '%quiet-zero%', limit: 5 });
|
||||
return rows.map(s => ({
|
||||
@@ -574,7 +597,7 @@ from `process.cwd()` against `sessions.project_path`, then from exact
|
||||
project_path,
|
||||
session_total,
|
||||
sessions: [
|
||||
{ id, title, project, project_path, started_at, ended_at, git_branch, message_count }
|
||||
{ id, title, project, project_path, started_at, ended_at, git_branch, message_count, source }
|
||||
],
|
||||
memory_total,
|
||||
memories: [
|
||||
@@ -592,7 +615,12 @@ from `process.cwd()` against `sessions.project_path`, then from exact
|
||||
recent_branches
|
||||
}
|
||||
],
|
||||
totals: { projects, sessions, memories }
|
||||
totals: {
|
||||
projects,
|
||||
sessions,
|
||||
memories,
|
||||
sources: [{ source: 'claude' | 'codex', session_count, last_session_at }]
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
@@ -629,6 +657,7 @@ Query sessions with filters. For backward compatibility, passing a number is tre
|
||||
| `opts.before` | `string` | ISO 8601 upper bound on `started_at` |
|
||||
| `opts.limit` | `number` | Max results (default 50) |
|
||||
| `opts.branch` | `string` | Filter by git branch (exact match) |
|
||||
| `opts.source` | `string` | Optional provider filter: `"claude"` or `"codex"` |
|
||||
| `opts.sessionId` | `string` | Restrict to one session |
|
||||
| `opts.sessions` | `string[]` | Restrict to a set of session IDs |
|
||||
|
||||
@@ -656,6 +685,7 @@ string is treated as `sessionId`, and passing a number is treated as `limit`.
|
||||
| `opts.after` | `string` | ISO 8601 lower bound on `created_at` |
|
||||
| `opts.before` | `string` | ISO 8601 upper bound on `created_at` |
|
||||
| `opts.branch` | `string` | Filter by source session git branch (exact match) |
|
||||
| `opts.source` | `string` | Optional provider filter through the source session: `"claude"` or `"codex"` |
|
||||
| `opts.limit` | `number` | Max results (default 50) |
|
||||
|
||||
**Returns:** `Array<memory_row & { rank?: number }>` with archived memories
|
||||
|
||||
+30
-3
@@ -6,14 +6,27 @@ const os = require('node:os');
|
||||
const { DatabaseSync } = require('node:sqlite');
|
||||
|
||||
const CLAUDE_DIR = path.join(os.homedir(), '.claude');
|
||||
const DB_PATH = path.join(CLAUDE_DIR, 'obelisk.sqlite');
|
||||
const CODEX_DIR = path.join(os.homedir(), '.codex');
|
||||
const OBELISK_DIR = path.join(os.homedir(), '.obelisk');
|
||||
const LEGACY_DB_PATH = path.join(CLAUDE_DIR, 'obelisk.sqlite');
|
||||
const DB_PATH = path.join(OBELISK_DIR, 'obelisk.sqlite');
|
||||
const TEXT_LIMIT = 10000;
|
||||
const SCHEMA = fs.readFileSync(new URL('./schema.sql', import.meta.url), 'utf8');
|
||||
|
||||
function migrateLegacyDbIfNeeded() {
|
||||
if (fs.existsSync(DB_PATH)) return;
|
||||
if (!fs.existsSync(LEGACY_DB_PATH)) return;
|
||||
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
|
||||
fs.copyFileSync(LEGACY_DB_PATH, DB_PATH);
|
||||
}
|
||||
|
||||
function openDb() {
|
||||
migrateLegacyDbIfNeeded();
|
||||
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
|
||||
const db = new DatabaseSync(DB_PATH);
|
||||
db.exec('PRAGMA journal_mode=WAL');
|
||||
db.exec('PRAGMA synchronous=NORMAL');
|
||||
migrateExistingColumns(db);
|
||||
db.exec(SCHEMA);
|
||||
migrateDb(db);
|
||||
return db;
|
||||
@@ -24,13 +37,27 @@ function ensureColumn(db, table, column, definition) {
|
||||
if (!columns.includes(column)) db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`);
|
||||
}
|
||||
|
||||
function migrateDb(db) {
|
||||
function tableExists(db, table) {
|
||||
return Boolean(db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name=?").get(table));
|
||||
}
|
||||
|
||||
function migrateExistingColumns(db) {
|
||||
if (tableExists(db, 'sessions')) ensureColumn(db, 'sessions', 'source', "TEXT DEFAULT 'claude'");
|
||||
if (tableExists(db, 'messages')) {
|
||||
ensureColumn(db, 'messages', 'content_type', 'TEXT');
|
||||
ensureColumn(db, 'messages', 'is_meta', 'INTEGER DEFAULT 0');
|
||||
ensureColumn(db, 'messages', 'source', "TEXT DEFAULT 'claude'");
|
||||
}
|
||||
if (tableExists(db, 'memories')) {
|
||||
ensureColumn(db, 'memories', 'anchors', 'TEXT');
|
||||
ensureColumn(db, 'memories', 'deleted_at', 'TEXT');
|
||||
ensureColumn(db, 'memories', 'deleted_reason', 'TEXT');
|
||||
}
|
||||
}
|
||||
|
||||
function migrateDb(db) {
|
||||
migrateExistingColumns(db);
|
||||
}
|
||||
|
||||
function rebuildMemoryFts(db) {
|
||||
db.exec("INSERT INTO memories_fts(memories_fts) VALUES('rebuild')");
|
||||
@@ -118,4 +145,4 @@ function readLines(filePath, callback) {
|
||||
}
|
||||
}
|
||||
|
||||
export { CLAUDE_DIR, DB_PATH, TEXT_LIMIT, openDb, rebuildMemoryFts, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines, fs, path, os };
|
||||
export { CLAUDE_DIR, CODEX_DIR, OBELISK_DIR, DB_PATH, TEXT_LIMIT, openDb, rebuildMemoryFts, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines, fs, path, os };
|
||||
|
||||
+491
-6
@@ -1,7 +1,8 @@
|
||||
import { CLAUDE_DIR, openDb, rebuildMemoryFts, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines, fs, path } from './db.mjs';
|
||||
import { CLAUDE_DIR, CODEX_DIR, openDb, rebuildMemoryFts, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines, fs, path } from './db.mjs';
|
||||
|
||||
const PROJECTS_DIR = path.join(CLAUDE_DIR, 'projects');
|
||||
const HISTORY_PATH = path.join(CLAUDE_DIR, 'history.jsonl');
|
||||
const CODEX_SESSIONS_DIR = path.join(CODEX_DIR, 'sessions');
|
||||
|
||||
function legacyProjectPathFromSlug(project) {
|
||||
if (!project) return null;
|
||||
@@ -13,6 +14,12 @@ function normalizeObservedCwd(cwd) {
|
||||
return path.normalize(cwd);
|
||||
}
|
||||
|
||||
function projectSlugFromPath(projectPath) {
|
||||
const normalized = normalizeObservedCwd(projectPath);
|
||||
if (!normalized) return null;
|
||||
return '-' + normalized.replace(/^[\\/]+/, '').replace(/[\\/]+/g, '-');
|
||||
}
|
||||
|
||||
function inferProjectPath(project, observedCwds = []) {
|
||||
const byPath = new Map();
|
||||
for (const cwd of observedCwds) {
|
||||
@@ -68,6 +75,25 @@ function discoverJsonlFiles() {
|
||||
return files;
|
||||
}
|
||||
|
||||
function discoverCodexJsonlFiles() {
|
||||
const files = [];
|
||||
if (!fs.existsSync(CODEX_SESSIONS_DIR)) return files;
|
||||
const walk = (dir) => {
|
||||
let entries;
|
||||
try { entries = fs.readdirSync(dir, { withFileTypes: true }); } catch { return; }
|
||||
for (const entry of entries) {
|
||||
const fp = path.join(dir, entry.name);
|
||||
if (entry.isDirectory()) {
|
||||
walk(fp);
|
||||
} else if (entry.isFile() && entry.name.endsWith('.jsonl')) {
|
||||
files.push({ path: fp, source: 'codex' });
|
||||
}
|
||||
}
|
||||
};
|
||||
walk(CODEX_SESSIONS_DIR);
|
||||
return files;
|
||||
}
|
||||
|
||||
function needsReindex(db, fp) {
|
||||
const mt = fs.statSync(fp).mtimeMs;
|
||||
const row = db.prepare('SELECT mtime, lines_processed FROM index_state WHERE jsonl_path = ?').get(fp);
|
||||
@@ -81,8 +107,8 @@ function indexJsonl(db, fi) {
|
||||
const mt = fs.statSync(fi.path).mtimeMs;
|
||||
|
||||
const ins = {
|
||||
ses: db.prepare('INSERT OR REPLACE INTO sessions (id,title,project,project_path,started_at,ended_at,git_branch,version,message_count,jsonl_path) VALUES (?,?,?,?,?,?,?,?,?,?)'),
|
||||
msg: db.prepare('INSERT OR REPLACE INTO messages (uuid,session_id,type,parent_uuid,timestamp,role,text,content_type,is_meta,model,is_sidechain,agent_id,input_tokens,output_tokens,cwd,skill) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)'),
|
||||
ses: db.prepare('INSERT OR REPLACE INTO sessions (id,title,project,project_path,started_at,ended_at,git_branch,version,message_count,jsonl_path,source) VALUES (?,?,?,?,?,?,?,?,?,?,?)'),
|
||||
msg: db.prepare('INSERT OR REPLACE INTO messages (uuid,session_id,type,parent_uuid,timestamp,role,text,content_type,is_meta,model,is_sidechain,agent_id,input_tokens,output_tokens,cwd,skill,source) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)'),
|
||||
tc: db.prepare('INSERT OR REPLACE INTO tool_calls (id,message_uuid,session_id,name,input_json,file_path) VALUES (?,?,?,?,?,?)'),
|
||||
tr: db.prepare('INSERT OR REPLACE INTO tool_results (tool_use_id,message_uuid,session_id,content,file_path,is_error) VALUES (?,?,?,?,?,?)'),
|
||||
sum: db.prepare('INSERT OR REPLACE INTO summaries (id,session_id,timestamp,source,content) VALUES (?,?,?,?,?)'),
|
||||
@@ -138,7 +164,7 @@ function indexJsonl(db, fi) {
|
||||
ins.msg.run(obj.uuid, sid, obj.type, obj.parentUuid || null, ts,
|
||||
msg.role || obj.type, text, contentType, isMeta, msg.model || null,
|
||||
obj.isSidechain ? 1 : 0, aid, usage.input_tokens || null, usage.output_tokens || null,
|
||||
obj.cwd || null, obj.attributionSkill || null);
|
||||
obj.cwd || null, obj.attributionSkill || null, 'claude');
|
||||
}
|
||||
|
||||
if (obj.type === 'assistant' && Array.isArray(msg.content)) {
|
||||
@@ -160,11 +186,462 @@ function indexJsonl(db, fi) {
|
||||
|
||||
if (!fi.isSubagent) {
|
||||
const pp = inferProjectPath(fi.project, sm.cwds);
|
||||
ins.ses.run(fi.sessionId, sm.title, fi.project, pp, sm.started_at, sm.ended_at, sm.git_branch, sm.version, sm.n, fi.path);
|
||||
ins.ses.run(fi.sessionId, sm.title, fi.project, pp, sm.started_at, sm.ended_at, sm.git_branch, sm.version, sm.n, fi.path, 'claude');
|
||||
}
|
||||
ins.idx.run(fi.path, mt, lineNum);
|
||||
}
|
||||
|
||||
function codexDbId(id) {
|
||||
if (!id) return null;
|
||||
const raw = String(id).replace(/^codex:/, '');
|
||||
return `codex:${raw}`;
|
||||
}
|
||||
|
||||
function codexRawId(id) {
|
||||
return id ? String(id).replace(/^codex:/, '') : null;
|
||||
}
|
||||
|
||||
function codexLineUuid(threadId, lineNum) {
|
||||
return `codex:${codexRawId(threadId)}:${String(lineNum).padStart(6, '0')}`;
|
||||
}
|
||||
|
||||
function codexCallId(callId) {
|
||||
if (!callId) return null;
|
||||
return `codex:${String(callId).replace(/^codex:/, '')}`;
|
||||
}
|
||||
|
||||
function codexParentThreadId(meta) {
|
||||
const subagent = meta?.source?.subagent;
|
||||
return subagent?.thread_spawn?.parent_thread_id
|
||||
|| meta?.forked_from_id
|
||||
|| subagent?.parent_thread_id
|
||||
|| null;
|
||||
}
|
||||
|
||||
function codexIsGuardianThread(meta, records = []) {
|
||||
const subagent = meta?.source?.subagent;
|
||||
if (subagent?.other === 'guardian') return true;
|
||||
if (meta?.thread_source !== 'subagent') return false;
|
||||
return records.some(({ obj }) => obj?.payload?.model === 'codex-auto-review' || obj?.model === 'codex-auto-review');
|
||||
}
|
||||
|
||||
function deleteCodexThreadRows(db, threadRawId) {
|
||||
const threadId = codexDbId(threadRawId);
|
||||
if (!threadId) return;
|
||||
db.prepare(`
|
||||
DELETE FROM tool_results
|
||||
WHERE session_id = ?
|
||||
OR message_uuid IN (SELECT uuid FROM messages WHERE session_id = ? OR agent_id = ?)
|
||||
`).run(threadId, threadId, threadId);
|
||||
db.prepare(`
|
||||
DELETE FROM tool_calls
|
||||
WHERE session_id = ?
|
||||
OR message_uuid IN (SELECT uuid FROM messages WHERE session_id = ? OR agent_id = ?)
|
||||
`).run(threadId, threadId, threadId);
|
||||
db.prepare('DELETE FROM messages WHERE session_id = ? OR agent_id = ?').run(threadId, threadId);
|
||||
db.prepare('DELETE FROM subagents WHERE agent_id = ? OR session_id = ?').run(threadId, threadId);
|
||||
db.prepare('DELETE FROM summaries WHERE session_id = ?').run(threadId);
|
||||
db.prepare('DELETE FROM sessions WHERE id = ?').run(threadId);
|
||||
}
|
||||
|
||||
function readCodexGuardianThreadInfo(filePath) {
|
||||
const records = [];
|
||||
let metaRecord = null;
|
||||
let lineNum = 0;
|
||||
readLines(filePath, (line) => {
|
||||
lineNum++;
|
||||
let obj;
|
||||
try {
|
||||
obj = JSON.parse(line);
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
records.push({ lineNum, obj });
|
||||
if (obj?.type === 'session_meta' && obj.payload?.id) {
|
||||
metaRecord = { lineNum, obj };
|
||||
if (obj.payload?.source?.subagent?.other === 'guardian') return false;
|
||||
if (obj.payload?.thread_source !== 'subagent') return false;
|
||||
}
|
||||
if (metaRecord && codexIsGuardianThread(metaRecord.obj.payload, records)) return false;
|
||||
});
|
||||
const meta = metaRecord?.obj?.payload;
|
||||
if (!meta || !codexIsGuardianThread(meta, records)) return null;
|
||||
return { threadRawId: codexRawId(meta.id), lineNum };
|
||||
}
|
||||
|
||||
function codexAgentNickname(meta) {
|
||||
return meta?.agent_nickname
|
||||
|| meta?.source?.subagent?.thread_spawn?.agent_nickname
|
||||
|| null;
|
||||
}
|
||||
|
||||
function codexAgentRole(meta) {
|
||||
return meta?.agent_role
|
||||
|| meta?.source?.subagent?.thread_spawn?.agent_role
|
||||
|| null;
|
||||
}
|
||||
|
||||
function parseCodexJsonInput(value) {
|
||||
if (value === null || value === undefined || value === '') return {};
|
||||
if (typeof value !== 'string') return value;
|
||||
try { return JSON.parse(value); } catch { return value; }
|
||||
}
|
||||
|
||||
function codexUsage(payload) {
|
||||
const usage = payload?.info?.last_token_usage || payload?.info?.total_token_usage || payload?.last_token_usage || null;
|
||||
if (!usage) return {};
|
||||
return {
|
||||
inputTokens: usage.input_tokens ?? null,
|
||||
outputTokens: usage.output_tokens ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
function codexEventText(payload) {
|
||||
if (typeof payload?.message === 'string') return payload.message;
|
||||
if (Array.isArray(payload?.text_elements) && payload.text_elements.length) {
|
||||
const parts = payload.text_elements.map(item => typeof item === 'string' ? item : item?.text).filter(Boolean);
|
||||
if (parts.length) return parts.join('\n');
|
||||
}
|
||||
if (typeof payload?.text === 'string') return payload.text;
|
||||
return null;
|
||||
}
|
||||
|
||||
function codexMessagePayloadText(payload) {
|
||||
if (!Array.isArray(payload?.content)) return null;
|
||||
const parts = [];
|
||||
for (const block of payload.content) {
|
||||
if (typeof block?.text === 'string') parts.push(block.text);
|
||||
}
|
||||
return parts.length ? parts.join('\n') : null;
|
||||
}
|
||||
|
||||
function codexVisibleMessageKey(role, text) {
|
||||
return `${role || ''}\u0000${text || ''}`;
|
||||
}
|
||||
|
||||
function codexToolInput(payload) {
|
||||
if (payload?.type === 'custom_tool_call') return parseCodexJsonInput(payload.input);
|
||||
if (payload?.type === 'tool_search_call') return parseCodexJsonInput(payload.arguments);
|
||||
if (payload?.type === 'web_search_call') return { action: payload.action || null };
|
||||
return parseCodexJsonInput(payload?.arguments);
|
||||
}
|
||||
|
||||
function codexToolOutput(payload) {
|
||||
if (typeof payload?.output === 'string') return payload.output;
|
||||
if (payload?.output !== undefined) return JSON.stringify(payload.output);
|
||||
if (payload?.tools !== undefined) return JSON.stringify(payload.tools);
|
||||
if (payload?.execution !== undefined) return JSON.stringify(payload.execution);
|
||||
return null;
|
||||
}
|
||||
|
||||
function upsertCodexSubagent(db, {
|
||||
agentId,
|
||||
sessionId,
|
||||
parentToolUseId = null,
|
||||
agentType = null,
|
||||
description = null,
|
||||
durationMs = null,
|
||||
totalTokens = null,
|
||||
} = {}) {
|
||||
if (!agentId || !sessionId) return;
|
||||
db.prepare(`
|
||||
INSERT INTO subagents (agent_id,session_id,parent_tool_use_id,agent_type,description,duration_ms,total_tokens)
|
||||
VALUES (?,?,?,?,?,?,?)
|
||||
ON CONFLICT(agent_id) DO UPDATE SET
|
||||
session_id=excluded.session_id,
|
||||
parent_tool_use_id=COALESCE(excluded.parent_tool_use_id, subagents.parent_tool_use_id),
|
||||
agent_type=COALESCE(excluded.agent_type, subagents.agent_type),
|
||||
description=COALESCE(excluded.description, subagents.description),
|
||||
duration_ms=COALESCE(excluded.duration_ms, subagents.duration_ms),
|
||||
total_tokens=COALESCE(excluded.total_tokens, subagents.total_tokens)
|
||||
`).run(agentId, sessionId, parentToolUseId, agentType, description, durationMs, totalTokens);
|
||||
}
|
||||
|
||||
function indexCodexJsonl(db, fi) {
|
||||
const state = needsReindex(db, fi.path);
|
||||
if (!state.needed) {
|
||||
const guardian = readCodexGuardianThreadInfo(fi.path);
|
||||
if (guardian) deleteCodexThreadRows(db, guardian.threadRawId);
|
||||
return;
|
||||
}
|
||||
const mt = fs.statSync(fi.path).mtimeMs;
|
||||
const records = [];
|
||||
let lineNum = 0;
|
||||
readLines(fi.path, (line) => {
|
||||
lineNum++;
|
||||
try {
|
||||
records.push({ lineNum, obj: JSON.parse(line) });
|
||||
} catch {}
|
||||
});
|
||||
|
||||
const metaRecord = records.find(r => r.obj?.type === 'session_meta' && r.obj.payload?.id);
|
||||
if (!metaRecord) {
|
||||
db.prepare('INSERT OR REPLACE INTO index_state (jsonl_path,mtime,lines_processed) VALUES (?,?,?)').run(fi.path, mt, lineNum);
|
||||
return;
|
||||
}
|
||||
|
||||
const meta = metaRecord.obj.payload;
|
||||
const threadRawId = codexRawId(meta.id);
|
||||
if (codexIsGuardianThread(meta, records)) {
|
||||
deleteCodexThreadRows(db, threadRawId);
|
||||
db.prepare('INSERT OR REPLACE INTO index_state (jsonl_path,mtime,lines_processed) VALUES (?,?,?)').run(fi.path, mt, lineNum);
|
||||
return;
|
||||
}
|
||||
const parentRawId = codexParentThreadId(meta);
|
||||
const sessionId = codexDbId(parentRawId || threadRawId);
|
||||
const agentId = parentRawId ? codexDbId(threadRawId) : null;
|
||||
const isSidechain = agentId ? 1 : 0;
|
||||
const projectPath = normalizeObservedCwd(meta.cwd);
|
||||
const project = projectSlugFromPath(projectPath);
|
||||
const sm = {
|
||||
started_at: meta.timestamp || metaRecord.obj.timestamp || null,
|
||||
ended_at: meta.timestamp || metaRecord.obj.timestamp || null,
|
||||
git_branch: meta.git?.branch || null,
|
||||
version: meta.cli_version || null,
|
||||
title: null,
|
||||
n: 0,
|
||||
cwds: projectPath ? [projectPath] : [],
|
||||
lastMessageUuid: null,
|
||||
lastTextAssistantUuid: null,
|
||||
totalInputTokens: 0,
|
||||
totalOutputTokens: 0,
|
||||
};
|
||||
|
||||
const ins = {
|
||||
ses: db.prepare('INSERT OR REPLACE INTO sessions (id,title,project,project_path,started_at,ended_at,git_branch,version,message_count,jsonl_path,source) VALUES (?,?,?,?,?,?,?,?,?,?,?)'),
|
||||
msg: db.prepare(`
|
||||
INSERT INTO messages (uuid,session_id,type,parent_uuid,timestamp,role,text,content_type,is_meta,model,is_sidechain,agent_id,input_tokens,output_tokens,cwd,skill,source)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(uuid) DO UPDATE SET
|
||||
session_id=excluded.session_id,
|
||||
type=excluded.type,
|
||||
parent_uuid=excluded.parent_uuid,
|
||||
timestamp=excluded.timestamp,
|
||||
role=excluded.role,
|
||||
text=excluded.text,
|
||||
content_type=excluded.content_type,
|
||||
is_meta=excluded.is_meta,
|
||||
model=excluded.model,
|
||||
is_sidechain=excluded.is_sidechain,
|
||||
agent_id=excluded.agent_id,
|
||||
input_tokens=excluded.input_tokens,
|
||||
output_tokens=excluded.output_tokens,
|
||||
cwd=excluded.cwd,
|
||||
skill=excluded.skill,
|
||||
source=excluded.source
|
||||
`),
|
||||
tc: db.prepare('INSERT OR REPLACE INTO tool_calls (id,message_uuid,session_id,name,input_json,file_path) VALUES (?,?,?,?,?,?)'),
|
||||
tr: db.prepare('INSERT OR REPLACE INTO tool_results (tool_use_id,message_uuid,session_id,content,file_path,is_error) VALUES (?,?,?,?,?,?)'),
|
||||
idx: db.prepare('INSERT OR REPLACE INTO index_state (jsonl_path,mtime,lines_processed) VALUES (?,?,?)'),
|
||||
dur: db.prepare('UPDATE messages SET turn_duration_ms=? WHERE uuid=?'),
|
||||
usage: db.prepare('UPDATE messages SET input_tokens=?, output_tokens=? WHERE uuid=?'),
|
||||
};
|
||||
|
||||
let currentCwd = projectPath;
|
||||
let currentModel = null;
|
||||
const eventMessageKeys = new Set();
|
||||
const callMessageUuids = new Map();
|
||||
|
||||
for (const { obj } of records) {
|
||||
if (obj?.type !== 'event_msg') continue;
|
||||
const payload = obj.payload || {};
|
||||
if (payload.type !== 'user_message' && payload.type !== 'agent_message') continue;
|
||||
const text = codexEventText(payload);
|
||||
if (text === null) continue;
|
||||
eventMessageKeys.add(codexVisibleMessageKey(payload.type === 'user_message' ? 'user' : 'assistant', text));
|
||||
}
|
||||
|
||||
const updateBounds = (ts) => {
|
||||
if (!ts) return;
|
||||
if (!sm.started_at || ts < sm.started_at) sm.started_at = ts;
|
||||
if (!sm.ended_at || ts > sm.ended_at) sm.ended_at = ts;
|
||||
};
|
||||
|
||||
const insertMessage = ({ uuid, type, role, text = null, contentType = 'text', timestamp, isMeta = 0 }) => {
|
||||
ins.msg.run(
|
||||
uuid,
|
||||
sessionId,
|
||||
type,
|
||||
sm.lastMessageUuid,
|
||||
timestamp || null,
|
||||
role,
|
||||
trunc(text),
|
||||
contentType,
|
||||
isMeta,
|
||||
currentModel,
|
||||
isSidechain,
|
||||
agentId,
|
||||
null,
|
||||
null,
|
||||
currentCwd,
|
||||
null,
|
||||
'codex',
|
||||
);
|
||||
sm.lastMessageUuid = uuid;
|
||||
if (!agentId) sm.n++;
|
||||
if (type === 'assistant' && contentType === 'text') sm.lastTextAssistantUuid = uuid;
|
||||
updateBounds(timestamp);
|
||||
return uuid;
|
||||
};
|
||||
|
||||
for (const { lineNum: currentLine, obj } of records) {
|
||||
const ts = obj.timestamp || null;
|
||||
if (obj.type === 'session_meta') {
|
||||
if (obj.payload?.cwd) {
|
||||
currentCwd = normalizeObservedCwd(obj.payload.cwd) || currentCwd;
|
||||
if (currentCwd) sm.cwds.push(currentCwd);
|
||||
}
|
||||
if (obj.payload?.git?.branch) sm.git_branch = obj.payload.git.branch;
|
||||
if (obj.payload?.cli_version) sm.version = obj.payload.cli_version;
|
||||
updateBounds(obj.payload?.timestamp || ts);
|
||||
continue;
|
||||
}
|
||||
if (obj.type === 'turn_context') {
|
||||
currentCwd = normalizeObservedCwd(obj.payload?.cwd) || currentCwd;
|
||||
currentModel = obj.payload?.model || currentModel;
|
||||
if (currentCwd) sm.cwds.push(currentCwd);
|
||||
updateBounds(ts);
|
||||
continue;
|
||||
}
|
||||
if (obj.type === 'event_msg') {
|
||||
const payload = obj.payload || {};
|
||||
if (payload.type === 'user_message' || payload.type === 'agent_message' || payload.type === 'agent_reasoning') {
|
||||
const text = codexEventText(payload);
|
||||
if (text === null) continue;
|
||||
const isReasoning = payload.type === 'agent_reasoning';
|
||||
insertMessage({
|
||||
uuid: codexLineUuid(threadRawId, currentLine),
|
||||
type: payload.type === 'user_message' ? 'user' : 'assistant',
|
||||
role: payload.type === 'user_message' ? 'user' : 'assistant',
|
||||
text,
|
||||
contentType: isReasoning ? 'thinking' : 'text',
|
||||
timestamp: ts,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
if (payload.type === 'collab_agent_spawn_end' && payload.call_id && payload.new_thread_id) {
|
||||
const uuid = insertMessage({
|
||||
uuid: codexLineUuid(threadRawId, currentLine),
|
||||
type: 'assistant',
|
||||
role: 'assistant',
|
||||
text: null,
|
||||
contentType: 'tool_use',
|
||||
timestamp: ts,
|
||||
});
|
||||
const toolId = codexCallId(payload.call_id);
|
||||
const description = payload.new_agent_nickname || payload.new_agent_role || 'Agent';
|
||||
const input = {
|
||||
description,
|
||||
subagent_type: payload.new_agent_role || 'Agent',
|
||||
prompt: payload.prompt || '',
|
||||
new_thread_id: payload.new_thread_id,
|
||||
model: payload.model || null,
|
||||
reasoning_effort: payload.reasoning_effort || null,
|
||||
};
|
||||
ins.tc.run(toolId, uuid, sessionId, 'Agent', truncJson(input), null);
|
||||
callMessageUuids.set(toolId, uuid);
|
||||
upsertCodexSubagent(db, {
|
||||
agentId: codexDbId(payload.new_thread_id),
|
||||
sessionId,
|
||||
parentToolUseId: toolId,
|
||||
agentType: payload.new_agent_role || null,
|
||||
description,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
if (payload.type === 'task_complete') {
|
||||
if (sm.lastTextAssistantUuid && payload.duration_ms !== undefined) {
|
||||
ins.dur.run(payload.duration_ms || null, sm.lastTextAssistantUuid);
|
||||
}
|
||||
updateBounds(ts);
|
||||
continue;
|
||||
}
|
||||
if (payload.type === 'token_count') {
|
||||
const usage = codexUsage(payload);
|
||||
if (usage.inputTokens !== null) sm.totalInputTokens = usage.inputTokens;
|
||||
if (usage.outputTokens !== null) sm.totalOutputTokens = usage.outputTokens;
|
||||
if (sm.lastTextAssistantUuid && (usage.inputTokens !== null || usage.outputTokens !== null)) {
|
||||
ins.usage.run(usage.inputTokens, usage.outputTokens, sm.lastTextAssistantUuid);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (payload.type === 'thread_name_updated' && payload.thread_name) {
|
||||
sm.title = payload.thread_name;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (obj.type !== 'response_item') continue;
|
||||
const payload = obj.payload || {};
|
||||
if (payload.type === 'message' && payload.role !== 'developer') {
|
||||
const text = codexMessagePayloadText(payload);
|
||||
const role = payload.role || 'assistant';
|
||||
if (text !== null && !eventMessageKeys.has(codexVisibleMessageKey(role, text))) {
|
||||
insertMessage({
|
||||
uuid: codexLineUuid(threadRawId, currentLine),
|
||||
type: role === 'user' ? 'user' : 'assistant',
|
||||
role,
|
||||
text,
|
||||
contentType: 'text',
|
||||
timestamp: ts,
|
||||
});
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (['function_call', 'custom_tool_call', 'tool_search_call', 'web_search_call'].includes(payload.type) && payload.call_id) {
|
||||
const uuid = insertMessage({
|
||||
uuid: codexLineUuid(threadRawId, currentLine),
|
||||
type: 'assistant',
|
||||
role: 'assistant',
|
||||
text: null,
|
||||
contentType: 'tool_use',
|
||||
timestamp: ts,
|
||||
});
|
||||
const name = payload.name || payload.tool || payload.type.replace(/_call$/, '');
|
||||
const toolId = codexCallId(payload.call_id);
|
||||
ins.tc.run(toolId, uuid, sessionId, name, truncJson(codexToolInput(payload)), null);
|
||||
callMessageUuids.set(toolId, uuid);
|
||||
continue;
|
||||
}
|
||||
if (['function_call_output', 'custom_tool_call_output', 'tool_search_output'].includes(payload.type) && payload.call_id) {
|
||||
const toolId = codexCallId(payload.call_id);
|
||||
ins.tr.run(toolId, callMessageUuids.get(toolId) || null, sessionId, trunc(codexToolOutput(payload) || ''), null, payload.is_error ? 1 : 0);
|
||||
}
|
||||
}
|
||||
|
||||
if (agentId) {
|
||||
const tokenTotal = (sm.totalInputTokens || 0) + (sm.totalOutputTokens || 0);
|
||||
const started = sm.started_at ? new Date(sm.started_at).getTime() : null;
|
||||
const ended = sm.ended_at ? new Date(sm.ended_at).getTime() : null;
|
||||
upsertCodexSubagent(db, {
|
||||
agentId,
|
||||
sessionId,
|
||||
agentType: codexAgentRole(meta),
|
||||
description: codexAgentNickname(meta),
|
||||
durationMs: started && ended ? ended - started : null,
|
||||
totalTokens: tokenTotal || null,
|
||||
});
|
||||
} else {
|
||||
const pp = inferProjectPath(project, sm.cwds);
|
||||
ins.ses.run(sessionId, sm.title, project, pp, sm.started_at, sm.ended_at, sm.git_branch, sm.version, sm.n, fi.path, 'codex');
|
||||
}
|
||||
ins.idx.run(fi.path, mt, lineNum);
|
||||
}
|
||||
|
||||
function indexCodexSessionIndex(db) {
|
||||
const indexPath = path.join(CODEX_DIR, 'session_index.jsonl');
|
||||
if (!fs.existsSync(indexPath)) return;
|
||||
readLines(indexPath, (line) => {
|
||||
try {
|
||||
const item = JSON.parse(line);
|
||||
if (!item.id || !item.thread_name) return;
|
||||
db.prepare('UPDATE sessions SET title=COALESCE(title, ?), ended_at=COALESCE(ended_at, ?) WHERE id=? AND source=?')
|
||||
.run(item.thread_name, item.updated_at || null, codexDbId(item.id), 'codex');
|
||||
} catch (e) {
|
||||
process.stderr.write(`Warning: malformed Codex session index line: ${e.message}\n`);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
function refreshSessionProjectPaths(db) {
|
||||
const sessions = db.prepare('SELECT id, project FROM sessions').all();
|
||||
const cwdStmt = db.prepare(`
|
||||
@@ -275,12 +752,19 @@ function buildIndex({ force = false } = {}) {
|
||||
db.prepare("DELETE FROM index_state WHERE jsonl_path != '__last_build__'").run();
|
||||
}
|
||||
|
||||
const files = discoverJsonlFiles();
|
||||
const files = [
|
||||
...discoverJsonlFiles(),
|
||||
...discoverCodexJsonlFiles(),
|
||||
];
|
||||
for (const f of files) {
|
||||
db.exec('BEGIN');
|
||||
try {
|
||||
if (f.source === 'codex') {
|
||||
indexCodexJsonl(db, f);
|
||||
} else {
|
||||
indexJsonl(db, f);
|
||||
indexSubagentMeta(db, f);
|
||||
}
|
||||
db.exec('COMMIT');
|
||||
} catch (e) {
|
||||
db.exec('ROLLBACK');
|
||||
@@ -292,6 +776,7 @@ function buildIndex({ force = false } = {}) {
|
||||
indexWorkflows(db);
|
||||
refreshSessionProjectPaths(db);
|
||||
indexHistory(db);
|
||||
indexCodexSessionIndex(db);
|
||||
db.exec("INSERT INTO messages_fts(messages_fts) VALUES('rebuild')");
|
||||
rebuildMemoryFts(db);
|
||||
db.prepare("INSERT OR REPLACE INTO index_state (jsonl_path, mtime, lines_processed) VALUES ('__last_build__', ?, 0)").run(Date.now());
|
||||
|
||||
+66
-17
@@ -19,6 +19,10 @@ function buildWhere(opts, aliases) {
|
||||
if (opts.after) { clauses.push(`${aliases.timestamp} > ?`); params.push(opts.after); }
|
||||
if (opts.before) { clauses.push(`${aliases.timestamp} < ?`); params.push(opts.before); }
|
||||
if (opts.branch) { clauses.push(`${aliases.branch} = ?`); params.push(opts.branch); }
|
||||
if (opts.source && opts.source !== 'all' && aliases.source) {
|
||||
clauses.push(`COALESCE(${aliases.source}, 'claude') = ?`);
|
||||
params.push(opts.source);
|
||||
}
|
||||
return { where: clauses.length ? clauses.join(' AND ') : '1=1', params };
|
||||
}
|
||||
|
||||
@@ -67,7 +71,7 @@ function createQueryApi(db) {
|
||||
};
|
||||
|
||||
const search = (text, opts = {}) => {
|
||||
const { limit = 20, sessionId, project, after, before, cwd, includeMeta = false } = opts;
|
||||
const { limit = 20, sessionId, project, after, before, cwd, source, includeMeta = false } = opts;
|
||||
let where = 'WHERE mf.text MATCH ?';
|
||||
const p = [text];
|
||||
if (sessionId) { where += ' AND mf.session_id=?'; p.push(sessionId); }
|
||||
@@ -75,22 +79,25 @@ function createQueryApi(db) {
|
||||
if (after) { where += ' AND m.timestamp>?'; p.push(after); }
|
||||
if (before) { where += ' AND m.timestamp<?'; p.push(before); }
|
||||
if (cwd) { where += ' AND m.cwd LIKE ?'; p.push(cwd); }
|
||||
if (source && source !== 'all') { where += " AND COALESCE(m.source, s.source, 'claude')=?"; p.push(source); }
|
||||
if (!includeMeta) where += ' AND COALESCE(m.is_meta,0)=0';
|
||||
p.push(limit);
|
||||
const rows = db.prepare(`
|
||||
SELECT m.uuid,m.session_id,m.text,m.content_type,m.is_meta,m.role,m.timestamp,m.model,m.cwd,
|
||||
SELECT m.uuid,m.session_id,m.text,m.content_type,m.is_meta,m.role,m.timestamp,m.model,m.cwd,m.source as m_source,
|
||||
s.id as s_id,s.title as s_title,s.project as s_project,s.started_at as s_started,
|
||||
s.source as s_source,
|
||||
rank
|
||||
FROM messages_fts mf JOIN messages m ON m.uuid=mf.uuid LEFT JOIN sessions s ON s.id=m.session_id
|
||||
${where} ORDER BY rank LIMIT ?`).all(...p);
|
||||
return rows.map(r => {
|
||||
const metaClause = includeMeta ? '' : 'AND COALESCE(is_meta,0)=0';
|
||||
const ctx = db.prepare(
|
||||
`SELECT uuid,text,content_type,is_meta,role,timestamp,model FROM messages WHERE session_id=? AND uuid!=? ${metaClause} ORDER BY ABS(JULIANDAY(timestamp)-JULIANDAY(?)) LIMIT 6`
|
||||
`SELECT uuid,text,content_type,is_meta,role,timestamp,model,COALESCE(source, 'claude') as source FROM messages WHERE session_id=? AND uuid!=? ${metaClause} ORDER BY ABS(JULIANDAY(timestamp)-JULIANDAY(?)) LIMIT 6`
|
||||
).all(r.session_id, r.uuid, r.timestamp).sort((a,b) => a.timestamp < b.timestamp ? -1 : 1);
|
||||
const sourceValue = r.m_source || r.s_source || 'claude';
|
||||
return {
|
||||
message: { uuid: r.uuid, text: r.text, content_type: r.content_type, is_meta: r.is_meta || 0, role: r.role, timestamp: r.timestamp, model: r.model, cwd: r.cwd },
|
||||
session: { id: r.s_id, title: r.s_title, project: r.s_project, started_at: r.s_started },
|
||||
message: { uuid: r.uuid, text: r.text, content_type: r.content_type, is_meta: r.is_meta || 0, role: r.role, timestamp: r.timestamp, model: r.model, cwd: r.cwd, source: sourceValue },
|
||||
session: { id: r.s_id, title: r.s_title, project: r.s_project, started_at: r.s_started, source: r.s_source || sourceValue },
|
||||
rank: r.rank,
|
||||
context: ctx,
|
||||
};
|
||||
@@ -129,8 +136,8 @@ function createQueryApi(db) {
|
||||
const subagents = (optsOrSid) => {
|
||||
const opts = normalizeOpts(optsOrSid);
|
||||
const { limit = 100 } = opts;
|
||||
const needsJoin = opts.project || opts.branch;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'sa.session_id', project: 's.project', timestamp: 'sa.session_id', branch: 's.git_branch' });
|
||||
const needsJoin = opts.project || opts.branch || opts.source;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'sa.session_id', project: 's.project', timestamp: 'sa.session_id', branch: 's.git_branch', source: 's.source' });
|
||||
params.push(limit);
|
||||
const join = needsJoin ? 'LEFT JOIN sessions s ON s.id=sa.session_id' : '';
|
||||
return db.prepare(`SELECT sa.* FROM subagents sa ${join} WHERE ${where} LIMIT ?`).all(...params).map(r => {
|
||||
@@ -142,8 +149,8 @@ function createQueryApi(db) {
|
||||
const workflows = (optsOrSid) => {
|
||||
const opts = normalizeOpts(optsOrSid);
|
||||
const { limit = 100 } = opts;
|
||||
const needsJoin = opts.project || opts.branch;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'w.session_id', project: 's.project', timestamp: 'w.timestamp', branch: 's.git_branch' });
|
||||
const needsJoin = opts.project || opts.branch || opts.source;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'w.session_id', project: 's.project', timestamp: 'w.timestamp', branch: 's.git_branch', source: 's.source' });
|
||||
params.push(limit);
|
||||
const join = needsJoin ? 'LEFT JOIN sessions s ON s.id=w.session_id' : '';
|
||||
return db.prepare(`SELECT w.* FROM workflows w ${join} WHERE ${where} ORDER BY w.timestamp DESC LIMIT ?`).all(...params);
|
||||
@@ -162,11 +169,12 @@ function createQueryApi(db) {
|
||||
};
|
||||
|
||||
const fileHistory = (fp, opts = {}) => {
|
||||
const { limit = 200, after, before } = opts;
|
||||
const { limit = 200, after, before, source } = opts;
|
||||
let where = 'tc.file_path=?';
|
||||
const params = [fp];
|
||||
if (after) { where += ' AND m.timestamp > ?'; params.push(after); }
|
||||
if (before) { where += ' AND m.timestamp < ?'; params.push(before); }
|
||||
if (source && source !== 'all') { where += " AND COALESCE(s.source, 'claude') = ?"; params.push(source); }
|
||||
params.push(limit);
|
||||
return db.prepare(
|
||||
`SELECT tc.*,s.title as s_title,s.project as s_project,m.timestamp as ts FROM tool_calls tc LEFT JOIN sessions s ON s.id=tc.session_id LEFT JOIN messages m ON m.uuid=tc.message_uuid WHERE ${where} ORDER BY m.timestamp LIMIT ?`
|
||||
@@ -180,8 +188,8 @@ function createQueryApi(db) {
|
||||
const failures = (optsOrSid) => {
|
||||
const opts = normalizeOpts(optsOrSid);
|
||||
const { limit = 50 } = opts;
|
||||
const needsJoin = opts.project || opts.branch;
|
||||
const { where, params: filterParams } = buildWhere(opts, { sessionId: 'tr.session_id', project: 's.project', timestamp: 'rm.timestamp', branch: 's.git_branch' });
|
||||
const needsJoin = opts.project || opts.branch || opts.source;
|
||||
const { where, params: filterParams } = buildWhere(opts, { sessionId: 'tr.session_id', project: 's.project', timestamp: 'rm.timestamp', branch: 's.git_branch', source: 's.source' });
|
||||
const join = needsJoin ? 'LEFT JOIN sessions s ON s.id=tr.session_id' : '';
|
||||
const errorCond = `(tr.is_error = 1 OR tr.content LIKE '${BASH_EXIT_PAT}')`;
|
||||
const allParams = [...filterParams, limit];
|
||||
@@ -198,7 +206,7 @@ function createQueryApi(db) {
|
||||
const sessions = (optsOrN) => {
|
||||
const opts = normalizeOpts(optsOrN, 'sessionId');
|
||||
const { limit = 50 } = opts;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 's.id', project: 's.project', timestamp: 's.started_at', branch: 's.git_branch' });
|
||||
const { where, params } = buildWhere(opts, { sessionId: 's.id', project: 's.project', timestamp: 's.started_at', branch: 's.git_branch', source: 's.source' });
|
||||
params.push(limit);
|
||||
return db.prepare(`SELECT * FROM sessions s WHERE ${where} ORDER BY ended_at DESC LIMIT ?`).all(...params);
|
||||
};
|
||||
@@ -208,7 +216,7 @@ function createQueryApi(db) {
|
||||
const summaries = (optsOrSid) => {
|
||||
const opts = normalizeOpts(optsOrSid);
|
||||
const { limit = 100 } = opts;
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'su.session_id', project: 's.project', timestamp: 'su.timestamp', branch: 's.git_branch' });
|
||||
const { where, params } = buildWhere(opts, { sessionId: 'su.session_id', project: 's.project', timestamp: 'su.timestamp', branch: 's.git_branch', source: 's.source' });
|
||||
params.push(limit);
|
||||
return db.prepare(`SELECT su.*, s.title as session_title, s.project FROM summaries su LEFT JOIN sessions s ON s.id=su.session_id WHERE ${where} ORDER BY su.timestamp DESC LIMIT ?`).all(...params);
|
||||
};
|
||||
@@ -330,7 +338,7 @@ function createQueryApi(db) {
|
||||
if (currentProject?.project) {
|
||||
const sessionTotal = db.prepare('SELECT COUNT(*) AS c FROM sessions WHERE project = ?').get(currentProject.project)?.c || 0;
|
||||
const sessionsForProject = db.prepare(`
|
||||
SELECT id, title, project, project_path, started_at, ended_at, git_branch, message_count
|
||||
SELECT id, title, project, project_path, started_at, ended_at, git_branch, message_count, COALESCE(source, 'claude') AS source
|
||||
FROM sessions
|
||||
WHERE project = ?
|
||||
ORDER BY COALESCE(ended_at, started_at) DESC
|
||||
@@ -364,6 +372,14 @@ function createQueryApi(db) {
|
||||
`).get()?.c || 0;
|
||||
const totalSessions = db.prepare('SELECT COUNT(*) AS c FROM sessions').get()?.c || 0;
|
||||
const totalMemories = db.prepare('SELECT COUNT(*) AS c FROM memories WHERE deleted_at IS NULL').get()?.c || 0;
|
||||
const sources = db.prepare(`
|
||||
SELECT COALESCE(source, 'claude') AS source,
|
||||
COUNT(*) AS session_count,
|
||||
MAX(COALESCE(ended_at, started_at)) AS last_session_at
|
||||
FROM sessions
|
||||
GROUP BY COALESCE(source, 'claude')
|
||||
ORDER BY last_session_at DESC
|
||||
`).all();
|
||||
|
||||
return {
|
||||
current: {
|
||||
@@ -376,13 +392,29 @@ function createQueryApi(db) {
|
||||
projects: totalProjects,
|
||||
sessions: totalSessions,
|
||||
memories: totalMemories,
|
||||
sources,
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
const resolveJsonlPath = (messageUuid) => {
|
||||
const msg = db.prepare('SELECT session_id, agent_id FROM messages WHERE uuid=?').get(messageUuid);
|
||||
const msg = db.prepare('SELECT session_id, agent_id, source FROM messages WHERE uuid=?').get(messageUuid);
|
||||
if (!msg) return null;
|
||||
if (msg.source === 'codex' || String(messageUuid).startsWith('codex:')) {
|
||||
const match = /^codex:([^:]+):(\d+)$/.exec(String(messageUuid));
|
||||
if (!match) return null;
|
||||
const rawThreadId = match[1];
|
||||
if (!msg.agent_id) {
|
||||
const ses = db.prepare('SELECT jsonl_path FROM sessions WHERE id=?').get(msg.session_id);
|
||||
if (ses?.jsonl_path) return ses.jsonl_path;
|
||||
}
|
||||
return db.prepare(`
|
||||
SELECT jsonl_path FROM index_state
|
||||
WHERE jsonl_path LIKE ? AND jsonl_path LIKE '%.jsonl'
|
||||
ORDER BY length(jsonl_path) ASC
|
||||
LIMIT 1
|
||||
`).get(`%${rawThreadId}.jsonl`)?.jsonl_path || null;
|
||||
}
|
||||
if (msg.agent_id) {
|
||||
const wa = db.prepare('SELECT agent_id, run_id, session_id FROM workflow_agents WHERE agent_id=?').get(msg.agent_id);
|
||||
if (wa) {
|
||||
@@ -401,8 +433,24 @@ function createQueryApi(db) {
|
||||
return null;
|
||||
};
|
||||
|
||||
const findCodexRawLine = (jsonlPath, uuid) => {
|
||||
const match = /^codex:[^:]+:(\d+)$/.exec(String(uuid));
|
||||
if (!match || !jsonlPath || !fs.existsSync(jsonlPath)) return null;
|
||||
const targetLine = Number(match[1]);
|
||||
let lineNum = 0;
|
||||
let found = null;
|
||||
readLines(jsonlPath, (line) => {
|
||||
lineNum++;
|
||||
if (lineNum !== targetLine) return;
|
||||
found = line;
|
||||
return false;
|
||||
});
|
||||
return found;
|
||||
};
|
||||
|
||||
const findRawLine = (jsonlPath, uuid) => {
|
||||
if (!jsonlPath || !fs.existsSync(jsonlPath)) return null;
|
||||
if (String(uuid).startsWith('codex:')) return findCodexRawLine(jsonlPath, uuid);
|
||||
let found = null;
|
||||
readLines(jsonlPath, (line) => {
|
||||
if (!line.includes(uuid)) return;
|
||||
@@ -429,12 +477,13 @@ function createQueryApi(db) {
|
||||
const opts = normalizeOpts(optsOrSid);
|
||||
const { limit = 50, query } = opts;
|
||||
assertEnglishMemoryText(query, 'memories() query');
|
||||
const needsJoin = opts.branch;
|
||||
const needsJoin = opts.branch || opts.source;
|
||||
const { where: baseWhere, params } = buildWhere(opts, {
|
||||
sessionId: 'mem.session_id',
|
||||
project: 'mem.project',
|
||||
timestamp: 'mem.created_at',
|
||||
branch: 's.git_branch',
|
||||
source: 's.source',
|
||||
});
|
||||
let where = baseWhere + ' AND mem.deleted_at IS NULL';
|
||||
const join = needsJoin ? 'LEFT JOIN sessions s ON s.id=mem.session_id' : '';
|
||||
|
||||
+5
-2
@@ -1,14 +1,15 @@
|
||||
CREATE TABLE IF NOT EXISTS sessions (
|
||||
id TEXT PRIMARY KEY, title TEXT, project TEXT, project_path TEXT,
|
||||
started_at TEXT, ended_at TEXT, git_branch TEXT, version TEXT,
|
||||
message_count INTEGER DEFAULT 0, jsonl_path TEXT);
|
||||
message_count INTEGER DEFAULT 0, jsonl_path TEXT, source TEXT DEFAULT 'claude');
|
||||
CREATE TABLE IF NOT EXISTS messages (
|
||||
uuid TEXT PRIMARY KEY, session_id TEXT, type TEXT, parent_uuid TEXT,
|
||||
timestamp TEXT, role TEXT, text TEXT, content_type TEXT,
|
||||
is_meta INTEGER DEFAULT 0, model TEXT,
|
||||
is_sidechain INTEGER DEFAULT 0, agent_id TEXT,
|
||||
input_tokens INTEGER, output_tokens INTEGER,
|
||||
cwd TEXT, skill TEXT, turn_duration_ms INTEGER);
|
||||
cwd TEXT, skill TEXT, turn_duration_ms INTEGER,
|
||||
source TEXT DEFAULT 'claude');
|
||||
CREATE TABLE IF NOT EXISTS tool_calls (
|
||||
id TEXT PRIMARY KEY, message_uuid TEXT, session_id TEXT,
|
||||
name TEXT, input_json TEXT, file_path TEXT);
|
||||
@@ -51,6 +52,8 @@ END;
|
||||
CREATE INDEX IF NOT EXISTS idx_messages_session ON messages(session_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_messages_agent ON messages(agent_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_messages_ts ON messages(session_id, timestamp);
|
||||
CREATE INDEX IF NOT EXISTS idx_sessions_source ON sessions(source);
|
||||
CREATE INDEX IF NOT EXISTS idx_messages_source ON messages(source);
|
||||
CREATE INDEX IF NOT EXISTS idx_tc_session_name ON tool_calls(session_id, name);
|
||||
CREATE INDEX IF NOT EXISTS idx_tc_file ON tool_calls(file_path);
|
||||
CREATE INDEX IF NOT EXISTS idx_sa_session ON subagents(session_id);
|
||||
|
||||
Reference in New Issue
Block a user