From b55a225ccc66d1eb306ee6f8db2e810deef54603 Mon Sep 17 00:00:00 2001 From: kdletters Date: Thu, 27 Aug 2026 21:12:38 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=20AI=20=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E4=B8=8E=E5=8E=86=E5=8F=B2=E6=B8=85=E7=90=86=E7=9A=84=E5=86=85?= =?UTF-8?q?=E5=AD=98=E8=BE=B9=E7=95=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 补齐外部生成历史清理 procedure 的 SATS tuple 解析与回归测试 限制 AI 流式文本每阶段 chunk 数量并分批释放明细 修正进程内失败收口的容量检查顺序 分批删除外部生成审计事件并同步更新契约文档 --- .../shared-memory/decision-log.md | 4 +- ...】server-rs与SpacetimeDB数据契约-2026-05-15.md | 9 +- scripts/spacetime-migration-common.mjs | 18 ++ scripts/spacetime-migration-common.test.ts | 31 ++++ .../crates/module-ai/src/application/store.rs | 157 +++++++++++++++++- server-rs/crates/module-ai/src/domain.rs | 3 +- .../crates/module-ai/src/domain/limits.rs | 2 + server-rs/crates/module-ai/src/lib.rs | 8 +- .../crates/spacetime-module/src/ai/stages.rs | 42 +++-- .../src/external_generation.rs | 46 ++++- 10 files changed, 287 insertions(+), 33 deletions(-) diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 82659cf5b..c569689d5 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -35,7 +35,7 @@ ## 2026-08-27 外部生成历史采用受控保留清理 - 背景:`external_generation_job`、`external_generation_job_summary` 与 `external_generation_job_event` 都是持久化表;摘要和 payload 边界收紧后,已确认的终态历史仍会继续占用 SpacetimeDB 常驻内存,且事件审计链会随任务数量增长。 -- 决策:新增仅 migration operator 可调用的 `prune_external_generation_job_history_and_return`。默认按 `source_module=editor-canvas`、30 天保留期和 `job_id` 游标分批运行;只删除主任务与摘要状态一致、属于 completed / failed / cancelled、摘要已有 `notification_acknowledged_at` 且终态时间达到 cutoff 的任务,并在同一事务内删除该任务的全部事件、摘要和主任务。默认 dry-run,必须固定 dry-run 返回的 cutoff 后再 apply;pending / running、未确认通知、摘要缺失或状态不一致的数据永不删除。其他 source module 必须显式指定并单独评估;资产对象和钱包流水不随任务历史删除;不新增自动定时器或 runtime 清理权限。 +- 决策:新增仅 migration operator 可调用的 `prune_external_generation_job_history_and_return`。默认按 `source_module=editor-canvas`、30 天保留期和 `job_id` 游标分批运行;只删除主任务与摘要状态一致、属于 completed / failed / cancelled、摘要已有 `notification_acknowledged_at` 且终态时间达到 cutoff 的任务。事件、摘要和主任务仍按同一事务顺序删除,但每次事务最多删除 256 条事件;事件未删完时保留任务与摘要并返回同一个 job cursor,维护脚本下一次继续,避免单个任务形成无界事务写集。默认 dry-run,必须固定 dry-run 返回的 cutoff 后再 apply;pending / running、未确认通知、摘要缺失或状态不一致的数据永不删除。其他 source module 必须显式指定并单独评估;资产对象和钱包流水不随任务历史删除;不新增自动定时器或 runtime 清理权限。 - 影响范围:`server-rs/crates/spacetime-module/src/external_generation.rs`、外部生成事件 job_id 单列索引、SpacetimeDB 生成 bindings、`scripts/spacetime-maintain-external-generation-jobs.mjs`、架构与生产运维文档。 - 验证方式:覆盖终态 / 活跃态 / 已确认与未确认摘要、状态或身份不一致、cutoff 边界测试;运行 SpacetimeDB module tests/check、bindings 生成、schema/encoding/diff 门禁,并在维护窗口先 dry-run 再 apply。 - 关联文档:`docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md`、`docs/【开发运维】本地开发验证与生产运维-2026-05-15.md`、PR #203。 @@ -7798,6 +7798,6 @@ CI 上 `background_agent_runtime_recovers_stale_running_before_pending_task` 在 - 决策:外部生成 worker 每轮主动 `try_join_next` 回收已完成 `JoinHandle`,避免持续有队列任务时只归还 semaphore permit 却让 `JoinSet` 句柄集合无界增长;超时脱管任务在 abort 后等待句柄结束,执行许可保持到 work 真正结束或被取消。 - 决策:`module-ai` 的阶段终态写入与流式增量统一受文本、结构化 JSON、warning 和全局 retained 工作集上限约束;认证投影恢复对过滤后的 refresh session 重新计数,超过 8192 条直接拒绝启动恢复,避免超限快照灌入内存。 -- 复审补充:`spacetime-module` 的 AI procedure 复用同一组任务元数据 / payload / 输出 / 结果引用上限,流式文本聚合超限时回滚事务,terminal task 收口后删除 `ai_text_chunk` 明细,避免真实持久化链路绕过内存边界。 +- 复审补充:`spacetime-module` 的 AI procedure 复用同一组任务元数据 / payload / 输出 / 结果引用上限,流式文本聚合超过 512 KiB 或每阶段 8192 个 chunk 时回滚事务,terminal task 收口后按有界批次删除 `ai_text_chunk` 明细,避免真实持久化链路绕过内存边界。 - 决策:备份脚本停库前写入受保护 `.spacetimedb-stopped` marker,正常恢复完成后清理;systemd 备份 service 通过 `MemoryHigh/MemoryMax/OOMPolicy` 和 `ExecStopPost` 在 Node OOM kill、无法执行 JS finally 时兜底拉起 SpacetimeDB、API、worker、controller,恢复不完整则保留 marker。 - 验证:worker/module-auth/module-ai 定向 Rust 测试、database-backup/production-ops/encoding 门禁和 `git diff --check` 必须在提交前通过;release 现场需按 archive-full 重新 provision 并核验旧 files-history drop-in 已删除、四个服务 active。 diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index dee419cf0..672d6c999 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -329,8 +329,8 @@ Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复 - Rust 结构体:`AiTask` - 源码:`server-rs/crates/spacetime-module/src/ai/tasks.rs` -- `module-ai` 的进程内热状态不是持久化真相:文本增量按阶段有序聚合并受单阶段 512 KiB 上限约束;terminal task 立即释放增量明细,内存工作集最多保留 1024 个任务。需要长期查询时必须读取 SpacetimeDB 的 `ai_task` / `ai_task_stage` 投影,不得依赖进程重启后仍存在的内存快照。 -- SpacetimeDB 的 AI 写入 procedure 必须复用同一组任务元数据、payload、文本、结构化输出、warning、失败消息和结果引用上限;流式聚合超过 512 KiB 时在事务内拒绝,terminal task 收口后删除 `ai_text_chunk` 明细,只保留阶段最终快照和结果引用。 +- `module-ai` 的进程内热状态不是持久化真相:文本增量按阶段有序聚合并受单阶段 512 KiB、每阶段 8192 个 chunk 上限约束;terminal task 立即释放增量明细,内存工作集最多保留 1024 个任务。需要长期查询时必须读取 SpacetimeDB 的 `ai_task` / `ai_task_stage` 投影,不得依赖进程重启后仍存在的内存快照。 +- SpacetimeDB 的 AI 写入 procedure 必须复用同一组任务元数据、payload、文本、结构化输出、warning、失败消息和结果引用上限;流式聚合超过 512 KiB 或 8192 个 chunk 时在事务内拒绝,terminal task 收口后分批删除 `ai_text_chunk` 明细,只保留阶段最终快照和结果引用。 ### `ai_task_event` @@ -362,19 +362,20 @@ Responses 的终态载荷既是工具调用的恢复源,也是正文的恢复 - 源码:`server-rs/crates/spacetime-module/src/external_generation.rs` - 用途:外部生成正式任务列表的轻量投影,按 `job_id` 保存 owner、来源、状态、可选 `phase`、价格、有界错误摘要、通知确认时间、各阶段时间和入队时提取的 `request_prompt`,不包含 request/result payload、worker lease 或 dedupe 内部字段。错误摘要统一拒绝内联媒体并限制为 2048 字符;列表在单次 owner 扫描中同时计数并只保留请求 limit 的固定大小 top-N,不得先收集全量历史再截断。enqueue、claim、renew、phase update、complete、fail 事务同步投影;acknowledge 只更新该轻量表并写审计事件,后续主任务同步必须保留已有确认时间,禁止为了写确认时间加载 / 重写大 payload 行。BFF 的列表、状态和确认只调用 summary procedure;`running + processing` 映射为“正在处理”,其它 running(含旧行 `phase=None`)映射为“正在生成”。历史终态任务由迁移操作员的游标分批 maintenance procedure 在压缩 payload 时同步回填摘要,正式列表不得为兼容旧数据回扫完整主表。 - 非阻断告警:摘要字段 `warning_message` 是展示投影,由完成任务的轻量 `result_payload_json.warning.reason` 原样提取,不等同于公开 inline / external v1 的原始结构化诊断字段。complete 和历史 backfill 共用同一构建路径;历史任务按其结果载荷中已写入的 `reason` 快照投影,不为格式升级重写或补前缀。单 job 状态和任务列表 BFF 以 `warning: string` 返回该可直接展示的完整文案,不再返回结构化 code,Web 不得再次补前缀或按字符串推断告警类型。错误与告警摘要都不复制内联媒体并限制为 2048 字符。`phase` 与 `warning_message` 分别表示当前执行阶段和成功降级提示,不得混用;worker / BFF / Web 必须同版本协调发布,不保证滚动混部或旧 Web 缓存下的字符串语义兼容。 -- 正式读取 procedure 为 `get_external_generation_job_summary_and_return`、`list_external_generation_job_summaries_and_return` 和 `acknowledge_external_generation_job_summaries_and_return`。历史维护 procedure 为 `compact_external_generation_job_payloads_and_return`、`backfill_external_generation_job_summaries_and_return` 与 `prune_external_generation_job_history_and_return`,仅 migration operator 可调用;运维入口统一使用 `npm run spacetime:external-generation:maintain -- ...`,默认 dry-run、单批最多 25 条。B-tree cursor 选择阶段最多反序列化 `limit + 1` 行,apply 再按主键逐条读取选中行;怀疑存在单行异常巨型 JSON 时必须先使用 `--limit 1`。payload 压缩额外固定使用 `source_module = editor-canvas` 的复合 cursor 索引,不得静默改写其它玩法历史任务。历史清理默认使用 `--prune-history`、`source_module = editor-canvas` 和 30 天保留期;只有主任务与摘要状态一致且属于 completed / failed / cancelled、摘要已有 `notification_acknowledged_at`、终态时间不晚于 cutoff 的记录才是候选。apply 在同一事务内按事件 → 摘要 → 主任务顺序删除,事件不得独立清理;pending / running、未确认通知、摘要缺失或状态不一致的记录永不删除。清理不触碰资产对象或钱包流水,其他 source module 必须显式指定并单独评估。 +- 正式读取 procedure 为 `get_external_generation_job_summary_and_return`、`list_external_generation_job_summaries_and_return` 和 `acknowledge_external_generation_job_summaries_and_return`。历史维护 procedure 为 `compact_external_generation_job_payloads_and_return`、`backfill_external_generation_job_summaries_and_return` 与 `prune_external_generation_job_history_and_return`,仅 migration operator 可调用;运维入口统一使用 `npm run spacetime:external-generation:maintain -- ...`,默认 dry-run、单批最多 25 条。B-tree cursor 选择阶段最多反序列化 `limit + 1` 行,apply 再按主键逐条读取选中行;怀疑存在单行异常巨型 JSON 时必须先使用 `--limit 1`。payload 压缩额外固定使用 `source_module = editor-canvas` 的复合 cursor 索引,不得静默改写其它玩法历史任务。历史清理默认使用 `--prune-history`、`source_module = editor-canvas` 和 30 天保留期;只有主任务与摘要状态一致且属于 completed / failed / cancelled、摘要已有 `notification_acknowledged_at`、终态时间不晚于 cutoff 的记录才是候选。apply 在同一事务内按事件 → 摘要 → 主任务顺序删除,事件不得独立清理;每次事务最多删除 256 条事件,若同一任务仍有事件则保留任务与摘要并返回同一个 `next_cursor_job_id`,下一次继续该任务,避免单个任务形成无界事务写集;pending / running、未确认通知、摘要缺失或状态不一致的记录永不删除。清理不触碰资产对象或钱包流水,其他 source module 必须显式指定并单独评估。 ### `external_generation_job_event` - Rust 结构体:`ExternalGenerationJobEvent` - 源码:`server-rs/crates/spacetime-module/src/external_generation.rs` - 用途:外部生成任务审计事件表,按 `job_id` 和 `owner_user_id` 记录 `enqueued`、`claimed`、`lease_renewed`、`completed`、`failed`、`acknowledged` 等状态转换事实。状态转换只能由 SpacetimeDB procedure 写入,不由前端或 worker 直接改表;该表用于追溯任务生命周期和排障,不替代 `external_generation_job` 当前状态。 -- 保留策略:事件只会随已确认通知的终态任务由 `prune_external_generation_job_history_and_return` 原子删除,不支持按事件单独清理,以保持任务、摘要和审计链一致。 +- 保留策略:事件只会随已确认通知的终态任务由 `prune_external_generation_job_history_and_return` 原子删除,不支持按事件单独清理,以保持任务、摘要和审计链一致。单次事务最多删除 256 条事件;若事件未删完,任务和摘要暂不删除,维护脚本用同一个 job cursor 重试剩余事件。 ### `ai_text_chunk` - Rust 结构体:`AiTextChunk` - 源码:`server-rs/crates/spacetime-module/src/ai/stages.rs` +- 单阶段最多保留 8192 个 chunk;聚合和终态清理均按有界批次处理,避免小 delta 堆积为无界行数或一次性 ID 列表。 ### `analytics_date_dimension` diff --git a/scripts/spacetime-migration-common.mjs b/scripts/spacetime-migration-common.mjs index 6cc7cbbe4..afe52a598 100644 --- a/scripts/spacetime-migration-common.mjs +++ b/scripts/spacetime-migration-common.mjs @@ -427,6 +427,24 @@ function normalizeSatsProduct(value, procedureName) { }; } + if ( + procedureName === 'prune_external_generation_job_history_and_return' && + value.length === 10 + ) { + return { + ok: normalizeSatsValue(value[0]), + dry_run: normalizeSatsValue(value[1]), + scanned_count: normalizeSatsValue(value[2]), + selected_count: normalizeSatsValue(value[3]), + deleted_job_count: normalizeSatsValue(value[4]), + deleted_summary_count: normalizeSatsValue(value[5]), + deleted_event_count: normalizeSatsValue(value[6]), + next_cursor_job_id: normalizeSatsOption(value[7]), + has_more: normalizeSatsValue(value[8]), + error_message: normalizeSatsOption(value[9]), + }; + } + if (value.length === 3) { return { ok: normalizeSatsValue(value[0]), diff --git a/scripts/spacetime-migration-common.test.ts b/scripts/spacetime-migration-common.test.ts index 6317bd5d9..689abe321 100644 --- a/scripts/spacetime-migration-common.test.ts +++ b/scripts/spacetime-migration-common.test.ts @@ -223,4 +223,35 @@ describe('SpacetimeDB CLI SATS option encoding', () => { expect(objectResult.batch_sha256).toBe('d'.repeat(64)); expect(objectResult).not.toHaveProperty('batch_sha_256'); }); + + it('normalizes external generation history prune tuple results', () => { + const result = parseProcedureResult( + JSON.stringify([ + true, + false, + 25, + 2, + 2, + 2, + 10, + [0, 'job-25'], + true, + [0, '清理失败'], + ]), + 'prune_external_generation_job_history_and_return', + ); + + expect(result).toEqual({ + ok: true, + dry_run: false, + scanned_count: 25, + selected_count: 2, + deleted_job_count: 2, + deleted_summary_count: 2, + deleted_event_count: 10, + next_cursor_job_id: 'job-25', + has_more: true, + error_message: '清理失败', + }); + }); }); diff --git a/server-rs/crates/module-ai/src/application/store.rs b/server-rs/crates/module-ai/src/application/store.rs index 6150dc397..35a9725a9 100644 --- a/server-rs/crates/module-ai/src/application/store.rs +++ b/server-rs/crates/module-ai/src/application/store.rs @@ -114,18 +114,25 @@ impl InMemoryAiTaskStore { .insert(task_id.trim().to_string(), previous_task); return Err(error); } + let released_text_chunks = if snapshot.status.is_terminal() { + state.text_chunks.remove(task_id.trim()) + } else { + None + }; let retained_output_bytes = retained_output_bytes(&state); if retained_output_bytes > MAX_AI_TASK_RETAINED_OUTPUT_BYTES { state .tasks .insert(task_id.trim().to_string(), previous_task); + if let Some(text_chunks) = released_text_chunks { + state + .text_chunks + .insert(task_id.trim().to_string(), text_chunks); + } return Err(AiTaskServiceError::Store( "AI 任务仓储输出工作集超过内存上限".to_string(), )); } - if snapshot.status.is_terminal() { - state.text_chunks.remove(task_id.trim()); - } Ok(snapshot) } @@ -166,6 +173,13 @@ impl InMemoryAiTaskStore { .get_mut(&chunk.task_id) .ok_or(AiTaskServiceError::TaskNotFound)?; let stage_chunks = chunks.entry(chunk.stage_kind).or_default(); + if !stage_chunks.contains_key(&chunk.sequence) + && stage_chunks.len() >= crate::MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE + { + return Err(AiTaskServiceError::Store( + "AI 任务文本 chunk 数量超过内存上限".to_string(), + )); + } let previous_chunk = stage_chunks.insert(chunk.sequence, chunk.delta_text.clone()); let aggregated_bytes = stage_chunks .values() @@ -352,3 +366,140 @@ fn validate_task_memory_limits(task: &AiTaskSnapshot) -> Result<(), AiTaskServic validate_ai_task_snapshot_memory_limits(task) .map_err(|message| AiTaskServiceError::Store(message.to_string())) } + +#[cfg(test)] +mod tests { + use super::*; + + fn build_running_task(task_id: &str) -> AiTaskSnapshot { + AiTaskSnapshot { + task_id: task_id.to_string(), + task_kind: crate::AiTaskKind::CharacterChat, + owner_user_id: "user-1".to_string(), + request_label: "测试任务".to_string(), + source_module: "test".to_string(), + source_entity_id: None, + request_payload_json: None, + status: AiTaskStatus::Running, + failure_message: None, + stages: vec![crate::AiTaskStageSnapshot { + stage_kind: crate::AiTaskStageKind::RequestModel, + label: "请求模型".to_string(), + detail: "测试".to_string(), + order: 0, + status: AiTaskStageStatus::Running, + text_output: None, + structured_payload_json: None, + warning_messages: Vec::new(), + started_at_micros: Some(1), + completed_at_micros: None, + }], + result_references: Vec::new(), + latest_text_output: None, + latest_structured_payload_json: None, + version: 1, + created_at_micros: 1, + started_at_micros: Some(1), + completed_at_micros: None, + updated_at_micros: 1, + } + } + + #[test] + fn append_text_chunk_rejects_excessive_chunk_count() { + let store = InMemoryAiTaskStore::default(); + let task_id = "task-chunk-limit"; + let task = build_running_task(task_id); + let mut state = store.inner.lock().expect("store lock should be available"); + state.tasks.insert(task_id.to_string(), task); + state.text_chunks.insert( + task_id.to_string(), + HashMap::from([( + crate::AiTaskStageKind::RequestModel, + (1..=crate::MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE as u32) + .map(|sequence| (sequence, "a".to_string())) + .collect(), + )]), + ); + drop(state); + + let error = store + .append_text_chunk(AiTextChunkSnapshot { + chunk_id: "chunk-overflow".to_string(), + task_id: task_id.to_string(), + stage_kind: crate::AiTaskStageKind::RequestModel, + sequence: crate::MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE as u32 + 1, + delta_text: "b".to_string(), + created_at_micros: 2, + }) + .expect_err("a new chunk beyond the count cap should fail"); + assert!(matches!( + error, + AiTaskServiceError::Store(message) if message.contains("chunk 数量") + )); + } + + #[test] + fn terminal_failure_releases_chunks_before_global_cap_check() { + let store = InMemoryAiTaskStore::default(); + let target_id = "task-terminal-release"; + let target = build_running_task(target_id); + let mut state = store.inner.lock().expect("store lock should be available"); + state.tasks.insert(target_id.to_string(), target); + state.text_chunks.insert( + target_id.to_string(), + HashMap::from([( + crate::AiTaskStageKind::RequestModel, + BTreeMap::from([(1, "t".repeat(crate::MAX_AI_TASK_TEXT_OUTPUT_BYTES))]), + )]), + ); + + let desired_retained = crate::MAX_AI_TASK_RETAINED_OUTPUT_BYTES + .saturating_sub(crate::MAX_AI_TASK_FAILURE_MESSAGE_BYTES) + .saturating_add(crate::MAX_AI_TASK_FAILURE_MESSAGE_BYTES / 2) + .saturating_sub(4 * 1024); + let mut filler_index = 0_u32; + while retained_output_bytes(&state) < desired_retained { + let remaining = desired_retained.saturating_sub(retained_output_bytes(&state)); + let bytes = remaining.min(crate::MAX_AI_TASK_TEXT_OUTPUT_BYTES); + if bytes == 0 { + break; + } + let filler_id = format!("task-filler-{filler_index}"); + filler_index += 1; + state + .tasks + .insert(filler_id.clone(), build_running_task(&filler_id)); + state.text_chunks.insert( + filler_id, + HashMap::from([( + crate::AiTaskStageKind::RequestModel, + BTreeMap::from([(1, "f".repeat(bytes))]), + )]), + ); + } + let retained_before_failure = retained_output_bytes(&state); + assert!(retained_before_failure <= crate::MAX_AI_TASK_RETAINED_OUTPUT_BYTES); + assert!( + retained_before_failure + crate::MAX_AI_TASK_FAILURE_MESSAGE_BYTES + > crate::MAX_AI_TASK_RETAINED_OUTPUT_BYTES + ); + drop(state); + + let failed = store + .update_task(target_id, |task| { + task.status = AiTaskStatus::Failed; + task.failure_message = Some("f".repeat(crate::MAX_AI_TASK_FAILURE_MESSAGE_BYTES)); + task.completed_at_micros = Some(2); + task.updated_at_micros = 2; + task.version += 1; + Ok(()) + }) + .expect("terminal transition should release chunks before checking the cap"); + assert_eq!(failed.status, AiTaskStatus::Failed); + assert_eq!( + failed.failure_message.as_deref().map(str::len), + Some(crate::MAX_AI_TASK_FAILURE_MESSAGE_BYTES) + ); + } +} diff --git a/server-rs/crates/module-ai/src/domain.rs b/server-rs/crates/module-ai/src/domain.rs index df225816e..abafd7616 100644 --- a/server-rs/crates/module-ai/src/domain.rs +++ b/server-rs/crates/module-ai/src/domain.rs @@ -16,7 +16,8 @@ pub use limits::{ MAX_AI_TASK_RESULT_REFERENCES, MAX_AI_TASK_RETAINED_OUTPUT_BYTES, MAX_AI_TASK_RETAINED_TASKS, MAX_AI_TASK_SOURCE_ENTITY_ID_BYTES, MAX_AI_TASK_SOURCE_MODULE_BYTES, MAX_AI_TASK_STAGE_DETAIL_BYTES, MAX_AI_TASK_STAGE_LABEL_BYTES, - MAX_AI_TASK_STRUCTURED_OUTPUT_BYTES, MAX_AI_TASK_TEXT_OUTPUT_BYTES, MAX_AI_TASK_WARNING_BYTES, + MAX_AI_TASK_STRUCTURED_OUTPUT_BYTES, MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE, + MAX_AI_TASK_TEXT_OUTPUT_BYTES, MAX_AI_TASK_WARNING_BYTES, validate_ai_task_snapshot_memory_limits, }; pub use types::{ diff --git a/server-rs/crates/module-ai/src/domain/limits.rs b/server-rs/crates/module-ai/src/domain/limits.rs index 759b4ac6a..16d3a46c2 100644 --- a/server-rs/crates/module-ai/src/domain/limits.rs +++ b/server-rs/crates/module-ai/src/domain/limits.rs @@ -9,6 +9,8 @@ pub const MAX_AI_TASK_SOURCE_ENTITY_ID_BYTES: usize = 512; pub const MAX_AI_TASK_STAGE_LABEL_BYTES: usize = 4 * 1024; pub const MAX_AI_TASK_STAGE_DETAIL_BYTES: usize = 8 * 1024; pub const MAX_AI_TASK_TEXT_OUTPUT_BYTES: usize = 512 * 1024; +// provider 产生大量细小流式增量时,限制行和索引开销。 +pub const MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE: usize = 8 * 1024; pub const MAX_AI_TASK_STRUCTURED_OUTPUT_BYTES: usize = 512 * 1024; pub const MAX_AI_TASK_WARNING_BYTES: usize = 64 * 1024; pub const MAX_AI_TASK_REQUEST_PAYLOAD_BYTES: usize = 512 * 1024; diff --git a/server-rs/crates/module-ai/src/lib.rs b/server-rs/crates/module-ai/src/lib.rs index 2eba88eb0..a3311e4d9 100644 --- a/server-rs/crates/module-ai/src/lib.rs +++ b/server-rs/crates/module-ai/src/lib.rs @@ -21,10 +21,10 @@ pub use domain::{ MAX_AI_TASK_RETAINED_OUTPUT_BYTES, MAX_AI_TASK_RETAINED_TASKS, MAX_AI_TASK_SOURCE_ENTITY_ID_BYTES, MAX_AI_TASK_SOURCE_MODULE_BYTES, MAX_AI_TASK_STAGE_DETAIL_BYTES, MAX_AI_TASK_STAGE_LABEL_BYTES, - MAX_AI_TASK_STRUCTURED_OUTPUT_BYTES, MAX_AI_TASK_TEXT_OUTPUT_BYTES, MAX_AI_TASK_WARNING_BYTES, - generate_ai_result_ref_id, generate_ai_task_id, generate_ai_task_stage_id, - generate_ai_text_chunk_id, normalize_optional_text, normalize_string_list, - validate_ai_task_snapshot_memory_limits, + MAX_AI_TASK_STRUCTURED_OUTPUT_BYTES, MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE, + MAX_AI_TASK_TEXT_OUTPUT_BYTES, MAX_AI_TASK_WARNING_BYTES, generate_ai_result_ref_id, + generate_ai_task_id, generate_ai_task_stage_id, generate_ai_text_chunk_id, + normalize_optional_text, normalize_string_list, validate_ai_task_snapshot_memory_limits, }; pub use errors::{AiTaskFieldError, AiTaskServiceError}; pub use events::AiTaskDomainEvent; diff --git a/server-rs/crates/spacetime-module/src/ai/stages.rs b/server-rs/crates/spacetime-module/src/ai/stages.rs index 1fa2cab3f..86bc26bc9 100644 --- a/server-rs/crates/spacetime-module/src/ai/stages.rs +++ b/server-rs/crates/spacetime-module/src/ai/stages.rs @@ -1,9 +1,12 @@ use crate::*; use module_ai::{ - MAX_AI_TASK_TEXT_OUTPUT_BYTES, generate_ai_result_ref_id, generate_ai_text_chunk_id, - normalize_optional_text, normalize_string_list, validate_ai_task_snapshot_memory_limits, + MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE, MAX_AI_TASK_TEXT_OUTPUT_BYTES, generate_ai_result_ref_id, + generate_ai_text_chunk_id, normalize_optional_text, normalize_string_list, + validate_ai_task_snapshot_memory_limits, }; +const AI_TEXT_CHUNK_DELETE_BATCH_SIZE: usize = 256; + #[spacetimedb::table( accessor = ai_task_stage, index(accessor = by_ai_task_stage_task_id, btree(columns = [task_id])), @@ -178,7 +181,7 @@ pub(crate) fn append_ai_text_chunk_tx( if input.sequence == 0 { return Err("ai_text_chunk.sequence 必须大于 0".to_string()); } - if input.delta_text.trim().len() > MAX_AI_TASK_TEXT_OUTPUT_BYTES { + if input.delta_text.len() > MAX_AI_TASK_TEXT_OUTPUT_BYTES { return Err("AI 任务文本输出超过内存上限".to_string()); } @@ -340,15 +343,25 @@ pub(crate) fn replace_ai_task_stages( } pub(crate) fn delete_ai_text_chunks_for_task(ctx: &ReducerContext, task_id: &str) { - let chunk_row_ids = ctx - .db - .ai_text_chunk() - .by_ai_text_chunk_task_id() - .filter(task_id) - .map(|row| row.text_chunk_row_id.clone()) - .collect::>(); - for row_id in chunk_row_ids { - ctx.db.ai_text_chunk().text_chunk_row_id().delete(&row_id); + loop { + let chunk_row_ids = ctx + .db + .ai_text_chunk() + .by_ai_text_chunk_task_id() + .filter(task_id) + .take(AI_TEXT_CHUNK_DELETE_BATCH_SIZE) + .map(|row| row.text_chunk_row_id.clone()) + .collect::>(); + if chunk_row_ids.is_empty() { + break; + } + let batch_len = chunk_row_ids.len(); + for row_id in chunk_row_ids { + ctx.db.ai_text_chunk().text_chunk_row_id().delete(&row_id); + } + if batch_len < AI_TEXT_CHUNK_DELETE_BATCH_SIZE { + break; + } } } @@ -358,6 +371,7 @@ pub(crate) fn collect_ai_stage_text_output( stage_kind: AiTaskStageKind, ) -> Result, String> { let mut chunks = Vec::new(); + let mut chunk_count = 0_usize; let mut aggregated_bytes = 0_usize; for row in ctx .db @@ -366,6 +380,10 @@ pub(crate) fn collect_ai_stage_text_output( .filter(task_id) .filter(|row| row.task_id == task_id && row.stage_kind == stage_kind) { + chunk_count = chunk_count.saturating_add(1); + if chunk_count > MAX_AI_TASK_TEXT_CHUNKS_PER_STAGE { + return Err("AI 任务文本 chunk 数量超过内存上限".to_string()); + } aggregated_bytes = aggregated_bytes.saturating_add(row.delta_text.len()); if aggregated_bytes > MAX_AI_TASK_TEXT_OUTPUT_BYTES { return Err("AI 任务文本输出超过内存上限".to_string()); diff --git a/server-rs/crates/spacetime-module/src/external_generation.rs b/server-rs/crates/spacetime-module/src/external_generation.rs index 2978f3cef..5e93b162c 100644 --- a/server-rs/crates/spacetime-module/src/external_generation.rs +++ b/server-rs/crates/spacetime-module/src/external_generation.rs @@ -23,6 +23,7 @@ const MAX_EXTERNAL_GENERATION_REQUEST_PROMPT_CHARS: usize = 2_048; const MAX_EXTERNAL_GENERATION_ERROR_MESSAGE_CHARS: usize = 2_048; const MAX_EXTERNAL_GENERATION_WARNING_MESSAGE_CHARS: usize = 2_048; const MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE: u32 = 25; +const EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE: usize = 256; const INLINE_MEDIA_REMOVED_PLACEHOLDER: &str = "[inline-media-removed]"; const INLINE_MEDIA_ERROR_REDACTED_MESSAGE: &str = "外部生成失败(错误详情含内联媒体引用,已省略)"; const INLINE_MEDIA_WARNING_REDACTED_MESSAGE: &str = @@ -1433,6 +1434,15 @@ fn prune_external_generation_job_history_tx( .clamp(1, MAX_EXTERNAL_GENERATION_MAINTENANCE_BATCH_SIZE) as usize; let cursor_range = external_generation_job_maintenance_cursor_range(cursor_job_id.as_deref()); let cursor_to_skip = cursor_job_id.clone(); + // 若 cursor 对应的任务仍存在,说明上一事务只删完了事件的一部分,或 dry-run + // 尚未执行 apply;下一次必须包含该任务继续清理,成功删除后它会自然消失。 + let include_existing_cursor = cursor_job_id.as_deref().is_some_and(|cursor| { + ctx.db + .external_generation_job() + .job_id() + .find(&cursor.to_string()) + .is_some() + }); let rows = ctx .db .external_generation_job() @@ -1441,7 +1451,7 @@ fn prune_external_generation_job_history_tx( .filter(move |row| { cursor_to_skip .as_deref() - .is_none_or(|cursor| row.job_id != cursor) + .is_none_or(|cursor| row.job_id != cursor || include_existing_cursor) }); let (job_ids, next_cursor_job_id, has_more, scanned_count) = select_external_generation_job_ids_for_maintenance(rows, limit, |row| { @@ -1462,6 +1472,7 @@ fn prune_external_generation_job_history_tx( let mut deleted_job_count = 0u32; let mut deleted_summary_count = 0u32; let mut deleted_event_count = 0u32; + let mut pending_event_cursor_job_id = None; if !input.dry_run { for job_id in &job_ids { let Some(row) = ctx.db.external_generation_job().job_id().find(job_id) else { @@ -1484,8 +1495,13 @@ fn prune_external_generation_job_history_tx( continue; } - deleted_event_count = deleted_event_count - .saturating_add(delete_external_generation_job_events_for_job(ctx, job_id)); + let (deleted_for_job, has_more_events) = + delete_external_generation_job_events_for_job(ctx, job_id); + deleted_event_count = deleted_event_count.saturating_add(deleted_for_job); + if has_more_events { + pending_event_cursor_job_id = Some(job_id.clone()); + break; + } ctx.db .external_generation_job_summary() .job_id() @@ -1496,6 +1512,10 @@ fn prune_external_generation_job_history_tx( } } + let (next_cursor_job_id, has_more) = pending_event_cursor_job_id + .map(|job_id| (Some(job_id), true)) + .unwrap_or((next_cursor_job_id, has_more)); + Ok(ExternalGenerationJobRetentionProcedureResult { ok: true, dry_run: input.dry_run, @@ -1996,22 +2016,34 @@ fn select_external_generation_job_ids_for_maintenance( ) } -fn delete_external_generation_job_events_for_job(ctx: &ReducerContext, job_id: &str) -> u32 { +fn delete_external_generation_job_events_for_job( + ctx: &ReducerContext, + job_id: &str, +) -> (u32, bool) { + // 每次 procedure 最多删除一个固定批次;若仍有事件,保留 job/summary,调用方 + // 通过同一个 job cursor 重试,避免单个任务把整段审计历史塞进一个事务写集。 let event_ids = ctx .db .external_generation_job_event() .by_external_generation_job_event_job_id_only() .filter(job_id) + .take(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE + 1) .map(|event| event.event_id.clone()) .collect::>(); - let deleted_count = event_ids.len() as u32; - for event_id in event_ids { + let has_more = event_ids.len() > EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE; + let deleted_count = event_ids + .len() + .min(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE) as u32; + for event_id in event_ids + .into_iter() + .take(EXTERNAL_GENERATION_EVENT_DELETE_BATCH_SIZE) + { ctx.db .external_generation_job_event() .event_id() .delete(&event_id); } - deleted_count + (deleted_count, has_more) } fn count_external_generation_job_summaries_for_owner( -- 2.52.0