Skip to content

Commit 15b22ac

Browse files
committed
feat: add PostgreSQL wake-ups
1 parent f0a5ef8 commit 15b22ac

12 files changed

Lines changed: 412 additions & 17 deletions

‎CHANGELOG.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,8 @@
3131
- Add PostgreSQL 14+ through the optional `pg` driver with pooled transactions,
3232
row-locked sequence allocation, portable schema and set queries, diagnostics,
3333
and real-server integration coverage.
34+
- Add opt-in role-specific PostgreSQL wake-ups with one event-driven listener
35+
per runtime and polling as the correctness fallback.
3436

3537
## 0.1.0 - 2026-08-13
3638

‎README.md‎

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,11 @@ const database = postgresql({
140140
connectionString,
141141
maximumConnections: 10,
142142
})
143+
144+
const runtime = configureSolidObjects({
145+
database,
146+
wakeUp: database.wakeUp(),
147+
})
143148
```
144149

145150
PostgreSQL uses a bounded `pg` pool, 64-bit database timestamps and sequences,
@@ -148,6 +153,13 @@ SQLite. Keep `pg` at 8.23 or newer within the supported major. Portable
148153
`DatabaseConnection` SQL uses `?` parameters; write `??` when a PostgreSQL query
149154
needs the literal JSON existence operator.
150155

156+
`database.wakeUp()` is opt-in. It uses one event-driven PostgreSQL client per
157+
runtime to listen on role-specific channels and wake every matching local
158+
waiter. Create it in every process that should send or receive notifications.
159+
Polling remains the fallback if a notification is missed or the listener
160+
reconnects. Because `LISTEN` is session-scoped, use a direct connection or
161+
session pooling rather than transaction pooling for this client.
162+
151163
Authorization is deny-by-default. Actor IDs identify actors; they are not
152164
capabilities.
153165

@@ -158,9 +170,10 @@ instances to finish, so no replacement can outlive the runtime.
158170

159171
The default in-process wake-up adapter interrupts role polling as soon as this
160172
runtime commits new work. Polling remains the correctness fallback, so a missed
161-
or failed signal costs latency rather than losing work. Multi-process hosts can
162-
provide a `WakeUpAdapter` backed by their existing notification system without
163-
adding a required broker to the default SQLite stack.
173+
or failed signal costs latency rather than losing work. Multi-process hosts not
174+
using PostgreSQL notifications can provide a `WakeUpAdapter` backed by their
175+
existing notification system without adding a required broker to the default
176+
SQLite stack.
164177

165178
## Call it like a local object
166179

‎docs/architecture.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,13 @@ transaction on one checked-out client, stores timestamps and sequences as
3535
64-bit integers, and locks an actor's instance row while allocating mailbox
3636
sequences. Both adapters use database time and the same fencing predicates.
3737

38+
PostgreSQL notifications are an opt-in latency layer. One event-driven client
39+
per runtime listens on role-specific channels before the worker checks durable
40+
state, which closes the listener-startup race without Ruby's connection per
41+
blocking thread. A notification advances a process-local role generation and
42+
wakes every matching waiter. Reconnection and notification loss fall back to
43+
the ordinary polling interval.
44+
3845
A worker claims one actor globally, then preferentially drains up to
3946
`maxMessagesPerActivationPass` ready turns for that instance. Reaching the cap
4047
moves only that actor's already-due ready memberships to current database time.

‎docs/parity.md‎

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -54,14 +54,15 @@ Reference: Ruby `solid_objects` 0.12.0 at commit `a01b6f5`.
5454

5555
## Databases and wake-up
5656

57-
| Capability | Status | TypeScript shape or remaining work |
58-
| ---------------------------- | ------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
59-
| SQLite | Native | Uses built-in `node:sqlite`, serialized process-local access, foreign keys, strict tables, and database time. |
60-
| PostgreSQL | Native | Optional `pg` 8.23 peer, bounded pooling, 64-bit schema, row-locked sequence allocation, portable set queries, server checks, and a real PostgreSQL integration suite. |
61-
| MySQL | Planned | Add a driver-neutral pool interface, MySQL schema, locking behavior, integration suite, engine/version checks, and client compatibility tests. |
62-
| Durable polling fallback | Native | Every role progresses without a notification service. |
63-
| In-process wake-up | Native | A generation-based default adapter prevents claim-to-wait signal loss; commits wake role-specific waiters and polling remains the fallback. |
64-
| PostgreSQL and Redis wake-up | Planned | Add optional adapters; neither may become a hidden required dependency. |
57+
| Capability | Status | TypeScript shape or remaining work |
58+
| ------------------------ | ------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
59+
| SQLite | Native | Uses built-in `node:sqlite`, serialized process-local access, foreign keys, strict tables, and database time. |
60+
| PostgreSQL | Native | Optional `pg` 8.23 peer, bounded pooling, 64-bit schema, row-locked sequence allocation, portable set queries, server checks, and a real PostgreSQL integration suite. |
61+
| MySQL | Planned | Add a driver-neutral pool interface, MySQL schema, locking behavior, integration suite, engine/version checks, and client compatibility tests. |
62+
| Durable polling fallback | Native | Every role progresses without a notification service. |
63+
| In-process wake-up | Native | A generation-based default adapter prevents claim-to-wait signal loss; commits wake role-specific waiters and polling remains the fallback. |
64+
| PostgreSQL wake-up | Native | `database.wakeUp()` uses one dedicated event-driven client, role-specific `LISTEN/NOTIFY`, generation fencing, reconnectable listeners, and durable polling fallback. |
65+
| Redis wake-up | Planned | Add an optional Pub/Sub adapter without making Redis a hidden required dependency. |
6566

6667
## Realtime and browser behavior
6768

‎src/broadcast-worker.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ export class BroadcastWorker {
3737
async run(signal: AbortSignal): Promise<void> {
3838
await this.ensureRegistered()
3939
while (!signal.aborted && !this.stopping) {
40-
const wakeUp = this.runtime.settings.wakeUp.watch("broadcasts")
40+
const wakeUp = await this.runtime.settings.wakeUp.watch("broadcasts")
4141
const processed = await this.runOnce()
4242
if (processed === 0) {
4343
await wakeUp.wait({

‎src/database/postgresql.ts‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,14 @@
11
import { Pool, TypeOverrides, types, type PoolClient, type PoolConfig } from "pg"
22
import { postgresqlSql } from "./postgresql-sql.js"
33
import type { Database, DatabaseConnection, RunResult } from "./types.js"
4+
import { PostgreSQLWakeUpAdapter, type PostgreSQLWakeUpFailure } from "../wake-up/postgresql.js"
5+
6+
export {
7+
PostgreSQLWakeUpAdapter,
8+
postgresqlWakeUp,
9+
type PostgreSQLWakeUpFailure,
10+
type PostgreSQLWakeUpOptions,
11+
} from "../wake-up/postgresql.js"
412

513
export interface PostgreSQLDatabaseOptions {
614
connectionString: string
@@ -11,6 +19,12 @@ export interface PostgreSQLDatabaseOptions {
1119
onPoolError?: (error: Error) => void
1220
}
1321

22+
export interface PostgreSQLDatabaseWakeUpOptions {
23+
channelPrefix?: string
24+
applicationName?: string
25+
onListenerError?: (failure: PostgreSQLWakeUpFailure) => void
26+
}
27+
1428
class PostgreSQLConnection implements DatabaseConnection {
1529
constructor(private readonly client: PoolClient) {}
1630

@@ -46,12 +60,14 @@ export class PostgreSQLDatabase implements Database {
4660
readonly family = "postgresql" as const
4761
readonly schemaIdentity = "solid-objects-node-v1"
4862
private readonly pool: Pool
63+
private readonly connectionString: string
4964
private closed = false
5065

5166
constructor(options: PostgreSQLDatabaseOptions) {
5267
if (options.connectionString.length === 0) {
5368
throw new TypeError("PostgreSQL connectionString must not be empty")
5469
}
70+
this.connectionString = options.connectionString
5571
if (
5672
options.maximumConnections !== undefined &&
5773
(!Number.isSafeInteger(options.maximumConnections) || options.maximumConnections < 1)
@@ -125,6 +141,13 @@ export class PostgreSQLDatabase implements Database {
125141
this.closed = true
126142
await this.pool.end()
127143
}
144+
145+
wakeUp(options: PostgreSQLDatabaseWakeUpOptions = {}): PostgreSQLWakeUpAdapter {
146+
return new PostgreSQLWakeUpAdapter({
147+
connectionString: this.connectionString,
148+
...options,
149+
})
150+
}
128151
}
129152

130153
function postgresqlTypes(): TypeOverrides {

‎src/effect-worker.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ export class EffectWorker {
3737
async run(signal: AbortSignal): Promise<void> {
3838
await this.ensureRegistered()
3939
while (!signal.aborted && !this.stopping) {
40-
const wakeUp = this.runtime.settings.wakeUp.watch("effects")
40+
const wakeUp = await this.runtime.settings.wakeUp.watch("effects")
4141
const processed = await this.runOnce()
4242
if (processed === 0) {
4343
await wakeUp.wait({

‎src/reminder-scheduler.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ export class ReminderScheduler {
5050
async run(signal: AbortSignal): Promise<void> {
5151
await this.ensureRegistered()
5252
while (!signal.aborted && !this.stopping) {
53-
const wakeUp = this.runtime.settings.wakeUp.watch("reminders")
53+
const wakeUp = await this.runtime.settings.wakeUp.watch("reminders")
5454
const processed = await this.runOnce()
5555
if (processed === 0) {
5656
await wakeUp.wait({

‎src/wake-up.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ export interface WakeUpWatch {
1010
}
1111

1212
export interface WakeUpAdapter {
13-
watch(role: WakeUpRole): WakeUpWatch
13+
watch(role: WakeUpRole): WakeUpWatch | Promise<WakeUpWatch>
1414
notify(role: WakeUpRole): void | Promise<void>
1515
close(): void | Promise<void>
1616
}

0 commit comments

Comments
 (0)