diff --git a/deploy/container/api-server.env.example b/deploy/container/api-server.env.example index 002bc4990..0cca35b21 100644 --- a/deploy/container/api-server.env.example +++ b/deploy/container/api-server.env.example @@ -58,6 +58,7 @@ AUTH_REFRESH_COOKIE_SECURE=false GENARRATIVE_SPACETIME_SERVER_URL=http://spacetimedb:3101 GENARRATIVE_SPACETIME_DATABASE=genarrative-loadtest GENARRATIVE_SPACETIME_TOKEN= +# HTTP 角色使用 8 条无 read-model 订阅的调用连接,并额外创建 1 条共享缓存读连接。 GENARRATIVE_SPACETIME_POOL_SIZE=8 GENARRATIVE_SPACETIME_PROCEDURE_TIMEOUT_SECONDS=45 diff --git a/deploy/env/api-server.env.example b/deploy/env/api-server.env.example index 2e4406d75..617553225 100644 --- a/deploy/env/api-server.env.example +++ b/deploy/env/api-server.env.example @@ -69,6 +69,7 @@ GENARRATIVE_SPACETIME_DATABASE=genarrative-prod GENARRATIVE_SPACETIME_TOKEN= # 只保存受保护文件路径;不要把 bootstrap secret 明文写入生产 env。 GENARRATIVE_SPACETIME_RUNTIME_SERVICE_BOOTSTRAP_SECRET_FILE=/var/lib/genarrative/spacetime/runtime-service-bootstrap-secret.txt +# HTTP 角色使用 8 条无 read-model 订阅的调用连接,并额外创建 1 条共享缓存读连接。 GENARRATIVE_SPACETIME_POOL_SIZE=8 GENARRATIVE_SPACETIME_PROCEDURE_TIMEOUT_SECONDS=45 diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 2bb7bef68..fefd40013 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -4160,3 +4160,12 @@ - 留存口径:筛选范围内 `profile_dashboard_state.created_at` 的北京时间注册日构成 cohort;在精确 `D+1` / `D+7` 存在有效登录 user scope 日聚合即留存。观察日必须早于今天;D1、D7 分别返回留存人数、可观察人数和四舍五入后的基点率,按人数加权汇总,零分母前端显示 `-`。 - 影响范围:SpacetimeDB Dashboard 聚合 procedure、`spacetime-client` facade、`/admin/api/dashboard` 与 shared contracts、`apps/admin-web` Dashboard 页面和运营文档。 - 验证方式:SpacetimeDB 聚合与 api-server 定向 Rust 测试、Dashboard Vitest、`npm run admin-web:typecheck`、`npm run check:spacetime-schema`、`npm run check:spacetime-runtime-access`、`npm run check:encoding`、`git diff --check`,并用桌面 / 移动浏览器核对留存、日期、每日 UV、零值与滚动同步。 + +## 2026-07-14 SpacetimeDB 调用池与缓存读连接分离 + +- 背景:HTTP 角色原先在每个 `GENARRATIVE_SPACETIME_POOL_SIZE` 槽位首次建连时订阅同一批 read model;release 配置为 8 时会保存 8 份相同行缓存和 subscription handles,放大 api-server 内存。 +- 决策:`pool_size` 条连接只承接 procedure / reducer 调用并保持无订阅;HTTP 角色额外创建且只创建 1 条共享缓存读连接,全部 `read_after_connect` 读取统一路由到该连接。缓存连接只在应用 facade 中作为只读用途,不宣称 SDK 或 identity 具备连接级只读权限。 +- 并发与恢复:缓存读连接通过 `Arc` 共享,读取不占用调用池 permit,也不使用单槽租约串行化;只有首次初始化和 broken 后重建使用单飞锁。首次建连与全部订阅共用一次总超时预算;required subscriptions 全部 applied、optional 阶段连接仍未 broken 后才发布新连接,旧连接由在途读取自然释放。 +- 就绪边界:HTTP `/readyz` 同时验证调用池握手和缓存读连接;required subscription 失败必须不就绪。非 HTTP worker / controller 不创建缓存读连接,继续使用 1 条调用连接和各自的队列窄订阅。 +- 运维口径:`GENARRATIVE_SPACETIME_POOL_SIZE=8` 表示 8 条调用连接,HTTP 基础拓扑另加 1 条缓存读连接;外部生成和充值过期监听的独立窄订阅不计入该值。读模型行缓存从 8 份降为 1 份,但 SDK 空 table metadata、8 条调用 socket 和 runner 仍存在,不承诺总 RSS 等比例降为八分之一。 +- 验证方式:`cargo test -p spacetime-client --manifest-path server-rs/Cargo.toml --lib`、`cargo check -p api-server --manifest-path server-rs/Cargo.toml`、`npm run check:encoding`、`git diff --check`;发布后在 8 个调用槽暖机后对比 api-server cgroup memory / PSS,并确认 `/readyz` 与代表性 gallery、公开详情、创作入口和用户标签读取正常。 diff --git a/docs/【后端架构】SpacetimeDB连接池租约Drop兜底与取消安全-2026-06-11.md b/docs/【后端架构】SpacetimeDB连接池租约Drop兜底与取消安全-2026-06-11.md index e72bec65a..c8045545f 100644 --- a/docs/【后端架构】SpacetimeDB连接池租约Drop兜底与取消安全-2026-06-11.md +++ b/docs/【后端架构】SpacetimeDB连接池租约Drop兜底与取消安全-2026-06-11.md @@ -30,11 +30,24 @@ ## 验收 -- `cargo test -p spacetime-client --manifest-path server-rs/Cargo.toml --lib`(35 通过,含上述新测试) +- `cargo test -p spacetime-client --manifest-path server-rs/Cargo.toml --lib`(44 通过,含上述连接池与缓存连接测试) - `cargo test -p api-server --manifest-path server-rs/Cargo.toml readyz`(2 通过) - `cargo check -p api-server --manifest-path server-rs/Cargo.toml` +## 2026-07-14 缓存读连接拆分 + +HTTP 角色的连接拓扑调整为“`pool_size` 条调用连接 + 1 条缓存读连接”: + +1. 调用连接池只承接 procedure / reducer 回调,不订阅 read model,也不装载订阅行;`GENARRATIVE_SPACETIME_POOL_SIZE=8` 表示 8 条调用连接,不包含额外的缓存读连接。 +2. 缓存读连接只由 `read_after_connect` 使用,并持有一份 required / optional read-model subscriptions。多个本地读取通过 `Arc` 共享同一 SDK `ClientCache`,初始化和断线重建使用单飞锁,读取本身不经过单槽 semaphore 串行化。 +3. 缓存读连接断线后,下一次读取只重建一条连接并等待 required subscriptions 全部 applied;首次建连与 required / optional 订阅共用一次总超时预算,optional 阶段发生断线时禁止发布 broken 连接;旧连接由在途读取持有到结束后再断开。 +4. `/readyz` 对 HTTP 角色同时检查调用池连接和缓存读连接;required subscription 失败时不得报告 ready。worker / controller 关闭 read-model cache 时不创建额外连接。 +5. 这里的“只读”是 facade 用途边界。SpacetimeDB Rust SDK 2.6 没有连接级 read-only builder;缓存连接仍使用相同 runtime identity,但代码不向 procedure / reducer 调用路径暴露它。 + +该拆分把相同 read-model 行缓存从 `pool_size` 份降为 1 份;SDK 每条连接仍会注册空 table metadata,总 WebSocket 数在 HTTP 角色中会从 `pool_size` 增加到 `pool_size + 1`,因此不能把总 RSS 简单承诺为原来的 `1 / pool_size`。 + ## 运维提示 - 此修复解决的是"取消导致的永久泄漏"。StDB 真慢时仍会出现成批 45 秒超时(连接被在途请求合法占用),那是容量/上游问题,应结合 `GENARRATIVE_SPACETIME_POOL_SIZE` 与 StDB 负载排查,不要再怀疑池泄漏。 - 健康检查 `/readyz` 在池被在途请求占满时仍可能短暂 503(stage=pool_acquire),恢复后自动转好,无需重启。 +- HTTP 角色排查连接数时按“调用池 + 1 条缓存读连接”计算;外部生成唤醒和充值过期监听还有各自的窄订阅连接,不属于调用池或缓存读连接。 diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md index 6cd9886b5..563c55b3f 100644 --- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md +++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md @@ -59,7 +59,9 @@ 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 读模型连接池。worker / controller 不执行模型定价 seed,启动时先调用受 runtime writer 鉴权的 queue-stats procedure 做只读预检,身份不匹配时 fail-fast;当前正式 systemd unit 通过共同加载 `/etc/genarrative/api-server.env` 继承同一 `GENARRATIVE_SPACETIME_TOKEN`,专属角色 env 示例不重复配置该 token。`GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS` 与 controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。 +HTTP 角色的 `GENARRATIVE_SPACETIME_POOL_SIZE` 只表示 procedure / reducer 调用池大小;池连接不订阅 read model。HTTP 角色会额外创建 1 条共享缓存读连接,统一保存公开读模型、近期游玩统计、创作入口配置和用户标签等订阅行,因此配置为 `8` 时基础连接拓扑是 8 条调用连接加 1 条缓存读连接。`/readyz` 同时检查调用池与缓存读连接的 required subscriptions,缓存未准备好时不能放量。 + +生产拆分角色时,`external-generation-worker` 和 `external-generation-controller` 的专属 env 示例会把 `GENARRATIVE_SPACETIME_POOL_SIZE` 覆盖为 `1`;非 HTTP 角色不创建 API 缓存读连接,只保留 `external_generation_job` 队列窄订阅作为响应式唤醒信号,实际抢占和扩缩容判断仍走 SpacetimeDB procedure。worker / controller 不执行模型定价 seed,启动时先调用受 runtime writer 鉴权的 queue-stats procedure 做只读预检,身份不匹配时 fail-fast;当前正式 systemd unit 通过共同加载 `/etc/genarrative/api-server.env` 继承同一 `GENARRATIVE_SPACETIME_TOKEN`,专属角色 env 示例不重复配置该 token。`GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS` 与 controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。 生产 worker 默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS=600`,只覆盖 worker 心跳抖动和短暂断连窗口,不再把 lease 当成完整任务时长;默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS=900`,角色动画 / 视频类长任务使用 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS=1800`。worker 在单次尝试超过执行预算后会停止当前尝试、写入失败 / 重试状态并释放 worker 槽位;如果 SpacetimeDB 当时不可写,当前租约最多再保留到 lease 过期,之后任务重新变为可领取。生产部署和 provision 脚本会给 `/etc/genarrative/api-server.env` 与 `/etc/genarrative/external-generation-worker.env` 补齐这些变量;已有自定义值不覆盖,只会把历史旧默认 `3600` 迁移为 `600`。 diff --git a/server-rs/crates/spacetime-client/README.md b/server-rs/crates/spacetime-client/README.md index 124c06c14..638f07223 100644 --- a/server-rs/crates/spacetime-client/README.md +++ b/server-rs/crates/spacetime-client/README.md @@ -24,7 +24,7 @@ 当前目录已不再只是占位。`WP-SC Spacetime Client` 在当前稳定 SpacetimeDB facade 范围内已经完成收尾: 1. 通过 `npm run spacetime:generate -- --rust-only` 生成并纳管公开 Rust bindings。 -2. `DbConnection` 连接池、握手等待、超时和断线清理已封装在 `SpacetimeClient` 内部。 +2. `DbConnection` 调用连接池、单一缓存读连接、握手等待、超时和断线清理已封装在 `SpacetimeClient` 内部;调用池不订阅 read model,HTTP 角色额外使用一条共享连接保存订阅行缓存。 3. 已稳定的 assets、auth、AI task、Big Fish、Custom World、Puzzle、Runtime/Profile/Save、Story session、combat、inventory、NPC facade 均通过 typed 方法对外暴露。 4. 生成绑定到 BFF record / module record 的 row snapshot mapper 已集中在 `mapper.rs`。 5. SDK 调用错误、reducer 业务错误、procedure 业务错误、缺快照错误和本地输入校验错误已统一收口到 `SpacetimeClientError` helper。 diff --git a/server-rs/crates/spacetime-client/src/lib.rs b/server-rs/crates/spacetime-client/src/lib.rs index 1c661eae2..e096b4ff4 100644 --- a/server-rs/crates/spacetime-client/src/lib.rs +++ b/server-rs/crates/spacetime-client/src/lib.rs @@ -370,6 +370,7 @@ pub struct AuthStoreProjectionSyncRecord { pub struct SpacetimeClient { config: SpacetimeClientConfig, pool: Arc, + cached_read_connection: Option>, health_state: Arc>, creation_entry_config_cache: Arc>>, feature_gate_config_cache: Arc>>>, @@ -454,11 +455,20 @@ struct PooledConnectionSlot { struct PooledConnection { connection: DbConnection, - _read_model_subscriptions: Vec, runner: Option>, broken: Arc, } +struct CachedReadConnectionState { + connection: RwLock>>, + initialize: tokio::sync::Mutex<()>, +} + +struct CachedReadConnection { + connection: PooledConnection, + _read_model_subscriptions: Vec, +} + struct PooledConnectionLease { slot_index: usize, connection: Option, @@ -505,10 +515,17 @@ impl SpacetimeClient { slots, permits: Arc::new(Semaphore::new(pool_size)), }); + let cached_read_connection = config.subscribe_cached_read_models.then(|| { + Arc::new(CachedReadConnectionState { + connection: RwLock::new(None), + initialize: tokio::sync::Mutex::new(()), + }) + }); Self { config, pool, + cached_read_connection, health_state: Arc::new(RwLock::new(SpacetimeClientHealthState::default())), creation_entry_config_cache: Arc::new(RwLock::new(None)), feature_gate_config_cache: Arc::new(RwLock::new(None)), @@ -655,11 +672,11 @@ impl SpacetimeClient { { let started_at = Instant::now(); let metrics_guard = telemetry::begin_read(read_name); - let lease = match self - .acquire_connection_with_timeout(self.config.procedure_timeout) + let connection = match self + .cached_read_connection_with_timeout(self.config.procedure_timeout) .await { - Ok(lease) => lease, + Ok(connection) => connection, Err(error) => { log_spacetime_client_failure( "read", @@ -673,14 +690,7 @@ impl SpacetimeClient { return final_result; } }; - let final_result = if let Some(connection) = lease.connection.as_ref() { - read(&connection.connection) - } else { - Err(SpacetimeClientError::Runtime( - "SpacetimeDB 连接租约缺少连接".to_string(), - )) - }; - self.release_connection(lease).await; + let final_result = read(&connection.connection.connection); if let Err(error) = &final_result { log_spacetime_client_failure( @@ -719,10 +729,22 @@ impl SpacetimeClient { }; let started_at = Instant::now(); let checked_at_micros = current_unix_micros(); - let result = self.acquire_connection_with_timeout(timeout).await; - match result { + let result = match self.acquire_connection_with_timeout(timeout).await { Ok(lease) => { self.release_connection(lease).await; + if self.cached_read_connection.is_some() { + let remaining = timeout.saturating_sub(started_at.elapsed()); + self.cached_read_connection_with_timeout(remaining) + .await + .map(|_| ()) + } else { + Ok(()) + } + } + Err(error) => Err(error), + }; + match result { + Ok(()) => { let mut health_state = self.health_state.write().await; health_state.last_success_at_micros = Some(checked_at_micros); health_state.last_error = None; @@ -826,6 +848,68 @@ impl SpacetimeClient { Ok(lease) } + async fn cached_read_connection_with_timeout( + &self, + operation_timeout: Duration, + ) -> Result, SpacetimeStageError> { + let started_at = Instant::now(); + let state = self.cached_read_connection.as_ref().ok_or_else(|| { + SpacetimeStageError::new( + SpacetimeClientStage::ReadCache, + SpacetimeClientError::Runtime("当前进程未启用 SpacetimeDB 缓存读连接".to_string()), + ) + })?; + + if let Some(connection) = state + .connection + .read() + .await + .clone() + .filter(|connection| !connection.is_broken()) + { + return Ok(connection); + } + + let remaining = operation_timeout.saturating_sub(started_at.elapsed()); + let _initialize = timeout(remaining, state.initialize.lock()) + .await + .map_err(|_| { + SpacetimeStageError::new( + SpacetimeClientStage::ReadCache, + SpacetimeClientError::Timeout(SpacetimeClientStage::ReadCache), + ) + })?; + + if let Some(connection) = state + .connection + .read() + .await + .clone() + .filter(|connection| !connection.is_broken()) + { + return Ok(connection); + } + + let remaining = operation_timeout.saturating_sub(started_at.elapsed()); + let connection = timeout(remaining, self.build_cached_read_connection(remaining)) + .await + .map_err(|_| { + SpacetimeStageError::new( + SpacetimeClientStage::ReadCache, + SpacetimeClientError::Timeout(SpacetimeClientStage::ReadCache), + ) + })??; + if connection.is_broken() { + return Err(SpacetimeStageError::new( + SpacetimeClientStage::ReadCache, + SpacetimeClientError::ConnectDropped, + )); + } + let connection = Arc::new(connection); + *state.connection.write().await = Some(connection.clone()); + Ok(connection) + } + async fn build_pooled_connection( &self, operation_timeout: Duration, @@ -894,21 +978,33 @@ impl SpacetimeClient { SpacetimeStageError::new(SpacetimeClientStage::ConnectHandshake, error) })?; - 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, - _read_model_subscriptions: read_model_subscriptions, runner: Some(runner), broken, }) } + async fn build_cached_read_connection( + &self, + operation_timeout: Duration, + ) -> Result { + let connection = self.build_pooled_connection(operation_timeout).await?; + let subscriptions = self + .subscribe_cached_read_models( + &connection.connection, + connection.broken.clone(), + operation_timeout, + ) + .await?; + Self::ensure_cached_read_connection_alive(connection.broken.as_ref())?; + + Ok(CachedReadConnection { + connection, + _read_model_subscriptions: subscriptions, + }) + } + async fn subscribe_cached_read_models( &self, connection: &DbConnection, @@ -917,6 +1013,7 @@ impl SpacetimeClient { ) -> Result, SpacetimeStageError> { let mut subscriptions = Vec::new(); for query in REQUIRED_CACHED_READ_MODEL_QUERIES { + Self::ensure_cached_read_connection_alive(broken.as_ref())?; let subscription = self .subscribe_cached_read_model_query( connection, @@ -930,7 +1027,8 @@ impl SpacetimeClient { } for query in OPTIONAL_CACHED_READ_MODEL_QUERIES { - if let Ok(subscription) = self + Self::ensure_cached_read_connection_alive(broken.as_ref())?; + match self .subscribe_cached_read_model_query( connection, broken.clone(), @@ -940,13 +1038,33 @@ impl SpacetimeClient { ) .await { - subscriptions.push(subscription); + Ok(subscription) => subscriptions.push(subscription), + Err(error) => { + Self::ensure_cached_read_connection_alive(broken.as_ref())?; + warn!( + query, + spacetime_stage = error.stage.as_str(), + error = %error.error, + "SpacetimeDB optional cached read-model subscription failed" + ); + } } } + Self::ensure_cached_read_connection_alive(broken.as_ref())?; Ok(subscriptions) } + fn ensure_cached_read_connection_alive(broken: &AtomicBool) -> Result<(), SpacetimeStageError> { + if broken.load(Ordering::SeqCst) { + return Err(SpacetimeStageError::new( + SpacetimeClientStage::ReadModelSubscribe, + SpacetimeClientError::ConnectDropped, + )); + } + Ok(()) + } + async fn subscribe_cached_read_model_query( &self, connection: &DbConnection, @@ -1105,6 +1223,12 @@ impl PooledConnection { } } +impl CachedReadConnection { + fn is_broken(&self) -> bool { + self.connection.is_broken() + } +} + impl Drop for PooledConnection { fn drop(&mut self) { let _ = self.connection.disconnect(); @@ -1119,6 +1243,10 @@ impl fmt::Debug for SpacetimeClient { f.debug_struct("SpacetimeClient") .field("config", &self.config) .field("pool_size", &self.pool.slots.len()) + .field( + "cached_read_connection_enabled", + &self.cached_read_connection.is_some(), + ) .finish() } } @@ -1230,6 +1358,94 @@ mod tests { }) } + fn test_client_with_cached_reads(pool_size: u32) -> SpacetimeClient { + SpacetimeClient::new(SpacetimeClientConfig { + server_url: "http://127.0.0.1:9".to_string(), + database: "pool-test".to_string(), + token: None, + pool_size, + procedure_timeout: Duration::from_millis(200), + subscribe_cached_read_models: true, + }) + } + + #[test] + fn cached_read_models_use_one_connection_outside_the_procedure_pool() { + let client = test_client_with_cached_reads(8); + + assert_eq!(client.pool.slots.len(), 8); + assert!(client.cached_read_connection.is_some()); + assert!( + client.pool.slots.iter().all(|slot| slot + .connection + .try_lock() + .expect("slot lock") + .is_none()) + ); + } + + #[test] + fn non_http_client_does_not_allocate_cached_read_connection_state() { + let client = test_client(1, Duration::from_millis(200)); + + assert_eq!(client.pool.slots.len(), 1); + assert!(client.cached_read_connection.is_none()); + } + + #[tokio::test] + async fn cached_read_fails_before_connect_when_role_disables_read_models() { + let client = test_client(1, Duration::from_millis(200)); + let result = client + .cached_read_connection_with_timeout(Duration::from_millis(200)) + .await; + + let error = match result { + Ok(_) => panic!("未启用 read model 的进程不应创建缓存读连接"), + Err(error) => error, + }; + assert_eq!(error.stage, SpacetimeClientStage::ReadCache); + assert!(matches!(error.error, SpacetimeClientError::Runtime(_))); + } + + #[tokio::test] + async fn cached_read_initialization_wait_respects_total_timeout() { + let client = test_client_with_cached_reads(1); + let state = client + .cached_read_connection + .as_ref() + .expect("cached read state") + .clone(); + let _initialize = state.initialize.lock().await; + let operation_timeout = Duration::from_millis(25); + let started_at = Instant::now(); + + let result = client + .cached_read_connection_with_timeout(operation_timeout) + .await; + let error = match result { + Ok(_) => panic!("singleflight wait should time out"), + Err(error) => error, + }; + + assert_eq!(error.stage, SpacetimeClientStage::ReadCache); + assert!(matches!( + error.error, + SpacetimeClientError::Timeout(SpacetimeClientStage::ReadCache) + )); + assert!(started_at.elapsed() < Duration::from_millis(150)); + } + + #[test] + fn broken_cached_read_connection_cannot_be_published() { + let broken = AtomicBool::new(true); + + let error = SpacetimeClient::ensure_cached_read_connection_alive(&broken) + .expect_err("broken cached read connection should fail"); + + assert_eq!(error.stage, SpacetimeClientStage::ReadModelSubscribe); + assert!(matches!(error.error, SpacetimeClientError::ConnectDropped)); + } + /// 复现线上故障机制:修复前请求 future 被取消时租约不会归还,槽位永久停留在 in_use, /// 后续 acquire 拿着 permit 空转挂死。修复后租约 Drop 必须同时复位槽位与 permit。 #[tokio::test]