/
fleisar
/
agent-timetracker
Обзор
Документация
Войти
/
fleisar
/
agent-timetracker
Код
Запросы
0
Задачи
Вики
Пакеты
1
Релизы
2
CI/CD
Аналитика
Безопасность
main
src/storage/sqliteStore.ts
475 строк
13 KB
Matvey Kuznetsov
feat: add project and chat task grouping
14 июл 2026, 05:54
14 июл 2026, 05:54
6b9f105
Код
Авторство
О чём код?
import { mkdirSync } from "node:fs"; import path from "node:path"; import Database from "better-sqlite3"; import type { Artifact, Task, WorkSession } from "../domain/types.js"; import { parseMetadata, parseNullableStringArray, stringifyMetadata, stringifyNullableStringArray } from "../util/json.js"; import { runMigrations } from "./migrations.js"; import type { NewArtifact, NewTask, SqliteTrackerStore, TaskListFilter, TrackerStore } from "./types.js"; type SqliteDatabase = InstanceType<typeof Database>; interface TaskRow { task_id: string; title: string | null; project_id: string | null; chat_id: string | null; status: Task["status"]; created_at: string; completed_at: string | null; total_work_seconds: number; last_activity_at: string | null; idle_timeout_seconds: number; allowed_agent_ids_json: string | null; metadata_json: string; } interface WorkSessionRow { session_id: string; task_id: string; agent_id: string; status: WorkSession["status"]; started_at: string; ended_at: string | null; last_activity_at: string; start_reason: WorkSession["start_reason"]; end_reason: WorkSession["end_reason"] | null; started_by: WorkSession["started_by"]; ended_by: WorkSession["ended_by"] | null; metadata_json: string; } interface ArtifactRow { artifact_id: string; task_id: string; session_id: string | null; agent_id: string; kind: Artifact["kind"]; title: string | null; content: string | null; uri: string | null; mime_type: string | null; created_at: string; metadata_json: string; idempotency_key: string; } function ensureParentDirectory(dbPath: string): void { if (dbPath === ":memory:" || dbPath.startsWith("file:")) { return; } mkdirSync(path.dirname(dbPath), { recursive: true }); } function taskFromRow(row: TaskRow, activeSessionIds: string[]): Task { return { task_id: row.task_id, title: row.title ?? undefined, ...(row.project_id == null ? {} : { project_id: row.project_id }), ...(row.chat_id == null ? {} : { chat_id: row.chat_id }), status: row.status, created_at: row.created_at, completed_at: row.completed_at, total_work_seconds: row.total_work_seconds, active_session_ids: activeSessionIds, last_activity_at: row.last_activity_at, idle_timeout_seconds: row.idle_timeout_seconds, allowed_agent_ids: parseNullableStringArray(row.allowed_agent_ids_json), metadata: parseMetadata(row.metadata_json) }; } function taskToRow(task: NewTask): Omit<TaskRow, "task_id"> & { task_id: string } { return { task_id: task.task_id, title: task.title ?? null, project_id: task.project_id ?? null, chat_id: task.chat_id ?? null, status: task.status, created_at: task.created_at, completed_at: task.completed_at ?? null, total_work_seconds: task.total_work_seconds, last_activity_at: task.last_activity_at ?? null, idle_timeout_seconds: task.idle_timeout_seconds, allowed_agent_ids_json: stringifyNullableStringArray(task.allowed_agent_ids), metadata_json: stringifyMetadata(task.metadata) }; } function sessionFromRow(row: WorkSessionRow): WorkSession { return { session_id: row.session_id, task_id: row.task_id, agent_id: row.agent_id, status: row.status, started_at: row.started_at, ended_at: row.ended_at, last_activity_at: row.last_activity_at, start_reason: row.start_reason, end_reason: row.end_reason, started_by: row.started_by, ended_by: row.ended_by, metadata: parseMetadata(row.metadata_json) }; } function sessionToRow(session: WorkSession): WorkSessionRow { return { session_id: session.session_id, task_id: session.task_id, agent_id: session.agent_id, status: session.status, started_at: session.started_at, ended_at: session.ended_at ?? null, last_activity_at: session.last_activity_at, start_reason: session.start_reason, end_reason: session.end_reason ?? null, started_by: session.started_by, ended_by: session.ended_by ?? null, metadata_json: stringifyMetadata(session.metadata) }; } function artifactFromRow(row: ArtifactRow): Artifact { return { artifact_id: row.artifact_id, task_id: row.task_id, session_id: row.session_id, agent_id: row.agent_id, kind: row.kind, title: row.title ?? undefined, content: row.content ?? undefined, uri: row.uri ?? undefined, mime_type: row.mime_type ?? undefined, created_at: row.created_at, metadata: parseMetadata(row.metadata_json) }; } function artifactToRow(artifact: NewArtifact): ArtifactRow { return { artifact_id: artifact.artifact_id, task_id: artifact.task_id, session_id: artifact.session_id ?? null, agent_id: artifact.agent_id, kind: artifact.kind, title: artifact.title ?? null, content: artifact.content ?? null, uri: artifact.uri ?? null, mime_type: artifact.mime_type ?? null, created_at: artifact.created_at, metadata_json: stringifyMetadata(artifact.metadata), idempotency_key: artifact.idempotency_key }; } function readTaskById(db: SqliteDatabase, taskId: string): Task | null { const taskRow = db.prepare("SELECT * FROM tasks WHERE task_id = ?").get(taskId) as TaskRow | undefined; if (taskRow == null) { return null; } const activeSessionRows = db .prepare("SELECT session_id FROM work_sessions WHERE task_id = ? AND status = 'active' ORDER BY started_at, session_id") .all(taskId) as Array<{ session_id: string }>; return taskFromRow( taskRow, activeSessionRows.map((row) => row.session_id) ); } class SqliteStore implements SqliteTrackerStore { constructor(private readonly db: SqliteDatabase) {} createTask(task: NewTask): Task { const row = taskToRow(task); this.db.prepare(` INSERT INTO tasks ( task_id, title, project_id, chat_id, status, created_at, completed_at, total_work_seconds, last_activity_at, idle_timeout_seconds, allowed_agent_ids_json, metadata_json ) VALUES ( @task_id, @title, @project_id, @chat_id, @status, @created_at, @completed_at, @total_work_seconds, @last_activity_at, @idle_timeout_seconds, @allowed_agent_ids_json, @metadata_json ) `).run(row); const createdTask = readTaskById(this.db, task.task_id); if (createdTask == null) { throw new Error("storage corruption: created task not found"); } return createdTask; } getTask(taskId: string): Task | null { return readTaskById(this.db, taskId); } listTasks(filter: TaskListFilter = {}): Task[] { const conditions: string[] = []; const params: Record<string, string> = {}; if (filter.project_id != null) { conditions.push("project_id = @project_id"); params.project_id = filter.project_id; } if (filter.chat_id != null) { conditions.push("chat_id = @chat_id"); params.chat_id = filter.chat_id; } const whereClause = conditions.length === 0 ? "" : ` WHERE ${conditions.join(" AND ")}`; const rows = this.db .prepare(`SELECT task_id FROM tasks${whereClause} ORDER BY created_at DESC, task_id DESC`) .all(params) as Array<{ task_id: string }>; return rows.flatMap((row) => { const task = readTaskById(this.db, row.task_id); return task == null ? [] : [task]; }); } updateTask(task: NewTask): Task { const row = taskToRow(task); const result = this.db.prepare(` UPDATE tasks SET title = @title, project_id = @project_id, chat_id = @chat_id, status = @status, created_at = @created_at, completed_at = @completed_at, total_work_seconds = @total_work_seconds, last_activity_at = @last_activity_at, idle_timeout_seconds = @idle_timeout_seconds, allowed_agent_ids_json = @allowed_agent_ids_json, metadata_json = @metadata_json WHERE task_id = @task_id `).run(row); if (result.changes === 0) { throw new Error(`task not found: ${task.task_id}`); } const updatedTask = readTaskById(this.db, task.task_id); if (updatedTask == null) { throw new Error("storage corruption: updated task not found"); } return updatedTask; } listTaskSessions(taskId: string): WorkSession[] { const rows = this.db .prepare("SELECT * FROM work_sessions WHERE task_id = ? ORDER BY started_at, session_id") .all(taskId) as WorkSessionRow[]; return rows.map(sessionFromRow); } getActiveSession(taskId: string, agentId: string): WorkSession | null { const row = this.db .prepare( "SELECT * FROM work_sessions WHERE task_id = ? AND agent_id = ? AND status = 'active' LIMIT 1" ) .get(taskId, agentId) as WorkSessionRow | undefined; return row == null ? null : sessionFromRow(row); } listActiveSessions(taskId: string): WorkSession[] { const rows = this.db .prepare("SELECT * FROM work_sessions WHERE task_id = ? AND status = 'active' ORDER BY started_at, session_id") .all(taskId) as WorkSessionRow[]; return rows.map(sessionFromRow); } createSession(session: WorkSession): WorkSession { const row = sessionToRow(session); this.db.prepare(` INSERT INTO work_sessions ( session_id, task_id, agent_id, status, started_at, ended_at, last_activity_at, start_reason, end_reason, started_by, ended_by, metadata_json ) VALUES ( @session_id, @task_id, @agent_id, @status, @started_at, @ended_at, @last_activity_at, @start_reason, @end_reason, @started_by, @ended_by, @metadata_json ) `).run(row); const createdSession = this.getSession(session.session_id); if (createdSession == null) { throw new Error("storage corruption: created session not found"); } return createdSession; } getSession(sessionId: string): WorkSession | null { const row = this.db.prepare("SELECT * FROM work_sessions WHERE session_id = ?").get(sessionId) as | WorkSessionRow | undefined; return row == null ? null : sessionFromRow(row); } updateSession(session: WorkSession): WorkSession { const row = sessionToRow(session); const result = this.db.prepare(` UPDATE work_sessions SET task_id = @task_id, agent_id = @agent_id, status = @status, started_at = @started_at, ended_at = @ended_at, last_activity_at = @last_activity_at, start_reason = @start_reason, end_reason = @end_reason, started_by = @started_by, ended_by = @ended_by, metadata_json = @metadata_json WHERE session_id = @session_id `).run(row); if (result.changes === 0) { throw new Error(`session not found: ${session.session_id}`); } const updatedSession = this.getSession(session.session_id); if (updatedSession == null) { throw new Error("storage corruption: updated session not found"); } return updatedSession; } createArtifact(artifact: NewArtifact): Artifact { const row = artifactToRow(artifact); this.db.prepare(` INSERT INTO artifacts ( artifact_id, task_id, session_id, agent_id, kind, title, content, uri, mime_type, created_at, metadata_json, idempotency_key ) VALUES ( @artifact_id, @task_id, @session_id, @agent_id, @kind, @title, @content, @uri, @mime_type, @created_at, @metadata_json, @idempotency_key ) `).run(row); const createdArtifact = this.getArtifact(artifact.artifact_id); if (createdArtifact == null) { throw new Error("storage corruption: created artifact not found"); } return createdArtifact; } getArtifact(artifactId: string): Artifact | null { const row = this.db.prepare("SELECT * FROM artifacts WHERE artifact_id = ?").get(artifactId) as | ArtifactRow | undefined; return row == null ? null : artifactFromRow(row); } findArtifactByIdempotencyKey(idempotencyKey: string): Artifact | null { const row = this.db .prepare("SELECT * FROM artifacts WHERE idempotency_key = ?") .get(idempotencyKey) as ArtifactRow | undefined; return row == null ? null : artifactFromRow(row); } listArtifacts(taskId: string): Artifact[] { const rows = this.db .prepare("SELECT * FROM artifacts WHERE task_id = ? ORDER BY created_at, artifact_id") .all(taskId) as ArtifactRow[]; return rows.map(artifactFromRow); } transaction<T>(fn: () => T): T { return this.db.transaction(fn)(); } close(): void { this.db.close(); } } export function createSqliteStore(dbPath: string): SqliteTrackerStore { ensureParentDirectory(dbPath); const db = new Database(dbPath); db.pragma("foreign_keys = ON"); runMigrations(db); return new SqliteStore(db); } export type { TrackerStore };