@@ -2,7 +2,12 @@ import { test } from "@fedify/fixture";
22import { PostgresMessageQueue } from "@fedify/postgres/mq" ;
33import { getRandomKey , testMessageQueue } from "@fedify/testing" ;
44import * as temporal from "@js-temporal/polyfill" ;
5- import { deepStrictEqual , rejects } from "node:assert/strict" ;
5+ import {
6+ deepStrictEqual ,
7+ notStrictEqual ,
8+ rejects ,
9+ strictEqual ,
10+ } from "node:assert/strict" ;
611import process from "node:process" ;
712import { test as nodeTest } from "node:test" ;
813import postgres from "postgres" ;
@@ -11,6 +16,107 @@ const Temporal = globalThis.Temporal ?? temporal.Temporal;
1116
1217const dbUrl = process . env . POSTGRES_URL ;
1318
19+ // Gate both attempts so concurrent callers cannot accidentally pass by running
20+ // separate initializations that happen to finish before the assertions.
21+ for (
22+ const { phase, failAt, initialized } of [
23+ { phase : "first DDL" , failAt : 1 , initialized : false } ,
24+ { phase : "partial DDL" , failAt : 2 , initialized : false } ,
25+ { phase : "JSON probe" , failAt : 4 , initialized : false } ,
26+ { phase : "JSON probe without DDL" , failAt : 1 , initialized : true } ,
27+ ]
28+ ) {
29+ test ( `PostgresMessageQueue.initialize() retries after ${ phase } failure` , async ( ) => {
30+ const failure = Object . assign ( new Error ( "statement timeout" ) , {
31+ code : "57014" ,
32+ } ) ;
33+ const failureStarted = Promise . withResolvers < void > ( ) ;
34+ const failedQuery = Promise . withResolvers < { test : string } [ ] > ( ) ;
35+ const retryQuery = Promise . withResolvers < { test : string } [ ] > ( ) ;
36+ let queries = 0 ;
37+ const statements : string [ ] = [ ] ;
38+ const result = [ { test : '{"foo":1}' } ] ;
39+ const sql = Object . assign (
40+ ( strings : TemplateStringsArray | string ) => {
41+ if ( ! Array . isArray ( strings ) ) return strings ;
42+ statements . push ( strings . join ( "" ) ) ;
43+ queries ++ ;
44+ if ( queries === failAt ) {
45+ failureStarted . resolve ( ) ;
46+ return failedQuery . promise ;
47+ }
48+ if ( queries === failAt + 1 ) return retryQuery . promise ;
49+ return Promise . resolve ( result ) ;
50+ } ,
51+ { json : ( value : unknown ) => value } ,
52+ ) as unknown as postgres . Sql ;
53+ const mq = new PostgresMessageQueue ( sql , { initialized } ) ;
54+ const first = mq . initialize ( ) ;
55+ const concurrent = mq . initialize ( ) ;
56+ strictEqual ( first , concurrent , "pending callers must share one promise" ) ;
57+ const outcomes = Promise . allSettled ( [ first , concurrent ] ) ;
58+ await failureStarted . promise ;
59+ const expectedStatement = phase === "first DDL"
60+ ? "CREATE TABLE"
61+ : phase === "partial DDL"
62+ ? "ALTER TABLE"
63+ : "SELECT" ;
64+ strictEqual ( statements . at ( - 1 ) ?. includes ( expectedStatement ) , true ) ;
65+ failedQuery . reject ( failure ) ;
66+ const rejected = await outcomes ;
67+ for ( const outcome of rejected ) {
68+ strictEqual ( outcome . status , "rejected" ) ;
69+ if ( outcome . status === "rejected" ) strictEqual ( outcome . reason , failure ) ;
70+ }
71+ await new Promise ( ( resolve ) => setTimeout ( resolve , 0 ) ) ;
72+ strictEqual ( queries , failAt , "a rejection must not automatically retry" ) ;
73+
74+ const retry = mq . initialize ( ) ;
75+ const concurrentRetry = mq . initialize ( ) ;
76+ notStrictEqual ( retry , first , "a later call must start a new attempt" ) ;
77+ strictEqual ( retry , concurrentRetry ) ;
78+ strictEqual ( queries , failAt + 1 ) ;
79+ retryQuery . resolve ( result ) ;
80+ await Promise . all ( [ retry , concurrentRetry ] ) ;
81+ const expectedQueries = failAt + ( initialized ? 1 : 4 ) ;
82+ strictEqual ( queries , expectedQueries ) ;
83+ await mq . initialize ( ) ;
84+ strictEqual (
85+ queries ,
86+ expectedQueries ,
87+ "successful initialization is cached" ,
88+ ) ;
89+ } ) ;
90+ }
91+
92+ test ( "PostgresMessageQueue.enqueue() recovers after initialization failure" , async ( ) => {
93+ const failure = new Error ( "database unavailable" ) ;
94+ let queries = 0 ;
95+ let notifications = 0 ;
96+ const sql = Object . assign (
97+ ( strings : TemplateStringsArray | string ) => {
98+ if ( ! Array . isArray ( strings ) ) return strings ;
99+ queries ++ ;
100+ if ( queries === 1 ) return Promise . reject ( failure ) ;
101+ return Promise . resolve ( [ { test : '{"foo":1}' } ] ) ;
102+ } ,
103+ {
104+ json : ( value : unknown ) => value ,
105+ notify : ( ) => {
106+ notifications ++ ;
107+ return Promise . resolve ( ) ;
108+ } ,
109+ } ,
110+ ) as unknown as postgres . Sql ;
111+ const mq = new PostgresMessageQueue ( sql ) ;
112+ await rejects ( mq . enqueue ( "first" ) , ( error : unknown ) => error === failure ) ;
113+ strictEqual ( queries , 1 ) ;
114+ strictEqual ( notifications , 0 ) ;
115+ await mq . enqueue ( "second" ) ;
116+ strictEqual ( queries , 6 , "retry runs four initialization queries and INSERT" ) ;
117+ strictEqual ( notifications , 1 ) ;
118+ } ) ;
119+
14120test ( "PostgresMessageQueue" , { ignore : dbUrl == null } , ( ) => {
15121 if ( dbUrl == null ) return ; // Bun does not support skip option
16122
0 commit comments