d3856cb52a
add a comprehensive gitignore file to exclude unnecessary files like dependencies, runtime data, uploads, local configs and debug temporary files
123 lines
4.2 KiB
JavaScript
123 lines
4.2 KiB
JavaScript
// ============================================================
|
||
// 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
|
||
} |