1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub struct ApplicationRelationIdentity {
98 pub audit_relation: ObjectId,
99 pub outbox_relation: ObjectId,
100}
101
102#[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#[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 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 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 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 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 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 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 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 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}