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]
|
||||
alkstore = { version = "0.1", path = "../alkstore" }
|
||||
serde_json = "1"
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { version = "1", features = ["macros", "rt"] }
|
||||
tokio = { version = "1", features = ["rt", "time"] }
|
||||
@@ -42,5 +42,6 @@ pub use properties::{
|
||||
in_tx_reads_see_own_writes, job_handle_validity_predicate,
|
||||
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,
|
||||
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
|
||||
//! family-wide no-panic rule here only).
|
||||
//!
|
||||
//! Determinism: no wall-clock-timing assertions and no two-task race
|
||||
//! windows — orchestration that would need one (runner spawning,
|
||||
//! cross-connection wakes) stays engine-test-side; rows keep to
|
||||
//! single-store sequential drives and the factory's isolation
|
||||
//! guarantee. Rows exercise a *close/dispose* arm by dropping the
|
||||
//! Determinism: no wall-clock-timing-value assertions and no two-task
|
||||
//! race windows; rows keep to single-store sequential drives and the
|
||||
//! factory's isolation guarantee.
|
||||
//!
|
||||
//! 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
|
||||
//! drop path closes receivers terminally); teardown-then-fresh-reopen
|
||||
//! rows would be factory-contract tests (ADR-022), not contract rows,
|
||||
//! so none exist.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use serde_json::json;
|
||||
|
||||
use alkstore::{
|
||||
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, JobState, RESERVED_PREFIX, Result,
|
||||
StreamEvent, encode_payload, validate_shared_name,
|
||||
BoxedFuture, Delivery, EnqueueOpts, Error, JobHandle, JobState, QueueOpts, RESERVED_PREFIX,
|
||||
Result, ScheduleOpts, StopToken, Store, StreamEvent, encode_payload, validate_shared_name,
|
||||
};
|
||||
|
||||
use crate::factory::StoreFactory;
|
||||
@@ -1547,3 +1562,377 @@ fn unix_now() -> i64 {
|
||||
.map(|d| d.as_secs() as i64)
|
||||
.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,
|
||||
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
||||
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;
|
||||
@@ -248,6 +250,21 @@ mod rows {
|
||||
sweep_no_stranded_rows,
|
||||
"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 {
|
||||
|
||||
@@ -71,7 +71,9 @@ mod rows {
|
||||
enqueue_opts_resolution, extent_clamp_semantics, in_tx_reads_see_own_writes,
|
||||
job_handle_validity_predicate, name_validation_rejects_empty_and_reserved,
|
||||
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;
|
||||
@@ -139,6 +141,21 @@ mod rows {
|
||||
async fn row_sweep_no_stranded_rows() {
|
||||
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 {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: suite-scheduler-rows
|
||||
name: Contract-suite rows — scheduler boundary/catch-up/leadership (+ determinism-posture extension)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: []
|
||||
scope: moderate
|
||||
risk: medium
|
||||
@@ -86,8 +86,93 @@ latency-posture in Notes and leave the engines as they are.
|
||||
|
||||
## 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
|
||||
|
||||
> 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