Files
memind/postgres-user-data-space-service.mjs
T

397 lines
21 KiB
JavaScript

import pg from 'pg';
import {
buildControlSchemaSql,
buildEnsureCurrentUserCanSetRoleSql,
deriveUserSpaceNames,
provisionUserSpace,
quotePgIdentifier,
} from './mindspace-userdata-postgres.mjs';
const pools = new Map();
const provisionedUsers = new Set();
const IDENTIFIER = /^[A-Za-z_][A-Za-z0-9_]*$/;
function ident(value, label = '标识符') {
const text = String(value ?? '').trim();
if (!IDENTIFIER.test(text)) throw Object.assign(new Error(`${label} 格式无效`), { code: 'invalid_identifier' });
return text;
}
function normalizeDataset(raw, fallbackName = null) {
const value = typeof raw === 'string' ? JSON.parse(raw) : raw;
if (!value || typeof value !== 'object') throw Object.assign(new Error('dataset 配置无效'), { code: 'invalid_dataset_config' });
const name = ident(value.name ?? fallbackName, 'dataset 名称');
const table = ident(value.table ?? value.table_name, 'dataset 表名');
const columns = value.columns && typeof value.columns === 'object' ? value.columns : {};
for (const [action, fields] of Object.entries(columns)) {
if (!Array.isArray(fields)) throw Object.assign(new Error(`dataset 字段白名单无效:${action}`), { code: 'invalid_dataset_config' });
columns[action] = fields.map((field) => ident(field, `${action} 字段`));
}
return {
name,
table,
description: String(value.description ?? ''),
actions: Array.isArray(value.actions) ? value.actions.map(String) : [],
columns,
limits: {
maxRowsPerRead: Number(value.limits?.maxRowsPerRead ?? 200),
maxInsertBytes: Number(value.limits?.maxInsertBytes ?? 8192),
},
};
}
function poolConfig(options) {
if (options.connectionString ?? process.env.MINDSPACE_USERDATA_PG_URL) {
return { connectionString: options.connectionString ?? process.env.MINDSPACE_USERDATA_PG_URL, max: 5 };
}
return {
host: options.pgHost ?? process.env.MINDSPACE_USERDATA_PG_HOST ?? '/tmp',
port: Number(options.pgPort ?? process.env.MINDSPACE_USERDATA_PG_PORT ?? 5433),
database: options.pgDatabase ?? process.env.MINDSPACE_USERDATA_PG_DATABASE ?? 'mindspace_userdata_dev',
user: options.pgUser ?? process.env.MINDSPACE_USERDATA_PG_USER ?? process.env.USER,
max: 5,
};
}
function sharedPool(options) {
if (options.pgPool) return options.pgPool;
const config = poolConfig(options);
const key = JSON.stringify(config);
if (!pools.has(key)) pools.set(key, new pg.Pool(config));
return pools.get(key);
}
function rejectSql(sql, { readonly = false } = {}) {
const cleaned = String(sql ?? '').replace(/--.*$/gm, '').replace(/\/\*[\s\S]*?\*\//g, '').trim();
if (!cleaned || cleaned.length > 20000) throw new Error('SQL 为空或过长');
if (readonly && !/^(SELECT|WITH)\b/i.test(cleaned)) throw new Error('private_data_query 只允许 SELECT/WITH');
if (/\b(CREATE|ALTER|DROP)\s+(ROLE|USER|DATABASE|TABLESPACE|EXTENSION)|ALTER\s+SYSTEM|COPY[\s\S]+PROGRAM|SECURITY\s+DEFINER|pg_(read|write)_file|lo_import|dblink/i.test(cleaned)) {
throw new Error('SQL 包含用户空间不允许的操作');
}
if (/\bSET\s+(ROLE|SESSION_AUTHORIZATION|search_path)\b/i.test(cleaned)) throw new Error('SQL 不允许改变安全上下文');
return cleaned;
}
function translateSqliteAgentDdl(sql) {
return sql
.replace(/INTEGER\s+PRIMARY\s+KEY\s+AUTOINCREMENT/gi, 'BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY')
.replace(/TEXT\s+NOT\s+NULL\s+DEFAULT\s*\(datetime\(\s*'now'\s*,\s*'(?:\+8 hours|localtime)'\s*\)\)/gi, 'TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP')
.replace(/TEXT\s+DEFAULT\s*\(datetime\(\s*'now'\s*,\s*'(?:\+8 hours|localtime)'\s*\)\)/gi, 'TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP');
}
export function createPostgresUserDataSpaceService(options = {}) {
const names = deriveUserSpaceNames(options.userId);
const pool = sharedPool(options);
const schema = quotePgIdentifier(names.schemaName);
const role = quotePgIdentifier(names.agentRole);
const maxRows = Number(options.maxRows ?? 200);
async function ensureProvisioned({ force = false } = {}) {
if (provisionedUsers.has(names.userId)) return;
const client = await pool.connect();
try {
await client.query(buildControlSchemaSql());
const existing = await client.query('SELECT migration_state FROM mindspace_control.user_spaces WHERE user_id=$1::uuid', [names.userId]);
if (!existing.rows[0]) {
if (!force && String(options.autoProvision ?? process.env.MINDSPACE_USERDATA_AUTO_PROVISION ?? '0') !== '1') {
throw Object.assign(new Error('用户 PG 空间尚未分配'), { code: 'user_space_not_provisioned' });
}
await provisionUserSpace(client, names.userId, { sourceSqlitePath: options.workspaceRoot ? `${options.workspaceRoot}/.mindspace/private-data.sqlite` : null });
await client.query('BEGIN');
try {
await client.query(`CREATE TABLE IF NOT EXISTS ${schema}.__page_data_datasets (name text PRIMARY KEY, table_name text NOT NULL, config_json text NOT NULL, created_at timestamptz NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at timestamptz NOT NULL DEFAULT CURRENT_TIMESTAMP)`);
await client.query(`ALTER TABLE ${schema}.__page_data_datasets OWNER TO ${quotePgIdentifier(names.ownerRole)}`);
await client.query(`GRANT SELECT,INSERT,UPDATE,DELETE ON ${schema}.__page_data_datasets TO ${role}`);
await client.query("UPDATE mindspace_control.user_spaces SET migration_state='cutover',cutover_at=CURRENT_TIMESTAMP,rollback_until=CURRENT_TIMESTAMP + INTERVAL '30 days',updated_at=CURRENT_TIMESTAMP WHERE user_id=$1::uuid", [names.userId]);
await client.query('COMMIT');
} catch (error) {
await client.query('ROLLBACK');
throw error;
}
} else {
// PostgreSQL 17 grants CREATEROLE owners ADMIN but not SET by default.
// Reconcile existing spaces before the runtime executes SET LOCAL ROLE.
await client.query(buildEnsureCurrentUserCanSetRoleSql(names.agentRole));
}
provisionedUsers.add(names.userId);
} finally {
client.release();
}
}
async function tenant(fn, { write = false } = {}) {
await ensureProvisioned();
const client = await pool.connect();
let quotaBytes = null;
try {
await client.query('BEGIN');
if (write) {
const quota = await client.query('SELECT quota_bytes FROM mindspace_control.user_spaces WHERE user_id=$1::uuid', [names.userId]);
quotaBytes = Number(quota.rows[0]?.quota_bytes ?? 0);
}
await client.query(`SET LOCAL ROLE ${role}`);
await client.query(`SET LOCAL search_path = ${schema}, pg_catalog`);
await client.query(`SET LOCAL statement_timeout = '${write ? 30 : 10}s'`);
const result = await fn(client);
if (write && quotaBytes > 0) {
const usage = await client.query(
`SELECT COALESCE(sum(pg_total_relation_size(c.oid)),0)::bigint AS bytes
FROM pg_class c JOIN pg_namespace n ON n.oid=c.relnamespace
WHERE n.nspname=$1 AND c.relkind IN ('r','p','m')`,
[names.schemaName],
);
if (Number(usage.rows[0].bytes) > quotaBytes) {
throw Object.assign(new Error(`用户 PG 空间超过配额:${usage.rows[0].bytes}/${quotaBytes}`), { code: 'quota_exceeded' });
}
}
await client.query(write ? 'COMMIT' : 'ROLLBACK');
if (write) {
await client.query(
`INSERT INTO mindspace_control.audit_events(user_id,actor_type,actor_id,action,result,detail_json)
VALUES($1::uuid,'agent',$2,'postgres_write','success',$3::jsonb)`,
[names.userId, names.agentRole, JSON.stringify({ schema: names.schemaName })],
);
}
return result;
} catch (error) {
await client.query('ROLLBACK').catch(() => {});
if (write) {
await client.query(
`INSERT INTO mindspace_control.audit_events(user_id,actor_type,actor_id,action,result,detail_json)
VALUES($1::uuid,'agent',$2,'postgres_write','failed',$3::jsonb)`,
[names.userId, names.agentRole, JSON.stringify({ schema: names.schemaName, code: error?.code ?? null })],
).catch(() => {});
}
throw error;
} finally {
client.release();
}
}
async function listTableColumns(tableName) {
const table = ident(tableName, '表名');
return tenant(async (client) => {
const result = await client.query(
`SELECT column_name AS name, data_type AS type,
(is_nullable = 'NO') AS "notNull",
(SELECT EXISTS (SELECT 1 FROM pg_index i JOIN pg_attribute a ON a.attrelid=i.indrelid AND a.attnum=ANY(i.indkey) WHERE i.indrelid=($1||'.'||$2)::regclass AND i.indisprimary AND a.attname=c.column_name)) AS pk
FROM information_schema.columns c WHERE table_schema=$1 AND table_name=$2 ORDER BY ordinal_position`,
[names.schemaName, table],
);
return result.rows;
});
}
async function getSchema() {
return tenant(async (client) => (await client.query(
`SELECT table_name, column_name, data_type AS type, (is_nullable='NO') AS not_null
FROM information_schema.columns WHERE table_schema=$1 ORDER BY table_name,ordinal_position`,
[names.schemaName],
)).rows);
}
async function querySql(sql) {
const cleaned = rejectSql(sql, { readonly: true }).replace(/;\s*$/, '');
return tenant(async (client) => (await client.query(`SELECT * FROM (${cleaned}) AS q LIMIT ${maxRows}`)).rows);
}
async function executeSql(sql) {
const cleaned = translateSqliteAgentDdl(rejectSql(sql));
return tenant(async (client) => {
const result = await client.query(cleaned);
return { backend: 'postgres', command: result.command, affectedRows: result.rowCount ?? 0, sizeBytes: null, deltaBytes: null, quotaSynced: false };
}, { write: true });
}
async function getDataset(name) {
const datasetName = ident(name, 'dataset 名称');
return tenant(async (client) => {
const result = await client.query('SELECT name,table_name,config_json,created_at,updated_at FROM __page_data_datasets WHERE name=$1 LIMIT 1', [datasetName]);
if (!result.rows[0]) return null;
return { ...normalizeDataset(result.rows[0].config_json, result.rows[0].name), createdAt: result.rows[0].created_at, updatedAt: result.rows[0].updated_at };
});
}
async function listDatasets() {
return tenant(async (client) => (await client.query('SELECT name,table_name,config_json,created_at,updated_at FROM __page_data_datasets ORDER BY name')).rows.map((row) => ({
...normalizeDataset(row.config_json, row.name), createdAt: row.created_at, updatedAt: row.updated_at,
})));
}
async function upsertDataset(input) {
const dataset = normalizeDataset(input);
if (!(await listTableColumns(dataset.table)).length) throw Object.assign(new Error(`dataset 对应表不存在:${dataset.table}`), { code: 'table_not_found' });
await tenant((client) => client.query(
`INSERT INTO __page_data_datasets(name,table_name,config_json,updated_at) VALUES($1,$2,$3,CURRENT_TIMESTAMP)
ON CONFLICT(name) DO UPDATE SET table_name=EXCLUDED.table_name,config_json=EXCLUDED.config_json,updated_at=CURRENT_TIMESTAMP`,
[dataset.name, dataset.table, JSON.stringify(dataset)],
), { write: true });
return getDataset(dataset.name);
}
function assertAction(dataset, action) {
if (!dataset.actions.includes(action)) throw Object.assign(new Error(`dataset 未授权 ${action}`), { code: 'action_not_allowed' });
}
async function readRowsForDataset(datasetInput, options = {}) {
const dataset = await Promise.resolve(datasetInput);
assertAction(dataset, 'read');
const columns = dataset.columns.read?.map((item) => ident(item, '字段')) ?? [];
if (!columns.length) throw Object.assign(new Error('dataset 未配置可读字段'), { code: 'columns_not_allowed' });
const limit = Math.min(Math.max(1, Number(options.limit ?? dataset.limits.maxRowsPerRead)), maxRows, dataset.limits.maxRowsPerRead);
const offset = Math.max(0, Number(options.offset ?? 0));
const order = options.orderBy ? ident(options.orderBy, '排序字段') : columns.includes('id') ? 'id' : null;
if (order && !columns.includes(order)) throw Object.assign(new Error('排序字段未授权'), { code: 'columns_not_allowed' });
return tenant(async (client) => {
const filters = [];
if ((await listTableColumns(dataset.table)).some((column) => column.name === 'deleted_at') && !options.includeDeleted) filters.push('deleted_at IS NULL');
if (options.rowScope?.whereClause) filters.push(`(${options.rowScope.whereClause})`);
const where = filters.length ? ` WHERE ${filters.join(' AND ')}` : '';
const result = await client.query(
`SELECT ${columns.map(quotePgIdentifier).join(',')} FROM ${quotePgIdentifier(dataset.table)}${where}${order ? ` ORDER BY ${quotePgIdentifier(order)} ${String(options.orderDir).toLowerCase() === 'asc' ? 'ASC' : 'DESC'}` : ''} LIMIT $1 OFFSET $2`,
[limit, offset],
);
return { dataset, rows: result.rows, limit, offset };
});
}
async function readDatasetRows(name, options) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return readRowsForDataset(dataset, options);
}
async function getDatasetStats(name) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
assertAction(dataset, 'read');
return tenant(async (client) => ({ dataset: { name: dataset.name, table: dataset.table }, total: Number((await client.query(`SELECT count(*)::bigint AS count FROM ${quotePgIdentifier(dataset.table)}`)).rows[0].count) }));
}
async function insertWithDataset(dataset, payload, meta = {}) {
assertAction(dataset, 'insert');
const allowed = new Set(dataset.columns.insert ?? []);
const keys = Object.keys(payload ?? {});
if (keys.some((key) => !allowed.has(key))) throw Object.assign(new Error('提交包含未授权字段'), { code: 'columns_not_allowed' });
return tenant(async (client) => {
const columns = await listTableColumns(dataset.table);
const values = { ...payload };
if (meta.rowScope?.ownerColumn && columns.some((c) => c.name === meta.rowScope.ownerColumn)) values[meta.rowScope.ownerColumn] ??= meta.rowScope.visitorUserId;
const insertKeys = Object.keys(values).map((key) => ident(key, '字段'));
const result = await client.query(
`INSERT INTO ${quotePgIdentifier(dataset.table)} (${insertKeys.map(quotePgIdentifier).join(',')}) VALUES (${insertKeys.map((_, i) => `$${i + 1}`).join(',')}) RETURNING *`,
insertKeys.map((key) => values[key]),
);
return { dataset, row: result.rows[0] };
}, { write: true });
}
async function insertDatasetRow(name, payload, meta = {}) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return insertWithDataset(dataset, payload, meta);
}
async function insertRowForDataset(datasetInput, payload, meta = {}) {
const dataset = await Promise.resolve(datasetInput);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return insertWithDataset(dataset, payload, meta);
}
async function updateRowForDataset(datasetInput, rowId, payload, meta = {}) {
const dataset = await Promise.resolve(datasetInput);
assertAction(dataset, 'update');
const id = Number(rowId);
if (!Number.isFinite(id) || id <= 0) throw Object.assign(new Error('行 id 无效'), { code: 'invalid_row_id' });
const allowed = new Set(dataset.columns.update ?? []);
const keys = Object.keys(payload ?? {});
if (!keys.length || keys.some((key) => !allowed.has(key))) throw Object.assign(new Error('更新包含未授权字段'), { code: 'columns_not_allowed' });
return tenant(async (client) => {
const result = await client.query(
`UPDATE ${quotePgIdentifier(dataset.table)} SET ${keys.map((key, i) => `${quotePgIdentifier(ident(key, '字段'))}=$${i + 1}`).join(',')} WHERE id=$${keys.length + 1}${meta.rowScope?.whereClause ? ` AND ${meta.rowScope.whereClause}` : ''} RETURNING *`,
[...keys.map((key) => payload[key]), id],
);
if (!result.rows[0]) throw Object.assign(new Error('无权更新该行或行不存在'), { code: 'row_not_allowed' });
return { dataset, row: result.rows[0] };
}, { write: true });
}
async function updateDatasetRow(name, rowId, payload, meta) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return updateRowForDataset(dataset, rowId, payload, meta);
}
async function softDeleteRowForDataset(datasetInput, rowId, meta = {}) {
const dataset = await Promise.resolve(datasetInput);
assertAction(dataset, 'soft_delete');
const columns = await listTableColumns(dataset.table);
if (!columns.some((item) => item.name === 'deleted_at')) throw Object.assign(new Error('目标表不支持软删除'), { code: 'soft_delete_unsupported' });
const id = Number(rowId);
return tenant(async (client) => {
const hasDeletedBy = columns.some((item) => item.name === 'deleted_by');
const result = await client.query(
`UPDATE ${quotePgIdentifier(dataset.table)} SET deleted_at=CURRENT_TIMESTAMP${hasDeletedBy ? ',deleted_by=$2' : ''} WHERE id=$1 AND deleted_at IS NULL${meta.rowScope?.whereClause ? ` AND ${meta.rowScope.whereClause}` : ''} RETURNING id,deleted_at`,
hasDeletedBy ? [id, meta.deletedBy ?? meta.updatedByLabel ?? 'public'] : [id],
);
if (!result.rows[0]) throw Object.assign(new Error('无权删除该行或行不存在'), { code: 'row_not_allowed' });
return { dataset, id, deleted: true };
}, { write: true });
}
async function softDeleteDatasetRow(name, rowId, meta) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return softDeleteRowForDataset(dataset, rowId, meta);
}
async function restoreSoftDeletedRow(name, rowId) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
const columns = await listTableColumns(dataset.table);
return tenant(async (client) => {
const result = await client.query(
`UPDATE ${quotePgIdentifier(dataset.table)} SET deleted_at=NULL${columns.some((item) => item.name === 'deleted_by') ? ',deleted_by=NULL' : ''} WHERE id=$1 AND deleted_at IS NOT NULL RETURNING *`,
[Number(rowId)],
);
if (!result.rows[0]) throw Object.assign(new Error('行不存在或未删除'), { code: 'row_not_allowed' });
return { dataset, row: result.rows[0], restored: true };
}, { write: true });
}
async function exportDatasetRows(name, { format = 'json', includeDeleted = false, limit = 1000 } = {}) {
const result = await readDatasetRows(name, { includeDeleted, limit: Math.min(Number(limit), 5000) });
if (String(format).toLowerCase() !== 'csv') return { dataset: name, format: 'json', rows: result.rows, rowCount: result.rows.length };
const columns = result.dataset.columns.read ?? [];
const escape = (value) => {
const text = value == null ? '' : String(value);
return /[",\n\r]/.test(text) ? `"${text.replaceAll('"', '""')}"` : text;
};
return { dataset: name, format: 'csv', content: `${columns.join(',')}\n${result.rows.map((row) => columns.map((column) => escape(row[column])).join(',')).join('\n')}\n`, rowCount: result.rows.length };
}
async function getDatasetSchema(name) {
const dataset = await getDataset(name);
if (!dataset) throw Object.assign(new Error('dataset 不存在'), { code: 'dataset_not_found' });
return { dataset, tableColumns: await listTableColumns(dataset.table) };
}
async function getInfo() {
return { name: '用户私有数据空间', backend: 'postgres', database: poolConfig(options).database, schema: names.schemaName, userId: names.userId };
}
async function ensureReady() {
await ensureProvisioned({ force: true });
return getInfo();
}
return {
backend: 'postgres', workspaceRoot: options.workspaceRoot, privateDataDb: null,
ensureReady, getInfo, getSchema, querySql, executeSql, listTableColumns, getDataset, listDatasets, upsertDataset,
readRowsForDataset, readDatasetRows, getDatasetStats, getStatsForDataset: async (dataset) => getDatasetStats((await Promise.resolve(dataset)).name),
getDatasetSchema, getSchemaForDataset: async (dataset) => ({ dataset: await Promise.resolve(dataset), tableColumns: await listTableColumns((await Promise.resolve(dataset)).table) }),
insertDatasetRow, insertRowForDataset, updateDatasetRow, updateRowForDataset,
softDeleteDatasetRow, softDeleteRowForDataset, restoreSoftDeletedRow, exportDatasetRows,
};
}