@@ -255,23 +255,29 @@ async function leaveOrganization(userId: string, organizationId: string) {
255255 . where ( and ( eq ( member . userId , userId ) , eq ( member . organizationId , organizationId ) ) )
256256}
257257
258+ /**
259+ * The sync event enqueued last for a subscription. `created_at` is the enqueuing transaction's
260+ * start time, compared at microsecond precision, so it orders events from different
261+ * transactions. No test enqueues two of one type for one subscription in a single transaction,
262+ * and a tie fails loudly rather than being broken arbitrarily: `outbox_event` has no per-insert
263+ * sequence to break it with.
264+ */
258265async function latestOutboxEventId ( eventType : string , subscriptionId : string ) {
259- const [ latest ] = await testDatabase
260- . select ( { id : outboxEvent . id } )
266+ const [ latest , previous ] = await testDatabase
267+ . select ( { id : outboxEvent . id , createdAt : sql < string > ` ${ outboxEvent . createdAt } ::text` } )
261268 . from ( outboxEvent )
262269 . where (
263270 and (
264271 eq ( outboxEvent . eventType , eventType ) ,
265272 sql `${ outboxEvent . payload } ->> 'subscriptionId' = ${ subscriptionId } `
266273 )
267274 )
268- . orderBy (
269- desc ( outboxEvent . createdAt ) ,
270- sql `(${ outboxEvent . payload } ->> 'committedAt')::numeric desc nulls last` ,
271- desc ( outboxEvent . id )
272- )
273- . limit ( 1 )
275+ . orderBy ( desc ( outboxEvent . createdAt ) )
276+ . limit ( 2 )
274277 if ( ! latest ) throw new Error ( `No ${ eventType } event for ${ subscriptionId } ` )
278+ if ( previous ?. createdAt === latest . createdAt ) {
279+ throw new Error ( `Two ${ eventType } events for ${ subscriptionId } share one enqueue time` )
280+ }
275281 return latest . id
276282}
277283
0 commit comments