Files
memind/mindspace-workspace-sync.mjs
T
john 637b32515a fix(mindspace): sync nested OA files and count workspace assets
Recursive workspace sync now imports files in oa/private/public subfolders,
references workspace paths without duplicating storage, and includes
workspace-sourced assets in category itemCount. Refresh counts when leaving
a category view.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-06-29 10:34:18 +08:00

473 lines
16 KiB
JavaScript

import crypto from 'node:crypto';
import fs from 'node:fs';
import fsPromises from 'node:fs/promises';
import path from 'node:path';
import {
UPLOAD_ZONE_CODES,
resolveUserWorkspaceRoot,
resolveZoneDir,
} from './user-space.mjs';
import { assetInternals } from './mindspace-assets.mjs';
import { runBasicFileScan } from './mindspace-scan.mjs';
import { buildWorkspaceStorageKey } from './workspace-storage.mjs';
const SKIP_FILENAMES = new Set(['index.html', '.tkmindhints', '.goosehints', '.ls_output']);
const SKIP_ZONE_DIR_NAMES = new Set([
'node_modules',
'.venv',
'.venv2',
'__pycache__',
'.git',
'dist',
'build',
'.next',
'target',
'backend',
'frontend',
]);
/** 子目录内常见的项目/脚手架文件,不同步到 OA 资产库 */
const NESTED_SKIP_BASENAMES = new Set([
'README.md',
'requirements.txt',
'package.json',
'package-lock.json',
'pnpm-lock.yaml',
'yarn.lock',
]);
function asNumber(value) {
return Number(value ?? 0);
}
export function normalizeWorkspaceRelativePath(relativePath) {
const normalized = String(relativePath ?? '')
.normalize('NFKC')
.trim()
.replace(/\\/g, '/');
if (!normalized || normalized.startsWith('/') || normalized.includes('\0')) return null;
const parts = normalized.split('/').filter(Boolean);
if (parts.some((part) => part === '.' || part === '..')) return null;
if (parts.some((part) => part.startsWith('.'))) return null;
return parts.join('/').slice(0, 255);
}
export function shouldSyncWorkspaceFilename(filename) {
const relativePath = normalizeWorkspaceRelativePath(filename);
if (!relativePath) return false;
const basename = path.posix.basename(relativePath);
if (!basename || basename.startsWith('.')) return false;
if (SKIP_FILENAMES.has(basename)) return false;
if (basename.endsWith('.thumbnail.svg')) return false;
if (basename.endsWith('.sh')) return false;
if (relativePath.includes('/') && NESTED_SKIP_BASENAMES.has(basename)) return false;
return assetInternals.expectedMimeType(basename) !== null;
}
async function walkWorkspaceZoneDir(zoneDir, relativePrefix, files) {
let entries;
try {
entries = await fsPromises.readdir(zoneDir, { withFileTypes: true });
} catch {
return;
}
for (const entry of entries) {
if (entry.name.startsWith('.')) continue;
const absolutePath = path.join(zoneDir, entry.name);
if (entry.isDirectory()) {
if (SKIP_ZONE_DIR_NAMES.has(entry.name)) continue;
const nextPrefix = relativePrefix ? `${relativePrefix}/${entry.name}` : entry.name;
await walkWorkspaceZoneDir(absolutePath, nextPrefix, files);
continue;
}
if (!entry.isFile()) continue;
const filename = relativePrefix ? `${relativePrefix}/${entry.name}` : entry.name;
if (!shouldSyncWorkspaceFilename(filename)) continue;
const stat = await fsPromises.stat(absolutePath);
files.push({
filename,
absolutePath,
sizeBytes: stat.size,
mtimeMs: stat.mtimeMs,
});
}
}
export async function listWorkspaceZoneFiles(workspaceRoot, categoryCode) {
if (!UPLOAD_ZONE_CODES.includes(categoryCode)) return [];
const zoneDir = resolveZoneDir(workspaceRoot, categoryCode);
const files = [];
await walkWorkspaceZoneDir(zoneDir, '', files);
return files;
}
export function createWorkspaceAssetSync({ pool, storageRoot, h5Root, maxFileBytes, idFactory }) {
void storageRoot;
const loadExistingAssets = async (userId, categoryId) => {
const [rows] = await pool.query(
`SELECT a.id, a.original_filename, a.checksum, a.size_bytes, a.current_version_id, a.status
FROM h5_assets a
WHERE a.user_id = ? AND a.category_id = ? AND a.status <> 'deleted'`,
[userId, categoryId],
);
return new Map(rows.map((row) => [row.original_filename, row]));
};
/** Skip re-importing workspace files the user already deleted (checksum unchanged). */
const loadDeletedWorkspaceChecksums = async (userId, categoryId) => {
const [rows] = await pool.query(
`SELECT original_filename, checksum
FROM h5_assets
WHERE user_id = ? AND category_id = ? AND status = 'deleted' AND source_type = 'workspace'`,
[userId, categoryId],
);
const byFilename = new Map();
for (const row of rows) {
const checksums = byFilename.get(row.original_filename) ?? new Set();
checksums.add(row.checksum);
byFilename.set(row.original_filename, checksums);
}
return byFilename;
};
const importWorkspaceFile = async (userId, category, file, buffer) => {
const conn = await pool.getConnection();
try {
await conn.beginTransaction();
const [spaces] = await conn.query(
`SELECT id, quota_bytes, used_bytes, reserved_bytes, status
FROM h5_user_spaces
WHERE id = ? AND user_id = ?
LIMIT 1
FOR UPDATE`,
[category.space_id, userId],
);
const space = spaces[0];
if (!space || space.status !== 'active') {
throw Object.assign(new Error('用户空间不可用'), { code: 'space_unavailable' });
}
const available =
asNumber(space.quota_bytes) - asNumber(space.used_bytes) - asNumber(space.reserved_bytes);
if (available < buffer.length) {
throw Object.assign(new Error('剩余空间不足'), { code: 'quota_exceeded' });
}
const detectedMimeType = assetInternals.detectMimeType(buffer, file.filename);
if (!detectedMimeType) {
throw Object.assign(new Error('无法确认文件类型'), { code: 'unsupported_file_type' });
}
const scan = runBasicFileScan(buffer, {
filename: file.filename,
mimeType: detectedMimeType,
});
const assetStatus = scan.scanStatus === 'passed' ? 'ready' : 'quarantined';
const versionScanStatus = scan.scanStatus === 'passed' ? 'passed' : 'blocked';
const checksum = crypto.createHash('sha256').update(buffer).digest('hex');
const assetId = idFactory();
const versionId = idFactory();
const finalStorageKey = buildWorkspaceStorageKey(
userId,
category.category_code,
file.filename,
);
const now = Date.now();
const visibility = category.category_code === 'public' ? 'public_candidate' : 'private';
await conn.query(
`INSERT INTO h5_assets
(id, user_id, space_id, category_id, asset_type, mime_type, original_filename,
display_name, current_version_id, size_bytes, checksum, risk_level, visibility,
status, source_type, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'workspace', ?, ?)`,
[
assetId,
userId,
category.space_id,
category.id,
assetInternals.assetTypeForMime(detectedMimeType),
detectedMimeType,
file.filename,
path.basename(file.filename, path.extname(file.filename)) || file.filename,
versionId,
buffer.length,
checksum,
scan.riskLevel,
visibility,
assetStatus,
now,
now,
],
);
await conn.query(
`INSERT INTO h5_asset_versions
(id, asset_id, version_no, storage_key, size_bytes, checksum, mime_type,
created_by, change_note, scan_status, created_at)
VALUES (?, ?, 1, ?, ?, ?, ?, ?, '从工作区同步', ?, ?)`,
[
versionId,
assetId,
finalStorageKey,
buffer.length,
checksum,
detectedMimeType,
userId,
versionScanStatus,
now,
],
);
await conn.query(
`UPDATE h5_user_spaces SET used_bytes = used_bytes + ?, updated_at = ?
WHERE id = ? AND user_id = ?`,
[buffer.length, now, category.space_id, userId],
);
await conn.commit();
return { action: 'imported', assetId, filename: file.filename, checksum };
} catch (error) {
await conn.rollback();
throw error;
} finally {
conn.release();
}
};
const updateWorkspaceFile = async (userId, category, existing, file, buffer) => {
const conn = await pool.getConnection();
try {
await conn.beginTransaction();
const [spaces] = await conn.query(
`SELECT id, quota_bytes, used_bytes, reserved_bytes, status
FROM h5_user_spaces
WHERE id = ? AND user_id = ?
LIMIT 1
FOR UPDATE`,
[category.space_id, userId],
);
const space = spaces[0];
if (!space || space.status !== 'active') {
throw Object.assign(new Error('用户空间不可用'), { code: 'space_unavailable' });
}
const sizeDelta = buffer.length - asNumber(existing.size_bytes);
const available =
asNumber(space.quota_bytes) - asNumber(space.used_bytes) - asNumber(space.reserved_bytes);
if (sizeDelta > 0 && available < sizeDelta) {
throw Object.assign(new Error('剩余空间不足'), { code: 'quota_exceeded' });
}
const detectedMimeType = assetInternals.detectMimeType(buffer, file.filename);
if (!detectedMimeType) {
throw Object.assign(new Error('无法确认文件类型'), { code: 'unsupported_file_type' });
}
const scan = runBasicFileScan(buffer, {
filename: file.filename,
mimeType: detectedMimeType,
});
const assetStatus = scan.scanStatus === 'passed' ? 'ready' : 'quarantined';
const versionScanStatus = scan.scanStatus === 'passed' ? 'passed' : 'blocked';
const checksum = crypto.createHash('sha256').update(buffer).digest('hex');
const [versionRows] = await conn.query(
`SELECT COALESCE(MAX(version_no), 0) AS max_version
FROM h5_asset_versions
WHERE asset_id = ?`,
[existing.id],
);
const versionNo = asNumber(versionRows[0]?.max_version) + 1;
const versionId = idFactory();
const finalStorageKey = buildWorkspaceStorageKey(
userId,
category.category_code,
file.filename,
);
const now = Date.now();
await conn.query(
`UPDATE h5_assets
SET current_version_id = ?, size_bytes = ?, checksum = ?, mime_type = ?,
asset_type = ?, risk_level = ?, status = ?, updated_at = ?
WHERE id = ? AND user_id = ?`,
[
versionId,
buffer.length,
checksum,
detectedMimeType,
assetInternals.assetTypeForMime(detectedMimeType),
scan.riskLevel,
assetStatus,
now,
existing.id,
userId,
],
);
await conn.query(
`INSERT INTO h5_asset_versions
(id, asset_id, version_no, storage_key, size_bytes, checksum, mime_type,
created_by, change_note, scan_status, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, '工作区文件更新', ?, ?)`,
[
versionId,
existing.id,
versionNo,
finalStorageKey,
buffer.length,
checksum,
detectedMimeType,
userId,
versionScanStatus,
now,
],
);
if (sizeDelta !== 0) {
await conn.query(
`UPDATE h5_user_spaces SET used_bytes = GREATEST(0, used_bytes + ?), updated_at = ?
WHERE id = ? AND user_id = ?`,
[sizeDelta, now, category.space_id, userId],
);
}
await conn.commit();
return { action: 'updated', assetId: existing.id, filename: file.filename, checksum };
} catch (error) {
await conn.rollback();
throw error;
} finally {
conn.release();
}
};
const syncCategory = async (userId, categoryCode) => {
if (!h5Root) return { imported: 0, updated: 0, skipped: 0 };
const workspaceRoot = resolveUserWorkspaceRoot(h5Root, { id: userId });
const files = await listWorkspaceZoneFiles(workspaceRoot, categoryCode);
const [categories] = await pool.query(
`SELECT c.id, c.space_id, c.category_code
FROM h5_space_categories c
WHERE c.user_id = ? AND c.category_code = ?
LIMIT 1`,
[userId, categoryCode],
);
const category = categories[0];
if (!category) return { imported: 0, updated: 0, skipped: 0 };
const existingByName = await loadExistingAssets(userId, category.id);
const deletedWorkspaceChecksums = await loadDeletedWorkspaceChecksums(userId, category.id);
let imported = 0;
let updated = 0;
let skipped = 0;
for (const file of files) {
if (file.sizeBytes <= 0 || file.sizeBytes > maxFileBytes) {
skipped += 1;
continue;
}
const buffer = await fsPromises.readFile(file.absolutePath);
const checksum = crypto.createHash('sha256').update(buffer).digest('hex');
if (deletedWorkspaceChecksums.get(file.filename)?.has(checksum)) {
skipped += 1;
continue;
}
const existing = existingByName.get(file.filename);
if (existing?.checksum === checksum) {
skipped += 1;
continue;
}
if (existing) {
await updateWorkspaceFile(userId, category, existing, file, buffer);
existing.checksum = checksum;
updated += 1;
} else {
await importWorkspaceFile(userId, category, file, buffer);
existingByName.set(file.filename, { checksum });
imported += 1;
}
}
return { imported, updated, skipped };
};
const syncUserWorkspace = async (userId, { categoryCode } = {}) => {
const codes = categoryCode ? [categoryCode] : UPLOAD_ZONE_CODES;
let imported = 0;
let updated = 0;
let skipped = 0;
for (const code of codes) {
const result = await syncCategory(userId, code);
imported += result.imported;
updated += result.updated;
skipped += result.skipped;
}
return { imported, updated, skipped };
};
return { syncUserWorkspace, syncCategory, listWorkspaceZoneFiles };
}
export function startWorkspaceAssetSyncWatcher({ publishRoot, syncUserWorkspaceByDirKey }) {
if (!publishRoot || !syncUserWorkspaceByDirKey) return () => {};
const pending = new Map();
const schedule = (dirKey, categoryCode) => {
const key = `${dirKey}:${categoryCode ?? 'all'}`;
const existing = pending.get(key);
if (existing) clearTimeout(existing);
pending.set(
key,
setTimeout(() => {
pending.delete(key);
void syncUserWorkspaceByDirKey(dirKey, categoryCode ? { categoryCode } : {}).catch(() => {});
}, 600),
);
};
const watchZoneDir = (dirKey, zoneDir, categoryCode) => {
try {
fs.watch(
zoneDir,
{ recursive: true },
(_event, filename) => {
if (!filename) {
schedule(dirKey, categoryCode);
return;
}
const normalized = String(filename).replace(/\\/g, '/');
if (shouldSyncWorkspaceFilename(normalized)) {
schedule(dirKey, categoryCode);
}
},
);
} catch {
// ignore unsupported watch targets
}
};
const attachUser = async (dirKey) => {
const userDir = path.join(publishRoot, dirKey);
try {
await fsPromises.access(userDir);
} catch {
return;
}
for (const categoryCode of UPLOAD_ZONE_CODES) {
watchZoneDir(dirKey, resolveZoneDir(userDir, categoryCode), categoryCode);
}
};
for (const entry of fs.readdirSync(publishRoot, { withFileTypes: true })) {
if (!entry.isDirectory() || entry.name === 'wiki' || entry.name.startsWith('.')) continue;
void attachUser(entry.name);
}
try {
fs.watch(publishRoot, (_event, filename) => {
if (!filename) return;
const userDir = path.join(publishRoot, filename);
if (fs.existsSync(userDir) && fs.statSync(userDir).isDirectory()) {
void attachUser(filename);
}
});
} catch {
// ignore
}
return () => {
for (const timer of pending.values()) clearTimeout(timer);
pending.clear();
};
}