Files
memind/mindspace-workspace-sync.mjs
T
john 7f8d692d16 fix(agent): recover stale runs and improve new-user OA delivery
Fix DEV logout cookie clearing, materialize selected MindSpace OA assets before agent runs, and recover zombie runs from synced workspace pages. Add client run wait timeout, harness retry limits, page-edit asset forwarding, and logout/john2 scenario tests.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-10 20:05:52 +08:00

585 lines
20 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 { resolveAssetWorkspaceRelativePath } from './mindspace-workspace-path.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);
}
function generatedArtifactKindForMime(mimeType, filename) {
const normalizedMime = String(mimeType ?? '').toLowerCase();
const normalizedName = String(filename ?? '').toLowerCase();
if (normalizedMime.startsWith('image/')) return 'generated_image';
if (normalizedName.endsWith('.docx') || normalizedName.endsWith('.doc')) return 'docx';
if (normalizedName.endsWith('.pdf')) return 'pdf';
return 'generated_file';
}
function workspaceAssetDownloadUrl(assetId) {
return `/api/mindspace/v1/assets/${encodeURIComponent(assetId)}/download`;
}
function scanOptionsForWorkspaceFile(categoryCode, filename, mimeType) {
const normalizedName = String(filename ?? '').toLowerCase();
if (
categoryCode === 'public' &&
mimeType === 'text/html' &&
normalizedName.endsWith('.html')
) {
return { htmlActiveContentPolicy: 'sandbox_warn' };
}
return {};
}
function assetStatusForScan(scan) {
return scan.scanStatus === 'blocked' ? 'quarantined' : 'ready';
}
function versionStatusForScan(scan) {
return scan.scanStatus === 'blocked' ? 'blocked' : scan.scanStatus;
}
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,
conversationPackageRegistry = null,
}) {
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,
...scanOptionsForWorkspaceFile(category.category_code, file.filename, detectedMimeType),
});
const assetStatus = assetStatusForScan(scan);
const versionScanStatus = versionStatusForScan(scan);
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';
const indexedWorkspacePath = resolveAssetWorkspaceRelativePath({
categoryCode: category.category_code,
originalFilename: file.filename,
});
await conn.query(
`INSERT INTO h5_assets
(id, user_id, space_id, category_id, asset_type, mime_type, original_filename,
display_name, workspace_relative_path, 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,
indexedWorkspacePath,
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,
mimeType: detectedMimeType,
sizeBytes: buffer.length,
storageKey: finalStorageKey,
categoryCode: category.category_code,
};
} 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,
...scanOptionsForWorkspaceFile(category.category_code, file.filename, detectedMimeType),
});
const assetStatus = assetStatusForScan(scan);
const versionScanStatus = versionStatusForScan(scan);
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();
const indexedWorkspacePath = resolveAssetWorkspaceRelativePath({
categoryCode: category.category_code,
originalFilename: file.filename,
});
await conn.query(
`UPDATE h5_assets
SET current_version_id = ?, size_bytes = ?, checksum = ?, mime_type = ?,
asset_type = ?, risk_level = ?, status = ?, workspace_relative_path = ?, updated_at = ?
WHERE id = ? AND user_id = ?`,
[
versionId,
buffer.length,
checksum,
detectedMimeType,
assetInternals.assetTypeForMime(detectedMimeType),
scan.riskLevel,
assetStatus,
indexedWorkspacePath,
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,
mimeType: detectedMimeType,
sizeBytes: buffer.length,
storageKey: finalStorageKey,
categoryCode: category.category_code,
};
} catch (error) {
await conn.rollback();
throw error;
} finally {
conn.release();
}
};
const registerWorkspaceArtifactForConversation = async (userId, source, result, now = Date.now()) => {
const sessionId = String(source?.sessionId ?? source?.sourceSessionId ?? '').trim();
if (!conversationPackageRegistry || !sessionId || !result?.assetId) return null;
try {
const packageRecord = await conversationPackageRegistry.ensurePackage({
userId,
sessionId,
title: source?.title ?? null,
now,
});
const artifact = await conversationPackageRegistry.recordArtifact({
id: `ca_workspace_${result.assetId}`,
packageId: packageRecord.id,
artifactKind: generatedArtifactKindForMime(result.mimeType, result.filename),
role: 'assistant',
assetId: result.assetId,
messageId: source?.messageId ?? source?.sourceMessageId ?? null,
displayName: result.filename,
mimeType: result.mimeType,
sizeBytes: result.sizeBytes,
storageKey: result.storageKey,
canonicalUrl: workspaceAssetDownloadUrl(result.assetId),
sortOrder: now,
now,
});
await conversationPackageRegistry.writeManifestForSession({ userId, sessionId });
return artifact;
} catch (error) {
console.warn('[MindSpace] workspace conversation artifact registration failed:', error?.message ?? error);
return null;
}
};
const syncCategory = async (userId, categoryCode, source = {}) => {
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) {
const result = await updateWorkspaceFile(userId, category, existing, file, buffer);
await registerWorkspaceArtifactForConversation(userId, source, result);
existing.checksum = checksum;
updated += 1;
} else {
const result = await importWorkspaceFile(userId, category, file, buffer);
await registerWorkspaceArtifactForConversation(userId, source, result);
existingByName.set(file.filename, { checksum });
imported += 1;
}
}
return { imported, updated, skipped };
};
const syncUserWorkspace = async (userId, { categoryCode, sourceSessionId, sourceMessageId, title } = {}) => {
const codes = categoryCode ? [categoryCode] : UPLOAD_ZONE_CODES;
let imported = 0;
let updated = 0;
let skipped = 0;
const source = { sessionId: sourceSessionId, messageId: sourceMessageId, title };
for (const code of codes) {
const result = await syncCategory(userId, code, source);
imported += result.imported;
updated += result.updated;
skipped += result.skipped;
}
return { imported, updated, skipped };
};
return { syncUserWorkspace, syncCategory, listWorkspaceZoneFiles };
}
export const workspaceSyncInternals = {
generatedArtifactKindForMime,
workspaceAssetDownloadUrl,
};
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();
};
}