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 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}