From fc289f7b420dcee0a1d908bc1f318b79545991af Mon Sep 17 00:00:00 2001 From: Codex Date: Mon, 29 Jun 2026 06:14:58 +0800 Subject: [PATCH] feat: recheck content assets --- .env.example | 36 ++ README.md | 17 +- apps/worker/package.json | 5 +- apps/worker/src/config.ts | 60 ++ apps/worker/src/index.ts | 9 + apps/worker/src/jobs/assets.ts | 715 +++++++++++++++++++++ apps/worker/src/types/ali-oss.d.ts | 25 + docs/refactor/backend-capability-status.md | 3 +- docs/refactor/backend-handoff-roadmap.md | 4 +- docs/refactor/backend-progress.md | 6 +- docs/refactor/next-development-todo.md | 4 +- docs/refactor/object-storage.md | 21 +- docs/refactor/taro-frontend-integration.md | 3 + package-lock.json | 2 + package.json | 1 + scripts/asset-worker-integration-test.js | 175 +++++ 16 files changed, 1073 insertions(+), 13 deletions(-) create mode 100644 apps/worker/src/jobs/assets.ts create mode 100644 apps/worker/src/types/ali-oss.d.ts create mode 100644 scripts/asset-worker-integration-test.js diff --git a/.env.example b/.env.example index 17680e9b..5c1a672f 100644 --- a/.env.example +++ b/.env.example @@ -55,3 +55,39 @@ WORKER_CRM_BACKOFF_SECONDS=5,30,120,600,1800 WORKER_CRM_REQUEST_TIMEOUT_MS=10000 # 仅本地 fake webhook 测试允许 http://127.0.0.1;生产建议 false WORKER_CRM_ALLOW_INSECURE_LOCALHOST=false + +# Worker 配置:支付/退款补偿 +WORKER_COMMERCE_BATCH_SIZE=20 +WORKER_COMMERCE_MIN_AGE_SECONDS=300 +WORKER_COMMERCE_REQUEST_TIMEOUT_MS=10000 + +# Worker 配置:内容资源复检。异常托管对象会被标记 failed 并从 active 退回 draft。 +WORKER_ASSET_BATCH_SIZE=50 +WORKER_ASSET_MIN_AGE_SECONDS=300 +WORKER_ASSET_RECHECK_INTERVAL_SECONDS=86400 +WORKER_ASSET_REQUEST_TIMEOUT_MS=10000 + +# 对象存储配置:API 和 assets worker 共用 +STORAGE_DEFAULT_PROVIDER=local_dev +STORAGE_DEFAULT_BUCKET=tenant-assets +STORAGE_PUBLIC_BASE_URL= +STORAGE_MAX_UPLOAD_BYTES=524288000 +STORAGE_ALLOWED_MIME_PREFIXES=image/,video/,audio/ +STORAGE_ALLOWED_MIME_TYPES=application/pdf,application/json,application/zip,application/x-zip-compressed,application/msword,application/vnd.openxmlformats-officedocument.wordprocessingml.document,application/vnd.ms-excel,application/vnd.openxmlformats-officedocument.spreadsheetml.sheet,application/vnd.ms-powerpoint,application/vnd.openxmlformats-officedocument.presentationml.presentation,application/octet-stream,text/plain,text/markdown,text/csv +STORAGE_REQUIRE_TENANT_PREFIX=true + +ALIYUN_OSS_REGION= +ALIYUN_OSS_ENDPOINT= +ALIYUN_OSS_ACCESS_KEY_ID= +ALIYUN_OSS_ACCESS_KEY_SECRET= +ALIYUN_OSS_STS_TOKEN= +ALIYUN_OSS_INTERNAL=false + +TENCENT_COS_REGION= +TENCENT_COS_APP_ID= +TENCENT_COS_SECRET_ID= +TENCENT_COS_SECRET_KEY= +TENCENT_COS_SECURITY_TOKEN= + +SUPABASE_STORAGE_URL= +SUPABASE_STORAGE_SERVICE_KEY= diff --git a/README.md b/README.md index 8aed5bdf..6d7aea36 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,7 @@ - 学生端能力:题库入口、分类树、题目集合、顺序/随机/模考 session 组卷快照、答题、错题本、收藏夹、背单词进度、个人中心、考试倒计时、签到积分、题目反馈、排行榜、分数线、题目视频、订单详情/状态轮询、优惠券领取/抵扣、权益、激活码预检查/兑换、资料下载。 - 平台后台能力:租户管理、SaaS 套餐、订阅、账单、服务费收款、用量记录。 - 销售/代理/CRM 增长链路:邀请码、扫码/分享事件、首绑客资保护、销售统计、团队关系、CRM 配置和队列。 -- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态,兜底漏通知订单。 +- `apps/worker` 后台任务进程:CRM webhook 队列消费、generic/钉钉/飞书/企微机器人发送、签名、失败重试和日志;commerce worker 可补偿查询微信/支付宝支付和退款状态;assets worker 可复检托管资源元数据并自动下架异常资源。 - 销售/代理分佣结算基础闭环:租户默认比例、成员比例、激活码批次比例、订单/激活码归因、结算单生成、审核、线下打款状态和权限隔离。 - 订单售后基础闭环:退款请求、审核、处理状态流、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、退款金额累计、部分/全额退款订单状态、全额退款权益撤销、退款事件和审计日志。 - PocketBase schema/数据导入器雏形和导入后校验脚本。 @@ -27,7 +27,7 @@ - Supabase Auth/JWT、租户角色模板、班级/教师/学生范围权限已可联调;生产前还要做真实云端 Auth/JWKS 回归和 RLS 深测。 - 阿里云/腾讯云短信、微信小程序登录、微信支付、支付宝主链路、微信/支付宝发起退款/查询确认/退款通知、支付/退款补偿 worker 已完成本地适配;微信网页登录、QQ 登录、手机号换绑、完整资金流水对账和真实生产账号联调还没接完。 -- OSS/COS/Supabase Storage 上传下载签名 provider 已接入;上传后校验、PDF 预览、防盗链和视频水印还没完成。 +- OSS/COS/Supabase Storage 上传下载签名 provider 已接入;上传后校验、PDF/图片预览和资源复检 worker 已完成,CDN 防盗链、杀毒扫描和视频动态水印还没完成。 - Excel/CSV 导入、分数线/视频批量导入和异步 worker 还没完成。 - 分佣真实打款、结算导出、发票/凭证、CRM 轮询/定向分配、富卡片模板、失败告警和销售转化看板还没完成。 - Taro 跨端前端还没开始 scaffold。 @@ -56,7 +56,7 @@ ```text apps/api/ Node.js 业务 API -apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、后续导入复检等 +apps/worker/ 后台异步任务:CRM webhook、支付/退款补偿、资源复检、后续导入复检等 packages/config/ 共享配置 packages/db/ PostgreSQL 连接池和查询封装 packages/domain/ 领域常量和共享类型 @@ -101,6 +101,12 @@ npm --workspace @tiku-saas/worker run crm:once npm --workspace @tiku-saas/worker run commerce:once ``` +单次运行内容资源复检 worker: + +```bash +npm --workspace @tiku-saas/worker run assets:once +``` + 默认本地数据库: ```text @@ -140,6 +146,7 @@ npm run pb:import:validate npm run test:api npm run test:worker:crm npm run test:worker:commerce +npm run test:worker:assets ``` ## API 模块 @@ -177,7 +184,7 @@ API 身份上下文: - 租户公开配置不能存放密钥。 - 商户密钥、短信密钥、OAuth app secret 等必须进入 `app_private.tenant_secrets`,或后续生产 KMS/Vault。 -- 资料、PDF、视频等资源必须先进入 `content_assets` 台账,再由 API 校验权限并下发签名 URL。 +- 资料、PDF、视频等资源必须先进入 `content_assets` 台账,再由 API 校验权限并下发签名 URL;生产环境应定时运行 assets worker 复检对象元数据,异常资源会被标记 failed 并退回 draft。 - 题库入口和分类使用 `content_entries/content_nodes`;题目列表和练习规则使用 `question_collections/practice_blueprints`,前端不要再把旧树字段当成唯一业务结构。 - 批量导入必须先写 `content_import_jobs/items/issues`,保留原始 payload、规范化 payload、逐行问题和审计记录。题目、单词、知识手册导入已走这套后台校验管线,前端只做预检查和预览展示。 - 支付 webhook 必须先设计幂等键和验签流程,再进入生产使用;生产环境还应定时运行 commerce worker 兜底供应商漏通知和处理中退款。 @@ -198,6 +205,6 @@ npm run check:refactor 1. 真实云端 Auth/JWKS 回归、RLS 深测和生产环境配置验收。 2. Taro 前端 scaffold,让 H5 和小程序共用同一套 API。 -3. 对象存储上传后校验、PDF 预览、防盗链和视频水印。 +3. 对象存储 CDN 防盗链、杀毒扫描、视频动态水印和生命周期策略。 4. Excel/CSV 以及分数线、视频批量导入;把现有 JSON 导入升级为可排队异步执行。 5. 微信网页/QQ 登录、完整资金流水对账、公共题库版本同步 worker、积分活动深化,以及排行榜防刷/预聚合。 diff --git a/apps/worker/package.json b/apps/worker/package.json index a29ad44d..f025b2f3 100644 --- a/apps/worker/package.json +++ b/apps/worker/package.json @@ -9,9 +9,12 @@ "build": "node -e \"fs.rmSync('dist',{recursive:true,force:true})\" && tsc -p tsconfig.json", "check": "tsc -p tsconfig.json --noEmit", "crm:once": "tsx src/index.ts --once --job crm", - "commerce:once": "tsx src/index.ts --once --job commerce" + "commerce:once": "tsx src/index.ts --once --job commerce", + "assets:once": "tsx src/index.ts --once --job assets" }, "dependencies": { + "@supabase/storage-js": "^2.108.2", + "ali-oss": "^6.23.0", "pg": "^8.16.3" }, "devDependencies": { diff --git a/apps/worker/src/config.ts b/apps/worker/src/config.ts index af2974b5..c64b5efe 100644 --- a/apps/worker/src/config.ts +++ b/apps/worker/src/config.ts @@ -13,6 +13,27 @@ export interface WorkerConfig { commerceBatchSize: number; commerceMinAgeSeconds: number; commerceRequestTimeoutMs: number; + assetBatchSize: number; + assetMinAgeSeconds: number; + assetRecheckIntervalSeconds: number; + assetRequestTimeoutMs: number; + storageMaxUploadBytes: number; + storageAllowedMimePrefixes: string[]; + storageAllowedMimeTypes: string[]; + storageRequireTenantPrefix: boolean; + aliyunOssRegion: string; + aliyunOssEndpoint: string; + aliyunOssAccessKeyId: string; + aliyunOssAccessKeySecret: string; + aliyunOssStsToken: string; + aliyunOssInternal: boolean; + tencentCosRegion: string; + tencentCosAppId: string; + tencentCosSecretId: string; + tencentCosSecretKey: string; + tencentCosSecurityToken: string; + supabaseStorageUrl: string; + supabaseStorageServiceKey: string; } export const config: WorkerConfig = { @@ -28,4 +49,43 @@ export const config: WorkerConfig = { commerceBatchSize: envNumber('WORKER_COMMERCE_BATCH_SIZE', 20), commerceMinAgeSeconds: envNumber('WORKER_COMMERCE_MIN_AGE_SECONDS', 300), commerceRequestTimeoutMs: envNumber('WORKER_COMMERCE_REQUEST_TIMEOUT_MS', 10_000), + assetBatchSize: envNumber('WORKER_ASSET_BATCH_SIZE', 50), + assetMinAgeSeconds: envNumber('WORKER_ASSET_MIN_AGE_SECONDS', 300), + assetRecheckIntervalSeconds: envNumber('WORKER_ASSET_RECHECK_INTERVAL_SECONDS', 60 * 60 * 24), + assetRequestTimeoutMs: envNumber('WORKER_ASSET_REQUEST_TIMEOUT_MS', 10_000), + storageMaxUploadBytes: envNumber('STORAGE_MAX_UPLOAD_BYTES', 1024 * 1024 * 500), + storageAllowedMimePrefixes: envList('STORAGE_ALLOWED_MIME_PREFIXES', 'image/,video/,audio/'), + storageAllowedMimeTypes: envList( + 'STORAGE_ALLOWED_MIME_TYPES', + [ + 'application/pdf', + 'application/json', + 'application/zip', + 'application/x-zip-compressed', + 'application/msword', + 'application/vnd.openxmlformats-officedocument.wordprocessingml.document', + 'application/vnd.ms-excel', + 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet', + 'application/vnd.ms-powerpoint', + 'application/vnd.openxmlformats-officedocument.presentationml.presentation', + 'application/octet-stream', + 'text/plain', + 'text/markdown', + 'text/csv', + ].join(','), + ), + storageRequireTenantPrefix: envBoolean('STORAGE_REQUIRE_TENANT_PREFIX', true), + aliyunOssRegion: envString('ALIYUN_OSS_REGION', ''), + aliyunOssEndpoint: envString('ALIYUN_OSS_ENDPOINT', ''), + aliyunOssAccessKeyId: envString('ALIYUN_OSS_ACCESS_KEY_ID', ''), + aliyunOssAccessKeySecret: envString('ALIYUN_OSS_ACCESS_KEY_SECRET', ''), + aliyunOssStsToken: envString('ALIYUN_OSS_STS_TOKEN', ''), + aliyunOssInternal: envBoolean('ALIYUN_OSS_INTERNAL', false), + tencentCosRegion: envString('TENCENT_COS_REGION', ''), + tencentCosAppId: envString('TENCENT_COS_APP_ID', ''), + tencentCosSecretId: envString('TENCENT_COS_SECRET_ID', ''), + tencentCosSecretKey: envString('TENCENT_COS_SECRET_KEY', ''), + tencentCosSecurityToken: envString('TENCENT_COS_SECURITY_TOKEN', ''), + supabaseStorageUrl: envString('SUPABASE_STORAGE_URL', ''), + supabaseStorageServiceKey: envString('SUPABASE_STORAGE_SERVICE_KEY', ''), }; diff --git a/apps/worker/src/index.ts b/apps/worker/src/index.ts index 6411c671..3a748f8f 100644 --- a/apps/worker/src/index.ts +++ b/apps/worker/src/index.ts @@ -2,6 +2,7 @@ import { closePool } from './db.js'; import { config } from './config.js'; import { processCrmBatch } from './jobs/crm.js'; import { processCommerceBatch } from './jobs/commerce.js'; +import { processAssetBatch } from './jobs/assets.js'; function hasArg(name: string) { return process.argv.includes(name); @@ -28,6 +29,14 @@ async function runOnce() { ); return; } + if (job === 'assets') { + const result = await processAssetBatch(); + console.log( + `[worker] assets batch processed=${result.processed}` + + ` verified=${result.verified} failed=${result.failed} skipped=${result.skipped} errors=${result.errors}`, + ); + return; + } throw new Error(`Unsupported worker job: ${job}`); } diff --git a/apps/worker/src/jobs/assets.ts b/apps/worker/src/jobs/assets.ts new file mode 100644 index 00000000..65a064d5 --- /dev/null +++ b/apps/worker/src/jobs/assets.ts @@ -0,0 +1,715 @@ +import crypto from 'node:crypto'; +import type pg from 'pg'; +import { pool } from '../db.js'; +import { config } from '../config.js'; + +type StorageProviderName = 'external_url' | 'supabase_storage' | 'aliyun_oss' | 'tencent_cos' | 'qiniu_kodo' | 'local_dev'; + +interface AssetCandidate { + id: string; + tenantId: string; + title: string | null; + status: string; + uploadStatus: string; + storageProvider: StorageProviderName; + bucket: string | null; + objectKey: string | null; + mimeType: string | null; + fileSizeBytes: number | string | null; + checksumSha256: string | null; + verifiedSizeBytes: number | string | null; + verifiedChecksumSha256: string | null; + verificationDetails: Record; + securityFlags: Record; +} + +interface StorageObjectMetadata { + provider: StorageProviderName; + bucket: string | null; + objectKey: string | null; + exists: boolean; + sizeBytes: number | null; + mimeType: string | null; + checksumSha256: string | null; + etag: string | null; + lastModified: string | null; + rawHeaders: Record; + verificationSource: string; +} + +interface AssetWorkerResult { + processed: number; + verified: number; + failed: number; + skipped: number; + errors: number; +} + +class AssetWorkerError extends Error { + readonly code: string; + + constructor(message: string, code = 'ASSET_WORKER_ERROR') { + super(message); + this.code = code; + } +} + +const UPLOADABLE_PROVIDERS = new Set(['local_dev', 'supabase_storage', 'aliyun_oss', 'tencent_cos']); +const SAFE_OBJECT_KEY_RE = /^[A-Za-z0-9][A-Za-z0-9._~!$&'()+,;=@/-]{0,1023}$/; + +function objectValue(value: unknown): Record { + return value && typeof value === 'object' && !Array.isArray(value) ? value as Record : {}; +} + +function nowIso() { + return new Date().toISOString(); +} + +function canonicalObjectKey(objectKey: string) { + return objectKey.replace(/^\/+/, '').replace(/\/{2,}/g, '/'); +} + +function validateObjectKey(tenantId: string, objectKey: string) { + const clean = canonicalObjectKey(objectKey); + if (!clean || clean.includes('..') || clean.includes('\\') || clean.includes('%2f') || clean.includes('%2F')) { + throw new AssetWorkerError('Invalid objectKey', 'INVALID_OBJECT_KEY'); + } + if (!SAFE_OBJECT_KEY_RE.test(clean)) { + throw new AssetWorkerError('objectKey contains unsafe characters', 'INVALID_OBJECT_KEY'); + } + if (config.storageRequireTenantPrefix && !clean.startsWith(`${tenantId}/`)) { + throw new AssetWorkerError('objectKey must be scoped by tenantId prefix', 'OBJECT_KEY_TENANT_PREFIX_REQUIRED'); + } + return clean; +} + +function normalizeSafeInteger(value: number | string | null) { + if (value === null) return null; + const numberValue = typeof value === 'number' ? value : Number(value); + if (!Number.isSafeInteger(numberValue)) { + throw new AssetWorkerError('fileSizeBytes must be a non-negative safe integer', 'INVALID_FILE_SIZE'); + } + return numberValue; +} + +function validateFileSize(fileSizeBytes: number | string | null) { + if (fileSizeBytes === null) return null; + const normalized = normalizeSafeInteger(fileSizeBytes); + if (normalized === null || normalized < 0) { + throw new AssetWorkerError('fileSizeBytes must be a non-negative safe integer', 'INVALID_FILE_SIZE'); + } + if (normalized > config.storageMaxUploadBytes) { + throw new AssetWorkerError('file exceeds STORAGE_MAX_UPLOAD_BYTES', 'FILE_TOO_LARGE'); + } + return normalized; +} + +function validateMimeType(mimeType: string | null) { + if (!mimeType) return null; + const normalized = mimeType.trim().toLowerCase(); + if (!normalized || normalized.length > 160 || normalized.includes('\r') || normalized.includes('\n')) { + throw new AssetWorkerError('Invalid mimeType', 'INVALID_MIME_TYPE'); + } + const exact = config.storageAllowedMimeTypes.map(item => item.toLowerCase()); + const prefixes = config.storageAllowedMimePrefixes.map(item => item.toLowerCase()); + if (!exact.includes(normalized) && !prefixes.some(prefix => normalized.startsWith(prefix))) { + throw new AssetWorkerError(`mimeType is not allowed: ${normalized}`, 'MIME_TYPE_NOT_ALLOWED'); + } + return normalized; +} + +function normalizeHeaderMap(headers: Headers | Record) { + const normalized: Record = {}; + if (headers instanceof Headers) { + headers.forEach((value, key) => { + normalized[key.toLowerCase()] = value; + }); + return normalized; + } + for (const [key, value] of Object.entries(headers)) { + if (value === undefined || value === null) continue; + normalized[key.toLowerCase()] = Array.isArray(value) ? String(value[0] || '') : String(value); + } + return normalized; +} + +function numberHeader(headers: Record, key: string) { + const value = Number(headers[key]); + return Number.isFinite(value) && value >= 0 ? Math.trunc(value) : null; +} + +function firstHeader(headers: Record, keys: string[]) { + for (const key of keys) { + const value = headers[key.toLowerCase()]; + if (value) return value; + } + return null; +} + +function cleanEtag(value: string | null) { + return value ? value.replace(/^"+|"+$/g, '') : null; +} + +function metadataFromHeaders(input: { + provider: StorageProviderName; + bucket: string | null; + objectKey: string | null; + headers: Record; + verificationSource: string; +}): StorageObjectMetadata { + const checksumSha256 = firstHeader(input.headers, [ + 'x-oss-meta-sha256', + 'x-oss-meta-checksum-sha256', + 'x-cos-meta-sha256', + 'x-cos-meta-checksum-sha256', + 'x-amz-meta-sha256', + 'x-amz-meta-checksum-sha256', + ]); + + return { + provider: input.provider, + bucket: input.bucket, + objectKey: input.objectKey, + exists: true, + sizeBytes: numberHeader(input.headers, 'content-length'), + mimeType: firstHeader(input.headers, ['content-type']), + checksumSha256: checksumSha256?.trim().toLowerCase() || null, + etag: cleanEtag(firstHeader(input.headers, ['etag'])), + lastModified: firstHeader(input.headers, ['last-modified']), + rawHeaders: input.headers, + verificationSource: input.verificationSource, + }; +} + +function hmacSha1Hex(key: string | Buffer, value: string) { + return crypto.createHmac('sha1', key).update(value).digest('hex'); +} + +function sha1Hex(value: string) { + return crypto.createHash('sha1').update(value).digest('hex'); +} + +function cosEncodePath(objectKey: string) { + return objectKey + .split('/') + .map(part => encodeURIComponent(part).replace(/[!'()*]/g, char => `%${char.charCodeAt(0).toString(16).toUpperCase()}`)) + .join('/'); +} + +function requireConfigured(condition: unknown, provider: StorageProviderName, missing: string) { + if (!condition) { + throw new AssetWorkerError(`${provider} is not configured: ${missing}`, 'STORAGE_PROVIDER_NOT_CONFIGURED'); + } +} + +function cosHost(bucket: string) { + requireConfigured(config.tencentCosRegion, 'tencent_cos', 'TENCENT_COS_REGION'); + const bucketWithAppId = config.tencentCosAppId && !bucket.endsWith(`-${config.tencentCosAppId}`) + ? `${bucket}-${config.tencentCosAppId}` + : bucket; + return `${bucketWithAppId}.cos.${config.tencentCosRegion}.myqcloud.com`; +} + +function nowSeconds() { + return Math.floor(Date.now() / 1000); +} + +function signTencentCosHead(input: { bucket: string; objectKey: string; expiresInSec: number }) { + requireConfigured(config.tencentCosSecretId, 'tencent_cos', 'TENCENT_COS_SECRET_ID'); + requireConfigured(config.tencentCosSecretKey, 'tencent_cos', 'TENCENT_COS_SECRET_KEY'); + const host = cosHost(input.bucket); + const start = nowSeconds(); + const end = start + input.expiresInSec; + const keyTime = `${start};${end}`; + const pathname = `/${cosEncodePath(input.objectKey)}`; + const signedHeaders: Record = { host }; + const headerKeys = Object.keys(signedHeaders).sort(); + const headerList = headerKeys.join(';'); + const httpHeaders = headerKeys + .map(key => `${encodeURIComponent(key)}=${encodeURIComponent(signedHeaders[key]).toLowerCase()}`) + .join('&'); + const signedQuery: Record = {}; + if (config.tencentCosSecurityToken) signedQuery['x-cos-security-token'] = config.tencentCosSecurityToken; + const queryKeys = Object.keys(signedQuery).sort(); + const urlParamList = queryKeys.join(';'); + const httpParameters = queryKeys + .map(key => `${encodeURIComponent(key)}=${encodeURIComponent(signedQuery[key])}`) + .join('&'); + const httpString = `head\n${pathname}\n${httpParameters}\n${httpHeaders}\n`; + const stringToSign = `sha1\n${keyTime}\n${sha1Hex(httpString)}\n`; + const signKey = hmacSha1Hex(config.tencentCosSecretKey, keyTime); + const signature = hmacSha1Hex(signKey, stringToSign); + const query = new URLSearchParams(); + query.set('q-sign-algorithm', 'sha1'); + query.set('q-ak', config.tencentCosSecretId); + query.set('q-sign-time', keyTime); + query.set('q-key-time', keyTime); + query.set('q-header-list', headerList); + query.set('q-url-param-list', urlParamList); + query.set('q-signature', signature); + for (const key of queryKeys) query.set(key, signedQuery[key]); + return `https://${host}${pathname}?${query.toString()}`; +} + +async function fetchWithTimeout(url: string, init: RequestInit, timeoutMs: number) { + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), timeoutMs); + try { + return await fetch(url, { ...init, signal: controller.signal }); + } finally { + clearTimeout(timeout); + } +} + +async function headAliyunOssObject(asset: AssetCandidate): Promise { + if (!asset.bucket || !asset.objectKey) { + throw new AssetWorkerError('Aliyun OSS asset requires bucket and objectKey', 'ASSET_OBJECT_LOCATION_REQUIRED'); + } + requireConfigured(config.aliyunOssAccessKeyId, 'aliyun_oss', 'ALIYUN_OSS_ACCESS_KEY_ID'); + requireConfigured(config.aliyunOssAccessKeySecret, 'aliyun_oss', 'ALIYUN_OSS_ACCESS_KEY_SECRET'); + requireConfigured(config.aliyunOssRegion || config.aliyunOssEndpoint, 'aliyun_oss', 'ALIYUN_OSS_REGION or ALIYUN_OSS_ENDPOINT'); + + const { default: OSS } = await import('ali-oss'); + const client = new OSS({ + region: config.aliyunOssRegion || undefined, + endpoint: config.aliyunOssEndpoint || undefined, + accessKeyId: config.aliyunOssAccessKeyId, + accessKeySecret: config.aliyunOssAccessKeySecret, + stsToken: config.aliyunOssStsToken || undefined, + bucket: asset.bucket, + internal: config.aliyunOssInternal, + secure: true, + }); + + try { + const response = await client.head(asset.objectKey); + const headers = normalizeHeaderMap(response.res?.headers || response); + return metadataFromHeaders({ + provider: 'aliyun_oss', + bucket: asset.bucket, + objectKey: asset.objectKey, + headers, + verificationSource: 'aliyun-oss-head-object', + }); + } catch (error) { + const status = Number((error as { status?: number; statusCode?: number }).status || (error as { statusCode?: number }).statusCode || 0); + if (status === 404) { + throw new AssetWorkerError('Object was not found in Aliyun OSS', 'STORAGE_OBJECT_NOT_FOUND'); + } + throw error; + } +} + +async function headTencentCosObject(asset: AssetCandidate): Promise { + if (!asset.bucket || !asset.objectKey) { + throw new AssetWorkerError('Tencent COS asset requires bucket and objectKey', 'ASSET_OBJECT_LOCATION_REQUIRED'); + } + const url = signTencentCosHead({ bucket: asset.bucket, objectKey: asset.objectKey, expiresInSec: 60 }); + const response = await fetchWithTimeout(url, { method: 'HEAD' }, config.assetRequestTimeoutMs); + if (response.status === 404) { + throw new AssetWorkerError('Object was not found in Tencent COS', 'STORAGE_OBJECT_NOT_FOUND'); + } + if (!response.ok) { + throw new AssetWorkerError(`Tencent COS object metadata check failed: ${response.status}`, 'STORAGE_HEAD_FAILED'); + } + return metadataFromHeaders({ + provider: 'tencent_cos', + bucket: asset.bucket, + objectKey: asset.objectKey, + headers: normalizeHeaderMap(response.headers), + verificationSource: 'tencent-cos-head-object', + }); +} + +async function headSupabaseStorageObject(asset: AssetCandidate): Promise { + if (!asset.bucket || !asset.objectKey) { + throw new AssetWorkerError('Supabase Storage asset requires bucket and objectKey', 'ASSET_OBJECT_LOCATION_REQUIRED'); + } + requireConfigured(config.supabaseStorageUrl, 'supabase_storage', 'SUPABASE_STORAGE_URL'); + requireConfigured(config.supabaseStorageServiceKey, 'supabase_storage', 'SUPABASE_STORAGE_SERVICE_KEY'); + const { StorageClient } = await import('@supabase/storage-js'); + const client = new StorageClient(config.supabaseStorageUrl.replace(/\/+$/, ''), { + apikey: config.supabaseStorageServiceKey, + authorization: `Bearer ${config.supabaseStorageServiceKey}`, + }); + const parts = asset.objectKey.split('/'); + const name = parts.pop() || ''; + const folder = parts.join('/'); + const response = await client.from(asset.bucket).list(folder || undefined, { + search: name, + limit: 20, + }); + if (response.error) { + throw new AssetWorkerError(response.error.message || 'Supabase Storage metadata check failed', 'STORAGE_HEAD_FAILED'); + } + const file = response.data?.find(item => item.name === name); + if (!file) { + throw new AssetWorkerError('Object was not found in Supabase Storage', 'STORAGE_OBJECT_NOT_FOUND'); + } + const metadata = (file as { metadata?: Record }).metadata || {}; + const headers: Record = {}; + for (const [key, value] of Object.entries(metadata)) { + if (value !== undefined && value !== null) headers[key.toLowerCase()] = String(value); + } + const size = Number(metadata.size ?? metadata.contentLength); + return { + provider: 'supabase_storage', + bucket: asset.bucket, + objectKey: asset.objectKey, + exists: true, + sizeBytes: Number.isFinite(size) && size >= 0 ? Math.trunc(size) : null, + mimeType: typeof metadata.mimetype === 'string' ? metadata.mimetype : typeof metadata.contentType === 'string' ? metadata.contentType : null, + checksumSha256: typeof metadata.sha256 === 'string' ? metadata.sha256.trim().toLowerCase() : null, + etag: typeof metadata.eTag === 'string' ? metadata.eTag : typeof metadata.etag === 'string' ? metadata.etag : null, + lastModified: typeof file.updated_at === 'string' ? file.updated_at : null, + rawHeaders: headers, + verificationSource: 'supabase-storage-list-metadata', + }; +} + +async function headStorageObject(asset: AssetCandidate): Promise { + if (!asset.objectKey) { + throw new AssetWorkerError('Asset requires objectKey', 'ASSET_LOCATION_REQUIRED'); + } + const objectKey = validateObjectKey(asset.tenantId, asset.objectKey); + const fileSizeBytes = validateFileSize(asset.fileSizeBytes); + validateFileSize(asset.verifiedSizeBytes); + const mimeType = validateMimeType(asset.mimeType); + const checksumSha256 = asset.checksumSha256?.trim().toLowerCase() || null; + + if (asset.storageProvider === 'local_dev') { + return { + provider: asset.storageProvider, + bucket: asset.bucket, + objectKey, + exists: true, + sizeBytes: fileSizeBytes, + mimeType, + checksumSha256, + etag: checksumSha256, + lastModified: nowIso(), + rawHeaders: {}, + verificationSource: 'local-dev-declared-metadata', + }; + } + if (asset.storageProvider === 'aliyun_oss') { + return headAliyunOssObject({ ...asset, objectKey }); + } + if (asset.storageProvider === 'tencent_cos') { + return headTencentCosObject({ ...asset, objectKey }); + } + if (asset.storageProvider === 'supabase_storage') { + return headSupabaseStorageObject({ ...asset, objectKey }); + } + throw new AssetWorkerError(`${asset.storageProvider} does not support managed upload recheck`, 'UPLOAD_RECHECK_PROVIDER_NOT_SUPPORTED'); +} + +function canonicalMime(value: string | null) { + return value?.split(';')[0]?.trim().toLowerCase() || null; +} + +function compareAssetMetadata(asset: AssetCandidate, metadata: StorageObjectMetadata) { + const issues: string[] = []; + const expectedSize = validateFileSize(asset.verifiedSizeBytes ?? asset.fileSizeBytes); + const expectedMime = canonicalMime(validateMimeType(asset.mimeType)); + const observedMime = canonicalMime(metadata.mimeType); + const expectedChecksum = (asset.verifiedChecksumSha256 || asset.checksumSha256 || '').trim().toLowerCase() || null; + + if (expectedSize !== null && metadata.sizeBytes !== null && expectedSize !== metadata.sizeBytes) { + issues.push('file_size_mismatch'); + } + if (expectedMime && observedMime && expectedMime !== observedMime) { + issues.push('mime_type_mismatch'); + } + if (expectedChecksum && metadata.checksumSha256 && expectedChecksum !== metadata.checksumSha256) { + issues.push('checksum_mismatch'); + } + + return { + issues, + expected: { + fileSizeBytes: expectedSize, + mimeType: expectedMime, + checksumSha256: expectedChecksum, + }, + observed: { + fileSizeBytes: metadata.sizeBytes, + mimeType: observedMime, + checksumSha256: metadata.checksumSha256, + etag: metadata.etag, + lastModified: metadata.lastModified, + verificationSource: metadata.verificationSource, + }, + checksumVerified: Boolean(expectedChecksum && metadata.checksumSha256 && expectedChecksum === metadata.checksumSha256), + checksumUnavailable: Boolean(expectedChecksum && !metadata.checksumSha256), + }; +} + +function eventError(error: unknown) { + if (error instanceof AssetWorkerError) { + return { code: error.code, message: error.message }; + } + if (error instanceof Error) { + return { code: 'ASSET_WORKER_ERROR', message: error.message }; + } + return { code: 'ASSET_WORKER_ERROR', message: String(error) }; +} + +function mergeVerificationDetails(asset: AssetCandidate, assetWorker: Record) { + return { + ...asset.verificationDetails, + assetWorker: { + ...objectValue(asset.verificationDetails.assetWorker), + ...assetWorker, + }, + }; +} + +async function recordAudit( + client: pg.PoolClient, + input: { + tenantId: string; + action: string; + targetId: string; + details: Record; + }, +) { + await client.query( + ` + insert into public.audit_logs (tenant_id, actor_user_id, action, target_type, target_id, details) + values ($1, null, $2, 'content_asset', $3, $4::jsonb) + `, + [input.tenantId, input.action, input.targetId, JSON.stringify(input.details)], + ); +} + +async function claimAssetCandidates(client: pg.PoolClient, limit: number, claimId: string) { + const result = await client.query( + ` + with candidates as ( + select id + from public.content_assets + where storage_provider = any($1::text[]) + and object_key is not null + and upload_status in ('pending', 'verified') + and ( + (upload_status = 'pending' and updated_at <= now() - ($2::int * interval '1 second')) + or ( + upload_status = 'verified' + and coalesce((verification_details #>> '{assetWorker,lastCheckedAt}')::timestamptz, verified_at, updated_at, created_at) + <= now() - ($3::int * interval '1 second') + ) + ) + and coalesce(verification_details #>> '{assetWorker,claimId}', '') <> $5 + order by + case when upload_status = 'pending' then 0 else 1 end, + coalesce((verification_details #>> '{assetWorker,lastCheckedAt}')::timestamptz, verified_at, updated_at, created_at) asc + limit $4 + for update skip locked + ) + update public.content_assets ca + set verification_details = jsonb_set( + coalesce(ca.verification_details, '{}'::jsonb), + '{assetWorker}', + coalesce(ca.verification_details->'assetWorker', '{}'::jsonb) || $6::jsonb, + true + ), + updated_at = now() + from candidates + where ca.id = candidates.id + returning ca.id, + ca.tenant_id as "tenantId", + ca.title, + ca.status, + ca.upload_status as "uploadStatus", + ca.storage_provider as "storageProvider", + ca.bucket, + ca.object_key as "objectKey", + ca.mime_type as "mimeType", + ca.file_size_bytes as "fileSizeBytes", + ca.checksum_sha256 as "checksumSha256", + ca.verified_size_bytes as "verifiedSizeBytes", + ca.verified_checksum_sha256 as "verifiedChecksumSha256", + ca.verification_details as "verificationDetails", + ca.security_flags as "securityFlags" + `, + [ + Array.from(UPLOADABLE_PROVIDERS), + config.assetMinAgeSeconds, + config.assetRecheckIntervalSeconds, + limit, + claimId, + JSON.stringify({ claimId, claimedAt: nowIso() }), + ], + ); + return result.rows; +} + +async function markAssetVerified(client: pg.PoolClient, asset: AssetCandidate, metadata: StorageObjectMetadata, comparison: ReturnType) { + const details = mergeVerificationDetails(asset, { + lastCheckedAt: nowIso(), + lastResult: 'verified', + issues: [], + observed: comparison.observed, + expected: comparison.expected, + checksumVerified: comparison.checksumVerified, + checksumUnavailable: comparison.checksumUnavailable, + }); + + await client.query( + ` + update public.content_assets + set upload_status = 'verified', + verified_at = coalesce(verified_at, now()), + verified_size_bytes = coalesce($3::bigint, verified_size_bytes), + verified_checksum_sha256 = coalesce($4, verified_checksum_sha256), + file_size_bytes = coalesce(file_size_bytes, $3::bigint), + mime_type = coalesce(mime_type, $5), + verification_details = $6::jsonb, + security_flags = coalesce(security_flags, '{}'::jsonb) - 'assetRecheckFailed', + updated_at = now() + where tenant_id = $1 and id = $2 + `, + [ + asset.tenantId, + asset.id, + metadata.sizeBytes ?? comparison.expected.fileSizeBytes, + metadata.checksumSha256 ?? comparison.expected.checksumSha256, + comparison.observed.mimeType || comparison.expected.mimeType, + JSON.stringify(details), + ], + ); + + await recordAudit(client, { + tenantId: asset.tenantId, + action: 'content.asset.rechecked', + targetId: asset.id, + details: { + provider: asset.storageProvider, + bucket: asset.bucket, + objectKey: asset.objectKey, + result: 'verified', + checksumUnavailable: comparison.checksumUnavailable, + }, + }); +} + +async function markAssetFailed( + client: pg.PoolClient, + asset: AssetCandidate, + issues: string[], + metadata: StorageObjectMetadata | null, + error: Record | null, +) { + const details = mergeVerificationDetails(asset, { + lastCheckedAt: nowIso(), + lastResult: 'failed', + issues, + observed: metadata + ? { + fileSizeBytes: metadata.sizeBytes, + mimeType: canonicalMime(metadata.mimeType), + checksumSha256: metadata.checksumSha256, + etag: metadata.etag, + lastModified: metadata.lastModified, + verificationSource: metadata.verificationSource, + } + : null, + error, + }); + + await client.query( + ` + update public.content_assets + set upload_status = 'failed', + status = case when status = 'active' then 'draft' else status end, + verification_details = $3::jsonb, + security_flags = coalesce(security_flags, '{}'::jsonb) + || jsonb_build_object('assetRecheckFailed', true, 'assetRecheckFailedAt', now()), + updated_at = now() + where tenant_id = $1 and id = $2 + `, + [asset.tenantId, asset.id, JSON.stringify(details)], + ); + + await recordAudit(client, { + tenantId: asset.tenantId, + action: 'content.asset.recheck_failed', + targetId: asset.id, + details: { + provider: asset.storageProvider, + bucket: asset.bucket, + objectKey: asset.objectKey, + result: 'failed', + issues, + error, + unpublished: asset.status === 'active', + }, + }); +} + +async function processAsset(asset: AssetCandidate) { + const client = await pool.connect(); + try { + const metadata = await headStorageObject(asset); + const comparison = compareAssetMetadata(asset, metadata); + await client.query('begin'); + if (comparison.issues.length) { + await markAssetFailed(client, asset, comparison.issues, metadata, null); + await client.query('commit'); + return 'failed'; + } + await markAssetVerified(client, asset, metadata, comparison); + await client.query('commit'); + return 'verified'; + } catch (error) { + await client.query('rollback').catch(() => {}); + try { + await client.query('begin'); + await markAssetFailed(client, asset, [eventError(error).code], null, eventError(error)); + await client.query('commit'); + } catch { + await client.query('rollback').catch(() => {}); + } + return 'error'; + } finally { + client.release(); + } +} + +export async function processAssetBatch(limit = config.assetBatchSize): Promise { + const claimId = crypto.randomUUID(); + const client = await pool.connect(); + let assets: AssetCandidate[] = []; + try { + await client.query('begin'); + assets = await claimAssetCandidates(client, limit, claimId); + await client.query('commit'); + } catch (error) { + await client.query('rollback'); + throw error; + } finally { + client.release(); + } + + const result: AssetWorkerResult = { + processed: assets.length, + verified: 0, + failed: 0, + skipped: 0, + errors: 0, + }; + + for (const asset of assets) { + if (!UPLOADABLE_PROVIDERS.has(asset.storageProvider)) { + result.skipped += 1; + continue; + } + const status = await processAsset(asset); + if (status === 'verified') result.verified += 1; + else if (status === 'failed') result.failed += 1; + else result.errors += 1; + } + + return result; +} diff --git a/apps/worker/src/types/ali-oss.d.ts b/apps/worker/src/types/ali-oss.d.ts new file mode 100644 index 00000000..badb5f68 --- /dev/null +++ b/apps/worker/src/types/ali-oss.d.ts @@ -0,0 +1,25 @@ +declare module 'ali-oss' { + interface ClientOptions { + region?: string; + endpoint?: string; + accessKeyId: string; + accessKeySecret: string; + stsToken?: string; + bucket?: string; + internal?: boolean; + secure?: boolean; + } + + interface HeadObjectResult { + res?: { + headers?: Record; + status?: number; + }; + [key: string]: unknown; + } + + export default class OSS { + constructor(options: ClientOptions); + head(name: string): Promise; + } +} diff --git a/docs/refactor/backend-capability-status.md b/docs/refactor/backend-capability-status.md index f5cea6a6..61758076 100644 --- a/docs/refactor/backend-capability-status.md +++ b/docs/refactor/backend-capability-status.md @@ -93,7 +93,8 @@ | Supabase Storage 签名 | 可联调 | `supabase_storage` provider | | 上传后对象校验 | 可联调 | `/api/tenant-content/assets/confirm-upload`;托管对象必须 verified 后才能发布/下载 | | PDF/图片预览签名 | 可联调 | `/api/catalog/assets/preview`、`/api/tenant-content/assets/sign-preview`;使用 inline 短期签名 | -| 深度防盗链/水印/杀毒 | 待补齐 | 商用上线前补 worker、CDN 防盗链、动态水印和安全扫描 | +| 托管资源 worker 复检 | 可联调 | `apps/worker --job assets` 定期复检 pending/verified 对象元数据;异常资源会标记 failed 并从 active 退回 draft,写入审计和 `security_flags` | +| 深度防盗链/水印/杀毒 | 待补齐 | 商用上线前继续补 CDN 防盗链、动态水印、安全扫描和对象生命周期策略 | ## 订单、会员、营销 diff --git a/docs/refactor/backend-handoff-roadmap.md b/docs/refactor/backend-handoff-roadmap.md index a7201aef..b20baa17 100644 --- a/docs/refactor/backend-handoff-roadmap.md +++ b/docs/refactor/backend-handoff-roadmap.md @@ -28,7 +28,7 @@ | 知识手册 | 可联调 | 科目、章节、条目、Markdown 内容、嵌套 JSON 导入 | 富文本资源、版本管理、附件/PDF 关联 | | 分数线 | 可联调 | 院校、专业、动态字段、记录、年份、趋势、后台维护 | 批量导入、复杂筛选、AI 择校上下文 | | 视频解析 | 部分完成 | 单题视频、批量查询、后台视频绑定 | 会员播放权限、播放次数扣减、签名 URL、防盗链、水印 | -| 资料下载 | 部分完成 | 资源台账、SVIP 权限校验、`local_dev`/阿里云 OSS/腾讯 COS/Supabase Storage 上传下载签名 | 上传后对象校验、PDF 预览、防盗链、视频水印 | +| 资料下载 | 部分完成 | 资源台账、SVIP 权限校验、`local_dev`/阿里云 OSS/腾讯 COS/Supabase Storage 上传下载签名、上传确认、PDF/图片预览签名、assets worker 复检异常下架 | PDF 渲染、CDN 防盗链、杀毒扫描、视频水印 | | 会员与订单 | 可联调 | 下单、订单详情/状态轮询、优惠券领取/抵扣、零元订单自动开通、手工确认权限保护、激活码预检查/兑换、微信支付、支付宝、微信/支付宝发起退款、微信/支付宝退款查询确认、微信/支付宝退款通知 webhook、支付/退款补偿 worker、权益发放 | 完整资金流水对账、异常订单运营台 | | 登录认证 | 迁移期可用 | 短信 mock、迁移期 session、OAuth 配置表 | 阿里云/腾讯云短信、微信小程序/网页登录、QQ 登录、Supabase Auth | | 销售/代理/CRM | 基础完成 | 邀请码、首绑保护、团队关系、销售统计、CRM 入队 | 小程序码真实生成、分佣结算、钉钉/飞书/企微 worker | @@ -76,7 +76,7 @@ ### P0:上云测试和前端主链路前必须处理 - 生产鉴权:API 已支持 Supabase Auth JWT;继续做真实云端 Auth/JWKS 回归、RLS 深测,并在生产关闭 `x-user-id` 与 `x-platform-admin-key` 兼容入口。 -- 对象存储:上传/下载签名已接入阿里云 OSS、腾讯云 COS、Supabase Storage;继续完成上传后对象校验、PDF 预览、视频播放签名、防盗链和水印。 +- 对象存储:上传/下载签名已接入阿里云 OSS、腾讯云 COS、Supabase Storage;上传确认、PDF/图片预览签名和 assets worker 复检已完成,继续补 PDF 渲染、视频播放防盗链、杀毒扫描和水印。 - 真实数据 dry-run:导出 PocketBase 用户、题库、单词、知识手册、分数线、订单、权益,跑迁移和校验报告。 - 生产环境配置:补 `.env` 模板、数据库迁移流程、备份恢复、日志、告警和 API 容器部署说明。 - Taro scaffold:建立 `apps/taro`,先完成租户解析、首页、题库、背单词、知识手册、个人中心主链路。 diff --git a/docs/refactor/backend-progress.md b/docs/refactor/backend-progress.md index c1525725..8175fa4b 100644 --- a/docs/refactor/backend-progress.md +++ b/docs/refactor/backend-progress.md @@ -30,6 +30,7 @@ - 已新增 `npm run test:api`,自动 seed、构建、启动临时 API,并断言核心学生端接口、内容导航/组卷、租户隔离、资源权限和题目导入。 - 已新增 `apps/worker` 和 `npm run test:worker:crm`,用于消费 CRM webhook 队列,验证本地 fake webhook、队列状态、日志和密钥不泄露。 - 已新增 commerce worker 和 `npm run test:worker:commerce`,用于补偿查询微信/支付宝支付、处理中退款和漏通知场景;支付成功会幂等更新订单/支付并开通权益,退款成功会幂等更新退款/订单/支付并在全额退款时撤销订单权益,测试覆盖密钥不泄露和重复执行不重复开通。 +- 已新增 assets worker 和 `npm run test:worker:assets`,用于复检 `content_assets` 托管对象元数据;正常资源会写入复检证据,异常资源会自动下架为 `draft`、标记 `upload_status=failed`,并记录审计与安全标记。 ## 已验证接口 @@ -237,14 +238,14 @@ GET /api/tenant-admin/audit-logs - 学生批量导入、批量分班、学生状态、备注和跟进任务都使用独立权限点;教师默认可为范围内学生写备注和跟进任务,但不能批量导入、禁用学生或放大可见班级。 - 销售/代理客资采用首绑保护:普通扫码/分享事件不会覆盖已有归属,只有具备 `referral:write` 的租户成员可手动强制补绑。 - CRM 当前完成配置、密钥入私密表、客资入队、队列查询和 `apps/worker` 消费;worker 支持 generic webhook、钉钉、飞书、企业微信机器人消息体/签名、失败重试和日志。 -- 内容资源当前完成台账、租户后台维护、学生端 SVIP 下载权限,以及 `local_dev`、阿里云 OSS、腾讯 COS、Supabase Storage 的上传/下载签名 provider。真实对象存在性校验、PDF 预览渲染、防盗链、水印和大文件上传后 worker 校验仍需继续补。 +- 内容资源当前完成台账、租户后台维护、学生端 SVIP 下载权限,以及 `local_dev`、阿里云 OSS、腾讯 COS、Supabase Storage 的上传/下载签名 provider;上传确认和 assets worker 已支持对象元数据校验/复检。PDF 预览渲染、防盗链、水印和安全扫描仍需继续补。 - 题库内容导航当前以 `content_entries/content_nodes` 为主模型,可表达“入口 -> 多级分类 -> 院校/专业/学科/销售意向标记”;题目集合和练习方式由 `question_collections/practice_blueprints` 管理,练习 session 会保存当次题目 ID 快照。 - 练习访问控制由 `content_entries/content_nodes/question_collections/practice_blueprints` 的 `accessRules` 合并决定;普通用户消耗 `practice_daily_usage`,事件写入 `practice_access_events`,SVIP/staff 不消耗免费额度。 - 批量导入当前支持题目、单词、知识手册 JSON 预览、逐行 issue、job/item 台账、执行导入、幂等跳过,并可落到新内容入口和分类节点。旧单词模板的 `vocabulary_units_示例数据` / `vocabulary_示例数据`、知识手册的书籍/章节/小节/知识点嵌套结构都由后端规范化。Excel/CSV、分数线/视频导入会继续复用同一套 `content_import_jobs` 管线。 ## 下一步 -1. 完善内容导入和文件上传:Excel/CSV、分数线、视频导入,对象存储上传后校验、PDF 预览、防盗链和视频水印。 +1. 完善内容导入和文件上传:Excel/CSV、分数线、视频导入,PDF 预览渲染、防盗链、杀毒扫描和视频水印。 2. 接入真实短信 provider:阿里云/腾讯云,密钥放 `app_private.tenant_secrets` 或生产 Vault。 3. 接入真实 OAuth provider:微信网页、微信小程序、QQ,并处理旧 PocketBase 身份映射。 4. 补完整资金流水对账、异常订单运营台和优惠券核销报表;支付/退款补偿、退款查询确认和退款通知主链路已完成。 @@ -257,5 +258,6 @@ GET /api/tenant-admin/audit-logs npm run test:api npm run test:worker:crm npm run test:worker:commerce +npm run test:worker:assets npm run check:refactor ``` diff --git a/docs/refactor/next-development-todo.md b/docs/refactor/next-development-todo.md index ac915ada..19b6193d 100644 --- a/docs/refactor/next-development-todo.md +++ b/docs/refactor/next-development-todo.md @@ -24,6 +24,7 @@ - 公共题库商业化基础闭环已完成:平台公共题库可由平台管理员按 SaaS 套餐/指定租户/全部活跃租户授权;租户内容管理员只能看到自己被授权的公共题库,并可采纳为本租户题库、内容入口、题目集合和题目快照,采纳后可直接进入练习 session。 - 租户后台数据看板已完成首版聚合 API:`GET /api/tenant-admin/dashboard`,支持租户/地区维度的收益、注册、学习、内容、激活码、反馈、趋势、24h 活跃、套餐销量和运营动态,前端可直接联调。 - 支付/退款补偿 worker 已完成:`apps/worker --job commerce` 可查询微信/支付宝支付和处理中退款,补偿漏通知订单,支付成功幂等开通权益,退款成功幂等更新退款/订单/支付并在全额退款时撤销订单权益。 +- 内容资源复检 worker 已完成:`apps/worker --job assets` 可复检 `content_assets` 中的托管对象元数据,正常资源写回复检证据,异常资源自动置为 `failed + draft` 并写入审计和安全标记。 - 本地验证:`npm run check:refactor` 已通过。 当前更适合进入前端联调前阅读的总览文档: @@ -43,7 +44,8 @@ 2. 对象存储 - 已接阿里云 OSS、腾讯云 COS、Supabase Storage 的上传/下载签名 provider。 - 已补上传后对象确认接口、托管对象发布前 verified 校验、PDF/图片 inline 预览签名。 - - 继续补视频深度防盗链、动态水印、worker 复检、杀毒扫描、CDN 刷新和对象生命周期策略。 + - 已补 assets worker 复检,异常托管对象会自动下架并记录审计。 + - 继续补视频深度防盗链、动态水印、杀毒扫描、CDN 刷新和对象生命周期策略。 - `content_assets` 继续作为资源台账,不允许前端绕过台账直接访问私有资源。 3. 真实导入 dry-run diff --git a/docs/refactor/object-storage.md b/docs/refactor/object-storage.md index f3a8cec6..d0dfd089 100644 --- a/docs/refactor/object-storage.md +++ b/docs/refactor/object-storage.md @@ -64,6 +64,12 @@ POST /api/tenant-content/assets/sign-download POST /api/tenant-content/assets/sign-preview ``` +后台资源复检 worker: + +```bash +npm --workspace @tiku-saas/worker run assets:once +``` + ## 标准上传流程 后台前端上传 PDF、图片、视频或资料包时必须走下面流程: @@ -74,6 +80,8 @@ POST /api/tenant-content/assets/sign-preview 4. 调用 `POST /api/tenant-content/assets/confirm-upload`,由后端读取对象元数据并比对大小、MIME、SHA-256。 5. 校验通过且 `publish=true` 时,后端将资源置为 `status=active`、`uploadStatus=verified`。 6. 学生端只能下载或预览 `active + verified` 的托管对象资源。 +7. 生产环境定时运行 assets worker,复检 `pending/verified` 托管对象的大小、MIME、SHA-256 等元数据。 +8. 如果复检发现对象丢失、跨租户 objectKey、大小/MIME/checksum 不一致,worker 会把资源置为 `uploadStatus=failed`,并将 `active` 资源退回 `draft`,同时写入 `security_flags.assetRecheckFailed=true` 和 `audit_logs`。 确认上传示例: @@ -94,6 +102,7 @@ POST /api/tenant-content/assets/sign-preview - `tencent_cos` 使用 COS HEAD Object 预签名请求读取元数据。 - `supabase_storage` 使用 Storage list metadata 做最小存在性/大小校验。 - 如果 provider 无法返回 SHA-256,后端会记录 `checksumUnavailable=true`,生产建议上传时写入对象自定义元数据,例如 `x-oss-meta-sha256` 或 `x-cos-meta-sha256`。 +- worker 的复检证据写入 `verification_details.assetWorker`,包含 `lastCheckedAt`、`lastResult`、`expected`、`observed`、`issues`,便于租户后台定位资源异常。 ## 安全规则 @@ -106,6 +115,7 @@ POST /api/tenant-content/assets/sign-preview - 下载和视频播放必须先经过 API 权限判断,再下发短期签名 URL。 - 云厂商 AccessKey、SecretKey、Service Role Key 只存在服务端环境变量,不返回前端。 - `content_assets` 是资源唯一台账,前端不得绕过台账直接访问私有 bucket。 +- 前端不能把 `uploadStatus=failed` 或 `status=draft` 的资源继续展示为可下载;列表仍返回时应展示“资料处理中”或“资源异常已下架”,真正下载/预览会被后端拒绝。 ## 环境变量 @@ -149,6 +159,15 @@ SUPABASE_STORAGE_URL=https://your-project.supabase.co/storage/v1 SUPABASE_STORAGE_SERVICE_KEY= ``` +assets worker: + +```text +WORKER_ASSET_BATCH_SIZE=50 +WORKER_ASSET_MIN_AGE_SECONDS=300 +WORKER_ASSET_RECHECK_INTERVAL_SECONDS=86400 +WORKER_ASSET_REQUEST_TIMEOUT_MS=10000 +``` + ## 生产建议 - 阿里云和腾讯云生产环境优先用 STS/临时密钥或 RAM/CAM 最小权限账号。 @@ -156,7 +175,7 @@ SUPABASE_STORAGE_SERVICE_KEY= - 图片、PDF、视频分别设置合理的 CORS,只允许前端域名和小程序业务域名访问。 - 开启对象版本控制、生命周期、跨区域复制或定时备份,满足后续容灾要求。 - 视频资源已接入 SVIP/播放次数校验、短期签名和播放日志;生产阶段继续补转码、动态水印、CDN 防盗链和播放统计。 -- 大文件上传已经支持 API 即时确认;后续可增加 worker 做异步复检、杀毒、转码、水印和 CDN 刷新。 +- 大文件上传已经支持 API 即时确认和 worker 元数据复检;后续继续补杀毒、转码、水印和 CDN 刷新。 ## 官方依据 diff --git a/docs/refactor/taro-frontend-integration.md b/docs/refactor/taro-frontend-integration.md index e6337d00..d908b30b 100644 --- a/docs/refactor/taro-frontend-integration.md +++ b/docs/refactor/taro-frontend-integration.md @@ -575,6 +575,7 @@ GET /api/catalog/assets/download?assetId= - `download.url` 是短期 attachment URL,只给下载动作使用。 - `ASSET_SVIP_REQUIRED`:提示开通对应地区/科目权益。 - `ASSET_UPLOAD_NOT_VERIFIED`:展示“资料正在处理中”,并上报前端日志。 +- `ASSET_NOT_FOUND` 或列表中资源从 `active` 消失:展示“资源异常已下架”或刷新列表,不要继续使用旧签名 URL。 - `ASSET_PREVIEW_NOT_SUPPORTED`:隐藏预览按钮,仅保留下载或提示不支持预览。 - `previewUrl` 字段只作为公开/托管预览提示,不代表可以绕过接口直接访问。 @@ -598,6 +599,8 @@ sign-upload -> 直传对象存储 -> PUT assets 登记草稿 -> confirm-upload - 托管对象在确认前会保持 `status=draft`、`uploadStatus=pending`,学生端不会看到。确认失败时后端返回 `UPLOAD_VERIFICATION_FAILED`,后台必须展示失败原因并允许重新上传,不能前端强行改为已发布。 +生产环境会定时运行 assets worker 复检对象存储元数据。复检发现对象丢失、跨租户 objectKey、大小/MIME/checksum 不一致时,后端会把资源置为 `uploadStatus=failed` 并从 `active` 退回 `draft`,同时写入 `securityFlags.assetRecheckFailed=true`。租户后台资源列表应对 failed 资源展示异常原因和重新上传入口;学生端不要缓存资料列表和签名 URL 作为长期状态。 + ## 考试倒计时、签到积分和反馈 首页可用 `GET /api/catalog/exam-dates?regionId=` 展示地区公开考试日期;个人中心优先用 `GET /api/profile/exam-countdowns`,后端会按学生当前 `regionId/selectedSchoolId` 返回匹配倒计时。 diff --git a/package-lock.json b/package-lock.json index 554793f4..9ec8df85 100644 --- a/package-lock.json +++ b/package-lock.json @@ -54,6 +54,8 @@ "name": "@tiku-saas/worker", "version": "0.1.0", "dependencies": { + "@supabase/storage-js": "^2.108.2", + "ali-oss": "^6.23.0", "pg": "^8.16.3" }, "devDependencies": { diff --git a/package.json b/package.json index 303f4cf4..08425839 100644 --- a/package.json +++ b/package.json @@ -35,6 +35,7 @@ "test:api": "npm run db:smoke-seed && npm run build:api && node scripts/api-integration-test.js --start-server", "test:worker:crm": "npm run db:smoke-seed && npm run build:worker && node scripts/crm-worker-integration-test.js", "test:worker:commerce": "npm run db:smoke-seed && npm run build:worker && node scripts/commerce-worker-integration-test.js", + "test:worker:assets": "npm run db:smoke-seed && npm run build:worker && node scripts/asset-worker-integration-test.js", "test:api:remote": "node scripts/api-integration-test.js", "pb:schema:summary": "npm --workspace @tiku-saas/import-pocketbase run schema:summary", "pb:schema:risk": "npm --workspace @tiku-saas/import-pocketbase run schema:risk", diff --git a/scripts/asset-worker-integration-test.js b/scripts/asset-worker-integration-test.js new file mode 100644 index 00000000..8afebfb9 --- /dev/null +++ b/scripts/asset-worker-integration-test.js @@ -0,0 +1,175 @@ +import assert from 'node:assert/strict'; +import pg from 'pg'; +import { spawn } from 'node:child_process'; + +const databaseUrl = process.env.DATABASE_URL || 'postgresql://postgres:postgres@127.0.0.1:54322/postgres'; +const tenantId = '00000000-0000-0000-0000-000000000001'; + +const ids = { + okAsset: '20000000-0000-0000-0000-000000000801', + badAsset: '20000000-0000-0000-0000-000000000802', +}; + +const checksumA = 'a'.repeat(64); +const checksumB = 'b'.repeat(64); + +async function runWorkerOnce() { + const child = spawn(process.execPath, ['apps/worker/dist/apps/worker/src/index.js', '--once', '--job', 'assets'], { + cwd: process.cwd(), + env: { + ...process.env, + DATABASE_URL: databaseUrl, + WORKER_ASSET_BATCH_SIZE: '20', + WORKER_ASSET_MIN_AGE_SECONDS: '0', + WORKER_ASSET_RECHECK_INTERVAL_SECONDS: '0', + WORKER_ASSET_REQUEST_TIMEOUT_MS: '5000', + STORAGE_REQUIRE_TENANT_PREFIX: 'true', + }, + stdio: ['ignore', 'pipe', 'pipe'], + windowsHide: true, + }); + let output = ''; + child.stdout.on('data', chunk => { + output += chunk.toString(); + }); + child.stderr.on('data', chunk => { + output += chunk.toString(); + }); + const code = await new Promise(resolve => child.on('exit', resolve)); + assert.equal(code, 0, `worker should exit 0\n${output}`); + assert.match(output, /assets batch processed=\d+/, 'worker output should include assets summary'); + return output; +} + +async function cleanup(pool) { + await pool.query( + ` + delete from public.audit_logs + where tenant_id = $1 + and target_type = 'content_asset' + and target_id in ($2, $3) + `, + [tenantId, ids.okAsset, ids.badAsset], + ); + await pool.query( + ` + delete from public.content_assets + where tenant_id = $1 and id in ($2::uuid, $3::uuid) + `, + [tenantId, ids.okAsset, ids.badAsset], + ); +} + +async function seed(pool) { + await pool.query( + ` + insert into public.content_assets ( + id, tenant_id, asset_key, title, asset_type, storage_provider, + bucket, object_key, file_name, mime_type, file_size_bytes, checksum_sha256, + visibility, status, upload_status, verified_at, verified_size_bytes, + verified_checksum_sha256, verification_details, security_flags, source, + created_at, updated_at + ) + values + ( + $1, $2, 'asset-worker-ok', '资源复检正常 PDF', 'pdf', 'local_dev', + 'tenant-assets', $3, 'ok.pdf', 'application/pdf', 4096, $4, + 'tenant', 'active', 'verified', now() - interval '2 days', 4096, + $4, '{"source":"asset-worker-test"}'::jsonb, '{}'::jsonb, 'integration-test', + now() - interval '2 days', now() - interval '2 days' + ), + ( + $5, $2, 'asset-worker-bad', '资源复检异常 PDF', 'pdf', 'local_dev', + 'tenant-assets', $6, 'bad.pdf', 'application/pdf', 1024, $7, + 'tenant', 'active', 'verified', now() - interval '2 days', 2048, + $7, '{"source":"asset-worker-test"}'::jsonb, '{}'::jsonb, 'integration-test', + now() - interval '2 days', now() - interval '2 days' + ) + `, + [ + ids.okAsset, + tenantId, + `${tenantId}/assets/worker-ok.pdf`, + checksumA, + ids.badAsset, + `${tenantId}/assets/worker-bad.pdf`, + checksumB, + ], + ); +} + +async function main() { + const pool = new pg.Pool({ connectionString: databaseUrl }); + let seeded = false; + try { + await pool.query('begin'); + await cleanup(pool); + await seed(pool); + await pool.query('commit'); + seeded = true; + + const output = await runWorkerOnce(); + assert.match(output, /verified=1/, 'worker should verify one asset'); + assert.match(output, /failed=1/, 'worker should fail one mismatched asset'); + + const assets = await pool.query( + ` + select id, status, upload_status, verification_details, security_flags + from public.content_assets + where tenant_id = $1 and id in ($2::uuid, $3::uuid) + order by id + `, + [tenantId, ids.okAsset, ids.badAsset], + ); + const okAsset = assets.rows.find(row => row.id === ids.okAsset); + const badAsset = assets.rows.find(row => row.id === ids.badAsset); + assert.equal(okAsset?.status, 'active', 'verified asset should remain active'); + assert.equal(okAsset?.upload_status, 'verified', 'verified asset should remain verified'); + assert.equal(okAsset?.verification_details?.assetWorker?.lastResult, 'verified', 'verified asset should record worker result'); + assert.equal(okAsset?.security_flags?.assetRecheckFailed, undefined, 'verified asset should not keep recheck failure flag'); + + assert.equal(badAsset?.status, 'draft', 'mismatched asset should be unpublished'); + assert.equal(badAsset?.upload_status, 'failed', 'mismatched asset should be marked failed'); + assert.equal(badAsset?.security_flags?.assetRecheckFailed, true, 'mismatched asset should record security flag'); + assert.deepEqual( + badAsset?.verification_details?.assetWorker?.issues, + ['file_size_mismatch'], + 'mismatched asset should record exact issue', + ); + + const audits = await pool.query( + ` + select action, details + from public.audit_logs + where tenant_id = $1 + and target_type = 'content_asset' + and target_id in ($2, $3) + order by created_at asc + `, + [tenantId, ids.okAsset, ids.badAsset], + ); + assert.ok( + audits.rows.some(row => row.action === 'content.asset.rechecked' && row.details?.result === 'verified'), + 'worker should write verified audit log', + ); + assert.ok( + audits.rows.some(row => row.action === 'content.asset.recheck_failed' && row.details?.issues?.includes('file_size_mismatch')), + 'worker should write failed audit log', + ); + + console.log('Asset worker integration test complete.'); + } catch (error) { + await pool.query('rollback').catch(() => {}); + throw error; + } finally { + if (seeded) { + await cleanup(pool).catch(() => {}); + } + await pool.end(); + } +} + +main().catch(error => { + console.error(error); + process.exit(1); +});