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