Documentation
¶
Overview ¶
Package postgres
Index ¶
- Constants
- func Open(ctx context.Context, dsn string, opts ...PoolOption) (*pgxpool.Pool, error)
- type Driver
- func (d *Driver) AcquireDeriveSessionLock(ctx context.Context, orgID, harnessID, harnessSessionID string) (release func(), err error)
- func (d *Driver) AggregateSpanStats(ctx context.Context, orgID string, since, until *time.Time, authSubject string) (storage.SpanStats, error)
- func (d *Driver) ClearDeriveDirty(ctx context.Context, e storage.DeriveQueueEntry) (bool, error)
- func (d *Driver) Close() error
- func (d *Driver) CountRawTurns(ctx context.Context) (int64, error)
- func (d *Driver) DB() *pgxpool.Pool
- func (d *Driver) DeleteSession(ctx context.Context, orgID, id string) (bool, error)
- func (d *Driver) DeriveQueueStats(ctx context.Context) (storage.DeriveQueueStats, error)
- func (d *Driver) GetDeriveDirty(ctx context.Context, orgID, harnessID, harnessSessionID string) (*storage.DeriveQueueEntry, error)
- func (d *Driver) GetSessionRecord(ctx context.Context, orgID, id string) (*storage.SessionRecord, error)
- func (d *Driver) GetSessionRecordByHarness(ctx context.Context, orgID string, harnessID string, harnessSessionID string) (*storage.SessionRecord, error)
- func (d *Driver) GetSpanRecord(ctx context.Context, orgID, traceID, spanID string) (*storage.SpanRecord, error)
- func (d *Driver) GetTraceSummary(ctx context.Context, orgID, traceID string) (*storage.TraceSummaryRecord, error)
- func (d *Driver) IngestTurn(ctx context.Context, req storage.IngestTurnRequest) (storage.IngestTurnResult, error)
- func (d *Driver) IterateSessionSpans(ctx context.Context, sessionID string, after storage.SpanCursor, ...) iter.Seq2[storage.SpanRecord, error]
- func (d *Driver) IterateTraceSpans(ctx context.Context, orgID, traceID string, after storage.SpanCursor, ...) iter.Seq2[storage.SpanRecord, error]
- func (d *Driver) ListDeriveDirty(ctx context.Context, dirtiedBefore, firstDirtiedBefore time.Time, limit int32) ([]storage.DeriveQueueEntry, error)
- func (d *Driver) ListRawTurnHeaders(ctx context.Context, orgID, harnessID, harnessSessionID string, afterID int64, ...) ([]storage.RawTurnHeader, error)
- func (d *Driver) ListRawTurns(ctx context.Context, afterID int64, pageSize int32) ([]storage.RawTurnRecord, error)
- func (d *Driver) ListSessionLinks(ctx context.Context, sessionID string) ([]storage.SpanLinkRecord, error)
- func (d *Driver) ListSessionRecords(ctx context.Context, orgID string, opts storage.SessionListOpts) ([]storage.SessionRecord, error)
- func (d *Driver) ListSessionRecordsByHarnessSessionID(ctx context.Context, orgID string, harnessSessionID string) ([]storage.SessionRecord, error)
- func (d *Driver) ListSessionSpanModel(ctx context.Context, sessionID string) ([]storage.SpanTurnRecord, []storage.SpanRecord, []storage.SpanLinkRecord, ...)
- func (d *Driver) ListSpansMissingPreviews(ctx context.Context, after storage.SpanBackfillCursor, sessionID string, ...) ([]storage.SpanBackfillRow, error)
- func (d *Driver) ListTraceLinks(ctx context.Context, orgID, traceID string) ([]storage.SpanLinkRecord, error)
- func (d *Driver) ListTraceSpans(ctx context.Context, orgID, traceID string) ([]storage.SpanRecord, error)
- func (d *Driver) ListTraceSummaries(ctx context.Context, sessionID string) ([]storage.TraceSummaryRecord, error)
- func (d *Driver) MarkDeriveDirty(ctx context.Context, orgID, harnessID, harnessSessionID string) error
- func (d *Driver) MatchesPublishedFilter(ctx context.Context, filter *storage.PublishedFilter, primitiveID string) (bool, error)
- func (d *Driver) Open(ctx context.Context) error
- func (d *Driver) ProbePublishedView(ctx context.Context, view storage.PublishedViewName, ...) error
- func (d *Driver) PutRawTurn(ctx context.Context, rec storage.RawTurnRecord) (bool, error)
- func (d *Driver) RederiveFromRaw(ctx context.Context, project string) (map[string]*derive.RederiveReport, error)
- func (d *Driver) RederiveSession(ctx context.Context, project, orgID, harnessID, harnessSessionID string) (*derive.RederiveReport, error)
- func (d *Driver) RederiveSessionLocked(ctx context.Context, project, orgID, harnessID, harnessSessionID string) (*derive.RederiveReport, error)
- func (d *Driver) RepairRawTurnAttribution(ctx context.Context, project string, ...) (storage.RawTurnAttributionRepairResult, error)
- func (d *Driver) SetSpanPreviews(ctx context.Context, updates []storage.SpanPreviewUpdate) error
- func (d *Driver) SweepDeriveDirty(ctx context.Context, activeSince time.Time) (int64, error)
- func (d *Driver) TryDeriveSessionLock(ctx context.Context, orgID, harnessID, harnessSessionID string) (release func(), acquired bool, err error)
- func (d *Driver) UpdateSessionDisplayName(ctx context.Context, orgID, id string, name *string) (int64, error)
- type PoolOption
Constants ¶
const ( // FidelityUnbacked marks a row with no raw turn behind it — synthetic or // transcript-derived. Not a gap: there were never wire bytes to keep. FidelityUnbacked = "" // FidelityRaw marks a row whose raw turn holds verbatim response bytes, // so it can be re-derived from what the upstream actually sent. FidelityRaw = "raw" // FidelityReduced marks a row whose raw turn holds only an adapter's // reduction. Re-derivation is bounded by what that adapter chose to keep. FidelityReduced = "reduced" // FidelityDegraded marks a row whose verbatim bytes existed and are gone: // either they arrived and exceeded the ingest cap, or the producer // captured them and withheld them to keep the envelope under the // transport limit. Distinct from reduced: this is a capture path that // HAD the bytes and a limit that took them, which is a tuning signal // rather than a deployment fact. // // Both causes share this tier on purpose. Fidelity answers what can be // re-derived from a row, and neither can — the reason a limit bit does // not change the answer. Which limit bit is an operations question, // answered by the ingest logs that record each cause separately; giving // it a tier would put a distinction nobody re-derives differently into // every rollup, and rollups take the worst tier, so the split would be // lost at the trace level anyway. FidelityDegraded = "degraded" )
Provenance tiers stamped on the projection. See the 1781470000 migration for what each one means and why 'reduced' and 'degraded' stay distinct.
Variables ¶
This section is empty.
Functions ¶
Types ¶
type Driver ¶
type Driver struct {
// contains filtered or unexported fields
}
func (*Driver) AcquireDeriveSessionLock ¶ added in v0.25.0
func (d *Driver) AcquireDeriveSessionLock(ctx context.Context, orgID, harnessID, harnessSessionID string) (release func(), err error)
AcquireDeriveSessionLock is the BLOCKING sibling of TryDeriveSessionLock: it waits until the per-session advisory lock is held rather than returning acquired=false. This is the right semantics for a manual/one-off re-derive that must serialize BEHIND the derive worker (finish deriving after the worker releases) instead of skipping. The worker itself keeps using the non-blocking probe, so it never waits on a manual re-derive — only the reverse — which keeps the two free of a lock cycle.
Like TryDeriveSessionLock the lock is connection-scoped: the pooled connection is pinned until release unlocks it (on a fresh background context, since the caller's ctx may already be canceled).
func (*Driver) AggregateSpanStats ¶ added in v0.16.0
func (d *Driver) AggregateSpanStats(ctx context.Context, orgID string, since, until *time.Time, authSubject string) (storage.SpanStats, error)
AggregateSpanStats sums the trace-grain rollups over a time window — the span-layer aggregate behind /v1/stats. A non-empty authSubject narrows every total to that subject's sessions. Implements storage.SpanStatsReader.
func (*Driver) ClearDeriveDirty ¶ added in v0.16.0
ClearDeriveDirty implements storage.DeriveQueue. The DELETE is guarded on dirtied_at equality so a raw turn landing mid-derive (which bumps dirtied_at) keeps the session queued.
func (*Driver) CountRawTurns ¶ added in v0.16.0
CountRawTurns implements storage.RawTurnStore.
func (*Driver) DeleteSession ¶ added in v0.22.0
DeleteSession removes a session, every subagent session beneath it, and the captured data behind all of them, and reports whether the session existed. A malformed id is a no-op delete, matching DeleteSkill.
The derived rows go through the session_id ON DELETE CASCADE foreign keys. raw_turns has no foreign key to sessions — it is keyed by the harness session — so the delete removes each subtree session's raw turns by their effective attribution, with the attribution corrections recorded against them and the sessions' derive_queue marks. Everything runs in one transaction: a failure leaves the session and its capture whole.
Before touching anything the transaction takes, for every session in the subtree:
- the per-session derive lock, as a transaction-scoped advisory lock in the same order attribution repair takes it. That serializes the delete behind an in-flight derive or repair, and keeps a derive that read the raw turns before the delete from writing them into a session row a later capture recreates.
- the per-session capture lock, exclusively. PutRawTurn holds it shared for its own transaction, so a capture either commits before the delete reads the raw turns (and is deleted with them) or waits for the delete to commit (and is new capture).
- the session row, FOR UPDATE. A child session's foreign-key check needs FOR KEY SHARE on its parent row, so no subagent session can attach to the subtree until the delete commits, and the cascade removes exactly the sessions whose capture the delete removed.
The session cannot come back from a re-derive: the deriver never creates a session row, and once the raw turns are gone there is nothing left to project. A turn captured for the same harness session after the delete commits is new capture: ingest creates a new session row with a new id, and it holds only what was captured after the delete.
func (*Driver) DeriveQueueStats ¶ added in v0.16.0
DeriveQueueStats implements storage.DeriveQueue.
func (*Driver) GetDeriveDirty ¶ added in v0.16.0
func (d *Driver) GetDeriveDirty(ctx context.Context, orgID, harnessID, harnessSessionID string) (*storage.DeriveQueueEntry, error)
GetDeriveDirty implements storage.DeriveQueue.
func (*Driver) GetSessionRecord ¶ added in v0.12.0
func (d *Driver) GetSessionRecord(ctx context.Context, orgID, id string) (*storage.SessionRecord, error)
GetSessionRecord returns a single session by its UUID, or nil if not found.
func (*Driver) GetSessionRecordByHarness ¶ added in v0.15.0
func (d *Driver) GetSessionRecordByHarness( ctx context.Context, orgID string, harnessID string, harnessSessionID string, ) (*storage.SessionRecord, error)
GetSessionRecordByHarness returns the single session matching the org-scoped natural key (org_id, harness_id, harness_session_id), or nil if no row matches. The lookup is an exact-match point read on the sessions_harness_uq unique index, mirroring the GetSessionRecord nil-on-no-rows contract.
func (*Driver) GetSpanRecord ¶ added in v0.16.0
func (d *Driver) GetSpanRecord(ctx context.Context, orgID, traceID, spanID string) (*storage.SpanRecord, error)
GetSpanRecord returns one span with full payloads. Implements storage.SpanModelReader.
func (*Driver) GetTraceSummary ¶ added in v0.49.0
func (d *Driver) GetTraceSummary(ctx context.Context, orgID, traceID string) (*storage.TraceSummaryRecord, error)
GetTraceSummary returns one turn header with its span count and no spans — what the streaming trace page writes before its first span is read. nil when the trace does not exist. Implements storage.SpanModelReader.
func (*Driver) IngestTurn ¶ added in v0.10.0
func (d *Driver) IngestTurn(ctx context.Context, req storage.IngestTurnRequest) (storage.IngestTurnResult, error)
IngestTurn implements storage.SessionIngester for the Postgres driver. The session-tracking flow runs in a single transaction: resolve / UPSERT a sessions row (keyed by the envelope's natural key or a synthetic harness_session_id from the turn's Merkle root), resolve the optional fork-parent FK (placeholder-inserting the parent when its own first turn hasn't landed yet), and fold the derived title. It writes session IDENTITY only: nodes are not persisted (the merkle layer is in-memory), and every derived rollup — counters, model_usage, tasks, kind_counts, and the chain-aware status/git/tool signals — is owned by the derive-time span fold.
func (*Driver) IterateSessionSpans ¶ added in v0.49.0
func (d *Driver) IterateSessionSpans(ctx context.Context, sessionID string, after storage.SpanCursor, mode storage.PayloadMode) iter.Seq2[storage.SpanRecord, error]
IterateSessionSpans streams one session's spans in composite order, starting strictly after the cursor (zero value: from the first span). Implements storage.SpanModelReader.
func (*Driver) IterateTraceSpans ¶ added in v0.49.0
func (d *Driver) IterateTraceSpans(ctx context.Context, orgID, traceID string, after storage.SpanCursor, mode storage.PayloadMode) iter.Seq2[storage.SpanRecord, error]
IterateTraceSpans streams one trace's spans in presentation order, starting strictly after the cursor (zero value: from the first span). Implements storage.SpanModelReader.
func (*Driver) ListDeriveDirty ¶ added in v0.16.0
func (d *Driver) ListDeriveDirty(ctx context.Context, dirtiedBefore, firstDirtiedBefore time.Time, limit int32) ([]storage.DeriveQueueEntry, error)
ListDeriveDirty implements storage.DeriveQueue.
func (*Driver) ListRawTurnHeaders ¶ added in v0.16.0
func (d *Driver) ListRawTurnHeaders(ctx context.Context, orgID, harnessID, harnessSessionID string, afterID int64, limit int) ([]storage.RawTurnHeader, error)
ListRawTurnHeaders returns one page of the wire log for one session: capture identity and payload sizes, no blobs, for up to limit rows in id order strictly after afterID. Implements storage.SpanModelReader.
func (*Driver) ListRawTurns ¶ added in v0.16.0
func (d *Driver) ListRawTurns(ctx context.Context, afterID int64, pageSize int32) ([]storage.RawTurnRecord, error)
ListRawTurns implements storage.RawTurnStore.
func (*Driver) ListSessionLinks ¶ added in v0.25.0
func (d *Driver) ListSessionLinks(ctx context.Context, sessionID string) ([]storage.SpanLinkRecord, error)
ListSessionLinks returns a session's dataflow links alone, in the same deterministic key order ListSessionSpanModel serves them. It is the payload-free half of that read, for the per-trace streaming export. Implements storage.SpanModelReader.
func (*Driver) ListSessionRecords ¶ added in v0.12.0
func (d *Driver) ListSessionRecords( ctx context.Context, orgID string, opts storage.SessionListOpts, ) ([]storage.SessionRecord, error)
ListSessionRecords returns a page of sessions for an org ordered by the requested sort column (default last_seen_at DESC), optionally windowed by activity (a turn started in the window, matching /v1/stats) and narrowed to one gateway-stamped JWT subject (exact match on the indexed column; empty lists every user's sessions). Pass zero-value opts to start from the beginning, unwindowed and unfiltered.
func (*Driver) ListSessionRecordsByHarnessSessionID ¶ added in v0.36.0
func (d *Driver) ListSessionRecordsByHarnessSessionID( ctx context.Context, orgID string, harnessSessionID string, ) ([]storage.SessionRecord, error)
ListSessionRecordsByHarnessSessionID returns every session in the org whose harness_session_id exactly matches, across all harnesses. The id is unique within a harness (sessions_harness_uq), so this returns at most one row per harness that has seen it — in practice zero or one. It serves the lone-harness_session_id form of the /v1/sessions filter; the paired form stays on GetSessionRecordByHarness's point lookup. No match is an empty slice, not an error.
func (*Driver) ListSessionSpanModel ¶ added in v0.16.0
func (d *Driver) ListSessionSpanModel(ctx context.Context, sessionID string) ([]storage.SpanTurnRecord, []storage.SpanRecord, []storage.SpanLinkRecord, error)
ListSessionSpanModel returns the stored span projection for one session: turns, spans, and links, each in stable presentation order. Implements storage.SpanModelReader.
func (*Driver) ListSpansMissingPreviews ¶ added in v0.49.0
func (d *Driver) ListSpansMissingPreviews(ctx context.Context, after storage.SpanBackfillCursor, sessionID string, limit int) ([]storage.SpanBackfillRow, error)
ListSpansMissingPreviews returns one keyset page of spans that carry a payload but no stored preview, in (session_id, trace_id, span_id) order starting strictly after `after`. Implements storage.PreviewBackfiller.
func (*Driver) ListTraceLinks ¶ added in v0.49.0
func (d *Driver) ListTraceLinks(ctx context.Context, orgID, traceID string) ([]storage.SpanLinkRecord, error)
ListTraceLinks returns the dataflow links touching one trace on either end. Implements storage.SpanModelReader.
func (*Driver) ListTraceSpans ¶ added in v0.25.0
func (d *Driver) ListTraceSpans(ctx context.Context, orgID, traceID string) ([]storage.SpanRecord, error)
ListTraceSpans returns one trace's spans in presentation order (seq ASC, matching ListSessionSpanModel restricted to the trace), with full payloads. Implements storage.SpanModelReader.
func (*Driver) ListTraceSummaries ¶ added in v0.16.0
func (d *Driver) ListTraceSummaries(ctx context.Context, sessionID string) ([]storage.TraceSummaryRecord, error)
ListTraceSummaries returns a session's turn headers with span counts — the lazy session-detail rows. Implements storage.SpanModelReader.
func (*Driver) MarkDeriveDirty ¶ added in v0.16.0
func (d *Driver) MarkDeriveDirty(ctx context.Context, orgID, harnessID, harnessSessionID string) error
MarkDeriveDirty implements storage.DeriveQueue.
func (*Driver) MatchesPublishedFilter ¶ added in v0.39.0
func (d *Driver) MatchesPublishedFilter( ctx context.Context, filter *storage.PublishedFilter, primitiveID string, ) (bool, error)
MatchesPublishedFilter reports whether one primitive id carries every value of the filter in the published view. It serves the point-lookup paths (the harness natural-key filter) where there is no list query to compose the predicate into: the evaluation still happens in SQL — one indexed EXISTS probe per value — never by fetching and filtering rows in Go. The same injection doctrine applies: the identifier comes from the opaque PublishedViewName's one quoting helper, every value binds.
func (*Driver) ProbePublishedView ¶ added in v0.40.0
func (d *Driver) ProbePublishedView( ctx context.Context, view storage.PublishedViewName, column storage.PublishedColumnName, ) error
ProbePublishedView verifies that this driver's own role can read the published view a filter claim names: one WHERE FALSE round trip that fails when the view is missing, when SELECT was never granted, or when the contract columns or the claim-declared value column are absent. WHERE FALSE keeps the probe planning-only — no rows are read, so its cost does not scale with the view's contents.
func (*Driver) PutRawTurn ¶ added in v0.16.0
PutRawTurn implements storage.RawTurnStore. The row is appended verbatim; a retried POST with the same (org_id, request_id) is a no-op per the partial unique index.
Session-keyed rows also mark the session dirty in derive_queue, in the same transaction, so the derive worker picks the session up, and hold the session's capture lock shared so a concurrent DeleteSession either removes the turn or runs entirely before it. Marking happens even when the row deduped: a re-POST of an existing turn is the natural "re-derive this session" signal, and a redundant mark only costs one idempotent derive.
func (*Driver) RederiveFromRaw ¶ added in v0.16.0
func (d *Driver) RederiveFromRaw(ctx context.Context, project string) (map[string]*derive.RederiveReport, error)
RederiveFromRaw rebuilds every persisted session from its effective raw turns. Sessions are enumerated from the read model rather than from current raw attribution so a repair source that has become empty is still covered and pruned.
Each session's full read-derive-write pass runs under the same advisory lock used by the derive worker and attribution repair. Holding one lock at a time avoids both stale-read races and lock cycles with repair's ordered two-lock acquisition. The pass is atomic per session, not per org; reports retain the existing per-org shape by aggregating session reports.
func (*Driver) RederiveSession ¶ added in v0.16.0
func (d *Driver) RederiveSession(ctx context.Context, project, orgID, harnessID, harnessSessionID string) (*derive.RederiveReport, error)
RederiveSession is the session-scoped sibling of RederiveFromRaw: re-derive ONE harness session from its raw turns and apply the result transactionally (upsert + prune scoped to that session). This is the derive worker's unit of work — memory stays bounded by one session's unique content, and the full rows stream through the deriver one at a time exactly like the full-org pass.
Same idempotence contract: re-running against an unchanged raw layer upserts the same set and prunes nothing.
func (*Driver) RederiveSessionLocked ¶ added in v0.25.0
func (d *Driver) RederiveSessionLocked(ctx context.Context, project, orgID, harnessID, harnessSessionID string) (*derive.RederiveReport, error)
RederiveSessionLocked is the externally-safe entry point for a session-scoped re-derive that may run WHILE the derive worker is live: it holds the per-session advisory lock across the whole read-derive-write pass, so it cannot interleave with the worker's derive of the same session and prune a turn the worker just wrote. The worker's own path (RederiveSession) is already called under this lock — it takes it in processEntry — so RederiveSession stays lock-free and this wrapper is the one non-worker callers use. Blocking: it waits out a concurrent worker derive rather than skipping, since a manual re-derive must actually run.
func (*Driver) RepairRawTurnAttribution ¶ added in v0.30.0
func (d *Driver) RepairRawTurnAttribution( ctx context.Context, project string, req storage.RawTurnAttributionRepairRequest, ) (storage.RawTurnAttributionRepairResult, error)
RepairRawTurnAttribution appends an attribution correction and rebuilds both affected projections. raw_turns remains byte-for-byte untouched.
func (*Driver) SetSpanPreviews ¶ added in v0.49.0
SetSpanPreviews fills input_preview / output_preview for every update in one transaction. The UPDATE names only those two columns, so the payload, content_hash and derive_seq are untouched. Implements storage.PreviewBackfiller.
func (*Driver) SweepDeriveDirty ¶ added in v0.16.0
SweepDeriveDirty implements storage.DeriveQueue.
func (*Driver) TryDeriveSessionLock ¶ added in v0.16.0
func (d *Driver) TryDeriveSessionLock(ctx context.Context, orgID, harnessID, harnessSessionID string) (release func(), acquired bool, err error)
TryDeriveSessionLock takes a session-scoped Postgres advisory lock so concurrent workers never double-derive one session. The lock is connection-scoped: the pooled connection is pinned for the lock's lifetime and returned by release. A false return with nil error means another holder has the session — skip it, this is not a failure.
The key hashes the harness triple to 64 bits (FNV-1a). A cross-triple collision only serializes two unrelated sessions' derives — safe, merely slower — so probabilistic keying is fine here.
func (*Driver) UpdateSessionDisplayName ¶ added in v0.27.0
func (d *Driver) UpdateSessionDisplayName(ctx context.Context, orgID, id string, name *string) (int64, error)
UpdateSessionDisplayName sets (or clears, when name is nil) the user-facing `display_name` column for a single session, scoped to the caller's org. It targets display_name — not name — so a Console rename survives ingest re-sending the harness slug every turn (PCC-970). The org_id predicate lives in the SQL (UpdateSessionDisplayName query), never just checked beforehand, so a cross-org id can never be updated. Returns the number of rows affected: 0 means no row matched (unknown id, or an id that exists but belongs to a different org), which the caller (the sessions handler) maps to a 404 rather than treating as success.
type PoolOption ¶ added in v0.16.0
PoolOption tunes the pgx pool configuration built from the DSN. Options are applied after the DSN is parsed, so they take precedence over pool parameters embedded in the connection string.
func WithConnectTimeout ¶ added in v0.16.0
func WithConnectTimeout(d time.Duration) PoolOption
WithConnectTimeout bounds each connection attempt. pgx has no built-in connect timeout, so an unreachable-but-blackholed host can otherwise hang a startup for the OS TCP timeout (minutes).
func WithMaxConns ¶ added in v0.16.0
func WithMaxConns(n int32) PoolOption
WithMaxConns caps the pool size. Single-purpose services (e.g. the derive worker, which processes one session at a time) should set a small cap instead of inheriting pgx's NumCPU-based default.