use gwk_domain::checkpoint::{CHECKPOINT_EVENT_INTERVAL, CHECKPOINT_INTERVAL_SECS};
use gwk_domain::envelope::{Actor, EventEnvelope, INLINE_PAYLOAD_MAX_BYTES, Origin, PayloadRef};
use gwk_domain::ids::{
AggregateId, CorrelationId, EventId, FenceToken, IdempotencyKey, ProjectId, Seq, Timestamp,
};
use gwk_domain::port::{AppendError, EventStore, MAX_READ_LIMIT, StorageError};
use secrecy::{ExposeSecret, SecretString};
use sqlx::postgres::PgPoolOptions;
use sqlx::{Executor, PgConnection, PgPool, Row, postgres::PgRow};
use tokio::sync::{Semaphore, SemaphorePermit};
use crate::blob::store::PgBlobStore;
use crate::checkpoint;
use crate::error::Result as KernelResult;
use crate::numeric::{from_numeric_text, to_numeric_text};
use crate::writer::claim_epoch;
pub const MAX_INFLIGHT_APPENDS: usize = 64;
pub const EVENT_CHANNEL: &str = "gwk_event";
const NO_FENCE_PRESENTED: u64 = 0;
macro_rules! event_columns {
() => {
"seq::text AS seq_text, event_id, project_id, aggregate_type, \
aggregate_id, aggregate_version, event_type, schema_version, \
to_json(occurred_at) #>> '{}' AS occurred_at, \
to_json(appended_at) #>> '{}' AS appended_at, \
actor, origin, causation_id, correlation_id, idempotency_key, payload, payload_ref"
};
}
pub async fn connect_pool(
database_url: &SecretString,
max_connections: u32,
) -> KernelResult<PgPool> {
let pool = PgPoolOptions::new()
.max_connections(max_connections)
.after_connect(|conn, _meta| {
Box::pin(async move {
conn.execute("SET TIME ZONE 'UTC'").await?;
Ok(())
})
})
.connect(database_url.expose_secret())
.await?;
Ok(pool)
}
pub struct PgEventStore {
pool: PgPool,
admission: Semaphore,
boot_epoch: i64,
blobs: Option<PgBlobStore>,
}
impl PgEventStore {
pub async fn open(pool: PgPool) -> KernelResult<Self> {
Self::with_capacity(pool, MAX_INFLIGHT_APPENDS).await
}
pub async fn with_capacity(pool: PgPool, inflight: usize) -> KernelResult<Self> {
let boot_epoch = claim_epoch(&pool).await?;
Ok(Self {
pool,
admission: Semaphore::new(inflight),
boot_epoch,
blobs: None,
})
}
pub fn open_reader(pool: PgPool) -> Self {
Self {
pool,
admission: Semaphore::new(MAX_INFLIGHT_APPENDS),
boot_epoch: 0,
blobs: None,
}
}
#[must_use]
pub fn with_blobs(mut self, blobs: PgBlobStore) -> Self {
self.blobs = Some(blobs);
self
}
pub fn blobs(&self) -> Option<&PgBlobStore> {
self.blobs.as_ref()
}
pub(crate) async fn checkpoint_if_due(
&self,
tx: &mut PgConnection,
writer: &WriterState,
through: Seq,
created_at: &Timestamp,
) -> Result<(), AppendError> {
let Some(blobs) = &self.blobs else {
return Ok(());
};
if !writer.checkpoint_due(through.value()) {
return Ok(());
}
checkpoint::snapshot(tx, blobs, through, created_at)
.await
.map(|_| ())
.map_err(|e| AppendError::Storage(format!("checkpoint: {e}")))
}
pub async fn checkpoint_at_watermark(&self) -> Result<Option<Seq>, AppendError> {
let Some(blobs) = &self.blobs else {
return Ok(None);
};
let mut tx = self
.pool
.begin()
.await
.map_err(|e| append_storage("begin a shutdown checkpoint", e))?;
let _writer = self.lock_writer(&mut tx).await?;
let at: (Option<String>, String) =
sqlx::query_as("SELECT max(seq)::text, to_json(now()) #>> '{}' FROM gwk.event")
.fetch_one(&mut *tx)
.await
.map_err(|e| append_storage("read the watermark", e))?;
let Some(watermark) = at.0 else {
return Ok(None);
};
let through =
Seq::new(from_numeric_text(&watermark).map_err(|e| append_storage("watermark", e))?);
checkpoint::snapshot(&mut tx, blobs, through, &Timestamp::new(at.1))
.await
.map_err(|e| AppendError::Storage(format!("shutdown checkpoint: {e}")))?;
tx.commit()
.await
.map_err(|e| append_storage("commit a shutdown checkpoint", e))?;
Ok(Some(through))
}
pub fn boot_epoch(&self) -> i64 {
self.boot_epoch
}
pub fn pool(&self) -> &PgPool {
&self.pool
}
pub(crate) fn admit(&self) -> Result<SemaphorePermit<'_>, AppendError> {
self.admission.try_acquire().map_err(|_| {
AppendError::Storage(format!(
"append queue is full ({MAX_INFLIGHT_APPENDS} in flight)"
))
})
}
pub(crate) async fn lock_writer(
&self,
tx: &mut PgConnection,
) -> Result<WriterState, AppendError> {
let writer = sqlx::query(
"SELECT epoch, fence_token::text AS fence_text, next_seq::text AS next_seq_text, \
checkpoint_seq::text AS checkpoint_seq_text, \
now() - checkpoint_at >= make_interval(secs => $1::double precision) \
AS checkpoint_overdue \
FROM gwk_internal.writer WHERE id = 1 FOR UPDATE",
)
.bind(CHECKPOINT_INTERVAL_SECS as f64)
.fetch_optional(&mut *tx)
.await
.map_err(|e| append_storage("lock writer row", e))?
.ok_or_else(|| {
AppendError::Storage(
"gwk_internal.writer has no singleton row: database not initialized".to_owned(),
)
})?;
let epoch: i64 = writer
.try_get("epoch")
.map_err(|e| append_storage("column epoch", e))?;
if epoch != self.boot_epoch {
return Err(AppendError::Storage(format!(
"writer epoch {} was superseded by {epoch}: another kernel took write authority",
self.boot_epoch
)));
}
let current_fence: Option<u64> = writer
.try_get::<Option<String>, _>("fence_text")
.map_err(|e| append_storage("column fence_token", e))?
.map(|text| from_numeric_text(&text))
.transpose()
.map_err(|e| append_storage("column fence_token", e))?;
let next_seq = from_numeric_text(
&writer
.try_get::<String, _>("next_seq_text")
.map_err(|e| append_storage("column next_seq", e))?,
)
.map_err(|e| append_storage("column next_seq", e))?;
let checkpoint_seq = from_numeric_text(
&writer
.try_get::<String, _>("checkpoint_seq_text")
.map_err(|e| append_storage("column checkpoint_seq", e))?,
)
.map_err(|e| append_storage("column checkpoint_seq", e))?;
let checkpoint_overdue: bool = writer
.try_get("checkpoint_overdue")
.map_err(|e| append_storage("column checkpoint_at", e))?;
Ok(WriterState {
next_seq,
current_fence,
checkpoint_seq,
checkpoint_overdue,
})
}
pub(crate) async fn append_locked(
&self,
tx: &mut PgConnection,
writer: &WriterState,
expected_version: u32,
fence: Option<FenceToken>,
events: &[EventEnvelope],
) -> Result<Appended, AppendError> {
validate_batch(expected_version, events)?;
let first = &events[0];
let aggregate_type = first.aggregate_type.clone();
let aggregate_id = first.aggregate_id.as_str().to_owned();
if let Some(current) = writer.current_fence {
let presented = fence.map(|f| f.value()).unwrap_or(NO_FENCE_PRESENTED);
if presented != current {
return Err(AppendError::Fenced {
presented: FenceToken::new(presented),
current: FenceToken::new(current),
});
}
}
let keys: Vec<String> = events
.iter()
.filter_map(|e| e.idempotency_key.as_ref())
.map(|k| k.as_str().to_owned())
.collect();
if !keys.is_empty() {
let stored = sqlx::query(concat!(
"SELECT ",
event_columns!(),
" FROM gwk.event \
WHERE aggregate_type = $1 AND aggregate_id = $2 AND idempotency_key = ANY($3) \
ORDER BY seq"
))
.bind(&aggregate_type)
.bind(&aggregate_id)
.bind(&keys)
.fetch_all(&mut *tx)
.await
.map_err(|e| append_storage("idempotency lookup", e))?;
if !stored.is_empty() {
if stored.len() == events.len() && keys.len() == events.len() {
let stored = stored
.iter()
.map(row_to_envelope)
.collect::<Result<Vec<_>, _>>()
.map_err(|e| AppendError::Storage(e.0))?;
if let Some((stored, presented)) = stored
.iter()
.zip(events)
.find(|(stored, presented)| !is_replay_of(stored, presented))
{
return Err(AppendError::MalformedBatch(format!(
"idempotency key {:?} already named {}/{} version {} in project {}: \
a retry must present the identical batch",
presented
.idempotency_key
.as_ref()
.map(|k| k.as_str())
.unwrap_or_default(),
stored.aggregate_type,
stored.aggregate_id,
stored.aggregate_version,
stored.project_id,
)));
}
return Ok(Appended {
events: stored,
replayed: true,
});
}
return Err(AppendError::MalformedBatch(format!(
"{} of {} idempotency keys already landed for {aggregate_type}/{aggregate_id}: \
a retry must present the identical batch",
stored.len(),
events.len()
)));
}
}
let actual = current_aggregate_version(&mut *tx, &aggregate_type, &aggregate_id)
.await
.map_err(|e| AppendError::Storage(e.0))?;
if actual != expected_version {
return Err(AppendError::VersionConflict {
actual,
expected: expected_version,
});
}
let mut appended = Vec::with_capacity(events.len());
for (offset, event) in events.iter().enumerate() {
let seq = writer
.next_seq
.checked_add(offset as u64)
.ok_or_else(|| AppendError::Storage("global sequence exhausted".to_owned()))?;
let row = sqlx::query(concat!(
"INSERT INTO gwk.event ( \
seq, event_id, project_id, aggregate_type, aggregate_id, aggregate_version, \
event_type, schema_version, occurred_at, appended_at, actor, origin, \
causation_id, correlation_id, idempotency_key, payload, payload_ref) \
VALUES ($1::numeric, $2, $3, $4, $5, $6, $7, $8, $9::timestamptz, now(), \
$10, $11, $12, $13, $14, $15, $16) \
RETURNING ",
event_columns!()
))
.bind(to_numeric_text(seq))
.bind(event.event_id.as_str())
.bind(event.project_id.as_str())
.bind(&event.aggregate_type)
.bind(event.aggregate_id.as_str())
.bind(i64::from(event.aggregate_version))
.bind(&event.event_type)
.bind(i64::from(event.schema_version))
.bind(event.occurred_at.as_str())
.bind(serde_json::to_value(&event.actor).map_err(|e| append_storage("actor", e))?)
.bind(serde_json::to_value(&event.origin).map_err(|e| append_storage("origin", e))?)
.bind(event.causation_id.as_ref().map(|v| v.as_str()))
.bind(event.correlation_id.as_ref().map(|v| v.as_str()))
.bind(event.idempotency_key.as_ref().map(|v| v.as_str()))
.bind(&event.payload)
.bind(
event
.payload_ref
.as_ref()
.map(serde_json::to_value)
.transpose()
.map_err(|e| append_storage("payload_ref", e))?,
)
.fetch_one(&mut *tx)
.await
.map_err(|e| append_storage("insert event", e))?;
appended.push(row_to_envelope(&row).map_err(|e| AppendError::Storage(e.0))?);
}
let advanced = writer
.next_seq
.checked_add(events.len() as u64)
.ok_or_else(|| AppendError::Storage("global sequence exhausted".to_owned()))?;
sqlx::query("UPDATE gwk_internal.writer SET next_seq = $1::numeric WHERE id = 1")
.bind(to_numeric_text(advanced))
.execute(&mut *tx)
.await
.map_err(|e| append_storage("advance sequence", e))?;
sqlx::query("SELECT pg_notify($1, $2)")
.bind(EVENT_CHANNEL)
.bind(to_numeric_text(advanced - 1))
.execute(&mut *tx)
.await
.map_err(|e| append_storage("notify", e))?;
Ok(Appended {
events: appended,
replayed: false,
})
}
}
pub(crate) struct WriterState {
pub(crate) next_seq: u64,
pub(crate) current_fence: Option<u64>,
pub(crate) checkpoint_seq: u64,
pub(crate) checkpoint_overdue: bool,
}
impl WriterState {
fn checkpoint_due(&self, through: u64) -> bool {
self.checkpoint_overdue
|| through.saturating_sub(self.checkpoint_seq) >= CHECKPOINT_EVENT_INTERVAL
}
}
pub(crate) struct Appended {
pub(crate) events: Vec<EventEnvelope>,
pub(crate) replayed: bool,
}
pub(crate) async fn events_for_key(
conn: &mut PgConnection,
project_id: &str,
aggregate_type: &str,
aggregate_id: &str,
key: &str,
) -> Result<Vec<EventEnvelope>, StorageError> {
let rows = sqlx::query(concat!(
"SELECT ",
event_columns!(),
" FROM gwk.event \
WHERE idempotency_key = $4 \
AND (project_id = $1 OR (aggregate_type = $2 AND aggregate_id = $3)) \
ORDER BY seq"
))
.bind(project_id)
.bind(aggregate_type)
.bind(aggregate_id)
.bind(key)
.fetch_all(conn)
.await
.map_err(|e| storage("idempotency lookup", e))?;
rows.iter().map(row_to_envelope).collect()
}
pub(crate) async fn current_aggregate_version(
conn: &mut PgConnection,
aggregate_type: &str,
aggregate_id: &str,
) -> Result<u32, StorageError> {
let actual: Option<i64> = sqlx::query_scalar(
"SELECT max(aggregate_version) FROM gwk.event \
WHERE aggregate_type = $1 AND aggregate_id = $2",
)
.bind(aggregate_type)
.bind(aggregate_id)
.fetch_one(conn)
.await
.map_err(|e| storage("read aggregate version", e))?;
u32::try_from(actual.unwrap_or(0)).map_err(|e| storage("aggregate_version out of range", e))
}
pub fn validate_batch(expected_version: u32, events: &[EventEnvelope]) -> Result<(), AppendError> {
let malformed = |reason: String| Err(AppendError::MalformedBatch(reason));
let Some(first) = events.first() else {
return malformed("empty batch".to_owned());
};
for (offset, event) in events.iter().enumerate() {
if event.aggregate_type != first.aggregate_type || event.aggregate_id != first.aggregate_id
{
return malformed(format!(
"batch mixes aggregates: {}/{} and {}/{}",
first.aggregate_type,
first.aggregate_id.as_str(),
event.aggregate_type,
event.aggregate_id.as_str()
));
}
let want = u64::from(expected_version) + 1 + offset as u64;
if want > u64::from(u32::MAX) {
return malformed(format!("aggregate_version would exceed {}", u32::MAX));
}
if u64::from(event.aggregate_version) != want {
return malformed(format!(
"aggregate_version {} at offset {offset}: expected {want} \
(contiguous from expected_version + 1)",
event.aggregate_version
));
}
if event.schema_version == 0 {
return malformed(format!(
"schema_version 0 at offset {offset} is not a version"
));
}
let inline = serde_json::to_vec(&event.payload)
.map(|bytes| bytes.len())
.unwrap_or(usize::MAX);
if inline > INLINE_PAYLOAD_MAX_BYTES {
return malformed(format!(
"inline payload at offset {offset} is {inline} bytes, over the \
{INLINE_PAYLOAD_MAX_BYTES} bound — use payload_ref"
));
}
if let Some(key) = &event.idempotency_key
&& events[..offset]
.iter()
.any(|earlier| earlier.idempotency_key.as_ref() == Some(key))
{
return malformed(format!(
"idempotency key {:?} appears twice in one batch",
key.as_str()
));
}
}
Ok(())
}
fn storage(context: &str, error: impl std::fmt::Display) -> StorageError {
StorageError(format!("{context}: {error}"))
}
fn append_storage(context: &str, error: impl std::fmt::Display) -> AppendError {
AppendError::Storage(format!("{context}: {error}"))
}
pub(crate) async fn read_page<'e, E>(
executor: E,
cursor: Option<Seq>,
limit: usize,
) -> Result<Vec<EventEnvelope>, StorageError>
where
E: sqlx::PgExecutor<'e>,
{
let limit = limit.min(MAX_READ_LIMIT) as i64;
let rows = sqlx::query(concat!(
"SELECT ",
event_columns!(),
" FROM gwk.event \
WHERE $1::numeric IS NULL OR seq > $1::numeric \
ORDER BY seq LIMIT $2"
))
.bind(cursor.map(|c| to_numeric_text(c.value())))
.bind(limit)
.fetch_all(executor)
.await
.map_err(|e| storage("read_from", e))?;
rows.iter().map(row_to_envelope).collect()
}
fn row_to_envelope(row: &PgRow) -> Result<EventEnvelope, StorageError> {
let get = |name: &str| -> Result<String, StorageError> {
row.try_get::<String, _>(name)
.map_err(|e| storage(&format!("column {name}"), e))
};
let opt = |name: &str| -> Result<Option<String>, StorageError> {
row.try_get::<Option<String>, _>(name)
.map_err(|e| storage(&format!("column {name}"), e))
};
let json = |name: &str| -> Result<serde_json::Value, StorageError> {
row.try_get::<serde_json::Value, _>(name)
.map_err(|e| storage(&format!("column {name}"), e))
};
let seq_text = get("seq_text")?;
let seq = from_numeric_text(&seq_text).map_err(|e| storage("column seq", e))?;
let aggregate_version: i64 = row
.try_get("aggregate_version")
.map_err(|e| storage("column aggregate_version", e))?;
let schema_version: i64 = row
.try_get("schema_version")
.map_err(|e| storage("column schema_version", e))?;
let payload_ref: Option<serde_json::Value> = row
.try_get("payload_ref")
.map_err(|e| storage("column payload_ref", e))?;
Ok(EventEnvelope {
event_id: EventId::new(get("event_id")?),
project_id: ProjectId::new(get("project_id")?),
aggregate_type: get("aggregate_type")?,
aggregate_id: AggregateId::new(get("aggregate_id")?),
aggregate_version: u32::try_from(aggregate_version)
.map_err(|e| storage("aggregate_version out of range", e))?,
event_type: get("event_type")?,
schema_version: u32::try_from(schema_version)
.map_err(|e| storage("schema_version out of range", e))?,
global_sequence: Seq::new(seq),
occurred_at: Timestamp::new(get("occurred_at")?),
appended_at: Timestamp::new(get("appended_at")?),
actor: serde_json::from_value::<Actor>(json("actor")?)
.map_err(|e| storage("column actor", e))?,
origin: serde_json::from_value::<Origin>(json("origin")?)
.map_err(|e| storage("column origin", e))?,
causation_id: opt("causation_id")?.map(EventId::new),
correlation_id: opt("correlation_id")?.map(CorrelationId::new),
idempotency_key: opt("idempotency_key")?.map(IdempotencyKey::new),
payload: json("payload")?,
payload_ref: payload_ref
.map(serde_json::from_value::<PayloadRef>)
.transpose()
.map_err(|e| storage("column payload_ref", e))?,
})
}
fn is_replay_of(stored: &EventEnvelope, presented: &EventEnvelope) -> bool {
let mut normalized = presented.clone();
normalized.global_sequence = stored.global_sequence;
normalized.appended_at = stored.appended_at.clone();
normalized.occurred_at = stored.occurred_at.clone();
normalized == *stored
}
impl EventStore for PgEventStore {
async fn append(
&self,
expected_version: u32,
fence: Option<FenceToken>,
events: Vec<EventEnvelope>,
) -> Result<Vec<EventEnvelope>, AppendError> {
let _permit = self.admit()?;
let mut tx = self
.pool
.begin()
.await
.map_err(|e| append_storage("begin", e))?;
let writer = self.lock_writer(&mut tx).await?;
let appended = self
.append_locked(&mut tx, &writer, expected_version, fence, &events)
.await?;
tx.commit().await.map_err(|e| append_storage("commit", e))?;
Ok(appended.events)
}
async fn read_from(
&self,
cursor: Option<Seq>,
limit: usize,
) -> Result<Vec<EventEnvelope>, StorageError> {
read_page(&self.pool, cursor, limit).await
}
async fn watermark(&self) -> Result<Option<Seq>, StorageError> {
let text: Option<String> = sqlx::query_scalar("SELECT max(seq)::text FROM gwk.event")
.fetch_one(&self.pool)
.await
.map_err(|e| storage("watermark", e))?;
text.map(|t| from_numeric_text(&t))
.transpose()
.map(|opt| opt.map(Seq::new))
.map_err(|e| storage("watermark", e))
}
async fn grant_fence(&self) -> Result<FenceToken, StorageError> {
let text: String = sqlx::query_scalar(
"UPDATE gwk_internal.writer \
SET fence_token = coalesce(fence_token, 0) + 1 \
WHERE id = 1 RETURNING fence_token::text",
)
.fetch_optional(&self.pool)
.await
.map_err(|e| storage("grant_fence", e))?
.ok_or_else(|| StorageError("gwk_internal.writer has no singleton row".to_owned()))?;
from_numeric_text(&text)
.map(FenceToken::new)
.map_err(|e| storage("grant_fence", e))
}
}
#[cfg(test)]
mod tests {
use gwk_domain::ids::Timestamp;
use super::*;
fn event(aggregate_id: &str, version: u32) -> EventEnvelope {
EventEnvelope {
event_id: EventId::new(format!("evt-{aggregate_id}-{version}")),
project_id: ProjectId::new("p"),
aggregate_type: "task".into(),
aggregate_id: AggregateId::new(aggregate_id),
aggregate_version: version,
event_type: "tick".into(),
schema_version: 1,
global_sequence: Seq::new(0),
occurred_at: Timestamp::new("2026-01-01T00:00:00Z"),
appended_at: Timestamp::new("2026-01-01T00:00:00Z"),
actor: Actor {
kind: "kernel".into(),
id: None,
},
origin: Origin {
system: "kernel".into(),
r#ref: None,
},
causation_id: None,
correlation_id: None,
idempotency_key: None,
payload: serde_json::json!({}),
payload_ref: None,
}
}
#[test]
fn a_batch_is_one_aggregate_contiguous_from_the_expected_version() {
validate_batch(0, &[event("a", 1), event("a", 2)]).expect("contiguous from 1");
validate_batch(7, &[event("a", 8)]).expect("contiguous from expected + 1");
for (expected, batch, why) in [
(0u32, vec![], "empty"),
(0, vec![event("a", 2)], "starts past expected + 1"),
(0, vec![event("a", 1), event("a", 3)], "gap"),
(0, vec![event("a", 1), event("a", 1)], "repeat"),
(5, vec![event("a", 1)], "ignores expected_version"),
] {
let err = validate_batch(expected, &batch).expect_err(why);
assert!(
matches!(err, AppendError::MalformedBatch(_)),
"{why}: {err:?}"
);
}
let mixed = vec![event("a", 1), event("b", 2)];
let err = validate_batch(0, &mixed).expect_err("mixed aggregates");
let AppendError::MalformedBatch(reason) = err else {
panic!("expected MalformedBatch");
};
assert!(reason.contains("mixes aggregates"), "{reason}");
}
#[test]
fn a_batch_may_not_run_past_the_aggregate_version_ceiling() {
let mut last = event("a", u32::MAX);
last.aggregate_version = u32::MAX;
validate_batch(u32::MAX - 1, &[last.clone()]).expect("the final version is reachable");
let err = validate_batch(u32::MAX, &[last]).expect_err("past the ceiling");
let AppendError::MalformedBatch(reason) = err else {
panic!("expected MalformedBatch");
};
assert!(reason.contains("would exceed"), "{reason}");
}
#[test]
fn the_column_carries_every_schema_version_the_envelope_can_declare() {
for version in [1, i32::MAX as u32, i32::MAX as u32 + 1, u32::MAX] {
let mut e = event("a", 1);
e.schema_version = version;
validate_batch(0, &[e]).unwrap_or_else(|err| panic!("schema_version {version}: {err}"));
}
let mut zero = event("a", 1);
zero.schema_version = 0;
let err = validate_batch(0, &[zero]).expect_err("0 is not a version");
let AppendError::MalformedBatch(reason) = err else {
panic!("expected MalformedBatch");
};
assert!(reason.contains("schema_version"), "{reason}");
}
#[test]
fn one_key_may_not_appear_twice_in_a_batch() {
let key = |k: &str| Some(gwk_domain::ids::IdempotencyKey::new(k));
let mut a = event("a", 1);
let mut b = event("a", 2);
a.idempotency_key = key("same");
b.idempotency_key = key("same");
let err = validate_batch(0, &[a.clone(), b.clone()]).expect_err("repeated key");
let AppendError::MalformedBatch(reason) = err else {
panic!("expected MalformedBatch");
};
assert!(reason.contains("twice in one batch"), "{reason}");
b.idempotency_key = key("other");
validate_batch(0, &[a.clone(), b.clone()]).expect("distinct keys");
b.idempotency_key = None;
validate_batch(0, &[a, b]).expect("a mix of keyed and unkeyed");
}
#[test]
fn an_oversized_inline_payload_is_refused_before_it_reaches_the_log() {
let mut big = event("a", 1);
big.payload = serde_json::json!({ "blob": "x".repeat(INLINE_PAYLOAD_MAX_BYTES) });
let err = validate_batch(0, &[big]).expect_err("oversized payload");
let AppendError::MalformedBatch(reason) = err else {
panic!("expected MalformedBatch");
};
assert!(reason.contains("payload_ref"), "{reason}");
}
}