补齐 AI 任务与历史清理的内存边界 #206

Merged
kdletters merged 1 commits from codex/merge-memory-boundaries-20260827 into master 2026-08-27 22:18:17 +08:00
10 changed files with 287 additions and 33 deletions
@@ -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 后再 applypending / 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 后再 applypending / 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。
@@ -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`
+18
View File
@@ -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]),
@@ -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: '清理失败',
});
});
});
@@ -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)
);
}
}
+2 -1
View File
@@ -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::{
@@ -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;
+4 -4
View File
@@ -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;
@@ -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::<Vec<_>>();
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::<Vec<_>>();
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<Option<String>, 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());
@@ -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::<Vec<_>>();
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(