378 lines
12 KiB
Rust
378 lines
12 KiB
Rust
//! The harness: runs the suite's property rows against a supplied
|
|
//! [`StoreFactory`]. This test target proves the harness end-to-end
|
|
//! with a trivial in-crate mock store — the acceptance criterion's
|
|
//! "exemplar property green against a trivial in-crate mock" — and is
|
|
//! the shape each engine crate's dev-dependency wiring copies.
|
|
|
|
use std::time::Duration;
|
|
|
|
use serde_json::json;
|
|
|
|
use alkstore::{
|
|
BoxedFuture, Delivery, EnqueueOpts, Error, EventReceiver, Job, JobHandle, Lock, Outbox, Queue,
|
|
QueueOpts, Result, Schedule, StopToken, Store, StreamEvent, StreamHandle, TxHandle, Wake,
|
|
WakeReceiver,
|
|
};
|
|
use alkstore_contract_suite::StoreFactory;
|
|
use alkstore_contract_suite::properties::name_validation_rejects_empty_and_reserved;
|
|
|
|
/// Trivial mock `TxHandle` — validation is the *engine's* job in real
|
|
/// engines; the mock leans on the shared helpers exactly as an engine
|
|
/// would, at each entry point.
|
|
struct MockTx;
|
|
|
|
impl TxHandle for MockTx {
|
|
fn enqueue_tx<'a>(
|
|
&'a mut self,
|
|
queue: &str,
|
|
_opts: EnqueueOpts,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_shared_name(queue).map(|_| 1);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn publish_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_shared_name(stream).map(|_| 1);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn publish_with_key_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
key: Option<String>,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_shared_name(stream).and_then(|_| match key {
|
|
Some(k) if k.trim().is_empty() => Err(Error::InvalidName { name: k }),
|
|
_ => Ok(1),
|
|
});
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn notify_tx<'a>(
|
|
&'a mut self,
|
|
channel: &str,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<()>> {
|
|
Box::pin(std::future::ready(alkstore::validate_shared_name(channel)))
|
|
}
|
|
fn save_offset_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
_offset: i64,
|
|
) -> BoxedFuture<'a, Result<()>> {
|
|
let r = alkstore::validate_shared_name(stream)
|
|
.and_then(|_| alkstore::validate_local_name(consumer));
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn get_job_tx<'a>(
|
|
&'a mut self,
|
|
queue: &str,
|
|
_job_id: i64,
|
|
) -> BoxedFuture<'a, Result<Option<Job>>> {
|
|
let r = alkstore::validate_shared_name(queue).map(|_| None);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn get_offset_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_shared_name(stream)
|
|
.and_then(|_| alkstore::validate_local_name(consumer))
|
|
.map(|_| 0);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn read_since_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
_offset: i64,
|
|
_limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
let r = alkstore::validate_shared_name(stream).map(|_| Vec::new());
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn read_from_consumer_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
_limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
let r = alkstore::validate_shared_name(stream)
|
|
.and_then(|_| alkstore::validate_local_name(consumer))
|
|
.map(|_| Vec::new());
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn outbox_enqueue_tx<'a>(
|
|
&'a mut self,
|
|
outbox: &str,
|
|
_opts: EnqueueOpts,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_shared_name(outbox).map(|_| 1);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn commit(self: Box<Self>) -> BoxedFuture<'static, Result<()>> {
|
|
Box::pin(std::future::ready(Ok(())))
|
|
}
|
|
}
|
|
|
|
struct MockQueue;
|
|
|
|
impl Queue for MockQueue {
|
|
fn name(&self) -> &str {
|
|
"mock"
|
|
}
|
|
fn enqueue<'a>(
|
|
&'a self,
|
|
_payload: serde_json::Value,
|
|
_opts: EnqueueOpts,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(1)))
|
|
}
|
|
fn claim_one<'a>(
|
|
&'a self,
|
|
worker_id: &str,
|
|
) -> BoxedFuture<'a, Result<Option<Box<dyn JobHandle>>>> {
|
|
let r = alkstore::validate_local_name(worker_id).map(|_| None);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn claim_batch<'a>(
|
|
&'a self,
|
|
worker_id: &str,
|
|
_n: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<Box<dyn JobHandle>>>> {
|
|
let r = alkstore::validate_local_name(worker_id).map(|_| Vec::new());
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn ack_batch<'a>(&'a self, _ids: &[i64]) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(0)))
|
|
}
|
|
fn cancel<'a>(&'a self, _job_id: i64) -> BoxedFuture<'a, Result<bool>> {
|
|
Box::pin(std::future::ready(Ok(false)))
|
|
}
|
|
fn get_job<'a>(&'a self, _job_id: i64) -> BoxedFuture<'a, Result<Option<Job>>> {
|
|
Box::pin(std::future::ready(Ok(None)))
|
|
}
|
|
fn sweep_expired<'a>(&'a self) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(0)))
|
|
}
|
|
}
|
|
|
|
struct MockStream;
|
|
|
|
impl StreamHandle for MockStream {
|
|
fn name(&self) -> &str {
|
|
"mock"
|
|
}
|
|
fn publish<'a>(&'a self, _payload: serde_json::Value) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(1)))
|
|
}
|
|
fn publish_with_key<'a>(
|
|
&'a self,
|
|
_key: Option<String>,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(1)))
|
|
}
|
|
fn read_since<'a>(
|
|
&'a self,
|
|
_offset: i64,
|
|
_limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
Box::pin(std::future::ready(Ok(Vec::new())))
|
|
}
|
|
fn read_from_consumer<'a>(
|
|
&'a self,
|
|
consumer: &str,
|
|
_limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
let r = alkstore::validate_local_name(consumer).map(|_| Vec::new());
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn save_offset<'a>(&'a self, consumer: &str, _offset: i64) -> BoxedFuture<'a, Result<()>> {
|
|
Box::pin(std::future::ready(alkstore::validate_local_name(consumer)))
|
|
}
|
|
fn get_offset<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result<i64>> {
|
|
let r = alkstore::validate_local_name(consumer).map(|_| 0);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn trim_to<'a>(&'a self, _horizon: i64) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(0)))
|
|
}
|
|
fn subscribe<'a>(&'a self, _consumer: &str) -> BoxedFuture<'a, Result<Box<dyn EventReceiver>>> {
|
|
Box::pin(std::future::ready(Err(Error::Database("mock".into()))))
|
|
}
|
|
}
|
|
|
|
struct MockWakeReceiver;
|
|
|
|
impl WakeReceiver for MockWakeReceiver {
|
|
fn recv<'a>(&'a mut self) -> BoxedFuture<'a, Option<Wake>> {
|
|
Box::pin(std::future::ready(None))
|
|
}
|
|
fn try_recv(&mut self) -> Result<Option<Wake>> {
|
|
Ok(None)
|
|
}
|
|
fn recv_timeout<'a>(&'a mut self, _timeout: Duration) -> BoxedFuture<'a, Result<Option<Wake>>> {
|
|
Box::pin(std::future::ready(Ok(None)))
|
|
}
|
|
}
|
|
|
|
/// `Schedule` is `#[non_exhaustive]` (ADR-017 §3) — the mock cannot
|
|
/// struct-construct it, so the happy-path `schedule()` spot-check is
|
|
/// delegated to core's own tests; the mock's `schedule()` returns a
|
|
/// `Database` stub and the row only asserts its *validation* outcomes
|
|
/// (which precede any construction).
|
|
struct MockStore;
|
|
|
|
impl Store for MockStore {
|
|
fn begin_tx(&self) -> BoxedFuture<'_, Result<Box<dyn TxHandle + Send>>> {
|
|
Box::pin(std::future::ready(Ok(
|
|
Box::new(MockTx) as Box<dyn TxHandle + Send>
|
|
)))
|
|
}
|
|
fn notify<'a>(
|
|
&'a self,
|
|
channel: &str,
|
|
_payload: serde_json::Value,
|
|
) -> BoxedFuture<'a, Result<()>> {
|
|
Box::pin(std::future::ready(alkstore::validate_shared_name(channel)))
|
|
}
|
|
fn listen<'a>(&'a self, channel: &str) -> BoxedFuture<'a, Result<Box<dyn WakeReceiver>>> {
|
|
let r = alkstore::validate_shared_name(channel)
|
|
.map(|_| Box::new(MockWakeReceiver) as Box<dyn WakeReceiver>);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn stream<'a>(&'a self, name: &str) -> BoxedFuture<'a, Result<Box<dyn StreamHandle>>> {
|
|
let r = alkstore::validate_shared_name(name)
|
|
.map(|_| Box::new(MockStream) as Box<dyn StreamHandle>);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn queue<'a>(
|
|
&'a self,
|
|
name: &str,
|
|
_opts: QueueOpts,
|
|
) -> BoxedFuture<'a, Result<Box<dyn Queue>>> {
|
|
let r = alkstore::validate_shared_name(name).map(|_| Box::new(MockQueue) as Box<dyn Queue>);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn outbox<'a>(&'a self, name: &str) -> BoxedFuture<'a, Result<Box<dyn Outbox>>> {
|
|
let r =
|
|
alkstore::validate_shared_name(name).map(|_| Box::new(MockOutbox) as Box<dyn Outbox>);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn try_lock<'a>(
|
|
&'a self,
|
|
name: &str,
|
|
owner: &str,
|
|
_ttl: i64,
|
|
) -> BoxedFuture<'a, Result<Option<Box<dyn Lock>>>> {
|
|
let r = alkstore::validate_shared_name(name)
|
|
.and_then(|_| alkstore::validate_local_name(owner))
|
|
.map(|_| None);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn schedule<'a>(
|
|
&'a self,
|
|
name: &str,
|
|
_spec: &str,
|
|
queue: &str,
|
|
_payload: serde_json::Value,
|
|
_opts: alkstore::ScheduleOpts,
|
|
) -> BoxedFuture<'a, Result<Schedule>> {
|
|
let r = alkstore::validate_shared_name(name)
|
|
.and_then(|_| alkstore::validate_shared_name(queue))
|
|
.and_then(|_| {
|
|
Err(Error::Database(
|
|
"mock: schedule read-back not exercised".into(),
|
|
))
|
|
});
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn unschedule<'a>(&'a self, name: &str) -> BoxedFuture<'a, Result<bool>> {
|
|
let r = alkstore::validate_shared_name(name).map(|_| false);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
fn run_schedules<'a>(&'a self, _stop: StopToken) -> BoxedFuture<'a, Result<()>> {
|
|
Box::pin(std::future::ready(Ok(())))
|
|
}
|
|
}
|
|
|
|
struct MockOutbox;
|
|
|
|
impl Outbox for MockOutbox {
|
|
fn name(&self) -> &str {
|
|
"mock"
|
|
}
|
|
fn enqueue<'a>(
|
|
&'a self,
|
|
_payload: serde_json::Value,
|
|
_opts: EnqueueOpts,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
Box::pin(std::future::ready(Ok(1)))
|
|
}
|
|
fn run_once<'a>(
|
|
&'a mut self,
|
|
worker_id: &str,
|
|
_delivery: &'a mut dyn Delivery,
|
|
) -> BoxedFuture<'a, Result<bool>> {
|
|
let r = alkstore::validate_local_name(worker_id).map(|_| false);
|
|
Box::pin(std::future::ready(r))
|
|
}
|
|
}
|
|
|
|
/// The mock factory: yields a fresh store over an isolated (in-memory,
|
|
/// stateless) backing store. Isolation is trivially true — the mock
|
|
/// holds no rows.
|
|
struct MockFactory;
|
|
|
|
impl StoreFactory for MockFactory {
|
|
fn open(&self) -> BoxedFuture<'_, Result<Box<dyn Store>>> {
|
|
Box::pin(std::future::ready(
|
|
Ok(Box::new(MockStore) as Box<dyn Store>),
|
|
))
|
|
}
|
|
fn teardown(&self) -> BoxedFuture<'_, Result<()>> {
|
|
Box::pin(std::future::ready(Ok(())))
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn exemplar_row_green_against_mock_factory() {
|
|
let factory = MockFactory;
|
|
name_validation_rejects_empty_and_reserved(&factory).await;
|
|
}
|
|
|
|
/// Happy-path spot-check behind the exemplar row: valid names construct
|
|
/// handles — proving the row's rejections are the validation's doing
|
|
/// (the mock's `schedule` happy path is a `Database` stub because
|
|
/// `Schedule` is `#[non_exhaustive]`; its validation outcomes are what
|
|
/// the row asserts).
|
|
#[tokio::test]
|
|
async fn mock_store_happy_paths_construct() {
|
|
let factory = MockFactory;
|
|
let store = factory.open().await.unwrap();
|
|
|
|
store.stream("work").await.expect("stream handle");
|
|
store
|
|
.queue("work", QueueOpts::default())
|
|
.await
|
|
.expect("queue handle");
|
|
store.outbox("work").await.expect("outbox handle");
|
|
|
|
factory.teardown().await.expect("teardown");
|
|
}
|
|
|
|
#[test]
|
|
fn stamp_marker_and_json_shape() {
|
|
assert_eq!(
|
|
alkstore_contract_suite::version_stamp::STAMP_MARKER,
|
|
"Contract stamp:"
|
|
);
|
|
let _ = json!(null);
|
|
}
|