chore: add and update gitignore rules for project files
add a comprehensive gitignore file to exclude unnecessary files like dependencies, runtime data, uploads, local configs and debug temporary files
This commit is contained in:
@@ -0,0 +1,85 @@
|
||||
// ============================================================
|
||||
// data_gateway/auth.js - 对外 API 鉴权模块
|
||||
// 职责:校验 X-API-Id + X-API-Key,匹配 skill config.json 中的哈希
|
||||
// ============================================================
|
||||
const crypto = require('crypto')
|
||||
const fs = require('fs')
|
||||
const path = require('path')
|
||||
|
||||
const SKILLS_DIR = path.join(__dirname, 'skills')
|
||||
|
||||
// 缓存已加载的 skill 配置
|
||||
const skillConfigCache = new Map()
|
||||
|
||||
const loadSkillConfig = (skillId) => {
|
||||
if (skillConfigCache.has(skillId)) return skillConfigCache.get(skillId)
|
||||
const cfgPath = path.join(SKILLS_DIR, skillId, 'config.json')
|
||||
if (!fs.existsSync(cfgPath)) return null
|
||||
try {
|
||||
const cfg = JSON.parse(fs.readFileSync(cfgPath, 'utf-8'))
|
||||
skillConfigCache.set(skillId, cfg)
|
||||
return cfg
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
// 刷新缓存(用于配置热更新)
|
||||
const refreshCache = (skillId) => {
|
||||
skillConfigCache.delete(skillId)
|
||||
return loadSkillConfig(skillId)
|
||||
}
|
||||
|
||||
// 枚举所有已启用的 skill
|
||||
const listEnabledSkills = () => {
|
||||
const result = []
|
||||
if (!fs.existsSync(SKILLS_DIR)) return result
|
||||
const dirs = fs.readdirSync(SKILLS_DIR, { withFileTypes: true })
|
||||
for (const d of dirs) {
|
||||
if (!d.isDirectory()) continue
|
||||
const cfg = loadSkillConfig(d.name)
|
||||
if (cfg && cfg.enabled !== false) {
|
||||
result.push({ id: cfg.id, name: cfg.name })
|
||||
}
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
// 校验外部 API 请求
|
||||
const verifyApiKey = (req) => {
|
||||
const apiId = String(req.headers['x-api-id'] || '').trim()
|
||||
const apiKey = String(req.headers['x-api-key'] || '').trim()
|
||||
if (!apiId || !apiKey) return { ok: false, error: 'missing api id or key' }
|
||||
|
||||
const cfg = loadSkillConfig(apiId)
|
||||
if (!cfg) return { ok: false, error: 'skill not found' }
|
||||
if (cfg.enabled === false) return { ok: false, error: 'skill disabled' }
|
||||
|
||||
const expectedHash = String(cfg.api_key_hash || '')
|
||||
if (!expectedHash) return { ok: false, error: 'skill not configured' }
|
||||
|
||||
const actualHash = 'sha256:' + crypto.createHash('sha256').update(apiKey).digest('hex')
|
||||
if (actualHash !== expectedHash) return { ok: false, error: 'invalid api key' }
|
||||
|
||||
return { ok: true, skill: cfg }
|
||||
}
|
||||
|
||||
// 校验内部 Tool 读取请求
|
||||
const verifyInternalToken = (req) => {
|
||||
const token = String(req.headers['x-internal-token'] || '').trim()
|
||||
const skillId = String(req.headers['x-skill-id'] || req.params?.skillId || '').trim()
|
||||
if (!token || !skillId) return { ok: false, error: 'missing token or skill id' }
|
||||
|
||||
const cfg = loadSkillConfig(skillId)
|
||||
if (!cfg) return { ok: false, error: 'skill not found' }
|
||||
|
||||
const expectedHash = String(cfg.read_token_hash || '')
|
||||
if (!expectedHash) return { ok: false, error: 'skill not configured for internal read' }
|
||||
|
||||
const actualHash = 'sha256:' + crypto.createHash('sha256').update(token).digest('hex')
|
||||
if (actualHash !== expectedHash) return { ok: false, error: 'invalid internal token' }
|
||||
|
||||
return { ok: true, skill: cfg }
|
||||
}
|
||||
|
||||
module.exports = { loadSkillConfig, refreshCache, listEnabledSkills, verifyApiKey, verifyInternalToken }
|
||||
@@ -0,0 +1,125 @@
|
||||
// ============================================================
|
||||
// data_gateway/index.js - 对外 API 数据网关主入口
|
||||
// 暴露两个路由:
|
||||
// POST /api/v1/ingest/:skillId 外部数据写入(公网,需 API Key)
|
||||
// GET /api/v1/data/:skillId/query 内部数据读取(仅 localhost,需 Internal Token)
|
||||
// ============================================================
|
||||
const { logJSON } = require('../logger')
|
||||
const { verifyApiKey, verifyInternalToken, listEnabledSkills } = require('./auth')
|
||||
const { insert, upsertLatest, queryLatest, queryList, queryByTimeRange, count } = require('./store')
|
||||
const { loadSkillConfig } = require('./auth')
|
||||
const { bindRoutes: bindExternalStorageRoutes } = require('./skills/external_storage')
|
||||
|
||||
const setNoCache = (res) => {
|
||||
try {
|
||||
res.set('Cache-Control', 'no-store, no-cache, must-revalidate, proxy-revalidate')
|
||||
res.set('Pragma', 'no-cache')
|
||||
res.set('Expires', '0')
|
||||
} catch {}
|
||||
}
|
||||
|
||||
const bindRoutes = (app) => {
|
||||
// ============================================================
|
||||
// 外部存储 skill(第三方 OSS 读写)
|
||||
// 路由挂载在 /api/v1/skills/external_storage/*
|
||||
// ============================================================
|
||||
try {
|
||||
bindExternalStorageRoutes(app)
|
||||
} catch (e) {
|
||||
console.error(`[data_gateway] Failed to bind external_storage routes: ${e.message}`)
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 外部写入接口(公网暴露)
|
||||
// POST /api/v1/ingest/:skillId
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// Body: { "source": "...", "data": { ... } }
|
||||
// ============================================================
|
||||
app.post('/api/v1/ingest/:skillId', (req, res) => {
|
||||
setNoCache(res)
|
||||
const skillId = String(req.params.skillId || '').trim()
|
||||
if (!skillId) return res.status(400).json({ ok: false, error: 'missing skill id' })
|
||||
|
||||
const auth = verifyApiKey(req)
|
||||
if (!auth.ok) {
|
||||
try { logJSON('data_gateway.auth.fail', { skillId, error: auth.error, ip: req.ip }, 'data_gateway') } catch {}
|
||||
return res.status(401).json({ ok: false, error: auth.error })
|
||||
}
|
||||
|
||||
const { source, data } = req.body || {}
|
||||
if (data === undefined && req.body && typeof req.body === 'object') {
|
||||
// 如果没传 data 字段,将整个 body 视为 data
|
||||
const bodyData = { ...req.body }
|
||||
delete bodyData.source
|
||||
const record = insert(skillId, { source: String(source || ''), data: Object.keys(bodyData).length > 0 ? bodyData : req.body })
|
||||
try { logJSON('data_gateway.ingest.ok', { skillId, source, recordId: record.id }, 'data_gateway') } catch {}
|
||||
return res.json({ ok: true, id: record.id, created_at: record.created_at })
|
||||
}
|
||||
|
||||
const record = insert(skillId, { source: String(source || ''), data: data || {} })
|
||||
try { logJSON('data_gateway.ingest.ok', { skillId, source, recordId: record.id }, 'data_gateway') } catch {}
|
||||
return res.json({ ok: true, id: record.id, created_at: record.created_at })
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 内部读取接口(仅 localhost)
|
||||
// GET /api/v1/data/:skillId/query?mode=latest|list&limit=100&offset=0
|
||||
// Headers: X-Internal-Token, X-Skill-Id
|
||||
// ============================================================
|
||||
app.get('/api/v1/data/:skillId/query', (req, res) => {
|
||||
setNoCache(res)
|
||||
const skillId = String(req.params.skillId || '').trim()
|
||||
if (!skillId) return res.status(400).json({ ok: false, error: 'missing skill id' })
|
||||
|
||||
const auth = verifyInternalToken(req)
|
||||
if (!auth.ok) return res.status(401).json({ ok: false, error: auth.error })
|
||||
|
||||
const mode = String(req.query.mode || 'latest').trim()
|
||||
const limit = parseInt(req.query.limit, 10) || 100
|
||||
const offset = parseInt(req.query.offset, 10) || 0
|
||||
const source = req.query.source || null
|
||||
const start = req.query.start || null
|
||||
const end = req.query.end || null
|
||||
|
||||
try {
|
||||
if (mode === 'list') {
|
||||
if (start || end) {
|
||||
const result = queryByTimeRange(skillId, { start, end, limit })
|
||||
return res.json({ ok: true, ...result })
|
||||
}
|
||||
const result = queryList(skillId, { limit, offset, source })
|
||||
return res.json({ ok: true, ...result })
|
||||
}
|
||||
|
||||
// mode === 'latest'(默认)
|
||||
const row = queryLatest(skillId)
|
||||
if (!row) return res.json({ ok: true, rows: [], total: 0 })
|
||||
return res.json({ ok: true, rows: [row], total: 1 })
|
||||
} catch (e) {
|
||||
try { logJSON('data_gateway.query.error', { skillId, error: String(e.message || e) }, 'data_gateway') } catch {}
|
||||
return res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 管理接口:列出所有已启用的 skill(仅 localhost)
|
||||
// GET /api/v1/data/skills
|
||||
// ============================================================
|
||||
app.get('/api/v1/data/skills', (req, res) => {
|
||||
setNoCache(res)
|
||||
const skills = listEnabledSkills()
|
||||
const result = skills.map(s => {
|
||||
const cfg = loadSkillConfig(s.id)
|
||||
return {
|
||||
id: s.id,
|
||||
name: s.name,
|
||||
mode: String(cfg && cfg.mode || 'append'),
|
||||
record_count: count(s.id),
|
||||
db_path: require('./store').getDbPath(s.id)
|
||||
}
|
||||
})
|
||||
res.json({ ok: true, skills: result })
|
||||
})
|
||||
}
|
||||
|
||||
module.exports = { bindRoutes, listEnabledSkills, loadSkillConfig }
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"_comment": "外部存储 skill - 供第三方程序读写 OSS external_storage 区域",
|
||||
"id": "external_storage",
|
||||
"name": "外部存储",
|
||||
"enabled": true,
|
||||
"api_key_comment": "这个发给对方,API Key,用于上传文件到外部存储区域",
|
||||
"api_key": "4a211084d059Pk0C6t8dlk3ew61d46",
|
||||
"api_key_hash": "sha256:17eb4c895064b955d80d007b15af992d7bb7bd1525ce5726f2f63f06229857c9",
|
||||
"oss_bucket": "android-o1-images",
|
||||
"oss_endpoint_backup1": "oss-accelerate.aliyuncs.com 传输加速域名(全地域上传下载加速)",
|
||||
"oss_endpoint_backup2": "oss-cn-shanghai.aliyuncs.com 上海区域",
|
||||
"oss_endpoint": "oss-accelerate.aliyuncs.com",
|
||||
"oss_prefix": "550e8400-e29b-41d4-a716-446655442082/external_storage"
|
||||
}
|
||||
@@ -0,0 +1,322 @@
|
||||
// ============================================================
|
||||
// external_storage/index.js - 外部存储 skill 路由入口
|
||||
// 提供第三方程序读写 OSS external_storage 区域的 API
|
||||
// ============================================================
|
||||
const express = require('express')
|
||||
const multer = require('multer')
|
||||
const path = require('path')
|
||||
const archiver = require('archiver')
|
||||
const { createOssClient } = require('./oss_client')
|
||||
const { sanitizePath, createSizeLimitMiddleware } = require('./middleware')
|
||||
const { loadSkillConfig } = require('../../auth')
|
||||
const { verifyApiKey } = require('../../auth')
|
||||
const { logJSON } = require('../../../logger')
|
||||
|
||||
const MAX_UPLOAD_BYTES = 50 * 1024 * 1024 // 50MB
|
||||
|
||||
// 内存存储上传文件
|
||||
const upload = multer({
|
||||
storage: multer.memoryStorage(),
|
||||
limits: { fileSize: MAX_UPLOAD_BYTES, files: 20 }
|
||||
})
|
||||
|
||||
// MIME 类型推断(支持中文文件名 UTF-8)
|
||||
const contentTypeByExt = (filename) => {
|
||||
const ext = path.extname(filename).toLowerCase().slice(1)
|
||||
const map = {
|
||||
jpg: 'image/jpeg', jpeg: 'image/jpeg', png: 'image/png', gif: 'image/gif', webp: 'image/webp', svg: 'image/svg+xml', ico: 'image/x-icon', bmp: 'image/bmp', avif: 'image/avif',
|
||||
txt: 'text/plain; charset=utf-8', md: 'text/plain; charset=utf-8', html: 'text/html; charset=utf-8',
|
||||
js: 'application/javascript', ts: 'application/javascript', css: 'text/css', json: 'application/json',
|
||||
xml: 'text/xml', log: 'text/plain; charset=utf-8', csv: 'text/csv',
|
||||
pdf: 'application/pdf',
|
||||
doc: 'application/msword', docx: 'application/vnd.openxmlformats-officedocument.wordprocessingml.document',
|
||||
xls: 'application/vnd.ms-excel', xlsx: 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
|
||||
ppt: 'application/vnd.ms-powerpoint', pptx: 'application/vnd.openxmlformats-officedocument.presentationml.presentation',
|
||||
mp4: 'video/mp4', mov: 'video/quicktime', avi: 'video/x-msvideo', webm: 'video/webm', mkv: 'video/x-matroska',
|
||||
mp3: 'audio/mpeg', wav: 'audio/wav', ogg: 'audio/ogg', flac: 'audio/flac',
|
||||
zip: 'application/zip', rar: 'application/x-rar-compressed', '7z': 'application/x-7z-compressed',
|
||||
tar: 'application/x-tar', gz: 'application/gzip',
|
||||
exe: 'application/x-msdownload', dmg: 'application/x-apple-diskimage',
|
||||
font: 'font/woff', woff: 'font/woff', woff2: 'font/woff2', ttf: 'font/ttf', otf: 'font/otf'
|
||||
}
|
||||
return map[ext] || 'application/octet-stream'
|
||||
}
|
||||
|
||||
// 安全文件名清理(保留中文,移除危险字符)
|
||||
const safeName = (name) => {
|
||||
let s = String(name || 'untitled')
|
||||
// 尝试修复编码(UTF-8 被误读为 latin1 的情况)
|
||||
try {
|
||||
const recovered = Buffer.from(s, 'latin1').toString('utf8')
|
||||
if (recovered !== s && !/\ufffd/.test(recovered)) s = recovered
|
||||
} catch {}
|
||||
// 移除危险字符,但保留中文、英文、数字、常见符号
|
||||
return s.replace(/[<>\":/\\|?*\x00-\x1F]/g, '_').trim() || 'untitled'
|
||||
}
|
||||
|
||||
// 设置无缓存响应头
|
||||
const setNoCache = (res) => {
|
||||
try {
|
||||
res.set('Cache-Control', 'no-store, no-cache, must-revalidate, proxy-revalidate')
|
||||
res.set('Pragma', 'no-cache')
|
||||
res.set('Expires', '0')
|
||||
} catch {}
|
||||
}
|
||||
|
||||
/**
|
||||
* 绑定路由到 Express 应用
|
||||
* @param {express.Application} app
|
||||
*/
|
||||
const bindRoutes = (app) => {
|
||||
const skillId = 'external_storage'
|
||||
const skillConfig = loadSkillConfig(skillId)
|
||||
if (!skillConfig) {
|
||||
console.warn(`[external_storage] config.json not found, skill disabled`)
|
||||
return
|
||||
}
|
||||
|
||||
let ossClient
|
||||
try {
|
||||
ossClient = createOssClient(skillConfig)
|
||||
} catch (e) {
|
||||
console.error(`[external_storage] Failed to init OSS client: ${e.message}`)
|
||||
return
|
||||
}
|
||||
|
||||
const router = express.Router()
|
||||
|
||||
// API Key 鉴权中间件
|
||||
router.use((req, res, next) => {
|
||||
const auth = verifyApiKey(req)
|
||||
if (!auth.ok) {
|
||||
try { logJSON('external_storage.auth.fail', { error: auth.error, ip: req.ip }, 'external_storage') } catch {}
|
||||
return res.status(401).json({ ok: false, error: auth.error })
|
||||
}
|
||||
next()
|
||||
})
|
||||
|
||||
// 全局无缓存
|
||||
router.use((req, res, next) => {
|
||||
setNoCache(res)
|
||||
next()
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 上传文件
|
||||
// POST /api/v1/skills/external_storage/upload?path=subdir/
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// Body: multipart/form-data, field: "files"
|
||||
// ============================================================
|
||||
router.post('/upload', createSizeLimitMiddleware(MAX_UPLOAD_BYTES), upload.array('files', 20), async (req, res) => {
|
||||
try {
|
||||
const basePath = sanitizePath(req.query.path || '')
|
||||
if (!req.files || !Array.isArray(req.files) || req.files.length === 0) {
|
||||
return res.status(400).json({ ok: false, error: '未收到文件' })
|
||||
}
|
||||
|
||||
const results = []
|
||||
for (const f of req.files) {
|
||||
const originalName = safeName(f.originalname || 'file')
|
||||
const keyPath = basePath ? `${basePath}/${originalName}` : originalName
|
||||
const contentType = f.mimetype || contentTypeByExt(originalName)
|
||||
|
||||
const { key, raw_url } = await ossClient.putObject(keyPath, f.buffer, contentType)
|
||||
results.push({
|
||||
path: keyPath,
|
||||
oss_key: key,
|
||||
raw_url,
|
||||
name: originalName,
|
||||
size: f.size,
|
||||
content_type: contentType
|
||||
})
|
||||
}
|
||||
|
||||
res.json({ ok: true, count: results.length, files: results })
|
||||
} catch (e) {
|
||||
if (String(e.message || '').includes('File too large')) {
|
||||
return res.status(413).json({ ok: false, error: '文件超过 50MB 上限' })
|
||||
}
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 下载文件(代理流)
|
||||
// GET /api/v1/skills/external_storage/download?path=xxx
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.get('/download', async (req, res) => {
|
||||
try {
|
||||
const filePath = sanitizePath(req.query.path)
|
||||
if (!filePath) return res.status(400).json({ ok: false, error: '缺少 path 参数' })
|
||||
|
||||
const displayName = path.basename(filePath)
|
||||
const { stream, headers } = await ossClient.getObjectStream(filePath)
|
||||
|
||||
// 使用 RFC 5987 格式支持 UTF-8 文件名
|
||||
res.set('Content-Disposition', `attachment; filename*=UTF-8''${encodeURIComponent(displayName)}`)
|
||||
if (headers['content-type']) res.set('Content-Type', headers['content-type'])
|
||||
if (headers['content-length']) res.set('Content-Length', headers['content-length'])
|
||||
|
||||
stream.pipe(res)
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 获取原始 URL
|
||||
// GET /api/v1/skills/external_storage/raw_url?path=xxx
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.get('/raw_url', async (req, res) => {
|
||||
try {
|
||||
const filePath = sanitizePath(req.query.path)
|
||||
if (!filePath) return res.status(400).json({ ok: false, error: '缺少 path 参数' })
|
||||
|
||||
const raw_url = ossClient.getRawUrl(filePath)
|
||||
res.json({ ok: true, path: filePath, raw_url })
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 列出文件
|
||||
// GET /api/v1/skills/external_storage/list?prefix=subdir/&recursive=true|false
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.get('/list', async (req, res) => {
|
||||
try {
|
||||
const prefix = sanitizePath(req.query.prefix || '')
|
||||
const recursive = req.query.recursive === 'true'
|
||||
|
||||
let result
|
||||
if (recursive) {
|
||||
const files = await ossClient.listRecursive(prefix)
|
||||
result = { files, folders: [] }
|
||||
} else {
|
||||
result = await ossClient.listObjects(prefix)
|
||||
}
|
||||
|
||||
res.json({ ok: true, prefix, recursive, ...result })
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 删除文件
|
||||
// DELETE /api/v1/skills/external_storage/delete?path=xxx
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.delete('/delete', async (req, res) => {
|
||||
try {
|
||||
const filePath = sanitizePath(req.query.path)
|
||||
if (!filePath) return res.status(400).json({ ok: false, error: '缺少 path 参数' })
|
||||
|
||||
await ossClient.deleteObject(filePath)
|
||||
res.json({ ok: true, path: filePath })
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 文件夹打包下载
|
||||
// GET /api/v1/skills/external_storage/folder/download?path=subdir/
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.get('/folder/download', async (req, res) => {
|
||||
try {
|
||||
const folderPath = sanitizePath(req.query.path)
|
||||
if (!folderPath) return res.status(400).json({ ok: false, error: '缺少 path 参数' })
|
||||
|
||||
const folderKey = folderPath.endsWith('/') ? folderPath : folderPath + '/'
|
||||
const folderName = path.basename(folderKey.replace(/\/+$/, '')) || 'folder'
|
||||
|
||||
const files = await ossClient.listRecursive(folderKey)
|
||||
|
||||
res.set('Content-Type', 'application/zip')
|
||||
res.set('Content-Disposition', `attachment; filename*=UTF-8''${encodeURIComponent(folderName + '.zip')}`)
|
||||
|
||||
const archive = archiver('zip', { zlib: { level: 6 } })
|
||||
archive.on('error', err => {
|
||||
try { res.status(500).end() } catch {}
|
||||
})
|
||||
archive.pipe(res)
|
||||
|
||||
const prefixClean = folderKey.replace(/\/+$/, '') + '/'
|
||||
for (const f of files) {
|
||||
try {
|
||||
// 提取相对于文件夹前缀的路径
|
||||
const fullRel = f.key
|
||||
const relPart = fullRel.slice(prefixClean.length)
|
||||
const { stream } = await ossClient.getObjectStream(relPart)
|
||||
archive.append(stream, { name: relPart })
|
||||
} catch {}
|
||||
}
|
||||
|
||||
archive.finalize()
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// ============================================================
|
||||
// 获取文件信息
|
||||
// GET /api/v1/skills/external_storage/info?path=xxx
|
||||
// Headers: X-API-Id, X-API-Key
|
||||
// ============================================================
|
||||
router.get('/info', async (req, res) => {
|
||||
try {
|
||||
const filePath = sanitizePath(req.query.path)
|
||||
if (!filePath) return res.status(400).json({ ok: false, error: '缺少 path 参数' })
|
||||
|
||||
const isFolder = filePath.endsWith('/')
|
||||
const fullKey = ossClient.fullKey(filePath)
|
||||
const raw_url = ossClient.getRawUrl(filePath)
|
||||
|
||||
let size = 0
|
||||
let lastModified = 0
|
||||
|
||||
if (!isFolder) {
|
||||
try {
|
||||
const parent = filePath.includes('/') ? filePath.split('/').slice(0, -1).join('/') : ''
|
||||
const result = await ossClient.listObjects(parent)
|
||||
const targetKey = fullKey
|
||||
const hit = result.files.find(f => f.key === targetKey)
|
||||
if (hit) {
|
||||
size = hit.size
|
||||
lastModified = hit.lastModified
|
||||
}
|
||||
} catch {}
|
||||
} else {
|
||||
try {
|
||||
const files = await ossClient.listRecursive(filePath)
|
||||
size = files.reduce((s, f) => s + (Number(f.size) || 0), 0)
|
||||
} catch {}
|
||||
}
|
||||
|
||||
res.json({
|
||||
ok: true,
|
||||
info: {
|
||||
path: filePath,
|
||||
oss_key: fullKey,
|
||||
raw_url,
|
||||
isFolder,
|
||||
size,
|
||||
lastModified: lastModified ? new Date(lastModified).toISOString() : null
|
||||
}
|
||||
})
|
||||
} catch (e) {
|
||||
res.status(500).json({ ok: false, error: String(e.message || e) })
|
||||
}
|
||||
})
|
||||
|
||||
// 挂载路由
|
||||
app.use('/api/v1/skills/external_storage', router)
|
||||
console.log('[external_storage] Routes mounted at /api/v1/skills/external_storage')
|
||||
}
|
||||
|
||||
module.exports = { bindRoutes }
|
||||
@@ -0,0 +1,83 @@
|
||||
// ============================================================
|
||||
// external_storage/middleware.js - 路径安全校验中间件
|
||||
// 职责:防止路径遍历攻击,确保所有操作限定在 external_storage 范围内
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* 校验并清理路径
|
||||
* - 拒绝包含 ../ 的路径遍历
|
||||
* - 拒绝绝对路径
|
||||
* - 拒绝控制字符
|
||||
* - 统一使用 UTF-8 编码,正确处理中文
|
||||
* @param {string} rawPath - 原始路径参数
|
||||
* @returns {string} 清理后的安全路径
|
||||
* @throws {Error} 如果路径不合法
|
||||
*/
|
||||
const sanitizePath = (rawPath) => {
|
||||
if (rawPath === undefined || rawPath === null) return ''
|
||||
|
||||
// 确保使用 UTF-8 字符串
|
||||
let p = String(rawPath)
|
||||
|
||||
// 拒绝控制字符(除了正常的换行和制表符)
|
||||
if (/[\x00-\x08\x0B\x0C\x0E-\x1F\x7F]/.test(p)) {
|
||||
throw new Error('路径包含非法控制字符')
|
||||
}
|
||||
|
||||
// 拒绝路径遍历
|
||||
if (p.includes('..')) {
|
||||
throw new Error('路径不允许包含 ".."')
|
||||
}
|
||||
|
||||
// 拒绝绝对路径
|
||||
if (p.startsWith('/') || p.startsWith('\\')) {
|
||||
throw new Error('路径不允许以 "/" 或 "\\" 开头')
|
||||
}
|
||||
|
||||
// 统一斜杠
|
||||
p = p.replace(/[\\]+/g, '/')
|
||||
|
||||
// 去除连续斜杠
|
||||
p = p.replace(/\/+/g, '/')
|
||||
|
||||
// 去除首尾斜杠
|
||||
p = p.replace(/^\/+|\/+$/g, '')
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
/**
|
||||
* 校验路径是否在允许的前缀范围内
|
||||
* @param {string} ossKey - 完整 OSS key
|
||||
* @param {string} allowedPrefix - 允许的前缀
|
||||
* @returns {boolean}
|
||||
*/
|
||||
const isPathWithinPrefix = (ossKey, allowedPrefix) => {
|
||||
const key = String(ossKey || '')
|
||||
const prefix = String(allowedPrefix || '').replace(/\/+$/, '')
|
||||
return key === prefix || key.startsWith(prefix + '/')
|
||||
}
|
||||
|
||||
/**
|
||||
* 文件大小校验中间件
|
||||
* @param {number} maxBytes - 最大字节数
|
||||
* @returns {Function} express 中间件
|
||||
*/
|
||||
const createSizeLimitMiddleware = (maxBytes = 50 * 1024 * 1024) => {
|
||||
return (req, res, next) => {
|
||||
const contentLength = parseInt(req.headers['content-length'] || '0', 10)
|
||||
if (contentLength > maxBytes) {
|
||||
return res.status(413).json({
|
||||
ok: false,
|
||||
error: `文件超过大小限制 (${Math.round(maxBytes / 1024 / 1024)}MB)`
|
||||
})
|
||||
}
|
||||
next()
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
sanitizePath,
|
||||
isPathWithinPrefix,
|
||||
createSizeLimitMiddleware
|
||||
}
|
||||
@@ -0,0 +1,321 @@
|
||||
// ============================================================
|
||||
// external_storage/oss_client.js - 阿里云 OSS 客户端封装
|
||||
// 职责:提供上传、下载、删除、列表、获取URL等基础操作
|
||||
// 所有操作限定在 config.json 配置的 oss_prefix 范围内
|
||||
// ============================================================
|
||||
const crypto = require('crypto')
|
||||
const https = require('https')
|
||||
const http = require('http')
|
||||
const fs = require('fs')
|
||||
const path = require('path')
|
||||
const { URLSearchParams } = require('url')
|
||||
|
||||
/**
|
||||
* 初始化 OSS 客户端配置
|
||||
* @param {Object} skillConfig - skill 的 config.json 内容
|
||||
* @returns {Object} 客户端实例
|
||||
*/
|
||||
const createOssClient = (skillConfig) => {
|
||||
const bucket = skillConfig.oss_bucket || 'android-o1-images'
|
||||
const endpoint = skillConfig.oss_endpoint || 'oss-cn-shanghai.aliyuncs.com'
|
||||
const prefix = (skillConfig.oss_prefix || '550e8400-e29b-41d4-a716-446655442082/external_storage').replace(/\/+$/, '')
|
||||
|
||||
// 从 process.env 读取(index.js 启动时已从 Toolbox_local_creds.env.local 加载)
|
||||
let accessKeyId = process.env.EXTERNAL_STORAGE_OSS_ACCESS_KEY_ID || ''
|
||||
let accessKeySecret = process.env.EXTERNAL_STORAGE_OSS_ACCESS_KEY_SECRET || ''
|
||||
|
||||
if (!accessKeyId || !accessKeySecret) {
|
||||
throw new Error('OSS credentials not found. Please set oss_access_key_id and oss_access_key_secret in config.json')
|
||||
}
|
||||
|
||||
/**
|
||||
* 规范化 OSS key:处理斜杠、去除危险字符
|
||||
*/
|
||||
const normalizeKey = (relativePath) => {
|
||||
let k = String(relativePath || '').replace(/[\\]+/g, '/').replace(/\/+/g, '/')
|
||||
if (k.startsWith('/')) k = k.slice(1)
|
||||
return k
|
||||
}
|
||||
|
||||
/**
|
||||
* 拼接完整 OSS key(自动加前缀)
|
||||
*/
|
||||
const fullKey = (relativePath) => {
|
||||
const rel = normalizeKey(relativePath)
|
||||
return rel ? `${prefix}/${rel}` : prefix
|
||||
}
|
||||
|
||||
/**
|
||||
* URL 编码 OSS key(用于 HTTP 路径)
|
||||
*/
|
||||
const encodeKey = (key) => String(key || '').split('/').map(s => encodeURIComponent(s)).join('/')
|
||||
|
||||
/**
|
||||
* 构建 OSS 资源标识
|
||||
*/
|
||||
const ossResource = (key) => `/${bucket}/${String(key || '').replace(/^\/+/, '')}`
|
||||
|
||||
/**
|
||||
* OSS 请求签名(HMAC-SHA1)
|
||||
*/
|
||||
const signRequest = ({ method, contentType = '', date, key, extraHeaders = {} }) => {
|
||||
const canonicalHeaders = Object.keys(extraHeaders)
|
||||
.filter(k => k && k.toLowerCase().startsWith('x-oss-'))
|
||||
.sort((a, b) => a.toLowerCase().localeCompare(b.toLowerCase()))
|
||||
.map(k => `${k.toLowerCase()}:${String(extraHeaders[k] || '').trim()}`)
|
||||
.join('\n')
|
||||
const resource = ossResource(key)
|
||||
const stringToSign = `${method.toUpperCase()}\n\n${contentType}\n${date}\n${canonicalHeaders ? canonicalHeaders + '\n' : ''}${resource}`
|
||||
return crypto.createHmac('sha1', accessKeySecret).update(stringToSign).digest('base64')
|
||||
}
|
||||
|
||||
/**
|
||||
* 发起 HTTP 请求
|
||||
*/
|
||||
const httpRequest = (options, body = null, timeoutMs = 60000) => new Promise((resolve, reject) => {
|
||||
const lib = options.protocol === 'https:' ? https : http
|
||||
const req = lib.request(options, (res) => {
|
||||
const chunks = []
|
||||
res.on('data', chunk => chunks.push(chunk))
|
||||
res.on('end', () => {
|
||||
const buf = Buffer.concat(chunks)
|
||||
resolve({ statusCode: res.statusCode, headers: res.headers, body: buf })
|
||||
})
|
||||
})
|
||||
req.on('error', reject)
|
||||
req.setTimeout(timeoutMs, () => req.destroy(new Error('request timeout')))
|
||||
if (body !== null && body !== undefined) {
|
||||
if (Buffer.isBuffer(body) || typeof body === 'string') req.write(body)
|
||||
}
|
||||
req.end()
|
||||
})
|
||||
|
||||
/**
|
||||
* 获取文件流
|
||||
*/
|
||||
const getObjectStream = (relativePath) => new Promise((resolve, reject) => {
|
||||
const key = fullKey(relativePath)
|
||||
const date = new Date().toUTCString()
|
||||
const authorization = `OSS ${accessKeyId}:${signRequest({ method: 'GET', date, key })}`
|
||||
const options = {
|
||||
hostname: `${bucket}.${endpoint}`,
|
||||
path: `/${encodeKey(key)}`,
|
||||
method: 'GET',
|
||||
headers: {
|
||||
Host: `${bucket}.${endpoint}`,
|
||||
Date: date,
|
||||
Authorization: authorization
|
||||
}
|
||||
}
|
||||
const req = https.request(options, (res) => {
|
||||
if (res.statusCode < 200 || res.statusCode >= 300) {
|
||||
const chunks = []
|
||||
res.on('data', c => chunks.push(c))
|
||||
res.on('end', () => reject(new Error(`OSS get failed: ${res.statusCode} ${Buffer.concat(chunks).toString('utf8').slice(0, 500)}`)))
|
||||
return
|
||||
}
|
||||
resolve({ stream: res, headers: res.headers })
|
||||
})
|
||||
req.on('error', reject)
|
||||
req.setTimeout(60000, () => req.destroy(new Error('request timeout')))
|
||||
req.end()
|
||||
})
|
||||
|
||||
/**
|
||||
* 上传文件
|
||||
* @param {string} relativePath - 相对路径(自动加前缀)
|
||||
* @param {Buffer} body - 文件内容
|
||||
* @param {string} contentType - MIME 类型
|
||||
* @returns {Object} { key, raw_url }
|
||||
*/
|
||||
const putObject = async (relativePath, body, contentType = 'application/octet-stream') => {
|
||||
const key = fullKey(relativePath)
|
||||
const date = new Date().toUTCString()
|
||||
const extraHeaders = { 'x-oss-object-acl': 'public-read' }
|
||||
const authorization = `OSS ${accessKeyId}:${signRequest({ method: 'PUT', contentType, date, key, extraHeaders })}`
|
||||
const options = {
|
||||
hostname: `${bucket}.${endpoint}`,
|
||||
path: `/${encodeKey(key)}`,
|
||||
method: 'PUT',
|
||||
headers: {
|
||||
Host: `${bucket}.${endpoint}`,
|
||||
Date: date,
|
||||
'Content-Type': contentType,
|
||||
'Content-Length': Buffer.isBuffer(body) ? body.length : Buffer.byteLength(body || ''),
|
||||
...extraHeaders,
|
||||
Authorization: authorization
|
||||
}
|
||||
}
|
||||
const resp = await httpRequest(options, body, 120000)
|
||||
if (resp.statusCode < 200 || resp.statusCode >= 300) {
|
||||
throw new Error(`OSS put failed: ${resp.statusCode} ${resp.body.toString('utf8').slice(0, 500)}`)
|
||||
}
|
||||
const raw_url = `https://${bucket}.${endpoint}/${encodeKey(key)}`
|
||||
return { key, raw_url }
|
||||
}
|
||||
|
||||
/**
|
||||
* 删除文件
|
||||
*/
|
||||
const deleteObject = async (relativePath) => {
|
||||
const key = fullKey(relativePath)
|
||||
const date = new Date().toUTCString()
|
||||
const authorization = `OSS ${accessKeyId}:${signRequest({ method: 'DELETE', date, key })}`
|
||||
const options = {
|
||||
hostname: `${bucket}.${endpoint}`,
|
||||
path: `/${encodeKey(key)}`,
|
||||
method: 'DELETE',
|
||||
headers: {
|
||||
Host: `${bucket}.${endpoint}`,
|
||||
Date: date,
|
||||
Authorization: authorization
|
||||
}
|
||||
}
|
||||
const resp = await httpRequest(options)
|
||||
if (resp.statusCode < 200 || resp.statusCode >= 300) {
|
||||
throw new Error(`OSS delete failed: ${resp.statusCode} ${resp.body.toString('utf8').slice(0, 500)}`)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 OSS ListObjects XML 响应
|
||||
*/
|
||||
const parseListXml = (xml) => {
|
||||
const extractTag = (tag, src) => {
|
||||
const open = `<${tag}>`
|
||||
const close = `</${tag}>`
|
||||
const result = []
|
||||
let idx = 0
|
||||
while (true) {
|
||||
const s = src.indexOf(open, idx)
|
||||
if (s < 0) break
|
||||
const e = src.indexOf(close, s + open.length)
|
||||
if (e < 0) break
|
||||
result.push(src.slice(s + open.length, e))
|
||||
idx = e + close.length
|
||||
}
|
||||
return result
|
||||
}
|
||||
const extractInSection = (tag, sectionOpen, sectionClose, src) => {
|
||||
const result = []
|
||||
let searchFrom = 0
|
||||
while (true) {
|
||||
const sOpen = src.indexOf(sectionOpen, searchFrom)
|
||||
if (sOpen < 0) break
|
||||
const sClose = src.indexOf(sectionClose, sOpen)
|
||||
if (sClose < 0) break
|
||||
const section = src.slice(sOpen, sClose)
|
||||
result.push(...extractTag(tag, section))
|
||||
searchFrom = sClose + sectionClose.length
|
||||
}
|
||||
return result
|
||||
}
|
||||
const folders = extractInSection('Prefix', '<CommonPrefixes>', '</CommonPrefixes>', xml).map(p => ({
|
||||
key: p,
|
||||
name: p.split('/').filter(Boolean).pop() + '/',
|
||||
isFolder: true,
|
||||
size: 0
|
||||
}))
|
||||
const fileKeys = extractTag('Key', xml)
|
||||
const fileSizes = extractTag('Size', xml)
|
||||
const fileLms = extractTag('LastModified', xml)
|
||||
const files = fileKeys.map((k, i) => ({
|
||||
key: k,
|
||||
name: k.split('/').filter(Boolean).pop(),
|
||||
isFolder: false,
|
||||
size: parseInt(fileSizes[i] || '0', 10) || 0,
|
||||
lastModified: fileLms[i] ? Date.parse(fileLms[i]) : 0
|
||||
})).filter(f => f.key && !f.key.endsWith('/'))
|
||||
return { folders, files }
|
||||
}
|
||||
|
||||
/**
|
||||
* 列出文件(递归)
|
||||
* @param {string} relativePrefix - 相对前缀
|
||||
* @returns {Array} 文件列表
|
||||
*/
|
||||
const listRecursive = async (relativePrefix = '') => {
|
||||
const all = []
|
||||
const stack = [relativePrefix]
|
||||
while (stack.length > 0) {
|
||||
const current = stack.pop()
|
||||
const keyPrefix = fullKey(current)
|
||||
const date = new Date().toUTCString()
|
||||
const authorization = `OSS ${accessKeyId}:${signRequest({ method: 'GET', date, key: '' })}`
|
||||
const query = new URLSearchParams()
|
||||
query.set('prefix', keyPrefix)
|
||||
query.set('delimiter', '/')
|
||||
query.set('max-keys', '1000')
|
||||
const options = {
|
||||
hostname: `${bucket}.${endpoint}`,
|
||||
path: `/?${query.toString()}`,
|
||||
method: 'GET',
|
||||
headers: {
|
||||
Host: `${bucket}.${endpoint}`,
|
||||
Date: date,
|
||||
Authorization: authorization
|
||||
}
|
||||
}
|
||||
const resp = await httpRequest(options)
|
||||
if (resp.statusCode < 200 || resp.statusCode >= 300) {
|
||||
throw new Error(`OSS list failed: ${resp.statusCode} ${resp.body.toString('utf8').slice(0, 500)}`)
|
||||
}
|
||||
const parsed = parseListXml(resp.body.toString('utf8'))
|
||||
all.push(...parsed.files)
|
||||
for (const f of parsed.folders) {
|
||||
stack.push(f.key.slice(keyPrefix.length > 0 ? (fullKey('').length) : 0))
|
||||
}
|
||||
}
|
||||
return all
|
||||
}
|
||||
|
||||
/**
|
||||
* 列出对象(单层)
|
||||
*/
|
||||
const listObjects = async (relativePrefix = '') => {
|
||||
const keyPrefix = fullKey(relativePrefix)
|
||||
const date = new Date().toUTCString()
|
||||
const authorization = `OSS ${accessKeyId}:${signRequest({ method: 'GET', date, key: '' })}`
|
||||
const query = new URLSearchParams()
|
||||
query.set('prefix', keyPrefix)
|
||||
query.set('delimiter', '/')
|
||||
query.set('max-keys', '1000')
|
||||
const options = {
|
||||
hostname: `${bucket}.${endpoint}`,
|
||||
path: `/?${query.toString()}`,
|
||||
method: 'GET',
|
||||
headers: {
|
||||
Host: `${bucket}.${endpoint}`,
|
||||
Date: date,
|
||||
Authorization: authorization
|
||||
}
|
||||
}
|
||||
const resp = await httpRequest(options)
|
||||
if (resp.statusCode < 200 || resp.statusCode >= 300) {
|
||||
throw new Error(`OSS list failed: ${resp.statusCode} ${resp.body.toString('utf8').slice(0, 500)}`)
|
||||
}
|
||||
return parseListXml(resp.body.toString('utf8'))
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取文件原始 URL(公共读可直接访问)
|
||||
*/
|
||||
const getRawUrl = (relativePath) => {
|
||||
const key = fullKey(relativePath)
|
||||
return `https://${bucket}.${endpoint}/${encodeKey(key)}`
|
||||
}
|
||||
|
||||
return {
|
||||
fullKey,
|
||||
normalizeKey,
|
||||
putObject,
|
||||
deleteObject,
|
||||
getObjectStream,
|
||||
listObjects,
|
||||
listRecursive,
|
||||
getRawUrl
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { createOssClient }
|
||||
@@ -0,0 +1,10 @@
|
||||
{
|
||||
"_comment": "天气数据 skill - 由外部 Hermes Agent 每30分钟写入一次",
|
||||
"id": "weather",
|
||||
"name": "天气数据",
|
||||
"enabled": true,
|
||||
"mode": "append",
|
||||
"mode_comment": "append: 增量追加保留所有历史;latest: 只保留最新一条",
|
||||
"api_key_hash": "sha256:7b4035f766bf66e39cff1a58a5e8c73f2fcbf75f3f2a368bf0003e1110a56255",
|
||||
"read_token_hash": "sha256:aa40833cd71a76b2fc558edb1384f4d993c9e8d966a730005f7c99d56ab19037"
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
// ============================================================
|
||||
// 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
|
||||
}
|
||||
Reference in New Issue
Block a user