Observability
Three layers, each opt-in: raw query events ($on), aggregated metrics persisted to Postgres ($observe), and a local dashboard over those metrics (npx turbine observe). All of it ships in the box, no agent, no SaaS, no extra dependency.
Query events, db.$on('query')#
Every query emits an event after execution, success or failure. Subscribe with $on, unsubscribe with $off:
db.$on('query', (event) => {
if (event.duration > 200) {
console.warn(`slow: ${event.model}.${event.action} took ${event.duration.toFixed(1)}ms`);
}
});The event shape:
interface QueryEvent {
sql: string; // the generated SQL text
params: unknown[]; // bound parameter values (redacted by default, see below)
duration: number; // wall-clock ms (performance.now)
model: string; // table name, e.g. 'users'
action: string; // operation, e.g. 'findMany'
rows: number; // rows returned / affected
timestamp: Date;
error?: Error; // present when the query failed
}A listener that throws never breaks the query; the error is swallowed (logged when logging: true).
$on differs from middleware: middleware wraps the call, it can observe args and transform results (it cannot change the query); $on is a passive tap that fires after execution with the final SQL, timing, and outcome.
Seeing the parameter values, logQueryParams#
Params are redacted by default. Every entry of event.params is replaced with '[REDACTED]' before any listener sees it, so you get the SQL shape and the timing without user data flowing into a log sink. This is a production-safety default, and it is deliberate: a query log is one of the easiest ways to leak email addresses, tokens, and anything else you bind into a where.
db.$on('query', (e) => console.log(e.sql, e.params));
await db.users.findUnique({ where: { email: 'alice@example.com' } });
// SELECT * FROM "users" WHERE "email" = $1 LIMIT 1 [ '[REDACTED]' ]To see the real parameter values instead of [REDACTED], opt in with logQueryParams: true on the client:
const db = turbine({
connectionString: process.env.DATABASE_URL,
logQueryParams: process.env.NODE_ENV !== 'production',
});
// SELECT * FROM "users" WHERE "email" = $1 LIMIT 1 [ 'alice@example.com' ]Only event.params changes. The SQL, timing, model, action, and row count are identical either way, so a listener written against redacted events keeps working.
Metrics, db.$observe()#
$observe turns the event stream into per-minute aggregates persisted to a _turbine_metrics table, typically in a separate database so metrics writes never touch your application data:
const handle = await db.$observe({
connectionString: process.env.TURBINE_OBSERVE_URL!, // metrics DB
flushIntervalMs: 60_000, // default: 60s
retentionDays: 30, // default: 30
});
// on shutdown
await handle.stop(); // final flush, then closes the metrics poolWhat it does, exactly:
- Buffers events in memory, keyed by
model:actionper minute bucket. - On each flush, computes count, avg, p50, p95, p99, and error count per key and upserts a row into
_turbine_metrics(ON CONFLICTmerges additively, so multiple app instances can flush into the same bucket). - Uses its own 1-connection pool against the metrics database, metric writes never contend with your application pool.
- Fire-and-forget: flush errors are swallowed, the flush timer is
unref()ed so it never keeps the process alive, and a failing metrics database cannot fail a query. - Creates the table on init (
CREATE TABLE IF NOT EXISTS _turbine_metrics) and prunes rows older thanretentionDayson each flush.
The table schema, if you want to query it yourself:
CREATE TABLE IF NOT EXISTS _turbine_metrics (
id BIGSERIAL PRIMARY KEY,
bucket TIMESTAMPTZ NOT NULL, -- minute bucket
model TEXT NOT NULL,
action TEXT NOT NULL,
count INTEGER NOT NULL DEFAULT 0,
avg_ms REAL NOT NULL DEFAULT 0,
p50_ms REAL NOT NULL DEFAULT 0,
p95_ms REAL NOT NULL DEFAULT 0,
p99_ms REAL NOT NULL DEFAULT 0,
error_count INTEGER NOT NULL DEFAULT 0,
UNIQUE(bucket, model, action)
);Calling $observe again replaces the previous engine (the old one is stopped and unsubscribed first).
Custom sinks: where the aggregates go#
By default $observe writes to _turbine_metrics in Postgres. The flush target is pluggable: pass a sink and the same per-minute aggregates flow wherever you point them. A sink is a small interface:
interface MetricsFlushBatch {
rows: Array<{
bucket: Date; // the minute bucket
model: string; // table name
action: string; // operation
count: number;
avg: number;
p50: number;
p95: number;
p99: number;
errors: number;
}>;
}
interface ObserveSink {
init?(): Promise<void>; // once at startup
flush(batch: MetricsFlushBatch): Promise<void>; // each aggregate batch
stop?(): Promise<void>; // on shutdown
}When you supply a sink, connectionString becomes optional (at least one of the two is required). The default Postgres sink is unchanged and byte-identical to the built-in writer when no sink is passed.
HttpJsonSink#
A generic sink that POSTs each batch as JSON, ready for Grafana, a self-hosted collector, or any HTTP endpoint:
import { HttpJsonSink } from 'turbine-orm';
const handle = await db.$observe({
sink: new HttpJsonSink({
url: 'https://metrics.internal/ingest',
headers: { authorization: `Bearer ${process.env.METRICS_TOKEN}` },
// fetchFunction: customFetch, // optional override
}),
});It is fire-and-forget: a failed request is swallowed and never throws, and there are no retries beyond the engine's next scheduled flush, so a down collector can never affect a query. The request body carries the aggregate batch and nothing else.
Privacy posture#
The aggregation boundary is deliberate and structural. A batch (to any sink, HTTP or Postgres) contains only:
- the model (table) name and the action (
findMany,create, …); - the count of queries and the count of errors;
- latency aggregates (avg, p50, p95, p99).
It carries no SQL text and no parameter values, those never leave the app process. Observability is opt-in and off by default: nothing is collected until you call $observe.
Dashboard, npx turbine observe#
A local, read-only dashboard over _turbine_metrics, latency over time and a per-model.action breakdown:
TURBINE_OBSERVE_URL=postgres://... npx turbine observe
npx turbine observe --port 5000 --no-openIt binds 127.0.0.1:4984 by default and shares Studio's security model: loopback by default, non-loopback refused without --allow-remote, a random per-process session token, and security headers on every response. It reads metrics only, it never connects to your application database.
See also#
- Studio, the read-only database UI with the same local security posture.
- API Reference, middleware, the active counterpart to
$on. - Typed Errors, what lands in
event.error.