Contract-suite scheduler rows (task suite-scheduler-rows): the determinism-posture extension in properties.rs's module doc (runner-driving rows admitted — spawned run_schedules(stop) tasks on a core StopToken against state outcomes only, elapsed-boundary-band tolerances, never timing-value assertions, never two-task race windows; ADR-017 §2 class 4) and three version-stamped rows: scheduler_boundary_fires (fired jobs are ordinary claimable work with ScheduleOpts over the plain-queue derived defaults 300/9/5/none + payload exact, clean-stop Ok(()), fired count inside the elapsed-boundary band — no double-fire per boundary while one leader runs; ADR-009 §3/§4, ADR-020 §3, ADR-019 §4), scheduler_bounded_catchup (runner-less downtime proves no fire without a runner, ≥3 elapsed boundaries replay boundary-by-boundary bounded below the 64-cap and inside the band; ADR-009 §4), scheduler_leadership_discipline (two spawned runners on one store: exactly one Ok(())/Err(LeadershipLost) pair by value, no duplicated fires inside the band; ADR-009 §1/§6, ADR-019 §4) — wired into both engines' suite targets (SQLite tests, pg harness_row!s), suite tokio dep added for the runner rows. Dispositions recorded in Notes: the beyond-cap skip-forward leg stays pinned engine-side (both engines' scheduler tests already backdate next_fire_at directly — catch_up_replays_up_to_the_cap_then_skips_forward twins), and the pg fire-wake parity gap is not demanded by these rows' shapes (claim-polling observation only; the one-call wake_tx disposition recorded for a later task). Verified: sqlite suite 16/16, pg suite 16/16 vs harness twice + solo re-runs of the three rows per engine (determinism), workspace build/test green, clippy -D warnings, fmt clean
This commit is contained in:
1 parent
5c9ae6a7fc
commit
7f749ac673
6 files changed
+523
-16
No files matched your search
@@ -9,6 +9,4 @@ publish = false
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
alkstore = { version = "0.1", path = "../alkstore" }
|
alkstore = { version = "0.1", path = "../alkstore" }
|
||||||
serde_json = "1"
|
serde_json = "1"
|
||||||
|
tokio = { version = "1", features = ["rt", "time"] }
|
||||||
[dev-dependencies]
|
|
||||||
tokio = { version = "1", features = ["macros", "rt"] }
|
|
||||||
@@ -42,5 +42,6 @@ pub use properties::{
|
|||||||
in_tx_reads_see_own_writes, job_handle_validity_predicate,
|
in_tx_reads_see_own_writes, job_handle_validity_predicate,
|
||||||
name_validation_rejects_empty_and_reserved, payload_round_trip_stores_exact_encoding,
|
name_validation_rejects_empty_and_reserved, payload_round_trip_stores_exact_encoding,
|
||||||
payload_too_large_never_produced_on_sqlite, payload_too_large_produced_on_pg,
|
payload_too_large_never_produced_on_sqlite, payload_too_large_produced_on_pg,
|
||||||
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, sweep_no_stranded_rows,
|
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, scheduler_boundary_fires,
|
||||||
|
scheduler_bounded_catchup, scheduler_leadership_discipline, sweep_no_stranded_rows,
|
||||||
};
|
};
|
||||||
@@ -7,21 +7,36 @@
|
|||||||
//! is the suite's purpose as a test-side artifact (relaxing the
|
//! is the suite's purpose as a test-side artifact (relaxing the
|
||||||
//! family-wide no-panic rule here only).
|
//! family-wide no-panic rule here only).
|
||||||
//!
|
//!
|
||||||
//! Determinism: no wall-clock-timing assertions and no two-task race
|
//! Determinism: no wall-clock-timing-value assertions and no two-task
|
||||||
//! windows — orchestration that would need one (runner spawning,
|
//! race windows; rows keep to single-store sequential drives and the
|
||||||
//! cross-connection wakes) stays engine-test-side; rows keep to
|
//! factory's isolation guarantee.
|
||||||
//! single-store sequential drives and the factory's isolation
|
//!
|
||||||
//! guarantee. Rows exercise a *close/dispose* arm by dropping the
|
//! Orchestration posture: cross-connection wake orchestration stays
|
||||||
|
//! engine-test-side. **Runner-driving rows are admitted as an
|
||||||
|
//! extension** (the backlog's scheduler row demands the tick machinery
|
||||||
|
//! in the suite, and it cannot be pinned sequentially —
|
||||||
|
//! `run_schedules` is the drive): a spawned `run_schedules(stop)`
|
||||||
|
//! task with a core [`StopToken`], driven against state outcomes only
|
||||||
|
//! (jobs fired, leadership returned), with tolerances bounded well
|
||||||
|
//! above the boundary resolution — elapsed-boundary bands, never
|
||||||
|
//! timing-value assertions, never two-task race windows. The wave-5
|
||||||
|
//! scheduler rows are the extension's riders; this is a suite-side
|
||||||
|
//! posture change (ADR-017 §2 class 4 — the suite tests the contract,
|
||||||
|
//! it is not the contract).
|
||||||
|
//!
|
||||||
|
//! Rows exercise a *close/dispose* arm by dropping the
|
||||||
//! store handle (the dispose carrier every engine implements — the
|
//! store handle (the dispose carrier every engine implements — the
|
||||||
//! drop path closes receivers terminally); teardown-then-fresh-reopen
|
//! drop path closes receivers terminally); teardown-then-fresh-reopen
|
||||||
//! rows would be factory-contract tests (ADR-022), not contract rows,
|
//! rows would be factory-contract tests (ADR-022), not contract rows,
|
||||||
//! so none exist.
|
//! so none exist.
|
||||||
|
|
||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
use serde_json::json;
|
use serde_json::json;
|
||||||
|
|
||||||
use alkstore::{
|
use alkstore::{
|
||||||
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, JobState, RESERVED_PREFIX, Result,
|
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, JobState, QueueOpts, RESERVED_PREFIX,
|
||||||
StreamEvent, encode_payload, validate_shared_name,
|
Result, ScheduleOpts, StopToken, Store, StreamEvent, encode_payload, validate_shared_name,
|
||||||
};
|
};
|
||||||
|
|
||||||
use crate::factory::StoreFactory;
|
use crate::factory::StoreFactory;
|
||||||
@@ -1547,3 +1562,377 @@ fn unix_now() -> i64 {
|
|||||||
.map(|d| d.as_secs() as i64)
|
.map(|d| d.as_secs() as i64)
|
||||||
.unwrap_or(0)
|
.unwrap_or(0)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type RunnerJoin = tokio::task::JoinHandle<Result<()>>;
|
||||||
|
|
||||||
|
/// Spawn a `run_schedules(stop)` runner driving off one shared store
|
||||||
|
/// handle (the runner-driving posture: the row holds the observing
|
||||||
|
/// Arc, the spawned task holds its clone; leadership is storage-level,
|
||||||
|
/// not handle-level). The token clone returned cancels only this
|
||||||
|
/// runner's stop.
|
||||||
|
fn spawn_runner(store: Arc<dyn Store>) -> (StopToken, RunnerJoin) {
|
||||||
|
let stop = StopToken::new();
|
||||||
|
let owned = store.clone();
|
||||||
|
let token = stop.clone();
|
||||||
|
let handle = tokio::spawn(async move { owned.run_schedules(token).await });
|
||||||
|
(stop, handle)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Cancel and join a runner: `Ok` of the runner's own return shape
|
||||||
|
/// (clean stop `Ok(())` — ADR-019 §4 — or `LeadershipLost`), or a
|
||||||
|
/// panic if the join itself fails or the stop is not prompt.
|
||||||
|
async fn join_runner(stop: StopToken, runner: RunnerJoin) -> Result<()> {
|
||||||
|
stop.cancel();
|
||||||
|
match tokio::time::timeout(std::time::Duration::from_secs(10), runner).await {
|
||||||
|
Ok(joined) => match joined {
|
||||||
|
Ok(outcome) => outcome,
|
||||||
|
Err(e) => panic!("the runner task join-failed: {e}"),
|
||||||
|
},
|
||||||
|
Err(_) => panic!("the stopped runner must terminate promptly, not hang"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Claim everything currently claimable from `queue` (one worker:
|
||||||
|
/// "suite-probe") and append the new fire ids to `fired`. Claiming is
|
||||||
|
/// the observation surface only — the fire rows are ordinary queue
|
||||||
|
/// work.
|
||||||
|
async fn collect_fires(store: &dyn Store, queue: &str, fired: &mut Vec<i64>) {
|
||||||
|
let q = store
|
||||||
|
.queue(queue, QueueOpts::default())
|
||||||
|
.await
|
||||||
|
.expect("queue handle for the fire probe");
|
||||||
|
loop {
|
||||||
|
let page = q
|
||||||
|
.claim_batch("suite-probe", 50)
|
||||||
|
.await
|
||||||
|
.expect("the fire probe claims");
|
||||||
|
let got = page.len() as i64;
|
||||||
|
for handle in page {
|
||||||
|
let id = handle.job().id;
|
||||||
|
if !fired.contains(&id) {
|
||||||
|
fired.push(id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if got < 50 {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Poll for fires until at least `min` distinct ids have landed, or
|
||||||
|
/// panic well-labelled after `deadline` (the tolerance-bounded wait:
|
||||||
|
/// the deadline is a safety bound, the assertions are on state
|
||||||
|
/// outcomes).
|
||||||
|
async fn await_fires(
|
||||||
|
store: &dyn Store,
|
||||||
|
queue: &str,
|
||||||
|
fired: &mut Vec<i64>,
|
||||||
|
min: usize,
|
||||||
|
deadline: std::time::Duration,
|
||||||
|
label: &str,
|
||||||
|
) {
|
||||||
|
let end = std::time::Instant::now() + deadline;
|
||||||
|
loop {
|
||||||
|
collect_fires(store, queue, fired).await;
|
||||||
|
if fired.len() >= min {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
assert!(
|
||||||
|
std::time::Instant::now() < end,
|
||||||
|
"{label}: only {} boundary fires within the wait window",
|
||||||
|
fired.len()
|
||||||
|
);
|
||||||
|
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The elapsed-boundary band's upper bound for a run started at
|
||||||
|
/// `started`: boundaries elapsed since start, plus slack for the
|
||||||
|
/// unix-second boundary arithmetic and the in-progress tick at
|
||||||
|
/// cancel (the engine tests' fire-window tolerance).
|
||||||
|
fn elapsed_boundary_band(started: &std::time::Instant) -> i64 {
|
||||||
|
started.elapsed().as_secs() as i64 + 2
|
||||||
|
}
|
||||||
|
|
||||||
|
/// **Scheduler boundary fires** — the registered tick machinery, run:
|
||||||
|
/// a spawned runner fires due `@every 1s` boundaries into the named
|
||||||
|
/// queue as ordinary claimable jobs stamped with the `ScheduleOpts`
|
||||||
|
/// over the plain-queue derived defaults (`priority`/`max_attempts`
|
||||||
|
/// ride; visibility/backoff/retention take the engine defaults
|
||||||
|
/// 300/5/none), the fired rows are `get_job`-visible queue work with
|
||||||
|
/// the schedule's payload exactly, and a cancelled stop terminates the
|
||||||
|
/// runner promptly with the clean `Ok(())` shape. No double-fire per
|
||||||
|
/// boundary while one leader runs (the fire + boundary-advance commit
|
||||||
|
/// atomically — ADR-009 §4's one-firer discipline): the total fired
|
||||||
|
/// count stays inside the elapsed-boundary band above.
|
||||||
|
///
|
||||||
|
/// Contract stamp: ADR-009 §3 (boundary fires enqueue the schedule's
|
||||||
|
/// payload with `ScheduleOpts` into the named queue — ordinary queue
|
||||||
|
/// work); ADR-009 §4 (one firer per boundary, fire + advance atomic);
|
||||||
|
/// ADR-020 §3 (the stamp source: `ScheduleOpts` over the engine's
|
||||||
|
/// plain-queue derived defaults); ADR-019 §4 (the clean-stop `Ok(())`
|
||||||
|
/// return shape).
|
||||||
|
pub async fn scheduler_boundary_fires(factory: &dyn StoreFactory) {
|
||||||
|
const LABEL: &str = "scheduler_boundary_fires";
|
||||||
|
let store: Arc<dyn Store> = Arc::from(factory.open().await.expect("factory opens a store"));
|
||||||
|
let started = std::time::Instant::now();
|
||||||
|
|
||||||
|
store
|
||||||
|
.schedule(
|
||||||
|
"suite-tick",
|
||||||
|
"@every 1s",
|
||||||
|
"work",
|
||||||
|
json!({"kind": "tick"}),
|
||||||
|
ScheduleOpts {
|
||||||
|
priority: 5,
|
||||||
|
max_attempts: Some(9),
|
||||||
|
expires: None,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("schedule registers");
|
||||||
|
|
||||||
|
let (stop, runner) = spawn_runner(store.clone());
|
||||||
|
let mut fired: Vec<i64> = Vec::new();
|
||||||
|
await_fires(
|
||||||
|
store.as_ref(),
|
||||||
|
"work",
|
||||||
|
&mut fired,
|
||||||
|
2,
|
||||||
|
std::time::Duration::from_secs(25),
|
||||||
|
LABEL,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
|
||||||
|
// ADR-020 §3's stamp source: ScheduleOpts applied over the
|
||||||
|
// plain-queue derived defaults (the fire resolves engine-side; the
|
||||||
|
// queue was never opened before the fire).
|
||||||
|
let q = store
|
||||||
|
.queue("work", QueueOpts::default())
|
||||||
|
.await
|
||||||
|
.expect("queue handle");
|
||||||
|
let job = q
|
||||||
|
.get_job(fired[0])
|
||||||
|
.await
|
||||||
|
.expect("get_job reads the fire row")
|
||||||
|
.expect("the fire row exists");
|
||||||
|
assert_eq!(
|
||||||
|
job.state,
|
||||||
|
JobState::Processing,
|
||||||
|
"the fire probe claimed this row"
|
||||||
|
);
|
||||||
|
assert_eq!(job.priority, 5, "ScheduleOpts priority rides the fire");
|
||||||
|
assert_eq!(
|
||||||
|
job.max_attempts, 9,
|
||||||
|
"ScheduleOpts max_attempts overrides the default"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
job.visibility_timeout_s, 300,
|
||||||
|
"visibility stamps the engine default"
|
||||||
|
);
|
||||||
|
assert_eq!(job.backoff_base_s, 5, "backoff stamps the engine default");
|
||||||
|
assert_eq!(
|
||||||
|
job.dead_letter_retention_s, None,
|
||||||
|
"retention stamps the engine default"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
job.payload,
|
||||||
|
json!({"kind": "tick"}).to_string().into_bytes(),
|
||||||
|
"the schedule's payload is the fire's payload, exactly"
|
||||||
|
);
|
||||||
|
|
||||||
|
let outcome = join_runner(stop, runner).await;
|
||||||
|
match outcome {
|
||||||
|
Ok(()) => {}
|
||||||
|
other => panic!("{LABEL}: a cancelled stop returns Ok(()), got {other:?}"),
|
||||||
|
}
|
||||||
|
|
||||||
|
// One firer per boundary: the total fired count stays inside the
|
||||||
|
// elapsed-boundary band (a double-fire per boundary would blow
|
||||||
|
// past it).
|
||||||
|
collect_fires(store.as_ref(), "work", &mut fired).await;
|
||||||
|
let band = elapsed_boundary_band(&started);
|
||||||
|
assert!(
|
||||||
|
fired.len() >= 2,
|
||||||
|
"{LABEL}: two boundaries must have fired, got {}",
|
||||||
|
fired.len()
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
(fired.len() as i64) <= band,
|
||||||
|
"{LABEL}: no double-fire per boundary while one leader runs — \
|
||||||
|
{} fires vs the elapsed-boundary band {band}",
|
||||||
|
fired.len()
|
||||||
|
);
|
||||||
|
|
||||||
|
factory.teardown().await.expect("factory teardown");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// **Scheduler bounded catch-up** — downtime replay is at-least-once
|
||||||
|
/// per elapsed boundary, bounded by the 64-cap: register, let several
|
||||||
|
/// boundaries elapse with no runner in sight (schedules never fire
|
||||||
|
/// without one — ADR-009 §1's opt-in posture), then start one: the
|
||||||
|
/// missed boundaries replay boundary-by-boundary — strictly more fires
|
||||||
|
/// land than a coalescing single-fire would leave — the count stays
|
||||||
|
/// bounded by the 64-cap and inside the elapsed-boundary band, and the
|
||||||
|
/// runner stops clean.
|
||||||
|
///
|
||||||
|
/// The beyond-cap skip-forward leg (>64 elapsed boundaries ≈ 65+ s of
|
||||||
|
/// downtime, plus a row-backdate the public surface does not express —
|
||||||
|
/// the schedule row's `next_fire_at` is engine-internal) pins
|
||||||
|
/// engine-side, where both engines' scheduler tests manipulate
|
||||||
|
/// `next_fire_at` directly
|
||||||
|
/// (`catch_up_replays_up_to_the_cap_then_skips_forward` on each); the
|
||||||
|
/// suite pins the ≤64 catch-up leg, the suite-observable one.
|
||||||
|
///
|
||||||
|
/// Contract stamp: ADR-009 §4 (at-least-once per elapsed boundary,
|
||||||
|
/// bounded catch-up — the fixed 64-per-tick contract constant, beyond
|
||||||
|
/// which missed boundaries skip forward); ADR-019 §4 (the clean-stop
|
||||||
|
/// `Ok(())` return shape).
|
||||||
|
pub async fn scheduler_bounded_catchup(factory: &dyn StoreFactory) {
|
||||||
|
const LABEL: &str = "scheduler_bounded_catchup";
|
||||||
|
let store: Arc<dyn Store> = Arc::from(factory.open().await.expect("factory opens a store"));
|
||||||
|
let started = std::time::Instant::now();
|
||||||
|
|
||||||
|
store
|
||||||
|
.schedule(
|
||||||
|
"suite-backlog",
|
||||||
|
"@every 1s",
|
||||||
|
"catchup",
|
||||||
|
json!("w"),
|
||||||
|
Default::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("schedule registers");
|
||||||
|
|
||||||
|
// Downtime posture: no runner exists and nothing fires for it.
|
||||||
|
let mut fired: Vec<i64> = Vec::new();
|
||||||
|
collect_fires(store.as_ref(), "catchup", &mut fired).await;
|
||||||
|
assert!(
|
||||||
|
fired.is_empty(),
|
||||||
|
"schedules never fire without a runner, got {} fires",
|
||||||
|
fired.len()
|
||||||
|
);
|
||||||
|
// Several boundaries elapse (the first sits one interval past
|
||||||
|
// registration; a 4 s runner-less pause leaves at least three
|
||||||
|
// elapsed) — still nothing, no runner in sight.
|
||||||
|
past_stamp_sleep(4_000);
|
||||||
|
collect_fires(store.as_ref(), "catchup", &mut fired).await;
|
||||||
|
assert!(
|
||||||
|
fired.is_empty(),
|
||||||
|
"elapsed boundaries pile up as schedule-row state, not fires, got {}",
|
||||||
|
fired.len()
|
||||||
|
);
|
||||||
|
|
||||||
|
// Now run: the first tick replays the elapsed boundaries
|
||||||
|
// boundary-by-boundary.
|
||||||
|
let (stop, runner) = spawn_runner(store.clone());
|
||||||
|
await_fires(
|
||||||
|
store.as_ref(),
|
||||||
|
"catchup",
|
||||||
|
&mut fired,
|
||||||
|
3,
|
||||||
|
std::time::Duration::from_secs(25),
|
||||||
|
LABEL,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let outcome = join_runner(stop, runner).await;
|
||||||
|
match outcome {
|
||||||
|
Ok(()) => {}
|
||||||
|
other => panic!("{LABEL}: a cancelled stop returns Ok(()), got {other:?}"),
|
||||||
|
}
|
||||||
|
collect_fires(store.as_ref(), "catchup", &mut fired).await;
|
||||||
|
|
||||||
|
// The bounded-catch-up property: at-least-once per elapsed
|
||||||
|
// boundary (strictly more than a coalescing fire), bounded by the
|
||||||
|
// 64-cap (a multi-second downtime stays far below it).
|
||||||
|
let count = fired.len() as i64;
|
||||||
|
assert!(
|
||||||
|
count >= 3,
|
||||||
|
"{LABEL}: multi-boundary downtime must replay boundary-by-boundary \
|
||||||
|
(at-least-once per elapsed boundary), got {count} fires"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
count <= 64,
|
||||||
|
"{LABEL}: the replay is bounded by the 64-cap, got {count}"
|
||||||
|
);
|
||||||
|
let band = elapsed_boundary_band(&started);
|
||||||
|
assert!(
|
||||||
|
count <= band,
|
||||||
|
"{LABEL}: fires stay inside the elapsed-boundary band ({count} vs {band})"
|
||||||
|
);
|
||||||
|
|
||||||
|
factory.teardown().await.expect("factory teardown");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// **Scheduler leadership discipline** — two concurrent runners on one
|
||||||
|
/// store: exactly one leader, the other surfaces `LeadershipLost`
|
||||||
|
/// rather than silently co-ticking (the acquire-time refusal is the
|
||||||
|
/// same return-shape outcome as a mid-run loss), the fired count stays
|
||||||
|
/// inside the elapsed-boundary band with both runners alive (fires are
|
||||||
|
/// not duplicated), and the winning leader's cancelled stop returns
|
||||||
|
/// the clean `Ok(())` shape (ADR-019 §4; the loser's return is the
|
||||||
|
/// matchable `Err(LeadershipLost)`, never an ordinary error or a
|
||||||
|
/// silent `Ok`).
|
||||||
|
///
|
||||||
|
/// Contract stamp: ADR-009 §1 (one leader at a time; leadership loss
|
||||||
|
/// ends the runner); ADR-009 §6 (`LeadershipLost` — the matchable
|
||||||
|
/// return-shape variant of the runner's stop); ADR-019 §4 (the
|
||||||
|
/// clean-stop `Ok(())` return shape under a cancelled token).
|
||||||
|
pub async fn scheduler_leadership_discipline(factory: &dyn StoreFactory) {
|
||||||
|
const LABEL: &str = "scheduler_leadership_discipline";
|
||||||
|
let store: Arc<dyn Store> = Arc::from(factory.open().await.expect("factory opens a store"));
|
||||||
|
let started = std::time::Instant::now();
|
||||||
|
|
||||||
|
store
|
||||||
|
.schedule(
|
||||||
|
"suite-lead",
|
||||||
|
"@every 1s",
|
||||||
|
"single",
|
||||||
|
json!(1),
|
||||||
|
Default::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("schedule registers");
|
||||||
|
|
||||||
|
// Two concurrent runners on the same store: the leadership lock is
|
||||||
|
// storage-level, so exactly one is granted and drives the tick.
|
||||||
|
let (stop_a, runner_a) = spawn_runner(store.clone());
|
||||||
|
let (stop_b, runner_b) = spawn_runner(store.clone());
|
||||||
|
let mut fired: Vec<i64> = Vec::new();
|
||||||
|
await_fires(
|
||||||
|
store.as_ref(),
|
||||||
|
"single",
|
||||||
|
&mut fired,
|
||||||
|
1,
|
||||||
|
std::time::Duration::from_secs(25),
|
||||||
|
LABEL,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
// A few more boundaries with both runners alive.
|
||||||
|
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
|
||||||
|
|
||||||
|
let outcome_a = join_runner(stop_a, runner_a).await;
|
||||||
|
let outcome_b = join_runner(stop_b, runner_b).await;
|
||||||
|
|
||||||
|
let clean = |o: &Result<()>| matches!(o, Ok(()));
|
||||||
|
let lost = |o: &Result<()>| matches!(o, Err(Error::LeadershipLost));
|
||||||
|
assert!(
|
||||||
|
(clean(&outcome_a) && lost(&outcome_b)) || (lost(&outcome_a) && clean(&outcome_b)),
|
||||||
|
"{LABEL}: exactly one leader — one clean stop and one \
|
||||||
|
`LeadershipLost`, got A {outcome_a:?} and B {outcome_b:?}"
|
||||||
|
);
|
||||||
|
|
||||||
|
// One firer per boundary while both runners were alive: the fired
|
||||||
|
// count stays inside the elapsed-boundary band (duplicated fires
|
||||||
|
// would blow past it).
|
||||||
|
collect_fires(store.as_ref(), "single", &mut fired).await;
|
||||||
|
let count = fired.len() as i64;
|
||||||
|
let band = elapsed_boundary_band(&started);
|
||||||
|
assert!(
|
||||||
|
count >= 1 && count <= band,
|
||||||
|
"{LABEL}: fires not duplicated under two runners — {count} fires \
|
||||||
|
vs the elapsed-boundary band {band}"
|
||||||
|
);
|
||||||
|
|
||||||
|
factory.teardown().await.expect("factory teardown");
|
||||||
|
}
|
||||||
@@ -167,7 +167,9 @@ mod rows {
|
|||||||
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
||||||
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
||||||
payload_round_trip_stores_exact_encoding, payload_too_large_produced_on_pg,
|
payload_round_trip_stores_exact_encoding, payload_too_large_produced_on_pg,
|
||||||
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, sweep_no_stranded_rows,
|
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms,
|
||||||
|
scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline,
|
||||||
|
sweep_no_stranded_rows,
|
||||||
};
|
};
|
||||||
|
|
||||||
use super::PgFactory;
|
use super::PgFactory;
|
||||||
@@ -248,6 +250,21 @@ mod rows {
|
|||||||
sweep_no_stranded_rows,
|
sweep_no_stranded_rows,
|
||||||
"row-sweep-stranded"
|
"row-sweep-stranded"
|
||||||
);
|
);
|
||||||
|
harness_row!(
|
||||||
|
row_scheduler_boundary_fires,
|
||||||
|
scheduler_boundary_fires,
|
||||||
|
"row-scheduler-boundary"
|
||||||
|
);
|
||||||
|
harness_row!(
|
||||||
|
row_scheduler_bounded_catchup,
|
||||||
|
scheduler_bounded_catchup,
|
||||||
|
"row-scheduler-catchup"
|
||||||
|
);
|
||||||
|
harness_row!(
|
||||||
|
row_scheduler_leadership_discipline,
|
||||||
|
scheduler_leadership_discipline,
|
||||||
|
"row-scheduler-leadership"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
mod factory_shape {
|
mod factory_shape {
|
||||||
|
|||||||
@@ -71,7 +71,9 @@ mod rows {
|
|||||||
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
||||||
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
||||||
payload_round_trip_stores_exact_encoding, payload_too_large_never_produced_on_sqlite,
|
payload_round_trip_stores_exact_encoding, payload_too_large_never_produced_on_sqlite,
|
||||||
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, sweep_no_stranded_rows,
|
queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms,
|
||||||
|
scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline,
|
||||||
|
sweep_no_stranded_rows,
|
||||||
};
|
};
|
||||||
|
|
||||||
use super::SqliteFactory;
|
use super::SqliteFactory;
|
||||||
@@ -139,6 +141,21 @@ mod rows {
|
|||||||
async fn row_sweep_no_stranded_rows() {
|
async fn row_sweep_no_stranded_rows() {
|
||||||
sweep_no_stranded_rows(&factory("row-sweep-stranded")).await;
|
sweep_no_stranded_rows(&factory("row-sweep-stranded")).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
async fn row_scheduler_boundary_fires() {
|
||||||
|
scheduler_boundary_fires(&factory("row-scheduler-boundary")).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
async fn row_scheduler_bounded_catchup() {
|
||||||
|
scheduler_bounded_catchup(&factory("row-scheduler-catchup")).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
async fn row_scheduler_leadership_discipline() {
|
||||||
|
scheduler_leadership_discipline(&factory("row-scheduler-leadership")).await;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
mod factory_shape {
|
mod factory_shape {
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
---
|
---
|
||||||
id: suite-scheduler-rows
|
id: suite-scheduler-rows
|
||||||
name: Contract-suite rows — scheduler boundary/catch-up/leadership (+ determinism-posture extension)
|
name: Contract-suite rows — scheduler boundary/catch-up/leadership (+ determinism-posture extension)
|
||||||
status: pending
|
status: completed
|
||||||
depends_on: []
|
depends_on: []
|
||||||
scope: moderate
|
scope: moderate
|
||||||
risk: medium
|
risk: medium
|
||||||
@@ -86,8 +86,93 @@ latency-posture in Notes and leave the engines as they are.
|
|||||||
|
|
||||||
## Notes
|
## Notes
|
||||||
|
|
||||||
> To be filled by implementation agent
|
- **Beyond-cap skip-forward leg: pinned engine-side, already done in
|
||||||
|
both engines.** The leg needs >64 elapsed boundaries (~65+ s of
|
||||||
|
downtime) and a direct schedule-row backdate the public surface does
|
||||||
|
not express (`next_fire_at` is engine-internal). Disposition: pinned
|
||||||
|
where `next_fire_at` is directly manipulable —
|
||||||
|
`catch_up_replays_up_to_the_cap_then_skips_forward` exists in
|
||||||
|
`alkstore-sqlite/src/store/scheduler_tests.rs` (substrate row
|
||||||
|
backdate, exactly 64 fires, skip-forward past now, skipped
|
||||||
|
boundaries never enqueue) and the twin in
|
||||||
|
`alkstore-postgres/src/store/scheduler_tests.rs` (admin-SQL backdate,
|
||||||
|
same shape). The suite row (`scheduler_bounded_catchup`) pins the
|
||||||
|
≤64 in-band catch-up leg — the suite-observable one — and states the
|
||||||
|
engine-side disposition in its doc text.
|
||||||
|
- **Fire-wake parity: not demanded by this row set; engines left
|
||||||
|
as-is.** The suite rows observe fire outcomes through ordinary
|
||||||
|
claim polling, never through wake latency, so the row shapes demand
|
||||||
|
no wake. The recorded cross-engine latency-parity gap stands (the pg
|
||||||
|
tick's in-tx enqueue fires no `pg_notify`; SQLite's watcher wakes on
|
||||||
|
those commits — `tasks/review-wave-4-fixes.md`'s recorded-for-later
|
||||||
|
note): correctness is covered by the queues row's pinned re-poll
|
||||||
|
safety net on both engines; the row set here pins state outcomes
|
||||||
|
with tolerances bounded well above the boundary resolution, so wake
|
||||||
|
latency is out of scope. If a later task demands the parity, the
|
||||||
|
one-call `wake_tx` inside the pg tick's tx is the shape to ride.
|
||||||
|
- **Runner-driving posture details** (the extension's riders, beyond
|
||||||
|
what the description pinned): the spawned runners share the row's
|
||||||
|
single store handle via `Arc<dyn Store>` clones — leadership is
|
||||||
|
storage-level, not handle-level, so two concurrent runners on one
|
||||||
|
handle still contend on the leadership lock (the SQLite engine's own
|
||||||
|
tests open a second handle on the same path; the suite's factory
|
||||||
|
gives one isolated store per row, and the shared-handle shape
|
||||||
|
exercises the same storage-level discipline). A leadership-lock
|
||||||
|
acquire through the writer slot is a short lease on both engines, so
|
||||||
|
the shared handle neither deadlocks nor serializes the tick away.
|
||||||
|
The two-concurrent-runners row takes whichever runner wins (no
|
||||||
|
identity assumption): exactly one `Ok(())` and one
|
||||||
|
`Err(LeadershipLost)` pair, checked by value.
|
||||||
|
- **Band tolerances** (the determinism posture's "bounded well above
|
||||||
|
the boundary resolution"): every fired-count assertion is bounded by
|
||||||
|
an elapsed-boundary band measured from just before registration
|
||||||
|
(`unix-second boundary arithmetic` slack + one in-progress tick at
|
||||||
|
cancel, +2 s — the engine tests' fire-window tolerance), never a
|
||||||
|
timing value. The catch-up row's lower bound (≥ 3 fires) is the
|
||||||
|
boundary-by-boundary pin — strictly more than a coalescing
|
||||||
|
single-fire would leave; its registration's first-fire grace of one
|
||||||
|
interval is counted in the band.
|
||||||
|
|
||||||
## Summary
|
## Summary
|
||||||
|
|
||||||
> To be filled on completion
|
**Landed:** the determinism-posture extension in
|
||||||
|
`alkstore-contract-suite/src/properties.rs`'s module doc (runner-driving
|
||||||
|
rows admitted — spawned `run_schedules(stop)` tasks with a core
|
||||||
|
`StopToken`, state outcomes only, elapsed-boundary-band tolerances,
|
||||||
|
never timing-value assertions, never two-task race windows; ADR-017 §2
|
||||||
|
class 4 posture note, cross-connection wake orchestration still
|
||||||
|
engine-test-side) and the three scheduler rows, version-stamped and
|
||||||
|
re-exported from the suite lib:
|
||||||
|
|
||||||
|
- `scheduler_boundary_fires` — stamps ADR-009 §3/§4, ADR-020 §3,
|
||||||
|
ADR-019 §4; spawned runner, ≥2 boundary fires observed via ordinary
|
||||||
|
claim polling, the fire row's stamp source pinned (ScheduleOpts
|
||||||
|
priority/max_attempts over the plain-queue derived defaults
|
||||||
|
300/5/none, payload exact), clean-stop `Ok(())` prompt, fired count
|
||||||
|
inside the elapsed-boundary band (no double-fire per boundary under
|
||||||
|
one leader).
|
||||||
|
- `scheduler_bounded_catchup` — stamp ADR-009 §4 (plus ADR-019 §4);
|
||||||
|
runner-less downtime proves schedules never fire without a runner,
|
||||||
|
a 4 s pause yields ≥3 replayed fires (boundary-by-boundary, not
|
||||||
|
coalesced) bounded below the 64-cap and inside the band; the
|
||||||
|
beyond-cap skip-forward disposition recorded in Notes (pinned
|
||||||
|
engine-side on both engines, already implemented there).
|
||||||
|
- `scheduler_leadership_discipline` — stamps ADR-009 §1/§6,
|
||||||
|
ADR-019 §4; two spawned runners on one store, exactly one
|
||||||
|
`Ok(())`/`Err(LeadershipLost)` pair by value, fired count inside the
|
||||||
|
elapsed-boundary band with both alive (no duplicated fires).
|
||||||
|
|
||||||
|
Wired into `alkstore-sqlite/tests/contract_suite.rs` (three direct
|
||||||
|
tests) and `alkstore-postgres/tests/contract_suite.rs` (three
|
||||||
|
`harness_row!` entries). Suite crate gained `tokio` (rt + time) as a
|
||||||
|
regular dependency for the runner-driving rows (the unused
|
||||||
|
dev-dependency dropped).
|
||||||
|
|
||||||
|
**Gates:** `cargo build`; `cargo test -p alkstore-sqlite` (189 lib +
|
||||||
|
16 suite rows green); `cargo test -p alkstore-postgres` serverless
|
||||||
|
green (pg rows skip cleanly); the full pg column green against the
|
||||||
|
harness server (`pglo-poc` :15432 — 121 lib incl. the previously
|
||||||
|
skipped rows + 16 suite rows + 9 schema); the three new rows re-run
|
||||||
|
twice solo on each engine (determinism check, green both columns both
|
||||||
|
runs); `cargo clippy --all-targets -- -D warnings`; `cargo fmt --check`
|
||||||
|
(clean after fmt).
|
||||||
Reference in new issue
Block a user