1use async_trait::async_trait;
2use chrono::{DateTime, TimeDelta, Utc};
3use minco_plugin_audit::{AuditError, AuditEvent, AuditSink};
4use minco_plugin_events::{DomainEvent, EventError, OutboxRecord, OutboxStatus, OutboxStore};
5use minco_plugin_idempotency::{
6 BeginOutcome, IdempotencyError, IdempotencyKey, IdempotencyLease, IdempotencyRecord,
7 IdempotencyStore, RequestFingerprint, validate_claim_timeout,
8};
9use minco_plugin_sessions::{
10 SessionError, SessionId, SessionRecord, SessionStore, SessionTokenHash,
11};
12use sqlx::{PgPool, Postgres, Row, Transaction};
13use uuid::Uuid;
14
15#[derive(Debug, Clone)]
16pub struct PostgresOutboxStore {
17 pool: PgPool,
18}
19
20impl PostgresOutboxStore {
21 pub const fn new(pool: PgPool) -> Self {
22 Self { pool }
23 }
24
25 pub async fn enqueue_in(
29 &self,
30 transaction: &mut Transaction<'_, Postgres>,
31 record: OutboxRecord,
32 ) -> Result<(), EventError> {
33 validate_event(&record.event)?;
34 if record.status != OutboxStatus::Pending {
35 return Err(EventError::InvalidOutboxState);
36 }
37 let event_id = record.event.id;
38 let attempt_count =
39 i32::try_from(record.attempt_count).map_err(event_infrastructure_error)?;
40 let metadata =
41 serde_json::to_value(&record.event.metadata).map_err(event_infrastructure_error)?;
42 let result = sqlx::query(
43 "INSERT INTO minco_outbox
44 (event_id, event_type, aggregate_type, aggregate_id, correlation_id, occurred_at,
45 payload, metadata, status, attempt_count, available_at, claimed_by,
46 claim_expires_at, last_error)
47 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, 'pending', $9, $10, $11, $12, $13)",
48 )
49 .bind(event_id)
50 .bind(record.event.event_type)
51 .bind(record.event.aggregate_type)
52 .bind(record.event.aggregate_id)
53 .bind(record.event.correlation_id)
54 .bind(record.event.occurred_at)
55 .bind(record.event.payload)
56 .bind(metadata)
57 .bind(attempt_count)
58 .bind(record.available_at)
59 .bind(record.claimed_by)
60 .bind(record.claim_expires_at)
61 .bind(record.last_error)
62 .execute(&mut **transaction)
63 .await;
64 match result {
65 Ok(_) => Ok(()),
66 Err(error) if is_unique_violation(&error) => Err(EventError::DuplicateEvent(event_id)),
67 Err(error) => Err(event_infrastructure_error(error)),
68 }
69 }
70}
71
72#[async_trait]
73impl OutboxStore for PostgresOutboxStore {
74 async fn enqueue(&self, record: OutboxRecord) -> Result<(), EventError> {
75 let mut transaction = self
76 .pool
77 .begin()
78 .await
79 .map_err(event_infrastructure_error)?;
80 self.enqueue_in(&mut transaction, record).await?;
81 transaction
82 .commit()
83 .await
84 .map_err(event_infrastructure_error)
85 }
86
87 async fn claim_pending(
88 &self,
89 worker_id: &str,
90 limit: usize,
91 claim_expires_at: DateTime<Utc>,
92 ) -> Result<Vec<OutboxRecord>, EventError> {
93 validate_claim(worker_id, claim_expires_at)?;
94 let now = Utc::now();
95 let limit = i64::try_from(limit).map_err(|_| EventError::InvalidClaim)?;
96 if limit == 0 {
97 return Err(EventError::InvalidClaim);
98 }
99 let rows = sqlx::query(
100 "WITH candidates AS (
101 SELECT event_id
102 FROM minco_outbox
103 WHERE status IN ('pending', 'failed') AND available_at <= $1
104 ORDER BY available_at, occurred_at, event_id
105 FOR UPDATE SKIP LOCKED
106 LIMIT $2
107 )
108 UPDATE minco_outbox AS outbox
109 SET status = 'claimed',
110 claimed_by = $3,
111 claim_expires_at = $4,
112 attempt_count = outbox.attempt_count + 1
113 FROM candidates
114 WHERE outbox.event_id = candidates.event_id
115 RETURNING outbox.*",
116 )
117 .bind(now)
118 .bind(limit)
119 .bind(worker_id)
120 .bind(claim_expires_at)
121 .fetch_all(&self.pool)
122 .await
123 .map_err(event_infrastructure_error)?;
124 rows.iter().map(decode_outbox).collect()
125 }
126
127 async fn claim_event(
128 &self,
129 event_id: Uuid,
130 worker_id: &str,
131 claim_expires_at: DateTime<Utc>,
132 ) -> Result<Option<OutboxRecord>, EventError> {
133 validate_claim(worker_id, claim_expires_at)?;
134 let now = Utc::now();
135 let row = sqlx::query(
136 "UPDATE minco_outbox
137 SET status = 'claimed', claimed_by = $2, claim_expires_at = $3,
138 attempt_count = attempt_count + 1
139 WHERE event_id = $1
140 AND status IN ('pending', 'failed')
141 AND available_at <= $4
142 RETURNING *",
143 )
144 .bind(event_id)
145 .bind(worker_id)
146 .bind(claim_expires_at)
147 .bind(now)
148 .fetch_optional(&self.pool)
149 .await
150 .map_err(event_infrastructure_error)?;
151 row.as_ref().map(decode_outbox).transpose()
152 }
153
154 async fn mark_published(&self, event_id: Uuid, worker_id: &str) -> Result<(), EventError> {
155 let result = sqlx::query(
156 "UPDATE minco_outbox
157 SET status = 'published', claimed_by = NULL, claim_expires_at = NULL,
158 last_error = NULL
159 WHERE event_id = $1 AND status = 'claimed' AND claimed_by = $2",
160 )
161 .bind(event_id)
162 .bind(worker_id)
163 .execute(&self.pool)
164 .await
165 .map_err(event_infrastructure_error)?;
166 require_claim_update(&self.pool, result.rows_affected(), event_id, worker_id).await
167 }
168
169 async fn mark_failed(
170 &self,
171 event_id: Uuid,
172 worker_id: &str,
173 error: String,
174 retry_at: DateTime<Utc>,
175 ) -> Result<(), EventError> {
176 let result = sqlx::query(
177 "UPDATE minco_outbox
178 SET status = 'failed', claimed_by = NULL, claim_expires_at = NULL,
179 available_at = $3, last_error = $4
180 WHERE event_id = $1 AND status = 'claimed' AND claimed_by = $2",
181 )
182 .bind(event_id)
183 .bind(worker_id)
184 .bind(retry_at)
185 .bind(error)
186 .execute(&self.pool)
187 .await
188 .map_err(event_infrastructure_error)?;
189 require_claim_update(&self.pool, result.rows_affected(), event_id, worker_id).await
190 }
191
192 async fn recover_expired_claims(&self, now: DateTime<Utc>) -> Result<usize, EventError> {
193 let result = sqlx::query(
194 "UPDATE minco_outbox
195 SET status = 'pending', claimed_by = NULL, claim_expires_at = NULL
196 WHERE status = 'claimed' AND claim_expires_at <= $1",
197 )
198 .bind(now)
199 .execute(&self.pool)
200 .await
201 .map_err(event_infrastructure_error)?;
202 usize::try_from(result.rows_affected()).map_err(event_infrastructure_error)
203 }
204}
205
206async fn require_claim_update(
207 pool: &PgPool,
208 rows_affected: u64,
209 event_id: Uuid,
210 worker_id: &str,
211) -> Result<(), EventError> {
212 if rows_affected == 1 {
213 return Ok(());
214 }
215 let exists: bool =
216 sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM minco_outbox WHERE event_id = $1)")
217 .bind(event_id)
218 .fetch_one(pool)
219 .await
220 .map_err(event_infrastructure_error)?;
221 if exists {
222 Err(EventError::ClaimOwnership {
223 event_id,
224 worker_id: worker_id.to_owned(),
225 })
226 } else {
227 Err(EventError::MissingEvent(event_id))
228 }
229}
230
231fn decode_outbox(row: &sqlx::postgres::PgRow) -> Result<OutboxRecord, EventError> {
232 let status: String = row.try_get("status").map_err(event_infrastructure_error)?;
233 let status = match status.as_str() {
234 "pending" => OutboxStatus::Pending,
235 "claimed" => OutboxStatus::Claimed,
236 "published" => OutboxStatus::Published,
237 "failed" => OutboxStatus::Failed,
238 _ => {
239 return Err(event_infrastructure_error(
240 "invalid outbox status in PostgreSQL",
241 ));
242 }
243 };
244 let metadata: serde_json::Value = row
245 .try_get("metadata")
246 .map_err(event_infrastructure_error)?;
247 let attempt_count: i32 = row
248 .try_get("attempt_count")
249 .map_err(event_infrastructure_error)?;
250 Ok(OutboxRecord {
251 event: DomainEvent {
252 id: row
253 .try_get("event_id")
254 .map_err(event_infrastructure_error)?,
255 event_type: row
256 .try_get("event_type")
257 .map_err(event_infrastructure_error)?,
258 aggregate_type: row
259 .try_get("aggregate_type")
260 .map_err(event_infrastructure_error)?,
261 aggregate_id: row
262 .try_get("aggregate_id")
263 .map_err(event_infrastructure_error)?,
264 correlation_id: row
265 .try_get("correlation_id")
266 .map_err(event_infrastructure_error)?,
267 occurred_at: row
268 .try_get("occurred_at")
269 .map_err(event_infrastructure_error)?,
270 payload: row.try_get("payload").map_err(event_infrastructure_error)?,
271 metadata: serde_json::from_value(metadata).map_err(event_infrastructure_error)?,
272 },
273 status,
274 attempt_count: u32::try_from(attempt_count).map_err(event_infrastructure_error)?,
275 available_at: row
276 .try_get("available_at")
277 .map_err(event_infrastructure_error)?,
278 claimed_by: row
279 .try_get("claimed_by")
280 .map_err(event_infrastructure_error)?,
281 claim_expires_at: row
282 .try_get("claim_expires_at")
283 .map_err(event_infrastructure_error)?,
284 last_error: row
285 .try_get("last_error")
286 .map_err(event_infrastructure_error)?,
287 })
288}
289
290#[derive(Debug, Clone)]
291pub struct PostgresSessionStore {
292 pool: PgPool,
293}
294
295impl PostgresSessionStore {
296 pub const fn new(pool: PgPool) -> Self {
297 Self { pool }
298 }
299}
300
301#[async_trait]
302impl SessionStore for PostgresSessionStore {
303 async fn create(
304 &self,
305 token_hash: SessionTokenHash,
306 session: SessionRecord,
307 ) -> Result<(), SessionError> {
308 let attributes = serde_json::to_value(&session.attributes).map_err(session_store_error)?;
309 let result = sqlx::query(
310 "INSERT INTO minco_sessions
311 (id, token_hash, subject, created_at, expires_at, revoked_at, attributes)
312 VALUES ($1, $2, $3, $4, $5, $6, $7)",
313 )
314 .bind(session.id.0)
315 .bind(token_hash.as_bytes().as_slice())
316 .bind(session.subject)
317 .bind(session.created_at)
318 .bind(session.expires_at)
319 .bind(session.revoked_at)
320 .bind(attributes)
321 .execute(&self.pool)
322 .await;
323 match result {
324 Ok(_) => Ok(()),
325 Err(error) if is_unique_violation(&error) => Err(SessionError::Duplicate),
326 Err(error) => Err(session_store_error(error)),
327 }
328 }
329
330 async fn find_by_token_hash(
331 &self,
332 token_hash: SessionTokenHash,
333 ) -> Result<Option<SessionRecord>, SessionError> {
334 let row = sqlx::query(
335 "SELECT id, subject, created_at, expires_at, revoked_at, attributes
336 FROM minco_sessions WHERE token_hash = $1",
337 )
338 .bind(token_hash.as_bytes().as_slice())
339 .fetch_optional(&self.pool)
340 .await
341 .map_err(session_store_error)?;
342 row.as_ref().map(decode_session).transpose()
343 }
344
345 async fn revoke(&self, id: SessionId, at: DateTime<Utc>) -> Result<bool, SessionError> {
346 let result = sqlx::query(
347 "UPDATE minco_sessions SET revoked_at = $2
348 WHERE id = $1 AND revoked_at IS NULL",
349 )
350 .bind(id.0)
351 .bind(at)
352 .execute(&self.pool)
353 .await
354 .map_err(session_store_error)?;
355 Ok(result.rows_affected() == 1)
356 }
357
358 async fn revoke_subject(
359 &self,
360 subject: &str,
361 at: DateTime<Utc>,
362 ) -> Result<usize, SessionError> {
363 let result = sqlx::query(
364 "UPDATE minco_sessions SET revoked_at = $2
365 WHERE subject = $1 AND revoked_at IS NULL",
366 )
367 .bind(subject)
368 .bind(at)
369 .execute(&self.pool)
370 .await
371 .map_err(session_store_error)?;
372 usize::try_from(result.rows_affected())
373 .map_err(|error| session_store_error(error.to_string()))
374 }
375}
376
377fn decode_session(row: &sqlx::postgres::PgRow) -> Result<SessionRecord, SessionError> {
378 let attributes: serde_json::Value = row.try_get("attributes").map_err(session_store_error)?;
379 Ok(SessionRecord {
380 id: SessionId(row.try_get("id").map_err(session_store_error)?),
381 subject: row.try_get("subject").map_err(session_store_error)?,
382 created_at: row.try_get("created_at").map_err(session_store_error)?,
383 expires_at: row.try_get("expires_at").map_err(session_store_error)?,
384 revoked_at: row.try_get("revoked_at").map_err(session_store_error)?,
385 attributes: serde_json::from_value(attributes).map_err(session_store_error)?,
386 })
387}
388
389#[derive(Debug, Clone)]
390pub struct PostgresIdempotencyStore {
391 pool: PgPool,
392}
393
394impl PostgresIdempotencyStore {
395 pub const fn new(pool: PgPool) -> Self {
396 Self { pool }
397 }
398}
399
400#[async_trait]
401impl IdempotencyStore for PostgresIdempotencyStore {
402 async fn get(
403 &self,
404 key: &IdempotencyKey,
405 ) -> Result<Option<IdempotencyRecord>, IdempotencyError> {
406 let row = sqlx::query(
407 "SELECT fingerprint, response, completed_at
408 FROM minco_idempotency WHERE key = $1 AND state = 'completed'",
409 )
410 .bind(key.as_str())
411 .fetch_optional(&self.pool)
412 .await
413 .map_err(idempotency_store_error)?;
414 row.as_ref().map(decode_completed).transpose()
415 }
416
417 async fn begin(
418 &self,
419 key: IdempotencyKey,
420 fingerprint: RequestFingerprint,
421 now: DateTime<Utc>,
422 stale_after: TimeDelta,
423 ) -> Result<BeginOutcome, IdempotencyError> {
424 validate_claim_timeout(stale_after)?;
425 let lease_id = Uuid::now_v7();
426 let mut transaction = self.pool.begin().await.map_err(idempotency_store_error)?;
427 let inserted = sqlx::query(
428 "INSERT INTO minco_idempotency
429 (key, fingerprint, state, lease_id, started_at)
430 VALUES ($1, $2, 'in_progress', $3, $4)
431 ON CONFLICT(key) DO NOTHING",
432 )
433 .bind(key.as_str())
434 .bind(fingerprint.as_str())
435 .bind(lease_id)
436 .bind(now)
437 .execute(&mut *transaction)
438 .await
439 .map_err(idempotency_store_error)?
440 .rows_affected()
441 == 1;
442 if inserted {
443 transaction
444 .commit()
445 .await
446 .map_err(idempotency_store_error)?;
447 return Ok(BeginOutcome::Started(IdempotencyLease {
448 key,
449 fingerprint,
450 lease_id,
451 started_at: now,
452 }));
453 }
454 let row = sqlx::query(
455 "SELECT fingerprint, state, started_at, response, completed_at
456 FROM minco_idempotency WHERE key = $1 FOR UPDATE",
457 )
458 .bind(key.as_str())
459 .fetch_one(&mut *transaction)
460 .await
461 .map_err(idempotency_store_error)?;
462 let stored_fingerprint = RequestFingerprint::parse(
463 row.try_get::<String, _>("fingerprint")
464 .map_err(idempotency_store_error)?,
465 )?;
466 if stored_fingerprint != fingerprint {
467 transaction
468 .commit()
469 .await
470 .map_err(idempotency_store_error)?;
471 return Ok(BeginOutcome::Conflict);
472 }
473 let state: String = row.try_get("state").map_err(idempotency_store_error)?;
474 if state == "completed" {
475 let record = decode_completed(&row)?;
476 transaction
477 .commit()
478 .await
479 .map_err(idempotency_store_error)?;
480 return Ok(BeginOutcome::Replay(record));
481 }
482 let started_at: DateTime<Utc> =
483 row.try_get("started_at").map_err(idempotency_store_error)?;
484 if started_at > now - stale_after {
485 transaction
486 .commit()
487 .await
488 .map_err(idempotency_store_error)?;
489 return Ok(BeginOutcome::InProgress { started_at });
490 }
491 sqlx::query(
492 "UPDATE minco_idempotency
493 SET lease_id = $2, started_at = $3, response = NULL, completed_at = NULL
494 WHERE key = $1 AND state = 'in_progress'",
495 )
496 .bind(key.as_str())
497 .bind(lease_id)
498 .bind(now)
499 .execute(&mut *transaction)
500 .await
501 .map_err(idempotency_store_error)?;
502 transaction
503 .commit()
504 .await
505 .map_err(idempotency_store_error)?;
506 Ok(BeginOutcome::Started(IdempotencyLease {
507 key,
508 fingerprint,
509 lease_id,
510 started_at: now,
511 }))
512 }
513
514 async fn complete(
515 &self,
516 lease: IdempotencyLease,
517 response: serde_json::Value,
518 completed_at: DateTime<Utc>,
519 ) -> Result<IdempotencyRecord, IdempotencyError> {
520 let result = sqlx::query(
521 "UPDATE minco_idempotency
522 SET state = 'completed', response = $4, completed_at = $5, lease_id = NULL
523 WHERE key = $1 AND fingerprint = $2 AND state = 'in_progress' AND lease_id = $3",
524 )
525 .bind(lease.key.as_str())
526 .bind(lease.fingerprint.as_str())
527 .bind(lease.lease_id)
528 .bind(response.clone())
529 .bind(completed_at)
530 .execute(&self.pool)
531 .await
532 .map_err(idempotency_store_error)?;
533 if result.rows_affected() != 1 {
534 return Err(IdempotencyError::InvalidLease);
535 }
536 Ok(IdempotencyRecord {
537 fingerprint: lease.fingerprint,
538 response,
539 created_at: completed_at,
540 })
541 }
542
543 async fn abort(&self, lease: &IdempotencyLease) -> Result<bool, IdempotencyError> {
544 let result = sqlx::query(
545 "DELETE FROM minco_idempotency
546 WHERE key = $1 AND fingerprint = $2 AND state = 'in_progress' AND lease_id = $3",
547 )
548 .bind(lease.key.as_str())
549 .bind(lease.fingerprint.as_str())
550 .bind(lease.lease_id)
551 .execute(&self.pool)
552 .await
553 .map_err(idempotency_store_error)?;
554 Ok(result.rows_affected() == 1)
555 }
556}
557
558fn decode_completed(row: &sqlx::postgres::PgRow) -> Result<IdempotencyRecord, IdempotencyError> {
559 Ok(IdempotencyRecord {
560 fingerprint: RequestFingerprint::parse(
561 row.try_get::<String, _>("fingerprint")
562 .map_err(idempotency_store_error)?,
563 )?,
564 response: row.try_get("response").map_err(idempotency_store_error)?,
565 created_at: row
566 .try_get("completed_at")
567 .map_err(idempotency_store_error)?,
568 })
569}
570
571#[derive(Debug, Clone)]
572pub struct PostgresAuditSink {
573 pool: PgPool,
574}
575
576impl PostgresAuditSink {
577 pub const fn new(pool: PgPool) -> Self {
578 Self { pool }
579 }
580}
581
582#[async_trait]
583impl AuditSink for PostgresAuditSink {
584 async fn append(&self, event: AuditEvent) -> Result<(), AuditError> {
585 if event.action.trim().is_empty() || event.resource_id.trim().is_empty() {
586 return Err(AuditError::InvalidEvent);
587 }
588 let metadata = serde_json::to_value(event.metadata)
589 .map_err(|error| AuditError::Append(error.to_string()))?;
590 sqlx::query(
591 "INSERT INTO minco_audit
592 (id, action, resource_type, resource_id, actor_subject, correlation_id,
593 occurred_at, metadata)
594 VALUES ($1, $2, $3, $4, $5, $6, $7, $8)",
595 )
596 .bind(event.id)
597 .bind(event.action)
598 .bind(event.resource_type)
599 .bind(event.resource_id)
600 .bind(event.actor_subject)
601 .bind(event.correlation_id)
602 .bind(event.occurred_at)
603 .bind(metadata)
604 .execute(&self.pool)
605 .await
606 .map_err(|error| AuditError::Append(error.to_string()))?;
607 Ok(())
608 }
609}
610
611pub async fn migrate_plugin_storage(pool: &PgPool) -> Result<(), sqlx::migrate::MigrateError> {
612 let mut migrator = sqlx::migrate!("migrations/plugins");
613 migrator.dangerous_set_table_name("_minco_plugin_storage_migrations");
614 migrator.run(pool).await
615}
616
617fn validate_event(event: &DomainEvent) -> Result<(), EventError> {
618 if event.event_type.trim().is_empty()
619 || event.aggregate_type.trim().is_empty()
620 || event.aggregate_id.trim().is_empty()
621 {
622 Err(EventError::InvalidEvent)
623 } else {
624 Ok(())
625 }
626}
627
628fn validate_claim(worker_id: &str, claim_expires_at: DateTime<Utc>) -> Result<(), EventError> {
629 if worker_id.trim().is_empty() || claim_expires_at <= Utc::now() {
630 Err(EventError::InvalidClaim)
631 } else {
632 Ok(())
633 }
634}
635
636fn event_infrastructure_error(error: impl std::fmt::Display) -> EventError {
637 EventError::Infrastructure(error.to_string())
638}
639
640fn session_store_error(error: impl std::fmt::Display) -> SessionError {
641 SessionError::Store(error.to_string())
642}
643
644fn idempotency_store_error(error: impl std::fmt::Display) -> IdempotencyError {
645 IdempotencyError::Store(error.to_string())
646}
647
648fn is_unique_violation(error: &sqlx::Error) -> bool {
649 error
650 .as_database_error()
651 .is_some_and(sqlx::error::DatabaseError::is_unique_violation)
652}
653
654#[cfg(test)]
655mod tests {
656 use super::*;
657 use minco_plugin_sessions::{CreateSession, SessionService};
658 use std::{
659 collections::{BTreeMap, BTreeSet},
660 sync::{Arc, OnceLock},
661 };
662
663 fn test_lock() -> &'static tokio::sync::Mutex<()> {
664 static LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
665 LOCK.get_or_init(|| tokio::sync::Mutex::new(()))
666 }
667
668 async fn pool() -> Option<PgPool> {
669 let url = std::env::var("MINCO_TEST_POSTGRES_URL").ok()?;
670 let pool = PgPool::connect(&url).await.ok()?;
671 migrate_plugin_storage(&pool).await.ok()?;
672 Some(pool)
673 }
674
675 #[tokio::test]
676 async fn postgres_plugin_stores_are_behavioral_when_database_is_configured() {
677 let _guard = test_lock().lock().await;
678 let Some(pool) = pool().await else {
679 eprintln!("MINCO_TEST_POSTGRES_URL not set; PostgreSQL adapter proof skipped");
680 return;
681 };
682 sqlx::raw_sql(
683 "TRUNCATE minco_outbox, minco_sessions, minco_idempotency, minco_audit RESTART IDENTITY",
684 )
685 .execute(&pool)
686 .await
687 .unwrap();
688 let migration_count: i64 =
689 sqlx::query_scalar("SELECT COUNT(*) FROM _minco_plugin_storage_migrations")
690 .fetch_one(&pool)
691 .await
692 .unwrap();
693 assert_eq!(migration_count, 2);
694
695 let outbox = PostgresOutboxStore::new(pool.clone());
696 let event = DomainEvent::new(
697 "feedback.created",
698 "feedback",
699 "one",
700 Uuid::now_v7(),
701 serde_json::json!({"id": "one"}),
702 );
703 outbox
704 .enqueue(OutboxRecord::pending(event.clone()))
705 .await
706 .unwrap();
707 let claimed = outbox
708 .claim_pending("worker-a", 10, Utc::now() + TimeDelta::minutes(1))
709 .await
710 .unwrap();
711 assert_eq!(claimed.len(), 1);
712 outbox.mark_published(event.id, "worker-a").await.unwrap();
713
714 let session_store = Arc::new(PostgresSessionStore::new(pool.clone()));
715 let sessions = SessionService::new(session_store.clone());
716 let issued = sessions
717 .issue(CreateSession {
718 subject: "postgres-subject".into(),
719 ttl: TimeDelta::minutes(5),
720 attributes: BTreeMap::new(),
721 })
722 .await
723 .unwrap();
724 assert_eq!(
725 sessions.resolve(&issued.token).await.unwrap().subject,
726 "postgres-subject"
727 );
728 assert_eq!(
729 session_store
730 .revoke_subject("postgres-subject", Utc::now())
731 .await
732 .unwrap(),
733 1
734 );
735 assert!(matches!(
736 sessions.resolve(&issued.token).await,
737 Err(SessionError::Unauthenticated)
738 ));
739
740 let idempotency = PostgresIdempotencyStore::new(pool.clone());
741 let key = IdempotencyKey::parse("postgres-request").unwrap();
742 let fingerprint =
743 RequestFingerprint::from_serializable(&serde_json::json!({"request": 1})).unwrap();
744 let BeginOutcome::Started(lease) = idempotency
745 .begin(
746 key.clone(),
747 fingerprint.clone(),
748 Utc::now(),
749 TimeDelta::minutes(5),
750 )
751 .await
752 .unwrap()
753 else {
754 panic!("expected idempotency lease");
755 };
756 idempotency
757 .complete(lease, serde_json::json!({"status": 201}), Utc::now())
758 .await
759 .unwrap();
760 assert!(matches!(
761 idempotency
762 .begin(key, fingerprint, Utc::now(), TimeDelta::minutes(5))
763 .await
764 .unwrap(),
765 BeginOutcome::Replay(_)
766 ));
767
768 PostgresAuditSink::new(pool.clone())
769 .append(AuditEvent::new(
770 "feedback.created",
771 "feedback",
772 "one",
773 Uuid::now_v7(),
774 ))
775 .await
776 .unwrap();
777 let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM minco_audit")
778 .fetch_one(&pool)
779 .await
780 .unwrap();
781 assert_eq!(count, 1);
782 }
783
784 #[tokio::test]
785 async fn enqueue_in_rolls_back_with_the_callers_transaction() {
786 let _guard = test_lock().lock().await;
787 let Some(pool) = pool().await else {
788 eprintln!("MINCO_TEST_POSTGRES_URL not set; PostgreSQL adapter proof skipped");
789 return;
790 };
791 sqlx::query("TRUNCATE minco_outbox")
792 .execute(&pool)
793 .await
794 .unwrap();
795
796 let store = PostgresOutboxStore::new(pool.clone());
797 let event = DomainEvent::new(
798 "feedback.created",
799 "feedback",
800 "rollback",
801 Uuid::now_v7(),
802 serde_json::json!({"id": "rollback"}),
803 );
804 let mut transaction = pool.begin().await.unwrap();
805 store
806 .enqueue_in(&mut transaction, OutboxRecord::pending(event.clone()))
807 .await
808 .unwrap();
809 transaction.rollback().await.unwrap();
810
811 let persisted: bool =
812 sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM minco_outbox WHERE event_id = $1)")
813 .bind(event.id)
814 .fetch_one(&pool)
815 .await
816 .unwrap();
817 assert!(!persisted);
818 }
819
820 #[tokio::test]
821 async fn concurrent_outbox_claims_are_disjoint() {
822 let _guard = test_lock().lock().await;
823 let Some(pool) = pool().await else {
824 eprintln!("MINCO_TEST_POSTGRES_URL not set; PostgreSQL adapter proof skipped");
825 return;
826 };
827 sqlx::query("TRUNCATE minco_outbox")
828 .execute(&pool)
829 .await
830 .unwrap();
831
832 let store = PostgresOutboxStore::new(pool.clone());
833 for aggregate_id in ["claim-one", "claim-two"] {
834 store
835 .enqueue(OutboxRecord::pending(DomainEvent::new(
836 "feedback.created",
837 "feedback",
838 aggregate_id,
839 Uuid::now_v7(),
840 serde_json::json!({"id": aggregate_id}),
841 )))
842 .await
843 .unwrap();
844 }
845 let expires_at = Utc::now() + TimeDelta::minutes(1);
846 let (first, second) = tokio::join!(
847 store.claim_pending("worker-one", 1, expires_at),
848 store.claim_pending("worker-two", 1, expires_at),
849 );
850 let first = first.unwrap();
851 let second = second.unwrap();
852 assert_eq!(first.len(), 1);
853 assert_eq!(second.len(), 1);
854 let claimed = first
855 .into_iter()
856 .chain(second)
857 .map(|record| record.event.id)
858 .collect::<BTreeSet<_>>();
859 assert_eq!(claimed.len(), 2);
860 }
861
862 #[tokio::test]
863 async fn concurrent_idempotency_begin_has_one_owner() {
864 let _guard = test_lock().lock().await;
865 let Some(pool) = pool().await else {
866 eprintln!("MINCO_TEST_POSTGRES_URL not set; PostgreSQL adapter proof skipped");
867 return;
868 };
869 sqlx::query("TRUNCATE minco_idempotency")
870 .execute(&pool)
871 .await
872 .unwrap();
873
874 let store = PostgresIdempotencyStore::new(pool);
875 let key = IdempotencyKey::parse("postgres-concurrent-request").unwrap();
876 let fingerprint =
877 RequestFingerprint::from_serializable(&serde_json::json!({"request": 2})).unwrap();
878 let now = Utc::now();
879 let (first, second) = tokio::join!(
880 store.begin(key.clone(), fingerprint.clone(), now, TimeDelta::minutes(5),),
881 store.begin(key, fingerprint, now, TimeDelta::minutes(5)),
882 );
883 let outcomes = [first.unwrap(), second.unwrap()];
884 assert_eq!(
885 outcomes
886 .iter()
887 .filter(|outcome| matches!(outcome, BeginOutcome::Started(_)))
888 .count(),
889 1
890 );
891 assert_eq!(
892 outcomes
893 .iter()
894 .filter(|outcome| matches!(outcome, BeginOutcome::InProgress { .. }))
895 .count(),
896 1
897 );
898 }
899}