postgres

package
v0.2.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 27, 2026 License: Apache-2.0 Imports: 110 Imported by: 0

Documentation

Overview

Package postgres provides PostgreSQL-backed repositories (pgx/v5) and applies migrations via goose. Used when SYNAPSE_DB_DSN is set; otherwise the server falls back to in-memory persistence for dev.

Index

Constants

This section is empty.

Variables

View Source
var ErrTenantCommit = errors.New("tenant transaction commit failed")

ErrTenantCommit marks a failure of the COMMIT that ends a WithTenant transaction, as opposed to a failure inside fn. Once COMMIT has been sent its durable outcome can be unknowable, so a caller that compensates on failure (by deleting an artifact it wrote outside the database, say) must not treat this as "the transaction did not happen".

Functions

func AcquireSingletonLock

func AcquireSingletonLock(ctx context.Context, pool *pgxpool.Pool, role string) (*pgxpool.Conn, bool, error)

AcquireSingletonLock takes a session-level advisory lock (keyed by role) on a DEDICATED connection the caller holds for the whole process lifetime – releasing it drops the lock. A second instance OF THE SAME ROLE gets ok=false so it can fail fast (the repos still ignore tenant_id, so two same-role writers would race). Returns the held connection (retain it; Release at shutdown), whether the lock was obtained, and any error.

func CheckDatabaseReady added in v0.2.0

func CheckDatabaseReady(ctx context.Context, pool *pgxpool.Pool) error

CheckDatabaseReady verifies the runtime pool can execute a trivial query.

func CheckMigrationsReady added in v0.2.0

func CheckMigrationsReady(ctx context.Context, pool *pgxpool.Pool) error

CheckMigrationsReady verifies the latest state of every embedded migration is applied. Goose retains down records, so the query considers only the newest row for each migration version.

func CheckRLSRuntimeRole added in v0.1.8

func CheckRLSRuntimeRole(ctx context.Context, pool *pgxpool.Pool) error

CheckRLSRuntimeRole reports whether the role the pool connects as can actually be constrained by Row Level Security. RLS is bypassed entirely by SUPERUSER and BYPASSRLS roles regardless of FORCE ROW LEVEL SECURITY, so if the runtime role holds either attribute the whole tenant isolation guarantee is silently a no-op. It returns a non-nil error naming the offending attribute when the role would bypass RLS.

This is intended to gate multi-tenant enablement (fail-closed): the caller that turns on RLS-protected tables must refuse to serve if this returns an error. It is exported and separate so that path can enforce it at startup, while single-tenant deployments that connect as a superuser and use no RLS-protected table are not forced to change their role.

func Connect

func Connect(ctx context.Context, dsn string) (*pgxpool.Pool, error)

Connect opens a pgx pool with default sizing (back-compat wrapper).

func ConnectPool

func ConnectPool(ctx context.Context, dsn string, pc PoolConfig) (*pgxpool.Pool, error)

ConnectPool opens a sized pgx pool and verifies connectivity.

func GrantRuntimePrivileges added in v0.1.8

func GrantRuntimePrivileges(ctx context.Context, adminDSN, runtimeDSN string, haltWriterDSNs ...string) error

GrantRuntimePrivileges grants ordinary runtime DML, then removes its ability to create response halt fences or dispatches. When configured, the halt writer receives only the small table/sequence privileges necessary for its one atomic transaction.

func Migrate

func Migrate(ctx context.Context, dsn string) error

Migrate applies all pending goose migrations (idempotent; tracked in goose_db_version).

func MigrateLocked added in v0.2.0

func MigrateLocked(ctx context.Context, dsn string) error

MigrateLocked applies the embedded migration set while holding a database-wide advisory lock. The lock and goose share the pool's sole connection, so the session lock is held for the complete migration run. Acquisition blocks until the caller's context expires.

func ValidateMigrationRoleSeparation added in v0.1.8

func ValidateMigrationRoleSeparation(migrationDSN, runtimeDSN string) error

ValidateMigrationRoleSeparation ensures migrations cannot run as the runtime role.

func ValidateResponseRoleSeparation added in v0.2.0

func ValidateResponseRoleSeparation(migrationDSN, runtimeDSN, haltWriterDSN string) error

ValidateResponseRoleSeparation ensures the DDL, normal runtime, and halt-writer identities cannot be confused in a governed response deployment.

func WithContextTenant added in v0.1.8

func WithContextTenant(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error

WithContextTenant runs fn under the immutable tenant previously bound to ctx.

func WithGlobalRead added in v0.1.8

func WithGlobalRead(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error

func WithGlobalWrite added in v0.1.8

func WithGlobalWrite(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error

func WithTenant added in v0.1.8

func WithTenant(ctx context.Context, pool *pgxpool.Pool, tenantID string, fn func(pgx.Tx) error) (err error)

WithTenant runs fn inside a transaction whose `app.current_tenant` session variable is set to tenantID for the life of that transaction only. Stores of Row-Level-Security-protected tables (see migration 0057 and its synapse_enable_tenant_rls procedure) MUST route reads and writes through this helper: the policy denies every row when the tenant resolves to NULL, so a query that runs outside WithTenant sees nothing rather than leaking across tenants. The setting is applied with set_config(..., is_local => true), which is transaction-scoped.

Fail-closed semantics: an empty tenantID resolves (via synapse_current_tenant's NULLIF) to NULL and therefore matches no row. Under RLS the empty string is DENY, not the default tenant, so callers of RLS-protected tables must pass a non-empty tenant id. This closes the placeholder-GUC reset hazard: app.current_tenant reverts to the empty string (not "unset") after a transaction, and mapping the empty string to NULL means a connection reused outside WithTenant still denies rather than exposing default-tenant rows.

Types

type AITriageReviewRepository added in v0.1.8

type AITriageReviewRepository struct {
	// contains filtered or unexported fields
}

func NewAITriageReviewRepository added in v0.1.8

func NewAITriageReviewRepository(pool *pgxpool.Pool) *AITriageReviewRepository

func (*AITriageReviewRepository) Get added in v0.1.8

func (*AITriageReviewRepository) List added in v0.1.8

func (*AITriageReviewRepository) SaveDecision added in v0.1.8

func (r *AITriageReviewRepository) SaveDecision(ctx context.Context, review aitriagereview.Review, expectedVersion int) error

func (*AITriageReviewRepository) SaveOwner added in v0.1.8

func (r *AITriageReviewRepository) SaveOwner(ctx context.Context, review aitriagereview.Review, expectedVersion int) error

func (*AITriageReviewRepository) UpsertPending added in v0.1.8

func (r *AITriageReviewRepository) UpsertPending(ctx context.Context, review aitriagereview.Review) error

type AUPStore

type AUPStore struct {
	// contains filtered or unexported fields
}

AUPStore persists Acceptable-Use-Policy acceptances to PostgreSQL.

func NewAUPStore

func NewAUPStore(pool *pgxpool.Pool) *AUPStore

NewAUPStore returns an AUP store backed by the given pool.

func (*AUPStore) Accepted

func (s *AUPStore) Accepted(ctx context.Context, version string) (bool, error)

Accepted reports whether the given policy version has been accepted.

func (*AUPStore) Save

func (s *AUPStore) Save(ctx context.Context, a aup.Acceptance) error

Save records an acceptance, idempotent per (actor, version) – this keeps per-actor history (the file dev sink keeps one record per version; both gate identically via Accepted's EXISTS-by-version). RBAC is enforced at the API edge. Should actor identifiers ever become attacker- influenced (e.g. a future external OIDC subject), key idempotency on a UNIQUE(actor, policy_version) constraint instead of a concatenated id.

type AccuracyRunRepository added in v0.2.0

type AccuracyRunRepository struct {
	// contains filtered or unexported fields
}

AccuracyRunRepository persists detection-accuracy regression runs. Runs are deployment-global engine data (the golden corpus is fixed), so there is no tenant scoping or RLS, matching the advisories repository.

func NewAccuracyRunRepository added in v0.2.0

func NewAccuracyRunRepository(pool *pgxpool.Pool) *AccuracyRunRepository

NewAccuracyRunRepository returns a repository backed by the given pool.

func (*AccuracyRunRepository) Recent added in v0.2.0

func (r *AccuracyRunRepository) Recent(ctx context.Context, limit int) ([]accuracy.Run, error)

Recent returns the most recent runs, newest first, capped at limit (a non-positive limit defaults to 100).

func (*AccuracyRunRepository) Save added in v0.2.0

Save inserts a run. A duplicate id is ignored (runs are immutable snapshots).

type AdvisoryMaterializer added in v0.1.8

type AdvisoryMaterializer struct {
	// contains filtered or unexported fields
}

func NewAdvisoryMaterializer added in v0.1.8

func NewAdvisoryMaterializer(pool *pgxpool.Pool) *AdvisoryMaterializer

func (*AdvisoryMaterializer) AdvisoryRevisionAt added in v0.1.8

func (r *AdvisoryMaterializer) AdvisoryRevisionAt(ctx context.Context, advisoryID string, snapshotAt time.Time) (ports.AdvisoryRevisionRef, error)

func (*AdvisoryMaterializer) ByCPE added in v0.1.8

func (r *AdvisoryMaterializer) ByCPE(ctx context.Context, part, vendor, product string) ([]advisory.Advisory, error)

func (*AdvisoryMaterializer) ByPackage added in v0.1.8

func (r *AdvisoryMaterializer) ByPackage(ctx context.Context, ecosystem, packageName string) ([]advisory.Advisory, error)

func (*AdvisoryMaterializer) CountVulnerabilityAdvisoriesChangedSince added in v0.1.8

func (r *AdvisoryMaterializer) CountVulnerabilityAdvisoriesChangedSince(ctx context.Context, since time.Time) (int64, error)

func (*AdvisoryMaterializer) CountVulnerabilityAdvisoryDailyImpact added in v0.2.0

func (r *AdvisoryMaterializer) CountVulnerabilityAdvisoryDailyImpact(ctx context.Context, since time.Time) (vulnerabilityintel.AdvisoryDailyImpact, error)

func (*AdvisoryMaterializer) CurrentRevision added in v0.1.8

func (r *AdvisoryMaterializer) CurrentRevision(ctx context.Context, id string) (int64, error)

func (*AdvisoryMaterializer) CurrentSourceRecordIDs added in v0.1.8

func (r *AdvisoryMaterializer) CurrentSourceRecordIDs(ctx context.Context, sourceID string, yield func(string) error) error

func (*AdvisoryMaterializer) CurrentSourceRecordIDsBounded added in v0.2.0

func (r *AdvisoryMaterializer) CurrentSourceRecordIDsBounded(ctx context.Context, sourceID string, limit int, yield func(string) error) error

func (*AdvisoryMaterializer) GetCanonical added in v0.1.8

func (r *AdvisoryMaterializer) GetCanonical(ctx context.Context, id string) (advisory.Canonical, error)

func (*AdvisoryMaterializer) GetCanonicalAtRevision added in v0.1.8

func (r *AdvisoryMaterializer) GetCanonicalAtRevision(ctx context.Context, id string, revision int64) (advisory.Canonical, error)

func (*AdvisoryMaterializer) ListAdvisoryRevisions added in v0.1.8

func (r *AdvisoryMaterializer) ListAdvisoryRevisions(ctx context.Context, after string, snapshotAt time.Time, limit int) (ports.AdvisoryRevisionPage, error)

func (*AdvisoryMaterializer) ListVulnerabilityAdvisories added in v0.1.8

func (r *AdvisoryMaterializer) ListVulnerabilityAdvisories(ctx context.Context, tenantID shared.ID, query vulnerabilityintel.AdvisoryQuery) (vulnerabilityintel.AdvisoryPage, error)

func (*AdvisoryMaterializer) ListVulnerabilityAdvisoryRevisions added in v0.1.8

func (*AdvisoryMaterializer) ListVulnerabilitySyncRunRevisions added in v0.1.8

func (r *AdvisoryMaterializer) ListVulnerabilitySyncRunRevisions(ctx context.Context, runIDs []shared.ID, limitPerRun int) (map[shared.ID]vulnerabilityintel.AdvisoryRevisionLinkPage, error)

func (*AdvisoryMaterializer) MarkAdvisoryEvaluated added in v0.1.8

func (r *AdvisoryMaterializer) MarkAdvisoryEvaluated(ctx context.Context, tenantID shared.ID, advisoryID string, revision int64, evaluatedAt time.Time) error

func (*AdvisoryMaterializer) Materialize added in v0.1.8

func (*AdvisoryMaterializer) MaterializeSourceSnapshot added in v0.2.0

func (r *AdvisoryMaterializer) MaterializeSourceSnapshot(ctx context.Context, records []advisory.ObservationRecord) ([]advisory.MaterializationResult, error)

MaterializeSourceSnapshot atomically publishes one complete source view. Independent records are materialized separately inside one transaction because they need not share identities, while any error rolls the full source change back.

func (*AdvisoryMaterializer) OldestUnevaluatedAdvisory added in v0.1.8

func (r *AdvisoryMaterializer) OldestUnevaluatedAdvisory(ctx context.Context, tenantID shared.ID) (*vulnerabilityintel.EvaluationLag, error)

func (*AdvisoryMaterializer) PublishSourceSnapshot added in v0.2.0

PublishSourceSnapshot atomically commits a complete source view with a receipt that preserves its provider checkpoint and exact revision results for durable recovery.

func (*AdvisoryMaterializer) PublishedSourceSnapshot added in v0.2.0

func (r *AdvisoryMaterializer) PublishedSourceSnapshot(ctx context.Context, syncRunID shared.ID) (ports.PublishedSourceSnapshot, bool, error)

func (*AdvisoryMaterializer) SummarizeVulnerabilityCoverage added in v0.2.0

func (*AdvisoryMaterializer) SupportsServerAdvisoryFilters added in v0.2.0

func (*AdvisoryMaterializer) SupportsServerAdvisoryFilters() bool

type AdvisoryRepository

type AdvisoryRepository struct {
	// contains filtered or unexported fields
}

AdvisoryRepository persists the OWNED normalized-advisory store to PostgreSQL. It is GLOBAL reference data (NOT tenant-scoped): the full advisory is a JSONB blob in `advisories`, with one `advisory_affects` row per affected (ecosystem, package) for the indexed ByPackage lookup.

func NewAdvisoryRepository

func NewAdvisoryRepository(pool *pgxpool.Pool) *AdvisoryRepository

NewAdvisoryRepository returns a repository backed by the given pool.

func (*AdvisoryRepository) AdvisoryAliasEdges added in v0.2.0

func (r *AdvisoryRepository) AdvisoryAliasEdges(ctx context.Context, ids []string) ([]advisory.AliasEdge, error)

AdvisoryAliasEdges returns the alias edges (alias id -> canonical id, the row id) for every advisory whose id is in ids OR whose stored Aliases array intersects ids. The `?|` intersection is backed by the advisories alias GIN index (migration 0145), so the query is bounded to the finding ids rather than a full-corpus scan. Ids are matched as stored; the caller normalizes for the alias graph. An empty ids slice returns no edges (no findings to expand).

func (*AdvisoryRepository) AdvisoryFreshness added in v0.2.0

func (r *AdvisoryRepository) AdvisoryFreshness(ctx context.Context) (time.Time, int, error)

AdvisoryFreshness reports the newest advisory timestamp and the corpus row count so a scan can warn when the owned advisory store is stale. Advisories are global reference data (not tenant-scoped), so the query is unfiltered. An empty corpus yields the zero time and count 0.

func (*AdvisoryRepository) ByCPE added in v0.1.8

func (r *AdvisoryRepository) ByCPE(ctx context.Context, part, vendor, product string) ([]advisory.Advisory, error)

func (*AdvisoryRepository) ByPackage

func (r *AdvisoryRepository) ByPackage(ctx context.Context, ecosystem, name string) ([]advisory.Advisory, error)

ByPackage returns the advisories that affect (ecosystem, name), decoded from their JSONB blobs. The caller runs advisory.Match to decide which actually hit the component's version. Deterministic id order.

func (*AdvisoryRepository) CoveredEcosystems added in v0.2.0

func (r *AdvisoryRepository) CoveredEcosystems(ctx context.Context) (map[string]bool, error)

CoveredEcosystems returns the distinct ecosystems the corpus has any affected-package row for, so the readiness guard can tell a genuinely-covered distro from a silent gap. One DISTINCT scan of the advisory_affects index, run once per scan.

func (*AdvisoryRepository) Upsert

Upsert inserts or replaces an advisory by id and rebuilds its (ecosystem, package) index rows, in one transaction. Idempotent – advisories are re-syncable reference data (a re-ingest REPLACES in place), not an append-only ledger. The affected (ecosystem, package) keys must be ingester-normalized per the ports.AdvisoryStore KEY CONTRACT. The full domain advisory round-trips through the JSONB `data` blob.

type AgentDecisionStore

type AgentDecisionStore struct {
	// contains filtered or unexported fields
}

AgentDecisionStore is the durable ports.DecisionStore on PostgreSQL (migration 0032). seq is a monotonic per-session counter (MAX+1). Idempotency is enforced by partial unique indexes – one row per (session_id, action_id) for steps and one stop per session – so a redelivered drive that re-records a decision is a no-op (ON CONFLICT DO NOTHING), never a forked log. Decisions are written by a single driver under the session run lock, so the MAX+1 read/insert pair is race-free.

func NewAgentDecisionStore

func NewAgentDecisionStore(pool *pgxpool.Pool) *AgentDecisionStore

NewAgentDecisionStore returns a Postgres-backed decision store.

func (*AgentDecisionStore) AppendDecision

func (s *AgentDecisionStore) AppendDecision(ctx context.Context, d agent.AgentDecision) error

func (*AgentDecisionStore) ListBySession

func (s *AgentDecisionStore) ListBySession(ctx context.Context, sessionID shared.ID) ([]agent.AgentDecision, error)

type AgentPlanStore

type AgentPlanStore struct {
	// contains filtered or unexported fields
}

AgentPlanStore is the durable ports.PlanStore on PostgreSQL (migration 0031). One plan per session (session_id UNIQUE → CreatePlan on a redelivery hits a unique violation → ErrConflict, preventing a forked second plan). SavePlan is a guarded UPDATE (… WHERE revision=$expected) that bumps the revision, so a node claim is an atomic compare-and-swap; a lost CAS (0 rows) returns ErrConflict.

agent_plans is RLS-protected (migration 0129). agent.Plan carries no tenant of its own, so every statement runs under the ambient tenant bound to ctx and keys on (tenant_id, session_id): a context without a tenant is a validation error, and a plan cannot be read or CAS-advanced from another tenant.

func NewAgentPlanStore

func NewAgentPlanStore(pool *pgxpool.Pool) *AgentPlanStore

NewAgentPlanStore returns a Postgres-backed plan store.

func (*AgentPlanStore) CreatePlan

func (s *AgentPlanStore) CreatePlan(ctx context.Context, p agent.Plan) error

func (*AgentPlanStore) GetBySession

func (s *AgentPlanStore) GetBySession(ctx context.Context, sessionID shared.ID) (agent.Plan, bool, error)

func (*AgentPlanStore) SavePlan

func (s *AgentPlanStore) SavePlan(ctx context.Context, p agent.Plan) error

type AgentSessionStore

type AgentSessionStore struct {
	// contains filtered or unexported fields
}

AgentSessionStore is the durable ports.AgentSessionStore on PostgreSQL: agent_sessions + agent_messages (migration 0027). The (session_id, seq) primary key is the transcript fork-guard - a duplicate seq is a unique violation, reported as ErrConflict.

agent_sessions is RLS-protected (migration 0129), so every statement runs inside a WithTenant transaction and keeps its explicit tenant_id predicate. agent_messages has no tenant column of its own; it is reached only through a join or subselect on agent_sessions, which the policy scopes. ListResumable is the one cross-tenant read: it fans out over the tenants table (which is not RLS-protected) and runs one tenant-bound query per tenant.

func NewAgentSessionStore

func NewAgentSessionStore(pool *pgxpool.Pool) *AgentSessionStore

NewAgentSessionStore returns a Postgres-backed agent session store.

func (*AgentSessionStore) AppendMessage

func (s *AgentSessionStore) AppendMessage(ctx context.Context, sessionID shared.ID, seq int, m agent.Message) error

func (*AgentSessionStore) GetSession

func (s *AgentSessionStore) GetSession(ctx context.Context, id shared.ID) (agent.Session, error)

func (*AgentSessionStore) ListByEngagement

func (s *AgentSessionStore) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]agent.Session, error)

func (*AgentSessionStore) ListResumable

func (s *AgentSessionStore) ListResumable(ctx context.Context, staleFor time.Duration, now time.Time, limit int) ([]agent.Session, error)

ListResumable is the startup reconciler's cross-tenant sweep, so it has no ambient tenant to run under. Rather than reading agent_sessions unscoped, it enumerates the tenants table (global reference data, not RLS-protected) and runs one tenant-bound query per tenant, exactly as the job queue's Claim and the engagement repository's reconciliation scan do. The caller binds each returned session's own tenant before acting on it.

func (*AgentSessionStore) Messages

func (s *AgentSessionStore) Messages(ctx context.Context, sessionID shared.ID) ([]agent.Message, error)

func (*AgentSessionStore) SaveSession

func (s *AgentSessionStore) SaveSession(ctx context.Context, e agent.Session) error

type AgentSigningKeyRepository added in v0.2.0

type AgentSigningKeyRepository struct {
	// contains filtered or unexported fields
}

AgentSigningKeyRepository is the Postgres-backed agent content-signing key registry (#607, A0.2). Every method runs through WithContextTenant so Row Level Security (migration 0105 via the 0057 procedure) isolates one tenant's keys from another's. A key's public half and window are immutable once written; revocation is monotonic. The primary key (tenant_id, agent_id, key_id) makes registration idempotent on identity and blocks re-pointing a KeyID at another key (anti-rollback).

func NewAgentSigningKeyRepository added in v0.2.0

func NewAgentSigningKeyRepository(pool *pgxpool.Pool) *AgentSigningKeyRepository

NewAgentSigningKeyRepository constructs the repository.

func (*AgentSigningKeyRepository) ListByAgent added in v0.2.0

ListByAgent returns every key for an agent under the ctx tenant, newest NotBefore first.

func (*AgentSigningKeyRepository) Register added in v0.2.0

Register stores a signing key, idempotent on identity and anti-rollback on KeyID.

func (*AgentSigningKeyRepository) ResolveSigningKey added in v0.2.0

func (r *AgentSigningKeyRepository) ResolveSigningKey(ctx context.Context, agentID shared.ID, keyID string) (fleetagent.AgentSigningKey, error)

ResolveSigningKey returns the (agent, keyID) key under the ctx tenant, or shared.ErrNotFound.

func (*AgentSigningKeyRepository) Revoke added in v0.2.0

func (r *AgentSigningKeyRepository) Revoke(ctx context.Context, agentID shared.ID, keyID string, at time.Time) error

Revoke marks (agentID, keyID) revoked at `at`, monotonic (an already-revoked key keeps its first RevokedAt), tenant-scoped. shared.ErrNotFound if the key is unknown.

type ApprovalStore

type ApprovalStore struct {
	// contains filtered or unexported fields
}

ApprovalStore is the durable ports.ApprovalStore on PostgreSQL: the HITL approval queue (migration 0028). Decide is a guarded UPDATE (... WHERE decision_state= 'pending') so the first decision wins; a second hits 0 rows and returns ErrConflict.

agent_approvals is RLS-protected (migration 0129). agent.ProposedAction carries no tenant of its own, so every statement runs under the ambient tenant bound to ctx and keys on (tenant_id, action_id). The guarded UPDATEs used to key on action_id alone, which meant a decision or a consume could be applied to another tenant's action id; the tenant is now part of every guard. EngagementsWithPending is the one cross-tenant read and fans out over the tenants table instead of scanning the queue unscoped.

func NewApprovalStore

func NewApprovalStore(pool *pgxpool.Pool) *ApprovalStore

NewApprovalStore returns a Postgres-backed approval store.

func (*ApprovalStore) Consume added in v0.1.8

func (s *ApprovalStore) Consume(ctx context.Context, actionID shared.ID) error

func (*ApprovalStore) Decide

func (*ApprovalStore) EngagementsWithPending

func (s *ApprovalStore) EngagementsWithPending(ctx context.Context) ([]ports.ApprovalSweepScope, error)

EngagementsWithPending is the fail-closed timeout sweeper's fan-out input. It has no ambient tenant, so it enumerates the tenants table and asks each tenant's partition for its pending engagements under a tenant-bound transaction. Each scope carries the tenant back so the caller binds it before sweeping, which is what keeps the rest of the sweep inside one tenant.

func (*ApprovalStore) Enqueue

func (*ApprovalStore) Get

func (*ApprovalStore) Pending

func (s *ApprovalStore) Pending(ctx context.Context, engagementID shared.ID) ([]agent.ProposedAction, error)

type AssessmentComparisonRepository added in v0.2.0

type AssessmentComparisonRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentComparisonRepository added in v0.2.0

func NewAssessmentComparisonRepository(pool *pgxpool.Pool) *AssessmentComparisonRepository

func (*AssessmentComparisonRepository) CreateQueued added in v0.2.0

func (*AssessmentComparisonRepository) Get added in v0.2.0

func (repository *AssessmentComparisonRepository) Get(ctx context.Context, tenantID, comparisonID shared.ID) (assessmentcomparison.Comparison, error)

func (*AssessmentComparisonRepository) GetAssessmentComparisonBacklog added in v0.2.0

func (repository *AssessmentComparisonRepository) GetAssessmentComparisonBacklog(ctx context.Context, tenantID shared.ID) (ports.AssessmentComparisonBacklog, error)

func (*AssessmentComparisonRepository) GetByInputHash added in v0.2.0

func (repository *AssessmentComparisonRepository) GetByInputHash(ctx context.Context, tenantID shared.ID, inputHash string) (assessmentcomparison.Comparison, error)

func (*AssessmentComparisonRepository) GetItem added in v0.2.0

func (repository *AssessmentComparisonRepository) GetItem(ctx context.Context, tenantID, comparisonID, itemID shared.ID) (assessmentcomparison.Item, error)

func (*AssessmentComparisonRepository) GetMetadata added in v0.2.0

func (repository *AssessmentComparisonRepository) GetMetadata(ctx context.Context, tenantID, comparisonID shared.ID) (assessmentcomparison.Comparison, error)

func (*AssessmentComparisonRepository) ListFailedAssessmentComparisons added in v0.2.0

func (repository *AssessmentComparisonRepository) ListFailedAssessmentComparisons(ctx context.Context, tenantID shared.ID, limit int) ([]assessmentcomparison.Comparison, error)

func (*AssessmentComparisonRepository) ListItems added in v0.2.0

func (*AssessmentComparisonRepository) ListMetadataByCycle added in v0.2.0

func (repository *AssessmentComparisonRepository) ListMetadataByCycle(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentcomparison.Comparison, error)

func (*AssessmentComparisonRepository) SummarizeItems added in v0.2.0

func (repository *AssessmentComparisonRepository) SummarizeItems(ctx context.Context, tenantID, comparisonID shared.ID, scope assessmentcomparison.Scope) (assessmentcomparison.Summary, error)

func (*AssessmentComparisonRepository) UpdateCAS added in v0.2.0

func (repository *AssessmentComparisonRepository) UpdateCAS(ctx context.Context, comparison assessmentcomparison.Comparison, expectedVersion int64) error

type AssessmentCycleBackfillRepository added in v0.2.0

type AssessmentCycleBackfillRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentCycleBackfillRepository added in v0.2.0

func NewAssessmentCycleBackfillRepository(pool *pgxpool.Pool) *AssessmentCycleBackfillRepository

func (*AssessmentCycleBackfillRepository) AcquireAssessmentCycleBackfillRun added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) AcquireAssessmentCycleBackfillRun(ctx context.Context, request ports.AssessmentCycleBackfillAcquireRequest) (run ports.AssessmentCycleBackfillRun, resumed bool, err error)

func (*AssessmentCycleBackfillRepository) AdvanceAssessmentCycleBackfillRun added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) AdvanceAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken, checkpoint shared.ID, now time.Time, leaseDuration time.Duration) (run ports.AssessmentCycleBackfillRun, err error)

func (*AssessmentCycleBackfillRepository) CommitAssessmentCycleBackfillItem added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) CommitAssessmentCycleBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, build func(context.Context) (ports.AssessmentCycleBackfillItem, error)) (item ports.AssessmentCycleBackfillItem, created bool, err error)

func (*AssessmentCycleBackfillRepository) FinishAssessmentCycleBackfillRun added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) FinishAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken shared.ID, state ports.AssessmentCycleBackfillState, now time.Time) (run ports.AssessmentCycleBackfillRun, err error)

func (*AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillItem added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillItem(ctx context.Context, tenantID, runID, assessmentID shared.ID) (item ports.AssessmentCycleBackfillItem, err error)

func (*AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillRun added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentCycleBackfillRun, err error)

type AssessmentCycleIntegrityRepository added in v0.2.0

type AssessmentCycleIntegrityRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentCycleIntegrityRepository added in v0.2.0

func NewAssessmentCycleIntegrityRepository(pool *pgxpool.Pool) *AssessmentCycleIntegrityRepository

func (*AssessmentCycleIntegrityRepository) AcquireAssessmentCycleIntegrityRun added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) AcquireAssessmentCycleIntegrityRun(ctx context.Context, request ports.AssessmentCycleIntegrityAcquireRequest) (run ports.AssessmentCycleIntegrityRun, resumed bool, err error)

func (*AssessmentCycleIntegrityRepository) AdvanceAssessmentCycleIntegrityRun added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) AdvanceAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken, checkpoint shared.ID, now time.Time, leaseDuration time.Duration) (run ports.AssessmentCycleIntegrityRun, err error)

func (*AssessmentCycleIntegrityRepository) AssessmentCycleIntegrityGeneration added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) AssessmentCycleIntegrityGeneration(ctx context.Context, tenantID shared.ID) (generation int64, err error)

func (*AssessmentCycleIntegrityRepository) CountAssessmentCycleIntegritySubjects added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) CountAssessmentCycleIntegritySubjects(ctx context.Context, tenantID shared.ID, snapshotAt time.Time) (eligible int, memberships int, err error)

func (*AssessmentCycleIntegrityRepository) FinishAssessmentCycleIntegrityRun added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) FinishAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken shared.ID, state ports.AssessmentCycleIntegrityState, now time.Time) (run ports.AssessmentCycleIntegrityRun, err error)

func (*AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegrityRun added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentCycleIntegrityRun, err error)

func (*AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegritySubject added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegritySubject(ctx context.Context, tenantID, runID, assessmentID shared.ID) (result ports.AssessmentCycleIntegritySubjectResult, err error)

func (*AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegrityFindings added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegrityFindings(ctx context.Context, tenantID, runID shared.ID) (findings []ports.AssessmentCycleIntegrityFinding, err error)

func (*AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegritySubjects added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegritySubjects(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, limit int) (subjects []ports.AssessmentCycleIntegritySubject, err error)

func (*AssessmentCycleIntegrityRepository) SaveAssessmentCycleIntegritySubject added in v0.2.0

func (repository *AssessmentCycleIntegrityRepository) SaveAssessmentCycleIntegritySubject(ctx context.Context, leaseToken shared.ID, now time.Time, result ports.AssessmentCycleIntegritySubjectResult, findings []ports.AssessmentCycleIntegrityFinding) (created bool, err error)

type AssessmentCycleRepository added in v0.2.0

type AssessmentCycleRepository struct {
	// contains filtered or unexported fields
}

AssessmentCycleRepository persists AssessmentCycle aggregates and their members to PostgreSQL.

func NewAssessmentCycleRepository added in v0.2.0

func NewAssessmentCycleRepository(pool *pgxpool.Pool) *AssessmentCycleRepository

NewAssessmentCycleRepository returns a repository backed by the given pool.

func (*AssessmentCycleRepository) CommitClosure added in v0.2.0

func (*AssessmentCycleRepository) CreateCycle added in v0.2.0

func (*AssessmentCycleRepository) CreateMember added in v0.2.0

func (r *AssessmentCycleRepository) CreateMember(ctx context.Context, member *assessmentcycle.Member) error

func (*AssessmentCycleRepository) DeleteCycle added in v0.2.0

func (r *AssessmentCycleRepository) DeleteCycle(ctx context.Context, tenantID, cycleID shared.ID) error

func (*AssessmentCycleRepository) DeleteMember added in v0.2.0

func (r *AssessmentCycleRepository) DeleteMember(ctx context.Context, tenantID, cycleID, assessmentID shared.ID) error

func (*AssessmentCycleRepository) GetActiveClosureManifest added in v0.2.0

func (r *AssessmentCycleRepository) GetActiveClosureManifest(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentclosure.Manifest, error)

func (*AssessmentCycleRepository) GetClosureManifest added in v0.2.0

func (r *AssessmentCycleRepository) GetClosureManifest(ctx context.Context, tenantID, cycleID, manifestID shared.ID) (*assessmentclosure.Manifest, error)

func (*AssessmentCycleRepository) GetClosureReport added in v0.2.0

func (r *AssessmentCycleRepository) GetClosureReport(ctx context.Context, tenantID, cycleID, manifestID shared.ID, rendererVersion string) (ports.AssessmentClosureReportArtifact, error)

func (*AssessmentCycleRepository) GetCycle added in v0.2.0

func (r *AssessmentCycleRepository) GetCycle(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentcycle.AssessmentCycle, error)

func (*AssessmentCycleRepository) GetCycleByAssessment added in v0.2.0

func (r *AssessmentCycleRepository) GetCycleByAssessment(ctx context.Context, tenantID, assessmentID shared.ID) (*assessmentcycle.AssessmentCycle, error)

func (*AssessmentCycleRepository) GetMember added in v0.2.0

func (r *AssessmentCycleRepository) GetMember(ctx context.Context, tenantID, cycleID, assessmentID shared.ID) (*assessmentcycle.Member, error)

func (*AssessmentCycleRepository) ListClosureManifests added in v0.2.0

func (r *AssessmentCycleRepository) ListClosureManifests(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentclosure.Manifest, error)

func (*AssessmentCycleRepository) ListCycles added in v0.2.0

func (*AssessmentCycleRepository) ListMembers added in v0.2.0

func (r *AssessmentCycleRepository) ListMembers(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentcycle.Member, error)

func (*AssessmentCycleRepository) ListMigrationPendingAssessments added in v0.2.0

func (*AssessmentCycleRepository) LockCycleForUpdate added in v0.2.0

func (r *AssessmentCycleRepository) LockCycleForUpdate(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentcycle.AssessmentCycle, error)

func (*AssessmentCycleRepository) NextManifestVersion added in v0.2.0

func (r *AssessmentCycleRepository) NextManifestVersion(ctx context.Context, tenantID, cycleID shared.ID) (int64, error)

func (*AssessmentCycleRepository) ReopenClosure added in v0.2.0

func (*AssessmentCycleRepository) SaveClosureReport added in v0.2.0

func (*AssessmentCycleRepository) UpdateCycleCAS added in v0.2.0

func (r *AssessmentCycleRepository) UpdateCycleCAS(ctx context.Context, cycle *assessmentcycle.AssessmentCycle, expectedVersion int64) error

func (*AssessmentCycleRepository) UpdateMemberCAS added in v0.2.0

func (r *AssessmentCycleRepository) UpdateMemberCAS(ctx context.Context, member *assessmentcycle.Member, expectedVersion int64) error

type AssessmentCycleRequestRepository added in v0.2.0

type AssessmentCycleRequestRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentCycleRequestRepository added in v0.2.0

func NewAssessmentCycleRequestRepository(pool *pgxpool.Pool) *AssessmentCycleRequestRepository

func (*AssessmentCycleRequestRepository) AbortAssessmentCycleRequest added in v0.2.0

func (repository *AssessmentCycleRequestRepository) AbortAssessmentCycleRequest(ctx context.Context, scope ports.AssessmentCycleRequestScope, requestHash string) error

func (*AssessmentCycleRequestRepository) BeginAssessmentCycleRequest added in v0.2.0

func (repository *AssessmentCycleRequestRepository) BeginAssessmentCycleRequest(ctx context.Context, request ports.AssessmentCycleRequest) (ports.AssessmentCycleRequest, bool, error)

func (*AssessmentCycleRequestRepository) CompleteAssessmentCycleRequest added in v0.2.0

func (repository *AssessmentCycleRequestRepository) CompleteAssessmentCycleRequest(ctx context.Context, scope ports.AssessmentCycleRequestScope, requestHash string, statusCode int, responseBody []byte, completedAt time.Time) error

type AssessmentRelationshipRepository added in v0.2.0

type AssessmentRelationshipRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentRelationshipRepository added in v0.2.0

func NewAssessmentRelationshipRepository(pool *pgxpool.Pool) *AssessmentRelationshipRepository

func (*AssessmentRelationshipRepository) CreateCandidate added in v0.2.0

func (repository *AssessmentRelationshipRepository) CreateCandidate(ctx context.Context, candidate assessmentrelationship.Candidate) (record assessmentrelationship.Record, created bool, err error)

func (*AssessmentRelationshipRepository) DecideCandidateCAS added in v0.2.0

func (*AssessmentRelationshipRepository) GetCandidate added in v0.2.0

func (repository *AssessmentRelationshipRepository) GetCandidate(ctx context.Context, tenantID, candidateID shared.ID) (record assessmentrelationship.Record, err error)

func (*AssessmentRelationshipRepository) ListCandidates added in v0.2.0

type AssessmentSnapshotBackfillRepository added in v0.2.0

type AssessmentSnapshotBackfillRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentSnapshotBackfillRepository added in v0.2.0

func NewAssessmentSnapshotBackfillRepository(pool *pgxpool.Pool) *AssessmentSnapshotBackfillRepository

func (*AssessmentSnapshotBackfillRepository) AcquireAssessmentSnapshotBackfillRun added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) AcquireAssessmentSnapshotBackfillRun(ctx context.Context, request ports.AssessmentSnapshotBackfillAcquireRequest) (run ports.AssessmentSnapshotBackfillRun, resumed bool, err error)

func (*AssessmentSnapshotBackfillRepository) AdvanceAssessmentSnapshotBackfillRun added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) AdvanceAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken, checkpoint shared.ID, now time.Time, leaseDuration time.Duration) (run ports.AssessmentSnapshotBackfillRun, err error)

func (*AssessmentSnapshotBackfillRepository) CommitAssessmentSnapshotBackfillItem added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) CommitAssessmentSnapshotBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, build func(context.Context) (ports.AssessmentSnapshotBackfillItem, error)) (item ports.AssessmentSnapshotBackfillItem, created bool, err error)

func (*AssessmentSnapshotBackfillRepository) FinishAssessmentSnapshotBackfillRun added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) FinishAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken shared.ID, state ports.AssessmentSnapshotBackfillState, now time.Time) (run ports.AssessmentSnapshotBackfillRun, err error)

func (*AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillItem added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillItem(ctx context.Context, tenantID, runID, assessmentID shared.ID) (item ports.AssessmentSnapshotBackfillItem, err error)

func (*AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillRun added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentSnapshotBackfillRun, err error)

type AssessmentSnapshotRepository added in v0.2.0

type AssessmentSnapshotRepository struct {
	// contains filtered or unexported fields
}

func NewAssessmentSnapshotRepository added in v0.2.0

func NewAssessmentSnapshotRepository(pool *pgxpool.Pool) *AssessmentSnapshotRepository

func (*AssessmentSnapshotRepository) CreateFinalizedCAS added in v0.2.0

func (repository *AssessmentSnapshotRepository) CreateFinalizedCAS(ctx context.Context, snapshot *assessmentsnapshot.Snapshot, expectedDefaultVersion int64) (*assessmentsnapshot.Snapshot, bool, error)

func (*AssessmentSnapshotRepository) CreateLegacyProjection added in v0.2.0

func (repository *AssessmentSnapshotRepository) CreateLegacyProjection(ctx context.Context, snapshot *assessmentsnapshot.Snapshot) (*assessmentsnapshot.Snapshot, bool, error)

func (*AssessmentSnapshotRepository) Get added in v0.2.0

func (repository *AssessmentSnapshotRepository) Get(ctx context.Context, tenantID, snapshotID shared.ID) (*assessmentsnapshot.Snapshot, error)

func (*AssessmentSnapshotRepository) GetByRequestKey added in v0.2.0

func (repository *AssessmentSnapshotRepository) GetByRequestKey(ctx context.Context, tenantID, assessmentID shared.ID, requestKey string) (*assessmentsnapshot.Snapshot, error)

func (*AssessmentSnapshotRepository) GetDefault added in v0.2.0

func (repository *AssessmentSnapshotRepository) GetDefault(ctx context.Context, tenantID, assessmentID shared.ID) (*assessmentsnapshot.Snapshot, ports.AssessmentSnapshotDefault, error)

func (*AssessmentSnapshotRepository) ListAssessmentSnapshots added in v0.2.0

func (*AssessmentSnapshotRepository) ListByAssessment added in v0.2.0

func (repository *AssessmentSnapshotRepository) ListByAssessment(ctx context.Context, tenantID, assessmentID shared.ID) ([]assessmentsnapshot.Snapshot, error)

type AssetRepository added in v0.1.8

type AssetRepository struct {
	// contains filtered or unexported fields
}

AssetRepository is the Postgres-backed fleet asset model. It is the first store to route every operation through WithTenant, so the Row Level Security policies on fleet_assets, fleet_asset_edges and fleet_business_services (migration 0058, using the 0057 procedure) enforce tenant isolation at the database. A query that bypassed WithTenant would resolve the tenant to NULL and see nothing.

func NewAssetRepository added in v0.1.8

func NewAssetRepository(pool *pgxpool.Pool) *AssetRepository

NewAssetRepository constructs the Postgres asset repository.

func (*AssetRepository) AssignEngagementBusinessAsset added in v0.1.8

func (r *AssetRepository) AssignEngagementBusinessAsset(ctx context.Context, tenantID, engagementID, assetID shared.ID) error

func (*AssetRepository) CountBusinessAssetsByCriticality added in v0.2.0

func (r *AssetRepository) CountBusinessAssetsByCriticality(ctx context.Context, tenantID shared.ID) (map[asset.Criticality]int, error)

CountBusinessAssetsByCriticality aggregates in the database rather than shipping rows. Postgres answers it from fleet_business_services without materialising the assets, so the cost does not grow with the size of the response the caller wanted.

func (*AssetRepository) CreateBusinessAsset added in v0.1.8

func (r *AssetRepository) CreateBusinessAsset(ctx context.Context, a *asset.BusinessAsset) error

func (*AssetRepository) GetAssetByID added in v0.2.0

func (r *AssetRepository) GetAssetByID(ctx context.Context, tenantID, id shared.ID) (*asset.Asset, error)

GetAssetByID returns one canonical technical asset by server-issued ID under tenant RLS. It is a narrow extension for control-plane policy admission; normal asset observation remains keyed by the natural (tenant, kind, key) identity. Cross-tenant and missing ids both resolve to ErrNotFound so a caller cannot distinguish another tenant's asset from an absent one.

func (*AssetRepository) GetAssetByKey added in v0.1.8

func (r *AssetRepository) GetAssetByKey(ctx context.Context, tenantID shared.ID, kind asset.Kind, key string) (*asset.Asset, error)

GetAssetByKey returns the asset for (tenantID, kind, key) or shared.ErrNotFound.

func (*AssetRepository) GetBusinessAssetByID added in v0.1.8

func (r *AssetRepository) GetBusinessAssetByID(ctx context.Context, tenantID, id shared.ID) (*asset.BusinessAsset, error)

func (*AssetRepository) GetBusinessAssetByKey added in v0.1.8

func (r *AssetRepository) GetBusinessAssetByKey(ctx context.Context, tenantID shared.ID, key string) (*asset.BusinessAsset, error)

func (*AssetRepository) ListAssets added in v0.1.8

func (r *AssetRepository) ListAssets(ctx context.Context, tenantID shared.ID) ([]*asset.Asset, error)

ListAssets returns the tenant's assets ordered by (kind, key).

func (*AssetRepository) ListBusinessAssetProjects added in v0.1.8

func (r *AssetRepository) ListBusinessAssetProjects(ctx context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)

func (*AssetRepository) ListBusinessAssetTechnicalAssets added in v0.1.8

func (r *AssetRepository) ListBusinessAssetTechnicalAssets(ctx context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)

func (*AssetRepository) ListBusinessAssets added in v0.1.8

func (r *AssetRepository) ListBusinessAssets(ctx context.Context, tenantID shared.ID) ([]*asset.BusinessAsset, error)

func (*AssetRepository) ListBusinessAssetsPage added in v0.2.0

func (r *AssetRepository) ListBusinessAssetsPage(ctx context.Context, tenantID shared.ID, query ports.BusinessAssetQuery) ([]*asset.BusinessAsset, int, error)

ListBusinessAssetsPage filters, orders and pages in the database. The count is of rows matching the filter before the page is applied, so a caller can report a total without fetching it.

The count rides along on the page query as a window aggregate rather than being read by a separate statement. WithTenant runs at READ COMMITTED, where each statement takes its own snapshot, so a separate COUNT could answer 1 while the SELECT beside it returned nothing, and the inventory would render "1 result" over an empty table.

Ordering is by byte value, not by the database's collation. The in-memory twin sorts Go strings, and a Postgres initialised with ICU orders "_infra" before "Billing" where Go does the reverse, which would put different rows on page N depending on which store answered.

func (*AssetRepository) ListEdges added in v0.1.8

func (r *AssetRepository) ListEdges(ctx context.Context, tenantID shared.ID) ([]*asset.Edge, error)

ListEdges returns the tenant's edges ordered by (from, to, kind, provenance).

func (*AssetRepository) ListEngagementsByBusinessAsset added in v0.1.8

func (r *AssetRepository) ListEngagementsByBusinessAsset(ctx context.Context, tenantID, assetID shared.ID) ([]*engagement.Engagement, error)

func (*AssetRepository) ReplaceBusinessAssetProjects added in v0.1.8

func (r *AssetRepository) ReplaceBusinessAssetProjects(ctx context.Context, tenantID, assetID shared.ID, links []asset.ComponentMembership) error

func (*AssetRepository) ReplaceBusinessAssetTechnicalAssets added in v0.1.8

func (r *AssetRepository) ReplaceBusinessAssetTechnicalAssets(ctx context.Context, tenantID, assetID shared.ID, links []asset.ComponentMembership) error

func (*AssetRepository) UpdateBusinessAsset added in v0.1.8

func (r *AssetRepository) UpdateBusinessAsset(ctx context.Context, a *asset.BusinessAsset, expectedVersion int) error

func (*AssetRepository) UpsertAsset added in v0.1.8

func (r *AssetRepository) UpsertAsset(ctx context.Context, a *asset.Asset) error

UpsertAsset inserts or updates by the (tenant_id, kind, key) natural key, preserving the id and created_at of an existing row so re-observation does not churn identity. A new host row that would take its reporting agent past the per-agent cap is refused by the fleet_assets trigger (migration 0132) and surfaces as shared.ErrForbidden.

func (*AssetRepository) UpsertEdge added in v0.1.8

func (r *AssetRepository) UpsertEdge(ctx context.Context, e *asset.Edge) error

UpsertEdge inserts the edge idempotently by its full natural key.

type AttackPathStore added in v0.1.8

type AttackPathStore struct {
	// contains filtered or unexported fields
}

AttackPathStore persists derived asset-to-finding links behind PostgreSQL RLS.

func NewAttackPathStore added in v0.1.8

func NewAttackPathStore(pool *pgxpool.Pool) *AttackPathStore

func (*AttackPathStore) ListBindings added in v0.1.8

func (s *AttackPathStore) ListBindings(ctx context.Context, tenantID shared.ID) ([]attackpath.Binding, error)

func (*AttackPathStore) ReplaceBindings added in v0.1.8

func (s *AttackPathStore) ReplaceBindings(ctx context.Context, tenantID, engagementID, producer shared.ID, bindings []attackpath.Binding) error

type AuditLog

type AuditLog struct {
	// contains filtered or unexported fields
}

AuditLog is an append-only, attributable audit log on PostgreSQL.

func NewAuditLog

func NewAuditLog(pool *pgxpool.Pool) *AuditLog

NewAuditLog returns an audit log backed by the given pool.

func (*AuditLog) List

func (l *AuditLog) List(ctx context.Context, limit int) (out []ports.AuditEntry, err error)

List returns the calling tenant's most recent audit entries (newest first), capped at limit.

func (*AuditLog) MigrationMetadata added in v0.2.0

func (l *AuditLog) MigrationMetadata(ctx context.Context) ([]ports.MigrationMetadata, error)

MigrationMetadata returns each migration's latest recorded state in version order. It makes no schema changes and is intended for offline restore verification.

func (*AuditLog) Record

func (l *AuditLog) Record(ctx context.Context, e ports.AuditEntry) error

func (*AuditLog) RecordOnce added in v0.1.8

func (l *AuditLog) RecordOnce(ctx context.Context, e ports.AuditEntry) error

Record appends an immutable audit entry (INSERT only – never update or delete), chaining it to the previous row. A transaction-scoped advisory lock serializes the read-head/insert so concurrent writers cannot fork the chain. The fork-guard unique index (migration 0033) is defense-in-depth on top of the lock: if the lock is ever bypassed, a concurrent append yields a 23505 unique violation – Record then re-reads the advanced head and re-chains (bounded), parity with the evidence store, rather than surfacing an opaque error. On the normal locked path the conflict is unreachable and the loop runs once. RecordOnce implements ports.IdempotentAuditLogger. Entries with a deterministic metadata idempotency_key are recovered without adding a second chain link.

func (*AuditLog) Verify

func (l *AuditLog) Verify(ctx context.Context) (audit.Report, error)

Verify examines only the calling tenant's rows. The historical audit chain is global, so omitted links make a tenant-only result unavailable rather than falsely claiming a complete integrity check. Server maintenance code may use VerifyGlobal.

func (*AuditLog) VerifyGlobal added in v0.1.8

func (l *AuditLog) VerifyGlobal(ctx context.Context) (audit.Report, error)

VerifyGlobal re-derives the complete, globally linked audit chain. It deliberately bypasses tenant visibility and must only be called by server-side maintenance code, never by a tenant HTTP endpoint.

type BaselineRepository added in v0.2.0

type BaselineRepository struct {
	// contains filtered or unexported fields
}

BaselineRepository is the Postgres tier for the behavioral-baseline projection (Phase D / D5). A baseline is a MUTABLE projection keyed by (tenant, group), upserted in place; integrity is enforced by the domain (baseline.NewBaselineFrom validates the accumulators on load) and the usecase. Every method runs under the authenticated ctx tenant via WithContextTenant (RLS) with an explicit tenant_id predicate as defense-in-depth. Reached only through ports.BaselineStore.

func NewBaselineRepository added in v0.2.0

func NewBaselineRepository(pool *pgxpool.Pool) *BaselineRepository

NewBaselineRepository constructs the baseline store over a pgx pool.

func (*BaselineRepository) Load added in v0.2.0

Load returns the record for a key, or shared.ErrNotFound.

func (*BaselineRepository) Save added in v0.2.0

Save upserts a baseline record for its (tenant, group) key.

type CloudObservationStore added in v0.1.8

type CloudObservationStore struct {
	// contains filtered or unexported fields
}

CloudObservationStore atomically replaces producer ownership only after a complete target snapshot.

func NewCloudObservationStore added in v0.1.8

func NewCloudObservationStore(pool *pgxpool.Pool) *CloudObservationStore

func (*CloudObservationStore) ReconcileCloudObservations added in v0.1.8

func (s *CloudObservationStore) ReconcileCloudObservations(ctx context.Context, tenantID, engagementID shared.ID, producer string, evidenceID shared.ID, assets, findings []shared.ID, edges []string, complete bool) error

type CloudRunStore added in v0.1.8

type CloudRunStore struct {
	// contains filtered or unexported fields
}

CloudRunStore persists the tenant-scoped CSPM lifecycle under RLS.

func NewCloudRunStore added in v0.1.8

func NewCloudRunStore(pool *pgxpool.Pool) *CloudRunStore

func (*CloudRunStore) EnqueueCloudRun added in v0.1.8

func (s *CloudRunStore) EnqueueCloudRun(ctx context.Context, run cloudposture.Run, kind string, payload []byte) error

func (*CloudRunStore) GetCloudRun added in v0.1.8

func (s *CloudRunStore) GetCloudRun(ctx context.Context, tenantID, id shared.ID) (out cloudposture.Run, err error)

func (*CloudRunStore) SaveCloudRun added in v0.1.8

func (s *CloudRunStore) SaveCloudRun(ctx context.Context, run cloudposture.Run) error

type CommentRepository

type CommentRepository struct {
	// contains filtered or unexported fields
}

CommentRepository persists the per-finding comment thread to PostgreSQL.

func NewCommentRepository

func NewCommentRepository(pool *pgxpool.Pool) *CommentRepository

NewCommentRepository returns a repository backed by the given pool.

func (*CommentRepository) Add

Add inserts a comment (append-only; comments are not edited or deleted in app code).

func (*CommentRepository) ListByEngagementFinding

func (r *CommentRepository) ListByEngagementFinding(ctx context.Context, engagementID, findingID shared.ID) (out []finding.Comment, err error)

ListByEngagementFinding returns a finding's comments oldest-first, scoped to the engagement (no cross-engagement read).

type ComponentInventoryStore added in v0.1.8

type ComponentInventoryStore struct {
	// contains filtered or unexported fields
}

func NewComponentInventoryStore added in v0.1.8

func NewComponentInventoryStore(pool *pgxpool.Pool) *ComponentInventoryStore

func (*ComponentInventoryStore) ClaimInventoryWork added in v0.2.0

func (s *ComponentInventoryStore) ClaimInventoryWork(ctx context.Context, tenantID shared.ID, owner string, at time.Time, lease time.Duration, limit int) ([]sbom.InventoryWork, error)

func (*ComponentInventoryStore) CompleteInventoryPublication added in v0.2.0

func (s *ComponentInventoryStore) CompleteInventoryPublication(ctx context.Context, publication sbom.InventoryPublication, at time.Time) error

func (*ComponentInventoryStore) FinishInventoryWork added in v0.2.0

func (s *ComponentInventoryStore) FinishInventoryWork(ctx context.Context, work sbom.InventoryWork, owner string, state sbom.InventoryWorkState, reason string, nextAttemptAt, at time.Time) error

func (*ComponentInventoryStore) GetCurrentInventoryPublication added in v0.2.0

func (s *ComponentInventoryStore) GetCurrentInventoryPublication(ctx context.Context, tenantID, engagementID shared.ID, scope string) (sbom.InventoryPublication, error)

func (*ComponentInventoryStore) ListCurrentComponents added in v0.1.8

func (s *ComponentInventoryStore) ListCurrentComponents(ctx context.Context, query sbom.ComponentQuery) (sbom.ComponentPage, error)

func (*ComponentInventoryStore) ListCurrentComponentsByEngagement added in v0.2.0

func (s *ComponentInventoryStore) ListCurrentComponentsByEngagement(ctx context.Context, tenantID, engagementID shared.ID) ([]sbom.ComponentRecord, error)

ListCurrentComponentsByEngagement returns the components of the engagement's latest SBOM (id + name + package only — enough to resolve a vulnerable ComponentID to a package name for running-vs-installed matching). Unlike ListCurrentComponents it takes no package/CPE key, so it can enumerate all components. Tenant-scoped via WithTenant (RLS) + an explicit tenant_id predicate.

func (*ComponentInventoryStore) ListCurrentInventoryPublications added in v0.2.0

func (s *ComponentInventoryStore) ListCurrentInventoryPublications(ctx context.Context, tenantID shared.ID, cursor sbom.InventoryCursor, limit int) (sbom.InventoryPublicationPage, error)

func (*ComponentInventoryStore) ListSnapshotComponents added in v0.2.0

func (s *ComponentInventoryStore) ListSnapshotComponents(ctx context.Context, query sbom.SnapshotQuery) (sbom.ComponentPage, error)

type CorrelationStateRepository added in v0.2.0

type CorrelationStateRepository struct {
	// contains filtered or unexported fields
}

func NewCorrelationStateRepository added in v0.2.0

func NewCorrelationStateRepository(pool *pgxpool.Pool) *CorrelationStateRepository

func (*CorrelationStateRepository) BeginCorrelationSnapshot added in v0.2.0

func (r *CorrelationStateRepository) BeginCorrelationSnapshot(ctx context.Context, engagementID shared.ID, expected uint64, upper correlation.SourcePosition, asOf time.Time, digest string) (correlation.Checkpoint, error)

func (*CorrelationStateRepository) CommitCorrelationConsume added in v0.2.0

func (r *CorrelationStateRepository) CommitCorrelationConsume(ctx context.Context, engagementID shared.ID, expected uint64, next correlation.Checkpoint, added []correlation.Assignment, active []correlation.ActiveSession, snapshot correlation.SourcePosition, complete bool) error

func (*CorrelationStateRepository) ListStagedCorrelationSignals added in v0.2.0

func (r *CorrelationStateRepository) ListStagedCorrelationSignals(ctx context.Context, engagementID shared.ID, snapshot correlation.SourcePosition, after correlation.SignalPosition, limit int) ([]correlation.Signal, bool, error)

func (*CorrelationStateRepository) LoadCorrelationState added in v0.2.0

func (r *CorrelationStateRepository) LoadCorrelationState(ctx context.Context, engagementID shared.ID, ids []shared.ID, activeLimit int) (correlation.State, error)

func (*CorrelationStateRepository) StageCorrelationSignals added in v0.2.0

func (r *CorrelationStateRepository) StageCorrelationSignals(ctx context.Context, engagementID shared.ID, expected uint64, next correlation.Checkpoint, signals []correlation.Signal) error

type CoverageWindowRepository added in v0.2.0

type CoverageWindowRepository struct {
	// contains filtered or unexported fields
}

func NewCoverageWindowRepository added in v0.2.0

func NewCoverageWindowRepository(pool *pgxpool.Pool) (*CoverageWindowRepository, error)

func (*CoverageWindowRepository) AppendCoverageWindow added in v0.2.0

func (*CoverageWindowRepository) ListCoverageWindows added in v0.2.0

func (*CoverageWindowRepository) ListCoverageWindowsBounded added in v0.2.0

func (r *CoverageWindowRepository) ListCoverageWindowsBounded(ctx context.Context, q ports.CoverageWindowQuery, limit int) ([]sensorstate.CoverageWindow, error)

type DASTRunStore added in v0.2.0

type DASTRunStore struct {
	// contains filtered or unexported fields
}

DASTRunStore persists the tenant-scoped DAST run lifecycle under RLS.

func NewDASTRunStore added in v0.2.0

func NewDASTRunStore(pool *pgxpool.Pool) *DASTRunStore

func (*DASTRunStore) EnqueueDASTRun added in v0.2.0

func (s *DASTRunStore) EnqueueDASTRun(ctx context.Context, run dastrun.Run, kind string, payload []byte) error

func (*DASTRunStore) FinishRun added in v0.2.0

func (s *DASTRunStore) FinishRun(ctx context.Context, tenantID shared.ID, run dastrun.Run) (won bool, err error)

FinishRun writes the terminal run only if the stored row is still 'running' (compare-and-set), and reports whether the row was updated. A redelivered or lease-overlapping worker that lost the race changes nothing.

func (*DASTRunStore) GetDASTRun added in v0.2.0

func (s *DASTRunStore) GetDASTRun(ctx context.Context, tenantID, id shared.ID) (out dastrun.Run, err error)

func (*DASTRunStore) SaveDASTRun added in v0.2.0

func (s *DASTRunStore) SaveDASTRun(ctx context.Context, run dastrun.Run) error

SaveDASTRun upserts the run, but NEVER moves a terminal row backward: the ON CONFLICT update is guarded so a stored 'succeeded' or 'failed' row is frozen. This is defense-in-depth against a stale worker whose lease was reclaimed (its claim context is cancelled, so its writes should already fail, but a write that slips through the cancellation window must not un-terminalize a recorded outcome). FinishRun is the CAS that records terminal outcomes; SaveDASTRun only advances a run toward one.

func (*DASTRunStore) StartRun added in v0.2.0

func (s *DASTRunStore) StartRun(ctx context.Context, tenantID shared.ID, run dastrun.Run) (won bool, err error)

StartRun moves a run 'queued' -> 'running' only if the stored row is still 'queued' (compare-and-set), reporting whether it won. A redelivered or lease-overlapping worker that finds the run already running or terminal loses (won=false) and must not execute the probe.

type DetectionProvenanceRepository added in v0.2.0

type DetectionProvenanceRepository struct {
	// contains filtered or unexported fields
}

DetectionProvenanceRepository retains the current read projection and append-only transition facts.

func NewDetectionProvenanceRepository added in v0.2.0

func NewDetectionProvenanceRepository(pool *pgxpool.Pool) (*DetectionProvenanceRepository, error)

func (*DetectionProvenanceRepository) AdmitPending added in v0.2.0

func (*DetectionProvenanceRepository) AppendTransition added in v0.2.0

func (*DetectionProvenanceRepository) Current added in v0.2.0

func (r *DetectionProvenanceRepository) Current(ctx context.Context, engagementID, detectionID shared.ID) (detectionprovenance.Current, bool, error)

func (*DetectionProvenanceRepository) ListCurrent added in v0.2.0

func (r *DetectionProvenanceRepository) ListCurrent(ctx context.Context, engagementID shared.ID) ([]detectionprovenance.Current, error)

func (*DetectionProvenanceRepository) ListPending added in v0.2.0

func (*DetectionProvenanceRepository) ListReceivedTransitions added in v0.2.0

func (r *DetectionProvenanceRepository) ListReceivedTransitions(ctx context.Context, engagementID shared.ID) ([]detectionprovenance.Transition, error)

func (*DetectionProvenanceRepository) ListTransitions added in v0.2.0

func (r *DetectionProvenanceRepository) ListTransitions(ctx context.Context, engagementID, detectionID shared.ID) ([]detectionprovenance.Transition, error)

func (*DetectionProvenanceRepository) LoadReceivedTransitions added in v0.2.0

func (r *DetectionProvenanceRepository) LoadReceivedTransitions(ctx context.Context, engagementID shared.ID, detectionIDs []shared.ID) ([]detectionprovenance.Transition, error)

type DetectionRecordRepository added in v0.1.8

type DetectionRecordRepository struct {
	// contains filtered or unexported fields
}

DetectionRecordRepository persists the detection-ledger projection (migration 0074), tenant-scoped via WithTenant so RLS isolates one tenant's detections from another's. The evidence-chain link each row references is the permanent ledger; these rows are the queryable, retention-bounded projection.

func NewDetectionRecordRepository added in v0.1.8

func NewDetectionRecordRepository(pool *pgxpool.Pool) *DetectionRecordRepository

NewDetectionRecordRepository constructs the repository.

func (*DetectionRecordRepository) AppendDetection added in v0.1.8

func (r *DetectionRecordRepository) AppendDetection(ctx context.Context, rec detection.Record) error

AppendDetection stores one projection row, idempotent on (tenant_id, engagement_id, id): a row is immutable once written (provenance), so a re-delivery of the same detection in the same engagement does not overwrite it, while the same id in a DIFFERENT engagement is a distinct row (not dropped).

func (*DetectionRecordRepository) ClassCountsByAsset added in v0.2.0

func (r *DetectionRecordRepository) ClassCountsByAsset(ctx context.Context, assetID shared.ID, since time.Time) (map[detection.Class]int, error)

ClassCountsByAsset counts the non-expired detections observed on an asset at or after a cutoff, grouped by telemetry class, tenant-scoped by RLS. It feeds the behavior baseline's runtime-anomaly features (#822): the network / privilege / file per-class rates a process snapshot cannot carry.

func (*DetectionRecordRepository) CorrelationHighWater added in v0.2.0

func (r *DetectionRecordRepository) CorrelationHighWater(ctx context.Context, engagementID shared.ID, completed correlation.SourcePosition, retentionAsOf time.Time) (correlation.SourcePosition, bool, error)

func (*DetectionRecordRepository) DeleteDetection added in v0.2.0

func (r *DetectionRecordRepository) DeleteDetection(ctx context.Context, engagementID, detectionID shared.ID) (bool, error)

DeleteDetection removes one exact projection row and leaves permanent evidence and provenance intact.

func (*DetectionRecordRepository) HasDetection added in v0.1.8

func (r *DetectionRecordRepository) HasDetection(ctx context.Context, engagementID, id shared.ID) (bool, error)

HasDetection reports whether a record with this id already exists in the given engagement (ctx tenant).

func (*DetectionRecordRepository) LastBatchSequence added in v0.1.8

func (r *DetectionRecordRepository) LastBatchSequence(ctx context.Context, agentID shared.ID) (uint64, error)

LastBatchSequence returns the highest batch sequence for an agent in the ctx tenant (0 = none yet).

func (*DetectionRecordRepository) ListCorrelationSourcePage added in v0.2.0

func (r *DetectionRecordRepository) ListCorrelationSourcePage(ctx context.Context, engagementID shared.ID, after, through correlation.SourcePosition, retentionAsOf time.Time, limit int) ([]detection.Record, bool, error)

func (*DetectionRecordRepository) ListDetections added in v0.1.8

func (r *DetectionRecordRepository) ListDetections(ctx context.Context, engagementID shared.ID) ([]detection.Record, error)

ListDetections returns the non-expired records for an engagement, oldest first, tenant-scoped by RLS.

func (*DetectionRecordRepository) ListExpiredDetections added in v0.2.0

func (r *DetectionRecordRepository) ListExpiredDetections(ctx context.Context, engagementID shared.ID, cutoff time.Time) ([]shared.ID, error)

ListExpiredDetections returns eligible record ids without mutating the projection.

type EmulationRunRepository added in v0.1.8

type EmulationRunRepository struct {
	// contains filtered or unexported fields
}

EmulationRunRepository persists emulation runs and coverage (migration 0073), tenant-scoped via WithTenant so RLS isolates one tenant's coverage from another's.

func NewEmulationRunRepository added in v0.1.8

func NewEmulationRunRepository(pool *pgxpool.Pool) *EmulationRunRepository

NewEmulationRunRepository constructs the repository.

func (*EmulationRunRepository) SaveRun added in v0.1.8

func (r *EmulationRunRepository) SaveRun(ctx context.Context, run demu.Run) error

SaveRun writes the run and all its coverage rows in one transaction, so a partially-written run can never present as complete coverage.

type EndpointProcessRepository added in v0.2.0

type EndpointProcessRepository struct {
	// contains filtered or unexported fields
}

EndpointProcessRepository is the Postgres tier for the per-host process snapshot projection (B5). A snapshot is a MUTABLE projection keyed by (tenant, asset, entity), upserted in place. Every method runs under the authenticated ctx tenant via WithContextTenant (RLS) with an explicit tenant_id predicate as defense-in-depth. Reached only through ports.EndpointProcessStore.

func NewEndpointProcessRepository added in v0.2.0

func NewEndpointProcessRepository(pool *pgxpool.Pool) *EndpointProcessRepository

NewEndpointProcessRepository constructs the process snapshot store over a pgx pool.

func (*EndpointProcessRepository) ListRunningByAsset added in v0.2.0

func (r *EndpointProcessRepository) ListRunningByAsset(ctx context.Context, assetID shared.ID) ([]ports.ProcessSnapshot, error)

ListRunningByAsset returns the running snapshots for an asset, ordered by entity_id (COLLATE "C" so the SQL order matches the memory twin's Go bytewise ordering).

func (*EndpointProcessRepository) ReplaceRunningProcesses added in v0.2.0

func (r *EndpointProcessRepository) ReplaceRunningProcesses(ctx context.Context, assetID shared.ID, snapshots []ports.ProcessSnapshot) error

ReplaceRunningProcesses makes the asset's running set exactly the reported snapshots in one transaction: it upserts them (multi-row unnest) and marks every other currently-running row for that asset that the report omitted as not-running. Without this a process that exits between reports stays running=true forever. An empty report clears the asset's running set.

func (*EndpointProcessRepository) SaveProcesses added in v0.2.0

func (r *EndpointProcessRepository) SaveProcesses(ctx context.Context, snapshots []ports.ProcessSnapshot) error

SaveProcesses upserts snapshots by (tenant, asset, entity) in one transaction.

type EndpointTimelineRepository added in v0.2.0

type EndpointTimelineRepository struct {
	// contains filtered or unexported fields
}

EndpointTimelineRepository is the Postgres tier for the durable endpoint State Timeline (Phase B / B7). Every method runs under the authenticated ctx tenant via WithContextTenant, so RLS binds the partition, and reads carry an explicit tenant_id predicate as defense-in-depth. Appends are idempotent by (tenant, asset, event_id). Reached only through ports.EndpointTimelineStore.

func NewEndpointTimelineRepository added in v0.2.0

func NewEndpointTimelineRepository(pool *pgxpool.Pool) *EndpointTimelineRepository

NewEndpointTimelineRepository constructs the endpoint-timeline store over a pgx pool.

func (*EndpointTimelineRepository) AppendTimeline added in v0.2.0

func (r *EndpointTimelineRepository) AppendTimeline(ctx context.Context, list []endpoint.TimelineEntry) error

AppendTimeline persists the transitions idempotently (ON CONFLICT DO NOTHING on the (tenant, asset, event_id) key). Every entry's TenantID must equal the context tenant.

func (*EndpointTimelineRepository) LoadTimelineEntries added in v0.2.0

func (r *EndpointTimelineRepository) LoadTimelineEntries(ctx context.Context, assetID shared.ID, eventIDs []shared.ID) ([]endpoint.TimelineEntry, error)

LoadTimelineEntries returns exact source events without widening the correlation read to a time range.

func (*EndpointTimelineRepository) QueryTimeline added in v0.2.0

QueryTimeline returns the stored transitions matching q, ordered by (occurred_at, event_id).

type EngagementRepository

type EngagementRepository struct {
	// contains filtered or unexported fields
}

EngagementRepository persists engagements and their scope to PostgreSQL.

func NewEngagementRepository

func NewEngagementRepository(pool *pgxpool.Pool) *EngagementRepository

NewEngagementRepository returns a repository backed by the given pool.

func (*EngagementRepository) Create

Create inserts the engagement and its scope targets in one transaction.

func (*EngagementRepository) Delete

func (r *EngagementRepository) Delete(ctx context.Context, id shared.ID) error

Delete removes an engagement; ON DELETE CASCADE drops its scope, findings, comments, evidence, recon runs, and retests. Retained scan-run provenance deliberately uses RESTRICT and is surfaced as a conflict. Idempotent (no error if absent). Used to roll back a partially-materialized import.

func (*EngagementRepository) GetByHostAssetID added in v0.2.0

func (r *EngagementRepository) GetByHostAssetID(ctx context.Context, tenantID, assetID shared.ID) (out *engagement.Engagement, err error)

GetByHostAssetID loads the hidden fleet host vulnerability context for one Kind=host asset.

func (*EngagementRepository) GetByID

func (r *EngagementRepository) GetByID(ctx context.Context, id shared.ID) (out *engagement.Engagement, err error)

GetByID returns the engagement with its full scope WITHOUT a tenant predicate. It is the INTERNAL execution-gate read (see ports.EngagementRepository.GetByID): the scope/window/RoE guard and the worker/agent execution paths, which act on an engagement a queued/authorized run already belongs to. User-facing access uses GetByIDInTenant (below), which adds the tenant predicate that blocks cross-tenant reads.

func (*EngagementRepository) GetByIDInTenant

func (r *EngagementRepository) GetByIDInTenant(ctx context.Context, tenantID, id shared.ID) (out *engagement.Engagement, err error)

GetByIDInTenant loads an engagement scoped to tenantID. Empty input normalizes to the non-empty default tenant and never becomes a wildcard; cross-tenant access returns ErrNotFound.

func (*EngagementRepository) GetByProjectID

func (r *EngagementRepository) GetByProjectID(ctx context.Context, tenantID, projectID shared.ID) (out *engagement.Engagement, err error)

func (*EngagementRepository) List

func (r *EngagementRepository) List(ctx context.Context, tenantID shared.ID) (out []*engagement.Engagement, err error)

List returns the tenant's engagements, each with its scope loaded (consistent with the in-memory repository; the UI and the scope gate both rely on scope).

func (*EngagementRepository) ListAssessmentCycleBackfillEngagements added in v0.2.0

func (r *EngagementRepository) ListAssessmentCycleBackfillEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, limit int) (out []*engagement.Engagement, err error)

func (*EngagementRepository) ListAssessmentSnapshotBackfillEngagements added in v0.2.0

func (r *EngagementRepository) ListAssessmentSnapshotBackfillEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, limit int) ([]*engagement.Engagement, error)

func (*EngagementRepository) ListHostEngagements added in v0.2.0

func (r *EngagementRepository) ListHostEngagements(ctx context.Context, tenantID shared.ID) ([]*engagement.Engagement, error)

ListHostEngagements returns the tenant's hidden fleet host vulnerability contexts for operational aggregation (advisory reconciliation). Normal engagement lists continue to hide these rows.

func (*EngagementRepository) ListProjectEngagements added in v0.1.8

func (r *EngagementRepository) ListProjectEngagements(ctx context.Context, tenantID shared.ID) ([]*engagement.Engagement, error)

ListProjectEngagements returns the tenant's hidden Project analysis contexts for operational aggregation. Normal engagement lists remain unchanged and continue to hide these rows.

func (*EngagementRepository) ListPromotionReconciliationScopes added in v0.1.8

func (r *EngagementRepository) ListPromotionReconciliationScopes(ctx context.Context) ([]ports.PromotionReconciliationScope, error)

ListPromotionReconciliationScopes returns every non-project tenant and engagement pair for process-local recovery. Each engagement query remains RLS-scoped.

func (*EngagementRepository) ListReconciliationEngagements added in v0.1.8

func (r *EngagementRepository) ListReconciliationEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, limit int) (ports.ReconciliationEngagementPage, error)

func (*EngagementRepository) ListTenantIDs added in v0.1.8

func (r *EngagementRepository) ListTenantIDs(ctx context.Context) ([]shared.ID, error)

func (*EngagementRepository) ProjectContexts

func (r *EngagementRepository) ProjectContexts(ctx context.Context, tenantID shared.ID, projectIDs []shared.ID) (out map[shared.ID]*engagement.Engagement, err error)

Update persists an existing engagement aggregate: the engagement row and its full scope target set, replaced atomically in one transaction (E1 scope CRUD + lifecycle). Returns shared.ErrNotFound if the engagement does not exist. Unlike Create's deterministic scope PKs, the replace path uses generated IDs.

func (*EngagementRepository) Update

type EngagementSourceRepository added in v0.2.0

type EngagementSourceRepository struct {
	// contains filtered or unexported fields
}

func NewEngagementSourceRepository added in v0.2.0

func NewEngagementSourceRepository(pool *pgxpool.Pool) *EngagementSourceRepository

func (*EngagementSourceRepository) Create added in v0.2.0

func (*EngagementSourceRepository) Delete added in v0.2.0

func (r *EngagementSourceRepository) Delete(ctx context.Context, tenantID, engagementID shared.ID) (sourcepackage.Package, bool, error)

func (*EngagementSourceRepository) Get added in v0.2.0

func (r *EngagementSourceRepository) Get(ctx context.Context, tenantID, engagementID shared.ID) (sourcepackage.Package, error)

func (*EngagementSourceRepository) GetByLocator added in v0.2.0

func (r *EngagementSourceRepository) GetByLocator(ctx context.Context, tenantID shared.ID, locator string) (sourcepackage.Package, error)

func (*EngagementSourceRepository) GetByVersion added in v0.2.0

func (r *EngagementSourceRepository) GetByVersion(ctx context.Context, tenantID, engagementID, versionID shared.ID) (sourcepackage.Package, error)

func (*EngagementSourceRepository) ObjectUnreferenced added in v0.2.0

func (r *EngagementSourceRepository) ObjectUnreferenced(ctx context.Context, tenantID shared.ID, objectKey string) (bool, error)

type EvidenceStore

type EvidenceStore struct {
	// contains filtered or unexported fields
}

EvidenceStore persists the per-engagement hash-chained evidence ledger.

func NewEvidenceStore

func NewEvidenceStore(pool *pgxpool.Pool) *EvidenceStore

NewEvidenceStore returns a store backed by the given pool.

func (*EvidenceStore) Append

func (r *EvidenceStore) Append(ctx context.Context, items []evidence.Evidence) error

Append inserts sealed evidence items in order, in one transaction (append-only).

func (*EvidenceStore) Head

func (r *EvidenceStore) Head(ctx context.Context, engagementID shared.ID) (string, error)

Head returns the most recent sealed hash for an engagement ("" if the chain is empty). A real query error is returned (NOT swallowed as "empty") so the caller never forks the append-only chain on a transient DB failure.

func (*EvidenceStore) ListByEngagement

func (r *EvidenceStore) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []evidence.Evidence, err error)

ListByEngagement returns the engagement's evidence in chain order (oldest first).

func (*EvidenceStore) LookupSealedForFinding added in v0.1.8

func (r *EvidenceStore) LookupSealedForFinding(ctx context.Context, engagementID, findingID shared.ID, kind string) (evidence.Evidence, bool, error)

LookupSealedForFinding returns the most recent sealed evidence link of the given kind for the specified finding, or (zero, false, nil) if none exists.

type ExploitationChainRepository added in v0.1.8

type ExploitationChainRepository struct {
	// contains filtered or unexported fields
}

ExploitationChainRepository persists attack chains and their steps (migration 0072). Every method runs through WithTenant so Row Level Security isolates one tenant's chains from another's.

func NewExploitationChainRepository added in v0.1.8

func NewExploitationChainRepository(pool *pgxpool.Pool) *ExploitationChainRepository

NewExploitationChainRepository constructs the repository.

func (*ExploitationChainRepository) SaveChain added in v0.1.8

func (r *ExploitationChainRepository) SaveChain(ctx context.Context, chain *dexploit.Chain) error

SaveChain upserts the chain and replaces its steps atomically. The whole chain — cursor, state, and every step's state and evidence — is written in one transaction, so a halt or a crash cannot leave the persisted chain disagreeing with itself about which steps still owe cleanup.

type FindingLineageBackfillRepository added in v0.2.0

type FindingLineageBackfillRepository struct {
	// contains filtered or unexported fields
}

func NewFindingLineageBackfillRepository added in v0.2.0

func NewFindingLineageBackfillRepository(pool *pgxpool.Pool) *FindingLineageBackfillRepository

func (*FindingLineageBackfillRepository) AcquireFindingLineageBackfillRun added in v0.2.0

func (repository *FindingLineageBackfillRepository) AcquireFindingLineageBackfillRun(ctx context.Context, request ports.FindingLineageBackfillAcquireRequest) (run ports.FindingLineageBackfillRun, resumed bool, err error)

func (*FindingLineageBackfillRepository) AdvanceFindingLineageBackfillRun added in v0.2.0

func (repository *FindingLineageBackfillRepository) AdvanceFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken, checkpoint shared.ID, now time.Time, leaseDuration time.Duration) (run ports.FindingLineageBackfillRun, err error)

func (*FindingLineageBackfillRepository) CommitFindingLineageBackfillItem added in v0.2.0

func (repository *FindingLineageBackfillRepository) CommitFindingLineageBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, build func(context.Context) (ports.FindingLineageBackfillItem, error)) (item ports.FindingLineageBackfillItem, created bool, err error)

func (*FindingLineageBackfillRepository) FinishFindingLineageBackfillRun added in v0.2.0

func (repository *FindingLineageBackfillRepository) FinishFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, leaseToken shared.ID, state ports.FindingLineageBackfillState, now time.Time) (run ports.FindingLineageBackfillRun, err error)

func (*FindingLineageBackfillRepository) GetFindingLineageBackfillItem added in v0.2.0

func (repository *FindingLineageBackfillRepository) GetFindingLineageBackfillItem(ctx context.Context, tenantID, runID, sourceFindingID shared.ID) (item ports.FindingLineageBackfillItem, err error)

func (*FindingLineageBackfillRepository) GetFindingLineageBackfillRun added in v0.2.0

func (repository *FindingLineageBackfillRepository) GetFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.FindingLineageBackfillRun, err error)

func (*FindingLineageBackfillRepository) ListFindingLineageBackfillSources added in v0.2.0

func (repository *FindingLineageBackfillRepository) ListFindingLineageBackfillSources(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, producerFilters []string, limit int) ([]ports.FindingLineageBackfillSourceRow, error)

type FindingLineageRepository added in v0.2.0

type FindingLineageRepository struct {
	// contains filtered or unexported fields
}

func NewFindingLineageRepository added in v0.2.0

func NewFindingLineageRepository(pool *pgxpool.Pool) *FindingLineageRepository

func (*FindingLineageRepository) AppendAlias added in v0.2.0

func (repository *FindingLineageRepository) AppendAlias(ctx context.Context, alias findinglineage.Alias) (bool, error)

func (*FindingLineageRepository) AppendObservation added in v0.2.0

func (repository *FindingLineageRepository) AppendObservation(ctx context.Context, observation findinglineage.Observation) error

func (*FindingLineageRepository) AppendOverrideCAS added in v0.2.0

func (*FindingLineageRepository) AppendSkip added in v0.2.0

func (*FindingLineageRepository) CreateCandidate added in v0.2.0

func (repository *FindingLineageRepository) CreateCandidate(ctx context.Context, candidate findinglineage.MatchCandidate, supersessionEventID shared.ID) (findinglineage.MatchCandidate, bool, error)

func (*FindingLineageRepository) CreateIdentityWithObservation added in v0.2.0

func (repository *FindingLineageRepository) CreateIdentityWithObservation(ctx context.Context, identity findinglineage.Identity, observation findinglineage.Observation) error

func (*FindingLineageRepository) FindIdentitiesByAlias added in v0.2.0

func (repository *FindingLineageRepository) FindIdentitiesByAlias(ctx context.Context, tenantID, cycleID shared.ID, producerKind, findingKind string, schemaVersion int, targetCanonical, fingerprint string) ([]findinglineage.Identity, error)

func (*FindingLineageRepository) FindIdentitiesByFingerprint added in v0.2.0

func (repository *FindingLineageRepository) FindIdentitiesByFingerprint(ctx context.Context, tenantID, cycleID shared.ID, producerKind, findingKind string, schemaVersion int, targetCanonical, fingerprint string) ([]findinglineage.Identity, error)

func (*FindingLineageRepository) FindIdentitiesByProducerID added in v0.2.0

func (repository *FindingLineageRepository) FindIdentitiesByProducerID(ctx context.Context, tenantID, cycleID shared.ID, producerKind, findingKind, targetCanonical, sourceID string) ([]findinglineage.Identity, error)

func (*FindingLineageRepository) GetActiveOverride added in v0.2.0

func (repository *FindingLineageRepository) GetActiveOverride(ctx context.Context, tenantID, cycleID, sourceObservationID shared.ID) (findinglineage.OverrideEvent, error)

func (*FindingLineageRepository) GetCandidate added in v0.2.0

func (repository *FindingLineageRepository) GetCandidate(ctx context.Context, tenantID, cycleID, candidateID shared.ID) (findinglineage.MatchCandidate, error)

func (*FindingLineageRepository) GetIdentity added in v0.2.0

func (repository *FindingLineageRepository) GetIdentity(ctx context.Context, tenantID, cycleID, identityID shared.ID) (findinglineage.Identity, error)

func (*FindingLineageRepository) GetObservation added in v0.2.0

func (repository *FindingLineageRepository) GetObservation(ctx context.Context, tenantID, cycleID, observationID shared.ID) (findinglineage.Observation, error)

func (*FindingLineageRepository) GetObservationBySource added in v0.2.0

func (repository *FindingLineageRepository) GetObservationBySource(ctx context.Context, tenantID, cycleID, snapshotID shared.ID, producerKind, findingKind, targetCanonical, sourceFindingID, sourceOccurrenceID string) (findinglineage.Observation, error)

func (*FindingLineageRepository) ListActiveOverridesBySnapshot added in v0.2.0

func (repository *FindingLineageRepository) ListActiveOverridesBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.OverrideEvent, error)

func (*FindingLineageRepository) ListCandidateResolutions added in v0.2.0

func (repository *FindingLineageRepository) ListCandidateResolutions(ctx context.Context, tenantID, cycleID, candidateID shared.ID) ([]findinglineage.ResolutionEvent, error)

func (*FindingLineageRepository) ListObservationsBySnapshot added in v0.2.0

func (repository *FindingLineageRepository) ListObservationsBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.Observation, error)

func (*FindingLineageRepository) ListOpenCandidatesBySnapshot added in v0.2.0

func (repository *FindingLineageRepository) ListOpenCandidatesBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.MatchCandidate, error)

func (*FindingLineageRepository) ListOverrideEvents added in v0.2.0

func (repository *FindingLineageRepository) ListOverrideEvents(ctx context.Context, tenantID, cycleID, sourceObservationID shared.ID) ([]findinglineage.OverrideEvent, error)

func (*FindingLineageRepository) ListSkipsBySnapshot added in v0.2.0

func (repository *FindingLineageRepository) ListSkipsBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.SkipRecord, error)

func (*FindingLineageRepository) LockCorrelationNamespace added in v0.2.0

func (repository *FindingLineageRepository) LockCorrelationNamespace(ctx context.Context, tenantID, cycleID shared.ID, producerKind, findingKind, targetCanonical, discriminator string) error

func (*FindingLineageRepository) ResolveCandidateCAS added in v0.2.0

type FindingRepository

type FindingRepository struct {
	// contains filtered or unexported fields
}

FindingRepository persists findings to PostgreSQL, deduped per engagement.

func NewFindingRepository

func NewFindingRepository(pool *pgxpool.Pool) *FindingRepository

NewFindingRepository returns a repository backed by the given pool.

func (*FindingRepository) CheckVulnerabilityPrimaryFinding added in v0.2.0

func (r *FindingRepository) CheckVulnerabilityPrimaryFinding(ctx context.Context, tenantID, engagementID shared.ID, inventoryScope, advisoryID string) error

func (*FindingRepository) ClaimFindingProjection added in v0.1.8

func (r *FindingRepository) ClaimFindingProjection(ctx context.Context, tenantID, engagementID, judgmentID shared.ID, mode ports.FindingProjectionMode) error

ClaimFindingProjection atomically reserves the SAST or legacy runtime DAST projection mode for a judgment.

func (*FindingRepository) GetByEngagementAndDedupKey added in v0.2.0

func (r *FindingRepository) GetByEngagementAndDedupKey(ctx context.Context, engagementID shared.ID, dedupKey string) (finding.Finding, error)

GetByEngagementAndDedupKey resolves the row selected by the repository's unique engagement/dedup constraint. Unlike a computed finding ID, this is stable across the historical component-to-target workflow migration.

func (*FindingRepository) GetByEngagementAndID added in v0.1.8

func (r *FindingRepository) GetByEngagementAndID(ctx context.Context, engagementID, findingID shared.ID) (finding.Finding, error)

GetByEngagementAndID loads a single finding by engagement and finding ID. Returns shared.ErrNotFound if no such finding exists in the engagement.

func (*FindingRepository) LinkVulnerabilityFindingOccurrence added in v0.2.0

func (r *FindingRepository) LinkVulnerabilityFindingOccurrence(ctx context.Context, tenantID, engagementID, findingID, occurrenceID shared.ID, at time.Time) error

func (*FindingRepository) ListByEngagement

func (r *FindingRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []finding.Finding, err error)

func (*FindingRepository) ListPublishableByEngagement

func (r *FindingRepository) ListPublishableByEngagement(ctx context.Context, engagementID shared.ID) ([]finding.Finding, error)

ListPublishableByEngagement returns only the engagement's findings that clear the evidence gate. It reuses ListByEngagement and the single domain rule finding.Publishable, so the publishability policy lives in exactly one place (the domain) rather than being re-encoded in SQL.

func (*FindingRepository) MapVulnerabilityPrimaryFinding added in v0.2.0

func (r *FindingRepository) MapVulnerabilityPrimaryFinding(ctx context.Context, tenantID, engagementID shared.ID, inventoryScope, advisoryID string, findingID shared.ID, at time.Time) error

func (*FindingRepository) SetAssignee

func (r *FindingRepository) SetAssignee(ctx context.Context, engagementID, findingID shared.ID, assignee string, expectedVersion int) (out finding.Finding, err error)

SetAssignee sets the assignee with the same optimistic-concurrency guard.

func (*FindingRepository) SetEvidenceScore

func (r *FindingRepository) SetEvidenceScore(ctx context.Context, engagementID, findingID shared.ID, score, expectedVersion int) (out finding.Finding, err error)

SetEvidenceScore sets a finding's evidence score with the same optimistic-concurrency guard as UpdateStatus (the adversarial-verdict path): the row is updated only if version matches, then version is bumped. Note the Upsert ON CONFLICT set deliberately omits evidence_score, so this is the only path that moves it for an already-stored finding.

func (*FindingRepository) SummarizeOpenFindingsByEngagements added in v0.2.0

func (r *FindingRepository) SummarizeOpenFindingsByEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.VulnerabilitySummary, error)

SummarizeOpenFindingsByEngagements aggregates open findings of every kind per engagement.

func (*FindingRepository) SummarizeVulnerabilitiesByEngagements added in v0.2.0

func (r *FindingRepository) SummarizeVulnerabilitiesByEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.VulnerabilitySummary, error)

ListByEngagement returns the engagement's findings, highest risk first (CISA KEV, then EPSS x CVSS, then severity). SummarizeVulnerabilitiesByEngagements aggregates open SCA vulnerability findings per engagement in one GROUP BY, for list views that would otherwise load every finding of every context.

func (*FindingRepository) UpdateStatus

func (r *FindingRepository) UpdateStatus(ctx context.Context, engagementID, findingID shared.ID, status finding.Status, expectedVersion int) (out finding.Finding, err error)

UpdateStatus sets the triage status with optimistic concurrency: the row is updated only if version matches expectedVersion, then version is bumped. On a miss it distinguishes ErrConflict (exists, version moved) from ErrNotFound.

func (*FindingRepository) Upsert

func (r *FindingRepository) Upsert(ctx context.Context, findings []finding.Finding) error

Upsert inserts or updates findings, deduped on (engagement_id, dedup_key). On conflict it updates machine-owned data, preserves id, status (triage), assignee, and created_at, and bumps version only when machine-owned data changes.

type FleetAgentRepository added in v0.1.8

type FleetAgentRepository struct {
	// contains filtered or unexported fields
}

FleetAgentRepository is the Postgres-backed fleet agent identity store. Every method runs through WithTenant so Row Level Security (migration 0060 via the 0057 procedure) isolates by tenant. The auth lookup is tenant-scoped because the agent credential carries a non-secret tenant prefix.

func NewFleetAgentRepository added in v0.1.8

func NewFleetAgentRepository(pool *pgxpool.Pool) *FleetAgentRepository

NewFleetAgentRepository constructs the Postgres fleet agent repository.

func (*FleetAgentRepository) ConsumeEnrolToken added in v0.1.8

func (r *FleetAgentRepository) ConsumeEnrolToken(ctx context.Context, tenantID shared.ID, hash string, now time.Time) (*fleetagent.EnrolToken, error)

ConsumeEnrolToken atomically marks a usable token used and returns it; shared.ErrNotFound if no unused, unexpired token with that hash exists for the tenant.

func (*FleetAgentRepository) CreateAgent added in v0.1.8

func (r *FleetAgentRepository) CreateAgent(ctx context.Context, a *fleetagent.Agent) error

func (*FleetAgentRepository) CreateEnrolToken added in v0.1.8

func (r *FleetAgentRepository) CreateEnrolToken(ctx context.Context, t *fleetagent.EnrolToken) error

func (*FleetAgentRepository) Decommission added in v0.1.8

func (r *FleetAgentRepository) Decommission(ctx context.Context, tenantID, id shared.ID, now time.Time) error

Decommission marks the agent decommissioned on its own report (#412). It is an idempotent no-op on a REVOKED agent: an operator revocation is the stronger terminal state and a self-report must not overwrite it (the CASE leaves a revoked row untouched but still matches it, so only a truly missing agent yields ErrNotFound — matching the memory store's contract).

func (*FleetAgentRepository) GetAgent added in v0.1.8

func (r *FleetAgentRepository) GetAgent(ctx context.Context, tenantID, id shared.ID) (*fleetagent.Agent, error)

func (*FleetAgentRepository) Heartbeat added in v0.1.8

func (r *FleetAgentRepository) Heartbeat(ctx context.Context, tenantID, id shared.ID, platform, osVersion, agentVersion string, capabilities []string, now time.Time) error

func (*FleetAgentRepository) ListAgents added in v0.1.8

func (r *FleetAgentRepository) ListAgents(ctx context.Context, tenantID shared.ID) ([]*fleetagent.Agent, error)

func (*FleetAgentRepository) Revoke added in v0.1.8

func (r *FleetAgentRepository) Revoke(ctx context.Context, tenantID, id, by shared.ID, reason string, now time.Time) error

func (*FleetAgentRepository) SetFingerprint added in v0.1.8

func (r *FleetAgentRepository) SetFingerprint(ctx context.Context, tenantID, id shared.ID, fingerprint string, now time.Time) error

type FleetAuditRepository added in v0.2.0

type FleetAuditRepository struct {
	// contains filtered or unexported fields
}

func NewFleetAuditRepository added in v0.2.0

func NewFleetAuditRepository(pool *pgxpool.Pool) (*FleetAuditRepository, error)

func (*FleetAuditRepository) AcknowledgeFleetAudit added in v0.2.0

func (r *FleetAuditRepository) AcknowledgeFleetAudit(ctx context.Context, id string) error

AcknowledgeFleetAudit marks one of the calling tenant's intentions delivered. It is monotonic (COALESCE keeps the first completion) and tenant-scoped in SQL, so it cannot retire another tenant's outstanding obligation.

func (*FleetAuditRepository) ListPendingFleetAudits added in v0.2.0

func (r *FleetAuditRepository) ListPendingFleetAudits(ctx context.Context) (out []ports.FleetAuditIntent, err error)

ListPendingFleetAudits returns the calling tenant's committed-but-undelivered audit obligations, oldest first. The tenant predicate is explicit rather than left to RLS alone: a runtime role holding BYPASSRLS would otherwise let one tenant's recovery sweep read another tenant's obligations and write them into its own audit chain.

type FleetDesiredRepository added in v0.2.0

type FleetDesiredRepository struct {
	// contains filtered or unexported fields
}

func NewFleetDesiredRepository added in v0.2.0

func NewFleetDesiredRepository(pool *pgxpool.Pool) *FleetDesiredRepository

func (*FleetDesiredRepository) Delete added in v0.2.0

func (r *FleetDesiredRepository) Delete(ctx context.Context, tenantID, assetID, expectedPolicyID shared.ID, expectedVersion int64) error

func (*FleetDesiredRepository) Get added in v0.2.0

func (r *FleetDesiredRepository) Get(ctx context.Context, tenantID, assetID shared.ID) (*fleetdesired.State, error)

func (*FleetDesiredRepository) List added in v0.2.0

func (r *FleetDesiredRepository) List(ctx context.Context, tenantID shared.ID) ([]*fleetdesired.State, error)

func (*FleetDesiredRepository) Put added in v0.2.0

Put applies lifecycle-aware CAS. Version 1 inserts a new PolicyID only when the asset has no current policy. Later versions must retain the same PolicyID and advance exactly one version.

type FleetRolloutRepository added in v0.1.8

type FleetRolloutRepository struct {
	// contains filtered or unexported fields
}

FleetRolloutRepository persists operator update-rollout plans (migration 0065).

Durability is the point: a plan held only in memory stops offering updates the moment the control plane restarts, and says nothing about why. Every method runs through WithTenant, so Row Level Security isolates by tenant — one tenant must never be able to read, still less move, another tenant's fleet.

func NewFleetRolloutRepository added in v0.1.8

func NewFleetRolloutRepository(pool *pgxpool.Pool) *FleetRolloutRepository

NewFleetRolloutRepository constructs the Postgres rollout repository.

func (*FleetRolloutRepository) Get added in v0.1.8

func (r *FleetRolloutRepository) Get(ctx context.Context, tenantID shared.ID, channel string) (*fleetrollout.Plan, error)

Get returns the plan for a channel, or shared.ErrNotFound when none is configured.

"No row" is a legitimate resting state, not an error condition: it means no rollout is in progress and therefore no agent is offered anything.

func (*FleetRolloutRepository) Put added in v0.1.8

Put upserts the plan for (tenant, channel).

It is an upsert rather than an insert-or-update decision at the call site because a plan is a single current STATE, not a history: "the fleet is moving to 1.4.0, canary first" replaces whatever came before it. The audit log is what carries the history of who decided what.

created_at is preserved on conflict so the plan keeps the moment the rollout began, while updated_at moves with each operator action.

type IdentityStore added in v0.2.0

type IdentityStore struct {
	// contains filtered or unexported fields
}

IdentityStore persists OIDC identity/session records under the tenant RLS boundary.

func NewIdentityStore added in v0.2.0

func NewIdentityStore(pool *pgxpool.Pool) (*IdentityStore, error)

NewIdentityStore returns a PostgreSQL-backed OIDC identity/session store.

func (*IdentityStore) ConsumeAuthorizationTransaction added in v0.2.0

func (s *IdentityStore) ConsumeAuthorizationTransaction(ctx context.Context, tenantID shared.ID, stateHash string, now time.Time) (transaction identity.AuthorizationTransaction, err error)

ConsumeAuthorizationTransaction atomically deletes and returns an unexpired transaction, making the state unusable by every concurrent request and every API replica after the first succeeds.

func (*IdentityStore) CreateAuthorizationTransaction added in v0.2.0

func (s *IdentityStore) CreateAuthorizationTransaction(ctx context.Context, transaction identity.AuthorizationTransaction) error

func (*IdentityStore) CreateExternalIdentity added in v0.2.0

func (s *IdentityStore) CreateExternalIdentity(ctx context.Context, external identity.ExternalIdentity) error

func (*IdentityStore) CreateSession added in v0.2.0

func (s *IdentityStore) CreateSession(ctx context.Context, session identity.Session) error

func (*IdentityStore) GetExternalIdentity added in v0.2.0

func (s *IdentityStore) GetExternalIdentity(ctx context.Context, issuer, subject string) (external identity.ExternalIdentity, err error)

func (*IdentityStore) GetSessionByTokenHash added in v0.2.0

func (s *IdentityStore) GetSessionByTokenHash(ctx context.Context, tokenHash string) (session identity.Session, err error)

func (*IdentityStore) RevokeSession added in v0.2.0

func (s *IdentityStore) RevokeSession(ctx context.Context, tenantID, sessionID shared.ID, now time.Time) error

func (*IdentityStore) RotateSession added in v0.2.0

func (s *IdentityStore) RotateSession(ctx context.Context, previousSessionID shared.ID, replacement identity.Session, now time.Time) error

RotateSession creates replacement and revokes the active previous session in one transaction.

type ImportReceiptRepository added in v0.2.0

type ImportReceiptRepository struct {
	// contains filtered or unexported fields
}

ImportReceiptRepository persists the tenant-scoped import ledger introduced by migration 0178.

func NewImportReceiptRepository added in v0.2.0

func NewImportReceiptRepository(pool *pgxpool.Pool) *ImportReceiptRepository

NewImportReceiptRepository constructs the PostgreSQL import-receipt repository. #1181 deliberately does not wire this repository into a production writer.

func (*ImportReceiptRepository) CreateOrGet added in v0.2.0

CreateOrGet is one transaction and one idempotency decision. INSERT ... ON CONFLICT DO NOTHING lets PostgreSQL serialize concurrent retries; the read in the same transaction returns whichever row won.

func (*ImportReceiptRepository) Finalize added in v0.2.0

func (r *ImportReceiptRepository) Finalize(ctx context.Context, tenantID, receiptID shared.ID, outcome importreceipt.Outcome, counters importreceipt.Counters, updatedAt time.Time) (importreceipt.Receipt, error)

Finalize transitions a pending receipt to one terminal outcome in a tenant transaction.

func (*ImportReceiptRepository) GetByIdentity added in v0.2.0

func (r *ImportReceiptRepository) GetByIdentity(ctx context.Context, tenantID, engagementID shared.ID, sourceIdentity, digest string) (importreceipt.Receipt, error)

GetByIdentity resolves one receipt without leaking whether another tenant owns the same logical key.

type ImportedFindingRepository added in v0.1.8

type ImportedFindingRepository struct {
	// contains filtered or unexported fields
}

ImportedFindingRepository persists third-party findings (migration 0064) to PostgreSQL.

Durability is not incidental here: the ingest writes an append-only audit entry claiming that N external results entered an engagement, and that claim is only true if the rows survive a restart. Every method runs through WithTenant so Row Level Security isolates by tenant.

func NewImportedFindingRepository added in v0.1.8

func NewImportedFindingRepository(pool *pgxpool.Pool) *ImportedFindingRepository

NewImportedFindingRepository constructs the Postgres imported-finding repository.

func (*ImportedFindingRepository) ExistsDigest added in v0.1.8

func (r *ImportedFindingRepository) ExistsDigest(ctx context.Context, tenantID, engagementID shared.ID, digest string) (bool, error)

ExistsDigest reports whether this tenant's engagement already ingested a document with this digest.

func (*ImportedFindingRepository) ListByEngagement added in v0.1.8

func (r *ImportedFindingRepository) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID) ([]importedfinding.ImportedFinding, error)

ListByEngagement returns the engagement's imported findings in a deterministic order.

func (*ImportedFindingRepository) Save added in v0.1.8

Save persists a batch ATOMICALLY: one transaction, so a failure anywhere leaves no partially ingested report and no recorded digest that would make a retry look like a clean deduplicated ingest.

Each row is inserted with ON CONFLICT DO NOTHING against the idempotency index, so re-posting the same document is a no-op rather than a duplicate, and the accepted/deduplicated split the caller reports is the database's own answer rather than a guess.

type ImportedSBOMStore

type ImportedSBOMStore struct {
	// contains filtered or unexported fields
}

ImportedSBOMStore persists the active imported SBOM per engagement.

func NewImportedSBOMStore

func NewImportedSBOMStore(pool *pgxpool.Pool) *ImportedSBOMStore

NewImportedSBOMStore returns a Postgres-backed imported-SBOM store.

func (*ImportedSBOMStore) LatestByEngagement

func (s *ImportedSBOMStore) LatestByEngagement(ctx context.Context, tenantID, engagementID shared.ID) (importedsbom.Record, error)

LatestByEngagement returns the active imported SBOM for a tenant-scoped engagement.

func (*ImportedSBOMStore) MetadataByEngagements added in v0.2.0

func (s *ImportedSBOMStore) MetadataByEngagements(ctx context.Context, tenantID shared.ID, engagementIDs []shared.ID) (map[shared.ID]importedsbom.Metadata, error)

MetadataByEngagements returns the active record's metadata for each listed engagement in one query, without the raw document.

func (*ImportedSBOMStore) SaveActive

func (s *ImportedSBOMStore) SaveActive(ctx context.Context, record importedsbom.Record) error

SaveActive upserts the active imported SBOM for an engagement.

type IncidentEventRepository added in v0.2.0

type IncidentEventRepository struct {
	// contains filtered or unexported fields
}

func NewIncidentEventRepository added in v0.2.0

func NewIncidentEventRepository(pool *pgxpool.Pool) *IncidentEventRepository

func (*IncidentEventRepository) AppendEvents added in v0.2.0

func (r *IncidentEventRepository) AppendEvents(ctx context.Context, incidentID shared.ID, expectedRevision int, events []incident.IncidentEvent) error

func (*IncidentEventRepository) ListIncidentIDs added in v0.2.0

func (r *IncidentEventRepository) ListIncidentIDs(ctx context.Context, q ports.IncidentQuery) ([]shared.ID, error)

func (*IncidentEventRepository) ListMergeEdges added in v0.2.0

func (r *IncidentEventRepository) ListMergeEdges(ctx context.Context, canonicalID shared.ID) ([]incident.MergeEdge, error)
func (r *IncidentEventRepository) ListPendingResponseLinks(ctx context.Context) ([]incident.ResponseLink, error)

func (*IncidentEventRepository) LoadEvents added in v0.2.0

func (r *IncidentEventRepository) LoadEvents(ctx context.Context, incidentID shared.ID) ([]incident.IncidentEvent, error)

func (*IncidentEventRepository) ResolveCanonicalID added in v0.2.0

func (r *IncidentEventRepository) ResolveCanonicalID(ctx context.Context, id shared.ID) (shared.ID, error)

type IntegrationStore added in v0.2.0

type IntegrationStore struct {
	// contains filtered or unexported fields
}

func NewIntegrationStore added in v0.2.0

func NewIntegrationStore(pool *pgxpool.Pool, cipher *vault.Cipher) *IntegrationStore

func (*IntegrationStore) ArchiveIntegration added in v0.2.0

func (store *IntegrationStore) ArchiveIntegration(ctx context.Context, id shared.ID, expectedVersion int, audit ports.AuditEntry) error

func (*IntegrationStore) BeginIntegrationOperation added in v0.2.0

func (store *IntegrationStore) BeginIntegrationOperation(ctx context.Context, id shared.ID, startedAt time.Time) (operation integration.Operation, execute bool, err error)

func (*IntegrationStore) CancelIntegrationOperation added in v0.2.0

func (store *IntegrationStore) CancelIntegrationOperation(ctx context.Context, id shared.ID, finishedAt time.Time, audit ports.AuditEntry) (operation integration.Operation, err error)

func (*IntegrationStore) CreateIntegration added in v0.2.0

func (store *IntegrationStore) CreateIntegration(ctx context.Context, item integration.Integration, audit ports.AuditEntry) error

func (*IntegrationStore) CreateIntegrationBinding added in v0.2.0

func (store *IntegrationStore) CreateIntegrationBinding(ctx context.Context, binding integration.Binding, audit ports.AuditEntry) error

func (*IntegrationStore) DeleteIntegrationBinding added in v0.2.0

func (store *IntegrationStore) DeleteIntegrationBinding(ctx context.Context, integrationID, bindingID shared.ID, audit ports.AuditEntry) error

func (*IntegrationStore) DeleteIntegrationCredential added in v0.2.0

func (store *IntegrationStore) DeleteIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, expectedVersion, expectedConnectionRevision int, audit ports.AuditEntry) error

func (*IntegrationStore) FinishIntegrationOperation added in v0.2.0

func (store *IntegrationStore) FinishIntegrationOperation(ctx context.Context, id shared.ID, state integration.OperationState, checkpoint string, counts integration.OperationCounts, errorsIn []string, pipelines []integration.Pipeline, finishedAt time.Time) (operation integration.Operation, err error)

func (*IntegrationStore) FinishIntegrationPoll added in v0.2.0

func (store *IntegrationStore) FinishIntegrationPoll(ctx context.Context, id shared.ID, state integration.OperationState, checkpoint string, counts integration.OperationCounts, errorsIn []string, runs []integration.ExternalRun, finishedAt time.Time) (operation integration.Operation, err error)

func (*IntegrationStore) GetIntegration added in v0.2.0

func (store *IntegrationStore) GetIntegration(ctx context.Context, id shared.ID) (item integration.Integration, err error)

func (*IntegrationStore) GetIntegrationOperation added in v0.2.0

func (store *IntegrationStore) GetIntegrationOperation(ctx context.Context, id shared.ID) (operation integration.Operation, err error)

func (*IntegrationStore) IntegrationCredentialConfigured added in v0.2.0

func (store *IntegrationStore) IntegrationCredentialConfigured(ctx context.Context, integrationID shared.ID, credentialID string) (configured bool, err error)

func (*IntegrationStore) ListDueIntegrations added in v0.2.0

func (store *IntegrationStore) ListDueIntegrations(ctx context.Context, now time.Time, limit int) (items []integration.Integration, err error)

func (*IntegrationStore) ListIntegrationBindings added in v0.2.0

func (store *IntegrationStore) ListIntegrationBindings(ctx context.Context, integrationID shared.ID) (bindings []integration.Binding, err error)

func (*IntegrationStore) ListIntegrationExternalRuns added in v0.2.0

func (store *IntegrationStore) ListIntegrationExternalRuns(ctx context.Context, integrationID shared.ID, limit int) (runs []integration.ExternalRun, err error)

func (*IntegrationStore) ListIntegrationOperations added in v0.2.0

func (store *IntegrationStore) ListIntegrationOperations(ctx context.Context, integrationID shared.ID, limit int) (operations []integration.Operation, err error)

func (*IntegrationStore) ListIntegrations added in v0.2.0

func (store *IntegrationStore) ListIntegrations(ctx context.Context, includeArchived bool) (items []integration.Integration, err error)

func (*IntegrationStore) PutIntegrationCredential added in v0.2.0

func (store *IntegrationStore) PutIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, plaintext []byte, expectedVersion, expectedConnectionRevision int, audit ports.AuditEntry) error

func (*IntegrationStore) ResolveIntegrationCredential added in v0.2.0

func (store *IntegrationStore) ResolveIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, expectedRevision int) (plaintext []byte, err error)

func (*IntegrationStore) SetIntegrationEnabled added in v0.2.0

func (store *IntegrationStore) SetIntegrationEnabled(ctx context.Context, id shared.ID, enabled bool, expectedVersion int, audit ports.AuditEntry) (updated integration.Integration, err error)

func (*IntegrationStore) StartIntegrationOperation added in v0.2.0

func (store *IntegrationStore) StartIntegrationOperation(ctx context.Context, operation integration.Operation, jobKind string, payload []byte, audit ports.AuditEntry) (integration.Operation, error)

func (*IntegrationStore) UpdateIntegration added in v0.2.0

func (store *IntegrationStore) UpdateIntegration(ctx context.Context, item integration.Integration, expectedVersion int, audit ports.AuditEntry) (updated integration.Integration, err error)

type JobQueue

type JobQueue struct {
	// contains filtered or unexported fields
}

JobQueue is the durable PostgreSQL queue with tenant-bound at-least-once delivery.

func NewJobQueue

func NewJobQueue(pool *pgxpool.Pool, ids ports.IDGenerator) *JobQueue

func (*JobQueue) AggregateJobQueueStats added in v0.2.0

func (q *JobQueue) AggregateJobQueueStats(ctx context.Context, kinds ...string) (ports.JobStats, error)

AggregateJobQueueStats aggregates each tenant's RLS-scoped queue statistics for the operator metrics collector. The tenant enumeration mirrors Claim so the forced-RLS per-tenant transaction remains the only way any job row is read; no tenant label is ever attached to the aggregated totals.

func (*JobQueue) Claim

func (q *JobQueue) Claim(ctx context.Context, visibility time.Duration, kinds ...string) (*ports.QueuedJob, error)

func (*JobQueue) Complete

func (q *JobQueue) Complete(ctx context.Context, id string, fence int64) error

func (*JobQueue) Deadletter

func (q *JobQueue) Deadletter(ctx context.Context, id string, fence int64) error

func (*JobQueue) Depth

func (q *JobQueue) Depth(ctx context.Context, kinds ...string) (count int, err error)

func (*JobQueue) Enqueue

func (q *JobQueue) Enqueue(ctx context.Context, kind string, payload []byte) (string, error)

func (*JobQueue) Fail

func (q *JobQueue) Fail(ctx context.Context, id string, fence int64, retryIn time.Duration) error

func (*JobQueue) Heartbeat

func (q *JobQueue) Heartbeat(ctx context.Context, id string, fence int64, extend time.Duration) error

func (*JobQueue) JobStatus added in v0.1.8

func (q *JobQueue) JobStatus(ctx context.Context, id string) (status ports.JobStatus, err error)

func (*JobQueue) Retry added in v0.2.0

func (q *JobQueue) Retry(ctx context.Context, id string, fence int64, retryIn time.Duration) error

Retry releases a contention delivery without charging it against the job's attempt budget.

func (*JobQueue) Stats added in v0.1.8

func (q *JobQueue) Stats(ctx context.Context, kinds ...string) (stats ports.JobStats, err error)

type JudgmentRepository

type JudgmentRepository struct {
	// contains filtered or unexported fields
}

JudgmentRepository persists AI judgments to PostgreSQL, engagement-scoped. All operations route through WithContextTenant so tenant isolation is enforced at the database level. Save validates that the engagement belongs to the context tenant before writing.

func NewJudgmentRepository

func NewJudgmentRepository(pool *pgxpool.Pool) *JudgmentRepository

NewJudgmentRepository returns a repository backed by the given pool.

func (*JudgmentRepository) AcknowledgeJudgmentAudit added in v0.1.8

func (r *JudgmentRepository) AcknowledgeJudgmentAudit(ctx context.Context, kind ports.JudgmentAuditKind, judgmentID shared.ID, version int) error

func (*JudgmentRepository) GetByID added in v0.2.0

func (r *JudgmentRepository) GetByID(ctx context.Context, engagementID, id shared.ID) (out judgment.Judgment, err error)

GetByID provides a bounded engagement-scoped lookup for migration/backfill consumers without loading every judgment into process memory.

func (*JudgmentRepository) ListByEngagement

func (r *JudgmentRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]judgment.Judgment, error)

ListByEngagement returns the engagement's judgments, oldest first (deterministic order). RLS scopes the query to the context tenant.

func (*JudgmentRepository) ListBySubject

func (r *JudgmentRepository) ListBySubject(ctx context.Context, engagementID, subjectID shared.ID) ([]judgment.Judgment, error)

ListBySubject returns the engagement's judgments about a given subject id, oldest first. RLS scopes the query to the context tenant.

func (*JudgmentRepository) ListPendingJudgmentAudits added in v0.1.8

func (r *JudgmentRepository) ListPendingJudgmentAudits(ctx context.Context, engagementID shared.ID) (out []ports.PendingJudgmentAudit, err error)

func (*JudgmentRepository) Save

Save inserts a proposed judgment (idempotent by id; never clobbers an existing row – score/state move only via SetScoreState). The typed claim is stored as its fail-closed discriminated envelope (JSONB). The tenant_id is resolved from context, and the engagement is validated to belong to that tenant before the insert.

func (*JudgmentRepository) SaveWithProposalAudit added in v0.1.8

func (r *JudgmentRepository) SaveWithProposalAudit(ctx context.Context, j judgment.Judgment, entry ports.AuditEntry) error

SaveWithProposalAudit persists a proposal with its immutable pending audit entry.

func (*JudgmentRepository) SetScoreState

func (r *JudgmentRepository) SetScoreState(ctx context.Context, engagementID, id shared.ID, score int, state judgment.State, expectedVersion int) (judgment.Judgment, error)

SetScoreState moves a judgment's evidence score + state under optimistic concurrency (the verify/accept path): the row updates only if version matches expectedVersion, then version is bumped. This is the ONLY path that moves a stored judgment's score/state, and it is deliberately off the broad ports.JudgmentStore interface (a read-only consumer cannot reach it). On a miss it distinguishes ErrConflict (exists, version moved) from ErrNotFound. RLS scopes the operation to the context tenant.

func (*JudgmentRepository) SetVerdictState added in v0.1.8

func (r *JudgmentRepository) SetVerdictState(ctx context.Context, engagementID, id shared.ID, score int, state judgment.State, verifiedBy, verdictRationale string, expectedVersion int) (judgment.Judgment, error)

SetVerdictState persists a verdict's sealed verifier and rationale with its score transition.

func (*JudgmentRepository) SetVerdictStateWithAudit added in v0.1.8

func (r *JudgmentRepository) SetVerdictStateWithAudit(ctx context.Context, engagementID, id shared.ID, score int, state judgment.State, verifiedBy, rationale string, expectedVersion int, entry ports.AuditEntry) (out judgment.Judgment, err error)

SetVerdictStateWithAudit commits the verdict and immutable pending audit entry atomically.

type LeaderStore added in v0.1.8

type LeaderStore struct {
	// contains filtered or unexported fields
}

LeaderStore is the Postgres fenced-lease implementation of leader election. It is global control-plane state (no tenant scoping, no RLS; see migration 0061).

func NewLeaderStore added in v0.1.8

func NewLeaderStore(pool *pgxpool.Pool) *LeaderStore

NewLeaderStore constructs the Postgres leader store.

func (*LeaderStore) Acquire added in v0.1.8

func (s *LeaderStore) Acquire(ctx context.Context, resource, holder string, term time.Duration, now time.Time) (bool, int64, error)

Acquire atomically takes or renews leadership in a single upsert: it takes the lease when the row is absent, already held by holder (renewal), or expired (takeover, which bumps the fence), and otherwise leaves a live foreign lease untouched. held is true iff holder owns it afterwards.

func (*LeaderStore) Resign added in v0.1.8

func (s *LeaderStore) Resign(ctx context.Context, resource, holder string, now time.Time) error

Resign releases the lease if held by holder. It EXPIRES the lease (clears the holder and sets the term to now) rather than deleting the row, so the fence survives: a graceful handover keeps the fence monotonic (the next acquirer takes over an expired row and bumps it), which a fresh INSERT with fence=1 would not. A challenger sees the expired row immediately and can take over without waiting out the term.

type LeaseRunLock

type LeaseRunLock struct {
	// contains filtered or unexported fields
}

LeaseRunLock implements ports.RunLocker via a jobs_run_lock ROW lease instead of a session advisory lock. Unlike RunLock it does NOT hold a pooled connection for the duration of the run – each operation borrows a connection transiently – so N concurrent runs cannot starve the pool (the ≤8-default-pool hazard the hostile review flagged). A background renewer extends the lease while the run is live; a crash lets the lease expire so another worker can reclaim it. Used for RECON; the agent SESSION lock stays the advisory RunLock (it must not expire mid-LLM-loop). owner is a process id; each acquisition derives a distinct owner token so concurrent loops cannot renew or release one another's row.

func NewLeaseRunLock

func NewLeaseRunLock(pool *pgxpool.Pool, owner string, lease time.Duration) *LeaseRunLock

NewLeaseRunLock returns a row-lease run locker. lease is the claim TTL (set it comfortably above the longest run, e.g. ReconTimeout + a minute); the renewer ticks at lease/4 so several renews fall inside one TTL.

func (*LeaseRunLock) TryLock

func (l *LeaseRunLock) TryLock(ctx context.Context, runID string) (func(), bool, error)

TryLock claims the lease (see TryLockLeased) and discards the lease-loss context – for callers that don't observe lease loss (e.g. the stale-run sweeper's liveness probe).

func (*LeaseRunLock) TryLockLeased

func (l *LeaseRunLock) TryLockLeased(ctx context.Context, runID string) (context.Context, func(), bool, error)

TryLockLeased claims the lease for runID if it is free or expired. On success it starts a renewer and returns: a leaseCtx cancelled when the lease is LOST (so the caller aborts the in-flight run), plus a release that stops the renewer, cancels leaseCtx, and deletes the owner's row. A claim held by a live owner returns ok=false (the at-least-once queue retries).

type LegalHoldRepository added in v0.2.0

type LegalHoldRepository struct {
	// contains filtered or unexported fields
}

LegalHoldRepository is the Postgres tier for legal holds (#635). Append-only history per (tenant, engagement); at most one active hold per engagement (partial unique index). Every method runs under the ctx tenant via WithContextTenant (RLS) with an explicit tenant_id predicate as defense-in-depth.

func NewLegalHoldRepository added in v0.2.0

func NewLegalHoldRepository(pool *pgxpool.Pool) *LegalHoldRepository

func (*LegalHoldRepository) IsHeld added in v0.2.0

func (r *LegalHoldRepository) IsHeld(ctx context.Context, engagementID shared.ID) (bool, error)

func (*LegalHoldRepository) ListActive added in v0.2.0

func (r *LegalHoldRepository) ListActive(ctx context.Context) ([]legalhold.Hold, error)

func (*LegalHoldRepository) Place added in v0.2.0

func (*LegalHoldRepository) Release added in v0.2.0

func (r *LegalHoldRepository) Release(ctx context.Context, engagementID shared.ID, releasedBy string, at time.Time) error

type MaterializingAdvisoryWriter added in v0.2.0

type MaterializingAdvisoryWriter struct {
	// contains filtered or unexported fields
}

MaterializingAdvisoryWriter adapts the bulk advisory ingester (ports.AdvisoryWriter, one advisory at a time) onto the observation materializer. Each ingested advisory is recorded as one observation of a named bulk source and re-materialized through advisory.Merge, so a CVE present in this feed AND in another source (the online NVD provider, a different feed dump) yields the UNION of their affected ranges instead of a flat last-writer-wins overwrite of advisories.data. This closes EPIC #860 D1.2: the CLI sync-advisories path is now a single provenance-tracked writer over the observation/canonical model, which the online sync already uses, rather than a second clobbering writer against the same table.

A re-ingest of the same feed supersedes that source's prior observation for the same record (the materializer keys on (source_id, record_id, content_hash) and marks the old content not-current), so a feed that narrows an advisory still takes effect for its own source while other sources' ranges remain.

Source identity is PER AUTHORITY (one osv, one csaf, one oval bulk source), which makes the union work across DIFFERENT authorities (an OSV dump and the NVD provider are separate sources, so their ranges are unioned) while a re-sync of the SAME authority supersedes rather than accumulating stale ranges. Two bounds follow and are intentional: two separate dumps of the SAME authority that disagree on one advisory resolve last-writer-wins (the newer snapshot is authoritative, not a union of stale + fresh); and this CLI osv bulk source is distinct from the ONLINE osv provider, so a deployment should feed OSV through one path (an offline dump OR the online provider), and a stale, un-re-synced offline dump can keep a range the online path later retracted until the dump is re-synced (the corpus-freshness warning surfaces staleness).

func NewMaterializingAdvisoryWriter added in v0.2.0

func NewMaterializingAdvisoryWriter(ctx context.Context, pool *pgxpool.Pool, sourceKey, displayName, adapterType string, onSkip func(id string, err error)) (*MaterializingAdvisoryWriter, error)

NewMaterializingAdvisoryWriter ensures a vulnerability_sources row for the given bulk feed exists (a disabled, full-sync source, since an offline dump is not scheduled) and returns a writer that materializes each advisory as one of that source's observations. sourceKey is the stable, unique feed identity (used as both the primary id and the source_key); adapterType is one of osv/csaf/oval. onSkip, if non-nil, is called for each advisory dropped on a per-record data error instead of aborting the run.

func (*MaterializingAdvisoryWriter) Upsert added in v0.2.0

Upsert records the advisory as one current observation of this source and re-materializes the canonical advisory, merging with every other source's current observation for the same identity. A withdrawn feed advisory is recorded with StatusWithdrawn so the canonical projection keeps it out of the matcher. No SyncRunID is set: an offline bulk ingest is not a scheduled, tenant-scoped provenance run, so the materializer's sync-run/tenant provenance checks are skipped.

type NotificationRepository added in v0.2.0

type NotificationRepository struct {
	// contains filtered or unexported fields
}

func NewNotificationRepository added in v0.2.0

func NewNotificationRepository(pool *pgxpool.Pool) *NotificationRepository

func (*NotificationRepository) BeginAttempt added in v0.2.0

func (r *NotificationRepository) BeginAttempt(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, aid shared.ID, at time.Time) (notification.Attempt, error)

func (*NotificationRepository) CancelDelivery added in v0.2.0

func (r *NotificationRepository) CancelDelivery(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, reason string) error

func (*NotificationRepository) CreateChannel added in v0.2.0

func (*NotificationRepository) CreateRule added in v0.2.0

func (*NotificationRepository) DeadLetterDelivery added in v0.2.0

func (r *NotificationRepository) DeadLetterDelivery(ctx context.Context, tenant, did shared.ID, reason string) error

func (*NotificationRepository) DeleteChannel added in v0.2.0

func (r *NotificationRepository) DeleteChannel(ctx context.Context, tenant, id shared.ID, revision int, at time.Time) error

func (*NotificationRepository) DeleteRule added in v0.2.0

func (r *NotificationRepository) DeleteRule(ctx context.Context, tenant, id shared.ID, revision int) error

func (*NotificationRepository) DeliveryStillRelevant added in v0.2.0

func (r *NotificationRepository) DeliveryStillRelevant(ctx context.Context, work ports.NotificationWork) (bool, error)

func (*NotificationRepository) FinishAttempt added in v0.2.0

func (r *NotificationRepository) FinishAttempt(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, aid shared.ID, at time.Time, outcome string, status int, errorCode string, next *time.Time) error

func (*NotificationRepository) GetChannel added in v0.2.0

func (r *NotificationRepository) GetChannel(ctx context.Context, tenant, id shared.ID) (notification.Channel, error)

func (*NotificationRepository) GetDelivery added in v0.2.0

func (r *NotificationRepository) GetDelivery(ctx context.Context, tenant, id shared.ID) (notification.Delivery, error)

func (*NotificationRepository) GetRule added in v0.2.0

func (r *NotificationRepository) GetRule(ctx context.Context, tenant, id shared.ID) (notification.Rule, error)

func (*NotificationRepository) ListAttempts added in v0.2.0

func (r *NotificationRepository) ListAttempts(ctx context.Context, tenant, did shared.ID) ([]notification.Attempt, error)

func (*NotificationRepository) ListChannels added in v0.2.0

func (r *NotificationRepository) ListChannels(ctx context.Context, tenant shared.ID) ([]notification.Channel, error)

func (*NotificationRepository) ListDeliveries added in v0.2.0

func (*NotificationRepository) ListRules added in v0.2.0

func (r *NotificationRepository) ListRules(ctx context.Context, tenant shared.ID) ([]notification.Rule, error)

func (*NotificationRepository) LoadWork added in v0.2.0

func (r *NotificationRepository) LoadWork(ctx context.Context, tenant, did shared.ID) (ports.NotificationWork, error)

func (*NotificationRepository) Publish added in v0.2.0

func (*NotificationRepository) PublishToChannel added in v0.2.0

func (r *NotificationRepository) PublishToChannel(ctx context.Context, e notification.Event, cid shared.ID) (shared.ID, error)

func (*NotificationRepository) UpdateChannel added in v0.2.0

func (r *NotificationRepository) UpdateChannel(ctx context.Context, c notification.Channel, sealed string, replace bool) (notification.Channel, error)

func (*NotificationRepository) UpdateRule added in v0.2.0

type NotificationSource added in v0.2.0

type NotificationSource struct {
	// contains filtered or unexported fields
}

NotificationSource projects durable source records into generic notification events. Every projection and its delivery jobs commit in one tenant transaction.

func NewNotificationSource added in v0.2.0

func NewNotificationSource(pool *pgxpool.Pool, repo *NotificationRepository, fleetStaleAfter time.Duration, incidentEnabled bool) *NotificationSource

func (*NotificationSource) Poll added in v0.2.0

func (s *NotificationSource) Poll(ctx context.Context, now time.Time, limit int) (int, error)

func (*NotificationSource) SetVulnerabilityEnabled added in v0.2.0

func (s *NotificationSource) SetVulnerabilityEnabled(enabled bool)

type OwnershipExecution added in v0.2.0

type OwnershipExecution struct {
	// contains filtered or unexported fields
}

func NewOwnershipExecution added in v0.2.0

func NewOwnershipExecution(repo *OwnershipRepository, ids ports.IDGenerator, clock ports.Clock) (*OwnershipExecution, error)

func (*OwnershipExecution) CommitOwnershipWork added in v0.2.0

func (r *OwnershipExecution) CommitOwnershipWork(ctx context.Context, job ports.QueuedJob, work ports.OwnershipWork, result ownership.Result, outcome string, notify bool) error

func (*OwnershipExecution) DispatchOwnership added in v0.2.0

func (r *OwnershipExecution) DispatchOwnership(ctx context.Context, mode string, limit int) (int, error)

DispatchOwnership enumerates tenant identities only outside RLS, then enters a separate tenant transaction. One job contains at most limit frozen inputs.

func (*OwnershipExecution) FailOwnershipWork added in v0.2.0

func (r *OwnershipExecution) FailOwnershipWork(ctx context.Context, job ports.QueuedJob) error

func (*OwnershipExecution) LoadOwnershipWork added in v0.2.0

func (r *OwnershipExecution) LoadOwnershipWork(ctx context.Context, job ports.QueuedJob, limit int) (out []ports.OwnershipWork, err error)

func (*OwnershipExecution) ReplayOwnershipRun added in v0.2.0

func (r *OwnershipExecution) ReplayOwnershipRun(ctx context.Context, actor, id shared.ID, revision int) error

func (*OwnershipExecution) StartOwnershipRun added in v0.2.0

func (r *OwnershipExecution) StartOwnershipRun(ctx context.Context, req ports.OwnershipRunRequest) (run ports.OwnershipRun, err error)

type OwnershipRepository added in v0.2.0

type OwnershipRepository struct {
	// contains filtered or unexported fields
}

func NewOwnershipRepository added in v0.2.0

func NewOwnershipRepository(pool *pgxpool.Pool) (*OwnershipRepository, error)

func (*OwnershipRepository) ActivatePolicy added in v0.2.0

func (*OwnershipRepository) AddMember added in v0.2.0

func (r *OwnershipRepository) AddMember(ctx context.Context, team, id shared.ID, at time.Time) error

func (*OwnershipRepository) AppendIntent added in v0.2.0

func (*OwnershipRepository) ApplyAssignment added in v0.2.0

func (r *OwnershipRepository) ApplyAssignment(ctx context.Context, m ports.OwnershipMutation) (out ownership.Decision, err error)

func (*OwnershipRepository) AuthorizeOwnership added in v0.2.0

func (r *OwnershipRepository) AuthorizeOwnership(ctx context.Context, actor shared.ID, permission user.Permission) error

func (*OwnershipRepository) CompleteIntent added in v0.2.0

func (r *OwnershipRepository) CompleteIntent(ctx context.Context, id shared.ID) error

func (*OwnershipRepository) CreatePolicyVersion added in v0.2.0

func (r *OwnershipRepository) CreatePolicyVersion(ctx context.Context, p ownership.PolicyVersion) error

func (*OwnershipRepository) CreateRun added in v0.2.0

func (r *OwnershipRepository) CreateRun(ctx context.Context, run ports.OwnershipRun) error

func (*OwnershipRepository) CreateSnapshot added in v0.2.0

func (r *OwnershipRepository) CreateSnapshot(ctx context.Context, s ownership.Snapshot) error

func (*OwnershipRepository) CreateTeam added in v0.2.0

func (*OwnershipRepository) DeleteAssetMapping added in v0.2.0

func (r *OwnershipRepository) DeleteAssetMapping(ctx context.Context, asset shared.ID, revision int) error

func (*OwnershipRepository) DeleteMapping added in v0.2.0

func (r *OwnershipRepository) DeleteMapping(ctx context.Context, eng shared.ID, repo, owner string, revision int) error

func (*OwnershipRepository) GetActivePolicy added in v0.2.0

func (r *OwnershipRepository) GetActivePolicy(ctx context.Context, eng shared.ID, repo string) (out ports.OwnershipPolicy, err error)

func (*OwnershipRepository) GetAssignment added in v0.2.0

func (r *OwnershipRepository) GetAssignment(ctx context.Context, eng, finding shared.ID) (out ports.OwnershipCurrent, err error)

func (*OwnershipRepository) GetOwnershipPolicy added in v0.2.0

func (r *OwnershipRepository) GetOwnershipPolicy(ctx context.Context, id shared.ID) (out ports.OwnershipPolicyHeader, err error)

func (*OwnershipRepository) GetPolicyVersion added in v0.2.0

func (r *OwnershipRepository) GetPolicyVersion(ctx context.Context, id shared.ID, version int) (p ownership.PolicyVersion, err error)

func (*OwnershipRepository) GetRun added in v0.2.0

func (r *OwnershipRepository) GetRun(ctx context.Context, id shared.ID) (run ports.OwnershipRun, err error)

func (*OwnershipRepository) GetSnapshot added in v0.2.0

func (r *OwnershipRepository) GetSnapshot(ctx context.Context, id shared.ID) (s ownership.Snapshot, err error)

func (*OwnershipRepository) GetTeam added in v0.2.0

func (r *OwnershipRepository) GetTeam(ctx context.Context, id shared.ID) (t ownership.Team, err error)

func (*OwnershipRepository) ListAssetMappings added in v0.2.0

func (r *OwnershipRepository) ListAssetMappings(ctx context.Context, after shared.ID, limit int) (out []ports.OwnershipAssetMapping, err error)

func (*OwnershipRepository) ListDecisions added in v0.2.0

func (r *OwnershipRepository) ListDecisions(ctx context.Context, eng, finding shared.ID, cursor ports.OwnershipHistoryCursor) (out []ownership.Decision, err error)

func (*OwnershipRepository) ListMappings added in v0.2.0

func (r *OwnershipRepository) ListMappings(ctx context.Context, eng shared.ID, repo, after string, limit int) (out []ports.OwnershipMapping, err error)

func (*OwnershipRepository) ListMembers added in v0.2.0

func (r *OwnershipRepository) ListMembers(ctx context.Context, team, after shared.ID, limit int) (out []ownership.Membership, err error)

func (*OwnershipRepository) ListOwnershipPolicies added in v0.2.0

func (r *OwnershipRepository) ListOwnershipPolicies(ctx context.Context, eng, after shared.ID, limit int) (out []ports.OwnershipPolicyHeader, err error)

func (*OwnershipRepository) ListOwnershipSnapshots added in v0.2.0

func (r *OwnershipRepository) ListOwnershipSnapshots(ctx context.Context, eng, after shared.ID, limit int) (out []ownership.Snapshot, err error)

func (*OwnershipRepository) ListPendingIntents added in v0.2.0

func (r *OwnershipRepository) ListPendingIntents(ctx context.Context, kind string, limit int) (out []ports.OwnershipIntent, err error)

func (*OwnershipRepository) ListRunItems added in v0.2.0

func (r *OwnershipRepository) ListRunItems(ctx context.Context, id, after shared.ID, limit int) (out []ports.OwnershipRunItem, err error)

func (*OwnershipRepository) ListTeams added in v0.2.0

func (r *OwnershipRepository) ListTeams(ctx context.Context, after shared.ID, limit int) (out []ownership.Team, err error)

func (*OwnershipRepository) MarkOwnershipSourceReady added in v0.2.0

func (r *OwnershipRepository) MarkOwnershipSourceReady(ctx context.Context, eng, id shared.ID) error

func (*OwnershipRepository) OwnershipInbox added in v0.2.0

func (*OwnershipRepository) RemoveMember added in v0.2.0

func (r *OwnershipRepository) RemoveMember(ctx context.Context, team, id shared.ID) error

func (*OwnershipRepository) ReserveOwnershipBulk added in v0.2.0

func (r *OwnershipRepository) ReserveOwnershipBulk(ctx context.Context, actor shared.ID, key, hash string, at time.Time) error

func (*OwnershipRepository) SaveAssetMapping added in v0.2.0

func (*OwnershipRepository) SaveMapping added in v0.2.0

func (r *OwnershipRepository) SaveMapping(ctx context.Context, eng shared.ID, record ports.OwnershipMapping) error

func (*OwnershipRepository) SaveOwnershipSource added in v0.2.0

func (r *OwnershipRepository) SaveOwnershipSource(ctx context.Context, source ports.OwnershipSourceRecord) error

func (*OwnershipRepository) SaveRunItems added in v0.2.0

func (r *OwnershipRepository) SaveRunItems(ctx context.Context, id shared.ID, revision int, items []ports.OwnershipRunItem) error

func (*OwnershipRepository) SetRunState added in v0.2.0

func (r *OwnershipRepository) SetRunState(ctx context.Context, id shared.ID, revision int, state string) error

func (*OwnershipRepository) UpdateTeam added in v0.2.0

func (r *OwnershipRepository) UpdateTeam(ctx context.Context, t ownership.Team, expected int) (ownership.Team, error)

func (*OwnershipRepository) VisibleOwnershipEngagement added in v0.2.0

func (r *OwnershipRepository) VisibleOwnershipEngagement(ctx context.Context, id shared.ID) error

type PoolConfig

type PoolConfig struct {
	MaxConns          int32
	MinConns          int32
	MaxConnLifetime   time.Duration
	MaxConnIdleTime   time.Duration
	HealthCheckPeriod time.Duration
}

PoolConfig sizes the pgx connection pool. Zero values get sane defaults. Sizing the pool explicitly (the default pgx cap is max(4, NumCPU) ≈ 8) is required now that the durable agent path holds a connection-bearing advisory lock per active run – an unsized pool would starve HTTP handlers at low-tens concurrency.

type PoolStatsSource added in v0.2.0

type PoolStatsSource struct {
	// contains filtered or unexported fields
}

PoolStatsSource exposes aggregate pool saturation through a driver-free port so the metrics adapter never depends on pgx.

func NewPoolStatsSource added in v0.2.0

func NewPoolStatsSource(pool *pgxpool.Pool) *PoolStatsSource

NewPoolStatsSource returns nil for a nil pool so a composition root can pass the result straight to an optional collector without constructing a typed-nil interface.

func (*PoolStatsSource) PoolStats added in v0.2.0

func (s *PoolStatsSource) PoolStats() ports.PoolStats

PoolStats snapshots the pool's aggregate counters; it carries no connection identity.

type PrivacyPolicyRepository added in v0.2.0

type PrivacyPolicyRepository struct {
	*FleetAuditRepository
	// contains filtered or unexported fields
}

PrivacyPolicyRepository persists immutable source-policy versions and the separately mutable active pointer under tenant RLS.

func NewPrivacyPolicyRepository added in v0.2.0

func NewPrivacyPolicyRepository(pool *pgxpool.Pool) (*PrivacyPolicyRepository, error)

func (*PrivacyPolicyRepository) ActivatePrivacyPolicy added in v0.2.0

func (r *PrivacyPolicyRepository) ActivatePrivacyPolicy(
	ctx context.Context,
	activation privacy.Activation,
) (privacy.Activation, error)

func (*PrivacyPolicyRepository) ActivatePrivacyPolicyWithAudit added in v0.2.0

func (r *PrivacyPolicyRepository) ActivatePrivacyPolicyWithAudit(
	ctx context.Context,
	activation privacy.Activation,
	intent ports.FleetAuditIntent,
) (privacy.Activation, ports.FleetAuditIntent, error)

func (*PrivacyPolicyRepository) ActivePrivacyPolicy added in v0.2.0

func (r *PrivacyPolicyRepository) ActivePrivacyPolicy(
	ctx context.Context,
	tenantID shared.ID,
) (privacy.Assignment, error)

func (*PrivacyPolicyRepository) PrivacyPolicyActivationHistory added in v0.2.0

func (r *PrivacyPolicyRepository) PrivacyPolicyActivationHistory(
	ctx context.Context,
	tenantID shared.ID,
) ([]privacy.Activation, error)

func (*PrivacyPolicyRepository) PrivacyPolicyByDigest added in v0.2.0

func (r *PrivacyPolicyRepository) PrivacyPolicyByDigest(
	ctx context.Context,
	tenantID shared.ID,
	digest string,
) (privacy.Assignment, error)

func (*PrivacyPolicyRepository) PrivacyPolicyHistory added in v0.2.0

func (r *PrivacyPolicyRepository) PrivacyPolicyHistory(
	ctx context.Context,
	tenantID shared.ID,
) ([]privacy.Assignment, error)

func (*PrivacyPolicyRepository) PutPrivacyPolicy added in v0.2.0

func (r *PrivacyPolicyRepository) PutPrivacyPolicy(
	ctx context.Context,
	assignment privacy.Assignment,
) (bool, error)

type ProjectAnalysisStore

type ProjectAnalysisStore struct {
	// contains filtered or unexported fields
}

ProjectAnalysisStore persists immutable Project analysis snapshots and, through the hotspot and issue files in this package, their projections. project_analyses and every projection table are RLS-protected (migration 0129), so each statement runs inside requireTenant and keeps its own tenant_id predicate.

The pre-0129 reads branched to a tenant-free WHERE whenever the caller passed a zero tenant, which widened the query instead of denying it. Those branches are gone: a missing tenant is a validation error.

func NewProjectAnalysisStore

func NewProjectAnalysisStore(pool *pgxpool.Pool) *ProjectAnalysisStore

func (*ProjectAnalysisStore) AttachSourceWithAudit added in v0.1.8

func (r *ProjectAnalysisStore) AttachSourceWithAudit(ctx context.Context, tenantID, projectID, analysisID shared.ID, capture projectanalysis.SourceCapture, audit ports.AuditEntry) error

func (*ProjectAnalysisStore) Branches added in v0.2.0

func (r *ProjectAnalysisStore) Branches(ctx context.Context, tenantID, projectID shared.ID) ([]string, error)

func (*ProjectAnalysisStore) CurrentAnalysisHotspotSummary

func (r *ProjectAnalysisStore) CurrentAnalysisHotspotSummary(ctx context.Context, tenantID, projectID, analysisID shared.ID, lens hotspot.Lens) (hotspot.Summary, error)

func (*ProjectAnalysisStore) CurrentFindingStatuses added in v0.1.8

func (r *ProjectAnalysisStore) CurrentFindingStatuses(ctx context.Context, tenantID, projectID shared.ID, keys []string) (map[string]string, error)

func (*ProjectAnalysisStore) Get

func (r *ProjectAnalysisStore) Get(ctx context.Context, tenantID, projectID, analysisID shared.ID) (projectanalysis.Analysis, error)

func (*ProjectAnalysisStore) GetHotspot

func (r *ProjectAnalysisStore) GetHotspot(ctx context.Context, tenantID, projectID, hotspotID shared.ID) (hotspot.Hotspot, error)

func (*ProjectAnalysisStore) GetIssue

func (r *ProjectAnalysisStore) GetIssue(ctx context.Context, tenantID, projectID, issueID shared.ID) (issue.Issue, error)

func (*ProjectAnalysisStore) HotspotHistory

func (r *ProjectAnalysisStore) HotspotHistory(ctx context.Context, tenantID, projectID, hotspotID shared.ID) ([]hotspot.ReviewEvent, error)

func (*ProjectAnalysisStore) IssueHistory

func (r *ProjectAnalysisStore) IssueHistory(ctx context.Context, tenantID, projectID, issueID shared.ID) ([]issue.ReviewEvent, error)

func (*ProjectAnalysisStore) LatestForProjects

func (r *ProjectAnalysisStore) LatestForProjects(ctx context.Context, tenantID shared.ID, projectIDs []shared.ID) (map[shared.ID]projectanalysis.Analysis, error)

func (*ProjectAnalysisStore) LatestWithResult

func (r *ProjectAnalysisStore) LatestWithResult(ctx context.Context, tenantID, projectID shared.ID, branch string) (projectanalysis.Analysis, []byte, error)

func (*ProjectAnalysisStore) List

func (r *ProjectAnalysisStore) List(ctx context.Context, tenantID, projectID shared.ID, branch string, limit int, beforeCreatedAt time.Time, beforeID shared.ID) ([]projectanalysis.Analysis, bool, error)

func (*ProjectAnalysisStore) ListAnalysisHotspots

func (r *ProjectAnalysisStore) ListAnalysisHotspots(ctx context.Context, tenantID, projectID, analysisID shared.ID, lens hotspot.Lens, filter hotspot.ListFilter) (hotspot.Page, hotspot.Summary, error)

func (*ProjectAnalysisStore) ListHotspots

func (r *ProjectAnalysisStore) ListHotspots(ctx context.Context, tenantID, projectID shared.ID, filter hotspot.ListFilter) (hotspot.Page, error)

func (*ProjectAnalysisStore) ListIssues

func (r *ProjectAnalysisStore) ListIssues(ctx context.Context, tenantID, projectID shared.ID, filter issue.ListFilter) (issue.Page, error)

func (*ProjectAnalysisStore) MatchIntegrationAnalysis added in v0.2.0

func (r *ProjectAnalysisStore) MatchIntegrationAnalysis(ctx context.Context, projectID shared.ID, revision string) (analysisID shared.ID, state integration.CorrelationState, err error)

func (*ProjectAnalysisStore) PruneBranchAnalyses added in v0.2.0

func (r *ProjectAnalysisStore) PruneBranchAnalyses(ctx context.Context, tenantID, projectID shared.ID, branch string, keep int) (int, error)

Branches returns the distinct branch values recorded for the project, sorted. PruneBranchAnalyses deletes all but the newest keep analyses on one branch, tenant-scoped. It uses the (tenant_id, project_id, branch, created_at DESC, id DESC) index from migration 0136, and orders id with COLLATE "C" to match that index and List's tie-break so exactly the newest keep rows survive.

func (*ProjectAnalysisStore) ResolvedIssueKeys

func (r *ProjectAnalysisStore) ResolvedIssueKeys(ctx context.Context, tenantID, projectID shared.ID) (map[string]bool, error)

func (*ProjectAnalysisStore) Save

func (*ProjectAnalysisStore) SaveWithResult

func (r *ProjectAnalysisStore) SaveWithResult(ctx context.Context, analysis projectanalysis.Analysis, result []byte) error

func (*ProjectAnalysisStore) SaveWithResultAndHotspots

func (r *ProjectAnalysisStore) SaveWithResultAndHotspots(ctx context.Context, analysis projectanalysis.Analysis, result []byte, candidates []hotspot.Candidate) error

SaveWithResultAndHotspots commits the immutable analysis and its Security Hotspot projection in one PostgreSQL transaction. It delegates to SaveWithResultAndProjections with no issue projection, so both write paths share the same single-transaction body: a projection write failure rolls the analysis back, and the scan worker cannot publish a successful analysis without its projections.

func (*ProjectAnalysisStore) SaveWithResultAndProjections

func (r *ProjectAnalysisStore) SaveWithResultAndProjections(ctx context.Context, analysis projectanalysis.Analysis, result []byte, hotspots []hotspot.Candidate, issues []issue.Candidate) error

SaveWithResultAndProjections commits the immutable analysis, its Security Hotspot projection, and its code-quality issue projection in a single PostgreSQL transaction. A projection write failure rolls the analysis back, so the scan worker cannot publish a successful analysis without both projections.

func (*ProjectAnalysisStore) TransitionHotspot

func (*ProjectAnalysisStore) TransitionIssue

type ProjectRepository

type ProjectRepository struct {
	// contains filtered or unexported fields
}

ProjectRepository persists the long-lived Project identity. projects is RLS-protected (migration 0129), so every statement runs inside requireTenant and keeps its own tenant_id predicate. The pre-0129 reads dropped that predicate entirely when the caller passed a zero tenant, which turned a missing tenant into a cross-tenant listing; a missing tenant is now a validation error.

func NewProjectRepository

func NewProjectRepository(pool *pgxpool.Pool) *ProjectRepository

func (*ProjectRepository) AssignProfile

func (r *ProjectRepository) AssignProfile(ctx context.Context, tenantID shared.ID, projectKey, language, profileKey string) error

AssignProfile sets or clears the quality profile for a language in the project's JSONB default_profile_by_lang map, atomically at the column level (no read-modify-write race).

func (*ProjectRepository) CountByGate

func (r *ProjectRepository) CountByGate(ctx context.Context, tenantID shared.ID, gateID string) (int, error)

func (*ProjectRepository) Create

func (*ProjectRepository) DeleteByKey

func (r *ProjectRepository) DeleteByKey(ctx context.Context, tenantID shared.ID, key string) error

func (*ProjectRepository) GetByID

func (r *ProjectRepository) GetByID(ctx context.Context, tenantID, projectID shared.ID) (*project.Project, error)

func (*ProjectRepository) GetByKey

func (r *ProjectRepository) GetByKey(ctx context.Context, tenantID shared.ID, key string) (*project.Project, error)

func (*ProjectRepository) List

func (r *ProjectRepository) List(ctx context.Context, tenantID shared.ID) ([]*project.Project, error)

func (*ProjectRepository) SetPullRequestDecoration added in v0.2.0

func (r *ProjectRepository) SetPullRequestDecoration(ctx context.Context, tenantID shared.ID, key string, enabled bool) error

func (*ProjectRepository) UpdateGate

func (r *ProjectRepository) UpdateGate(ctx context.Context, tenantID shared.ID, key, gateID string) error

type PromotionStore added in v0.1.8

type PromotionStore struct {
	// contains filtered or unexported fields
}

PromotionStore persists promotion lifecycle events to PostgreSQL. Every operation runs inside a single RLS-scoped transaction via WithContextTenant so tenant isolation is enforced at the database level.

Apply atomically:

  1. Acquires sorted advisory transaction locks on (judgment, fingerprint) to serialize concurrent idempotency checks (prevents deadlocks).
  2. Checks judgment-level idempotency (tenant+judgmentID).
  3. Checks fingerprint-level idempotency (tenant+fingerprint).
  4. Locks the finding FOR UPDATE, verifies CAS (priority + version).
  5. Binds command metadata to CAS state.
  6. Validates exact reversal for corroborating_signal_loss.
  7. Constructs and validates the PromotionEvent.
  8. Mutates the finding (priority + version) for escalating/de-escalating effects.
  9. Appends the event to the append-only table.

func NewPromotionStore added in v0.1.8

func NewPromotionStore(pool *pgxpool.Pool) (*PromotionStore, error)

NewPromotionStore returns a repository backed by the given pool.

func (*PromotionStore) Apply added in v0.1.8

func (r *PromotionStore) Apply(ctx context.Context, engagementID, findingID shared.ID, cmd ports.PromotionCommand) (out finding.Finding, err error)

Apply constructs a PromotionEvent from the command, persists it, and atomically moves the finding's priority. Returns the existing event on exact replay (same judgmentID), or shared.ErrConflict on semantic conflicts.

func (*PromotionStore) FindByJudgment added in v0.1.8

func (r *PromotionStore) FindByJudgment(ctx context.Context, engagementID, findingID, judgmentID shared.ID) (evt promotion.PromotionEvent, ok bool, err error)

FindByJudgment returns an event scoped to its tenant, engagement, and finding.

func (*PromotionStore) LatestByFinding added in v0.1.8

func (r *PromotionStore) LatestByFinding(ctx context.Context, engagementID, findingID shared.ID) (evt promotion.PromotionEvent, ok bool, err error)

LatestByFinding returns the most recent promotion event for a finding, or (zero, false) if none exist.

func (*PromotionStore) ListByFinding added in v0.1.8

func (r *PromotionStore) ListByFinding(ctx context.Context, engagementID, findingID shared.ID) (out []promotion.PromotionEvent, err error)

ListByFinding returns all promotion events for a finding, oldest first.

func (*PromotionStore) ListPendingAudits added in v0.1.8

func (r *PromotionStore) ListPendingAudits(ctx context.Context, engagementID shared.ID) (out []promotion.PromotionEvent, err error)

ListPendingAudits returns applied events whose required audit record has not been acknowledged. The status row is created atomically with each event.

func (*PromotionStore) MarkAuditComplete added in v0.1.8

func (r *PromotionStore) MarkAuditComplete(ctx context.Context, eventID shared.ID) error

MarkAuditComplete acknowledges an event's required audit record. Repeating the acknowledgement is idempotent.

type PurpleRepository added in v0.1.8

type PurpleRepository struct {
	// contains filtered or unexported fields
}

PurpleRepository persists purple-team coverage (migration 0077), tenant-scoped via WithTenant/RLS so one tenant's coverage is never visible to another.

func NewPurpleRepository added in v0.1.8

func NewPurpleRepository(pool *pgxpool.Pool) *PurpleRepository

NewPurpleRepository constructs the repository.

func (*PurpleRepository) ListByEngagement added in v0.1.8

func (r *PurpleRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]pcdom.Coverage, error)

ListByEngagement returns all coverage for an engagement in the ctx tenant, oldest first, so a trend across runs is queryable.

func (*PurpleRepository) ListByRun added in v0.1.8

func (r *PurpleRepository) ListByRun(ctx context.Context, runID shared.ID) ([]pcdom.Coverage, error)

ListByRun returns one run's coverage in the ctx tenant, ordered by technique.

func (*PurpleRepository) SaveCoverage added in v0.1.8

func (r *PurpleRepository) SaveCoverage(ctx context.Context, records []pcdom.Coverage) error

SaveCoverage upserts a run's coverage under the authenticated tenant, keyed (run, technique), in one transaction so a re-computation of a run is atomic.

type QualityGateMutator

type QualityGateMutator struct {
	// contains filtered or unexported fields
}

QualityGateMutator commits managed-gate writes and their audit records together. quality_gates and projects are RLS-protected (migration 0129), so every mutation runs inside requireTenant. That replaces the hand-rolled transaction plus manual set_config this used before: the tenant is now fail-closed when empty, and a nested tenant transaction is checked for a mismatch instead of being silently rebound. The audit_log policy reads the same binding, which is why the audit append still lands in the same transaction as the write it records.

func NewQualityGateMutator

func NewQualityGateMutator(pool *pgxpool.Pool) *QualityGateMutator

func (*QualityGateMutator) AssignProjectGate

func (m *QualityGateMutator) AssignProjectGate(ctx context.Context, tenantID shared.ID, projectKey, gateID string, audit ports.AuditEntry) error

func (*QualityGateMutator) CreateGate

func (m *QualityGateMutator) CreateGate(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate, audit ports.AuditEntry) error

func (*QualityGateMutator) CreateProjectWithGate

func (m *QualityGateMutator) CreateProjectWithGate(ctx context.Context, p *project.Project) error

func (*QualityGateMutator) DeleteGate

func (m *QualityGateMutator) DeleteGate(ctx context.Context, tenantID shared.ID, key string, audit ports.AuditEntry) error

func (*QualityGateMutator) UpdateGate

func (m *QualityGateMutator) UpdateGate(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate, audit ports.AuditEntry) error

type QualityGateStore

type QualityGateStore struct {
	// contains filtered or unexported fields
}

QualityGateStore persists tenant-scoped custom quality gates. quality_gates is RLS-protected (migration 0129), so every statement runs inside requireTenant: the transaction binds app.current_tenant and the SQL keeps its own tenant_id predicate as defense in depth.

func NewQualityGateStore

func NewQualityGateStore(pool *pgxpool.Pool) *QualityGateStore

func (*QualityGateStore) Create

func (s *QualityGateStore) Create(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate) error

func (*QualityGateStore) Delete

func (s *QualityGateStore) Delete(ctx context.Context, tenantID shared.ID, key string) error

func (*QualityGateStore) DeleteIfUnassigned

func (s *QualityGateStore) DeleteIfUnassigned(ctx context.Context, tenantID shared.ID, key string) error

func (*QualityGateStore) Get

func (s *QualityGateStore) Get(ctx context.Context, tenantID shared.ID, key string) (qualitygate.Gate, error)

func (*QualityGateStore) List

func (s *QualityGateStore) List(ctx context.Context, tenantID shared.ID) ([]qualitygate.Gate, error)

func (*QualityGateStore) Update

func (s *QualityGateStore) Update(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate) error

type QualityProfileStore

type QualityProfileStore struct {
	// contains filtered or unexported fields
}

QualityProfileStore persists tenant-scoped custom quality profiles. Built-in profiles are generated from the rule catalog and never stored here. quality_profiles is RLS-protected (migration 0129), so every statement runs inside requireTenant and keeps its own tenant_id predicate.

func NewQualityProfileStore

func NewQualityProfileStore(pool *pgxpool.Pool) *QualityProfileStore

func (*QualityProfileStore) Create

func (s *QualityProfileStore) Create(ctx context.Context, tenantID shared.ID, profile qualityprofile.Profile) error

func (*QualityProfileStore) Delete

func (s *QualityProfileStore) Delete(ctx context.Context, tenantID shared.ID, key string) error

func (*QualityProfileStore) Get

func (*QualityProfileStore) List

func (*QualityProfileStore) Update

func (s *QualityProfileStore) Update(ctx context.Context, tenantID shared.ID, profile qualityprofile.Profile) error

type ReconRunStore

type ReconRunStore struct {
	// contains filtered or unexported fields
}

ReconRunStore persists recon-run records.

func NewReconRunStore

func NewReconRunStore(pool *pgxpool.Pool) *ReconRunStore

NewReconRunStore returns a store backed by the given pool.

func (*ReconRunStore) Get

func (r *ReconRunStore) Get(ctx context.Context, id shared.ID) (recon.Run, error)

Get returns a run by id, or shared.ErrNotFound.

func (*ReconRunStore) ListByEngagement

func (r *ReconRunStore) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]recon.Run, error)

ListByEngagement returns an engagement's runs, newest first.

func (*ReconRunStore) ListStaleRunning

func (r *ReconRunStore) ListStaleRunning(ctx context.Context, olderThan time.Time, limit int) ([]recon.Run, error)

ListStaleRunning returns runs still 'running' that started before olderThan (≤ limit), oldest first – the stale-run sweeper's input.

func (*ReconRunStore) Save

func (r *ReconRunStore) Save(ctx context.Context, run recon.Run) error

Save upserts a run (used on create and on every stage/status update).

type ResponseHaltWriterRepository added in v0.2.0

type ResponseHaltWriterRepository struct {
	// contains filtered or unexported fields
}

ResponseHaltWriterRepository owns the only PostgreSQL pool permitted to manufacture response halt fences and immutable dispatch obligations. Its role is intentionally not usable for ordinary response, evidence, signing, or work-order mutations.

func NewResponseHaltWriterRepository added in v0.2.0

func NewResponseHaltWriterRepository(pool *pgxpool.Pool) *ResponseHaltWriterRepository

NewResponseHaltWriterRepository constructs the dedicated halt-writer store.

func (*ResponseHaltWriterRepository) AdvanceHaltGenerationWithAudit added in v0.2.0

AdvanceHaltGenerationWithAudit atomically inserts mandatory audit intent, locks and advances the tenant fence, and inserts the immutable dispatch obligation.

func (*ResponseHaltWriterRepository) CurrentHaltGeneration added in v0.2.0

func (r *ResponseHaltWriterRepository) CurrentHaltGeneration(ctx context.Context) (int64, error)

CurrentHaltGeneration reads a tenant fence without creating one. A missing fence is the unhalted generation zero state; only AdvanceHaltGenerationWithAudit may create it.

type ResponseObserverBindingRepository added in v0.2.0

type ResponseObserverBindingRepository struct {
	*FleetAuditRepository
	// contains filtered or unexported fields
}

func NewResponseObserverBindingRepository added in v0.2.0

func NewResponseObserverBindingRepository(pool *pgxpool.Pool) (*ResponseObserverBindingRepository, error)

func (*ResponseObserverBindingRepository) GetResponseObserverBinding added in v0.2.0

func (r *ResponseObserverBindingRepository) GetResponseObserverBinding(ctx context.Context, agentID shared.ID) (fleetagent.ResponseObserverBinding, error)

func (*ResponseObserverBindingRepository) ListResponseObserverBindings added in v0.2.0

func (r *ResponseObserverBindingRepository) ListResponseObserverBindings(ctx context.Context, assetID shared.ID) ([]fleetagent.ResponseObserverBinding, error)

func (*ResponseObserverBindingRepository) SaveResponseObserverBindingWithAudit added in v0.2.0

type ResponseRepository added in v0.1.8

type ResponseRepository struct {
	// contains filtered or unexported fields
}

ResponseRepository persists governed response actions (migration 0076), tenant-scoped via WithTenant/RLS so one tenant's actions are never visible to another.

func NewResponseRepository added in v0.1.8

func NewResponseRepository(pool *pgxpool.Pool) *ResponseRepository

NewResponseRepository constructs the repository.

func (*ResponseRepository) AcknowledgeResponseAudit added in v0.2.0

func (r *ResponseRepository) AcknowledgeResponseAudit(ctx context.Context, id string) error

AcknowledgeResponseAudit monotonically marks an obligation delivered to the audit chain.

func (*ResponseRepository) AcknowledgeResponseHaltDispatch added in v0.2.0

func (r *ResponseRepository) AcknowledgeResponseHaltDispatch(ctx context.Context, generation int64) error

AcknowledgeResponseHaltDispatch verifies the immutable delivery obligation still exists. Completion is deliberately not persisted: the application role can issue raw SQL, so it cannot receive a database-side capability to forge a delivery acknowledgement. Replaying the signed endpoint fence remains idempotent and is therefore the fail-safe recovery path.

func (*ResponseRepository) AdvanceHaltGenerationWithAudit added in v0.2.0

AdvanceHaltGenerationWithAudit is intentionally unavailable to the normal runtime role. PostgreSQL composition passes ResponseHaltWriterRepository to the response service instead.

func (*ResponseRepository) AttemptStillCurrent added in v0.2.0

func (r *ResponseRepository) AttemptStillCurrent(ctx context.Context, idempotencyKey string, state responsesaga.SagaState, at time.Time) (bool, error)

AttemptStillCurrent checks attempt state and halt generation in one tenant-scoped query.

func (*ResponseRepository) ClaimAttempt added in v0.2.0

func (r *ResponseRepository) ClaimAttempt(ctx context.Context, idempotencyKey string, from, to responsesaga.SagaState, at time.Time) (responsesaga.ResponseAttempt, bool, error)

ClaimAttempt uses one conditional UPDATE to ensure only one concurrent delivery owns execution.

func (*ResponseRepository) CurrentHaltGeneration added in v0.2.0

func (r *ResponseRepository) CurrentHaltGeneration(ctx context.Context) (int64, error)

CurrentHaltGeneration reads a tenant's fence without creating it. A missing row is the unhalted zero generation; only ResponseHaltWriterRepository can initialize a fence.

func (*ResponseRepository) EnqueueResponseAudit added in v0.2.0

EnqueueResponseAudit persists an idempotent response summary obligation.

func (*ResponseRepository) Get added in v0.1.8

Get returns the record for an id in the ctx tenant.

func (*ResponseRepository) GetAttempt added in v0.2.0

func (r *ResponseRepository) GetAttempt(ctx context.Context, idempotencyKey string) (responsesaga.ResponseAttempt, bool, error)

GetAttempt returns a tenant-scoped execution journal entry by idempotency key.

func (*ResponseRepository) ListAttemptsByState added in v0.2.0

func (r *ResponseRepository) ListAttemptsByState(ctx context.Context, states ...responsesaga.SagaState) ([]responsesaga.ResponseAttempt, error)

ListAttemptsByState returns matching tenant attempts in deterministic execution order.

func (*ResponseRepository) ListByState added in v0.1.8

func (r *ResponseRepository) ListByState(ctx context.Context, state rdom.State) ([]rdom.Record, error)

ListByState returns the ctx tenant's records in a state, deterministically ordered by id.

func (*ResponseRepository) ListPendingResponseAudits added in v0.2.0

func (r *ResponseRepository) ListPendingResponseAudits(ctx context.Context) ([]ports.ResponseAuditIntent, error)

ListPendingResponseAudits returns the calling tenant's pending obligations in stable order.

func (*ResponseRepository) ListPendingResponseHaltDispatches added in v0.2.0

func (r *ResponseRepository) ListPendingResponseHaltDispatches(ctx context.Context) ([]ports.ResponseHaltDispatch, error)

ListPendingResponseHaltDispatches returns durable executor-fence obligations in generation order.

func (*ResponseRepository) Put added in v0.1.8

Put creates a response record or refreshes mutable fields without changing state or immutable identity.

func (*ResponseRepository) StartAttempt added in v0.2.0

StartAttempt durably inserts the pre-side-effect execution journal entry. Concurrent redeliveries race on the idempotency primary key and all observe the same immutable attempt identity.

func (*ResponseRepository) Transition added in v0.2.0

func (r *ResponseRepository) Transition(ctx context.Context, rec rdom.Record, from rdom.State) (bool, error)

Transition atomically replaces mutable fields and state when the persisted state still matches from.

func (*ResponseRepository) TransitionAttempt added in v0.2.0

TransitionAttempt atomically persists a state transition and its outcome/provenance payload.

func (*ResponseRepository) TransitionAttemptWithAudit added in v0.2.0

TransitionAttemptWithAudit atomically commits an attempt transition and its audit obligation.

func (*ResponseRepository) TransitionWithAudit added in v0.2.0

func (r *ResponseRepository) TransitionWithAudit(ctx context.Context, rec rdom.Record, from rdom.State, intent ports.ResponseAuditIntent) (bool, ports.ResponseAuditIntent, error)

TransitionWithAudit atomically commits a response transition and its audit obligation.

type ResponseVerificationRepository added in v0.2.0

type ResponseVerificationRepository struct {
	*FleetAuditRepository
	// contains filtered or unexported fields
}

func NewResponseVerificationRepository added in v0.2.0

func NewResponseVerificationRepository(pool *pgxpool.Pool) (*ResponseVerificationRepository, error)

func (*ResponseVerificationRepository) AppendResponseTargetEvidenceReceipt added in v0.2.0

func (*ResponseVerificationRepository) AppendResponseVerificationWithAudit added in v0.2.0

func (r *ResponseVerificationRepository) AppendResponseVerificationWithAudit(ctx context.Context, observation ports.AcceptedResponseVerification, intent ports.FleetAuditIntent) (ports.FleetAuditIntent, error)

func (*ResponseVerificationRepository) GetResponseTargetEvidenceReceipt added in v0.2.0

func (r *ResponseVerificationRepository) GetResponseTargetEvidenceReceipt(ctx context.Context, attemptKey string) (fleetagent.ResponseTargetEvidenceReceipt, bool, error)

func (*ResponseVerificationRepository) GetResponseVerification added in v0.2.0

func (r *ResponseVerificationRepository) GetResponseVerification(ctx context.Context, attemptKey string) (ports.AcceptedResponseVerification, bool, error)

type RestoreEvidenceChain added in v0.2.0

type RestoreEvidenceChain = ports.RestoreEvidenceChain

RestoreEvidenceChain is one engagement's restored evidence ledger in chain order.

type RestoreEvidenceStore added in v0.2.0

type RestoreEvidenceStore struct {
	// contains filtered or unexported fields
}

RestoreEvidenceStore enumerates a restored evidence ledger. It has no write methods and is constructed only after the database identity is proven able to bypass tenant RLS for recovery verification.

func NewReadOnlyEvidenceStore added in v0.2.0

func NewReadOnlyEvidenceStore(ctx context.Context, pool *pgxpool.Pool) (*RestoreEvidenceStore, error)

NewReadOnlyEvidenceStore constructs a recovery-only reader. The connection identity must have rolbypassrls, or be the explicitly verified database owner; normal application identities fail closed rather than relying on row_security. The returned reader opens every enumeration transaction READ ONLY.

func (*RestoreEvidenceStore) ListEvidenceChains added in v0.2.0

func (r *RestoreEvidenceStore) ListEvidenceChains(ctx context.Context) ([]RestoreEvidenceChain, error)

ListEvidenceChains enumerates every restored chain only inside a READ ONLY transaction. It neither sets row_security nor assumes that a normal RLS role can disable it; construction has already required the recovery identity.

type RetestRepository

type RetestRepository struct {
	// contains filtered or unexported fields
}

RetestRepository persists the per-finding retest history to PostgreSQL.

func NewRetestRepository

func NewRetestRepository(pool *pgxpool.Pool) *RetestRepository

NewRetestRepository returns a repository backed by the given pool.

func (*RetestRepository) Add

Add inserts a retest (append-only; retests are not edited or deleted in app code).

func (*RetestRepository) LatestByEngagementFindings added in v0.2.0

func (r *RetestRepository) LatestByEngagementFindings(ctx context.Context, engagementID shared.ID, findingIDs []shared.ID) (map[shared.ID]finding.Retest, error)

LatestByEngagementFindings projects only the effective decision, not arbitrary notes.

func (*RetestRepository) ListByEngagementFinding

func (r *RetestRepository) ListByEngagementFinding(ctx context.Context, engagementID, findingID shared.ID) (out []finding.Retest, err error)

ListByEngagementFinding returns a finding's retests oldest-first, scoped to the engagement (no cross-engagement read).

func (*RetestRepository) RetestHistories added in v0.2.0

func (r *RetestRepository) RetestHistories(ctx context.Context, engagementID shared.ID, findingIDs []shared.ID) (map[shared.ID][]finding.Retest, error)

type RunLock

type RunLock struct {
	// contains filtered or unexported fields
}

RunLock implements ports.RunLocker with a PostgreSQL session advisory lock keyed by the run id (F9). The lock is held on a dedicated pooled connection for the duration of the execution and released (with the connection) afterwards – so across the API + the worker, at most one delivery of a given run executes at a time. A redelivery that finds the lock held gets ok=false and skips, preventing a duplicate live scan.

func NewRunLock

func NewRunLock(pool *pgxpool.Pool) *RunLock

NewRunLock returns a Postgres-backed run locker.

func (*RunLock) TryLock

func (l *RunLock) TryLock(ctx context.Context, runID string) (func(), bool, error)

type SCMConnectorRepository added in v0.2.0

type SCMConnectorRepository struct {
	// contains filtered or unexported fields
}

SCMConnectorRepository persists tenant-scoped source-control connectors on PostgreSQL. The token is sealed with the vault cipher (AES-256-GCM, AAD-bound to tenant+id) so the column holds only ciphertext; the plaintext is returned ONLY by ResolveGitCredential at clone time. scm_connectors is RLS-protected (migration 0133): every statement runs inside a WithTenant transaction and carries an explicit tenant_id predicate as defense in depth.

func NewSCMConnectorRepository added in v0.2.0

func NewSCMConnectorRepository(pool *pgxpool.Pool, cipher *vault.Cipher) (*SCMConnectorRepository, error)

NewSCMConnectorRepository returns a repository over the pool and the vault cipher. The cipher is required: without it a token could only be stored in the clear, which the contract forbids.

func (*SCMConnectorRepository) Delete added in v0.2.0

func (r *SCMConnectorRepository) Delete(ctx context.Context, id shared.ID) error

Delete removes one connector; ErrNotFound when absent under the tenant.

func (*SCMConnectorRepository) Get added in v0.2.0

Get returns one connector's metadata; ErrNotFound when absent under the tenant.

func (*SCMConnectorRepository) List added in v0.2.0

List returns the tenant's connectors as metadata (never the token), ordered by host then name.

func (*SCMConnectorRepository) Put added in v0.2.0

Put upserts a connector under the caller's tenant, sealing the token bound to the connector's (tenant, id, host, username) identity. A token is always required. The UNIQUE(tenant_id, host) constraint makes a second connector for the same host a conflict.

func (*SCMConnectorRepository) ResolveGitCredential added in v0.2.0

func (r *SCMConnectorRepository) ResolveGitCredential(ctx context.Context, host string) (ports.GitCredential, bool, error)

ResolveGitCredential returns the credential for a normalized clone-URL host under the caller's tenant, opening the sealed token. ok=false when the tenant has no connector for that host.

type SLAStore added in v0.1.8

type SLAStore struct {
	// contains filtered or unexported fields
}

SLAStore is the PostgreSQL adapter for versioned SLA policy, immutable assessments, and the human remediation lifecycle. Every operation is routed through WithTenant so RLS remains the final isolation boundary even when an application-level predicate regresses.

func NewSLAStore added in v0.1.8

func NewSLAStore(pool *pgxpool.Pool) *SLAStore

func (*SLAStore) ActivePolicy added in v0.1.8

func (s *SLAStore) ActivePolicy(ctx context.Context, tenantID shared.ID) (sla.Policy, error)

func (*SLAStore) AssessmentHistory added in v0.1.8

func (s *SLAStore) AssessmentHistory(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.Assessment, error)

func (*SLAStore) Current added in v0.1.8

func (s *SLAStore) Current(ctx context.Context, tenantID, engagementID, findingID shared.ID) (sla.Current, error)

func (*SLAStore) LifecycleEvents added in v0.1.8

func (s *SLAStore) LifecycleEvents(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.LifecycleEvent, error)

func (*SLAStore) ListCurrent added in v0.1.8

func (s *SLAStore) ListCurrent(ctx context.Context, tenantID, engagementID shared.ID) ([]sla.Current, error)

func (*SLAStore) PolicyHistory added in v0.1.8

func (s *SLAStore) PolicyHistory(ctx context.Context, tenantID shared.ID) ([]sla.Policy, error)

func (*SLAStore) PutPolicy added in v0.1.8

func (s *SLAStore) PutPolicy(ctx context.Context, policy sla.Policy, activate bool) (bool, error)

func (*SLAStore) SLAHistories added in v0.2.0

func (s *SLAStore) SLAHistories(ctx context.Context, tenantID, engagementID shared.ID, findingIDs []shared.ID) (ports.SLAHistoryBatch, error)

func (*SLAStore) SaveTransition added in v0.1.8

func (s *SLAStore) SaveTransition(ctx context.Context, next sla.Lifecycle, event sla.LifecycleEvent) error

func (*SLAStore) UpsertAssessment added in v0.1.8

func (s *SLAStore) UpsertAssessment(ctx context.Context, assessment sla.Assessment) (sla.AssessmentUpsertResult, error)

type ScanJobStore

type ScanJobStore struct {
	// contains filtered or unexported fields
}

ScanJobStore persists asynchronous scan-job status.

func NewScanJobStore

func NewScanJobStore(pool *pgxpool.Pool) *ScanJobStore

NewScanJobStore returns a store backed by the given pool.

func (*ScanJobStore) CreateRunning

func (r *ScanJobStore) CreateRunning(ctx context.Context, j ports.ScanJob) error

func (*ScanJobStore) GetJob

func (r *ScanJobStore) GetJob(ctx context.Context, id string) (ports.ScanJob, error)

GetJob returns a scan job by its own id, or ErrNotFound.

func (*ScanJobStore) LatestForEngagement

func (r *ScanJobStore) LatestForEngagement(ctx context.Context, engagementID shared.ID) (ports.ScanJob, error)

func (*ScanJobStore) LatestForEngagements

func (r *ScanJobStore) LatestForEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.ScanJob, error)

LatestForEngagement returns the engagement's most recent scan job, or ErrNotFound.

func (*ScanJobStore) ListStaleRunning

func (r *ScanJobStore) ListStaleRunning(ctx context.Context, olderThan time.Time, limit int) ([]ports.ScanJob, error)

ListStaleRunning returns scan jobs still 'running' that started before olderThan (≤ limit), oldest first – the stale-scan sweeper's input.

func (*ScanJobStore) Save

func (r *ScanJobStore) Save(ctx context.Context, j ports.ScanJob) error

Save upserts a scan job (used on create and on every stage/status update).

type ScanRepository

type ScanRepository struct {
	// contains filtered or unexported fields
}

ScanRepository persists SCA scans (SBOM + components + vulnerabilities).

func NewScanRepository

func NewScanRepository(pool *pgxpool.Pool) *ScanRepository

NewScanRepository returns a repository backed by the given pool.

func (*ScanRepository) AdmitInventory added in v0.2.0

func (r *ScanRepository) AdmitInventory(ctx context.Context, engagementID shared.ID, scope string, admittedAt time.Time) (sbom.InventoryAdmission, error)

func (*ScanRepository) SaveScan

func (r *ScanRepository) SaveScan(ctx context.Context, engagementID shared.ID, doc *sbom.SBOM, vulns []vulnerability.Vulnerability, snap ports.ScanSnapshot) (ports.ScanSaveResult, error)

SaveScan stores the SBOM, its components, and the vulnerabilities found against them in one transaction – a new immutable snapshot per scan. It returns the number of vulns that could not be linked to a component in this SBOM (skipped, never orphaned); the caller surfaces a non-zero count on the audit log so a dropped advisory is never invisible on a chain-of-custody tool.

type ScanResultStore

type ScanResultStore struct {
	// contains filtered or unexported fields
}

ScanResultStore caches the latest full scan result (JSON) per engagement.

func NewScanResultStore

func NewScanResultStore(pool *pgxpool.Pool) *ScanResultStore

NewScanResultStore returns a store backed by the given pool.

func (*ScanResultStore) LatestResult

func (r *ScanResultStore) LatestResult(ctx context.Context, engagementID shared.ID) ([]byte, error)

LatestResult returns the engagement's cached scan result, or shared.ErrNotFound.

func (*ScanResultStore) SaveResult

func (r *ScanResultStore) SaveResult(ctx context.Context, engagementID shared.ID, result []byte) error

SaveResult upserts the engagement's latest scan result JSON.

type ScanRunStore

type ScanRunStore struct {
	// contains filtered or unexported fields
}

ScanRunStore persists scan-run manifests, finding keys, and tenant-scoped sealed provenance.

func NewScanRunStore

func NewScanRunStore(pool *pgxpool.Pool) *ScanRunStore

NewScanRunStore returns a store backed by the given pool.

func (*ScanRunStore) Get

func (r *ScanRunStore) Get(ctx context.Context, runID string) (ports.ScanRun, error)

Get returns a legacy scan run by ID.

func (*ScanRunStore) GetScanRun added in v0.2.0

func (r *ScanRunStore) GetScanRun(ctx context.Context, tenantID shared.ID, runID string) (scanrun.ScanRun, error)

GetScanRun retrieves a full scan run aggregate including lanes, versions, and stages.

func (*ScanRunStore) GetScanRunEvidence added in v0.2.0

func (r *ScanRunStore) GetScanRunEvidence(ctx context.Context, tenantID shared.ID, runID string) (ports.ScanRunEvidence, error)

func (*ScanRunStore) List

func (r *ScanRunStore) List(ctx context.Context, engagementID shared.ID) ([]ports.ScanRun, error)

List returns legacy scan runs for an engagement.

func (*ScanRunStore) ListScanRuns added in v0.2.0

func (r *ScanRunStore) ListScanRuns(ctx context.Context, tenantID, engagementID shared.ID) ([]scanrun.ScanRun, error)

ListScanRuns returns all scan runs for an engagement, ordered newest first.

func (*ScanRunStore) Save

func (r *ScanRunStore) Save(ctx context.Context, run ports.ScanRun) error

Save records a legacy scan run for backwards compatibility.

func (*ScanRunStore) SaveScanRun added in v0.2.0

func (r *ScanRunStore) SaveScanRun(ctx context.Context, run scanrun.ScanRun) error

SaveScanRun persists a tenant-owned native or legacy scan run.

func (*ScanRunStore) SaveScanRunEvidence added in v0.2.0

func (r *ScanRunStore) SaveScanRunEvidence(ctx context.Context, item ports.ScanRunEvidence) error

func (*ScanRunStore) SealScanRun added in v0.2.0

func (r *ScanRunStore) SealScanRun(ctx context.Context, command ports.SealScanRunCommand) error

SealScanRun atomically acquires a row lock on scan_runs, validates seal state, and inserts lanes/versions/stages.

type ScannedImageStore added in v0.1.8

type ScannedImageStore struct {
	// contains filtered or unexported fields
}

ScannedImageStore is the Postgres-backed scanned-image digest index (#446). Every operation routes through WithTenant so the Row Level Security policy on scanned_image (migration 0063) enforces tenant isolation at the database — a query that bypassed WithTenant would resolve the tenant to NULL and see nothing.

func NewScannedImageStore added in v0.1.8

func NewScannedImageStore(pool *pgxpool.Pool) *ScannedImageStore

NewScannedImageStore constructs the Postgres scanned-image store.

func (*ScannedImageStore) MarkScanned added in v0.1.8

func (s *ScannedImageStore) MarkScanned(ctx context.Context, tenantID shared.ID, digest string, at time.Time) error

MarkScanned records digest as scanned for the tenant. Idempotent by (tenant, digest): a repeat keeps the earliest first_scanned_at.

func (*ScannedImageStore) ScannedDigests added in v0.1.8

func (s *ScannedImageStore) ScannedDigests(ctx context.Context, tenantID shared.ID) (map[string]bool, error)

ScannedDigests returns the set of scanned digests for the tenant.

type SensorStateRepository added in v0.2.0

type SensorStateRepository struct {
	*FleetAuditRepository
	// contains filtered or unexported fields
}

func NewSensorStateRepository added in v0.2.0

func NewSensorStateRepository(pool *pgxpool.Pool) (*SensorStateRepository, error)

func (*SensorStateRepository) AppendSensorState added in v0.2.0

func (r *SensorStateRepository) AppendSensorState(ctx context.Context, observation sensorstate.Observation) error

func (*SensorStateRepository) AppendSensorStateWithAudit added in v0.2.0

func (r *SensorStateRepository) AppendSensorStateWithAudit(
	ctx context.Context,
	observation sensorstate.Observation,
	intent ports.FleetAuditIntent,
) (ports.FleetAuditIntent, error)

func (*SensorStateRepository) ListCoverageSensorStates added in v0.2.0

func (*SensorStateRepository) ListSensorStates added in v0.2.0

type SyncRunStore added in v0.1.8

type SyncRunStore struct {
	// contains filtered or unexported fields
}

func NewSyncRunStore added in v0.1.8

func NewSyncRunStore(pool *pgxpool.Pool, ids ports.IDGenerator) *SyncRunStore

func (*SyncRunStore) Advance added in v0.1.8

func (s *SyncRunStore) Advance(ctx context.Context, id shared.ID, expectedCheckpoint, nextCheckpoint []byte, counts vulnerabilitysync.Counts, errors []string) (vulnerabilitysync.Run, error)

func (*SyncRunStore) Finish added in v0.1.8

func (*SyncRunStore) Get added in v0.1.8

func (*SyncRunStore) GetByDurableJobID added in v0.1.8

func (s *SyncRunStore) GetByDurableJobID(ctx context.Context, jobID string) (vulnerabilitysync.Run, error)

func (*SyncRunStore) GetVulnerabilitySyncRun added in v0.1.8

func (s *SyncRunStore) GetVulnerabilitySyncRun(ctx context.Context, tenantID, id shared.ID) (vulnerabilityintel.SyncRunItem, error)

func (*SyncRunStore) LatestForSource added in v0.1.8

func (s *SyncRunStore) LatestForSource(ctx context.Context, sourceID shared.ID, states []vulnerabilitysync.State) (vulnerabilitysync.Run, error)

func (*SyncRunStore) LatestSuccessfulVulnerabilitySync added in v0.1.8

func (s *SyncRunStore) LatestSuccessfulVulnerabilitySync(ctx context.Context, tenantID shared.ID) (*time.Time, error)

func (*SyncRunStore) ListStale added in v0.1.8

func (s *SyncRunStore) ListStale(ctx context.Context, olderThan time.Time, limit int) ([]vulnerabilitysync.Run, error)

func (*SyncRunStore) ListVulnerabilitySyncRuns added in v0.1.8

func (s *SyncRunStore) ListVulnerabilitySyncRuns(ctx context.Context, query vulnerabilityintel.SyncRunQuery) (vulnerabilityintel.SyncRunPage, error)

func (*SyncRunStore) MarkRunning added in v0.1.8

func (s *SyncRunStore) MarkRunning(ctx context.Context, id shared.ID) error

func (*SyncRunStore) RecoverStale added in v0.1.8

func (s *SyncRunStore) RecoverStale(ctx context.Context, staleRunID shared.ID, staleBefore time.Time, request ports.SyncRunStart) (vulnerabilitysync.Run, bool, error)

func (*SyncRunStore) Start added in v0.1.8

func (*SyncRunStore) Supersede added in v0.1.8

func (s *SyncRunStore) Supersede(ctx context.Context, id shared.ID) error

type TelemetryRepository added in v0.1.8

type TelemetryRepository struct {
	// contains filtered or unexported fields
}

TelemetryRepository is the CE columnar-tier store (migration 0075) behind ports.TelemetryStore, tenant-scoped via WithTenant/RLS. It is reached only through the port and appears in no domain type.

func NewTelemetryRepository added in v0.1.8

func NewTelemetryRepository(pool *pgxpool.Pool, hot, warm time.Duration) *TelemetryRepository

NewTelemetryRepository constructs the store with the hot/warm tier boundaries (ADR 0001 config).

func (*TelemetryRepository) Footprint added in v0.1.8

Footprint reports the GLOBAL store size — an operator spend metric across all tenants, not a per-tenant figure. Both numbers are global and coherent: an estimated row count from the planner statistics (pg_class.reltuples) paired with the real on-disk size (pg_total_relation_size). These are catalog reads, unaffected by RLS, and carry only counts/bytes — never tenant data — so a global scope leaks nothing. reltuples is an estimate (refreshed by ANALYZE/autovacuum), which is the right shape for "predict spend, don't discover it".

func (*TelemetryRepository) Ingest added in v0.1.8

Ingest bulk-inserts the batch's events, idempotent on (tenant, host, class, seq, idx).

func (*TelemetryRepository) LastSequence added in v0.1.8

func (r *TelemetryRepository) LastSequence(ctx context.Context, hostID shared.ID, class detection.Class) (uint64, error)

LastSequence returns the highest stored seq for a (host, class) in the ctx tenant.

func (*TelemetryRepository) Query added in v0.1.8

Query runs a retro-hunt over a window and reports completeness honestly (sampled + sequence gaps).

func (*TelemetryRepository) RecordLoss added in v0.2.0

func (r *TelemetryRepository) RecordLoss(ctx context.Context, loss ports.TelemetryLoss) error

RecordLoss persists a Truncated/Dropped loss record for the ctx tenant, idempotent on (host, class, seq, disposition) via the primary key, so a re-ingest records one loss.

func (*TelemetryRepository) RetentionSweep added in v0.1.8

func (r *TelemetryRepository) RetentionSweep(ctx context.Context, now time.Time) (ports.SweepReport, error)

RetentionSweep down-samples the warm window (drops 1-in-2 by idx for reduced resolution) and expires the past-warm window, for the ctx tenant. Returns the counts so the caller can audit the expiry.

type TelemetryTransportRepository added in v0.2.0

type TelemetryTransportRepository struct {
	*FleetAuditRepository
	// contains filtered or unexported fields
}

TelemetryTransportRepository is the Postgres tier for A3 transport state. The ACK snapshot remains the sequencing source of truth; migration 0110 materializes its open gaps transactionally and stores the authenticated agent->asset binding.

func NewTelemetryTransportRepository added in v0.2.0

func NewTelemetryTransportRepository(pool *pgxpool.Pool) (*TelemetryTransportRepository, error)

func (*TelemetryTransportRepository) AcceptAgentGapRevision added in v0.2.0

func (r *TelemetryTransportRepository) AcceptAgentGapRevision(ctx context.Context, revision ports.TelemetryAgentGapRevision) error

AcceptAgentGapRevision appends the exact signed report and advances the mutable current projection in one tenant-scoped PostgreSQL transaction. It preserves nanosecond signed-time identity in dedicated integer columns because PostgreSQL timestamptz normalizes its queryable projection to microsecond precision.

func (*TelemetryTransportRepository) AcceptAgentGapRevisionWithAudit added in v0.2.0

func (r *TelemetryTransportRepository) AcceptAgentGapRevisionWithAudit(
	ctx context.Context,
	revision ports.TelemetryAgentGapRevision,
	intent ports.FleetAuditIntent,
) (ports.FleetAuditIntent, error)

func (*TelemetryTransportRepository) AgentGapRevisions added in v0.2.0

AgentGapRevisions returns immutable signed reports in revision order. It reconstructs signed timestamps from their exact Unix-nanosecond identity.

func (*TelemetryTransportRepository) BindTelemetryAsset added in v0.2.0

func (r *TelemetryTransportRepository) BindTelemetryAsset(ctx context.Context, binding ports.TelemetryAssetBinding) error

func (*TelemetryTransportRepository) CommitBatch added in v0.2.0

CommitBatch atomically claims one delivery coordinate for the exact signed batch identity. Concurrent identical replays converge; any equivocation at the same (tenant, agent, stream, epoch, sequence) fails closed before ACK classification.

func (*TelemetryTransportRepository) CommitBatchWithAudit added in v0.2.0

func (*TelemetryTransportRepository) CountBatchEvents added in v0.2.0

func (r *TelemetryTransportRepository) CountBatchEvents(ctx context.Context, agentID, streamID shared.ID, epoch, sequence uint64) (int, error)

func (*TelemetryTransportRepository) IngestBatchEvents added in v0.2.0

func (r *TelemetryTransportRepository) IngestBatchEvents(ctx context.Context, batch ports.TelemetryEventBatch) (int, error)

func (*TelemetryTransportRepository) ListCoverageGapFacts added in v0.2.0

ListCoverageGapFacts exposes exact loss provenance for deterministic coverage revisions. Agent-origin and inferred facts remain separate even when they describe the same delivery coordinate.

func (*TelemetryTransportRepository) ListGapChanges added in v0.2.0

func (r *TelemetryTransportRepository) ListGapChanges(
	ctx context.Context,
	agentID, streamID shared.ID,
	epoch, sequence uint64,
) ([]ports.TelemetryGap, error)

ListGapChanges returns both open and resolved inferred gaps affected by one sequence so exact source retries can repair failed coverage materialization.

func (*TelemetryTransportRepository) ListGaps added in v0.2.0

func (r *TelemetryTransportRepository) ListGaps(ctx context.Context, agentID, streamID shared.ID) ([]ports.TelemetryGap, error)

func (*TelemetryTransportRepository) ListTelemetryAssetBindings added in v0.2.0

func (r *TelemetryTransportRepository) ListTelemetryAssetBindings(ctx context.Context) ([]ports.TelemetryAssetBinding, error)

ListTelemetryAssetBindings returns the tenant's current agent→asset bindings (#633 desired-vs-observed), tenant-scoped from ctx and ordered by agent id.

func (*TelemetryTransportRepository) MaxEpoch added in v0.2.0

func (r *TelemetryTransportRepository) MaxEpoch(ctx context.Context, agentID, streamID shared.ID) (uint64, error)

func (*TelemetryTransportRepository) QueryAgentGaps added in v0.2.0

QueryAgentGaps returns durable agent-origin loss whose observed span overlaps the requested hunt window. It never consults/resolves delivery ACK state: these rows describe local loss facts that remain provenance even if sequence holes fill.

func (*TelemetryTransportRepository) QueryDeliveryGaps added in v0.2.0

func (*TelemetryTransportRepository) QueryTelemetryBatchAccounting added in v0.2.0

func (*TelemetryTransportRepository) RecordAgentGap added in v0.2.0

func (*TelemetryTransportRepository) ResolveTelemetryAsset added in v0.2.0

func (r *TelemetryTransportRepository) ResolveTelemetryAsset(ctx context.Context, agentID shared.ID) (shared.ID, error)

func (*TelemetryTransportRepository) ResolveTelemetryReferences added in v0.2.0

func (r *TelemetryTransportRepository) ResolveTelemetryReferences(ctx context.Context, agentID, assetID shared.ID, redactionPolicyDigest string, refs []fleetagent.TelemetryReference) (ports.TelemetryReferenceStatus, error)

TelemetryReferencesDurable resolves a detection's causal references from telemetry_batch_events, the existing accepted raw-telemetry fact store. Missing or mismatched facts are pending, never inferred.

func (*TelemetryTransportRepository) SaveStreamState added in v0.2.0

func (*TelemetryTransportRepository) StreamState added in v0.2.0

func (r *TelemetryTransportRepository) StreamState(ctx context.Context, agentID, streamID shared.ID, epoch uint64) (ports.TelemetryStreamState, error)

type TenantTransactionRunner added in v0.1.8

type TenantTransactionRunner struct {
	// contains filtered or unexported fields
}

func NewTenantTransactionRunner added in v0.1.8

func NewTenantTransactionRunner(pool *pgxpool.Pool) *TenantTransactionRunner

func (*TenantTransactionRunner) Run added in v0.1.8

func (runner *TenantTransactionRunner) Run(ctx context.Context, tenantID shared.ID, fn func(context.Context) error) error

type ThreatModelRepository

type ThreatModelRepository struct {
	// contains filtered or unexported fields
}

ThreatModelRepository persists the architecture-input threat model per engagement to PostgreSQL: one row per engagement, the validated domain model stored as a JSONB blob. threat_models is RLS-protected (migration 0129): both statements run inside a WithTenant transaction and carry their own tenant_id predicate, so a read can no longer reach another tenant's engagement id even if the upstream route gate is bypassed.

func NewThreatModelRepository

func NewThreatModelRepository(pool *pgxpool.Pool) *ThreatModelRepository

NewThreatModelRepository returns a repository backed by the given pool.

func (*ThreatModelRepository) Get

func (r *ThreatModelRepository) Get(ctx context.Context, engagementID shared.ID) (threatmodel.Model, bool, error)

Get decodes the engagement's model from its JSONB blob; ok=false when none has been ingested. The port carries no tenant argument, so the read runs under the ambient tenant bound to ctx; a context without one is a validation error, never a cross-tenant read.

func (*ThreatModelRepository) Save

func (r *ThreatModelRepository) Save(ctx context.Context, engagementID, tenantID shared.ID, m threatmodel.Model) error

Save upserts the engagement's model (the usecase has already bounded size + validated it), bumping version on each re-ingest. The model round-trips through the JSONB `data` blob.

type TimestampStore

type TimestampStore struct {
	// contains filtered or unexported fields
}

TimestampStore persists external RFC-3161 tokens for chain heads on PostgreSQL, out-of-band from the report. One token per (chain, engagement, head).

func NewTimestampStore

func NewTimestampStore(pool *pgxpool.Pool) *TimestampStore

NewTimestampStore returns a timestamp store backed by the given pool.

func (*TimestampStore) Get

func (s *TimestampStore) Get(ctx context.Context, chain string, eng shared.ID, head string) (*ports.TimestampToken, error)

Get returns the stored token for a head, or nil if it is not yet anchored.

func (*TimestampStore) LatestHead

func (s *TimestampStore) LatestHead(ctx context.Context, chain string, eng shared.ID) (string, bool, error)

LatestHead returns the most-recently-anchored head for a chain (ok=false if none) – the retained head for out-of-band tail-truncation detection.

func (*TimestampStore) Put

func (s *TimestampStore) Put(ctx context.Context, chain string, eng shared.ID, head string, token ports.TimestampToken) error

Put stores a token for a head, idempotent per (chain, engagement, head).

type UserRepository

type UserRepository struct {
	// contains filtered or unexported fields
}

UserRepository persists operator identities to PostgreSQL.

func NewUserRepository

func NewUserRepository(pool *pgxpool.Pool) *UserRepository

NewUserRepository returns a repository backed by the given pool.

func (*UserRepository) Bootstrap added in v0.2.0

func (r *UserRepository) Bootstrap(ctx context.Context, u *user.User, auditEntry ports.AuditEntry) error

Bootstrap atomically seeds or refreshes the bootstrap administrator. The audit row is appended in the same transaction only when the user is first inserted, so concurrent API startups cannot create duplicate bootstrap audit events.

func (*UserRepository) Create

func (r *UserRepository) Create(ctx context.Context, u *user.User) error

func (*UserRepository) GetByAPIKeyHash

func (r *UserRepository) GetByAPIKeyHash(ctx context.Context, hash string) (*user.User, error)

GetByAPIKeyHash is the authentication path: the tenant is unknown until the presented token resolves to a user, so this is the one user lookup without a tenant predicate. The key is the SHA-256 digest of a 192-bit random secret, and the resolved user's own tenant scopes every subsequent read and write.

func (*UserRepository) GetByID

func (r *UserRepository) GetByID(ctx context.Context, tenantID, id shared.ID) (*user.User, error)

func (*UserRepository) List

func (r *UserRepository) List(ctx context.Context, tenantID shared.ID) ([]*user.User, error)

func (*UserRepository) ListForUpdate added in v0.2.0

func (r *UserRepository) ListForUpdate(ctx context.Context, tenantID shared.ID) ([]*user.User, error)

ListForUpdate is List with the tenant's rows locked for the rest of the caller's transaction, so the last-admin guard's count cannot be invalidated by a concurrent demotion between the count and the write. Outside a transaction the lock is released immediately and this is just List.

func (*UserRepository) Update added in v0.2.0

func (r *UserRepository) Update(ctx context.Context, tenantID shared.ID, u *user.User) error

Update writes the mutable fields of a user that already exists in tenantID. tenant_id is absent from the SET list, so an update can never move a user between tenants, and the tenant predicate means a cross-tenant id updates nothing.

func (*UserRepository) Upsert

func (r *UserRepository) Upsert(ctx context.Context, u *user.User) error

type VEXStatementRepository added in v0.2.0

type VEXStatementRepository struct {
	// contains filtered or unexported fields
}

VEXStatementRepository persists ingested VEX statements (migration 0173) to PostgreSQL so they can be re-applied after a rescan. Every method runs through WithTenant, so Row Level Security isolates by tenant.

func NewVEXStatementRepository added in v0.2.0

func NewVEXStatementRepository(pool *pgxpool.Pool) *VEXStatementRepository

NewVEXStatementRepository constructs the Postgres VEX-statement repository.

func (*VEXStatementRepository) ListByEngagement added in v0.2.0

func (r *VEXStatementRepository) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID) ([]vex.StoredStatement, error)

ListByEngagement returns the engagement's persisted statements in insertion order (the seq identity), so a re-apply walking them applies the most-recent assertion last.

func (*VEXStatementRepository) Save added in v0.2.0

func (r *VEXStatementRepository) Save(ctx context.Context, tenantID, engagementID shared.ID, statements []vex.StoredStatement) error

Save persists a batch ATOMICALLY: one transaction, so a partial failure leaves no half-persisted policy. Each row is inserted ON CONFLICT DO NOTHING against the (tenant, engagement, digest) primary key, so re-importing an identical assertion is a no-op rather than a duplicate.

type VulnerabilityActionStore added in v0.1.8

type VulnerabilityActionStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilityActionStore added in v0.1.8

func NewVulnerabilityActionStore(pool *pgxpool.Pool) *VulnerabilityActionStore

func (*VulnerabilityActionStore) AcknowledgeAction added in v0.1.8

func (s *VulnerabilityActionStore) AcknowledgeAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) ClaimOutbox added in v0.1.8

func (s *VulnerabilityActionStore) ClaimOutbox(ctx context.Context, tenantID shared.ID, now, lockedUntil time.Time, limit int) ([]vulnerabilityaction.OutboxEvent, error)

func (*VulnerabilityActionStore) CompleteOutbox added in v0.1.8

func (s *VulnerabilityActionStore) CompleteOutbox(ctx context.Context, tenantID, eventID shared.ID, at time.Time) error

func (*VulnerabilityActionStore) CountPendingVulnerabilityActions added in v0.1.8

func (s *VulnerabilityActionStore) CountPendingVulnerabilityActions(ctx context.Context, tenantID shared.ID) (int64, error)

func (*VulnerabilityActionStore) GetAction added in v0.1.8

func (s *VulnerabilityActionStore) GetAction(ctx context.Context, tenantID, actionID shared.ID) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) ListActions added in v0.1.8

func (*VulnerabilityActionStore) ListVulnerabilityTransitions added in v0.1.8

func (*VulnerabilityActionStore) RecordChange added in v0.1.8

func (*VulnerabilityActionStore) ResolveAction added in v0.1.8

func (s *VulnerabilityActionStore) ResolveAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) RetryOutbox added in v0.1.8

func (s *VulnerabilityActionStore) RetryOutbox(ctx context.Context, tenantID, eventID shared.ID, at, availableAt time.Time, lastError string, terminal bool) error

func (*VulnerabilityActionStore) SummarizeVulnerabilityActions added in v0.1.8

func (s *VulnerabilityActionStore) SummarizeVulnerabilityActions(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryActionSummary, error)

type VulnerabilityOccurrenceStore added in v0.1.8

type VulnerabilityOccurrenceStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilityOccurrenceStore added in v0.1.8

func NewVulnerabilityOccurrenceStore(pool *pgxpool.Pool) *VulnerabilityOccurrenceStore

func (*VulnerabilityOccurrenceStore) CountActiveVulnerabilityOccurrences added in v0.1.8

func (s *VulnerabilityOccurrenceStore) CountActiveVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryID string) (int64, error)

func (*VulnerabilityOccurrenceStore) CountNewlyAffectedAssets added in v0.2.0

func (s *VulnerabilityOccurrenceStore) CountNewlyAffectedAssets(ctx context.Context, tenantID shared.ID, since time.Time) (int64, error)

func (*VulnerabilityOccurrenceStore) Get added in v0.1.8

func (s *VulnerabilityOccurrenceStore) Get(ctx context.Context, tenantID, engagementID shared.ID, advisoryID, componentFingerprint string) (vulnerabilityoccurrence.Occurrence, error)

func (*VulnerabilityOccurrenceStore) ListByEngagement added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID, states []vulnerabilityoccurrence.State) ([]vulnerabilityoccurrence.Occurrence, error)

func (*VulnerabilityOccurrenceStore) ListEvents added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListEvents(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityoccurrence.Event, error)

func (*VulnerabilityOccurrenceStore) ListUnreconciled added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListUnreconciled(ctx context.Context, tenantID, runID shared.ID, advisoryID string, after shared.ID, snapshotAt time.Time, limit int) (ports.VulnerabilityOccurrenceReconciliationPage, error)

func (*VulnerabilityOccurrenceStore) ListVulnerabilityOccurrences added in v0.1.8

func (*VulnerabilityOccurrenceStore) SummarizeVulnerabilityOccurrences added in v0.1.8

func (s *VulnerabilityOccurrenceStore) SummarizeVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryIDs []string, affectedAsset string, states []vulnerabilityoccurrence.State) (map[string]vulnerabilityintel.AdvisoryOccurrenceSummary, error)

func (*VulnerabilityOccurrenceStore) Upsert added in v0.1.8

type VulnerabilityReconcileRunStore added in v0.1.8

type VulnerabilityReconcileRunStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilityReconcileRunStore added in v0.1.8

func NewVulnerabilityReconcileRunStore(pool *pgxpool.Pool, ids ports.IDGenerator) *VulnerabilityReconcileRunStore

func (*VulnerabilityReconcileRunStore) Advance added in v0.1.8

func (s *VulnerabilityReconcileRunStore) Advance(ctx context.Context, id shared.ID, expectedCheckpoint, nextCheckpoint []byte, counts vulnerabilityreconcile.Counts, samples []string) (vulnerabilityreconcile.Run, error)

func (*VulnerabilityReconcileRunStore) Finish added in v0.1.8

func (*VulnerabilityReconcileRunStore) Get added in v0.1.8

func (*VulnerabilityReconcileRunStore) GetByDurableJobID added in v0.1.8

func (*VulnerabilityReconcileRunStore) HasReconciliationMatch added in v0.1.8

func (s *VulnerabilityReconcileRunStore) HasReconciliationMatch(ctx context.Context, tenantID, runID, engagementID shared.ID, advisoryID, componentFingerprint string) (bool, error)

func (*VulnerabilityReconcileRunStore) ListReconciliationDiffs added in v0.1.8

func (*VulnerabilityReconcileRunStore) MarkRunning added in v0.1.8

func (s *VulnerabilityReconcileRunStore) MarkRunning(ctx context.Context, id shared.ID) error

func (*VulnerabilityReconcileRunStore) RecordReconciliationDiff added in v0.1.8

func (s *VulnerabilityReconcileRunStore) RecordReconciliationDiff(ctx context.Context, diff vulnerabilityreconcile.Diff) (bool, error)

func (*VulnerabilityReconcileRunStore) Start added in v0.1.8

func (*VulnerabilityReconcileRunStore) SummarizeReconciliationDiffs added in v0.1.8

func (s *VulnerabilityReconcileRunStore) SummarizeReconciliationDiffs(ctx context.Context, tenantID, runID shared.ID) (vulnerabilityreconcile.Counts, error)

type VulnerabilityRetentionStore added in v0.2.0

type VulnerabilityRetentionStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilityRetentionStore added in v0.2.0

func NewVulnerabilityRetentionStore(pool *pgxpool.Pool) *VulnerabilityRetentionStore

func (*VulnerabilityRetentionStore) RunVulnerabilityRetention added in v0.2.0

func (s *VulnerabilityRetentionStore) RunVulnerabilityRetention(ctx context.Context, runID shared.ID, policy vulnerabilitymaintenance.Policy, dryRun bool, at time.Time) (vulnerabilitymaintenance.Run, error)

type VulnerabilityRiskAssessmentStore added in v0.1.8

type VulnerabilityRiskAssessmentStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilityRiskAssessmentStore added in v0.1.8

func NewVulnerabilityRiskAssessmentStore(pool *pgxpool.Pool) *VulnerabilityRiskAssessmentStore

func (*VulnerabilityRiskAssessmentStore) CountOpenHighCriticalVulnerabilityExposure added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) CountOpenHighCriticalVulnerabilityExposure(ctx context.Context, tenantID shared.ID) (int64, error)

func (*VulnerabilityRiskAssessmentStore) Current added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) Current(ctx context.Context, tenantID, occurrenceID shared.ID) (vulnerabilityrisk.Assessment, error)

func (*VulnerabilityRiskAssessmentStore) History added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) History(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityrisk.Assessment, error)

func (*VulnerabilityRiskAssessmentStore) ListVulnerabilityAssessments added in v0.1.8

func (*VulnerabilityRiskAssessmentStore) SummarizeVulnerabilityRisk added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) SummarizeVulnerabilityRisk(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryRiskSummary, error)

func (*VulnerabilityRiskAssessmentStore) Upsert added in v0.1.8

type VulnerabilitySourceStore added in v0.1.8

type VulnerabilitySourceStore struct {
	// contains filtered or unexported fields
}

func NewVulnerabilitySourceStore added in v0.1.8

func NewVulnerabilitySourceStore(pool *pgxpool.Pool) *VulnerabilitySourceStore

func (*VulnerabilitySourceStore) Archive added in v0.1.8

func (s *VulnerabilitySourceStore) Archive(ctx context.Context, id shared.ID, expectedVersion int) error

func (*VulnerabilitySourceStore) Create added in v0.1.8

func (*VulnerabilitySourceStore) Get added in v0.1.8

func (*VulnerabilitySourceStore) List added in v0.1.8

func (s *VulnerabilitySourceStore) List(ctx context.Context, includeArchived bool) ([]vulnerabilitysource.Source, error)

func (*VulnerabilitySourceStore) SetEnabled added in v0.1.8

func (s *VulnerabilitySourceStore) SetEnabled(ctx context.Context, id shared.ID, enabled bool, expectedVersion int) (vulnerabilitysource.Source, error)

func (*VulnerabilitySourceStore) Update added in v0.1.8

func (s *VulnerabilitySourceStore) Update(ctx context.Context, source vulnerabilitysource.Source, expectedVersion int) error

type WorkOrderRepository added in v0.1.8

type WorkOrderRepository struct {
	// contains filtered or unexported fields
}

WorkOrderRepository is the Postgres-backed fleet work order store. Every method runs through WithTenant so Row Level Security (migration 0059 via the 0057 procedure) isolates by tenant.

func NewWorkOrderRepository added in v0.1.8

func NewWorkOrderRepository(pool *pgxpool.Pool) *WorkOrderRepository

NewWorkOrderRepository constructs the Postgres work order repository.

func (*WorkOrderRepository) AcknowledgeFleetAudit added in v0.2.0

func (r *WorkOrderRepository) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*WorkOrderRepository) CancelForAgent added in v0.1.8

func (r *WorkOrderRepository) CancelForAgent(ctx context.Context, tenantID, agentID shared.ID, reason string, now time.Time) (int, error)

CancelForAgent cancels every live order addressed to agentID (used on agent revocation).

func (*WorkOrderRepository) CancelResponsesBelowGeneration added in v0.2.0

func (r *WorkOrderRepository) CancelResponsesBelowGeneration(ctx context.Context, tenantID shared.ID, generation int64, reason string, now time.Time) (int, error)

func (*WorkOrderRepository) CancelResponsesBelowGenerationWithAudit added in v0.2.0

func (r *WorkOrderRepository) CancelResponsesBelowGenerationWithAudit(ctx context.Context, tenantID shared.ID, generation int64, reason string, now time.Time, actor string) (int, ports.FleetAuditIntent, error)

func (*WorkOrderRepository) Claim added in v0.1.8

func (r *WorkOrderRepository) Claim(ctx context.Context, tenantID, agentID shared.ID, max int, now time.Time, leaseID string, leaseUntil time.Time) ([]*workorder.WorkOrder, error)

Claim atomically moves up to max unexpired issued orders addressed to agentID into claimed and returns them, using FOR UPDATE SKIP LOCKED so concurrent claimers never double-claim.

func (*WorkOrderRepository) ClaimWithAudit added in v0.2.0

func (r *WorkOrderRepository) ClaimWithAudit(ctx context.Context, tenantID, agentID shared.ID, max int, now time.Time, leaseID string, leaseUntil time.Time, actor string) ([]*workorder.WorkOrder, []ports.FleetAuditIntent, error)

func (*WorkOrderRepository) CompleteResponse added in v0.2.0

func (r *WorkOrderRepository) CompleteResponse(ctx context.Context, tenantID, id shared.ID, result fleetagent.ResponseExecutionResult, reason string, now time.Time) (bool, error)

func (*WorkOrderRepository) CompleteResponseWithAudit added in v0.2.0

func (r *WorkOrderRepository) CompleteResponseWithAudit(ctx context.Context, tenantID, id shared.ID, result fleetagent.ResponseExecutionResult, reason string, now time.Time, actor string) (bool, ports.FleetAuditIntent, error)

func (*WorkOrderRepository) GetByID added in v0.1.8

func (r *WorkOrderRepository) GetByID(ctx context.Context, tenantID, id shared.ID) (*workorder.WorkOrder, error)

GetByID returns the order for (tenantID, id) or shared.ErrNotFound.

func (*WorkOrderRepository) GetByIdempotencyKey added in v0.2.0

func (r *WorkOrderRepository) GetByIdempotencyKey(ctx context.Context, tenantID shared.ID, idempotencyKey string) (*workorder.WorkOrder, error)

func (*WorkOrderRepository) Issue added in v0.1.8

Issue inserts wo. It is idempotent by (tenant, idempotency key): a duplicate returns the existing order. A second LIVE order for the same (tenant, asset, capability, time bucket) returns shared.ErrConflict (the partial unique index).

func (*WorkOrderRepository) IssueWithAudit added in v0.2.0

func (*WorkOrderRepository) ListByTenant added in v0.1.8

func (r *WorkOrderRepository) ListByTenant(ctx context.Context, tenantID shared.ID) ([]*workorder.WorkOrder, error)

ListByTenant returns every work order for the tenant, ordered deterministically. Read-only, used by the coverage projection (#413); routed through WithTenant so RLS scopes it to the tenant.

func (*WorkOrderRepository) ListPendingFleetAudits added in v0.2.0

func (r *WorkOrderRepository) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

IssueWithAudit atomically persists a signed order and its normalized audit obligation.

func (*WorkOrderRepository) Transition added in v0.1.8

func (r *WorkOrderRepository) Transition(ctx context.Context, tenantID, id shared.ID, to workorder.State, reason string, expected workorder.State, now time.Time) error

Transition applies to with an optimistic expected-state check. It returns shared.ErrConflict when no row matched (the state changed concurrently or the order does not exist under this tenant).

func (*WorkOrderRepository) TransitionLeased added in v0.2.0

func (r *WorkOrderRepository) TransitionLeased(ctx context.Context, tenantID, id shared.ID, leaseID string, to workorder.State, reason string, expected workorder.State, now time.Time) error

func (*WorkOrderRepository) TransitionLeasedWithAudit added in v0.2.0

func (r *WorkOrderRepository) TransitionLeasedWithAudit(ctx context.Context, tenantID, id shared.ID, leaseID string, to workorder.State, reason string, expected workorder.State, now time.Time, actor string) (ports.FleetAuditIntent, error)

type WriteupDraftRepository

type WriteupDraftRepository struct {
	// contains filtered or unexported fields
}

WriteupDraftRepository persists AI-proposed, human-gated finding write-up drafts to PostgreSQL, engagement-scoped.

func NewWriteupDraftRepository

func NewWriteupDraftRepository(pool *pgxpool.Pool) *WriteupDraftRepository

NewWriteupDraftRepository returns a repository backed by the given pool.

func (*WriteupDraftRepository) Get

func (r *WriteupDraftRepository) Get(ctx context.Context, engagementID, id shared.ID) (out writeupdraft.Draft, err error)

Get returns the engagement's draft by id, or shared.ErrNotFound.

func (*WriteupDraftRepository) ListByEngagement

func (r *WriteupDraftRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []writeupdraft.Draft, err error)

ListByEngagement returns the engagement's drafts, oldest first (deterministic order).

func (*WriteupDraftRepository) Save

Save upserts a draft by id. Unlike a Judgment (insert-only), a Draft is mutable working data, so on conflict the mutable fields (text, state, decided_by, updated_at) are replaced; the immutable fields (engagement_id, finding_id, proposed_by, created_at) are never moved. tenant_id is written as the empty-string default tenant (mirrors judgments/findings); reads are tenant-isolated via the engagement gate, and the column is present so the P5/E22 row-scoping sweep covers it.

Source Files

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL