Skip to main content

radixdb_executor/
application.rs

1//! System-owned application relations built on the ordinary MVCC path.
2//!
3//! Audit and outbox records deliberately use normal catalog tables and the
4//! caller's storage transaction.  This module owns only their typed contract
5//! and lifecycle; it does not introduce another log, commit marker, or
6//! recovery authority.
7
8use chrono::{DateTime, Duration, Utc};
9pub use radixdb_catalog::ObjectId;
10use radixdb_catalog::{CatalogGeneration, CatalogPayload, ConstraintPayload, ObjectKind};
11use radixdb_core::{DataType, Error, Result, Row, Value};
12pub use radixdb_procedural::{AuditEvent, OutboxMessage};
13use radixdb_storage::traits::{QueryResult, Table};
14
15use crate::context::ExecutionContext;
16use crate::procedural::transaction_visible_catalog;
17use crate::Executor;
18
19pub const AUDIT_RELATION_NAME: &str = "audit.event";
20pub const OUTBOX_RELATION_NAME: &str = "outbox.message";
21
22const MAX_AUDIT_METADATA_ENTRIES: usize = 32;
23const MAX_AUDIT_METADATA_KEY_BYTES: usize = 128;
24const MAX_AUDIT_METADATA_BYTES: usize = 16 * 1024;
25const MAX_OUTBOX_IDEMPOTENCY_KEY_BYTES: usize = 512;
26const MAX_OUTBOX_PAYLOAD_BYTES: usize = 1024 * 1024;
27const MAX_WORKER_ID_BYTES: usize = 256;
28const MAX_OUTBOX_ERROR_BYTES: usize = 4096;
29const MAX_CLAIM_BATCH: usize = 1024;
30const MAX_CLAIM_SCAN_ROWS: usize = 65_536;
31const MAX_OUTBOX_ATTEMPTS: u32 = 1_000;
32const MAX_RETENTION_BATCH: usize = 10_000;
33
34const AUDIT_COLUMNS: &[(&str, DataType, bool)] = &[
35    ("event_id", DataType::Uuid, false),
36    ("occurred_at", DataType::Timestamp, false),
37    ("transaction_id", DataType::Integer, false),
38    ("session_principal", DataType::Uuid, false),
39    ("effective_principal", DataType::Uuid, false),
40    ("object_id", DataType::Uuid, false),
41    ("command_fingerprint", DataType::Bytes, false),
42    ("outcome", DataType::Text, false),
43    ("metadata", DataType::Json, false),
44];
45
46const OUTBOX_COLUMNS: &[(&str, DataType, bool)] = &[
47    ("message_id", DataType::Uuid, false),
48    ("idempotency_key", DataType::Text, false),
49    ("schema_version", DataType::Integer, false),
50    ("payload", DataType::Json, false),
51    ("state", DataType::Text, false),
52    ("created_at", DataType::Timestamp, false),
53    ("available_at", DataType::Timestamp, false),
54    ("lease_owner", DataType::Text, true),
55    ("lease_token", DataType::Uuid, true),
56    ("lease_expires_at", DataType::Timestamp, true),
57    ("attempt_count", DataType::Integer, false),
58    ("completed_at", DataType::Timestamp, true),
59    ("last_error", DataType::Text, true),
60    ("dead_lettered_at", DataType::Timestamp, true),
61];
62
63const INSTALL_SQL: &str = r#"
64BEGIN;
65CREATE TABLE IF NOT EXISTS "audit.event" (
66    event_id UUID PRIMARY KEY,
67    occurred_at TIMESTAMP NOT NULL,
68    transaction_id INTEGER NOT NULL,
69    session_principal UUID NOT NULL,
70    effective_principal UUID NOT NULL,
71    object_id UUID NOT NULL,
72    command_fingerprint BYTES NOT NULL,
73    outcome TEXT NOT NULL,
74    metadata JSON NOT NULL
75);
76CREATE TABLE IF NOT EXISTS "outbox.message" (
77    message_id UUID PRIMARY KEY,
78    idempotency_key TEXT NOT NULL UNIQUE,
79    schema_version INTEGER NOT NULL,
80    payload JSON NOT NULL,
81    state TEXT NOT NULL,
82    created_at TIMESTAMP NOT NULL,
83    available_at TIMESTAMP NOT NULL,
84    lease_owner TEXT,
85    lease_token UUID,
86    lease_expires_at TIMESTAMP,
87    attempt_count INTEGER NOT NULL,
88    completed_at TIMESTAMP,
89    last_error TEXT,
90    dead_lettered_at TIMESTAMP
91);
92COMMIT;
93"#;
94
95/// Stable per-database identities of the two typed ordinary relations.
96#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub struct ApplicationRelationIdentity {
98    pub audit_relation: ObjectId,
99    pub outbox_relation: ObjectId,
100}
101
102/// One durable lease returned only after its claim transaction commits.
103#[derive(Debug, Clone, PartialEq)]
104pub struct OutboxClaim {
105    pub message_id: [u8; 16],
106    pub idempotency_key: String,
107    pub schema_version: u32,
108    pub payload: Value,
109    pub lease_token: [u8; 16],
110    pub lease_expires_at: DateTime<Utc>,
111    pub attempt: u32,
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq)]
115pub enum OutboxCompletion {
116    Completed,
117    AlreadyCompleted,
118}
119
120#[derive(Debug, Clone, Copy, PartialEq, Eq)]
121pub enum OutboxRetryDisposition {
122    Retried,
123    DeadLettered,
124}
125
126/// Bounded retention work performed by one maintenance transaction.
127#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128pub struct ApplicationRetentionPolicy {
129    pub audit_retention: Duration,
130    pub outbox_retention: Duration,
131    pub max_rows_per_relation: usize,
132}
133
134impl Default for ApplicationRetentionPolicy {
135    fn default() -> Self {
136        Self {
137            audit_retention: Duration::days(90),
138            outbox_retention: Duration::days(30),
139            max_rows_per_relation: 1024,
140        }
141    }
142}
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq)]
145pub struct ApplicationRetentionOutcome {
146    pub audit_rows_deleted: usize,
147    pub outbox_rows_deleted: usize,
148}
149
150impl Executor {
151    /// Install both relation contracts in one ordinary transaction.
152    ///
153    /// Existing names are accepted only when owner and ordered schema match
154    /// exactly. Opening a database never invokes this method implicitly.
155    pub fn install_application_relations(&self) -> Result<ApplicationRelationIdentity> {
156        if self.has_active_transaction() {
157            return Err(Error::invalid_argument(
158                "application relations cannot be installed inside an active transaction",
159            ));
160        }
161
162        let catalog = self.engine.pin_catalog()?;
163        validate_relation_if_present(catalog.as_ref(), AUDIT_RELATION_NAME, AUDIT_COLUMNS)?;
164        validate_relation_if_present(catalog.as_ref(), OUTBOX_RELATION_NAME, OUTBOX_COLUMNS)?;
165        drop(catalog);
166
167        let mut result = self.execute(INSTALL_SQL)?;
168        drain_result(result.as_mut())?;
169        result.close()?;
170        self.application_relation_identity()
171    }
172
173    /// Resolve and validate the durable catalog identity without mutating it.
174    pub fn application_relation_identity(&self) -> Result<ApplicationRelationIdentity> {
175        let (catalog, _) = transaction_visible_catalog(self)?;
176        Ok(ApplicationRelationIdentity {
177            audit_relation: validate_relation(
178                catalog.as_ref(),
179                AUDIT_RELATION_NAME,
180                AUDIT_COLUMNS,
181            )?,
182            outbox_relation: validate_relation(
183                catalog.as_ref(),
184                OUTBOX_RELATION_NAME,
185                OUTBOX_COLUMNS,
186            )?,
187        })
188    }
189
190    /// Append one immutable audit success record to the caller transaction.
191    /// Parameter values are never captured implicitly.
192    pub fn append_audit_event(
193        &self,
194        context: &ExecutionContext,
195        event: AuditEvent,
196    ) -> Result<[u8; 16]> {
197        self.application_relation_identity()?;
198        let metadata = encode_audit_metadata(&event.metadata)?;
199        let transaction_id = self
200            .active_transaction_id()
201            .ok_or(Error::TransactionNotStarted)?;
202        let event_id = *uuid::Uuid::now_v7().as_bytes();
203        let occurred_at = Utc::now();
204        let row = Row::from_values(vec![
205            Value::uuid(event_id),
206            Value::timestamp(occurred_at),
207            Value::Integer(transaction_id),
208            Value::uuid(context.principal_id().into_bytes()),
209            Value::uuid(context.effective_principal_id().into_bytes()),
210            Value::uuid(event.object_id.into_bytes()),
211            Value::bytes(event.command_fingerprint.to_vec()),
212            Value::text("success"),
213            Value::json(metadata),
214        ]);
215        self.insert_active_system_row(AUDIT_RELATION_NAME, row)?;
216        Ok(event_id)
217    }
218
219    /// Append one typed external side-effect intent to the caller transaction.
220    pub fn append_outbox_message(&self, message: OutboxMessage) -> Result<[u8; 16]> {
221        self.application_relation_identity()?;
222        validate_outbox_message(&message)?;
223        let message_id = *uuid::Uuid::now_v7().as_bytes();
224        let created_at = Utc::now();
225        let payload = message
226            .payload
227            .as_json()
228            .expect("validated JSON payload")
229            .to_owned();
230        let row = Row::from_values(vec![
231            Value::uuid(message_id),
232            Value::text(message.idempotency_key),
233            Value::Integer(i64::from(message.schema_version)),
234            Value::json(payload),
235            Value::text("pending"),
236            Value::timestamp(created_at),
237            Value::timestamp(created_at),
238            Value::Null(DataType::Text),
239            Value::Null(DataType::Uuid),
240            Value::Null(DataType::Timestamp),
241            Value::Integer(0),
242            Value::Null(DataType::Timestamp),
243            Value::Null(DataType::Text),
244            Value::Null(DataType::Timestamp),
245        ]);
246        self.insert_active_system_row(OUTBOX_RELATION_NAME, row)?;
247        Ok(message_id)
248    }
249
250    /// Atomically claim pending or expired-lease messages for one worker.
251    ///
252    /// The scan and returned batch are both bounded. A conflicting claimant
253    /// loses at the ordinary MVCC write-claim/commit boundary and receives no
254    /// unpublished lease records.
255    pub fn claim_outbox(
256        &self,
257        worker_id: &str,
258        now: DateTime<Utc>,
259        lease_duration: Duration,
260        limit: usize,
261        max_attempts: u32,
262    ) -> Result<Vec<OutboxClaim>> {
263        validate_claim_request(worker_id, lease_duration, limit, max_attempts)?;
264        self.with_owned_application_transaction(|executor| {
265            executor.claim_outbox_inside(worker_id, now, lease_duration, limit, max_attempts)
266        })
267    }
268
269    /// Durably mark delivery complete. Repeating the same token is idempotent.
270    pub fn complete_outbox(
271        &self,
272        message_id: [u8; 16],
273        lease_token: [u8; 16],
274        completed_at: DateTime<Utc>,
275    ) -> Result<OutboxCompletion> {
276        self.with_owned_application_transaction(|executor| {
277            executor.complete_outbox_inside(message_id, lease_token, completed_at)
278        })
279    }
280
281    /// Release a failed delivery for retry or move it to the dead-letter state.
282    pub fn retry_outbox(
283        &self,
284        message_id: [u8; 16],
285        lease_token: [u8; 16],
286        failed_at: DateTime<Utc>,
287        retry_at: DateTime<Utc>,
288        error: &str,
289        max_attempts: u32,
290    ) -> Result<OutboxRetryDisposition> {
291        if error.len() > MAX_OUTBOX_ERROR_BYTES {
292            return Err(Error::invalid_argument(format!(
293                "outbox delivery error exceeds {MAX_OUTBOX_ERROR_BYTES} bytes"
294            )));
295        }
296        if max_attempts == 0 || max_attempts > MAX_OUTBOX_ATTEMPTS {
297            return Err(Error::invalid_argument(
298                "outbox max_attempts is outside the supported range",
299            ));
300        }
301        self.with_owned_application_transaction(|executor| {
302            executor.retry_outbox_inside(
303                message_id,
304                lease_token,
305                failed_at,
306                retry_at,
307                error,
308                max_attempts,
309            )
310        })
311    }
312
313    /// Delete only expired immutable audit history and terminal outbox rows.
314    pub fn prune_application_history(
315        &self,
316        now: DateTime<Utc>,
317        policy: ApplicationRetentionPolicy,
318    ) -> Result<ApplicationRetentionOutcome> {
319        validate_retention_policy(policy)?;
320        self.with_owned_application_transaction(|executor| {
321            executor.prune_application_history_inside(now, policy)
322        })
323    }
324
325    fn with_owned_application_transaction<T>(
326        &self,
327        operation: impl FnOnce(&Self) -> Result<T>,
328    ) -> Result<T> {
329        if self.has_active_transaction() {
330            return Err(Error::invalid_argument(
331                "outbox worker lifecycle requires a connection without an active transaction",
332            ));
333        }
334        self.application_relation_identity()?;
335        let boundary = self.begin_procedural_boundary()?;
336        match operation(self) {
337            Ok(value) => match self.complete_procedural_boundary(&boundary) {
338                Ok(()) => Ok(value),
339                Err(error) => {
340                    let _ = self.abort_procedural_boundary(&boundary);
341                    Err(error)
342                }
343            },
344            Err(error) => {
345                let _ = self.abort_procedural_boundary(&boundary);
346                Err(error)
347            }
348        }
349    }
350
351    fn claim_outbox_inside(
352        &self,
353        worker_id: &str,
354        now: DateTime<Utc>,
355        lease_duration: Duration,
356        limit: usize,
357        max_attempts: u32,
358    ) -> Result<Vec<OutboxClaim>> {
359        let mut table = self.active_system_table(OUTBOX_RELATION_NAME)?;
360        let rows = scan_bounded(&*table, MAX_CLAIM_SCAN_ROWS)?;
361        let mut candidates = Vec::new();
362        let mut exhausted = Vec::new();
363        for (row_id, row) in rows {
364            if !outbox_row_is_claimable(&row, now)? {
365                continue;
366            }
367            let attempt = row_u32(&row, 10, "attempt_count")?;
368            if attempt >= max_attempts {
369                exhausted.push(row_id);
370            } else {
371                candidates.push((row_id, row));
372            }
373        }
374        candidates.sort_by(|left, right| {
375            row_timestamp(&left.1, 6, "available_at")
376                .expect("validated candidate timestamp")
377                .cmp(
378                    &row_timestamp(&right.1, 6, "available_at")
379                        .expect("validated candidate timestamp"),
380                )
381                .then_with(|| left.0.cmp(&right.0))
382        });
383        candidates.truncate(limit);
384
385        let mut claimed = Vec::with_capacity(candidates.len());
386        let mut row_ids = exhausted.clone();
387        row_ids.extend(candidates.iter().map(|(row_id, _)| *row_id));
388        row_ids.sort_unstable();
389        row_ids.dedup();
390        table.try_claim_rows(&row_ids)?;
391
392        for row_id in exhausted {
393            let mut setter = |mut row: Row| -> Result<(Row, bool)> {
394                if !outbox_row_is_claimable(&row, now)?
395                    || row_u32(&row, 10, "attempt_count")? < max_attempts
396                {
397                    return Ok((row, false));
398                }
399                row.set(4, Value::text("dead_lettered"))?;
400                row.set(7, Value::Null(DataType::Text))?;
401                row.set(8, Value::Null(DataType::Uuid))?;
402                row.set(9, Value::Null(DataType::Timestamp))?;
403                row.set(13, Value::timestamp(now))?;
404                Ok((row, true))
405            };
406            table.update_by_row_ids(&[row_id], &mut setter)?;
407        }
408
409        for (row_id, original) in candidates {
410            let lease_token = *uuid::Uuid::now_v7().as_bytes();
411            let lease_expires_at = now + lease_duration;
412            let attempt = row_u32(&original, 10, "attempt_count")?
413                .checked_add(1)
414                .ok_or_else(|| Error::invalid_argument("outbox attempt counter overflow"))?;
415            let mut setter = |mut row: Row| -> Result<(Row, bool)> {
416                if !outbox_row_is_claimable(&row, now)? {
417                    return Ok((row, false));
418                }
419                row.set(4, Value::text("claimed"))?;
420                row.set(7, Value::text(worker_id))?;
421                row.set(8, Value::uuid(lease_token))?;
422                row.set(9, Value::timestamp(lease_expires_at))?;
423                row.set(10, Value::Integer(i64::from(attempt)))?;
424                Ok((row, true))
425            };
426            if table.update_by_row_ids(&[row_id], &mut setter)? != 1 {
427                return Err(Error::TransactionSerializationConflict { row_id });
428            }
429            claimed.push(OutboxClaim {
430                message_id: row_uuid(&original, 0, "message_id")?,
431                idempotency_key: row_text(&original, 1, "idempotency_key")?.to_owned(),
432                schema_version: row_u32(&original, 2, "schema_version")?,
433                payload: original
434                    .get(3)
435                    .cloned()
436                    .ok_or_else(|| Error::internal("outbox payload column is missing"))?,
437                lease_token,
438                lease_expires_at,
439                attempt,
440            });
441        }
442        Ok(claimed)
443    }
444
445    fn complete_outbox_inside(
446        &self,
447        message_id: [u8; 16],
448        lease_token: [u8; 16],
449        completed_at: DateTime<Utc>,
450    ) -> Result<OutboxCompletion> {
451        let mut table = self.active_system_table(OUTBOX_RELATION_NAME)?;
452        let row_id = find_unique_row_id(&*table, "message_id", Value::uuid(message_id))?
453            .ok_or_else(|| Error::invalid_argument("outbox message does not exist"))?;
454        table.try_claim_rows(&[row_id])?;
455        let mut disposition = None;
456        let mut setter = |mut row: Row| -> Result<(Row, bool)> {
457            let state = row_text(&row, 4, "state")?;
458            let token = row_optional_uuid(&row, 8, "lease_token")?;
459            if state == "completed" && token == Some(lease_token) {
460                disposition = Some(OutboxCompletion::AlreadyCompleted);
461                return Ok((row, false));
462            }
463            if state != "claimed" || token != Some(lease_token) {
464                return Err(Error::invalid_argument(
465                    "outbox completion does not own the active lease",
466                ));
467            }
468            let expires = row_timestamp(&row, 9, "lease_expires_at")?;
469            if completed_at > expires {
470                return Err(Error::invalid_argument(
471                    "outbox completion lease has expired",
472                ));
473            }
474            row.set(4, Value::text("completed"))?;
475            row.set(9, Value::Null(DataType::Timestamp))?;
476            row.set(11, Value::timestamp(completed_at))?;
477            disposition = Some(OutboxCompletion::Completed);
478            Ok((row, true))
479        };
480        table.update_by_row_ids(&[row_id], &mut setter)?;
481        disposition.ok_or_else(|| Error::internal("outbox completion produced no disposition"))
482    }
483
484    fn retry_outbox_inside(
485        &self,
486        message_id: [u8; 16],
487        lease_token: [u8; 16],
488        failed_at: DateTime<Utc>,
489        retry_at: DateTime<Utc>,
490        error: &str,
491        max_attempts: u32,
492    ) -> Result<OutboxRetryDisposition> {
493        let mut table = self.active_system_table(OUTBOX_RELATION_NAME)?;
494        let row_id = find_unique_row_id(&*table, "message_id", Value::uuid(message_id))?
495            .ok_or_else(|| Error::invalid_argument("outbox message does not exist"))?;
496        table.try_claim_rows(&[row_id])?;
497        let mut disposition = None;
498        let mut setter = |mut row: Row| -> Result<(Row, bool)> {
499            if row_text(&row, 4, "state")? != "claimed"
500                || row_optional_uuid(&row, 8, "lease_token")? != Some(lease_token)
501            {
502                return Err(Error::invalid_argument(
503                    "outbox retry does not own the active lease",
504                ));
505            }
506            if failed_at > row_timestamp(&row, 9, "lease_expires_at")? {
507                return Err(Error::invalid_argument("outbox retry lease has expired"));
508            }
509            let attempts = row_u32(&row, 10, "attempt_count")?;
510            let dead = attempts >= max_attempts;
511            row.set(
512                4,
513                Value::text(if dead { "dead_lettered" } else { "pending" }),
514            )?;
515            row.set(7, Value::Null(DataType::Text))?;
516            row.set(8, Value::Null(DataType::Uuid))?;
517            row.set(9, Value::Null(DataType::Timestamp))?;
518            row.set(12, Value::text(error))?;
519            if dead {
520                row.set(13, Value::timestamp(failed_at))?;
521                disposition = Some(OutboxRetryDisposition::DeadLettered);
522            } else {
523                row.set(6, Value::timestamp(retry_at))?;
524                disposition = Some(OutboxRetryDisposition::Retried);
525            }
526            Ok((row, true))
527        };
528        table.update_by_row_ids(&[row_id], &mut setter)?;
529        disposition.ok_or_else(|| Error::internal("outbox retry produced no disposition"))
530    }
531
532    fn prune_application_history_inside(
533        &self,
534        now: DateTime<Utc>,
535        policy: ApplicationRetentionPolicy,
536    ) -> Result<ApplicationRetentionOutcome> {
537        let audit_cutoff = now - policy.audit_retention;
538        let outbox_cutoff = now - policy.outbox_retention;
539        let mut audit = self.active_system_table(AUDIT_RELATION_NAME)?;
540        let audit_rows = scan_bounded(&*audit, policy.max_rows_per_relation)?;
541        let audit_ids = audit_rows
542            .into_iter()
543            .filter_map(|(row_id, row)| {
544                row_timestamp(&row, 1, "occurred_at")
545                    .map(|timestamp| (timestamp < audit_cutoff).then_some(row_id))
546                    .transpose()
547            })
548            .collect::<Result<Vec<_>>>()?;
549        audit.try_claim_rows_for_delete(&audit_ids)?;
550        let audit_rows_deleted = usize::try_from(audit.delete_by_row_ids(&audit_ids)?)
551            .map_err(|_| Error::internal("negative audit retention delete count"))?;
552
553        let mut outbox = self.active_system_table(OUTBOX_RELATION_NAME)?;
554        let outbox_rows = scan_bounded(&*outbox, policy.max_rows_per_relation)?;
555        let outbox_ids = outbox_rows
556            .into_iter()
557            .filter_map(|(row_id, row)| match terminal_outbox_timestamp(&row) {
558                Ok(Some(timestamp)) if timestamp < outbox_cutoff => Some(Ok(row_id)),
559                Ok(_) => None,
560                Err(error) => Some(Err(error)),
561            })
562            .collect::<Result<Vec<_>>>()?;
563        outbox.try_claim_rows_for_delete(&outbox_ids)?;
564        let outbox_rows_deleted = usize::try_from(outbox.delete_by_row_ids(&outbox_ids)?)
565            .map_err(|_| Error::internal("negative outbox retention delete count"))?;
566        Ok(ApplicationRetentionOutcome {
567            audit_rows_deleted,
568            outbox_rows_deleted,
569        })
570    }
571
572    fn active_system_table(&self, relation: &str) -> Result<Box<dyn Table>> {
573        let mut active = self.active_transaction.lock().unwrap();
574        let state = active.as_mut().ok_or(Error::TransactionNotStarted)?;
575        let table = state.transaction.get_table(relation)?;
576        if !state.tables.contains_key(relation) {
577            state
578                .tables
579                .insert(relation.to_owned(), state.transaction.get_table(relation)?);
580        }
581        Ok(table)
582    }
583
584    fn insert_active_system_row(&self, relation: &str, row: Row) -> Result<()> {
585        let mut active = self.active_transaction.lock().unwrap();
586        let state = active.as_mut().ok_or(Error::TransactionNotStarted)?;
587        let mut table = state.transaction.get_table(relation)?;
588        if !state.tables.contains_key(relation) {
589            state
590                .tables
591                .insert(relation.to_owned(), state.transaction.get_table(relation)?);
592        }
593        drop(active);
594        table.insert_discard(row)
595    }
596}
597
598fn validate_claim_request(
599    worker_id: &str,
600    lease_duration: Duration,
601    limit: usize,
602    max_attempts: u32,
603) -> Result<()> {
604    if worker_id.is_empty() || worker_id.len() > MAX_WORKER_ID_BYTES {
605        return Err(Error::invalid_argument(format!(
606            "outbox worker ID must contain 1..={MAX_WORKER_ID_BYTES} bytes"
607        )));
608    }
609    if lease_duration < Duration::seconds(1) || lease_duration > Duration::hours(24) {
610        return Err(Error::invalid_argument(
611            "outbox lease duration must be between 1 second and 24 hours",
612        ));
613    }
614    if limit == 0 || limit > MAX_CLAIM_BATCH {
615        return Err(Error::invalid_argument(format!(
616            "outbox claim batch must contain 1..={MAX_CLAIM_BATCH} rows"
617        )));
618    }
619    if max_attempts == 0 || max_attempts > MAX_OUTBOX_ATTEMPTS {
620        return Err(Error::invalid_argument(
621            "outbox max_attempts is outside the supported range",
622        ));
623    }
624    Ok(())
625}
626
627fn validate_retention_policy(policy: ApplicationRetentionPolicy) -> Result<()> {
628    if policy.audit_retention < Duration::zero()
629        || policy.outbox_retention < Duration::zero()
630        || policy.max_rows_per_relation == 0
631        || policy.max_rows_per_relation > MAX_RETENTION_BATCH
632    {
633        return Err(Error::invalid_argument(format!(
634            "retention must be non-negative and each pass must delete 1..={MAX_RETENTION_BATCH} rows per relation"
635        )));
636    }
637    Ok(())
638}
639
640fn scan_bounded(table: &dyn Table, limit: usize) -> Result<Vec<(i64, Row)>> {
641    let projection = (0..table.schema().columns.len()).collect::<Vec<_>>();
642    let mut scanner = table.scan(&projection, None)?;
643    let mut rows = Vec::new();
644    while rows.len() < limit && scanner.next() {
645        rows.push(scanner.take_row_with_id()?);
646    }
647    if let Some(error) = scanner.err().cloned() {
648        let _ = scanner.close();
649        return Err(error);
650    }
651    scanner.close()?;
652    Ok(rows)
653}
654
655fn find_unique_row_id(table: &dyn Table, column: &str, value: Value) -> Result<Option<i64>> {
656    if let Some(row_ids) =
657        table.collect_row_ids_by_index_values(column, std::slice::from_ref(&value))
658    {
659        let row_ids = row_ids?;
660        return match row_ids.as_slice() {
661            [] => Ok(None),
662            [row_id] => Ok(Some(*row_id)),
663            _ => Err(Error::internal(format!(
664                "unique system relation column '{column}' returned multiple rows"
665            ))),
666        };
667    }
668    let column_index = table
669        .schema()
670        .find_column(column)
671        .map(|(index, _)| index)
672        .ok_or_else(|| Error::internal(format!("system relation column '{column}' is missing")))?;
673    let mut scanner = table.scan(&[column_index], None)?;
674    let mut found = None;
675    while scanner.next() {
676        let (row_id, row) = scanner.take_row_with_id()?;
677        if row.get(0) == Some(&value) && found.replace(row_id).is_some() {
678            let _ = scanner.close();
679            return Err(Error::internal(format!(
680                "unique system relation column '{column}' returned multiple rows"
681            )));
682        }
683    }
684    if let Some(error) = scanner.err().cloned() {
685        let _ = scanner.close();
686        return Err(error);
687    }
688    scanner.close()?;
689    Ok(found)
690}
691
692fn outbox_row_is_claimable(row: &Row, now: DateTime<Utc>) -> Result<bool> {
693    match row_text(row, 4, "state")? {
694        "pending" => Ok(row_timestamp(row, 6, "available_at")? <= now),
695        "claimed" => Ok(row_timestamp(row, 9, "lease_expires_at")? <= now),
696        "completed" | "dead_lettered" => Ok(false),
697        state => Err(Error::invalid_argument(format!(
698            "outbox row contains unknown state '{state}'"
699        ))),
700    }
701}
702
703fn terminal_outbox_timestamp(row: &Row) -> Result<Option<DateTime<Utc>>> {
704    match row_text(row, 4, "state")? {
705        "completed" => row_optional_timestamp(row, 11, "completed_at"),
706        "dead_lettered" => row_optional_timestamp(row, 13, "dead_lettered_at"),
707        "pending" | "claimed" => Ok(None),
708        state => Err(Error::invalid_argument(format!(
709            "outbox row contains unknown state '{state}'"
710        ))),
711    }
712}
713
714fn row_text<'a>(row: &'a Row, index: usize, column: &str) -> Result<&'a str> {
715    match row.get(index) {
716        Some(Value::Text(value)) => Ok(value.as_str()),
717        _ => Err(Error::internal(format!(
718            "system relation column '{column}' is not TEXT"
719        ))),
720    }
721}
722
723fn row_u32(row: &Row, index: usize, column: &str) -> Result<u32> {
724    match row.get(index) {
725        Some(Value::Integer(value)) => u32::try_from(*value).map_err(|_| {
726            Error::invalid_argument(format!("system relation column '{column}' is outside u32"))
727        }),
728        _ => Err(Error::internal(format!(
729            "system relation column '{column}' is not INTEGER"
730        ))),
731    }
732}
733
734fn row_uuid(row: &Row, index: usize, column: &str) -> Result<[u8; 16]> {
735    row.get(index)
736        .and_then(Value::as_uuid_bytes)
737        .ok_or_else(|| Error::internal(format!("system relation column '{column}' is not UUID")))
738}
739
740fn row_optional_uuid(row: &Row, index: usize, column: &str) -> Result<Option<[u8; 16]>> {
741    match row.get(index) {
742        Some(Value::Null(_)) => Ok(None),
743        Some(value) => value.as_uuid_bytes().map(Some).ok_or_else(|| {
744            Error::internal(format!("system relation column '{column}' is not UUID"))
745        }),
746        None => Err(Error::internal(format!(
747            "system relation column '{column}' is missing"
748        ))),
749    }
750}
751
752fn row_timestamp(row: &Row, index: usize, column: &str) -> Result<DateTime<Utc>> {
753    match row.get(index) {
754        Some(Value::Timestamp(value)) => Ok(*value),
755        _ => Err(Error::internal(format!(
756            "system relation column '{column}' is not TIMESTAMP"
757        ))),
758    }
759}
760
761fn row_optional_timestamp(row: &Row, index: usize, column: &str) -> Result<Option<DateTime<Utc>>> {
762    match row.get(index) {
763        Some(Value::Null(_)) => Ok(None),
764        Some(Value::Timestamp(value)) => Ok(Some(*value)),
765        _ => Err(Error::internal(format!(
766            "system relation column '{column}' is not TIMESTAMP"
767        ))),
768    }
769}
770
771fn drain_result(result: &mut dyn QueryResult) -> Result<()> {
772    while result.next() {
773        drop(result.take_row());
774    }
775    if let Some(error) = result.last_error() {
776        return Err(error);
777    }
778    Ok(())
779}
780
781fn validate_relation_if_present(
782    catalog: &CatalogGeneration,
783    name: &str,
784    columns: &[(&str, DataType, bool)],
785) -> Result<()> {
786    if catalog
787        .find_relation(ObjectId::BOOTSTRAP_NAMESPACE, name)
788        .map_err(catalog_error)?
789        .is_some()
790    {
791        validate_relation(catalog, name, columns)?;
792    }
793    Ok(())
794}
795
796fn validate_relation(
797    catalog: &CatalogGeneration,
798    name: &str,
799    columns: &[(&str, DataType, bool)],
800) -> Result<ObjectId> {
801    let relation = catalog
802        .find_relation(ObjectId::BOOTSTRAP_NAMESPACE, name)
803        .map_err(catalog_error)?
804        .ok_or_else(|| Error::TableNotFound(name.to_owned()))?;
805    if relation.kind() != ObjectKind::Table
806        || relation.owner_principal_id() != ObjectId::BOOTSTRAP_OWNER
807    {
808        return Err(Error::invalid_argument(format!(
809            "system relation '{name}' has an invalid kind or owner"
810        )));
811    }
812    let CatalogPayload::Table(payload) = relation.payload() else {
813        return Err(Error::internal("table catalog payload kind mismatch"));
814    };
815    if payload.column_ids().len() != columns.len() {
816        return Err(Error::invalid_argument(format!(
817            "system relation '{name}' has an incompatible column count"
818        )));
819    }
820    for (column_id, (expected_name, expected_type, expected_nullable)) in
821        payload.column_ids().iter().zip(columns)
822    {
823        let column = catalog
824            .object(*column_id)
825            .ok_or_else(|| Error::internal("system relation column is missing"))?;
826        let CatalogPayload::Column(column_payload) = column.payload() else {
827            return Err(Error::internal("system relation child is not a column"));
828        };
829        if column.name().normalized().as_str() != *expected_name
830            || column_payload.data_type().logical_type() != *expected_type
831            || column_payload.nullable() != *expected_nullable
832        {
833            return Err(Error::invalid_argument(format!(
834                "system relation '{name}' column contract mismatch at '{expected_name}'"
835            )));
836        }
837    }
838    validate_relation_keys(catalog, name, payload)?;
839    Ok(relation.id())
840}
841
842fn validate_relation_keys(
843    catalog: &CatalogGeneration,
844    name: &str,
845    table: &radixdb_catalog::TablePayload,
846) -> Result<()> {
847    let expected_primary_key = table
848        .column_ids()
849        .first()
850        .copied()
851        .ok_or_else(|| Error::internal("system relation has no identity column"))?;
852    let primary_key = table
853        .primary_key_constraint_id()
854        .and_then(|id| catalog.object(id))
855        .and_then(|object| match object.payload() {
856            CatalogPayload::Constraint(ConstraintPayload::PrimaryKey { local_column_ids }) => {
857                Some(local_column_ids.as_slice())
858            }
859            _ => None,
860        });
861    if primary_key != Some(std::slice::from_ref(&expected_primary_key)) {
862        return Err(Error::invalid_argument(format!(
863            "system relation '{name}' has an incompatible primary key"
864        )));
865    }
866
867    if name == OUTBOX_RELATION_NAME {
868        let expected_idempotency_key = table.column_ids()[1];
869        let has_unique_idempotency_key = table.constraint_ids().iter().any(|id| {
870            catalog.object(*id).is_some_and(|object| {
871                matches!(
872                    object.payload(),
873                    CatalogPayload::Constraint(ConstraintPayload::Unique { local_column_ids })
874                        if local_column_ids.as_slice()
875                            == std::slice::from_ref(&expected_idempotency_key)
876                )
877            })
878        });
879        if !has_unique_idempotency_key {
880            return Err(Error::invalid_argument(format!(
881                "system relation '{name}' lacks its idempotency-key uniqueness contract"
882            )));
883        }
884    }
885    Ok(())
886}
887
888fn validate_outbox_message(message: &OutboxMessage) -> Result<()> {
889    if message.idempotency_key.is_empty()
890        || message.idempotency_key.len() > MAX_OUTBOX_IDEMPOTENCY_KEY_BYTES
891    {
892        return Err(Error::invalid_argument(format!(
893            "outbox idempotency key must contain 1..={MAX_OUTBOX_IDEMPOTENCY_KEY_BYTES} bytes"
894        )));
895    }
896    if message.schema_version == 0 {
897        return Err(Error::invalid_argument(
898            "outbox payload schema version must be positive",
899        ));
900    }
901    let payload = message
902        .payload
903        .as_json()
904        .ok_or_else(|| Error::invalid_argument("outbox payload must be a validated JSON value"))?;
905    if payload.len() > MAX_OUTBOX_PAYLOAD_BYTES {
906        return Err(Error::invalid_argument(format!(
907            "outbox payload exceeds {MAX_OUTBOX_PAYLOAD_BYTES} bytes"
908        )));
909    }
910    Ok(())
911}
912
913fn encode_audit_metadata(metadata: &Value) -> Result<String> {
914    let encoded = metadata
915        .as_json()
916        .ok_or_else(|| Error::invalid_argument("audit metadata must be a validated JSON object"))?;
917    let parsed: serde_json::Value = serde_json::from_str(encoded)
918        .map_err(|error| Error::invalid_argument(format!("invalid audit JSON: {error}")))?;
919    let object = parsed
920        .as_object()
921        .ok_or_else(|| Error::invalid_argument("audit metadata must be a JSON object"))?;
922    if object.len() > MAX_AUDIT_METADATA_ENTRIES {
923        return Err(Error::invalid_argument(format!(
924            "audit metadata exceeds {MAX_AUDIT_METADATA_ENTRIES} entries"
925        )));
926    }
927    for key in object.keys() {
928        if key.is_empty() || key.len() > MAX_AUDIT_METADATA_KEY_BYTES {
929            return Err(Error::invalid_argument(format!(
930                "audit metadata key must contain 1..={MAX_AUDIT_METADATA_KEY_BYTES} bytes"
931            )));
932        }
933        let lowered = key.to_ascii_lowercase();
934        if ["password", "secret", "token", "credential", "authorization"]
935            .iter()
936            .any(|needle| lowered.contains(needle))
937        {
938            return Err(Error::invalid_argument(format!(
939                "audit metadata key '{key}' is secret-bearing"
940            )));
941        }
942    }
943    let encoded = serde_json::Value::Object(object.clone()).to_string();
944    if encoded.len() > MAX_AUDIT_METADATA_BYTES {
945        return Err(Error::invalid_argument(format!(
946            "audit metadata exceeds {MAX_AUDIT_METADATA_BYTES} encoded bytes"
947        )));
948    }
949    Ok(encoded)
950}
951
952fn catalog_error(error: impl std::fmt::Display) -> Error {
953    Error::invalid_argument(format!("invalid system relation catalog: {error}"))
954}
955
956#[cfg(test)]
957mod tests {
958    use std::sync::Arc;
959    use std::thread;
960
961    use chrono::Duration;
962    use radixdb_procedural::{AuditEvent, OutboxMessage};
963    use radixdb_storage::config::Config;
964    use radixdb_storage::mvcc::engine::MVCCEngine;
965
966    use super::*;
967
968    fn executor() -> Executor {
969        let engine = MVCCEngine::in_memory();
970        engine.open_engine().unwrap();
971        Executor::new(Arc::new(engine))
972    }
973
974    fn persistent_executor(path: &std::path::Path) -> (Executor, Arc<MVCCEngine>) {
975        let mut config = Config::with_path(path.to_string_lossy().to_string());
976        config.persistence.checkpoint_on_close = true;
977        let engine = Arc::new(MVCCEngine::new_with_composition_binders(
978            config,
979            crate::mutation::partial_index::bind_from_sql,
980            crate::mutation::row_validation::bind,
981            crate::mutation::view_binding::bind_from_sql,
982            radixdb_storage::mvcc::engine::CatalogRuntimeBinder::new(
983                crate::catalog::bind_runtime_catalog,
984            ),
985        ));
986        engine.open_engine().unwrap();
987        (Executor::new(Arc::clone(&engine)), engine)
988    }
989
990    fn scalar_count(executor: &Executor, relation: &str) -> i64 {
991        let mut result = executor
992            .execute(&format!("SELECT COUNT(*) FROM \"{relation}\""))
993            .unwrap();
994        assert!(result.next());
995        let row = result.take_row();
996        let Value::Integer(count) = row.get(0).unwrap() else {
997            panic!("COUNT must return INTEGER")
998        };
999        *count
1000    }
1001
1002    fn enqueue(executor: &Executor, key: &str) -> [u8; 16] {
1003        let boundary = executor.begin_procedural_boundary().unwrap();
1004        let id = executor
1005            .append_outbox_message(OutboxMessage {
1006                idempotency_key: key.to_owned(),
1007                schema_version: 1,
1008                payload: Value::json(format!(r#"{{"key":"{key}"}}"#)),
1009            })
1010            .unwrap();
1011        executor.complete_procedural_boundary(&boundary).unwrap();
1012        id
1013    }
1014
1015    #[test]
1016    fn relations_install_atomically_and_keep_stable_catalog_identity() {
1017        let executor = executor();
1018        let first = executor.install_application_relations().unwrap();
1019        let second = executor.install_application_relations().unwrap();
1020        assert_eq!(first, second);
1021        assert_eq!(scalar_count(&executor, AUDIT_RELATION_NAME), 0);
1022        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 0);
1023
1024        let mut unquoted = executor
1025            .execute("SELECT COUNT(*) FROM audit.event")
1026            .unwrap();
1027        assert!(unquoted.next());
1028        assert_eq!(unquoted.row().get(0), Some(&Value::Integer(0)));
1029        unquoted.close().unwrap();
1030    }
1031
1032    #[test]
1033    fn installer_rejects_lookalike_relations_without_required_keys() {
1034        let executor = executor();
1035        executor
1036            .execute(
1037                r#"CREATE TABLE "outbox.message" (
1038                    message_id UUID PRIMARY KEY,
1039                    idempotency_key TEXT NOT NULL,
1040                    schema_version INTEGER NOT NULL,
1041                    payload JSON NOT NULL,
1042                    state TEXT NOT NULL,
1043                    created_at TIMESTAMP NOT NULL,
1044                    available_at TIMESTAMP NOT NULL,
1045                    lease_owner TEXT,
1046                    lease_token UUID,
1047                    lease_expires_at TIMESTAMP,
1048                    attempt_count INTEGER NOT NULL,
1049                    completed_at TIMESTAMP,
1050                    last_error TEXT,
1051                    dead_lettered_at TIMESTAMP
1052                )"#,
1053            )
1054            .unwrap();
1055
1056        let error = executor.install_application_relations().unwrap_err();
1057        assert!(error.to_string().contains("idempotency-key uniqueness"));
1058        assert!(executor
1059            .engine()
1060            .pin_catalog()
1061            .unwrap()
1062            .find_relation(ObjectId::BOOTSTRAP_NAMESPACE, AUDIT_RELATION_NAME)
1063            .unwrap()
1064            .is_none());
1065    }
1066
1067    #[test]
1068    fn system_relation_schema_and_immutable_audit_history_are_protected() {
1069        let executor = executor();
1070        executor.install_application_relations().unwrap();
1071
1072        for statement in [
1073            "INSERT INTO audit.event VALUES (1)",
1074            "INSERT INTO outbox.message VALUES (1)",
1075            "UPDATE audit.event SET outcome = 'forged'",
1076            "DELETE FROM audit.event",
1077            "TRUNCATE TABLE audit.event",
1078            "TRUNCATE TABLE outbox.message",
1079            "DROP TABLE audit.event",
1080            "DROP TABLE outbox.message",
1081            "ALTER TABLE \"audit.event\" ADD COLUMN forged TEXT",
1082            "ALTER TABLE \"outbox.message\" ADD COLUMN forged TEXT",
1083        ] {
1084            assert!(
1085                matches!(
1086                    executor.execute(statement),
1087                    Err(Error::AuthorizationDenied(_))
1088                ),
1089                "system relation mutation unexpectedly passed: {statement}"
1090            );
1091        }
1092    }
1093
1094    #[test]
1095    fn audit_and_outbox_share_caller_commit_and_rollback() {
1096        let executor = executor();
1097        executor.install_application_relations().unwrap();
1098        let context = ExecutionContext::new();
1099
1100        let rollback = executor.begin_procedural_boundary().unwrap();
1101        executor
1102            .append_audit_event(
1103                &context,
1104                AuditEvent {
1105                    object_id: ObjectId::BOOTSTRAP_NAMESPACE,
1106                    command_fingerprint: [7; 32],
1107                    metadata: Value::json(r#"{"command":"update"}"#),
1108                },
1109            )
1110            .unwrap();
1111        executor
1112            .append_outbox_message(OutboxMessage {
1113                idempotency_key: "rollback-key".to_owned(),
1114                schema_version: 1,
1115                payload: Value::json(r#"{"kind":"rollback"}"#),
1116            })
1117            .unwrap();
1118        executor.abort_procedural_boundary(&rollback).unwrap();
1119        assert_eq!(scalar_count(&executor, AUDIT_RELATION_NAME), 0);
1120        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 0);
1121
1122        let commit = executor.begin_procedural_boundary().unwrap();
1123        executor
1124            .append_audit_event(
1125                &context,
1126                AuditEvent {
1127                    object_id: ObjectId::BOOTSTRAP_NAMESPACE,
1128                    command_fingerprint: [8; 32],
1129                    metadata: Value::json("{}"),
1130                },
1131            )
1132            .unwrap();
1133        executor
1134            .append_outbox_message(OutboxMessage {
1135                idempotency_key: "commit-key".to_owned(),
1136                schema_version: 1,
1137                payload: Value::json(r#"{"kind":"commit"}"#),
1138            })
1139            .unwrap();
1140        executor.complete_procedural_boundary(&commit).unwrap();
1141        assert_eq!(scalar_count(&executor, AUDIT_RELATION_NAME), 1);
1142        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 1);
1143    }
1144
1145    #[test]
1146    fn metadata_and_payload_contracts_fail_closed() {
1147        let executor = executor();
1148        executor.install_application_relations().unwrap();
1149        let boundary = executor.begin_procedural_boundary().unwrap();
1150        let error = executor
1151            .append_audit_event(
1152                &ExecutionContext::new(),
1153                AuditEvent {
1154                    object_id: ObjectId::BOOTSTRAP_NAMESPACE,
1155                    command_fingerprint: [0; 32],
1156                    metadata: Value::json(r#"{"password":"do-not-store"}"#),
1157                },
1158            )
1159            .unwrap_err();
1160        assert!(error.to_string().contains("secret-bearing"));
1161        let error = executor
1162            .append_outbox_message(OutboxMessage {
1163                idempotency_key: "not-json".to_owned(),
1164                schema_version: 1,
1165                payload: Value::text("not-json"),
1166            })
1167            .unwrap_err();
1168        assert!(error.to_string().contains("validated JSON"));
1169        executor.abort_procedural_boundary(&boundary).unwrap();
1170    }
1171
1172    #[test]
1173    fn claim_expiry_and_completion_are_transactional_and_idempotent() {
1174        let executor = executor();
1175        executor.install_application_relations().unwrap();
1176        let message_id = enqueue(&executor, "lease-key");
1177        let now = Utc::now();
1178
1179        let first = executor
1180            .claim_outbox("worker-a", now, Duration::minutes(1), 4, 3)
1181            .unwrap();
1182        assert_eq!(first.len(), 1);
1183        assert_eq!(first[0].message_id, message_id);
1184        assert_eq!(first[0].attempt, 1);
1185        assert!(executor
1186            .claim_outbox(
1187                "worker-b",
1188                now + Duration::seconds(30),
1189                Duration::minutes(1),
1190                4,
1191                3,
1192            )
1193            .unwrap()
1194            .is_empty());
1195
1196        assert!(executor
1197            .retry_outbox(
1198                message_id,
1199                first[0].lease_token,
1200                now + Duration::minutes(2),
1201                now + Duration::minutes(3),
1202                "late worker",
1203                3,
1204            )
1205            .is_err());
1206
1207        let reclaimed = executor
1208            .claim_outbox(
1209                "worker-b",
1210                now + Duration::minutes(2),
1211                Duration::minutes(1),
1212                4,
1213                3,
1214            )
1215            .unwrap();
1216        assert_eq!(reclaimed.len(), 1);
1217        assert_eq!(reclaimed[0].attempt, 2);
1218        assert_ne!(reclaimed[0].lease_token, first[0].lease_token);
1219        assert!(executor
1220            .complete_outbox(message_id, first[0].lease_token, now + Duration::minutes(2),)
1221            .is_err());
1222        assert_eq!(
1223            executor
1224                .complete_outbox(
1225                    message_id,
1226                    reclaimed[0].lease_token,
1227                    now + Duration::minutes(2),
1228                )
1229                .unwrap(),
1230            OutboxCompletion::Completed
1231        );
1232        assert_eq!(
1233            executor
1234                .complete_outbox(
1235                    message_id,
1236                    reclaimed[0].lease_token,
1237                    now + Duration::minutes(2),
1238                )
1239                .unwrap(),
1240            OutboxCompletion::AlreadyCompleted
1241        );
1242    }
1243
1244    #[test]
1245    fn retry_releases_lease_and_last_attempt_dead_letters() {
1246        let executor = executor();
1247        executor.install_application_relations().unwrap();
1248        let message_id = enqueue(&executor, "retry-key");
1249        let now = Utc::now();
1250        let first = executor
1251            .claim_outbox("worker", now, Duration::minutes(1), 1, 2)
1252            .unwrap()
1253            .remove(0);
1254        assert_eq!(
1255            executor
1256                .retry_outbox(
1257                    message_id,
1258                    first.lease_token,
1259                    now + Duration::seconds(1),
1260                    now + Duration::seconds(5),
1261                    "temporary",
1262                    2,
1263                )
1264                .unwrap(),
1265            OutboxRetryDisposition::Retried
1266        );
1267        let second = executor
1268            .claim_outbox(
1269                "worker",
1270                now + Duration::seconds(5),
1271                Duration::minutes(1),
1272                1,
1273                2,
1274            )
1275            .unwrap()
1276            .remove(0);
1277        assert_eq!(second.attempt, 2);
1278        assert_eq!(
1279            executor
1280                .retry_outbox(
1281                    message_id,
1282                    second.lease_token,
1283                    now + Duration::seconds(6),
1284                    now + Duration::seconds(7),
1285                    "permanent",
1286                    2,
1287                )
1288                .unwrap(),
1289            OutboxRetryDisposition::DeadLettered
1290        );
1291        assert!(executor
1292            .claim_outbox(
1293                "worker",
1294                now + Duration::hours(1),
1295                Duration::minutes(1),
1296                1,
1297                2,
1298            )
1299            .unwrap()
1300            .is_empty());
1301    }
1302
1303    #[test]
1304    fn business_write_and_outbox_share_one_commit_and_idempotency_is_unique() {
1305        let executor = executor();
1306        executor.install_application_relations().unwrap();
1307        executor
1308            .execute("CREATE TABLE business_event (id INTEGER PRIMARY KEY)")
1309            .unwrap();
1310
1311        let rollback = executor.begin_procedural_boundary().unwrap();
1312        executor
1313            .execute("INSERT INTO business_event VALUES (1)")
1314            .unwrap();
1315        executor
1316            .append_outbox_message(OutboxMessage {
1317                idempotency_key: "business-1".to_owned(),
1318                schema_version: 1,
1319                payload: Value::json(r#"{"id":1}"#),
1320            })
1321            .unwrap();
1322        executor.abort_procedural_boundary(&rollback).unwrap();
1323        assert_eq!(scalar_count(&executor, "business_event"), 0);
1324        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 0);
1325
1326        let commit = executor.begin_procedural_boundary().unwrap();
1327        executor
1328            .execute("INSERT INTO business_event VALUES (1)")
1329            .unwrap();
1330        executor
1331            .append_outbox_message(OutboxMessage {
1332                idempotency_key: "business-1".to_owned(),
1333                schema_version: 1,
1334                payload: Value::json(r#"{"id":1}"#),
1335            })
1336            .unwrap();
1337        executor.complete_procedural_boundary(&commit).unwrap();
1338        assert_eq!(scalar_count(&executor, "business_event"), 1);
1339        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 1);
1340
1341        let duplicate = executor.begin_procedural_boundary().unwrap();
1342        assert!(executor
1343            .append_outbox_message(OutboxMessage {
1344                idempotency_key: "business-1".to_owned(),
1345                schema_version: 1,
1346                payload: Value::json(r#"{"id":1}"#),
1347            })
1348            .is_err());
1349        executor.abort_procedural_boundary(&duplicate).unwrap();
1350        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 1);
1351    }
1352
1353    #[test]
1354    fn retention_is_bounded_and_removes_only_terminal_history() {
1355        let executor = executor();
1356        executor.install_application_relations().unwrap();
1357        let context = ExecutionContext::new();
1358        let boundary = executor.begin_procedural_boundary().unwrap();
1359        executor
1360            .append_audit_event(
1361                &context,
1362                AuditEvent {
1363                    object_id: ObjectId::BOOTSTRAP_NAMESPACE,
1364                    command_fingerprint: [1; 32],
1365                    metadata: Value::json("{}"),
1366                },
1367            )
1368            .unwrap();
1369        executor.complete_procedural_boundary(&boundary).unwrap();
1370        let terminal_id = enqueue(&executor, "terminal");
1371        enqueue(&executor, "pending");
1372        let now = Utc::now();
1373        let claim = executor
1374            .claim_outbox("worker", now, Duration::minutes(1), 1, 3)
1375            .unwrap()
1376            .remove(0);
1377        assert_eq!(claim.message_id, terminal_id);
1378        executor
1379            .complete_outbox(terminal_id, claim.lease_token, now)
1380            .unwrap();
1381
1382        let outcome = executor
1383            .prune_application_history(
1384                now + Duration::seconds(1),
1385                ApplicationRetentionPolicy {
1386                    audit_retention: Duration::zero(),
1387                    outbox_retention: Duration::zero(),
1388                    max_rows_per_relation: 16,
1389                },
1390            )
1391            .unwrap();
1392        assert_eq!(outcome.audit_rows_deleted, 1);
1393        assert_eq!(outcome.outbox_rows_deleted, 1);
1394        assert_eq!(scalar_count(&executor, AUDIT_RELATION_NAME), 0);
1395        assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 1);
1396    }
1397
1398    #[test]
1399    fn concurrent_claimers_never_publish_two_leases_for_one_message() {
1400        let engine = Arc::new(MVCCEngine::in_memory());
1401        engine.open_engine().unwrap();
1402        let setup = Executor::new(Arc::clone(&engine));
1403        setup.install_application_relations().unwrap();
1404        let message_id = enqueue(&setup, "concurrent-key");
1405        let barrier = Arc::new(std::sync::Barrier::new(3));
1406        let now = Utc::now();
1407        let mut handles = Vec::new();
1408        for worker in ["worker-a", "worker-b"] {
1409            let engine = Arc::clone(&engine);
1410            let barrier = Arc::clone(&barrier);
1411            handles.push(thread::spawn(move || {
1412                let executor = Executor::new(engine);
1413                barrier.wait();
1414                executor.claim_outbox(worker, now, Duration::minutes(1), 1, 3)
1415            }));
1416        }
1417        barrier.wait();
1418        let outcomes = handles
1419            .into_iter()
1420            .map(|handle| handle.join().unwrap())
1421            .collect::<Vec<_>>();
1422        let claims = outcomes
1423            .iter()
1424            .filter_map(|outcome| outcome.as_ref().ok())
1425            .flatten()
1426            .collect::<Vec<_>>();
1427        assert_eq!(claims.len(), 1, "one message must have one durable lease");
1428        assert_eq!(claims[0].message_id, message_id);
1429    }
1430
1431    #[test]
1432    fn commit_claim_delivery_and_completion_survive_each_reopen_boundary() {
1433        let directory = tempfile::tempdir().unwrap();
1434        let message_id = {
1435            let (executor, engine) = persistent_executor(directory.path());
1436            executor.install_application_relations().unwrap();
1437            let id = enqueue(&executor, "reopen-key");
1438            engine.close_engine().unwrap();
1439            id
1440        };
1441        let now = Utc::now();
1442        let claim = {
1443            let (executor, engine) = persistent_executor(directory.path());
1444            assert_eq!(scalar_count(&executor, OUTBOX_RELATION_NAME), 1);
1445            let claim = executor
1446                .claim_outbox("worker-a", now, Duration::minutes(1), 1, 3)
1447                .unwrap()
1448                .remove(0);
1449            engine.close_engine().unwrap();
1450            claim
1451        };
1452        let reclaimed = {
1453            let (executor, engine) = persistent_executor(directory.path());
1454            assert!(executor
1455                .claim_outbox(
1456                    "worker-b",
1457                    now + Duration::seconds(30),
1458                    Duration::minutes(1),
1459                    1,
1460                    3,
1461                )
1462                .unwrap()
1463                .is_empty());
1464            let reclaimed = executor
1465                .claim_outbox(
1466                    "worker-b",
1467                    now + Duration::minutes(2),
1468                    Duration::minutes(1),
1469                    1,
1470                    3,
1471                )
1472                .unwrap()
1473                .remove(0);
1474            assert_eq!(reclaimed.message_id, message_id);
1475            assert_ne!(reclaimed.lease_token, claim.lease_token);
1476            engine.close_engine().unwrap();
1477            reclaimed
1478        };
1479        {
1480            let (executor, engine) = persistent_executor(directory.path());
1481            assert_eq!(
1482                executor
1483                    .complete_outbox(
1484                        message_id,
1485                        reclaimed.lease_token,
1486                        now + Duration::minutes(2),
1487                    )
1488                    .unwrap(),
1489                OutboxCompletion::Completed
1490            );
1491            engine.close_engine().unwrap();
1492        }
1493        {
1494            let (executor, engine) = persistent_executor(directory.path());
1495            assert_eq!(
1496                executor
1497                    .complete_outbox(
1498                        message_id,
1499                        reclaimed.lease_token,
1500                        now + Duration::minutes(2),
1501                    )
1502                    .unwrap(),
1503                OutboxCompletion::AlreadyCompleted
1504            );
1505            assert!(executor
1506                .claim_outbox(
1507                    "worker-c",
1508                    now + Duration::hours(1),
1509                    Duration::minutes(1),
1510                    1,
1511                    3,
1512                )
1513                .unwrap()
1514                .is_empty());
1515            engine.close_engine().unwrap();
1516        }
1517    }
1518}