From 01e3c4ae542a625dfb3c04c1a9714d15033c1d34 Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Mon, 24 Aug 2026 14:09:20 -0400 Subject: [PATCH 1/5] fix(schedule): wait for preload completion --- docs/development/recovery-internals.md | 5 ++++ src/domains/schedule/sink/domain_sink_impl.rs | 13 +++++++++- .../schedule/sink/mailbox_sink_impl.rs | 5 ++++ src/domains/schedule/sink/model.rs | 5 ++++ .../sink/tests/lifecycle_and_admin.rs | 24 +++++++++++++++++++ 5 files changed, 51 insertions(+), 1 deletion(-) diff --git a/docs/development/recovery-internals.md b/docs/development/recovery-internals.md index eef36a69..2c67c03c 100644 --- a/docs/development/recovery-internals.md +++ b/docs/development/recovery-internals.md @@ -18,6 +18,11 @@ Recovery behavior is designed to preserve committed durability state and reject `/targetz` is intentionally weaker than data-plane readiness. It can return `200` once the HTTP listener is usable and the process is not draining, even while storage or domain preload is still pending. It is only for a separate orchestration path; a customer-facing ALB target group must use `/healthz`. WebSocket upgrades and TCP sessions still reject data-plane traffic until the strict readiness gate passes. +Schedule preload waits for the actor-owned preload result rather than applying +an aggregate one-second reply deadline. Backend operations retain their own +deadlines, and actor failure disconnects the reply channel so boot still fails +closed instead of waiting indefinitely on a dead actor. + Live session state is never recovered during startup. Notice subscriptions, Stream live subscriptions and append sessions, KV open transactions, Queue inflight ownership tokens, RPC worker registrations and pending calls, Lease ownership, and Schedule subscriptions are rebuilt only by reconnecting clients when their domain contract permits it. ## Persistent Domain Partial-State Policy diff --git a/src/domains/schedule/sink/domain_sink_impl.rs b/src/domains/schedule/sink/domain_sink_impl.rs index aa28f795..e93b3479 100644 --- a/src/domains/schedule/sink/domain_sink_impl.rs +++ b/src/domains/schedule/sink/domain_sink_impl.rs @@ -185,6 +185,17 @@ impl ScheduleDomainSink { self.actor.stop(); } + #[cfg(test)] + pub(super) fn block_actor_for_tests( + &self, + entered: crossbeam_channel::Sender<()>, + release: crossbeam_channel::Receiver<()>, + ) { + self.actor + .try_send_high_priority(ScheduleDomainCommand::BlockForTests(entered, release)) + .expect("enqueue Schedule actor test block"); + } + /// # Errors /// /// Returns an error when listing column families or preloading a persisted @@ -199,7 +210,7 @@ impl ScheduleDomainSink { } reply_rx - .recv_timeout(std::time::Duration::from_secs(1)) + .recv() .map_err(|error| format!("schedule preload reply failed: {error}"))? } diff --git a/src/domains/schedule/sink/mailbox_sink_impl.rs b/src/domains/schedule/sink/mailbox_sink_impl.rs index e392fa46..6bd16b0c 100644 --- a/src/domains/schedule/sink/mailbox_sink_impl.rs +++ b/src/domains/schedule/sink/mailbox_sink_impl.rs @@ -60,6 +60,11 @@ impl Actor for ScheduleDomainActor { ScheduleDomainCommand::PanicForTests => { panic!("test Schedule domain actor panic"); } + #[cfg(test)] + ScheduleDomainCommand::BlockForTests(entered, release) => { + let _ = entered.send(()); + let _ = release.recv(); + } } } } diff --git a/src/domains/schedule/sink/model.rs b/src/domains/schedule/sink/model.rs index abcb2414..c25f94c3 100644 --- a/src/domains/schedule/sink/model.rs +++ b/src/domains/schedule/sink/model.rs @@ -213,6 +213,11 @@ pub(super) enum ScheduleDomainCommand { ForceDueScanForTests(usize, crossbeam_channel::Sender<()>), #[cfg(test)] PanicForTests, + #[cfg(test)] + BlockForTests( + crossbeam_channel::Sender<()>, + crossbeam_channel::Receiver<()>, + ), } pub(super) struct ScheduleDomainActor { diff --git a/src/domains/schedule/sink/tests/lifecycle_and_admin.rs b/src/domains/schedule/sink/tests/lifecycle_and_admin.rs index 5a318f2a..9bcb9a08 100644 --- a/src/domains/schedule/sink/tests/lifecycle_and_admin.rs +++ b/src/domains/schedule/sink/tests/lifecycle_and_admin.rs @@ -194,6 +194,30 @@ fn should_route_schedule_preload_through_actor_command() { assert_eq!(actor_count, 0); } +#[test] +fn should_wait_for_schedule_preload_reply_beyond_one_second() { + // Arrange + let store = crate::testkit::create_test_engine_with_cfs(vec![1]); + let router = Arc::new(Router::new()); + let admin_read_model = crate::control::admin::read_model::AdminReadModel::new(); + let sink = ScheduleDomainSink::new(store, router, admin_read_model); + let (entered_tx, entered_rx) = crossbeam_channel::bounded(1); + let (release_tx, release_rx) = crossbeam_channel::bounded(1); + sink.block_actor_for_tests(entered_tx, release_rx); + entered_rx.recv().expect("Schedule actor should block"); + + // Act + let preload_result = std::thread::scope(|scope| { + let preload = scope.spawn(|| sink.preload_persisted_families()); + std::thread::sleep(Duration::from_millis(1_100)); + release_tx.send(()).expect("release Schedule actor"); + preload.join().expect("join schedule preload") + }); + + // Assert + assert!(preload_result.is_ok()); +} + #[test] fn should_route_schedule_admin_refresh_through_actor_command() { // Arrange From b0800f906fdc403b01e44ffddac50896ea4cd306 Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Mon, 24 Aug 2026 14:32:59 -0400 Subject: [PATCH 2/5] fix(schedule): bound preload startup wait --- docs/development/recovery-internals.md | 11 +-- docs/user-guides/vars.md | 1 + src/boot/domains.rs | 7 +- src/boot/mod.rs | 1 + src/boot/runtime/config.rs | 29 ++++++++ src/boot/runtime/config/env.rs | 10 ++- .../config/tests/base_auth_and_network.rs | 27 +++++++ .../runtime/config/tests/cloud_storage.rs | 18 +++++ src/boot/stats/core.rs | 2 + src/domains/schedule/sink/domain_sink_impl.rs | 72 ++++++++++++++++++- src/domains/schedule/sink/mod.rs | 1 + .../sink/tests/lifecycle_and_admin.rs | 29 ++++++++ src/testkit/transport/server.rs | 4 ++ 13 files changed, 203 insertions(+), 9 deletions(-) diff --git a/docs/development/recovery-internals.md b/docs/development/recovery-internals.md index 2c67c03c..b11ed4f7 100644 --- a/docs/development/recovery-internals.md +++ b/docs/development/recovery-internals.md @@ -18,10 +18,13 @@ Recovery behavior is designed to preserve committed durability state and reject `/targetz` is intentionally weaker than data-plane readiness. It can return `200` once the HTTP listener is usable and the process is not draining, even while storage or domain preload is still pending. It is only for a separate orchestration path; a customer-facing ALB target group must use `/healthz`. WebSocket upgrades and TCP sessions still reject data-plane traffic until the strict readiness gate passes. -Schedule preload waits for the actor-owned preload result rather than applying -an aggregate one-second reply deadline. Backend operations retain their own -deadlines, and actor failure disconnects the reply channel so boot still fails -closed instead of waiting indefinitely on a dead actor. +Schedule preload waits for the actor-owned preload result within the +`FITZ_SCHEDULE_PRELOAD_TIMEOUT_SECS` startup watchdog. The default 120-second +deadline replaces the former one-second actor reply deadline while preserving a +bounded, diagnosable startup failure. Preload logs its configured deadline, +discovered family count, per-family progress at debug level, elapsed completion +time, and timeout. Actor failure also disconnects the reply channel so boot +fails closed before the watchdog expires. Live session state is never recovered during startup. Notice subscriptions, Stream live subscriptions and append sessions, KV open transactions, Queue inflight ownership tokens, RPC worker registrations and pending calls, Lease ownership, and Schedule subscriptions are rebuilt only by reconnecting clients when their domain contract permits it. diff --git a/docs/user-guides/vars.md b/docs/user-guides/vars.md index 426d3484..8ec636ac 100644 --- a/docs/user-guides/vars.md +++ b/docs/user-guides/vars.md @@ -131,6 +131,7 @@ For the auth and browser-perimeter checklist, see | FITZ_STORAGE_CACHE_PATH | Filesystem path | ./.fitz-cloud-cache | Local cache path for cloud-backed storage mode. | | FITZ_STORAGE_CLOUD_DURABILITY | background or strict | background | Cloud policy for broker-selected durable writes and client sync intent. `background` completes at Midge's local cloud commit barrier and uploads asynchronously; `strict` waits for provider acknowledgement. | | FITZ_STORAGE_MEMTABLE_BYTES | Unsigned integer byte count | Auto | Optional explicit memtable size override for embedded engine. | +| FITZ_SCHEDULE_PRELOAD_TIMEOUT_SECS | Positive integer second count | 120 | Maximum aggregate wait for required Schedule actor preload during startup. Expiry fails startup with an explicit timeout rather than leaving the broker wedged indefinitely. | | FITZ_QUEUE_WRITE_POLICY | fast, buffered, or strict | fast | Queue mutation write policy. `fast` skips WAL and flushes in the background; `buffered` uses local buffered WAL or cloud asynchronous durability; `strict` waits for local sync or cloud provider acknowledgement. | | FITZ_QUEUE_LOSS_WINDOW_MS | Positive integer millisecond count | 100 | Target background flush interval for fast queue writes. Accepted recent queue mutations can be lost before this window closes. | | FITZ_KV_IDLE_TRANSACTION_TTL_SECS | Positive integer second count | 300 | Maximum inactivity for an open KV transaction before the broker force-rolls it back and releases its broker-local resource lock. | diff --git a/src/boot/domains.rs b/src/boot/domains.rs index 6faaa495..5d8669a1 100644 --- a/src/boot/domains.rs +++ b/src/boot/domains.rs @@ -421,6 +421,7 @@ pub struct DomainSetupOptions { pub rpc_request_timeout: Option, pub stream_storage_layout: crate::domains::stream::StreamStorageLayout, pub kv_idle_transaction_ttl: std::time::Duration, + pub schedule_preload_timeout: std::time::Duration, } fn provisioned_route_families(options: &DomainSetupOptions) -> Vec { @@ -546,7 +547,7 @@ pub fn setup( ); register_domain_sink(DomainKind::Schedule, router, schedule_sink.clone()); schedule_sink - .preload_persisted_families() + .preload_persisted_families_with_timeout(options.schedule_preload_timeout) .map_err(|error| format!("schedule preload failed: {error}"))?; tracing::info!( "All {} domain sinks registered with router", @@ -595,6 +596,8 @@ mod tests { rpc_request_timeout: None, stream_storage_layout: crate::domains::stream::StreamStorageLayout::default(), kv_idle_transaction_ttl: std::time::Duration::from_mins(5), + schedule_preload_timeout: + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT, } } @@ -611,6 +614,8 @@ mod tests { rpc_request_timeout: None, stream_storage_layout: crate::domains::stream::StreamStorageLayout::default(), kv_idle_transaction_ttl: std::time::Duration::from_mins(5), + schedule_preload_timeout: + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT, } } diff --git a/src/boot/mod.rs b/src/boot/mod.rs index 512510a9..0945f854 100644 --- a/src/boot/mod.rs +++ b/src/boot/mod.rs @@ -211,6 +211,7 @@ fn register_domains_stage( rpc_request_timeout: None, stream_storage_layout: config.stream_storage_layout, kv_idle_transaction_ttl: Duration::from_secs(config.kv_idle_transaction_ttl_seconds), + schedule_preload_timeout: config.schedule_preload_timeout(), }; match domains::setup(router, store, &runtime.admin_read_model(), &options) { Ok(handles) => BootStage::Continue(handles), diff --git a/src/boot/runtime/config.rs b/src/boot/runtime/config.rs index ae4cc069..d0153f5a 100644 --- a/src/boot/runtime/config.rs +++ b/src/boot/runtime/config.rs @@ -15,6 +15,7 @@ const ENV_STORAGE_MEMTABLE_BYTES: &str = "FITZ_STORAGE_MEMTABLE_BYTES"; const ENV_QUEUE_WRITE_POLICY: &str = "FITZ_QUEUE_WRITE_POLICY"; const ENV_QUEUE_LOSS_WINDOW_MS: &str = "FITZ_QUEUE_LOSS_WINDOW_MS"; const ENV_KV_IDLE_TRANSACTION_TTL_SECS: &str = "FITZ_KV_IDLE_TRANSACTION_TTL_SECS"; +const ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS: &str = "FITZ_SCHEDULE_PRELOAD_TIMEOUT_SECS"; const ENV_DRAIN_GRACE_SECONDS: &str = "FITZ_DRAIN_GRACE_SECONDS"; const ENV_DRAIN_CLOSE_REASON: &str = "FITZ_DRAIN_CLOSE_REASON"; const DEFAULT_QUEUE_LOSS_WINDOW_MS: u64 = 100; @@ -327,6 +328,7 @@ mod env; use env::{ drain_close_reason_from_env, drain_grace_seconds_from_env, env_non_empty, kv_idle_transaction_ttl_seconds_from_env, queue_loss_window_ms_from_env, required_env, + schedule_preload_timeout_seconds_from_env, }; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -434,6 +436,14 @@ impl<'a> StorageConfig<'a> { format!("{ENV_KV_IDLE_TRANSACTION_TTL_SECS} must be greater than 0").into(), ); } + if let Some(error) = &config.schedule_preload_timeout_error { + return Err(error.clone().into()); + } + if config.schedule_preload_timeout_seconds == 0 { + return Err( + format!("{ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS} must be greater than 0").into(), + ); + } config .storage_memtable .validate() @@ -519,6 +529,9 @@ pub struct BootConfig { /// Maximum inactivity before an open KV transaction is force-rolled back. pub kv_idle_transaction_ttl_seconds: u64, pub(crate) kv_idle_transaction_ttl_error: Option, + /// Maximum wait for required Schedule preload during broker startup. + pub schedule_preload_timeout_seconds: u64, + pub(crate) schedule_preload_timeout_error: Option, /// Whether an external TLS terminator is explicitly protecting public listeners. pub assume_external_tls: bool, pub(crate) local_listener_exposure: LocalListenerExposure, @@ -599,6 +612,11 @@ impl BootConfig { .then(|| Duration::from_millis(self.queue_loss_window_ms)) } + #[must_use] + pub fn schedule_preload_timeout(&self) -> Duration { + Duration::from_secs(self.schedule_preload_timeout_seconds) + } + #[must_use] pub fn queue_write_policy_defaulted_fast(&self) -> bool { self.queue_write_policy_source.is_defaulted() @@ -656,6 +674,8 @@ impl Default for BootConfig { let (queue_loss_window_ms, queue_loss_window_error) = queue_loss_window_ms_from_env(); let (kv_idle_transaction_ttl_seconds, kv_idle_transaction_ttl_error) = kv_idle_transaction_ttl_seconds_from_env(); + let (schedule_preload_timeout_seconds, schedule_preload_timeout_error) = + schedule_preload_timeout_seconds_from_env(); let (queue_write_policy, queue_write_policy_source) = QueueWritePolicy::from_env_with_source(); let drain_close_reason = drain_close_reason_from_env(); @@ -690,6 +710,8 @@ impl Default for BootConfig { queue_loss_window_error, kv_idle_transaction_ttl_seconds, kv_idle_transaction_ttl_error, + schedule_preload_timeout_seconds, + schedule_preload_timeout_error, assume_external_tls, local_listener_exposure, ws_allowed_origins, @@ -800,6 +822,13 @@ impl BootConfig { self } + #[must_use] + pub fn with_schedule_preload_timeout_seconds(mut self, seconds: u64) -> Self { + self.schedule_preload_timeout_seconds = seconds; + self.schedule_preload_timeout_error = None; + self + } + #[must_use] pub fn with_drain_close_reason(mut self, reason: impl Into) -> Self { self.drain_close_reason = reason.into(); diff --git a/src/boot/runtime/config/env.rs b/src/boot/runtime/config/env.rs index 620d954e..8e617e47 100644 --- a/src/boot/runtime/config/env.rs +++ b/src/boot/runtime/config/env.rs @@ -1,7 +1,7 @@ use super::{ DEFAULT_DRAIN_CLOSE_REASON, DEFAULT_DRAIN_GRACE_SECONDS, DEFAULT_KV_IDLE_TRANSACTION_TTL_SECS, DEFAULT_QUEUE_LOSS_WINDOW_MS, ENV_DRAIN_CLOSE_REASON, ENV_DRAIN_GRACE_SECONDS, - ENV_KV_IDLE_TRANSACTION_TTL_SECS, ENV_QUEUE_LOSS_WINDOW_MS, + ENV_KV_IDLE_TRANSACTION_TTL_SECS, ENV_QUEUE_LOSS_WINDOW_MS, ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, }; pub(super) fn env_non_empty(key: &str) -> Option { @@ -84,6 +84,14 @@ pub(super) fn kv_idle_transaction_ttl_seconds_from_env() -> (u64, Option } } +pub(super) fn schedule_preload_timeout_seconds_from_env() -> (u64, Option) { + positive_u64_from_env( + ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT.as_secs(), + "second count", + ) +} + pub(super) fn drain_close_reason_from_env() -> String { env_non_empty(ENV_DRAIN_CLOSE_REASON).unwrap_or_else(|| DEFAULT_DRAIN_CLOSE_REASON.to_string()) } diff --git a/src/boot/runtime/config/tests/base_auth_and_network.rs b/src/boot/runtime/config/tests/base_auth_and_network.rs index 9159fb98..3cc37bf2 100644 --- a/src/boot/runtime/config/tests/base_auth_and_network.rs +++ b/src/boot/runtime/config/tests/base_auth_and_network.rs @@ -47,6 +47,7 @@ pub(super) fn with_auth_env(values: &[(&str, &str)], test: impl FnOnce() -> T ENV_QUEUE_WRITE_POLICY, ENV_QUEUE_LOSS_WINDOW_MS, ENV_KV_IDLE_TRANSACTION_TTL_SECS, + ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, ENV_DRAIN_GRACE_SECONDS, ENV_DRAIN_CLOSE_REASON, ]; @@ -97,6 +98,7 @@ pub(super) fn with_storage_env(values: &[(&str, &str)], test: impl FnOnce() - ENV_QUEUE_WRITE_POLICY, ENV_QUEUE_LOSS_WINDOW_MS, ENV_KV_IDLE_TRANSACTION_TTL_SECS, + ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, "AWS_REGION", "AWS_DEFAULT_REGION", "AZURE_STORAGE_ACCOUNT_NAME", @@ -174,6 +176,10 @@ pub(super) fn should_create_default_boot_config() { ); assert!(config.queue_write_policy_defaulted_fast()); assert_eq!(config.queue_loss_window_ms, DEFAULT_QUEUE_LOSS_WINDOW_MS); + assert_eq!( + config.schedule_preload_timeout(), + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT + ); assert!(config.queue_write_options().is_best_effort()); assert_eq!(config.drain_grace_seconds, DEFAULT_DRAIN_GRACE_SECONDS); assert_eq!(config.drain_close_reason, DEFAULT_DRAIN_CLOSE_REASON); @@ -278,6 +284,27 @@ pub(super) fn should_read_drain_config_from_environment() { ); } +#[test] +#[serial] +pub(super) fn should_read_schedule_preload_timeout_from_environment() { + with_auth_env( + &[ + ("FITZ_AUTH_REQUIRED", "false"), + (ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, "75"), + ], + || { + // Arrange + + // Act + let config = BootConfig::default(); + + // Assert + assert_eq!(config.schedule_preload_timeout(), Duration::from_secs(75)); + assert!(config.validate().is_ok()); + }, + ); +} + #[test] #[serial] pub(super) fn should_reject_invalid_drain_grace_from_environment() { diff --git a/src/boot/runtime/config/tests/cloud_storage.rs b/src/boot/runtime/config/tests/cloud_storage.rs index 539a41b3..7d36ac07 100644 --- a/src/boot/runtime/config/tests/cloud_storage.rs +++ b/src/boot/runtime/config/tests/cloud_storage.rs @@ -261,6 +261,24 @@ fn should_reject_invalid_kv_idle_transaction_ttl() { }); } +#[test] +#[serial] +fn should_reject_invalid_schedule_preload_timeout() { + with_storage_env(&[(ENV_SCHEDULE_PRELOAD_TIMEOUT_SECS, "0")], || { + // Arrange + + // Act + let result = BootConfig::new().validate(); + + // Assert + assert!(result.is_err()); + assert!(result + .unwrap_err() + .to_string() + .contains("FITZ_SCHEDULE_PRELOAD_TIMEOUT_SECS must be greater than 0")); + }); +} + #[test] fn should_keep_non_cloud_sync_write_options_local() { // Arrange diff --git a/src/boot/stats/core.rs b/src/boot/stats/core.rs index bd187658..dcb31e25 100644 --- a/src/boot/stats/core.rs +++ b/src/boot/stats/core.rs @@ -166,6 +166,8 @@ impl Runtime { rpc_request_timeout: None, stream_storage_layout: crate::domains::stream::StreamStorageLayout::default(), kv_idle_transaction_ttl: Duration::from_mins(5), + schedule_preload_timeout: + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT, }, ) .expect("setup domains"); diff --git a/src/domains/schedule/sink/domain_sink_impl.rs b/src/domains/schedule/sink/domain_sink_impl.rs index e93b3479..b2ef9ffa 100644 --- a/src/domains/schedule/sink/domain_sink_impl.rs +++ b/src/domains/schedule/sink/domain_sink_impl.rs @@ -10,6 +10,8 @@ use crate::dispatch::protocol::frame_context::FrameContext; use crate::runtime::routing::{Route, RouteAddress, RouteFamily}; type PendingAckRetryMap = HashMap>; +pub(crate) const DEFAULT_SCHEDULE_PRELOAD_TIMEOUT: std::time::Duration = + std::time::Duration::from_secs(120); type LivePublishCandidate = ( crate::runtime::routing::RouteFamily, u64, @@ -201,6 +203,20 @@ impl ScheduleDomainSink { /// Returns an error when listing column families or preloading a persisted /// schedule actor fails. pub fn preload_persisted_families(&self) -> Result<(), String> { + self.preload_persisted_families_with_timeout(DEFAULT_SCHEDULE_PRELOAD_TIMEOUT) + } + + /// # Errors + /// + /// Returns an error when the actor cannot be reached, preload fails, or the + /// actor does not reply before `timeout`. + pub(crate) fn preload_persisted_families_with_timeout( + &self, + timeout: std::time::Duration, + ) -> Result<(), String> { + let started_at = std::time::Instant::now(); + let timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX); + tracing::info!(domain = "schedule", timeout_ms, "Schedule preload started"); let (reply_tx, reply_rx) = crossbeam_channel::bounded(1); if let Err(error) = self .actor @@ -209,9 +225,33 @@ impl ScheduleDomainSink { return Err(format!("schedule preload enqueue failed: {error}")); } - reply_rx - .recv() - .map_err(|error| format!("schedule preload reply failed: {error}"))? + match reply_rx.recv_timeout(timeout) { + Ok(result) => { + result?; + tracing::info!( + domain = "schedule", + elapsed_ms = + u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + "Schedule preload completed" + ); + Ok(()) + } + Err(crossbeam_channel::RecvTimeoutError::Timeout) => { + tracing::error!( + domain = "schedule", + timeout_ms, + elapsed_ms = + u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + "Schedule preload timed out" + ); + Err(format!( + "schedule preload reply timed out after {timeout_ms}ms" + )) + } + Err(crossbeam_channel::RecvTimeoutError::Disconnected) => { + Err("schedule preload reply failed: actor reply channel disconnected".to_string()) + } + } } pub(crate) fn is_active(&self) -> bool { @@ -407,13 +447,24 @@ impl ScheduleDomainRuntime<'_> { /// Returns an error when listing column families or preloading a persisted /// schedule actor fails. pub(super) fn preload_persisted_families(&self) -> Result<(), String> { + let started_at = std::time::Instant::now(); let column_families = self .core .store .list_column_families() .map_err(|e| format!("list schedule column families failed: {e}"))?; + let persisted_family_count = column_families + .iter() + .filter(|column_family| column_family.id() != 0) + .count(); + tracing::info!( + domain = "schedule", + persisted_family_count, + "Schedule preload discovered persisted families" + ); let mut actors = self.core.actors.lock(); + let mut preloaded_family_count = 0_usize; for column_family in column_families { if column_family.id() == 0 { continue; @@ -430,6 +481,14 @@ impl ScheduleDomainRuntime<'_> { self.core.write_options, )?; actors.insert(family, actor); + preloaded_family_count = preloaded_family_count.saturating_add(1); + tracing::debug!( + domain = "schedule", + route_family = family.id(), + preloaded_family_count, + persisted_family_count, + "Schedule persisted family preloaded" + ); } // Seed the rolling-window acknowledgement counter from persisted @@ -449,6 +508,13 @@ impl ScheduleDomainRuntime<'_> { drop(actors); self.schedule_admin_snapshot(true); + tracing::info!( + domain = "schedule", + preloaded_family_count, + persisted_family_count, + elapsed_ms = u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + "Schedule actor projection preload completed" + ); Ok(()) } diff --git a/src/domains/schedule/sink/mod.rs b/src/domains/schedule/sink/mod.rs index 79b3788c..fe92c4c8 100644 --- a/src/domains/schedule/sink/mod.rs +++ b/src/domains/schedule/sink/mod.rs @@ -6,6 +6,7 @@ mod model; mod test_helpers; pub use domain_sink_impl::ScheduleObservability; +pub(crate) use domain_sink_impl::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT; pub use model::ScheduleDomainSink; #[cfg(test)] diff --git a/src/domains/schedule/sink/tests/lifecycle_and_admin.rs b/src/domains/schedule/sink/tests/lifecycle_and_admin.rs index 9bcb9a08..c4341437 100644 --- a/src/domains/schedule/sink/tests/lifecycle_and_admin.rs +++ b/src/domains/schedule/sink/tests/lifecycle_and_admin.rs @@ -218,6 +218,35 @@ fn should_wait_for_schedule_preload_reply_beyond_one_second() { assert!(preload_result.is_ok()); } +#[test] +fn should_timeout_schedule_preload_when_actor_does_not_reply_before_deadline() { + // Arrange + let store = crate::testkit::create_test_engine_with_cfs(vec![1]); + let router = Arc::new(Router::new()); + let admin_read_model = crate::control::admin::read_model::AdminReadModel::new(); + let sink = ScheduleDomainSink::new(store, router, admin_read_model); + let (entered_tx, entered_rx) = crossbeam_channel::bounded(1); + let (release_tx, release_rx) = crossbeam_channel::bounded(1); + sink.block_actor_for_tests(entered_tx, release_rx); + entered_rx.recv().expect("Schedule actor should block"); + + // Act + let preload_result = std::thread::scope(|scope| { + let release = scope.spawn(|| { + std::thread::sleep(Duration::from_millis(100)); + release_tx.send(()).expect("release Schedule actor"); + }); + let result = sink.preload_persisted_families_with_timeout(Duration::from_millis(20)); + release.join().expect("join Schedule actor release"); + result + }); + + // Assert + assert!(preload_result + .expect_err("Schedule preload should time out") + .contains("timed out")); +} + #[test] fn should_route_schedule_admin_refresh_through_actor_command() { // Arrange diff --git a/src/testkit/transport/server.rs b/src/testkit/transport/server.rs index d62bc2c3..3afc7466 100644 --- a/src/testkit/transport/server.rs +++ b/src/testkit/transport/server.rs @@ -368,6 +368,9 @@ impl TestServer { queue_loss_window_error: None, kv_idle_transaction_ttl_seconds: 300, kv_idle_transaction_ttl_error: None, + schedule_preload_timeout_seconds: + crate::domains::schedule::sink::DEFAULT_SCHEDULE_PRELOAD_TIMEOUT.as_secs(), + schedule_preload_timeout_error: None, assume_external_tls: false, local_listener_exposure: crate::boot::runtime::LocalListenerExposure::Direct, ws_allowed_origins, @@ -407,6 +410,7 @@ impl TestServer { kv_idle_transaction_ttl: std::time::Duration::from_secs( boot_config.kv_idle_transaction_ttl_seconds, ), + schedule_preload_timeout: boot_config.schedule_preload_timeout(), }, )?; runtime.attach_domains(domains); From db92d703378e50cfff8532a6b20d9e828622cde8 Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Mon, 24 Aug 2026 14:33:09 -0400 Subject: [PATCH 3/5] chore: satisfy Rust 1.98 clippy --- src/api/admin/auth.rs | 3 +-- src/domains/stream/storage/compact_page_values.rs | 4 +++- src/observability/global.rs | 14 ++++++-------- 3 files changed, 10 insertions(+), 11 deletions(-) diff --git a/src/api/admin/auth.rs b/src/api/admin/auth.rs index 1332ce67..1c9aa504 100644 --- a/src/api/admin/auth.rs +++ b/src/api/admin/auth.rs @@ -465,8 +465,7 @@ impl AdminAuth { AdminRouteFamilyAccess::Explicit(values) => values.iter().all(|value| { value .parse::() - .ok() - .is_some_and(|family| provisioned.contains(&family)) + .is_ok_and(|family| provisioned.contains(&family)) }), } } diff --git a/src/domains/stream/storage/compact_page_values.rs b/src/domains/stream/storage/compact_page_values.rs index 232f2982..d65dd327 100644 --- a/src/domains/stream/storage/compact_page_values.rs +++ b/src/domains/stream/storage/compact_page_values.rs @@ -451,7 +451,9 @@ impl PostingPageValue { return Err("decode posting page value: invalid offset payload".to_string()); } let entries = bytes[6..] - .chunks_exact(24) + .as_chunks::<24>() + .0 + .iter() .map(|chunk| { let mut offset = [0u8; 8]; let mut parent = [0u8; 8]; diff --git a/src/observability/global.rs b/src/observability/global.rs index 14b0bb04..7c29df23 100644 --- a/src/observability/global.rs +++ b/src/observability/global.rs @@ -78,14 +78,12 @@ pub fn gauge_dec(name: &str) { /// collecting attribution data. pub fn hot_path_metrics_enabled() -> bool { *HOT_PATH_METRICS_ENABLED.get_or_init(|| { - std::env::var("FITZ_HOT_PATH_METRICS") - .ok() - .is_some_and(|value| { - matches!( - value.trim().to_ascii_lowercase().as_str(), - "1" | "true" | "yes" | "on" - ) - }) + std::env::var("FITZ_HOT_PATH_METRICS").is_ok_and(|value| { + matches!( + value.trim().to_ascii_lowercase().as_str(), + "1" | "true" | "yes" | "on" + ) + }) }) } From ec85e4ddc9ef7fd62a05dcc569397f23b109a70f Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Mon, 24 Aug 2026 14:37:48 -0400 Subject: [PATCH 4/5] refactor(schedule): centralize duration logging --- src/domains/schedule/sink/domain_sink_impl.rs | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/src/domains/schedule/sink/domain_sink_impl.rs b/src/domains/schedule/sink/domain_sink_impl.rs index b2ef9ffa..ca4f079a 100644 --- a/src/domains/schedule/sink/domain_sink_impl.rs +++ b/src/domains/schedule/sink/domain_sink_impl.rs @@ -12,6 +12,11 @@ use crate::runtime::routing::{Route, RouteAddress, RouteFamily}; type PendingAckRetryMap = HashMap>; pub(crate) const DEFAULT_SCHEDULE_PRELOAD_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(120); + +fn duration_millis(duration: std::time::Duration) -> u64 { + u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) +} + type LivePublishCandidate = ( crate::runtime::routing::RouteFamily, u64, @@ -215,7 +220,7 @@ impl ScheduleDomainSink { timeout: std::time::Duration, ) -> Result<(), String> { let started_at = std::time::Instant::now(); - let timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX); + let timeout_ms = duration_millis(timeout); tracing::info!(domain = "schedule", timeout_ms, "Schedule preload started"); let (reply_tx, reply_rx) = crossbeam_channel::bounded(1); if let Err(error) = self @@ -230,8 +235,7 @@ impl ScheduleDomainSink { result?; tracing::info!( domain = "schedule", - elapsed_ms = - u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + elapsed_ms = duration_millis(started_at.elapsed()), "Schedule preload completed" ); Ok(()) @@ -240,8 +244,7 @@ impl ScheduleDomainSink { tracing::error!( domain = "schedule", timeout_ms, - elapsed_ms = - u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + elapsed_ms = duration_millis(started_at.elapsed()), "Schedule preload timed out" ); Err(format!( @@ -512,7 +515,7 @@ impl ScheduleDomainRuntime<'_> { domain = "schedule", preloaded_family_count, persisted_family_count, - elapsed_ms = u64::try_from(started_at.elapsed().as_millis()).unwrap_or(u64::MAX), + elapsed_ms = duration_millis(started_at.elapsed()), "Schedule actor projection preload completed" ); Ok(()) From 3f7e73230179c4a6d71c1a22d329fef7104b595c Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Mon, 24 Aug 2026 14:42:24 -0400 Subject: [PATCH 5/5] ci: cancel superseded backend runs --- .github/workflows/ci-backend.yml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/.github/workflows/ci-backend.yml b/.github/workflows/ci-backend.yml index 0611259f..908dd976 100644 --- a/.github/workflows/ci-backend.yml +++ b/.github/workflows/ci-backend.yml @@ -9,6 +9,10 @@ on: branches: - main +concurrency: + group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + permissions: contents: read