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 = {}, { onlyRelativePaths = null } = {}) => { if (!h5Root) return { imported: 0, updated: 0, skipped: 0 }; const workspaceRoot = resolveUserWorkspaceRoot(h5Root, { id: userId }); let files = await listWorkspaceZoneFiles(workspaceRoot, categoryCode); if (onlyRelativePaths != null) { const allowed = new Set( onlyRelativePaths .map((item) => normalizeWorkspaceRelativePath(item)) .filter(Boolean), ); files = files.filter((file) => allowed.has(normalizeWorkspaceRelativePath(`${categoryCode}/${file.filename}`)), ); } 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, onlyRelativePaths = null } = {}, ) => { 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, { onlyRelativePaths }); 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(); }; }