453 lines
16 KiB
Rust
453 lines
16 KiB
Rust
//! The transactional seam's SQLite arm — `begin_tx`, the
|
|
//! [`SqliteTxHandle`] (the writer-slot lease, ADR-007), and all eleven
|
|
//! `*_tx` methods (`sqlite-engine-seam-tx`).
|
|
//!
|
|
//! `begin_tx` acquires the writer slot and opens `BEGIN IMMEDIATE` on
|
|
//! the acquired connection; the handle owns connection + slot. Every
|
|
//! op round-trips the `spawn_blocking` seam to that connection —
|
|
//! holding the handle across `await` points parks the writer, which
|
|
//! is the design (WAL single-writer, honest).
|
|
//!
|
|
//! Drop = rollback (ADR-021 §4): dropping the handle without commit
|
|
//! issues `ROLLBACK` (ignoring failure — no panic may escape drop) and
|
|
//! releases the writer slot; a panic unwinding through any scope
|
|
//! holding the handle reaches the same rollback. The no-ghosts
|
|
//! property holds through the RAII path.
|
|
//!
|
|
//! Errors: substrate errors map to
|
|
//! [`Error::Database`](alkstore::Error::Database) with the source chain
|
|
//! preserved (ADR-008 §5); `encode_payload`'s typed `Codec` error
|
|
//! propagates from the enqueue/publish paths (ADR-023 §1). Every
|
|
//! name-bearing method validates at the entry point, before any round
|
|
//! trip (ADR-008 §4) — the commit-atomic path cannot bypass it.
|
|
|
|
use serde_json::Value;
|
|
use std::sync::Arc;
|
|
|
|
use alkstore::{BoxedFuture, EnqueueOpts, Error, Job, Result, StreamEvent, TxHandle};
|
|
use alkstore::{validate_local_name, validate_shared_name};
|
|
|
|
use crate::queue::job_from_json;
|
|
use crate::resolution::{
|
|
outbox_backing_queue_name, outbox_default_stamps, plain_queue_default_stamps,
|
|
resolve_enqueue_opts, stamps_with_override,
|
|
};
|
|
use crate::seam::{blocking, sqlite_error};
|
|
use crate::stream::stream_events_from_json_page;
|
|
use crate::substrate::{
|
|
Writer, enqueue, get_job, stream_get_offset, stream_publish, stream_read_since,
|
|
stream_save_offset,
|
|
};
|
|
|
|
/// The caller-held SQLite transaction handle: the writer-slot lease.
|
|
/// Constructed only through [`crate::SqliteStore`]'s `begin_tx` —
|
|
/// never held across an engine `Store` trait object (the trait is what
|
|
/// consumers see).
|
|
pub struct SqliteTxHandle {
|
|
conn: Option<rusqlite::Connection>,
|
|
writer: Arc<Writer>,
|
|
reopen: Arc<dyn Fn() -> Result<rusqlite::Connection> + Send + Sync>,
|
|
/// This store's commit-fault arm (`sqlite-commit-error-arm`) —
|
|
/// read once per commit by the test-only fault branch; dead in
|
|
/// production builds.
|
|
#[cfg(test)]
|
|
commit_fault: crate::seam::commit_fault::Flag,
|
|
}
|
|
|
|
impl std::fmt::Debug for SqliteTxHandle {
|
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
f.debug_struct("SqliteTxHandle")
|
|
.field("open", &self.conn.is_some())
|
|
.finish_non_exhaustive()
|
|
}
|
|
}
|
|
|
|
impl SqliteTxHandle {
|
|
pub(crate) fn begin(
|
|
writer: Arc<Writer>,
|
|
reopen: Arc<dyn Fn() -> Result<rusqlite::Connection> + Send + Sync>,
|
|
#[cfg(test)] commit_fault: crate::seam::commit_fault::Flag,
|
|
) -> BoxedFuture<'static, Result<Box<dyn TxHandle + Send>>> {
|
|
Box::pin(async move {
|
|
let writer2 = writer.clone();
|
|
let conn = blocking(move || {
|
|
let conn = writer2.acquire().ok_or_else(|| {
|
|
alkstore::Error::database(std::io::Error::other(
|
|
"the store is closed: writer slot unavailable",
|
|
))
|
|
})?;
|
|
match conn.execute_batch("BEGIN IMMEDIATE").map_err(sqlite_error) {
|
|
Ok(()) => Ok(conn),
|
|
Err(e) => {
|
|
writer2.release(conn);
|
|
Err(e)
|
|
}
|
|
}
|
|
})
|
|
.await?;
|
|
Ok(Box::new(SqliteTxHandle {
|
|
conn: Some(conn),
|
|
writer,
|
|
reopen,
|
|
#[cfg(test)]
|
|
commit_fault,
|
|
}) as _)
|
|
})
|
|
}
|
|
|
|
/// The private bridge: every op routes through the same connection
|
|
/// via `spawn_blocking`. The connection crosses the seam inside an
|
|
/// RAII [`ConnLease`]: a future cancelled mid-op drops the lease,
|
|
/// which closes the connection (SQLite rolls back the uncommitted
|
|
/// transaction) and replenishes the writer slot with a fresh
|
|
/// connection — the lease cannot leak, and the slot stays
|
|
/// grantable. If `commit`/drop already consumed the handle (`conn`
|
|
/// taken), the op fails closed with `Closed` — one-shot handles
|
|
/// cannot op past their terminal transition.
|
|
fn with_conn<T, F>(&mut self, f: F) -> BoxedFuture<'_, Result<T>>
|
|
where
|
|
T: Send + 'static,
|
|
F: FnOnce(&rusqlite::Connection) -> Result<T> + Send + 'static,
|
|
{
|
|
let Some(conn) = self.conn.take() else {
|
|
return Box::pin(async { Err(Error::Closed) });
|
|
};
|
|
let writer = self.writer.clone();
|
|
let reopen = self.reopen.clone();
|
|
let fut = blocking(move || {
|
|
let lease = ConnLease {
|
|
conn: Some(conn),
|
|
writer: writer.clone(),
|
|
reopen: reopen.clone(),
|
|
};
|
|
let out = match lease.conn() {
|
|
Some(conn) => f(conn),
|
|
None => Err(Error::database(std::io::Error::other(
|
|
"the store is closed: writer slot unavailable",
|
|
))),
|
|
};
|
|
Ok((out, lease))
|
|
});
|
|
Box::pin(async move {
|
|
let (out, lease) = fut.await?;
|
|
self.conn = lease.release();
|
|
out
|
|
})
|
|
}
|
|
}
|
|
|
|
/// The writer connection inside the `with_conn` bridge: an RAII guard
|
|
/// so the connection crossing `spawn_blocking` always has an owner
|
|
/// that accounts for the writer slot. Handled normally, the future's
|
|
/// completion defuses the guard (`release`) and the connection returns
|
|
/// to the handle; on a mid-op-future cancellation the guard drops
|
|
/// instead — dropping the connection (its uncommitted transaction
|
|
/// rolls back at the SQLite layer) and replenishing the writer slot
|
|
/// with a fresh connection, so the next lease proceeds.
|
|
struct ConnLease {
|
|
conn: Option<rusqlite::Connection>,
|
|
writer: std::sync::Arc<Writer>,
|
|
reopen: Arc<dyn Fn() -> Result<rusqlite::Connection> + Send + Sync>,
|
|
}
|
|
|
|
impl ConnLease {
|
|
fn conn(&self) -> Option<&rusqlite::Connection> {
|
|
self.conn.as_ref()
|
|
}
|
|
|
|
fn release(mut self) -> Option<rusqlite::Connection> {
|
|
self.conn.take()
|
|
}
|
|
}
|
|
|
|
impl Drop for ConnLease {
|
|
fn drop(&mut self) {
|
|
if let Some(conn) = self.conn.take() {
|
|
if let Err(rollback) = conn.execute_batch("ROLLBACK") {
|
|
eprintln!(
|
|
"alkstore: error: rollback on a cancelled transaction op \
|
|
left the connection in an unknown state: {rollback}"
|
|
);
|
|
}
|
|
drop(conn);
|
|
if let Ok(fresh) = (self.reopen)() {
|
|
self.writer.release(fresh);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
impl TxHandle for SqliteTxHandle {
|
|
fn enqueue_tx<'a>(
|
|
&'a mut self,
|
|
queue: &str,
|
|
opts: EnqueueOpts,
|
|
payload: Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
if let Err(e) = validate_shared_name(queue) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let queue = queue.to_string();
|
|
self.with_conn(move |conn| {
|
|
let bytes = alkstore::encode_payload(&payload)?;
|
|
let resolved = resolve_enqueue_opts(conn, &opts)?;
|
|
let stamps = stamps_with_override(plain_queue_default_stamps(), &opts);
|
|
enqueue(
|
|
conn,
|
|
&queue,
|
|
std::str::from_utf8(&bytes)
|
|
.map_err(|e| Error::Codec(format!("row payload must be utf-8: {e}")))?,
|
|
resolved.run_at,
|
|
opts.priority,
|
|
resolved.expires_at,
|
|
stamps,
|
|
)
|
|
.map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn publish_tx<'a>(&'a mut self, stream: &str, payload: Value) -> BoxedFuture<'a, Result<i64>> {
|
|
let stream = stream.to_string();
|
|
Box::pin(async move { self.publish_with_key_tx(&stream, None, payload).await })
|
|
}
|
|
|
|
fn publish_with_key_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
key: Option<String>,
|
|
payload: Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
if let Err(e) = validate_shared_name(stream) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let key = match key {
|
|
Some(k) if k.is_empty() || k.trim().is_empty() => {
|
|
return Box::pin(async move { Err(Error::InvalidName { name: k }) });
|
|
}
|
|
other => other,
|
|
};
|
|
let stream = stream.to_string();
|
|
self.with_conn(move |conn| {
|
|
let bytes = alkstore::encode_payload(&payload)?;
|
|
stream_publish(
|
|
conn,
|
|
&stream,
|
|
key.as_deref(),
|
|
std::str::from_utf8(&bytes)
|
|
.map_err(|e| Error::Codec(format!("row payload must be utf-8: {e}")))?,
|
|
)
|
|
.map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn notify_tx<'a>(&'a mut self, channel: &str, payload: Value) -> BoxedFuture<'a, Result<()>> {
|
|
if let Err(e) = validate_shared_name(channel) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let channel = channel.to_string();
|
|
self.with_conn(move |conn| {
|
|
let bytes = alkstore::encode_payload(&payload)?;
|
|
conn.query_row(
|
|
"SELECT notify(?1, ?2)",
|
|
rusqlite::params![
|
|
channel,
|
|
String::from_utf8(bytes).map_err(|e| Error::Codec(format!(
|
|
"notify payload must be utf-8 text: {e}"
|
|
)))?,
|
|
],
|
|
|_| Ok(()),
|
|
)
|
|
.map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn save_offset_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
offset: i64,
|
|
) -> BoxedFuture<'a, Result<()>> {
|
|
if let Err(e) = validate_shared_name(stream) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
if let Err(e) = validate_local_name(consumer) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let stream = stream.to_string();
|
|
let consumer = consumer.to_string();
|
|
self.with_conn(move |conn| {
|
|
stream_save_offset(conn, &consumer, &stream, offset)
|
|
.map(|_| ())
|
|
.map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn get_job_tx<'a>(
|
|
&'a mut self,
|
|
queue: &str,
|
|
job_id: i64,
|
|
) -> BoxedFuture<'a, Result<Option<Job>>> {
|
|
if let Err(e) = validate_shared_name(queue) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let queue = queue.to_string();
|
|
self.with_conn(move |conn| {
|
|
let row = get_job(conn, job_id).map_err(sqlite_error)?;
|
|
let Some(json) = row else {
|
|
return Ok(None);
|
|
};
|
|
let value: serde_json::Value = serde_json::from_str(&json)
|
|
.map_err(|e| Error::Codec(format!("job row must be json: {e}")))?;
|
|
if value["queue"].as_str() != Some(queue.as_str()) {
|
|
return Ok(None);
|
|
}
|
|
job_from_json(&value)
|
|
})
|
|
}
|
|
|
|
fn get_offset_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
if let Err(e) = validate_shared_name(stream) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
if let Err(e) = validate_local_name(consumer) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let stream = stream.to_string();
|
|
let consumer = consumer.to_string();
|
|
self.with_conn(move |conn| {
|
|
stream_get_offset(conn, &consumer, &stream).map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn read_since_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
offset: i64,
|
|
limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
if let Err(e) = validate_shared_name(stream) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let stream = stream.to_string();
|
|
self.with_conn(move |conn| {
|
|
if limit <= 0 {
|
|
return Ok(Vec::new());
|
|
}
|
|
let page = stream_read_since(conn, &stream, offset, limit).map_err(sqlite_error)?;
|
|
stream_events_from_json_page(&page, stream)
|
|
})
|
|
}
|
|
|
|
fn read_from_consumer_tx<'a>(
|
|
&'a mut self,
|
|
stream: &str,
|
|
consumer: &str,
|
|
limit: i64,
|
|
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
|
if let Err(e) = validate_shared_name(stream) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
if let Err(e) = validate_local_name(consumer) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let stream = stream.to_string();
|
|
let consumer = consumer.to_string();
|
|
self.with_conn(move |conn| {
|
|
if limit <= 0 {
|
|
return Ok(Vec::new());
|
|
}
|
|
let from = stream_get_offset(conn, &consumer, &stream).map_err(sqlite_error)?;
|
|
let page = stream_read_since(conn, &stream, from, limit).map_err(sqlite_error)?;
|
|
stream_events_from_json_page(&page, stream)
|
|
})
|
|
}
|
|
|
|
fn outbox_enqueue_tx<'a>(
|
|
&'a mut self,
|
|
outbox: &str,
|
|
opts: EnqueueOpts,
|
|
payload: Value,
|
|
) -> BoxedFuture<'a, Result<i64>> {
|
|
if let Err(e) = validate_shared_name(outbox) {
|
|
return Box::pin(async move { Err(e) });
|
|
}
|
|
let outbox = outbox.to_string();
|
|
self.with_conn(move |conn| {
|
|
let bytes = alkstore::encode_payload(&payload)?;
|
|
let resolved = resolve_enqueue_opts(conn, &opts)?;
|
|
let stamps = stamps_with_override(outbox_default_stamps(), &opts);
|
|
enqueue(
|
|
conn,
|
|
&outbox_backing_queue_name(&outbox),
|
|
std::str::from_utf8(&bytes)
|
|
.map_err(|e| Error::Codec(format!("row payload must be utf-8: {e}")))?,
|
|
resolved.run_at,
|
|
opts.priority,
|
|
resolved.expires_at,
|
|
stamps,
|
|
)
|
|
.map_err(sqlite_error)
|
|
})
|
|
}
|
|
|
|
fn commit(mut self: Box<Self>) -> BoxedFuture<'static, Result<()>> {
|
|
let Some(conn) = self.conn.take() else {
|
|
return Box::pin(async { Err(Error::Closed) });
|
|
};
|
|
let writer = self.writer.clone();
|
|
let writer_reopen = self.reopen.clone();
|
|
#[cfg(test)]
|
|
let commit_fault = self.commit_fault.clone();
|
|
Box::pin(async move {
|
|
blocking(move || {
|
|
#[cfg(test)]
|
|
let committed = if crate::seam::commit_fault::take(&commit_fault) {
|
|
Err(crate::seam::commit_fault::failure())
|
|
} else {
|
|
conn.execute_batch("COMMIT")
|
|
};
|
|
#[cfg(not(test))]
|
|
let committed = conn.execute_batch("COMMIT");
|
|
match committed.map_err(sqlite_error) {
|
|
Ok(()) => {
|
|
writer.release(conn);
|
|
Ok(())
|
|
}
|
|
Err(e) => {
|
|
drop(conn);
|
|
if let Ok(fresh) = (writer_reopen)() {
|
|
writer.release(fresh);
|
|
}
|
|
Err(e)
|
|
}
|
|
}
|
|
})
|
|
.await
|
|
})
|
|
}
|
|
}
|
|
|
|
impl Drop for SqliteTxHandle {
|
|
fn drop(&mut self) {
|
|
if let Some(conn) = self.conn.take() {
|
|
match conn.execute_batch("ROLLBACK") {
|
|
Ok(()) => self.writer.release(conn),
|
|
Err(rollback) => {
|
|
eprintln!(
|
|
"alkstore: error: rollback on a dropped transaction handle \
|
|
left the connection in an unknown state: {rollback}"
|
|
);
|
|
drop(conn);
|
|
if let Ok(fresh) = (self.reopen)() {
|
|
self.writer.release(fresh);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|