refactor(core): share sqlite boundary types
This commit is contained in:
@@ -14,6 +14,7 @@
|
|||||||
"./providers/codex": "./dist/providers/codex.js",
|
"./providers/codex": "./dist/providers/codex.js",
|
||||||
"./providers/types": "./dist/providers/types.js",
|
"./providers/types": "./dist/providers/types.js",
|
||||||
"./query": "./dist/query.js",
|
"./query": "./dist/query.js",
|
||||||
|
"./sqlite-types": "./dist/sqlite-types.js",
|
||||||
"./tx": "./dist/tx.js",
|
"./tx": "./dist/tx.js",
|
||||||
"./write-coordinator": "./dist/write-coordinator.js",
|
"./write-coordinator": "./dist/write-coordinator.js",
|
||||||
"./writer-lease": "./dist/writer-lease.js"
|
"./writer-lease": "./dist/writer-lease.js"
|
||||||
|
|||||||
@@ -2,14 +2,13 @@
|
|||||||
import { createRequire } from 'node:module';
|
import { createRequire } from 'node:module';
|
||||||
import { CLAUDE_DIR, CODEX_DIR, TEXT_LIMIT, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines } from './parsing.ts';
|
import { CLAUDE_DIR, CODEX_DIR, TEXT_LIMIT, trunc, truncJson, extractText, extractContentType, extractMessageIsMeta, filePath, isDir, readLines } from './parsing.ts';
|
||||||
import { configureConnection } from './tx.ts';
|
import { configureConnection } from './tx.ts';
|
||||||
|
import type { NodeSqliteDb, SqliteDb } from './sqlite-types.ts';
|
||||||
const require = createRequire(import.meta.url);
|
const require = createRequire(import.meta.url);
|
||||||
const fs = require('node:fs');
|
const fs = require('node:fs');
|
||||||
const path = require('node:path');
|
const path = require('node:path');
|
||||||
const os = require('node:os');
|
const os = require('node:os');
|
||||||
const { DatabaseSync } = require('node:sqlite');
|
const { DatabaseSync } = require('node:sqlite');
|
||||||
|
|
||||||
type SqliteDb = any;
|
|
||||||
|
|
||||||
const OBELISK_DIR = path.join(os.homedir(), '.obelisk');
|
const OBELISK_DIR = path.join(os.homedir(), '.obelisk');
|
||||||
const LEGACY_DB_PATH = path.join(CLAUDE_DIR, 'obelisk.sqlite');
|
const LEGACY_DB_PATH = path.join(CLAUDE_DIR, 'obelisk.sqlite');
|
||||||
const DB_PATH = path.join(OBELISK_DIR, 'obelisk.sqlite');
|
const DB_PATH = path.join(OBELISK_DIR, 'obelisk.sqlite');
|
||||||
@@ -22,7 +21,7 @@ function migrateLegacyDbIfNeeded() {
|
|||||||
fs.copyFileSync(LEGACY_DB_PATH, DB_PATH);
|
fs.copyFileSync(LEGACY_DB_PATH, DB_PATH);
|
||||||
}
|
}
|
||||||
|
|
||||||
function openDb() {
|
function openDb(): NodeSqliteDb {
|
||||||
migrateLegacyDbIfNeeded();
|
migrateLegacyDbIfNeeded();
|
||||||
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
|
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
|
||||||
const db = new DatabaseSync(DB_PATH);
|
const db = new DatabaseSync(DB_PATH);
|
||||||
@@ -35,18 +34,18 @@ function openDb() {
|
|||||||
|
|
||||||
// Queries and daemon-arbitration checks must never migrate/configure the index.
|
// Queries and daemon-arbitration checks must never migrate/configure the index.
|
||||||
// The caller is responsible for ensuring the database exists first.
|
// The caller is responsible for ensuring the database exists first.
|
||||||
function openReadDb() {
|
function openReadDb(): NodeSqliteDb {
|
||||||
const db = new DatabaseSync(DB_PATH, { readOnly: true });
|
const db = new DatabaseSync(DB_PATH, { readOnly: true });
|
||||||
db.exec('PRAGMA busy_timeout=250');
|
db.exec('PRAGMA busy_timeout=250');
|
||||||
return db;
|
return db;
|
||||||
}
|
}
|
||||||
|
|
||||||
function openWriterLeaseDb(lockPath: string): SqliteDb {
|
function openWriterLeaseDb(lockPath: string): NodeSqliteDb {
|
||||||
return new DatabaseSync(lockPath);
|
return new DatabaseSync(lockPath);
|
||||||
}
|
}
|
||||||
|
|
||||||
function ensureColumn(db: SqliteDb, table: string, column: string, definition: string): void {
|
function ensureColumn(db: SqliteDb, table: string, column: string, definition: string): void {
|
||||||
const columns = db.prepare(`PRAGMA table_info(${table})`).all().map((c: { name: string }) => c.name);
|
const columns = db.prepare(`PRAGMA table_info(${table})`).all().map(c => c.name);
|
||||||
if (!columns.includes(column)) db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`);
|
if (!columns.includes(column)) db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -11,21 +11,13 @@ import { runRetryableWriteTransaction, isBeginBusyFailure, hasUnusableTransactio
|
|||||||
import { parse as claudeParse } from './providers/claude.ts';
|
import { parse as claudeParse } from './providers/claude.ts';
|
||||||
import { parse as codexParse } from './providers/codex.ts';
|
import { parse as codexParse } from './providers/codex.ts';
|
||||||
import type { Cursor, IndexRecord } from './providers/types.ts';
|
import type { Cursor, IndexRecord } from './providers/types.ts';
|
||||||
|
import type { ClaudeJsonlFile } from './parsing.ts';
|
||||||
|
import type { NodeSqliteDb, SqliteRow } from './sqlite-types.ts';
|
||||||
|
|
||||||
const HISTORY_PATH = path.join(CLAUDE_DIR, 'history.jsonl');
|
const HISTORY_PATH = path.join(CLAUDE_DIR, 'history.jsonl');
|
||||||
|
|
||||||
type SqliteDb = any;
|
|
||||||
type JsonRecord = Record<string, any>;
|
type JsonRecord = Record<string, any>;
|
||||||
|
|
||||||
interface ClaudeFileInfo {
|
|
||||||
path: string;
|
|
||||||
sessionId: string;
|
|
||||||
project: string;
|
|
||||||
isSubagent: boolean;
|
|
||||||
agentId?: string;
|
|
||||||
workflowRunId?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
interface SkippedFile {
|
interface SkippedFile {
|
||||||
path: string;
|
path: string;
|
||||||
error: string;
|
error: string;
|
||||||
@@ -42,7 +34,7 @@ function errorMessage(error: unknown): string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
function needsReindex(db: SqliteDb, fp: string) {
|
function needsReindex(db: NodeSqliteDb, fp: string) {
|
||||||
const mt = fs.statSync(fp).mtimeMs;
|
const mt = fs.statSync(fp).mtimeMs;
|
||||||
const row = db.prepare('SELECT mtime, lines_processed FROM index_state WHERE jsonl_path = ?').get(fp);
|
const row = db.prepare('SELECT mtime, lines_processed FROM index_state WHERE jsonl_path = ?').get(fp);
|
||||||
if (!row) return { needed: true, skip: 0 };
|
if (!row) return { needed: true, skip: 0 };
|
||||||
@@ -50,7 +42,7 @@ function needsReindex(db: SqliteDb, fp: string) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
function indexCodexSessionIndex(db: SqliteDb): void {
|
function indexCodexSessionIndex(db: NodeSqliteDb): void {
|
||||||
const indexPath = path.join(CODEX_DIR, 'session_index.jsonl');
|
const indexPath = path.join(CODEX_DIR, 'session_index.jsonl');
|
||||||
if (!fs.existsSync(indexPath)) return;
|
if (!fs.existsSync(indexPath)) return;
|
||||||
readLines(indexPath, (line) => {
|
readLines(indexPath, (line) => {
|
||||||
@@ -67,7 +59,7 @@ function indexCodexSessionIndex(db: SqliteDb): void {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
function refreshSessionProjectPaths(db: SqliteDb): void {
|
function refreshSessionProjectPaths(db: NodeSqliteDb): void {
|
||||||
const sessions = db.prepare('SELECT id, project FROM sessions').all();
|
const sessions = db.prepare('SELECT id, project FROM sessions').all();
|
||||||
const cwdStmt = db.prepare(`
|
const cwdStmt = db.prepare(`
|
||||||
SELECT cwd
|
SELECT cwd
|
||||||
@@ -77,13 +69,13 @@ function refreshSessionProjectPaths(db: SqliteDb): void {
|
|||||||
`);
|
`);
|
||||||
const update = db.prepare('UPDATE sessions SET project_path = ? WHERE id = ?');
|
const update = db.prepare('UPDATE sessions SET project_path = ? WHERE id = ?');
|
||||||
for (const session of sessions) {
|
for (const session of sessions) {
|
||||||
const cwds = cwdStmt.all(session.id).map((row: JsonRecord) => row.cwd);
|
const cwds = cwdStmt.all(session.id).map((row: SqliteRow) => row.cwd);
|
||||||
const projectPath = inferProjectPath(session.project, cwds);
|
const projectPath = inferProjectPath(session.project, cwds);
|
||||||
if (projectPath) update.run(projectPath, session.id);
|
if (projectPath) update.run(projectPath, session.id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function indexSubagentMeta(db: SqliteDb, fi: ClaudeFileInfo): void {
|
function indexSubagentMeta(db: NodeSqliteDb, fi: ClaudeJsonlFile): void {
|
||||||
if (!fi.isSubagent) return;
|
if (!fi.isSubagent) return;
|
||||||
const mp = fi.path.replace('.jsonl', '.meta.json');
|
const mp = fi.path.replace('.jsonl', '.meta.json');
|
||||||
if (!fs.existsSync(mp)) return;
|
if (!fs.existsSync(mp)) return;
|
||||||
@@ -104,7 +96,7 @@ function indexSubagentMeta(db: SqliteDb, fi: ClaudeFileInfo): void {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function indexWorkflows(db: SqliteDb): void {
|
function indexWorkflows(db: NodeSqliteDb): void {
|
||||||
if (!fs.existsSync(PROJECTS_DIR)) return;
|
if (!fs.existsSync(PROJECTS_DIR)) return;
|
||||||
let projects;
|
let projects;
|
||||||
try { projects = fs.readdirSync(PROJECTS_DIR); } catch { return; }
|
try { projects = fs.readdirSync(PROJECTS_DIR); } catch { return; }
|
||||||
@@ -145,7 +137,7 @@ function indexWorkflows(db: SqliteDb): void {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function indexHistory(db: SqliteDb): void {
|
function indexHistory(db: NodeSqliteDb): void {
|
||||||
if (!fs.existsSync(HISTORY_PATH)) return;
|
if (!fs.existsSync(HISTORY_PATH)) return;
|
||||||
readLines(HISTORY_PATH, (line) => {
|
readLines(HISTORY_PATH, (line) => {
|
||||||
let item: JsonRecord;
|
let item: JsonRecord;
|
||||||
@@ -162,7 +154,7 @@ function indexHistory(db: SqliteDb): void {
|
|||||||
const BUILD_DEBOUNCE_MS = 30000;
|
const BUILD_DEBOUNCE_MS = 30000;
|
||||||
const APP_HEARTBEAT_FRESH_MS = 60000;
|
const APP_HEARTBEAT_FRESH_MS = 60000;
|
||||||
|
|
||||||
function shouldSkipBuild(db: SqliteDb, { now = Date.now(), ignoreRecentBuild = false }: BuildCheckOptions = {}) {
|
function shouldSkipBuild(db: NodeSqliteDb, { now = Date.now(), ignoreRecentBuild = false }: BuildCheckOptions = {}) {
|
||||||
const appHeartbeat = db.prepare("SELECT mtime FROM index_state WHERE jsonl_path='__app_heartbeat__'").get();
|
const appHeartbeat = db.prepare("SELECT mtime FROM index_state WHERE jsonl_path='__app_heartbeat__'").get();
|
||||||
if (appHeartbeat && now - appHeartbeat.mtime < APP_HEARTBEAT_FRESH_MS) {
|
if (appHeartbeat && now - appHeartbeat.mtime < APP_HEARTBEAT_FRESH_MS) {
|
||||||
return { skip: true, reason: 'daemon_active' };
|
return { skip: true, reason: 'daemon_active' };
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ const TEXT_LIMIT = 10000;
|
|||||||
type JsonRecord = Record<string, any>;
|
type JsonRecord = Record<string, any>;
|
||||||
type JsonValue = any;
|
type JsonValue = any;
|
||||||
|
|
||||||
interface ClaudeJsonlFile {
|
export interface ClaudeJsonlFile {
|
||||||
path: string;
|
path: string;
|
||||||
sessionId: string;
|
sessionId: string;
|
||||||
project: string;
|
project: string;
|
||||||
@@ -27,7 +27,7 @@ interface ClaudeJsonlFile {
|
|||||||
source?: 'claude';
|
source?: 'claude';
|
||||||
}
|
}
|
||||||
|
|
||||||
interface CodexJsonlFile {
|
export interface CodexJsonlFile {
|
||||||
path: string;
|
path: string;
|
||||||
source: 'codex';
|
source: 'codex';
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,7 +10,7 @@
|
|||||||
// that consumes the records and writes them (index_state, FTS, upsert).
|
// that consumes the records and writes them (index_state, FTS, upsert).
|
||||||
//
|
//
|
||||||
// This file defines only the shapes crossing that boundary. Record fields mirror
|
// This file defines only the shapes crossing that boundary. Record fields mirror
|
||||||
// the columns in scripts/schema.sql; keep them in sync. Types only — no runtime
|
// the columns in packages/core/src/schema.sql; keep them in sync. Types only — no runtime
|
||||||
// code — so consumers must import with `import type`.
|
// code — so consumers must import with `import type`.
|
||||||
|
|
||||||
// Opaque per-unit resume/watermark token. The orchestration stores it verbatim
|
// Opaque per-unit resume/watermark token. The orchestration stores it verbatim
|
||||||
@@ -47,7 +47,7 @@ export interface DiscoverContext {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/** Discriminated union of everything an adapter's parse can emit. Each record
|
/** Discriminated union of everything an adapter's parse can emit. Each record
|
||||||
* kind maps to one schema table (see scripts/schema.sql); `delete-session` is
|
* kind maps to one schema table (see packages/core/src/schema.sql); `delete-session` is
|
||||||
* the exception — a retraction op, not a table. Sources without a table
|
* the exception — a retraction op, not a table. Sources without a table
|
||||||
* (history.jsonl, codex session_index.jsonl) are not records: adapters fold them
|
* (history.jsonl, codex session_index.jsonl) are not records: adapters fold them
|
||||||
* into the SessionRecord they already emit. */
|
* into the SessionRecord they already emit. */
|
||||||
|
|||||||
@@ -1,8 +1,8 @@
|
|||||||
// Query and attune sandbox helpers for the Core package.
|
// Query and attune sandbox helpers for the Core package.
|
||||||
import { readLines, fs, path } from './db.ts';
|
import { readLines, fs, path } from './db.ts';
|
||||||
|
import type { SqliteDb, SqliteRow } from './sqlite-types.ts';
|
||||||
|
|
||||||
type SqliteDb = any;
|
type DbRow = SqliteRow;
|
||||||
type DbRow = Record<string, any>;
|
|
||||||
|
|
||||||
interface QueryOptions extends Record<string, any> {
|
interface QueryOptions extends Record<string, any> {
|
||||||
limit?: number;
|
limit?: number;
|
||||||
@@ -162,7 +162,7 @@ function createQueryApi(db: SqliteDb) {
|
|||||||
if (!msg) return null;
|
if (!msg) return null;
|
||||||
const session = db.prepare('SELECT * FROM sessions WHERE id=?').get(msg.session_id);
|
const session = db.prepare('SELECT * FROM sessions WHERE id=?').get(msg.session_id);
|
||||||
const chain: DbRow[] = [];
|
const chain: DbRow[] = [];
|
||||||
let cur = msg;
|
let cur: DbRow | undefined = msg;
|
||||||
while (cur?.parent_uuid) { cur = db.prepare('SELECT * FROM messages WHERE uuid=?').get(cur.parent_uuid); if (cur) chain.unshift(cur); }
|
while (cur?.parent_uuid) { cur = db.prepare('SELECT * FROM messages WHERE uuid=?').get(cur.parent_uuid); if (cur) chain.unshift(cur); }
|
||||||
const subagent = msg.agent_id ? db.prepare('SELECT * FROM subagents WHERE agent_id=?').get(msg.agent_id) : null;
|
const subagent = msg.agent_id ? db.prepare('SELECT * FROM subagents WHERE agent_id=?').get(msg.agent_id) : null;
|
||||||
let workflow = null;
|
let workflow = null;
|
||||||
@@ -176,7 +176,7 @@ function createQueryApi(db: SqliteDb) {
|
|||||||
const trace = (uuid: string) => {
|
const trace = (uuid: string) => {
|
||||||
const chain: DbRow[] = [];
|
const chain: DbRow[] = [];
|
||||||
let cur = db.prepare('SELECT * FROM messages WHERE uuid=?').get(uuid);
|
let cur = db.prepare('SELECT * FROM messages WHERE uuid=?').get(uuid);
|
||||||
while (cur) { chain.unshift(cur); cur = cur.parent_uuid ? db.prepare('SELECT * FROM messages WHERE uuid=?').get(cur.parent_uuid) : null; }
|
while (cur) { chain.unshift(cur); cur = cur.parent_uuid ? db.prepare('SELECT * FROM messages WHERE uuid=?').get(cur.parent_uuid) : undefined; }
|
||||||
return chain;
|
return chain;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,21 @@
|
|||||||
|
// Minimal structural types shared by node:sqlite and better-sqlite3 consumers.
|
||||||
|
// SQLite rows and bindings are dynamic at this boundary; domain records become
|
||||||
|
// strongly typed after parsing, in providers/types.ts.
|
||||||
|
|
||||||
|
export type SqliteRow = Record<string, any>;
|
||||||
|
|
||||||
|
export interface SqliteStatement {
|
||||||
|
all(...bindings: any[]): SqliteRow[];
|
||||||
|
get(...bindings: any[]): SqliteRow | undefined;
|
||||||
|
run(...bindings: any[]): unknown;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface SqliteDb {
|
||||||
|
exec(sql: string): unknown;
|
||||||
|
prepare(sql: string): SqliteStatement;
|
||||||
|
close(): void;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface NodeSqliteDb extends SqliteDb {
|
||||||
|
readonly isTransaction: boolean;
|
||||||
|
}
|
||||||
@@ -32,7 +32,8 @@ test('build:skill produces a runnable, readable, .ts-free skill artifact', () =>
|
|||||||
'package.json', 'SKILL.md', 'references/api-reference.md',
|
'package.json', 'SKILL.md', 'references/api-reference.md',
|
||||||
'scripts/core.js', 'scripts/persist.js', 'scripts/providers/claude.js',
|
'scripts/core.js', 'scripts/persist.js', 'scripts/providers/claude.js',
|
||||||
'scripts/providers/codex.js', 'scripts/runtime.js', 'scripts/indexer.js',
|
'scripts/providers/codex.js', 'scripts/runtime.js', 'scripts/indexer.js',
|
||||||
'scripts/db.js', 'scripts/parsing.js', 'scripts/query.js', 'scripts/schema.sql',
|
'scripts/db.js', 'scripts/parsing.js', 'scripts/query.js',
|
||||||
|
'scripts/sqlite-types.js', 'scripts/schema.sql',
|
||||||
]) {
|
]) {
|
||||||
assert.ok(existsSync(join(skillDir, rel)), `artifact missing ${rel}`);
|
assert.ok(existsSync(join(skillDir, rel)), `artifact missing ${rel}`);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user