From f1ca3e79c150627b36dcc78458ee464b2639220e Mon Sep 17 00:00:00 2001 From: kdletters Date: Wed, 8 Jul 2026 22:16:41 +0800 Subject: [PATCH] =?UTF-8?q?=E8=BD=BB=E9=87=8F=E5=8C=96=E5=A4=96=E9=83=A8?= =?UTF-8?q?=E7=94=9F=E6=88=90=E9=98=9F=E5=88=97=E8=BF=9B=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit worker 和 controller 改为订阅 external_generation_job 队列变更来唤醒 claim 与扩缩容评估 非 HTTP 角色关闭 API 读模型订阅并将 SpacetimeDB 连接池收敛为 1 更新 worker/controller 环境模板和运维记忆,明确 poll interval 只作兜底 --- ...external-generation-controller.env.example | 2 + .../external-generation-worker.env.example | 2 + .../shared-memory/decision-log.md | 1 + ...发运维】本地开发验证与生产运维-2026-05-15.md | 2 + .../src/external_generation_worker.rs | 61 ++++++- .../external_generation_worker_controller.rs | 48 ++++- server-rs/crates/api-server/src/state.rs | 44 +++-- .../crates/api-server/src/tracking_outbox.rs | 1 + .../api-server/src/wallet_refund_outbox.rs | 1 + .../src/external_generation.rs | 164 ++++++++++++++++++ server-rs/crates/spacetime-client/src/lib.rs | 14 +- 11 files changed, 316 insertions(+), 24 deletions(-) diff --git a/deploy/env/external-generation-controller.env.example b/deploy/env/external-generation-controller.env.example index e6d0cecaa..9b922a941 100644 --- a/deploy/env/external-generation-controller.env.example +++ b/deploy/env/external-generation-controller.env.example @@ -1,7 +1,9 @@ # 复制到 /etc/genarrative/external-generation-controller.env 后按机器容量调整。 # controller 只管理 systemd worker 实例;SpacetimeDB、外部 provider 密钥继续复用 api-server.env。 # systemd unit 会强制设置 GENARRATIVE_PROCESS_ROLE=external-generation-controller。 +# 非 HTTP 角色只保留队列窄订阅唤醒,不需要 API 读模型连接池。 +GENARRATIVE_SPACETIME_POOL_SIZE=1 GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_MIN_WORKERS=1 GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_MAX_WORKERS=8 GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_TARGET_JOBS_PER_WORKER=2 diff --git a/deploy/env/external-generation-worker.env.example b/deploy/env/external-generation-worker.env.example index 3ddd83726..04a85f13b 100644 --- a/deploy/env/external-generation-worker.env.example +++ b/deploy/env/external-generation-worker.env.example @@ -2,7 +2,9 @@ # 该文件只覆盖 worker 专属参数;SpacetimeDB、外部 provider 密钥继续复用 api-server.env。 # systemd 模板会强制设置 GENARRATIVE_PROCESS_ROLE=external-generation-worker # 和 GENARRATIVE_EXTERNAL_GENERATION_WORKER_ID=%H-%i,避免多实例 ID 冲突。 +# 非 HTTP 角色只保留队列窄订阅唤醒,不需要 API 读模型连接池。 +GENARRATIVE_SPACETIME_POOL_SIZE=1 GENARRATIVE_EXTERNAL_GENERATION_WORKER_CONCURRENCY=2 GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS=2000 # 单次 lease 会由 worker 自动续租;该值覆盖心跳抖动窗口即可。 diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 9e3f9c608..244dfaa0f 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -1214,6 +1214,7 @@ - 决策:外部生成任务统一进入 SpacetimeDB `external_generation_job` 持久队列,由 `api-server` 的 `external-generation-worker` 进程角色 claim lease 后执行;HTTP 角色只做鉴权、表单/状态初始化、入队和返回 `queued/running/completed/failed` 操作状态。生产通过 systemd worker 模板增加实例数或提高 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_CONCURRENCY` 动态扩缩容,`GENARRATIVE_PROCESS_ROLE=all` 仅用于本地 smoke。拼图 `compile_puzzle_draft`、结果页 `generate_puzzle_images` 与 `generate_puzzle_ui_background` 已接入 worker;业务写回必须在 SpacetimeDB transaction 内校验 `external_generation_job` 的 `job_id + worker_id + lease_token`、job kind、owner 和 source entity,其中首图 worker 的前置 `compile_puzzle_agent_draft` 也必须带 guard。worker 核心业务写回失败不能返回内存快照并把 job 标成 completed;失败态业务写回成功后才能把 job 标成 failed,失败态未写回则保留租约等待后续重领。拼图业务失败不自动重试,只保留 lease 过期后的崩溃重领,避免钱包扣退费幂等漂移。生产发布会启用默认 `genarrative-external-generation-worker@1.service` 并等待 worker active,worker 停机时停止 claim 新任务并 drain 当前任务。 - 2026-06-07 追加:`GENARRATIVE_EXTERNAL_GENERATION_MODE` 使用 `queue|inline` 显式策略;生产和容器扩缩容验证保持 `queue`。本地开发若需要同步等待结果,应通过 `.env.local` 或本机环境显式配置为 `inline`,由 HTTP handler 复用同一 worker executor 直接返回 `completed`,不创建 `external_generation_job`,不支持 worker 动态扩缩容;脚本不得硬编码该策略。拼图写回 guard 字段改为可选,queue 路径仍必须完整校验 `job_id + worker_id + lease_token`;inline 路径只允许三项同时为空,半空 guard 仍拒绝。 - 2026-06-11 追加:生产新增固定 `external-generation-controller` 进程角色和 `genarrative-external-generation-controller.service`。controller 只读取 `get_external_generation_queue_stats_and_return` 队列统计并管理 `genarrative-external-generation-worker@N.service`,不监听 HTTP、不执行外部生成任务;默认保留 `@1`,按 `claimable_pending + running_active + expired_running` 计算目标实例数,上限由 `GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_MAX_WORKERS` 控制,缩容需要连续空闲轮数且每轮只停最高编号一个实例。 +- 2026-07-08 追加:生产 worker/controller 作为轻量 SpacetimeDB 客户端运行,专属 env 示例默认 `GENARRATIVE_SPACETIME_POOL_SIZE=1`;非 HTTP 角色只保留 `external_generation_job` 队列窄订阅作为响应式唤醒信号,实际抢占和扩缩容判断仍走 SpacetimeDB procedure,且不再订阅 API 读模型连接池。worker/controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。 - 影响范围:`server-rs/crates/spacetime-module/src/external_generation.rs`、`server-rs/crates/spacetime-client/src/external_generation.rs`、`server-rs/crates/api-server/src/external_generation_worker.rs`、`server-rs/crates/api-server/src/external_generation_worker_controller.rs`、`deploy/systemd/genarrative-external-generation-worker@.service`、`deploy/systemd/genarrative-external-generation-controller.service`、`deploy/env/external-generation-controller.env.example`、`scripts/deploy/production-api-deploy.sh`、`scripts/jenkins-server-provision.sh`、拼图 `compile_puzzle_draft`、拼图 `generate_puzzle_images`、拼图 `generate_puzzle_ui_background`、生产 env 模板和运维文档。 - 验证方式:`npm run spacetime:generate`、`npm run check:spacetime-schema`、`npm run check:server-rs-ddd`、`cargo check -p api-server --manifest-path server-rs/Cargo.toml`,并在 queue 模式下用 `GENARRATIVE_PROCESS_ROLE=all npm run dev` smoke 至少一次 queued -> worker 完成链路;本地 inline 排查只确认不创建 `external_generation_job`。 - 关联文档:`docs/technical/【后端架构】外部生成Worker化方案-2026-06-03.md`、`docs/【开发运维】本地开发验证与生产运维-2026-05-15.md`、`docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md`。 diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md index 535477266..1efc27ee2 100644 --- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md +++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md @@ -57,6 +57,8 @@ Windows 本地如果已在 `%LOCALAPPDATA%\Genarrative\ffmpeg\bin` 安装 FFmpeg 本地排查外部内容生成 worker 队列时,默认同一 Rust 进程同时监听 HTTP 并消费 `external_generation_job` 队列;更接近生产的验证应分别启动 `api`、`external-generation-worker` 和 `external-generation-controller`。生产默认 `GENARRATIVE_PROCESS_ROLE=api`,外部生成任务由独立 `GENARRATIVE_PROCESS_ROLE=external-generation-worker` 进程消费;生产与容器扩缩容验证保持 `queue`。当前进入持久队列的外部生成动作包括:拼图 `compile_puzzle_draft` / `generate_puzzle_images` / `generate_puzzle_ui_background`,跳一跳 `compile-draft` / `regenerate-tiles`,拼消消 `compile-draft` / `regenerate-atlas`,敲木鱼 `compile-draft` / `regenerate-hit-object`,以及图片画布 `editor_image_generation` / `editor_image_edit` / `editor_background_removal` / `editor_icon_spritesheet_generation` / `editor_ui_design_asset_extraction` / `editor_character_animation_generation` / `editor_video_generation` / `editor_sound_effect_generation` / `editor_background_music_generation`。非外部 provider 生成动作继续 inline,不进入队列。显式把本地进程角色设为 `api` 且没有 worker 时,HTTP 只返回 queued/running,不会兜底执行外部 provider。 +生产拆分角色时,`external-generation-worker` 和 `external-generation-controller` 的专属 env 示例会把 `GENARRATIVE_SPACETIME_POOL_SIZE` 覆盖为 `1`;非 HTTP 角色只保留 `external_generation_job` 队列窄订阅作为响应式唤醒信号,实际抢占和扩缩容判断仍走 SpacetimeDB procedure,且不再订阅 API 读模型连接池。`GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS` 与 controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。 + `我的` 页签或排障面板展示队列等待时,只读取 BFF 队列接口:`GET /api/runtime/external-generation/queue-overview` 查看当前用户可见队列概览,`GET /api/runtime/external-generation/jobs/{jobId}` 查看单 job 状态。生成页 / 进度页不承接队列概览,只展示当前玩法业务进度;队列接口只提供等待 / 运行 / 失败 / 完成状态补充,最终草稿、作品和结果页仍要轮询对应玩法 session/detail 接口收敛到 ready 或 failed;不要直接查询 `external_generation_job` private table,也不要把 worker 内部 payload 暴露到前端。 需要验证“更新 API 不停 worker”和“worker 是否持续消费队列”时,优先使用隔离容器 smoke:`npm run container:worker-smoke -- smoke`。该脚本生成 gitignored 的 `deploy/container/worker-smoke/api-server.env`,启动独立 compose project 与独立 SpacetimeDB,发布当前 `spacetime-module` 后写入 `worker_smoke_unsupported` 测试 job;预期 worker claim 后执行 unsupported 失败分支,再执行 API-only recreate 并确认 worker 容器 ID 不变,最后再次入队验证 API 更新后队列仍可消费。`external_generation_job` 是 private table,脚本通过 worker 日志确认 job_id 被消费,不用 CLI SQL 查询私表。该 smoke 不读取 `.env.local`,也不依赖真实 VectorEngine / OSS 密钥;真实生图链路联调再在本地私有 env 中补齐 provider 配置。worker-smoke 默认把本机 `spacetime` CLI 打成轻量 SpacetimeDB 镜像,避免本机首次 smoke 依赖官方大镜像下载。若容器内 Cargo 拉取 crates.io 依赖不稳定,可用 `npm run container:worker-smoke -- smoke --local-binary` 让容器内 Cargo 复用本机 Cargo 缓存构建当前二进制,再打入 Debian bookworm smoke runtime 临时镜像;可用 `GENARRATIVE_WORKER_SMOKE_LOCAL_BASE_IMAGE` 覆盖运行时基础镜像;若隔离端口或库数据需要重建,追加 `--force`。完成 queue 链路验证时,还要用队列概览 BFF 和单 job 状态接口确认 job 从 queued/running 收敛,并用对应玩法 session/detail 接口确认业务状态同步完成。 diff --git a/server-rs/crates/api-server/src/external_generation_worker.rs b/server-rs/crates/api-server/src/external_generation_worker.rs index e775e637d..491da8e78 100644 --- a/server-rs/crates/api-server/src/external_generation_worker.rs +++ b/server-rs/crates/api-server/src/external_generation_worker.rs @@ -6,7 +6,7 @@ use shared_kernel::offset_datetime_to_unix_micros; use spacetime_client::{ ExternalGenerationJobClaimRecordInput, ExternalGenerationJobCompleteRecordInput, ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord, - ExternalGenerationJobRenewLeaseRecordInput, + ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription, }; use tokio::{ task::JoinSet, @@ -70,6 +70,7 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<() let lease = state.config.external_generation_worker_lease; let mut tasks = JoinSet::new(); let mut shutdown = external_generation_worker_shutdown_signal(); + let mut queue_wake = None; info!( worker_id, @@ -80,6 +81,8 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<() ); loop { + ensure_external_generation_queue_wake_subscription(&state, &mut queue_wake).await; + while tasks.len() >= concurrency { if await_worker_task_or_shutdown(&mut tasks, &mut shutdown).await { drain_external_generation_worker_tasks(&mut tasks).await; @@ -110,8 +113,9 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<() Ok(jobs) => jobs, Err(error) => { error!(error = %error, "领取外部生成任务失败,等待下一轮重试"); - if await_one_task_or_sleep_or_shutdown( + if await_one_task_or_queue_wake_or_sleep_or_shutdown( &mut tasks, + &mut queue_wake, sleep(poll_interval), &mut shutdown, ) @@ -125,8 +129,13 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<() }; if jobs.is_empty() { - if await_one_task_or_sleep_or_shutdown(&mut tasks, sleep(poll_interval), &mut shutdown) - .await + if await_one_task_or_queue_wake_or_sleep_or_shutdown( + &mut tasks, + &mut queue_wake, + sleep(poll_interval), + &mut shutdown, + ) + .await { drain_external_generation_worker_tasks(&mut tasks).await; return Ok(()); @@ -148,6 +157,28 @@ pub(crate) async fn run_external_generation_worker(state: AppState) -> Result<() } } +async fn ensure_external_generation_queue_wake_subscription( + state: &AppState, + queue_wake: &mut Option, +) { + if queue_wake.is_some() { + return; + } + match state + .spacetime_client() + .subscribe_external_generation_queue_wake() + .await + { + Ok(subscription) => { + *queue_wake = Some(subscription); + info!("external generation worker 已订阅队列变更唤醒"); + } + Err(error) => { + warn!(error = %error, "external generation worker 订阅队列变更失败,暂时使用间隔兜底"); + } + } +} + type ExternalGenerationShutdownSignal = Pin + Send>>; fn external_generation_worker_shutdown_signal() -> ExternalGenerationShutdownSignal { @@ -202,20 +233,25 @@ async fn await_worker_task_or_shutdown( } } -async fn await_one_task_or_sleep_or_shutdown( +async fn await_one_task_or_queue_wake_or_sleep_or_shutdown( tasks: &mut JoinSet<()>, + queue_wake: &mut Option, sleeper: impl Future, shutdown: &mut ExternalGenerationShutdownSignal, ) -> bool { + let queue_changed = await_external_generation_queue_wake(queue_wake); + tokio::pin!(queue_changed); tokio::pin!(sleeper); if tasks.is_empty() { tokio::select! { _ = shutdown.as_mut() => true, + _ = &mut queue_changed => false, _ = &mut sleeper => false, } } else { tokio::select! { _ = shutdown.as_mut() => true, + _ = &mut queue_changed => false, _ = &mut sleeper => false, result = tasks.join_next() => { if let Some(Err(error)) = result { @@ -227,6 +263,21 @@ async fn await_one_task_or_sleep_or_shutdown( } } +async fn await_external_generation_queue_wake( + queue_wake: &mut Option, +) { + let result = if let Some(subscription) = queue_wake.as_mut() { + subscription.changed().await + } else { + std::future::pending().await + }; + if let Err(error) = result { + warn!(error = %error, "external generation worker 队列变更订阅已失效,等待兜底间隔后重连"); + *queue_wake = None; + std::future::pending::<()>().await; + } +} + async fn drain_external_generation_worker_tasks(tasks: &mut JoinSet<()>) { info!( in_flight_jobs = tasks.len(), diff --git a/server-rs/crates/api-server/src/external_generation_worker_controller.rs b/server-rs/crates/api-server/src/external_generation_worker_controller.rs index 3c4e588cc..087eb5c1b 100644 --- a/server-rs/crates/api-server/src/external_generation_worker_controller.rs +++ b/server-rs/crates/api-server/src/external_generation_worker_controller.rs @@ -1,6 +1,8 @@ use std::{collections::BTreeSet, future::Future, io, pin::Pin, process::Stdio, time::Duration}; -use spacetime_client::ExternalGenerationQueueStatsRecord; +use spacetime_client::{ + ExternalGenerationQueueStatsRecord, ExternalGenerationQueueWakeSubscription, +}; use tokio::{ process::Command, time::{Instant, sleep}, @@ -38,6 +40,7 @@ pub(crate) async fn run_external_generation_worker_controller( let config = ExternalGenerationWorkerControllerConfig::from_state(&state); let mut controller_state = ExternalGenerationWorkerControllerState::default(); let mut shutdown = external_generation_controller_shutdown_signal(); + let mut queue_wake = None; info!( min_workers = config.min_workers, @@ -51,6 +54,9 @@ pub(crate) async fn run_external_generation_worker_controller( ); loop { + ensure_external_generation_controller_queue_wake_subscription(&state, &mut queue_wake) + .await; + let tick = run_external_generation_controller_tick(&state, &config, &mut controller_state); tokio::select! { _ = shutdown.as_mut() => { @@ -65,17 +71,57 @@ pub(crate) async fn run_external_generation_worker_controller( } let next_tick = sleep(config.poll_interval); + let queue_changed = await_external_generation_controller_queue_wake(&mut queue_wake); tokio::pin!(next_tick); + tokio::pin!(queue_changed); tokio::select! { _ = shutdown.as_mut() => { info!("external generation worker controller 收到停机信号"); return Ok(()); } + _ = &mut queue_changed => {} _ = &mut next_tick => {} } } } +async fn ensure_external_generation_controller_queue_wake_subscription( + state: &AppState, + queue_wake: &mut Option, +) { + if queue_wake.is_some() { + return; + } + match state + .spacetime_client() + .subscribe_external_generation_queue_wake() + .await + { + Ok(subscription) => { + *queue_wake = Some(subscription); + info!("external generation worker controller 已订阅队列变更唤醒"); + } + Err(error) => { + warn!(error = %error, "external generation worker controller 订阅队列变更失败,暂时使用间隔兜底"); + } + } +} + +async fn await_external_generation_controller_queue_wake( + queue_wake: &mut Option, +) { + let result = if let Some(subscription) = queue_wake.as_mut() { + subscription.changed().await + } else { + std::future::pending().await + }; + if let Err(error) = result { + warn!(error = %error, "external generation worker controller 队列变更订阅已失效,等待兜底间隔后重连"); + *queue_wake = None; + std::future::pending::<()>().await; + } +} + async fn run_external_generation_controller_tick( state: &AppState, config: &ExternalGenerationWorkerControllerConfig, diff --git a/server-rs/crates/api-server/src/state.rs b/server-rs/crates/api-server/src/state.rs index 9fd400d49..243133fe5 100644 --- a/server-rs/crates/api-server/src/state.rs +++ b/server-rs/crates/api-server/src/state.rs @@ -404,13 +404,7 @@ impl AppState { RefreshSessionService::new(auth_store.clone(), config.refresh_session_ttl_days); // AI 编排服务当前先挂接内存态 store,后续再按 task table / procedure 接到 SpacetimeDB 真相源。 let ai_task_service = AiTaskService::new(InMemoryAiTaskStore::default()); - let spacetime_client = SpacetimeClient::new(SpacetimeClientConfig { - server_url: config.spacetime_server_url.clone(), - database: config.spacetime_database.clone(), - token: config.spacetime_token.clone(), - pool_size: config.spacetime_pool_size, - procedure_timeout: config.spacetime_procedure_timeout, - }); + let spacetime_client = SpacetimeClient::new(spacetime_client_config_for_process(&config)); let tracking_outbox = TrackingOutbox::from_config(&config, spacetime_client.clone()); let wallet_refund_outbox = WalletRefundOutbox::from_config(&config, spacetime_client.clone()); @@ -970,13 +964,8 @@ impl AppState { pub async fn try_restore_auth_store_from_spacetime( config: AppConfig, ) -> Result { - let spacetime_client = SpacetimeClient::new(SpacetimeClientConfig { - server_url: config.spacetime_server_url.clone(), - database: config.spacetime_database.clone(), - token: config.spacetime_token.clone(), - pool_size: config.spacetime_pool_size, - procedure_timeout: config.spacetime_procedure_timeout, - }); + let spacetime_client = + SpacetimeClient::new(spacetime_client_config_for_startup_restore(&config)); let mut spacetime_restore_available = false; let mut restore_errors = Vec::new(); @@ -1461,6 +1450,33 @@ fn auth_store_candidate_from_projection_view( })) } +fn spacetime_client_config_for_process(config: &AppConfig) -> SpacetimeClientConfig { + let runs_http = config.process_role.runs_http(); + SpacetimeClientConfig { + server_url: config.spacetime_server_url.clone(), + database: config.spacetime_database.clone(), + token: config.spacetime_token.clone(), + pool_size: if runs_http { + config.spacetime_pool_size + } else { + 1 + }, + procedure_timeout: config.spacetime_procedure_timeout, + subscribe_cached_read_models: runs_http, + } +} + +fn spacetime_client_config_for_startup_restore(config: &AppConfig) -> SpacetimeClientConfig { + SpacetimeClientConfig { + server_url: config.spacetime_server_url.clone(), + database: config.spacetime_database.clone(), + token: config.spacetime_token.clone(), + pool_size: 1, + procedure_timeout: config.spacetime_procedure_timeout, + subscribe_cached_read_models: false, + } +} + impl fmt::Display for AppStateInitError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { diff --git a/server-rs/crates/api-server/src/tracking_outbox.rs b/server-rs/crates/api-server/src/tracking_outbox.rs index eb04762b5..5c1e29b7b 100644 --- a/server-rs/crates/api-server/src/tracking_outbox.rs +++ b/server-rs/crates/api-server/src/tracking_outbox.rs @@ -525,6 +525,7 @@ mod tests { token: None, pool_size: 1, procedure_timeout: Duration::from_millis(10), + subscribe_cached_read_models: false, }), ) .expect("outbox should be enabled") diff --git a/server-rs/crates/api-server/src/wallet_refund_outbox.rs b/server-rs/crates/api-server/src/wallet_refund_outbox.rs index 25e7f9060..7c27580f3 100644 --- a/server-rs/crates/api-server/src/wallet_refund_outbox.rs +++ b/server-rs/crates/api-server/src/wallet_refund_outbox.rs @@ -391,6 +391,7 @@ mod tests { token: None, pool_size: 1, procedure_timeout: Duration::from_millis(10), + subscribe_cached_read_models: false, }), ) .expect("outbox should be enabled") diff --git a/server-rs/crates/spacetime-client/src/external_generation.rs b/server-rs/crates/spacetime-client/src/external_generation.rs index dfa77805b..fd217fb35 100644 --- a/server-rs/crates/spacetime-client/src/external_generation.rs +++ b/server-rs/crates/spacetime-client/src/external_generation.rs @@ -1,7 +1,162 @@ use super::*; use crate::mapper::*; +use spacetimedb_sdk::TableWithPrimaryKey; +use tokio::sync::watch; + +const EXTERNAL_GENERATION_QUEUE_WAKE_SUBSCRIPTION_QUERIES: [&str; 2] = [ + "SELECT * FROM external_generation_job WHERE status = 'pending'", + "SELECT * FROM external_generation_job WHERE status = 'running'", +]; + +pub struct ExternalGenerationQueueWakeSubscription { + connection: DbConnection, + _subscriptions: Vec, + _insert_callback: ExternalGenerationJobInsertCallbackId, + _update_callback: ExternalGenerationJobUpdateCallbackId, + _delete_callback: ExternalGenerationJobDeleteCallbackId, + runner: Option>, + receiver: watch::Receiver, +} + +impl ExternalGenerationQueueWakeSubscription { + pub async fn changed(&mut self) -> Result<(), SpacetimeClientError> { + self.receiver + .changed() + .await + .map_err(|_| SpacetimeClientError::ConnectDropped) + } +} + +impl Drop for ExternalGenerationQueueWakeSubscription { + fn drop(&mut self) { + let _ = self.connection.disconnect(); + if let Some(runner) = self.runner.take() { + drop(runner); + } + } +} impl SpacetimeClient { + pub async fn subscribe_external_generation_queue_wake( + &self, + ) -> Result { + let config = self.config.clone(); + let operation_timeout = config.procedure_timeout; + let (connected_sender, connected_receiver) = + oneshot::channel::>(); + let connected_sender = Arc::new(Mutex::new(Some(connected_sender))); + let connect_sender = connected_sender.clone(); + let disconnect_sender = connected_sender.clone(); + let connection = timeout( + operation_timeout, + tokio::task::spawn_blocking(move || { + DbConnection::builder() + .with_uri(config.server_url) + .with_database_name(config.database) + .with_token(config.token) + .on_connect(move |_, _, _| { + send_connect_once(&connect_sender, Ok(())); + }) + .on_disconnect(move |_, error| { + let message = error + .map(|error| error.to_string()) + .unwrap_or_else(|| "SpacetimeDB 队列订阅连接已断开".to_string()); + send_connect_once( + &disconnect_sender, + Err(SpacetimeClientError::Procedure(message)), + ); + }) + .build() + .map_err(|error| SpacetimeClientError::Build(error.to_string())) + }), + ) + .await + .map_err(|_| SpacetimeClientError::Timeout(SpacetimeClientStage::ConnectBuild))? + .map_err(|error| SpacetimeClientError::Runtime(error.to_string()))??; + + let runner = connection.run_threaded(); + timeout(operation_timeout, connected_receiver) + .await + .map_err(|_| SpacetimeClientError::Timeout(SpacetimeClientStage::ConnectHandshake))? + .map_err(|_| SpacetimeClientError::ConnectDropped)??; + + let (wake_sender, wake_receiver) = watch::channel(0u64); + let wake_counter = Arc::new(AtomicU64::new(0)); + let insert_sender = wake_sender.clone(); + let insert_counter = wake_counter.clone(); + let update_sender = wake_sender.clone(); + let update_counter = wake_counter.clone(); + let delete_sender = wake_sender.clone(); + let delete_counter = wake_counter.clone(); + let insert_callback = connection + .db() + .external_generation_job() + .on_insert(move |_, row| { + if external_generation_queue_row_should_wake(row) { + send_external_generation_queue_wake(&insert_sender, &insert_counter); + } + }); + let update_callback = + connection + .db() + .external_generation_job() + .on_update(move |_, old, new| { + if external_generation_queue_row_should_wake(old) + || external_generation_queue_row_should_wake(new) + { + send_external_generation_queue_wake(&update_sender, &update_counter); + } + }); + let delete_callback = connection + .db() + .external_generation_job() + .on_delete(move |_, row| { + if external_generation_queue_row_should_wake(row) { + send_external_generation_queue_wake(&delete_sender, &delete_counter); + } + }); + + let mut subscriptions = Vec::new(); + for query in EXTERNAL_GENERATION_QUEUE_WAKE_SUBSCRIPTION_QUERIES { + let (applied_sender, applied_receiver) = + oneshot::channel::>(); + let applied_sender = Arc::new(Mutex::new(Some(applied_sender))); + let on_applied_sender = applied_sender.clone(); + let on_error_sender = applied_sender.clone(); + let subscription = connection + .subscription_builder() + .on_applied(move |_| { + send_connect_once(&on_applied_sender, Ok(())); + }) + .on_error(move |_, error| { + send_connect_once( + &on_error_sender, + Err(SpacetimeClientError::Procedure(error.to_string())), + ); + }) + .subscribe(query); + + timeout(operation_timeout, applied_receiver) + .await + .map_err(|_| { + SpacetimeClientError::Timeout(SpacetimeClientStage::ReadModelSubscribe) + })? + .map_err(|_| SpacetimeClientError::ConnectDropped)??; + subscriptions.push(subscription); + } + send_external_generation_queue_wake(&wake_sender, &wake_counter); + + Ok(ExternalGenerationQueueWakeSubscription { + connection, + _subscriptions: subscriptions, + _insert_callback: insert_callback, + _update_callback: update_callback, + _delete_callback: delete_callback, + runner: Some(runner), + receiver: wake_receiver, + }) + } + pub async fn enqueue_external_generation_job( &self, input: ExternalGenerationJobEnqueueRecordInput, @@ -221,3 +376,12 @@ impl SpacetimeClient { .await } } + +fn external_generation_queue_row_should_wake(row: &ExternalGenerationJob) -> bool { + matches!(row.status.as_str(), "pending" | "running") +} + +fn send_external_generation_queue_wake(sender: &watch::Sender, counter: &AtomicU64) { + let next = counter.fetch_add(1, Ordering::Relaxed).saturating_add(1); + let _ = sender.send(next); +} diff --git a/server-rs/crates/spacetime-client/src/lib.rs b/server-rs/crates/spacetime-client/src/lib.rs index aa2977862..1daa4e410 100644 --- a/server-rs/crates/spacetime-client/src/lib.rs +++ b/server-rs/crates/spacetime-client/src/lib.rs @@ -143,6 +143,7 @@ pub mod editor_agent; pub mod editor_project; pub mod external_api_key; pub mod external_generation; +pub use external_generation::ExternalGenerationQueueWakeSubscription; pub mod inventory; pub mod jump_hop; @@ -162,7 +163,7 @@ use std::{ collections::HashMap, error::Error, fmt, - sync::atomic::{AtomicBool, Ordering}, + sync::atomic::{AtomicBool, AtomicU64, Ordering}, sync::{Arc, Mutex}, thread::JoinHandle, time::{Duration, Instant}, @@ -291,6 +292,7 @@ pub struct SpacetimeClientConfig { pub token: Option, pub pool_size: u32, pub procedure_timeout: Duration, + pub subscribe_cached_read_models: bool, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -853,9 +855,12 @@ impl SpacetimeClient { SpacetimeStageError::new(SpacetimeClientStage::ConnectHandshake, error) })?; - let read_model_subscriptions = self - .subscribe_cached_read_models(&connection, broken.clone(), operation_timeout) - .await?; + let read_model_subscriptions = if self.config.subscribe_cached_read_models { + self.subscribe_cached_read_models(&connection, broken.clone(), operation_timeout) + .await? + } else { + Vec::new() + }; Ok(PooledConnection { connection, @@ -1200,6 +1205,7 @@ mod tests { token: None, pool_size, procedure_timeout, + subscribe_cached_read_models: false, }) }