Files

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);
}