Skip to main content

minco_sqlx_postgres/
plugin_adapters.rs

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    /// Inserts an outbox record into an application adapter's existing
26    /// transaction so the domain mutation and publication intent commit
27    /// atomically.
28    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}