@fuguejs/pg
v0.5.1
Published
PostgreSQL capability adapter for Fugue workflows.
Readme
@fuguejs/pg
PostgreSQL capability adapter for Fugue workflows.
Installation
bun add @fuguejs/pg pg zodpg and zod are peer dependencies — pg provides the connection pool,
zod the schemas every query/queryOne call validates against.
Usage
Register with the host
import { createPgAdapter } from "@fuguejs/pg";
const pgHandle = createPgAdapter({
connectionString: process.env.DATABASE_URL!,
poolSize: 20,
statementTimeoutMs: 15_000,
});
// Pass to SharedInfra capabilities:
const sharedInfra = {
// ... other infra ...
capabilities: [pgHandle],
};Use in a node
import { createFetchNode } from "@fuguejs/framework";
import { z } from "zod";
const UserSchema = z.object({
id: z.string(),
name: z.string(),
email: z.string().email(),
});
const fetchUser = createFetchNode({
id: "fetch-user",
inputSchema: z.object({ userId: z.string() }),
outputSchema: UserSchema,
requires: ["db"] as const,
fetch: async (input, ctx) => {
// ctx.db is typed as PgCapability — non-null, schema-validated
return ctx.db.queryOne(
UserSchema,
"SELECT id, name, email FROM users WHERE id = $1",
[input.userId],
);
},
});Testing with the fake
import { createFakePgCapability } from "@fuguejs/pg";
const fakeDb = createFakePgCapability({
"SELECT * FROM users WHERE id": [
{ id: "1", name: "Alice", email: "[email protected]" },
],
"INSERT INTO orders": { rowCount: 1 },
});
// Use in tests via makeNodeContext:
const ctx = makeNodeContext({
runId: "test-run",
dagId: "test-dag",
capabilities: { db: fakeDb.client },
});API
PgCapability
| Method | Description |
|--------|-------------|
| query<T>(schema, sql, params?) | Execute query, validate all rows against Zod schema |
| queryOne<T>(schema, sql, params?) | Execute query, validate first row (or null) |
| execute(sql, params?) | Execute write, return { rowCount } |
| queryRaw(sql, params?) | Escape hatch: raw unknown[] rows, no validation — prefer query with a schema |
All methods return Result<T, FrameworkError> — no exceptions escape.
createPgAdapter(config)
Creates a CapabilityHandle<"db"> with lifecycle management:
connect(): validates connectivity with SELECT 1close(): drains the connection poolhealthCheck(): SELECT 1, racing a 5s timeout (a hung pool reports unhealthy)
Config options: poolSize (default 10), statementTimeoutMs (default 30000, applied as the pool's statement_timeout), connectionTimeoutMs (default 5000), idleTimeoutMs (default 10000).
createFakePgCapability(routes)
In-memory fake for testing. Routes are matched by exact SQL or longest prefix match.
