Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .changeset/quiet-databases-close.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
'@electric-sql/pglite': patch
---

Prevent `close()` from hanging when called while a query or transaction is in
flight. Work started before `close()` is allowed to finish, while later
operations are rejected.
33 changes: 31 additions & 2 deletions packages/pglite/src/pglite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ export class PGlite
#ready = false
#closing = false
#closed = false
#closePromise?: Promise<void>
#relaxedDurability = false

readonly waitReady: Promise<void>
Expand Down Expand Up @@ -784,10 +785,32 @@ export class PGlite
* Close the database
* @returns A promise that resolves when the database is closed
*/
async close() {
await this._checkReady()
close() {
if (!this.#closePromise) {
this.#closePromise = this.#close()
}
return this.#closePromise
}

async #close() {
if (!this.#ready) {
await this.waitReady
}
// Claim the lifecycle transition synchronously once initialization is
// complete so later operations cannot queue behind close().
this.#closing = true

// Let operations that passed _checkReady() before close() was called
// enqueue on the mutexes. Operations started after close() are rejected
// synchronously by the #closing flag set above.
await Promise.resolve()

await this._runExclusiveTransaction(() =>
this._runExclusiveQuery(() => this.#closeExclusive()),
)
}

async #closeExclusive() {
// Close all extensions
for (const closeFn of this.#extensionsClose) {
await closeFn()
Expand Down Expand Up @@ -893,6 +916,12 @@ export class PGlite
// Starting the database can take a while and it might not be ready yet
// We'll wait for it to be ready before continuing
await this.waitReady
if (this.#closing) {
throw new Error('PGlite is closing')
}
if (this.#closed) {
throw new Error('PGlite is closed')
}
}
}

Expand Down
110 changes: 110 additions & 0 deletions packages/pglite/tests/targets/runtimes/node-close.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
import { spawn } from 'node:child_process'
import { describe, expect, it } from 'vitest'

const pgliteUrl = new URL('../../../dist/index.js', import.meta.url).href

async function expectChildToExitCleanly(script) {
const result = await new Promise((resolve) => {
const child = spawn(
process.execPath,
['--input-type=module', '--eval', script, pgliteUrl],
{ stdio: ['ignore', 'pipe', 'pipe'] },
)
let stderr = ''
let settled = false

const finish = (result) => {
if (settled) return
settled = true
clearTimeout(timeout)
resolve(result)
}

child.stderr.setEncoding('utf8')
child.stderr.on('data', (chunk) => {
stderr += chunk
stderr = stderr.slice(-4_000)
})

const timeout = setTimeout(() => {
child.kill('SIGKILL')
finish({ code: null, stderr: 'PGlite close timed out' })
}, 5_000)

child.on('error', (error) => {
finish({ code: null, stderr: error.message })
})
child.on('exit', (code) => {
finish({ code, stderr })
})
})

expect(result).toEqual({ code: 0, stderr: '' })
}

describe('close', () => {
it('waits for an in-flight query before shutting down', async () => {
await expectChildToExitCleanly(`
const { PGlite } = await import(process.argv[1])
const db = new PGlite()

await db.exec('CREATE TABLE t (workflow_name TEXT, run_id TEXT)')
await db.exec("INSERT INTO t VALUES ('agentic-loop', 'run-1')")

const query = db.query(
'DELETE FROM t WHERE workflow_name = $1 AND run_id = $2',
['agentic-loop', 'run-1'],
)
const close = db.close()

await Promise.all([query, close])
`)
}, 10_000)

it('waits for an active transaction before shutting down', async () => {
await expectChildToExitCleanly(`
const { PGlite } = await import(process.argv[1])
const db = new PGlite()

await db.exec('CREATE TABLE t (value INTEGER)')

let markTransactionStarted
const transactionStarted = new Promise((resolve) => {
markTransactionStarted = resolve
})
let resumeTransaction
const transactionGate = new Promise((resolve) => {
resumeTransaction = resolve
})
const events = []

const transaction = db.transaction(async (tx) => {
markTransactionStarted()
await transactionGate
await tx.query('INSERT INTO t VALUES (1)')
events.push('transaction')
})

await transactionStarted
const close = db.close().then(() => events.push('close'))
resumeTransaction()
await Promise.all([transaction, close])

if (events.join(',') !== 'transaction,close') {
throw new Error('close did not wait for the active transaction')
}
`)
}, 10_000)

it('closes once during initialization and rejects later queries', async () => {
const { PGlite } = await import(pgliteUrl)
const db = new PGlite()

const firstClose = db.close()
expect(db.close()).toBe(firstClose)
await expect(db.query('SELECT 1')).rejects.toThrow('PGlite is closing')

await firstClose
expect(db.closed).toBe(true)
}, 10_000)
})