轻量化外部生成队列进程

worker 和 controller 改为订阅 external_generation_job 队列变更来唤醒 claim 与扩缩容评估
非 HTTP 角色关闭 API 读模型订阅并将 SpacetimeDB 连接池收敛为 1
更新 worker/controller 环境模板和运维记忆,明确 poll interval 只作兜底
This commit is contained in:
2026-07-08 22:16:41 +08:00
parent 542fb09365
commit f1ca3e79c1
11 changed files with 316 additions and 24 deletions
+2
View File
@@ -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
+2
View File
@@ -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 自动续租;该值覆盖心跳抖动窗口即可。
@@ -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 activeworker 停机时停止 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`
@@ -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 接口确认业务状态同步完成。
@@ -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,7 +129,12 @@ 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)
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;
@@ -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<ExternalGenerationQueueWakeSubscription>,
) {
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<Box<dyn Future<Output = ()> + 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<ExternalGenerationQueueWakeSubscription>,
sleeper: impl Future<Output = ()>,
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<ExternalGenerationQueueWakeSubscription>,
) {
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(),
@@ -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<ExternalGenerationQueueWakeSubscription>,
) {
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<ExternalGenerationQueueWakeSubscription>,
) {
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,
+30 -14
View File
@@ -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<Self, AppStateInitError> {
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 {
@@ -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")
@@ -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")
@@ -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<SubscriptionHandle>,
_insert_callback: ExternalGenerationJobInsertCallbackId,
_update_callback: ExternalGenerationJobUpdateCallbackId,
_delete_callback: ExternalGenerationJobDeleteCallbackId,
runner: Option<JoinHandle<()>>,
receiver: watch::Receiver<u64>,
}
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<ExternalGenerationQueueWakeSubscription, SpacetimeClientError> {
let config = self.config.clone();
let operation_timeout = config.procedure_timeout;
let (connected_sender, connected_receiver) =
oneshot::channel::<Result<(), SpacetimeClientError>>();
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::<Result<(), SpacetimeClientError>>();
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<u64>, counter: &AtomicU64) {
let next = counter.fetch_add(1, Ordering::Relaxed).saturating_add(1);
let _ = sender.send(next);
}
+10 -4
View File
@@ -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<String>,
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,
})
}