// ============================================================ // data_gateway/store.js - 数据存储模块 // 职责:动态创建/管理每个 skill 的独立 SQLite 数据库 // 表结构:id | source | data(JSON) | created_at // 支持 mode: "append"(增量追加)和 "latest"(覆盖最新) // ============================================================ const Database = require('better-sqlite3') const fs = require('fs') const path = require('path') const DATA_DIR = path.join(process.cwd(), 'data', 'ingest') if (!fs.existsSync(DATA_DIR)) fs.mkdirSync(DATA_DIR, { recursive: true }) // 已打开的数据库连接缓存 const dbCache = new Map() const getDbPath = (skillId) => path.join(DATA_DIR, `${skillId}.db`) const getOrOpenDb = (skillId) => { if (dbCache.has(skillId)) return dbCache.get(skillId) const dbPath = getDbPath(skillId) const db = new Database(dbPath) // 启用 WAL 模式提高并发写入性能 db.pragma('journal_mode = WAL') db.pragma('synchronous = NORMAL') db.exec( `CREATE TABLE IF NOT EXISTS records ( id INTEGER PRIMARY KEY AUTOINCREMENT, source TEXT NOT NULL DEFAULT '', data TEXT NOT NULL, created_at TEXT NOT NULL DEFAULT (datetime('now')) )` ) db.exec(`CREATE INDEX IF NOT EXISTS idx_records_created ON records(created_at DESC)`) dbCache.set(skillId, db) return db } // 写入一条记录 const insert = (skillId, { source = '', data = {} }) => { const db = getOrOpenDb(skillId) const dataStr = JSON.stringify(data) const info = db.prepare(`INSERT INTO records (source, data) VALUES (?, ?)`).run( String(source || ''), dataStr ) return { id: info.lastInsertRowid, source, data, created_at: new Date().toISOString() } } // 覆盖写入(mode: "latest")- 先删后插,保证只有一条 const upsertLatest = (skillId, { source = '', data = {} }) => { const db = getOrOpenDb(skillId) const dataStr = JSON.stringify(data) const tx = db.transaction(() => { db.prepare(`DELETE FROM records`).run() return db.prepare(`INSERT INTO records (source, data) VALUES (?, ?)`).run( String(source || ''), dataStr ) }) const info = tx() return { id: info.lastInsertRowid, source, data, created_at: new Date().toISOString() } } // 查询最新一条记录 const queryLatest = (skillId) => { const db = getOrOpenDb(skillId) const row = db.prepare(`SELECT * FROM records ORDER BY id DESC LIMIT 1`).get() if (!row) return null try { row.data = JSON.parse(row.data) } catch { /* keep as string */ } return row } // 查询记录列表(分页) const queryList = (skillId, { limit = 100, offset = 0, source = null } = {}) => { const db = getOrOpenDb(skillId) let sql = `SELECT * FROM records` const params = [] if (source) { sql += ` WHERE source = ?` params.push(String(source)) } sql += ` ORDER BY id DESC LIMIT ? OFFSET ?` params.push(Math.min(Number(limit) || 100, 1000), Math.max(0, Number(offset) || 0)) const rows = db.prepare(sql).all(...params) for (const row of rows) { try { row.data = JSON.parse(row.data) } catch { /* keep as string */ } } const total = db.prepare(`SELECT COUNT(*) as cnt FROM records`).get().cnt return { rows, total } } // 按时间范围查询 const queryByTimeRange = (skillId, { start, end, limit = 100 } = {}) => { const db = getOrOpenDb(skillId) let sql = `SELECT * FROM records WHERE 1=1` const params = [] if (start) { sql += ` AND created_at >= ?`; params.push(String(start)) } if (end) { sql += ` AND created_at <= ?`; params.push(String(end)) } sql += ` ORDER BY id DESC LIMIT ?` params.push(Math.min(Number(limit) || 100, 1000)) const rows = db.prepare(sql).all(...params) for (const row of rows) { try { row.data = JSON.parse(row.data) } catch { /* keep as string */ } } return { rows, total: rows.length } } // 获取记录总数 const count = (skillId) => { const db = getOrOpenDb(skillId) return db.prepare(`SELECT COUNT(*) as cnt FROM records`).get().cnt } // 关闭所有数据库连接 const closeAll = () => { for (const [skillId, db] of dbCache) { try { db.close() } catch {} dbCache.delete(skillId) } } module.exports = { insert, upsertLatest, queryLatest, queryList, queryByTimeRange, count, closeAll, getDbPath, getOrOpenDb }