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 ¶
- Variables
- func AcquireSingletonLock(ctx context.Context, pool *pgxpool.Pool, role string) (*pgxpool.Conn, bool, error)
- func CheckDatabaseReady(ctx context.Context, pool *pgxpool.Pool) error
- func CheckMigrationsReady(ctx context.Context, pool *pgxpool.Pool) error
- func CheckRLSRuntimeRole(ctx context.Context, pool *pgxpool.Pool) error
- func Connect(ctx context.Context, dsn string) (*pgxpool.Pool, error)
- func ConnectPool(ctx context.Context, dsn string, pc PoolConfig) (*pgxpool.Pool, error)
- func GrantRuntimePrivileges(ctx context.Context, adminDSN, runtimeDSN string, haltWriterDSNs ...string) error
- func Migrate(ctx context.Context, dsn string) error
- func MigrateLocked(ctx context.Context, dsn string) error
- func ValidateMigrationRoleSeparation(migrationDSN, runtimeDSN string) error
- func ValidateResponseRoleSeparation(migrationDSN, runtimeDSN, haltWriterDSN string) error
- func WithContextTenant(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error
- func WithGlobalRead(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error
- func WithGlobalWrite(ctx context.Context, pool *pgxpool.Pool, fn func(pgx.Tx) error) error
- func WithTenant(ctx context.Context, pool *pgxpool.Pool, tenantID string, ...) (err error)
- type AITriageReviewRepository
- func (r *AITriageReviewRepository) Get(ctx context.Context, tenantID, id shared.ID) (aitriagereview.Review, error)
- func (r *AITriageReviewRepository) List(ctx context.Context, tenantID shared.ID, filter ports.AITriageReviewFilter) ([]aitriagereview.Review, error)
- func (r *AITriageReviewRepository) SaveDecision(ctx context.Context, review aitriagereview.Review, expectedVersion int) error
- func (r *AITriageReviewRepository) SaveOwner(ctx context.Context, review aitriagereview.Review, expectedVersion int) error
- func (r *AITriageReviewRepository) UpsertPending(ctx context.Context, review aitriagereview.Review) error
- type AUPStore
- type AccuracyRunRepository
- type AdvisoryMaterializer
- func (r *AdvisoryMaterializer) AdvisoryRevisionAt(ctx context.Context, advisoryID string, snapshotAt time.Time) (ports.AdvisoryRevisionRef, error)
- func (r *AdvisoryMaterializer) ByCPE(ctx context.Context, part, vendor, product string) ([]advisory.Advisory, error)
- func (r *AdvisoryMaterializer) ByPackage(ctx context.Context, ecosystem, packageName string) ([]advisory.Advisory, error)
- func (r *AdvisoryMaterializer) CountVulnerabilityAdvisoriesChangedSince(ctx context.Context, since time.Time) (int64, error)
- func (r *AdvisoryMaterializer) CountVulnerabilityAdvisoryDailyImpact(ctx context.Context, since time.Time) (vulnerabilityintel.AdvisoryDailyImpact, error)
- func (r *AdvisoryMaterializer) CurrentRevision(ctx context.Context, id string) (int64, error)
- func (r *AdvisoryMaterializer) CurrentSourceRecordIDs(ctx context.Context, sourceID string, yield func(string) error) error
- func (r *AdvisoryMaterializer) CurrentSourceRecordIDsBounded(ctx context.Context, sourceID string, limit int, yield func(string) error) error
- func (r *AdvisoryMaterializer) GetCanonical(ctx context.Context, id string) (advisory.Canonical, error)
- func (r *AdvisoryMaterializer) GetCanonicalAtRevision(ctx context.Context, id string, revision int64) (advisory.Canonical, error)
- func (r *AdvisoryMaterializer) ListAdvisoryRevisions(ctx context.Context, after string, snapshotAt time.Time, limit int) (ports.AdvisoryRevisionPage, error)
- func (r *AdvisoryMaterializer) ListVulnerabilityAdvisories(ctx context.Context, tenantID shared.ID, ...) (vulnerabilityintel.AdvisoryPage, error)
- func (r *AdvisoryMaterializer) ListVulnerabilityAdvisoryRevisions(ctx context.Context, query vulnerabilityintel.AdvisoryRevisionQuery) (vulnerabilityintel.AdvisoryRevisionPage, error)
- func (r *AdvisoryMaterializer) ListVulnerabilitySyncRunRevisions(ctx context.Context, runIDs []shared.ID, limitPerRun int) (map[shared.ID]vulnerabilityintel.AdvisoryRevisionLinkPage, error)
- func (r *AdvisoryMaterializer) MarkAdvisoryEvaluated(ctx context.Context, tenantID shared.ID, advisoryID string, revision int64, ...) error
- func (r *AdvisoryMaterializer) Materialize(ctx context.Context, records []advisory.ObservationRecord) (advisory.MaterializationResult, error)
- func (r *AdvisoryMaterializer) MaterializeSourceSnapshot(ctx context.Context, records []advisory.ObservationRecord) ([]advisory.MaterializationResult, error)
- func (r *AdvisoryMaterializer) OldestUnevaluatedAdvisory(ctx context.Context, tenantID shared.ID) (*vulnerabilityintel.EvaluationLag, error)
- func (r *AdvisoryMaterializer) PublishSourceSnapshot(ctx context.Context, publication ports.SourceSnapshotPublication, ...) ([]advisory.MaterializationResult, error)
- func (r *AdvisoryMaterializer) PublishedSourceSnapshot(ctx context.Context, syncRunID shared.ID) (ports.PublishedSourceSnapshot, bool, error)
- func (r *AdvisoryMaterializer) SummarizeVulnerabilityCoverage(ctx context.Context, tenantID shared.ID, ...) (map[string]vulnerabilityintel.AdvisoryCoverageSummary, error)
- func (*AdvisoryMaterializer) SupportsServerAdvisoryFilters() bool
- type AdvisoryRepository
- func (r *AdvisoryRepository) AdvisoryAliasEdges(ctx context.Context, ids []string) ([]advisory.AliasEdge, error)
- func (r *AdvisoryRepository) AdvisoryFreshness(ctx context.Context) (time.Time, int, error)
- func (r *AdvisoryRepository) ByCPE(ctx context.Context, part, vendor, product string) ([]advisory.Advisory, error)
- func (r *AdvisoryRepository) ByPackage(ctx context.Context, ecosystem, name string) ([]advisory.Advisory, error)
- func (r *AdvisoryRepository) CoveredEcosystems(ctx context.Context) (map[string]bool, error)
- func (r *AdvisoryRepository) Upsert(ctx context.Context, a advisory.Advisory) error
- type AgentDecisionStore
- type AgentPlanStore
- type AgentSessionStore
- func (s *AgentSessionStore) AppendMessage(ctx context.Context, sessionID shared.ID, seq int, m agent.Message) error
- func (s *AgentSessionStore) GetSession(ctx context.Context, id shared.ID) (agent.Session, error)
- func (s *AgentSessionStore) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]agent.Session, error)
- func (s *AgentSessionStore) ListResumable(ctx context.Context, staleFor time.Duration, now time.Time, limit int) ([]agent.Session, error)
- func (s *AgentSessionStore) Messages(ctx context.Context, sessionID shared.ID) ([]agent.Message, error)
- func (s *AgentSessionStore) SaveSession(ctx context.Context, e agent.Session) error
- type AgentSigningKeyRepository
- func (r *AgentSigningKeyRepository) ListByAgent(ctx context.Context, agentID shared.ID) ([]fleetagent.AgentSigningKey, error)
- func (r *AgentSigningKeyRepository) Register(ctx context.Context, key fleetagent.AgentSigningKey) error
- func (r *AgentSigningKeyRepository) ResolveSigningKey(ctx context.Context, agentID shared.ID, keyID string) (fleetagent.AgentSigningKey, error)
- func (r *AgentSigningKeyRepository) Revoke(ctx context.Context, agentID shared.ID, keyID string, at time.Time) error
- type ApprovalStore
- func (s *ApprovalStore) Consume(ctx context.Context, actionID shared.ID) error
- func (s *ApprovalStore) Decide(ctx context.Context, d agent.ApprovalDecision) error
- func (s *ApprovalStore) EngagementsWithPending(ctx context.Context) ([]ports.ApprovalSweepScope, error)
- func (s *ApprovalStore) Enqueue(ctx context.Context, a agent.ProposedAction) error
- func (s *ApprovalStore) Get(ctx context.Context, actionID shared.ID) (agent.ProposedAction, agent.ApprovalDecision, error)
- func (s *ApprovalStore) Pending(ctx context.Context, engagementID shared.ID) ([]agent.ProposedAction, error)
- type AssessmentComparisonRepository
- func (repository *AssessmentComparisonRepository) CreateQueued(ctx context.Context, comparison assessmentcomparison.Comparison) (assessmentcomparison.Comparison, bool, error)
- func (repository *AssessmentComparisonRepository) Get(ctx context.Context, tenantID, comparisonID shared.ID) (assessmentcomparison.Comparison, error)
- func (repository *AssessmentComparisonRepository) GetAssessmentComparisonBacklog(ctx context.Context, tenantID shared.ID) (ports.AssessmentComparisonBacklog, error)
- func (repository *AssessmentComparisonRepository) GetByInputHash(ctx context.Context, tenantID shared.ID, inputHash string) (assessmentcomparison.Comparison, error)
- func (repository *AssessmentComparisonRepository) GetItem(ctx context.Context, tenantID, comparisonID, itemID shared.ID) (assessmentcomparison.Item, error)
- func (repository *AssessmentComparisonRepository) GetMetadata(ctx context.Context, tenantID, comparisonID shared.ID) (assessmentcomparison.Comparison, error)
- func (repository *AssessmentComparisonRepository) ListFailedAssessmentComparisons(ctx context.Context, tenantID shared.ID, limit int) ([]assessmentcomparison.Comparison, error)
- func (repository *AssessmentComparisonRepository) ListItems(ctx context.Context, tenantID, comparisonID shared.ID, ...) (ports.AssessmentComparisonItemPage, error)
- func (repository *AssessmentComparisonRepository) ListMetadataByCycle(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentcomparison.Comparison, error)
- func (repository *AssessmentComparisonRepository) SummarizeItems(ctx context.Context, tenantID, comparisonID shared.ID, ...) (assessmentcomparison.Summary, error)
- func (repository *AssessmentComparisonRepository) UpdateCAS(ctx context.Context, comparison assessmentcomparison.Comparison, ...) error
- type AssessmentCycleBackfillRepository
- func (repository *AssessmentCycleBackfillRepository) AcquireAssessmentCycleBackfillRun(ctx context.Context, request ports.AssessmentCycleBackfillAcquireRequest) (run ports.AssessmentCycleBackfillRun, resumed bool, err error)
- func (repository *AssessmentCycleBackfillRepository) AdvanceAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentCycleBackfillRun, err error)
- func (repository *AssessmentCycleBackfillRepository) CommitAssessmentCycleBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, ...) (item ports.AssessmentCycleBackfillItem, created bool, err error)
- func (repository *AssessmentCycleBackfillRepository) FinishAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentCycleBackfillRun, err error)
- func (repository *AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillItem(ctx context.Context, tenantID, runID, assessmentID shared.ID) (item ports.AssessmentCycleBackfillItem, err error)
- func (repository *AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentCycleBackfillRun, err error)
- type AssessmentCycleIntegrityRepository
- func (repository *AssessmentCycleIntegrityRepository) AcquireAssessmentCycleIntegrityRun(ctx context.Context, request ports.AssessmentCycleIntegrityAcquireRequest) (run ports.AssessmentCycleIntegrityRun, resumed bool, err error)
- func (repository *AssessmentCycleIntegrityRepository) AdvanceAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentCycleIntegrityRun, err error)
- func (repository *AssessmentCycleIntegrityRepository) AssessmentCycleIntegrityGeneration(ctx context.Context, tenantID shared.ID) (generation int64, err error)
- func (repository *AssessmentCycleIntegrityRepository) CountAssessmentCycleIntegritySubjects(ctx context.Context, tenantID shared.ID, snapshotAt time.Time) (eligible int, memberships int, err error)
- func (repository *AssessmentCycleIntegrityRepository) FinishAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentCycleIntegrityRun, err error)
- func (repository *AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegrityRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentCycleIntegrityRun, err error)
- func (repository *AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegritySubject(ctx context.Context, tenantID, runID, assessmentID shared.ID) (result ports.AssessmentCycleIntegritySubjectResult, err error)
- func (repository *AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegrityFindings(ctx context.Context, tenantID, runID shared.ID) (findings []ports.AssessmentCycleIntegrityFinding, err error)
- func (repository *AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegritySubjects(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, ...) (subjects []ports.AssessmentCycleIntegritySubject, err error)
- func (repository *AssessmentCycleIntegrityRepository) SaveAssessmentCycleIntegritySubject(ctx context.Context, leaseToken shared.ID, now time.Time, ...) (created bool, err error)
- type AssessmentCycleRepository
- func (r *AssessmentCycleRepository) CommitClosure(ctx context.Context, commit ports.AssessmentClosureCommit) error
- func (r *AssessmentCycleRepository) CreateCycle(ctx context.Context, cycle *assessmentcycle.AssessmentCycle) error
- func (r *AssessmentCycleRepository) CreateMember(ctx context.Context, member *assessmentcycle.Member) error
- func (r *AssessmentCycleRepository) DeleteCycle(ctx context.Context, tenantID, cycleID shared.ID) error
- func (r *AssessmentCycleRepository) DeleteMember(ctx context.Context, tenantID, cycleID, assessmentID shared.ID) error
- func (r *AssessmentCycleRepository) GetActiveClosureManifest(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentclosure.Manifest, error)
- func (r *AssessmentCycleRepository) GetClosureManifest(ctx context.Context, tenantID, cycleID, manifestID shared.ID) (*assessmentclosure.Manifest, error)
- func (r *AssessmentCycleRepository) GetClosureReport(ctx context.Context, tenantID, cycleID, manifestID shared.ID, ...) (ports.AssessmentClosureReportArtifact, error)
- func (r *AssessmentCycleRepository) GetCycle(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentcycle.AssessmentCycle, error)
- func (r *AssessmentCycleRepository) GetCycleByAssessment(ctx context.Context, tenantID, assessmentID shared.ID) (*assessmentcycle.AssessmentCycle, error)
- func (r *AssessmentCycleRepository) GetMember(ctx context.Context, tenantID, cycleID, assessmentID shared.ID) (*assessmentcycle.Member, error)
- func (r *AssessmentCycleRepository) ListClosureManifests(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentclosure.Manifest, error)
- func (r *AssessmentCycleRepository) ListCycles(ctx context.Context, query ports.AssessmentCycleListQuery) ([]ports.AssessmentCycleListRecord, error)
- func (r *AssessmentCycleRepository) ListMembers(ctx context.Context, tenantID, cycleID shared.ID) ([]assessmentcycle.Member, error)
- func (r *AssessmentCycleRepository) ListMigrationPendingAssessments(ctx context.Context, query ports.AssessmentCycleListQuery) ([]ports.AssessmentCycleMigrationPendingRecord, int, error)
- func (r *AssessmentCycleRepository) LockCycleForUpdate(ctx context.Context, tenantID, cycleID shared.ID) (*assessmentcycle.AssessmentCycle, error)
- func (r *AssessmentCycleRepository) NextManifestVersion(ctx context.Context, tenantID, cycleID shared.ID) (int64, error)
- func (r *AssessmentCycleRepository) ReopenClosure(ctx context.Context, reopen ports.AssessmentClosureReopen) error
- func (r *AssessmentCycleRepository) SaveClosureReport(ctx context.Context, report ports.AssessmentClosureReportArtifact) (ports.AssessmentClosureReportArtifact, bool, error)
- func (r *AssessmentCycleRepository) UpdateCycleCAS(ctx context.Context, cycle *assessmentcycle.AssessmentCycle, ...) error
- func (r *AssessmentCycleRepository) UpdateMemberCAS(ctx context.Context, member *assessmentcycle.Member, expectedVersion int64) error
- type AssessmentCycleRequestRepository
- func (repository *AssessmentCycleRequestRepository) AbortAssessmentCycleRequest(ctx context.Context, scope ports.AssessmentCycleRequestScope, ...) error
- func (repository *AssessmentCycleRequestRepository) BeginAssessmentCycleRequest(ctx context.Context, request ports.AssessmentCycleRequest) (ports.AssessmentCycleRequest, bool, error)
- func (repository *AssessmentCycleRequestRepository) CompleteAssessmentCycleRequest(ctx context.Context, scope ports.AssessmentCycleRequestScope, ...) error
- type AssessmentRelationshipRepository
- func (repository *AssessmentRelationshipRepository) CreateCandidate(ctx context.Context, candidate assessmentrelationship.Candidate) (record assessmentrelationship.Record, created bool, err error)
- func (repository *AssessmentRelationshipRepository) DecideCandidateCAS(ctx context.Context, decision assessmentrelationship.Decision, ...) (record assessmentrelationship.Record, replayed bool, err error)
- func (repository *AssessmentRelationshipRepository) GetCandidate(ctx context.Context, tenantID, candidateID shared.ID) (record assessmentrelationship.Record, err error)
- func (repository *AssessmentRelationshipRepository) ListCandidates(ctx context.Context, tenantID shared.ID, ...) (records []assessmentrelationship.Record, err error)
- type AssessmentSnapshotBackfillRepository
- func (repository *AssessmentSnapshotBackfillRepository) AcquireAssessmentSnapshotBackfillRun(ctx context.Context, request ports.AssessmentSnapshotBackfillAcquireRequest) (run ports.AssessmentSnapshotBackfillRun, resumed bool, err error)
- func (repository *AssessmentSnapshotBackfillRepository) AdvanceAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentSnapshotBackfillRun, err error)
- func (repository *AssessmentSnapshotBackfillRepository) CommitAssessmentSnapshotBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, ...) (item ports.AssessmentSnapshotBackfillItem, created bool, err error)
- func (repository *AssessmentSnapshotBackfillRepository) FinishAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.AssessmentSnapshotBackfillRun, err error)
- func (repository *AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillItem(ctx context.Context, tenantID, runID, assessmentID shared.ID) (item ports.AssessmentSnapshotBackfillItem, err error)
- func (repository *AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.AssessmentSnapshotBackfillRun, err error)
- type AssessmentSnapshotRepository
- func (repository *AssessmentSnapshotRepository) CreateFinalizedCAS(ctx context.Context, snapshot *assessmentsnapshot.Snapshot, ...) (*assessmentsnapshot.Snapshot, bool, error)
- func (repository *AssessmentSnapshotRepository) CreateLegacyProjection(ctx context.Context, snapshot *assessmentsnapshot.Snapshot) (*assessmentsnapshot.Snapshot, bool, error)
- func (repository *AssessmentSnapshotRepository) Get(ctx context.Context, tenantID, snapshotID shared.ID) (*assessmentsnapshot.Snapshot, error)
- func (repository *AssessmentSnapshotRepository) GetByRequestKey(ctx context.Context, tenantID, assessmentID shared.ID, requestKey string) (*assessmentsnapshot.Snapshot, error)
- func (repository *AssessmentSnapshotRepository) GetDefault(ctx context.Context, tenantID, assessmentID shared.ID) (*assessmentsnapshot.Snapshot, ports.AssessmentSnapshotDefault, error)
- func (repository *AssessmentSnapshotRepository) ListAssessmentSnapshots(ctx context.Context, query ports.AssessmentSnapshotListQuery) (ports.AssessmentSnapshotPage, error)
- func (repository *AssessmentSnapshotRepository) ListByAssessment(ctx context.Context, tenantID, assessmentID shared.ID) ([]assessmentsnapshot.Snapshot, error)
- type AssetRepository
- func (r *AssetRepository) AssignEngagementBusinessAsset(ctx context.Context, tenantID, engagementID, assetID shared.ID) error
- func (r *AssetRepository) CountBusinessAssetsByCriticality(ctx context.Context, tenantID shared.ID) (map[asset.Criticality]int, error)
- func (r *AssetRepository) CreateBusinessAsset(ctx context.Context, a *asset.BusinessAsset) error
- func (r *AssetRepository) GetAssetByID(ctx context.Context, tenantID, id shared.ID) (*asset.Asset, error)
- func (r *AssetRepository) GetAssetByKey(ctx context.Context, tenantID shared.ID, kind asset.Kind, key string) (*asset.Asset, error)
- func (r *AssetRepository) GetBusinessAssetByID(ctx context.Context, tenantID, id shared.ID) (*asset.BusinessAsset, error)
- func (r *AssetRepository) GetBusinessAssetByKey(ctx context.Context, tenantID shared.ID, key string) (*asset.BusinessAsset, error)
- func (r *AssetRepository) ListAssets(ctx context.Context, tenantID shared.ID) ([]*asset.Asset, error)
- func (r *AssetRepository) ListBusinessAssetProjects(ctx context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)
- func (r *AssetRepository) ListBusinessAssetTechnicalAssets(ctx context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)
- func (r *AssetRepository) ListBusinessAssets(ctx context.Context, tenantID shared.ID) ([]*asset.BusinessAsset, error)
- func (r *AssetRepository) ListBusinessAssetsPage(ctx context.Context, tenantID shared.ID, query ports.BusinessAssetQuery) ([]*asset.BusinessAsset, int, error)
- func (r *AssetRepository) ListEdges(ctx context.Context, tenantID shared.ID) ([]*asset.Edge, error)
- func (r *AssetRepository) ListEngagementsByBusinessAsset(ctx context.Context, tenantID, assetID shared.ID) ([]*engagement.Engagement, error)
- func (r *AssetRepository) ReplaceBusinessAssetProjects(ctx context.Context, tenantID, assetID shared.ID, ...) error
- func (r *AssetRepository) ReplaceBusinessAssetTechnicalAssets(ctx context.Context, tenantID, assetID shared.ID, ...) error
- func (r *AssetRepository) UpdateBusinessAsset(ctx context.Context, a *asset.BusinessAsset, expectedVersion int) error
- func (r *AssetRepository) UpsertAsset(ctx context.Context, a *asset.Asset) error
- func (r *AssetRepository) UpsertEdge(ctx context.Context, e *asset.Edge) error
- type AttackPathStore
- type AuditLog
- func (l *AuditLog) List(ctx context.Context, limit int) (out []ports.AuditEntry, err error)
- func (l *AuditLog) MigrationMetadata(ctx context.Context) ([]ports.MigrationMetadata, error)
- func (l *AuditLog) Record(ctx context.Context, e ports.AuditEntry) error
- func (l *AuditLog) RecordOnce(ctx context.Context, e ports.AuditEntry) error
- func (l *AuditLog) Verify(ctx context.Context) (audit.Report, error)
- func (l *AuditLog) VerifyGlobal(ctx context.Context) (audit.Report, error)
- type BaselineRepository
- type CloudObservationStore
- type CloudRunStore
- func (s *CloudRunStore) EnqueueCloudRun(ctx context.Context, run cloudposture.Run, kind string, payload []byte) error
- func (s *CloudRunStore) GetCloudRun(ctx context.Context, tenantID, id shared.ID) (out cloudposture.Run, err error)
- func (s *CloudRunStore) SaveCloudRun(ctx context.Context, run cloudposture.Run) error
- type CommentRepository
- type ComponentInventoryStore
- func (s *ComponentInventoryStore) ClaimInventoryWork(ctx context.Context, tenantID shared.ID, owner string, at time.Time, ...) ([]sbom.InventoryWork, error)
- func (s *ComponentInventoryStore) CompleteInventoryPublication(ctx context.Context, publication sbom.InventoryPublication, at time.Time) error
- func (s *ComponentInventoryStore) FinishInventoryWork(ctx context.Context, work sbom.InventoryWork, owner string, ...) error
- func (s *ComponentInventoryStore) GetCurrentInventoryPublication(ctx context.Context, tenantID, engagementID shared.ID, scope string) (sbom.InventoryPublication, error)
- func (s *ComponentInventoryStore) ListCurrentComponents(ctx context.Context, query sbom.ComponentQuery) (sbom.ComponentPage, error)
- func (s *ComponentInventoryStore) ListCurrentComponentsByEngagement(ctx context.Context, tenantID, engagementID shared.ID) ([]sbom.ComponentRecord, error)
- func (s *ComponentInventoryStore) ListCurrentInventoryPublications(ctx context.Context, tenantID shared.ID, cursor sbom.InventoryCursor, ...) (sbom.InventoryPublicationPage, error)
- func (s *ComponentInventoryStore) ListSnapshotComponents(ctx context.Context, query sbom.SnapshotQuery) (sbom.ComponentPage, error)
- type CorrelationStateRepository
- func (r *CorrelationStateRepository) BeginCorrelationSnapshot(ctx context.Context, engagementID shared.ID, expected uint64, ...) (correlation.Checkpoint, error)
- func (r *CorrelationStateRepository) CommitCorrelationConsume(ctx context.Context, engagementID shared.ID, expected uint64, ...) error
- func (r *CorrelationStateRepository) ListStagedCorrelationSignals(ctx context.Context, engagementID shared.ID, ...) ([]correlation.Signal, bool, error)
- func (r *CorrelationStateRepository) LoadCorrelationState(ctx context.Context, engagementID shared.ID, ids []shared.ID, activeLimit int) (correlation.State, error)
- func (r *CorrelationStateRepository) StageCorrelationSignals(ctx context.Context, engagementID shared.ID, expected uint64, ...) error
- type CoverageWindowRepository
- func (r *CoverageWindowRepository) AppendCoverageWindow(ctx context.Context, window sensorstate.CoverageWindow) (sensorstate.CoverageWindow, error)
- func (r *CoverageWindowRepository) ListCoverageWindows(ctx context.Context, q ports.CoverageWindowQuery) ([]sensorstate.CoverageWindow, error)
- func (r *CoverageWindowRepository) ListCoverageWindowsBounded(ctx context.Context, q ports.CoverageWindowQuery, limit int) ([]sensorstate.CoverageWindow, error)
- type DASTRunStore
- func (s *DASTRunStore) EnqueueDASTRun(ctx context.Context, run dastrun.Run, kind string, payload []byte) error
- func (s *DASTRunStore) FinishRun(ctx context.Context, tenantID shared.ID, run dastrun.Run) (won bool, err error)
- func (s *DASTRunStore) GetDASTRun(ctx context.Context, tenantID, id shared.ID) (out dastrun.Run, err error)
- func (s *DASTRunStore) SaveDASTRun(ctx context.Context, run dastrun.Run) error
- func (s *DASTRunStore) StartRun(ctx context.Context, tenantID shared.ID, run dastrun.Run) (won bool, err error)
- type DetectionProvenanceRepository
- func (r *DetectionProvenanceRepository) AdmitPending(ctx context.Context, current detectionprovenance.Current, ...) error
- func (r *DetectionProvenanceRepository) AppendTransition(ctx context.Context, transition detectionprovenance.Transition) error
- func (r *DetectionProvenanceRepository) Current(ctx context.Context, engagementID, detectionID shared.ID) (detectionprovenance.Current, bool, error)
- func (r *DetectionProvenanceRepository) ListCurrent(ctx context.Context, engagementID shared.ID) ([]detectionprovenance.Current, error)
- func (r *DetectionProvenanceRepository) ListPending(ctx context.Context) ([]detectionprovenance.Current, error)
- func (r *DetectionProvenanceRepository) ListReceivedTransitions(ctx context.Context, engagementID shared.ID) ([]detectionprovenance.Transition, error)
- func (r *DetectionProvenanceRepository) ListTransitions(ctx context.Context, engagementID, detectionID shared.ID) ([]detectionprovenance.Transition, error)
- func (r *DetectionProvenanceRepository) LoadReceivedTransitions(ctx context.Context, engagementID shared.ID, detectionIDs []shared.ID) ([]detectionprovenance.Transition, error)
- type DetectionRecordRepository
- func (r *DetectionRecordRepository) AppendDetection(ctx context.Context, rec detection.Record) error
- func (r *DetectionRecordRepository) ClassCountsByAsset(ctx context.Context, assetID shared.ID, since time.Time) (map[detection.Class]int, error)
- func (r *DetectionRecordRepository) CorrelationHighWater(ctx context.Context, engagementID shared.ID, ...) (correlation.SourcePosition, bool, error)
- func (r *DetectionRecordRepository) DeleteDetection(ctx context.Context, engagementID, detectionID shared.ID) (bool, error)
- func (r *DetectionRecordRepository) HasDetection(ctx context.Context, engagementID, id shared.ID) (bool, error)
- func (r *DetectionRecordRepository) LastBatchSequence(ctx context.Context, agentID shared.ID) (uint64, error)
- func (r *DetectionRecordRepository) ListCorrelationSourcePage(ctx context.Context, engagementID shared.ID, ...) ([]detection.Record, bool, error)
- func (r *DetectionRecordRepository) ListDetections(ctx context.Context, engagementID shared.ID) ([]detection.Record, error)
- func (r *DetectionRecordRepository) ListExpiredDetections(ctx context.Context, engagementID shared.ID, cutoff time.Time) ([]shared.ID, error)
- type EmulationRunRepository
- type EndpointProcessRepository
- func (r *EndpointProcessRepository) ListRunningByAsset(ctx context.Context, assetID shared.ID) ([]ports.ProcessSnapshot, error)
- func (r *EndpointProcessRepository) ReplaceRunningProcesses(ctx context.Context, assetID shared.ID, snapshots []ports.ProcessSnapshot) error
- func (r *EndpointProcessRepository) SaveProcesses(ctx context.Context, snapshots []ports.ProcessSnapshot) error
- type EndpointTimelineRepository
- func (r *EndpointTimelineRepository) AppendTimeline(ctx context.Context, list []endpoint.TimelineEntry) error
- func (r *EndpointTimelineRepository) LoadTimelineEntries(ctx context.Context, assetID shared.ID, eventIDs []shared.ID) ([]endpoint.TimelineEntry, error)
- func (r *EndpointTimelineRepository) QueryTimeline(ctx context.Context, q ports.EndpointTimelineQuery) ([]endpoint.TimelineEntry, error)
- type EngagementRepository
- func (r *EngagementRepository) Create(ctx context.Context, e *engagement.Engagement) error
- func (r *EngagementRepository) Delete(ctx context.Context, id shared.ID) error
- func (r *EngagementRepository) GetByHostAssetID(ctx context.Context, tenantID, assetID shared.ID) (out *engagement.Engagement, err error)
- func (r *EngagementRepository) GetByID(ctx context.Context, id shared.ID) (out *engagement.Engagement, err error)
- func (r *EngagementRepository) GetByIDInTenant(ctx context.Context, tenantID, id shared.ID) (out *engagement.Engagement, err error)
- func (r *EngagementRepository) GetByProjectID(ctx context.Context, tenantID, projectID shared.ID) (out *engagement.Engagement, err error)
- func (r *EngagementRepository) List(ctx context.Context, tenantID shared.ID) (out []*engagement.Engagement, err error)
- func (r *EngagementRepository) ListAssessmentCycleBackfillEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, ...) (out []*engagement.Engagement, err error)
- func (r *EngagementRepository) ListAssessmentSnapshotBackfillEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, ...) ([]*engagement.Engagement, error)
- func (r *EngagementRepository) ListHostEngagements(ctx context.Context, tenantID shared.ID) ([]*engagement.Engagement, error)
- func (r *EngagementRepository) ListProjectEngagements(ctx context.Context, tenantID shared.ID) ([]*engagement.Engagement, error)
- func (r *EngagementRepository) ListPromotionReconciliationScopes(ctx context.Context) ([]ports.PromotionReconciliationScope, error)
- func (r *EngagementRepository) ListReconciliationEngagements(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, ...) (ports.ReconciliationEngagementPage, error)
- func (r *EngagementRepository) ListTenantIDs(ctx context.Context) ([]shared.ID, error)
- func (r *EngagementRepository) ProjectContexts(ctx context.Context, tenantID shared.ID, projectIDs []shared.ID) (out map[shared.ID]*engagement.Engagement, err error)
- func (r *EngagementRepository) Update(ctx context.Context, e *engagement.Engagement) error
- type EngagementSourceRepository
- func (r *EngagementSourceRepository) Create(ctx context.Context, item sourcepackage.Package) (sourcepackage.Package, bool, error)
- func (r *EngagementSourceRepository) Delete(ctx context.Context, tenantID, engagementID shared.ID) (sourcepackage.Package, bool, error)
- func (r *EngagementSourceRepository) Get(ctx context.Context, tenantID, engagementID shared.ID) (sourcepackage.Package, error)
- func (r *EngagementSourceRepository) GetByLocator(ctx context.Context, tenantID shared.ID, locator string) (sourcepackage.Package, error)
- func (r *EngagementSourceRepository) GetByVersion(ctx context.Context, tenantID, engagementID, versionID shared.ID) (sourcepackage.Package, error)
- func (r *EngagementSourceRepository) ObjectUnreferenced(ctx context.Context, tenantID shared.ID, objectKey string) (bool, error)
- type EvidenceStore
- func (r *EvidenceStore) Append(ctx context.Context, items []evidence.Evidence) error
- func (r *EvidenceStore) Head(ctx context.Context, engagementID shared.ID) (string, error)
- func (r *EvidenceStore) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []evidence.Evidence, err error)
- func (r *EvidenceStore) LookupSealedForFinding(ctx context.Context, engagementID, findingID shared.ID, kind string) (evidence.Evidence, bool, error)
- type ExploitationChainRepository
- type FindingLineageBackfillRepository
- func (repository *FindingLineageBackfillRepository) AcquireFindingLineageBackfillRun(ctx context.Context, request ports.FindingLineageBackfillAcquireRequest) (run ports.FindingLineageBackfillRun, resumed bool, err error)
- func (repository *FindingLineageBackfillRepository) AdvanceFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.FindingLineageBackfillRun, err error)
- func (repository *FindingLineageBackfillRepository) CommitFindingLineageBackfillItem(ctx context.Context, tenantID, runID, leaseToken shared.ID, now time.Time, ...) (item ports.FindingLineageBackfillItem, created bool, err error)
- func (repository *FindingLineageBackfillRepository) FinishFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID, leaseOwner string, ...) (run ports.FindingLineageBackfillRun, err error)
- func (repository *FindingLineageBackfillRepository) GetFindingLineageBackfillItem(ctx context.Context, tenantID, runID, sourceFindingID shared.ID) (item ports.FindingLineageBackfillItem, err error)
- func (repository *FindingLineageBackfillRepository) GetFindingLineageBackfillRun(ctx context.Context, tenantID, runID shared.ID) (run ports.FindingLineageBackfillRun, err error)
- func (repository *FindingLineageBackfillRepository) ListFindingLineageBackfillSources(ctx context.Context, tenantID, after shared.ID, snapshotAt time.Time, ...) ([]ports.FindingLineageBackfillSourceRow, error)
- type FindingLineageRepository
- func (repository *FindingLineageRepository) AppendAlias(ctx context.Context, alias findinglineage.Alias) (bool, error)
- func (repository *FindingLineageRepository) AppendObservation(ctx context.Context, observation findinglineage.Observation) error
- func (repository *FindingLineageRepository) AppendOverrideCAS(ctx context.Context, event findinglineage.OverrideEvent) (findinglineage.OverrideEvent, bool, error)
- func (repository *FindingLineageRepository) AppendSkip(ctx context.Context, record findinglineage.SkipRecord) (findinglineage.SkipRecord, bool, error)
- func (repository *FindingLineageRepository) CreateCandidate(ctx context.Context, candidate findinglineage.MatchCandidate, ...) (findinglineage.MatchCandidate, bool, error)
- func (repository *FindingLineageRepository) CreateIdentityWithObservation(ctx context.Context, identity findinglineage.Identity, ...) error
- func (repository *FindingLineageRepository) FindIdentitiesByAlias(ctx context.Context, tenantID, cycleID shared.ID, ...) ([]findinglineage.Identity, error)
- func (repository *FindingLineageRepository) FindIdentitiesByFingerprint(ctx context.Context, tenantID, cycleID shared.ID, ...) ([]findinglineage.Identity, error)
- func (repository *FindingLineageRepository) FindIdentitiesByProducerID(ctx context.Context, tenantID, cycleID shared.ID, ...) ([]findinglineage.Identity, error)
- func (repository *FindingLineageRepository) GetActiveOverride(ctx context.Context, tenantID, cycleID, sourceObservationID shared.ID) (findinglineage.OverrideEvent, error)
- func (repository *FindingLineageRepository) GetCandidate(ctx context.Context, tenantID, cycleID, candidateID shared.ID) (findinglineage.MatchCandidate, error)
- func (repository *FindingLineageRepository) GetIdentity(ctx context.Context, tenantID, cycleID, identityID shared.ID) (findinglineage.Identity, error)
- func (repository *FindingLineageRepository) GetObservation(ctx context.Context, tenantID, cycleID, observationID shared.ID) (findinglineage.Observation, error)
- func (repository *FindingLineageRepository) GetObservationBySource(ctx context.Context, tenantID, cycleID, snapshotID shared.ID, ...) (findinglineage.Observation, error)
- func (repository *FindingLineageRepository) ListActiveOverridesBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.OverrideEvent, error)
- func (repository *FindingLineageRepository) ListCandidateResolutions(ctx context.Context, tenantID, cycleID, candidateID shared.ID) ([]findinglineage.ResolutionEvent, error)
- func (repository *FindingLineageRepository) ListObservationsBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.Observation, error)
- func (repository *FindingLineageRepository) ListOpenCandidatesBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.MatchCandidate, error)
- func (repository *FindingLineageRepository) ListOverrideEvents(ctx context.Context, tenantID, cycleID, sourceObservationID shared.ID) ([]findinglineage.OverrideEvent, error)
- func (repository *FindingLineageRepository) ListSkipsBySnapshot(ctx context.Context, tenantID, cycleID, snapshotID shared.ID) ([]findinglineage.SkipRecord, error)
- func (repository *FindingLineageRepository) LockCorrelationNamespace(ctx context.Context, tenantID, cycleID shared.ID, ...) error
- func (repository *FindingLineageRepository) ResolveCandidateCAS(ctx context.Context, updated findinglineage.MatchCandidate, ...) (findinglineage.MatchCandidate, findinglineage.ResolutionEvent, bool, error)
- type FindingRepository
- func (r *FindingRepository) CheckVulnerabilityPrimaryFinding(ctx context.Context, tenantID, engagementID shared.ID, ...) error
- func (r *FindingRepository) ClaimFindingProjection(ctx context.Context, tenantID, engagementID, judgmentID shared.ID, ...) error
- func (r *FindingRepository) GetByEngagementAndDedupKey(ctx context.Context, engagementID shared.ID, dedupKey string) (finding.Finding, error)
- func (r *FindingRepository) GetByEngagementAndID(ctx context.Context, engagementID, findingID shared.ID) (finding.Finding, error)
- func (r *FindingRepository) LinkVulnerabilityFindingOccurrence(ctx context.Context, tenantID, engagementID, findingID, occurrenceID shared.ID, ...) error
- func (r *FindingRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []finding.Finding, err error)
- func (r *FindingRepository) ListPublishableByEngagement(ctx context.Context, engagementID shared.ID) ([]finding.Finding, error)
- func (r *FindingRepository) MapVulnerabilityPrimaryFinding(ctx context.Context, tenantID, engagementID shared.ID, ...) error
- func (r *FindingRepository) SetAssignee(ctx context.Context, engagementID, findingID shared.ID, assignee string, ...) (out finding.Finding, err error)
- func (r *FindingRepository) SetEvidenceScore(ctx context.Context, engagementID, findingID shared.ID, ...) (out finding.Finding, err error)
- func (r *FindingRepository) SummarizeOpenFindingsByEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.VulnerabilitySummary, error)
- func (r *FindingRepository) SummarizeVulnerabilitiesByEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.VulnerabilitySummary, error)
- func (r *FindingRepository) UpdateStatus(ctx context.Context, engagementID, findingID shared.ID, status finding.Status, ...) (out finding.Finding, err error)
- func (r *FindingRepository) Upsert(ctx context.Context, findings []finding.Finding) error
- type FleetAgentRepository
- func (r *FleetAgentRepository) ConsumeEnrolToken(ctx context.Context, tenantID shared.ID, hash string, now time.Time) (*fleetagent.EnrolToken, error)
- func (r *FleetAgentRepository) CreateAgent(ctx context.Context, a *fleetagent.Agent) error
- func (r *FleetAgentRepository) CreateEnrolToken(ctx context.Context, t *fleetagent.EnrolToken) error
- func (r *FleetAgentRepository) Decommission(ctx context.Context, tenantID, id shared.ID, now time.Time) error
- func (r *FleetAgentRepository) GetAgent(ctx context.Context, tenantID, id shared.ID) (*fleetagent.Agent, error)
- func (r *FleetAgentRepository) Heartbeat(ctx context.Context, tenantID, id shared.ID, ...) error
- func (r *FleetAgentRepository) ListAgents(ctx context.Context, tenantID shared.ID) ([]*fleetagent.Agent, error)
- func (r *FleetAgentRepository) Revoke(ctx context.Context, tenantID, id, by shared.ID, reason string, now time.Time) error
- func (r *FleetAgentRepository) SetFingerprint(ctx context.Context, tenantID, id shared.ID, fingerprint string, now time.Time) error
- type FleetAuditRepository
- type FleetDesiredRepository
- func (r *FleetDesiredRepository) Delete(ctx context.Context, tenantID, assetID, expectedPolicyID shared.ID, ...) error
- func (r *FleetDesiredRepository) Get(ctx context.Context, tenantID, assetID shared.ID) (*fleetdesired.State, error)
- func (r *FleetDesiredRepository) List(ctx context.Context, tenantID shared.ID) ([]*fleetdesired.State, error)
- func (r *FleetDesiredRepository) Put(ctx context.Context, state *fleetdesired.State) error
- type FleetRolloutRepository
- type IdentityStore
- func (s *IdentityStore) ConsumeAuthorizationTransaction(ctx context.Context, tenantID shared.ID, stateHash string, now time.Time) (transaction identity.AuthorizationTransaction, err error)
- func (s *IdentityStore) CreateAuthorizationTransaction(ctx context.Context, transaction identity.AuthorizationTransaction) error
- func (s *IdentityStore) CreateExternalIdentity(ctx context.Context, external identity.ExternalIdentity) error
- func (s *IdentityStore) CreateSession(ctx context.Context, session identity.Session) error
- func (s *IdentityStore) GetExternalIdentity(ctx context.Context, issuer, subject string) (external identity.ExternalIdentity, err error)
- func (s *IdentityStore) GetSessionByTokenHash(ctx context.Context, tokenHash string) (session identity.Session, err error)
- func (s *IdentityStore) RevokeSession(ctx context.Context, tenantID, sessionID shared.ID, now time.Time) error
- func (s *IdentityStore) RotateSession(ctx context.Context, previousSessionID shared.ID, replacement identity.Session, ...) error
- type ImportReceiptRepository
- func (r *ImportReceiptRepository) CreateOrGet(ctx context.Context, tenantID shared.ID, receipt importreceipt.Receipt) (importreceipt.Receipt, bool, error)
- func (r *ImportReceiptRepository) Finalize(ctx context.Context, tenantID, receiptID shared.ID, ...) (importreceipt.Receipt, error)
- func (r *ImportReceiptRepository) GetByIdentity(ctx context.Context, tenantID, engagementID shared.ID, ...) (importreceipt.Receipt, error)
- type ImportedFindingRepository
- func (r *ImportedFindingRepository) ExistsDigest(ctx context.Context, tenantID, engagementID shared.ID, digest string) (bool, error)
- func (r *ImportedFindingRepository) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID) ([]importedfinding.ImportedFinding, error)
- func (r *ImportedFindingRepository) Save(ctx context.Context, tenantID shared.ID, ...) (int, int, error)
- type ImportedSBOMStore
- func (s *ImportedSBOMStore) LatestByEngagement(ctx context.Context, tenantID, engagementID shared.ID) (importedsbom.Record, error)
- func (s *ImportedSBOMStore) MetadataByEngagements(ctx context.Context, tenantID shared.ID, engagementIDs []shared.ID) (map[shared.ID]importedsbom.Metadata, error)
- func (s *ImportedSBOMStore) SaveActive(ctx context.Context, record importedsbom.Record) error
- type IncidentEventRepository
- func (r *IncidentEventRepository) AppendEvents(ctx context.Context, incidentID shared.ID, expectedRevision int, ...) error
- func (r *IncidentEventRepository) ListIncidentIDs(ctx context.Context, q ports.IncidentQuery) ([]shared.ID, error)
- func (r *IncidentEventRepository) ListMergeEdges(ctx context.Context, canonicalID shared.ID) ([]incident.MergeEdge, error)
- func (r *IncidentEventRepository) ListPendingResponseLinks(ctx context.Context) ([]incident.ResponseLink, error)
- func (r *IncidentEventRepository) LoadEvents(ctx context.Context, incidentID shared.ID) ([]incident.IncidentEvent, error)
- func (r *IncidentEventRepository) ResolveCanonicalID(ctx context.Context, id shared.ID) (shared.ID, error)
- type IntegrationStore
- func (store *IntegrationStore) ArchiveIntegration(ctx context.Context, id shared.ID, expectedVersion int, audit ports.AuditEntry) error
- func (store *IntegrationStore) BeginIntegrationOperation(ctx context.Context, id shared.ID, startedAt time.Time) (operation integration.Operation, execute bool, err error)
- func (store *IntegrationStore) CancelIntegrationOperation(ctx context.Context, id shared.ID, finishedAt time.Time, ...) (operation integration.Operation, err error)
- func (store *IntegrationStore) CreateIntegration(ctx context.Context, item integration.Integration, audit ports.AuditEntry) error
- func (store *IntegrationStore) CreateIntegrationBinding(ctx context.Context, binding integration.Binding, audit ports.AuditEntry) error
- func (store *IntegrationStore) DeleteIntegrationBinding(ctx context.Context, integrationID, bindingID shared.ID, ...) error
- func (store *IntegrationStore) DeleteIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, ...) error
- func (store *IntegrationStore) FinishIntegrationOperation(ctx context.Context, id shared.ID, state integration.OperationState, ...) (operation integration.Operation, err error)
- func (store *IntegrationStore) FinishIntegrationPoll(ctx context.Context, id shared.ID, state integration.OperationState, ...) (operation integration.Operation, err error)
- func (store *IntegrationStore) GetIntegration(ctx context.Context, id shared.ID) (item integration.Integration, err error)
- func (store *IntegrationStore) GetIntegrationOperation(ctx context.Context, id shared.ID) (operation integration.Operation, err error)
- func (store *IntegrationStore) IntegrationCredentialConfigured(ctx context.Context, integrationID shared.ID, credentialID string) (configured bool, err error)
- func (store *IntegrationStore) ListDueIntegrations(ctx context.Context, now time.Time, limit int) (items []integration.Integration, err error)
- func (store *IntegrationStore) ListIntegrationBindings(ctx context.Context, integrationID shared.ID) (bindings []integration.Binding, err error)
- func (store *IntegrationStore) ListIntegrationExternalRuns(ctx context.Context, integrationID shared.ID, limit int) (runs []integration.ExternalRun, err error)
- func (store *IntegrationStore) ListIntegrationOperations(ctx context.Context, integrationID shared.ID, limit int) (operations []integration.Operation, err error)
- func (store *IntegrationStore) ListIntegrations(ctx context.Context, includeArchived bool) (items []integration.Integration, err error)
- func (store *IntegrationStore) PutIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, ...) error
- func (store *IntegrationStore) ResolveIntegrationCredential(ctx context.Context, integrationID shared.ID, credentialID string, ...) (plaintext []byte, err error)
- func (store *IntegrationStore) SetIntegrationEnabled(ctx context.Context, id shared.ID, enabled bool, expectedVersion int, ...) (updated integration.Integration, err error)
- func (store *IntegrationStore) StartIntegrationOperation(ctx context.Context, operation integration.Operation, jobKind string, ...) (integration.Operation, error)
- func (store *IntegrationStore) UpdateIntegration(ctx context.Context, item integration.Integration, expectedVersion int, ...) (updated integration.Integration, err error)
- type JobQueue
- func (q *JobQueue) AggregateJobQueueStats(ctx context.Context, kinds ...string) (ports.JobStats, error)
- func (q *JobQueue) Claim(ctx context.Context, visibility time.Duration, kinds ...string) (*ports.QueuedJob, error)
- func (q *JobQueue) Complete(ctx context.Context, id string, fence int64) error
- func (q *JobQueue) Deadletter(ctx context.Context, id string, fence int64) error
- func (q *JobQueue) Depth(ctx context.Context, kinds ...string) (count int, err error)
- func (q *JobQueue) Enqueue(ctx context.Context, kind string, payload []byte) (string, error)
- func (q *JobQueue) Fail(ctx context.Context, id string, fence int64, retryIn time.Duration) error
- func (q *JobQueue) Heartbeat(ctx context.Context, id string, fence int64, extend time.Duration) error
- func (q *JobQueue) JobStatus(ctx context.Context, id string) (status ports.JobStatus, err error)
- func (q *JobQueue) Retry(ctx context.Context, id string, fence int64, retryIn time.Duration) error
- func (q *JobQueue) Stats(ctx context.Context, kinds ...string) (stats ports.JobStats, err error)
- type JudgmentRepository
- func (r *JudgmentRepository) AcknowledgeJudgmentAudit(ctx context.Context, kind ports.JudgmentAuditKind, judgmentID shared.ID, ...) error
- func (r *JudgmentRepository) GetByID(ctx context.Context, engagementID, id shared.ID) (out judgment.Judgment, err error)
- func (r *JudgmentRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]judgment.Judgment, error)
- func (r *JudgmentRepository) ListBySubject(ctx context.Context, engagementID, subjectID shared.ID) ([]judgment.Judgment, error)
- func (r *JudgmentRepository) ListPendingJudgmentAudits(ctx context.Context, engagementID shared.ID) (out []ports.PendingJudgmentAudit, err error)
- func (r *JudgmentRepository) Save(ctx context.Context, j judgment.Judgment) error
- func (r *JudgmentRepository) SaveWithProposalAudit(ctx context.Context, j judgment.Judgment, entry ports.AuditEntry) error
- func (r *JudgmentRepository) SetScoreState(ctx context.Context, engagementID, id shared.ID, score int, ...) (judgment.Judgment, error)
- func (r *JudgmentRepository) SetVerdictState(ctx context.Context, engagementID, id shared.ID, score int, ...) (judgment.Judgment, error)
- func (r *JudgmentRepository) SetVerdictStateWithAudit(ctx context.Context, engagementID, id shared.ID, score int, ...) (out judgment.Judgment, err error)
- type LeaderStore
- type LeaseRunLock
- type LegalHoldRepository
- func (r *LegalHoldRepository) IsHeld(ctx context.Context, engagementID shared.ID) (bool, error)
- func (r *LegalHoldRepository) ListActive(ctx context.Context) ([]legalhold.Hold, error)
- func (r *LegalHoldRepository) Place(ctx context.Context, h legalhold.Hold) (legalhold.Hold, error)
- func (r *LegalHoldRepository) Release(ctx context.Context, engagementID shared.ID, releasedBy string, at time.Time) error
- type MaterializingAdvisoryWriter
- type NotificationRepository
- func (r *NotificationRepository) BeginAttempt(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, ...) (notification.Attempt, error)
- func (r *NotificationRepository) CancelDelivery(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, ...) error
- func (r *NotificationRepository) CreateChannel(ctx context.Context, c notification.Channel, sealed string) (notification.Channel, error)
- func (r *NotificationRepository) CreateRule(ctx context.Context, rule notification.Rule) (notification.Rule, error)
- func (r *NotificationRepository) DeadLetterDelivery(ctx context.Context, tenant, did shared.ID, reason string) error
- func (r *NotificationRepository) DeleteChannel(ctx context.Context, tenant, id shared.ID, revision int, at time.Time) error
- func (r *NotificationRepository) DeleteRule(ctx context.Context, tenant, id shared.ID, revision int) error
- func (r *NotificationRepository) DeliveryStillRelevant(ctx context.Context, work ports.NotificationWork) (bool, error)
- func (r *NotificationRepository) FinishAttempt(ctx context.Context, tenant, did shared.ID, jobID string, fence int64, ...) error
- func (r *NotificationRepository) GetChannel(ctx context.Context, tenant, id shared.ID) (notification.Channel, error)
- func (r *NotificationRepository) GetDelivery(ctx context.Context, tenant, id shared.ID) (notification.Delivery, error)
- func (r *NotificationRepository) GetRule(ctx context.Context, tenant, id shared.ID) (notification.Rule, error)
- func (r *NotificationRepository) ListAttempts(ctx context.Context, tenant, did shared.ID) ([]notification.Attempt, error)
- func (r *NotificationRepository) ListChannels(ctx context.Context, tenant shared.ID) ([]notification.Channel, error)
- func (r *NotificationRepository) ListDeliveries(ctx context.Context, f ports.NotificationDeliveryFilter) (notification.Page, error)
- func (r *NotificationRepository) ListRules(ctx context.Context, tenant shared.ID) ([]notification.Rule, error)
- func (r *NotificationRepository) LoadWork(ctx context.Context, tenant, did shared.ID) (ports.NotificationWork, error)
- func (r *NotificationRepository) Publish(ctx context.Context, e notification.Event) ([]shared.ID, error)
- func (r *NotificationRepository) PublishToChannel(ctx context.Context, e notification.Event, cid shared.ID) (shared.ID, error)
- func (r *NotificationRepository) UpdateChannel(ctx context.Context, c notification.Channel, sealed string, replace bool) (notification.Channel, error)
- func (r *NotificationRepository) UpdateRule(ctx context.Context, rule notification.Rule) (notification.Rule, error)
- type NotificationSource
- type OwnershipExecution
- func (r *OwnershipExecution) CommitOwnershipWork(ctx context.Context, job ports.QueuedJob, work ports.OwnershipWork, ...) error
- func (r *OwnershipExecution) DispatchOwnership(ctx context.Context, mode string, limit int) (int, error)
- func (r *OwnershipExecution) FailOwnershipWork(ctx context.Context, job ports.QueuedJob) error
- func (r *OwnershipExecution) LoadOwnershipWork(ctx context.Context, job ports.QueuedJob, limit int) (out []ports.OwnershipWork, err error)
- func (r *OwnershipExecution) ReplayOwnershipRun(ctx context.Context, actor, id shared.ID, revision int) error
- func (r *OwnershipExecution) StartOwnershipRun(ctx context.Context, req ports.OwnershipRunRequest) (run ports.OwnershipRun, err error)
- type OwnershipRepository
- func (r *OwnershipRepository) ActivatePolicy(ctx context.Context, a ports.OwnershipActivation) error
- func (r *OwnershipRepository) AddMember(ctx context.Context, team, id shared.ID, at time.Time) error
- func (r *OwnershipRepository) AppendIntent(ctx context.Context, i ports.OwnershipIntent) error
- func (r *OwnershipRepository) ApplyAssignment(ctx context.Context, m ports.OwnershipMutation) (out ownership.Decision, err error)
- func (r *OwnershipRepository) AuthorizeOwnership(ctx context.Context, actor shared.ID, permission user.Permission) error
- func (r *OwnershipRepository) CompleteIntent(ctx context.Context, id shared.ID) error
- func (r *OwnershipRepository) CreatePolicyVersion(ctx context.Context, p ownership.PolicyVersion) error
- func (r *OwnershipRepository) CreateRun(ctx context.Context, run ports.OwnershipRun) error
- func (r *OwnershipRepository) CreateSnapshot(ctx context.Context, s ownership.Snapshot) error
- func (r *OwnershipRepository) CreateTeam(ctx context.Context, t ownership.Team) (ownership.Team, error)
- func (r *OwnershipRepository) DeleteAssetMapping(ctx context.Context, asset shared.ID, revision int) error
- func (r *OwnershipRepository) DeleteMapping(ctx context.Context, eng shared.ID, repo, owner string, revision int) error
- func (r *OwnershipRepository) GetActivePolicy(ctx context.Context, eng shared.ID, repo string) (out ports.OwnershipPolicy, err error)
- func (r *OwnershipRepository) GetAssignment(ctx context.Context, eng, finding shared.ID) (out ports.OwnershipCurrent, err error)
- func (r *OwnershipRepository) GetOwnershipPolicy(ctx context.Context, id shared.ID) (out ports.OwnershipPolicyHeader, err error)
- func (r *OwnershipRepository) GetPolicyVersion(ctx context.Context, id shared.ID, version int) (p ownership.PolicyVersion, err error)
- func (r *OwnershipRepository) GetRun(ctx context.Context, id shared.ID) (run ports.OwnershipRun, err error)
- func (r *OwnershipRepository) GetSnapshot(ctx context.Context, id shared.ID) (s ownership.Snapshot, err error)
- func (r *OwnershipRepository) GetTeam(ctx context.Context, id shared.ID) (t ownership.Team, err error)
- func (r *OwnershipRepository) ListAssetMappings(ctx context.Context, after shared.ID, limit int) (out []ports.OwnershipAssetMapping, err error)
- func (r *OwnershipRepository) ListDecisions(ctx context.Context, eng, finding shared.ID, ...) (out []ownership.Decision, err error)
- func (r *OwnershipRepository) ListMappings(ctx context.Context, eng shared.ID, repo, after string, limit int) (out []ports.OwnershipMapping, err error)
- func (r *OwnershipRepository) ListMembers(ctx context.Context, team, after shared.ID, limit int) (out []ownership.Membership, err error)
- func (r *OwnershipRepository) ListOwnershipPolicies(ctx context.Context, eng, after shared.ID, limit int) (out []ports.OwnershipPolicyHeader, err error)
- func (r *OwnershipRepository) ListOwnershipSnapshots(ctx context.Context, eng, after shared.ID, limit int) (out []ownership.Snapshot, err error)
- func (r *OwnershipRepository) ListPendingIntents(ctx context.Context, kind string, limit int) (out []ports.OwnershipIntent, err error)
- func (r *OwnershipRepository) ListRunItems(ctx context.Context, id, after shared.ID, limit int) (out []ports.OwnershipRunItem, err error)
- func (r *OwnershipRepository) ListTeams(ctx context.Context, after shared.ID, limit int) (out []ownership.Team, err error)
- func (r *OwnershipRepository) MarkOwnershipSourceReady(ctx context.Context, eng, id shared.ID) error
- func (r *OwnershipRepository) OwnershipInbox(ctx context.Context, f ports.OwnershipInboxFilter) (out ports.OwnershipInboxPage, err error)
- func (r *OwnershipRepository) RemoveMember(ctx context.Context, team, id shared.ID) error
- func (r *OwnershipRepository) ReserveOwnershipBulk(ctx context.Context, actor shared.ID, key, hash string, at time.Time) error
- func (r *OwnershipRepository) SaveAssetMapping(ctx context.Context, m ports.OwnershipAssetMapping) error
- func (r *OwnershipRepository) SaveMapping(ctx context.Context, eng shared.ID, record ports.OwnershipMapping) error
- func (r *OwnershipRepository) SaveOwnershipSource(ctx context.Context, source ports.OwnershipSourceRecord) error
- func (r *OwnershipRepository) SaveRunItems(ctx context.Context, id shared.ID, revision int, ...) error
- func (r *OwnershipRepository) SetRunState(ctx context.Context, id shared.ID, revision int, state string) error
- func (r *OwnershipRepository) UpdateTeam(ctx context.Context, t ownership.Team, expected int) (ownership.Team, error)
- func (r *OwnershipRepository) VisibleOwnershipEngagement(ctx context.Context, id shared.ID) error
- type PoolConfig
- type PoolStatsSource
- type PrivacyPolicyRepository
- func (r *PrivacyPolicyRepository) ActivatePrivacyPolicy(ctx context.Context, activation privacy.Activation) (privacy.Activation, error)
- func (r *PrivacyPolicyRepository) ActivatePrivacyPolicyWithAudit(ctx context.Context, activation privacy.Activation, ...) (privacy.Activation, ports.FleetAuditIntent, error)
- func (r *PrivacyPolicyRepository) ActivePrivacyPolicy(ctx context.Context, tenantID shared.ID) (privacy.Assignment, error)
- func (r *PrivacyPolicyRepository) PrivacyPolicyActivationHistory(ctx context.Context, tenantID shared.ID) ([]privacy.Activation, error)
- func (r *PrivacyPolicyRepository) PrivacyPolicyByDigest(ctx context.Context, tenantID shared.ID, digest string) (privacy.Assignment, error)
- func (r *PrivacyPolicyRepository) PrivacyPolicyHistory(ctx context.Context, tenantID shared.ID) ([]privacy.Assignment, error)
- func (r *PrivacyPolicyRepository) PutPrivacyPolicy(ctx context.Context, assignment privacy.Assignment) (bool, error)
- type ProjectAnalysisStore
- func (r *ProjectAnalysisStore) AttachSourceWithAudit(ctx context.Context, tenantID, projectID, analysisID shared.ID, ...) error
- func (r *ProjectAnalysisStore) Branches(ctx context.Context, tenantID, projectID shared.ID) ([]string, error)
- func (r *ProjectAnalysisStore) CurrentAnalysisHotspotSummary(ctx context.Context, tenantID, projectID, analysisID shared.ID, ...) (hotspot.Summary, error)
- func (r *ProjectAnalysisStore) CurrentFindingStatuses(ctx context.Context, tenantID, projectID shared.ID, keys []string) (map[string]string, error)
- func (r *ProjectAnalysisStore) Get(ctx context.Context, tenantID, projectID, analysisID shared.ID) (projectanalysis.Analysis, error)
- func (r *ProjectAnalysisStore) GetHotspot(ctx context.Context, tenantID, projectID, hotspotID shared.ID) (hotspot.Hotspot, error)
- func (r *ProjectAnalysisStore) GetIssue(ctx context.Context, tenantID, projectID, issueID shared.ID) (issue.Issue, error)
- func (r *ProjectAnalysisStore) HotspotHistory(ctx context.Context, tenantID, projectID, hotspotID shared.ID) ([]hotspot.ReviewEvent, error)
- func (r *ProjectAnalysisStore) IssueHistory(ctx context.Context, tenantID, projectID, issueID shared.ID) ([]issue.ReviewEvent, error)
- func (r *ProjectAnalysisStore) LatestForProjects(ctx context.Context, tenantID shared.ID, projectIDs []shared.ID) (map[shared.ID]projectanalysis.Analysis, error)
- func (r *ProjectAnalysisStore) LatestWithResult(ctx context.Context, tenantID, projectID shared.ID, branch string) (projectanalysis.Analysis, []byte, error)
- func (r *ProjectAnalysisStore) List(ctx context.Context, tenantID, projectID shared.ID, branch string, limit int, ...) ([]projectanalysis.Analysis, bool, error)
- func (r *ProjectAnalysisStore) ListAnalysisHotspots(ctx context.Context, tenantID, projectID, analysisID shared.ID, ...) (hotspot.Page, hotspot.Summary, error)
- func (r *ProjectAnalysisStore) ListHotspots(ctx context.Context, tenantID, projectID shared.ID, filter hotspot.ListFilter) (hotspot.Page, error)
- func (r *ProjectAnalysisStore) ListIssues(ctx context.Context, tenantID, projectID shared.ID, filter issue.ListFilter) (issue.Page, error)
- func (r *ProjectAnalysisStore) MatchIntegrationAnalysis(ctx context.Context, projectID shared.ID, revision string) (analysisID shared.ID, state integration.CorrelationState, err error)
- func (r *ProjectAnalysisStore) PruneBranchAnalyses(ctx context.Context, tenantID, projectID shared.ID, branch string, keep int) (int, error)
- func (r *ProjectAnalysisStore) ResolvedIssueKeys(ctx context.Context, tenantID, projectID shared.ID) (map[string]bool, error)
- func (r *ProjectAnalysisStore) Save(ctx context.Context, analysis projectanalysis.Analysis) error
- func (r *ProjectAnalysisStore) SaveWithResult(ctx context.Context, analysis projectanalysis.Analysis, result []byte) error
- func (r *ProjectAnalysisStore) SaveWithResultAndHotspots(ctx context.Context, analysis projectanalysis.Analysis, result []byte, ...) error
- func (r *ProjectAnalysisStore) SaveWithResultAndProjections(ctx context.Context, analysis projectanalysis.Analysis, result []byte, ...) error
- func (r *ProjectAnalysisStore) TransitionHotspot(ctx context.Context, cmd hotspot.TransitionCommand) (hotspot.Hotspot, hotspot.ReviewEvent, error)
- func (r *ProjectAnalysisStore) TransitionIssue(ctx context.Context, cmd issue.TransitionCommand) (issue.Issue, issue.ReviewEvent, error)
- type ProjectRepository
- func (r *ProjectRepository) AssignProfile(ctx context.Context, tenantID shared.ID, ...) error
- func (r *ProjectRepository) CountByGate(ctx context.Context, tenantID shared.ID, gateID string) (int, error)
- func (r *ProjectRepository) Create(ctx context.Context, p *project.Project) error
- func (r *ProjectRepository) DeleteByKey(ctx context.Context, tenantID shared.ID, key string) error
- func (r *ProjectRepository) GetByID(ctx context.Context, tenantID, projectID shared.ID) (*project.Project, error)
- func (r *ProjectRepository) GetByKey(ctx context.Context, tenantID shared.ID, key string) (*project.Project, error)
- func (r *ProjectRepository) List(ctx context.Context, tenantID shared.ID) ([]*project.Project, error)
- func (r *ProjectRepository) SetPullRequestDecoration(ctx context.Context, tenantID shared.ID, key string, enabled bool) error
- func (r *ProjectRepository) UpdateGate(ctx context.Context, tenantID shared.ID, key, gateID string) error
- type PromotionStore
- func (r *PromotionStore) Apply(ctx context.Context, engagementID, findingID shared.ID, ...) (out finding.Finding, err error)
- func (r *PromotionStore) FindByJudgment(ctx context.Context, engagementID, findingID, judgmentID shared.ID) (evt promotion.PromotionEvent, ok bool, err error)
- func (r *PromotionStore) LatestByFinding(ctx context.Context, engagementID, findingID shared.ID) (evt promotion.PromotionEvent, ok bool, err error)
- func (r *PromotionStore) ListByFinding(ctx context.Context, engagementID, findingID shared.ID) (out []promotion.PromotionEvent, err error)
- func (r *PromotionStore) ListPendingAudits(ctx context.Context, engagementID shared.ID) (out []promotion.PromotionEvent, err error)
- func (r *PromotionStore) MarkAuditComplete(ctx context.Context, eventID shared.ID) error
- type PurpleRepository
- func (r *PurpleRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]pcdom.Coverage, error)
- func (r *PurpleRepository) ListByRun(ctx context.Context, runID shared.ID) ([]pcdom.Coverage, error)
- func (r *PurpleRepository) SaveCoverage(ctx context.Context, records []pcdom.Coverage) error
- type QualityGateMutator
- func (m *QualityGateMutator) AssignProjectGate(ctx context.Context, tenantID shared.ID, projectKey, gateID string, ...) error
- func (m *QualityGateMutator) CreateGate(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate, ...) error
- func (m *QualityGateMutator) CreateProjectWithGate(ctx context.Context, p *project.Project) error
- func (m *QualityGateMutator) DeleteGate(ctx context.Context, tenantID shared.ID, key string, audit ports.AuditEntry) error
- func (m *QualityGateMutator) UpdateGate(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate, ...) error
- type QualityGateStore
- func (s *QualityGateStore) Create(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate) error
- func (s *QualityGateStore) Delete(ctx context.Context, tenantID shared.ID, key string) error
- func (s *QualityGateStore) DeleteIfUnassigned(ctx context.Context, tenantID shared.ID, key string) error
- func (s *QualityGateStore) Get(ctx context.Context, tenantID shared.ID, key string) (qualitygate.Gate, error)
- func (s *QualityGateStore) List(ctx context.Context, tenantID shared.ID) ([]qualitygate.Gate, error)
- func (s *QualityGateStore) Update(ctx context.Context, tenantID shared.ID, gate qualitygate.Gate) error
- type QualityProfileStore
- func (s *QualityProfileStore) Create(ctx context.Context, tenantID shared.ID, profile qualityprofile.Profile) error
- func (s *QualityProfileStore) Delete(ctx context.Context, tenantID shared.ID, key string) error
- func (s *QualityProfileStore) Get(ctx context.Context, tenantID shared.ID, key string) (qualityprofile.Profile, error)
- func (s *QualityProfileStore) List(ctx context.Context, tenantID shared.ID) ([]qualityprofile.Profile, error)
- func (s *QualityProfileStore) Update(ctx context.Context, tenantID shared.ID, profile qualityprofile.Profile) error
- type ReconRunStore
- func (r *ReconRunStore) Get(ctx context.Context, id shared.ID) (recon.Run, error)
- func (r *ReconRunStore) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]recon.Run, error)
- func (r *ReconRunStore) ListStaleRunning(ctx context.Context, olderThan time.Time, limit int) ([]recon.Run, error)
- func (r *ReconRunStore) Save(ctx context.Context, run recon.Run) error
- type ResponseHaltWriterRepository
- func (r *ResponseHaltWriterRepository) AdvanceHaltGenerationWithAudit(ctx context.Context, expected int64, intent ports.ResponseAuditIntent) (int64, ports.ResponseAuditIntent, ports.ResponseHaltDispatch, error)
- func (r *ResponseHaltWriterRepository) CurrentHaltGeneration(ctx context.Context) (int64, error)
- type ResponseObserverBindingRepository
- func (r *ResponseObserverBindingRepository) GetResponseObserverBinding(ctx context.Context, agentID shared.ID) (fleetagent.ResponseObserverBinding, error)
- func (r *ResponseObserverBindingRepository) ListResponseObserverBindings(ctx context.Context, assetID shared.ID) ([]fleetagent.ResponseObserverBinding, error)
- func (r *ResponseObserverBindingRepository) SaveResponseObserverBindingWithAudit(ctx context.Context, binding fleetagent.ResponseObserverBinding, ...) (fleetagent.ResponseObserverBinding, ports.FleetAuditIntent, error)
- type ResponseRepository
- func (r *ResponseRepository) AcknowledgeResponseAudit(ctx context.Context, id string) error
- func (r *ResponseRepository) AcknowledgeResponseHaltDispatch(ctx context.Context, generation int64) error
- func (r *ResponseRepository) AdvanceHaltGenerationWithAudit(context.Context, int64, ports.ResponseAuditIntent) (int64, ports.ResponseAuditIntent, ports.ResponseHaltDispatch, error)
- func (r *ResponseRepository) AttemptStillCurrent(ctx context.Context, idempotencyKey string, state responsesaga.SagaState, ...) (bool, error)
- func (r *ResponseRepository) ClaimAttempt(ctx context.Context, idempotencyKey string, from, to responsesaga.SagaState, ...) (responsesaga.ResponseAttempt, bool, error)
- func (r *ResponseRepository) CurrentHaltGeneration(ctx context.Context) (int64, error)
- func (r *ResponseRepository) EnqueueResponseAudit(ctx context.Context, intent ports.ResponseAuditIntent) (ports.ResponseAuditIntent, error)
- func (r *ResponseRepository) Get(ctx context.Context, id shared.ID) (rdom.Record, bool, error)
- func (r *ResponseRepository) GetAttempt(ctx context.Context, idempotencyKey string) (responsesaga.ResponseAttempt, bool, error)
- func (r *ResponseRepository) ListAttemptsByState(ctx context.Context, states ...responsesaga.SagaState) ([]responsesaga.ResponseAttempt, error)
- func (r *ResponseRepository) ListByState(ctx context.Context, state rdom.State) ([]rdom.Record, error)
- func (r *ResponseRepository) ListPendingResponseAudits(ctx context.Context) ([]ports.ResponseAuditIntent, error)
- func (r *ResponseRepository) ListPendingResponseHaltDispatches(ctx context.Context) ([]ports.ResponseHaltDispatch, error)
- func (r *ResponseRepository) Put(ctx context.Context, rec rdom.Record) error
- func (r *ResponseRepository) StartAttempt(ctx context.Context, a responsesaga.ResponseAttempt) (responsesaga.ResponseAttempt, bool, error)
- func (r *ResponseRepository) Transition(ctx context.Context, rec rdom.Record, from rdom.State) (bool, error)
- func (r *ResponseRepository) TransitionAttempt(ctx context.Context, a responsesaga.ResponseAttempt, ...) (responsesaga.ResponseAttempt, bool, error)
- func (r *ResponseRepository) TransitionAttemptWithAudit(ctx context.Context, attempt responsesaga.ResponseAttempt, ...) (responsesaga.ResponseAttempt, bool, ports.ResponseAuditIntent, error)
- func (r *ResponseRepository) TransitionWithAudit(ctx context.Context, rec rdom.Record, from rdom.State, ...) (bool, ports.ResponseAuditIntent, error)
- type ResponseVerificationRepository
- func (r *ResponseVerificationRepository) AppendResponseTargetEvidenceReceipt(ctx context.Context, receipt fleetagent.ResponseTargetEvidenceReceipt) (fleetagent.ResponseTargetEvidenceReceipt, error)
- func (r *ResponseVerificationRepository) AppendResponseVerificationWithAudit(ctx context.Context, observation ports.AcceptedResponseVerification, ...) (ports.FleetAuditIntent, error)
- func (r *ResponseVerificationRepository) GetResponseTargetEvidenceReceipt(ctx context.Context, attemptKey string) (fleetagent.ResponseTargetEvidenceReceipt, bool, error)
- func (r *ResponseVerificationRepository) GetResponseVerification(ctx context.Context, attemptKey string) (ports.AcceptedResponseVerification, bool, error)
- type RestoreEvidenceChain
- type RestoreEvidenceStore
- type RetestRepository
- func (r *RetestRepository) Add(ctx context.Context, rt finding.Retest) error
- func (r *RetestRepository) LatestByEngagementFindings(ctx context.Context, engagementID shared.ID, findingIDs []shared.ID) (map[shared.ID]finding.Retest, error)
- func (r *RetestRepository) ListByEngagementFinding(ctx context.Context, engagementID, findingID shared.ID) (out []finding.Retest, err error)
- func (r *RetestRepository) RetestHistories(ctx context.Context, engagementID shared.ID, findingIDs []shared.ID) (map[shared.ID][]finding.Retest, error)
- type RunLock
- type SCMConnectorRepository
- func (r *SCMConnectorRepository) Delete(ctx context.Context, id shared.ID) error
- func (r *SCMConnectorRepository) Get(ctx context.Context, id shared.ID) (ports.SCMConnectorMeta, error)
- func (r *SCMConnectorRepository) List(ctx context.Context) ([]ports.SCMConnectorMeta, error)
- func (r *SCMConnectorRepository) Put(ctx context.Context, c scmconnector.Connector, token []byte) error
- func (r *SCMConnectorRepository) ResolveGitCredential(ctx context.Context, host string) (ports.GitCredential, bool, error)
- type SLAStore
- func (s *SLAStore) ActivePolicy(ctx context.Context, tenantID shared.ID) (sla.Policy, error)
- func (s *SLAStore) AssessmentHistory(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.Assessment, error)
- func (s *SLAStore) Current(ctx context.Context, tenantID, engagementID, findingID shared.ID) (sla.Current, error)
- func (s *SLAStore) LifecycleEvents(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.LifecycleEvent, error)
- func (s *SLAStore) ListCurrent(ctx context.Context, tenantID, engagementID shared.ID) ([]sla.Current, error)
- func (s *SLAStore) PolicyHistory(ctx context.Context, tenantID shared.ID) ([]sla.Policy, error)
- func (s *SLAStore) PutPolicy(ctx context.Context, policy sla.Policy, activate bool) (bool, error)
- func (s *SLAStore) SLAHistories(ctx context.Context, tenantID, engagementID shared.ID, findingIDs []shared.ID) (ports.SLAHistoryBatch, error)
- func (s *SLAStore) SaveTransition(ctx context.Context, next sla.Lifecycle, event sla.LifecycleEvent) error
- func (s *SLAStore) UpsertAssessment(ctx context.Context, assessment sla.Assessment) (sla.AssessmentUpsertResult, error)
- type ScanJobStore
- func (r *ScanJobStore) CreateRunning(ctx context.Context, j ports.ScanJob) error
- func (r *ScanJobStore) GetJob(ctx context.Context, id string) (ports.ScanJob, error)
- func (r *ScanJobStore) LatestForEngagement(ctx context.Context, engagementID shared.ID) (ports.ScanJob, error)
- func (r *ScanJobStore) LatestForEngagements(ctx context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.ScanJob, error)
- func (r *ScanJobStore) ListStaleRunning(ctx context.Context, olderThan time.Time, limit int) ([]ports.ScanJob, error)
- func (r *ScanJobStore) Save(ctx context.Context, j ports.ScanJob) error
- type ScanRepository
- type ScanResultStore
- type ScanRunStore
- func (r *ScanRunStore) Get(ctx context.Context, runID string) (ports.ScanRun, error)
- func (r *ScanRunStore) GetScanRun(ctx context.Context, tenantID shared.ID, runID string) (scanrun.ScanRun, error)
- func (r *ScanRunStore) GetScanRunEvidence(ctx context.Context, tenantID shared.ID, runID string) (ports.ScanRunEvidence, error)
- func (r *ScanRunStore) List(ctx context.Context, engagementID shared.ID) ([]ports.ScanRun, error)
- func (r *ScanRunStore) ListScanRuns(ctx context.Context, tenantID, engagementID shared.ID) ([]scanrun.ScanRun, error)
- func (r *ScanRunStore) Save(ctx context.Context, run ports.ScanRun) error
- func (r *ScanRunStore) SaveScanRun(ctx context.Context, run scanrun.ScanRun) error
- func (r *ScanRunStore) SaveScanRunEvidence(ctx context.Context, item ports.ScanRunEvidence) error
- func (r *ScanRunStore) SealScanRun(ctx context.Context, command ports.SealScanRunCommand) error
- type ScannedImageStore
- type SensorStateRepository
- func (r *SensorStateRepository) AppendSensorState(ctx context.Context, observation sensorstate.Observation) error
- func (r *SensorStateRepository) AppendSensorStateWithAudit(ctx context.Context, observation sensorstate.Observation, ...) (ports.FleetAuditIntent, error)
- func (r *SensorStateRepository) ListCoverageSensorStates(ctx context.Context, q ports.CoverageSensorStateQuery) ([]sensorstate.Observation, error)
- func (r *SensorStateRepository) ListSensorStates(ctx context.Context, q ports.SensorStateQuery) ([]sensorstate.Observation, error)
- type SyncRunStore
- func (s *SyncRunStore) Advance(ctx context.Context, id shared.ID, expectedCheckpoint, nextCheckpoint []byte, ...) (vulnerabilitysync.Run, error)
- func (s *SyncRunStore) Finish(ctx context.Context, id shared.ID, state vulnerabilitysync.State, ...) (vulnerabilitysync.Run, error)
- func (s *SyncRunStore) Get(ctx context.Context, id shared.ID) (vulnerabilitysync.Run, error)
- func (s *SyncRunStore) GetByDurableJobID(ctx context.Context, jobID string) (vulnerabilitysync.Run, error)
- func (s *SyncRunStore) GetVulnerabilitySyncRun(ctx context.Context, tenantID, id shared.ID) (vulnerabilityintel.SyncRunItem, error)
- func (s *SyncRunStore) LatestForSource(ctx context.Context, sourceID shared.ID, states []vulnerabilitysync.State) (vulnerabilitysync.Run, error)
- func (s *SyncRunStore) LatestSuccessfulVulnerabilitySync(ctx context.Context, tenantID shared.ID) (*time.Time, error)
- func (s *SyncRunStore) ListStale(ctx context.Context, olderThan time.Time, limit int) ([]vulnerabilitysync.Run, error)
- func (s *SyncRunStore) ListVulnerabilitySyncRuns(ctx context.Context, query vulnerabilityintel.SyncRunQuery) (vulnerabilityintel.SyncRunPage, error)
- func (s *SyncRunStore) MarkRunning(ctx context.Context, id shared.ID) error
- func (s *SyncRunStore) RecoverStale(ctx context.Context, staleRunID shared.ID, staleBefore time.Time, ...) (vulnerabilitysync.Run, bool, error)
- func (s *SyncRunStore) Start(ctx context.Context, request ports.SyncRunStart) (vulnerabilitysync.Run, bool, error)
- func (s *SyncRunStore) Supersede(ctx context.Context, id shared.ID) error
- type TelemetryRepository
- func (r *TelemetryRepository) Footprint(ctx context.Context) (ports.TelemetryFootprint, error)
- func (r *TelemetryRepository) Ingest(ctx context.Context, batch ports.TelemetryBatch) error
- func (r *TelemetryRepository) LastSequence(ctx context.Context, hostID shared.ID, class detection.Class) (uint64, error)
- func (r *TelemetryRepository) Query(ctx context.Context, q ports.HuntQuery) (ports.HuntResult, error)
- func (r *TelemetryRepository) RecordLoss(ctx context.Context, loss ports.TelemetryLoss) error
- func (r *TelemetryRepository) RetentionSweep(ctx context.Context, now time.Time) (ports.SweepReport, error)
- type TelemetryTransportRepository
- func (r *TelemetryTransportRepository) AcceptAgentGapRevision(ctx context.Context, revision ports.TelemetryAgentGapRevision) error
- func (r *TelemetryTransportRepository) AcceptAgentGapRevisionWithAudit(ctx context.Context, revision ports.TelemetryAgentGapRevision, ...) (ports.FleetAuditIntent, error)
- func (r *TelemetryTransportRepository) AgentGapRevisions(ctx context.Context, gapID shared.ID) ([]ports.TelemetryAgentGapRevision, error)
- func (r *TelemetryTransportRepository) BindTelemetryAsset(ctx context.Context, binding ports.TelemetryAssetBinding) error
- func (r *TelemetryTransportRepository) CommitBatch(ctx context.Context, batch ports.TelemetryEventBatch) error
- func (r *TelemetryTransportRepository) CommitBatchWithAudit(ctx context.Context, batch ports.TelemetryEventBatch, ...) (ports.FleetAuditIntent, error)
- func (r *TelemetryTransportRepository) CountBatchEvents(ctx context.Context, agentID, streamID shared.ID, epoch, sequence uint64) (int, error)
- func (r *TelemetryTransportRepository) IngestBatchEvents(ctx context.Context, batch ports.TelemetryEventBatch) (int, error)
- func (r *TelemetryTransportRepository) ListCoverageGapFacts(ctx context.Context, q ports.CoverageGapQuery) ([]ports.CoverageGapFact, error)
- func (r *TelemetryTransportRepository) ListGapChanges(ctx context.Context, agentID, streamID shared.ID, epoch, sequence uint64) ([]ports.TelemetryGap, error)
- func (r *TelemetryTransportRepository) ListGaps(ctx context.Context, agentID, streamID shared.ID) ([]ports.TelemetryGap, error)
- func (r *TelemetryTransportRepository) ListTelemetryAssetBindings(ctx context.Context) ([]ports.TelemetryAssetBinding, error)
- func (r *TelemetryTransportRepository) MaxEpoch(ctx context.Context, agentID, streamID shared.ID) (uint64, error)
- func (r *TelemetryTransportRepository) QueryAgentGaps(ctx context.Context, q ports.TelemetryGapQuery) ([]ports.TelemetryGap, error)
- func (r *TelemetryTransportRepository) QueryDeliveryGaps(ctx context.Context, q ports.TelemetryGapQuery) ([]ports.TelemetryGap, error)
- func (r *TelemetryTransportRepository) QueryTelemetryBatchAccounting(ctx context.Context, q ports.TelemetryBatchAccountingQuery) ([]ports.TelemetryBatchAccounting, error)
- func (r *TelemetryTransportRepository) RecordAgentGap(ctx context.Context, gap ports.TelemetryAgentGap) error
- func (r *TelemetryTransportRepository) ResolveTelemetryAsset(ctx context.Context, agentID shared.ID) (shared.ID, error)
- func (r *TelemetryTransportRepository) ResolveTelemetryReferences(ctx context.Context, agentID, assetID shared.ID, redactionPolicyDigest string, ...) (ports.TelemetryReferenceStatus, error)
- func (r *TelemetryTransportRepository) SaveStreamState(ctx context.Context, state ports.TelemetryStreamState) error
- func (r *TelemetryTransportRepository) StreamState(ctx context.Context, agentID, streamID shared.ID, epoch uint64) (ports.TelemetryStreamState, error)
- type TenantTransactionRunner
- type ThreatModelRepository
- type TimestampStore
- func (s *TimestampStore) Get(ctx context.Context, chain string, eng shared.ID, head string) (*ports.TimestampToken, error)
- func (s *TimestampStore) LatestHead(ctx context.Context, chain string, eng shared.ID) (string, bool, error)
- func (s *TimestampStore) Put(ctx context.Context, chain string, eng shared.ID, head string, ...) error
- type UserRepository
- func (r *UserRepository) Bootstrap(ctx context.Context, u *user.User, auditEntry ports.AuditEntry) error
- func (r *UserRepository) Create(ctx context.Context, u *user.User) error
- func (r *UserRepository) GetByAPIKeyHash(ctx context.Context, hash string) (*user.User, error)
- func (r *UserRepository) GetByID(ctx context.Context, tenantID, id shared.ID) (*user.User, error)
- func (r *UserRepository) List(ctx context.Context, tenantID shared.ID) ([]*user.User, error)
- func (r *UserRepository) ListForUpdate(ctx context.Context, tenantID shared.ID) ([]*user.User, error)
- func (r *UserRepository) Update(ctx context.Context, tenantID shared.ID, u *user.User) error
- func (r *UserRepository) Upsert(ctx context.Context, u *user.User) error
- type VEXStatementRepository
- type VulnerabilityActionStore
- func (s *VulnerabilityActionStore) AcknowledgeAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)
- func (s *VulnerabilityActionStore) ClaimOutbox(ctx context.Context, tenantID shared.ID, now, lockedUntil time.Time, limit int) ([]vulnerabilityaction.OutboxEvent, error)
- func (s *VulnerabilityActionStore) CompleteOutbox(ctx context.Context, tenantID, eventID shared.ID, at time.Time) error
- func (s *VulnerabilityActionStore) CountPendingVulnerabilityActions(ctx context.Context, tenantID shared.ID) (int64, error)
- func (s *VulnerabilityActionStore) GetAction(ctx context.Context, tenantID, actionID shared.ID) (vulnerabilityaction.Action, error)
- func (s *VulnerabilityActionStore) ListActions(ctx context.Context, query vulnerabilityaction.ActionQuery) (vulnerabilityaction.ActionPage, error)
- func (s *VulnerabilityActionStore) ListVulnerabilityTransitions(ctx context.Context, query vulnerabilityintel.TransitionQuery) (vulnerabilityintel.TransitionPage, error)
- func (s *VulnerabilityActionStore) RecordChange(ctx context.Context, change vulnerabilityaction.Change) (bool, error)
- func (s *VulnerabilityActionStore) ResolveAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)
- func (s *VulnerabilityActionStore) RetryOutbox(ctx context.Context, tenantID, eventID shared.ID, at, availableAt time.Time, ...) error
- func (s *VulnerabilityActionStore) SummarizeVulnerabilityActions(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryActionSummary, error)
- type VulnerabilityOccurrenceStore
- func (s *VulnerabilityOccurrenceStore) CountActiveVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryID string) (int64, error)
- func (s *VulnerabilityOccurrenceStore) CountNewlyAffectedAssets(ctx context.Context, tenantID shared.ID, since time.Time) (int64, error)
- func (s *VulnerabilityOccurrenceStore) Get(ctx context.Context, tenantID, engagementID shared.ID, ...) (vulnerabilityoccurrence.Occurrence, error)
- func (s *VulnerabilityOccurrenceStore) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID, ...) ([]vulnerabilityoccurrence.Occurrence, error)
- func (s *VulnerabilityOccurrenceStore) ListEvents(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityoccurrence.Event, error)
- func (s *VulnerabilityOccurrenceStore) ListUnreconciled(ctx context.Context, tenantID, runID shared.ID, advisoryID string, ...) (ports.VulnerabilityOccurrenceReconciliationPage, error)
- func (s *VulnerabilityOccurrenceStore) ListVulnerabilityOccurrences(ctx context.Context, query vulnerabilityintel.OccurrenceQuery) (vulnerabilityintel.OccurrencePage, error)
- func (s *VulnerabilityOccurrenceStore) SummarizeVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryIDs []string, ...) (map[string]vulnerabilityintel.AdvisoryOccurrenceSummary, error)
- func (s *VulnerabilityOccurrenceStore) Upsert(ctx context.Context, occurrence vulnerabilityoccurrence.Occurrence) (vulnerabilityoccurrence.UpsertResult, error)
- type VulnerabilityReconcileRunStore
- func (s *VulnerabilityReconcileRunStore) Advance(ctx context.Context, id shared.ID, expectedCheckpoint, nextCheckpoint []byte, ...) (vulnerabilityreconcile.Run, error)
- func (s *VulnerabilityReconcileRunStore) Finish(ctx context.Context, id shared.ID, state vulnerabilityreconcile.State, ...) (vulnerabilityreconcile.Run, error)
- func (s *VulnerabilityReconcileRunStore) Get(ctx context.Context, id shared.ID) (vulnerabilityreconcile.Run, error)
- func (s *VulnerabilityReconcileRunStore) GetByDurableJobID(ctx context.Context, jobID string) (vulnerabilityreconcile.Run, error)
- func (s *VulnerabilityReconcileRunStore) HasReconciliationMatch(ctx context.Context, tenantID, runID, engagementID shared.ID, ...) (bool, error)
- func (s *VulnerabilityReconcileRunStore) ListReconciliationDiffs(ctx context.Context, query vulnerabilityreconcile.DiffQuery) (vulnerabilityreconcile.DiffPage, error)
- func (s *VulnerabilityReconcileRunStore) MarkRunning(ctx context.Context, id shared.ID) error
- func (s *VulnerabilityReconcileRunStore) RecordReconciliationDiff(ctx context.Context, diff vulnerabilityreconcile.Diff) (bool, error)
- func (s *VulnerabilityReconcileRunStore) Start(ctx context.Context, request ports.VulnerabilityReconcileStart) (vulnerabilityreconcile.Run, bool, error)
- func (s *VulnerabilityReconcileRunStore) SummarizeReconciliationDiffs(ctx context.Context, tenantID, runID shared.ID) (vulnerabilityreconcile.Counts, error)
- type VulnerabilityRetentionStore
- type VulnerabilityRiskAssessmentStore
- func (s *VulnerabilityRiskAssessmentStore) CountOpenHighCriticalVulnerabilityExposure(ctx context.Context, tenantID shared.ID) (int64, error)
- func (s *VulnerabilityRiskAssessmentStore) Current(ctx context.Context, tenantID, occurrenceID shared.ID) (vulnerabilityrisk.Assessment, error)
- func (s *VulnerabilityRiskAssessmentStore) History(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityrisk.Assessment, error)
- func (s *VulnerabilityRiskAssessmentStore) ListVulnerabilityAssessments(ctx context.Context, query vulnerabilityintel.AssessmentQuery) (vulnerabilityintel.AssessmentPage, error)
- func (s *VulnerabilityRiskAssessmentStore) SummarizeVulnerabilityRisk(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryRiskSummary, error)
- func (s *VulnerabilityRiskAssessmentStore) Upsert(ctx context.Context, assessment vulnerabilityrisk.Assessment) (vulnerabilityrisk.Result, error)
- type VulnerabilitySourceStore
- func (s *VulnerabilitySourceStore) Archive(ctx context.Context, id shared.ID, expectedVersion int) error
- func (s *VulnerabilitySourceStore) Create(ctx context.Context, source vulnerabilitysource.Source) error
- func (s *VulnerabilitySourceStore) Get(ctx context.Context, id shared.ID) (vulnerabilitysource.Source, error)
- func (s *VulnerabilitySourceStore) List(ctx context.Context, includeArchived bool) ([]vulnerabilitysource.Source, error)
- func (s *VulnerabilitySourceStore) SetEnabled(ctx context.Context, id shared.ID, enabled bool, expectedVersion int) (vulnerabilitysource.Source, error)
- func (s *VulnerabilitySourceStore) Update(ctx context.Context, source vulnerabilitysource.Source, expectedVersion int) error
- type WorkOrderRepository
- func (r *WorkOrderRepository) AcknowledgeFleetAudit(ctx context.Context, id string) error
- func (r *WorkOrderRepository) CancelForAgent(ctx context.Context, tenantID, agentID shared.ID, reason string, now time.Time) (int, error)
- func (r *WorkOrderRepository) CancelResponsesBelowGeneration(ctx context.Context, tenantID shared.ID, generation int64, reason string, ...) (int, error)
- func (r *WorkOrderRepository) CancelResponsesBelowGenerationWithAudit(ctx context.Context, tenantID shared.ID, generation int64, reason string, ...) (int, ports.FleetAuditIntent, error)
- func (r *WorkOrderRepository) Claim(ctx context.Context, tenantID, agentID shared.ID, max int, now time.Time, ...) ([]*workorder.WorkOrder, error)
- func (r *WorkOrderRepository) ClaimWithAudit(ctx context.Context, tenantID, agentID shared.ID, max int, now time.Time, ...) ([]*workorder.WorkOrder, []ports.FleetAuditIntent, error)
- func (r *WorkOrderRepository) CompleteResponse(ctx context.Context, tenantID, id shared.ID, ...) (bool, error)
- func (r *WorkOrderRepository) CompleteResponseWithAudit(ctx context.Context, tenantID, id shared.ID, ...) (bool, ports.FleetAuditIntent, error)
- func (r *WorkOrderRepository) GetByID(ctx context.Context, tenantID, id shared.ID) (*workorder.WorkOrder, error)
- func (r *WorkOrderRepository) GetByIdempotencyKey(ctx context.Context, tenantID shared.ID, idempotencyKey string) (*workorder.WorkOrder, error)
- func (r *WorkOrderRepository) Issue(ctx context.Context, wo *workorder.WorkOrder) (*workorder.WorkOrder, error)
- func (r *WorkOrderRepository) IssueWithAudit(ctx context.Context, wo *workorder.WorkOrder, intent ports.FleetAuditIntent) (*workorder.WorkOrder, ports.FleetAuditIntent, error)
- func (r *WorkOrderRepository) ListByTenant(ctx context.Context, tenantID shared.ID) ([]*workorder.WorkOrder, error)
- func (r *WorkOrderRepository) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)
- func (r *WorkOrderRepository) Transition(ctx context.Context, tenantID, id shared.ID, to workorder.State, reason string, ...) error
- func (r *WorkOrderRepository) TransitionLeased(ctx context.Context, tenantID, id shared.ID, leaseID string, ...) error
- func (r *WorkOrderRepository) TransitionLeasedWithAudit(ctx context.Context, tenantID, id shared.ID, leaseID string, ...) (ports.FleetAuditIntent, error)
- type WriteupDraftRepository
- func (r *WriteupDraftRepository) Get(ctx context.Context, engagementID, id shared.ID) (out writeupdraft.Draft, err error)
- func (r *WriteupDraftRepository) ListByEngagement(ctx context.Context, engagementID shared.ID) (out []writeupdraft.Draft, err error)
- func (r *WriteupDraftRepository) Save(ctx context.Context, d writeupdraft.Draft) error
Constants ¶
This section is empty.
Variables ¶
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
CheckDatabaseReady verifies the runtime pool can execute a trivial query.
func CheckMigrationsReady ¶ added in v0.2.0
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
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 ConnectPool ¶
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 ¶
Migrate applies all pending goose migrations (idempotent; tracked in goose_db_version).
func MigrateLocked ¶ added in v0.2.0
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
ValidateMigrationRoleSeparation ensures migrations cannot run as the runtime role.
func ValidateResponseRoleSeparation ¶ added in v0.2.0
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
WithContextTenant runs fn under the immutable tenant previously bound to ctx.
func WithGlobalRead ¶ added in v0.1.8
func WithGlobalWrite ¶ added in v0.1.8
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 (r *AITriageReviewRepository) Get(ctx context.Context, tenantID, id shared.ID) (aitriagereview.Review, error)
func (*AITriageReviewRepository) List ¶ added in v0.1.8
func (r *AITriageReviewRepository) List(ctx context.Context, tenantID shared.ID, filter ports.AITriageReviewFilter) ([]aitriagereview.Review, error)
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 ¶
NewAUPStore returns an AUP store backed by the given pool.
func (*AUPStore) Save ¶
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.
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) CountVulnerabilityAdvisoriesChangedSince ¶ added in v0.1.8
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 (*AdvisoryMaterializer) CurrentSourceRecordIDs ¶ added in v0.1.8
func (*AdvisoryMaterializer) CurrentSourceRecordIDsBounded ¶ added in v0.2.0
func (*AdvisoryMaterializer) GetCanonical ¶ added in v0.1.8
func (*AdvisoryMaterializer) GetCanonicalAtRevision ¶ added in v0.1.8
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 (r *AdvisoryMaterializer) ListVulnerabilityAdvisoryRevisions(ctx context.Context, query vulnerabilityintel.AdvisoryRevisionQuery) (vulnerabilityintel.AdvisoryRevisionPage, error)
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 (*AdvisoryMaterializer) Materialize ¶ added in v0.1.8
func (r *AdvisoryMaterializer) Materialize(ctx context.Context, records []advisory.ObservationRecord) (advisory.MaterializationResult, error)
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
func (r *AdvisoryMaterializer) PublishSourceSnapshot(ctx context.Context, publication ports.SourceSnapshotPublication, records []advisory.ObservationRecord) ([]advisory.MaterializationResult, error)
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 (r *AdvisoryMaterializer) SummarizeVulnerabilityCoverage(ctx context.Context, tenantID shared.ID, requests []vulnerabilityintel.AdvisoryCoverageRequest) (map[string]vulnerabilityintel.AdvisoryCoverageSummary, error)
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
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) 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
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 (*AgentPlanStore) GetBySession ¶
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 (*AgentSessionStore) GetSession ¶
func (*AgentSessionStore) ListByEngagement ¶
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) SaveSession ¶
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
func (r *AgentSigningKeyRepository) ListByAgent(ctx context.Context, agentID shared.ID) ([]fleetagent.AgentSigningKey, error)
ListByAgent returns every key for an agent under the ctx tenant, newest NotBefore first.
func (*AgentSigningKeyRepository) Register ¶ added in v0.2.0
func (r *AgentSigningKeyRepository) Register(ctx context.Context, key fleetagent.AgentSigningKey) error
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) Decide ¶
func (s *ApprovalStore) Decide(ctx context.Context, d agent.ApprovalDecision) error
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 (s *ApprovalStore) Enqueue(ctx context.Context, a agent.ProposedAction) error
func (*ApprovalStore) Get ¶
func (s *ApprovalStore) Get(ctx context.Context, actionID shared.ID) (agent.ProposedAction, agent.ApprovalDecision, error)
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 (repository *AssessmentComparisonRepository) CreateQueued(ctx context.Context, comparison assessmentcomparison.Comparison) (assessmentcomparison.Comparison, bool, error)
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 (repository *AssessmentComparisonRepository) ListItems(ctx context.Context, tenantID, comparisonID shared.ID, filter ports.AssessmentComparisonItemFilter) (ports.AssessmentComparisonItemPage, error)
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 (*AssessmentCycleBackfillRepository) CommitAssessmentCycleBackfillItem ¶ added in v0.2.0
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 (*AssessmentCycleIntegrityRepository) AssessmentCycleIntegrityGeneration ¶ added in v0.2.0
func (*AssessmentCycleIntegrityRepository) CountAssessmentCycleIntegritySubjects ¶ added in v0.2.0
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 (*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 (r *AssessmentCycleRepository) CommitClosure(ctx context.Context, commit ports.AssessmentClosureCommit) error
func (*AssessmentCycleRepository) CreateCycle ¶ added in v0.2.0
func (r *AssessmentCycleRepository) CreateCycle(ctx context.Context, cycle *assessmentcycle.AssessmentCycle) error
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 (*AssessmentCycleRepository) DeleteMember ¶ added in v0.2.0
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 (r *AssessmentCycleRepository) ListCycles(ctx context.Context, query ports.AssessmentCycleListQuery) ([]ports.AssessmentCycleListRecord, error)
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 (r *AssessmentCycleRepository) ListMigrationPendingAssessments(ctx context.Context, query ports.AssessmentCycleListQuery) ([]ports.AssessmentCycleMigrationPendingRecord, int, error)
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 (*AssessmentCycleRepository) ReopenClosure ¶ added in v0.2.0
func (r *AssessmentCycleRepository) ReopenClosure(ctx context.Context, reopen ports.AssessmentClosureReopen) error
func (*AssessmentCycleRepository) SaveClosureReport ¶ added in v0.2.0
func (r *AssessmentCycleRepository) SaveClosureReport(ctx context.Context, report ports.AssessmentClosureReportArtifact) (ports.AssessmentClosureReportArtifact, bool, error)
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
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 (repository *AssessmentRelationshipRepository) DecideCandidateCAS(ctx context.Context, decision assessmentrelationship.Decision, plan *assessmentrelationship.RepairPlan) (record assessmentrelationship.Record, replayed bool, err error)
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
func (repository *AssessmentRelationshipRepository) ListCandidates(ctx context.Context, tenantID shared.ID, filter ports.AssessmentRelationshipCandidateFilter) (records []assessmentrelationship.Record, err error)
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 (*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 (repository *AssessmentSnapshotRepository) ListAssessmentSnapshots(ctx context.Context, query ports.AssessmentSnapshotListQuery) (ports.AssessmentSnapshotPage, error)
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 (*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
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
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
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 ¶
NewAuditLog returns an audit log backed by the given pool.
func (*AuditLog) List ¶
List returns the calling tenant's most recent audit entries (newest first), capped at limit.
func (*AuditLog) MigrationMetadata ¶ added in v0.2.0
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) RecordOnce ¶ added in v0.1.8
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 ¶
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
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
func (r *BaselineRepository) Load(ctx context.Context, key baseline.Key) (ports.BaselineRecord, error)
Load returns the record for a key, or shared.ErrNotFound.
func (*BaselineRepository) Save ¶ added in v0.2.0
func (r *BaselineRepository) Save(ctx context.Context, rec ports.BaselineRecord) error
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
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 (*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 (*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 (r *CoverageWindowRepository) AppendCoverageWindow(ctx context.Context, window sensorstate.CoverageWindow) (sensorstate.CoverageWindow, error)
func (*CoverageWindowRepository) ListCoverageWindows ¶ added in v0.2.0
func (r *CoverageWindowRepository) ListCoverageWindows(ctx context.Context, q ports.CoverageWindowQuery) ([]sensorstate.CoverageWindow, error)
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 (*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 (*DASTRunStore) SaveDASTRun ¶ added in v0.2.0
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 (r *DetectionProvenanceRepository) AdmitPending(ctx context.Context, current detectionprovenance.Current, received detectionprovenance.Transition) error
func (*DetectionProvenanceRepository) AppendTransition ¶ added in v0.2.0
func (r *DetectionProvenanceRepository) AppendTransition(ctx context.Context, transition detectionprovenance.Transition) error
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 (r *DetectionProvenanceRepository) ListPending(ctx context.Context) ([]detectionprovenance.Current, error)
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.
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.
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
func (r *EndpointTimelineRepository) QueryTimeline(ctx context.Context, q ports.EndpointTimelineQuery) ([]endpoint.TimelineEntry, error)
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 ¶
func (r *EngagementRepository) Create(ctx context.Context, e *engagement.Engagement) error
Create inserts the engagement and its scope targets in one transaction.
func (*EngagementRepository) Delete ¶
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 (*EngagementRepository) ListTenantIDs ¶ added in v0.1.8
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 ¶
func (r *EngagementRepository) Update(ctx context.Context, e *engagement.Engagement) error
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 (r *EngagementSourceRepository) Create(ctx context.Context, item sourcepackage.Package) (sourcepackage.Package, bool, error)
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
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 ¶
Append inserts sealed evidence items in order, in one transaction (append-only).
func (*EvidenceStore) Head ¶
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
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 (*FindingLineageBackfillRepository) CommitFindingLineageBackfillItem ¶ added in v0.2.0
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
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 (repository *FindingLineageRepository) AppendOverrideCAS(ctx context.Context, event findinglineage.OverrideEvent) (findinglineage.OverrideEvent, bool, error)
func (*FindingLineageRepository) AppendSkip ¶ added in v0.2.0
func (repository *FindingLineageRepository) AppendSkip(ctx context.Context, record findinglineage.SkipRecord) (findinglineage.SkipRecord, bool, error)
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 (*FindingLineageRepository) FindIdentitiesByFingerprint ¶ added in v0.2.0
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 (*FindingLineageRepository) ResolveCandidateCAS ¶ added in v0.2.0
func (repository *FindingLineageRepository) ResolveCandidateCAS(ctx context.Context, updated findinglineage.MatchCandidate, event findinglineage.ResolutionEvent) (findinglineage.MatchCandidate, findinglineage.ResolutionEvent, bool, error)
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 (*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 (*FindingRepository) ListByEngagement ¶
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 (*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.
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) ListAgents ¶ added in v0.1.8
func (r *FleetAgentRepository) ListAgents(ctx context.Context, tenantID shared.ID) ([]*fleetagent.Agent, 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) 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
func (r *FleetDesiredRepository) Put(ctx context.Context, state *fleetdesired.State) error
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
func (r *FleetRolloutRepository) Put(ctx context.Context, plan *fleetrollout.Plan) error
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 (*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 (*IdentityStore) RevokeSession ¶ added in v0.2.0
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
func (r *ImportReceiptRepository) CreateOrGet(ctx context.Context, tenantID shared.ID, receipt importreceipt.Receipt) (importreceipt.Receipt, bool, error)
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
func (r *ImportedFindingRepository) Save(ctx context.Context, tenantID shared.ID, findings []importedfinding.ImportedFinding) (int, int, error)
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 (*IncidentEventRepository) ListPendingResponseLinks ¶ added in v0.2.0
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
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 (*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 (*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 (*IntegrationStore) ResolveIntegrationCredential ¶ added in v0.2.0
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) Deadletter ¶
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
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 ¶
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 ¶
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) ListActive ¶ added in v0.2.0
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 (*NotificationRepository) CancelDelivery ¶ added in v0.2.0
func (*NotificationRepository) CreateChannel ¶ added in v0.2.0
func (r *NotificationRepository) CreateChannel(ctx context.Context, c notification.Channel, sealed string) (notification.Channel, error)
func (*NotificationRepository) CreateRule ¶ added in v0.2.0
func (r *NotificationRepository) CreateRule(ctx context.Context, rule notification.Rule) (notification.Rule, error)
func (*NotificationRepository) DeadLetterDelivery ¶ added in v0.2.0
func (*NotificationRepository) DeleteChannel ¶ added in v0.2.0
func (*NotificationRepository) DeleteRule ¶ added in v0.2.0
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 (*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 (r *NotificationRepository) ListDeliveries(ctx context.Context, f ports.NotificationDeliveryFilter) (notification.Page, error)
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 (r *NotificationRepository) Publish(ctx context.Context, e notification.Event) ([]shared.ID, error)
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
func (r *NotificationRepository) UpdateRule(ctx context.Context, rule notification.Rule) (notification.Rule, error)
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) 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 (*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 (*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 (*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 (r *OwnershipRepository) ActivatePolicy(ctx context.Context, a ports.OwnershipActivation) error
func (*OwnershipRepository) AppendIntent ¶ added in v0.2.0
func (r *OwnershipRepository) AppendIntent(ctx context.Context, i ports.OwnershipIntent) error
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 (*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 (*OwnershipRepository) CreateTeam ¶ added in v0.2.0
func (*OwnershipRepository) DeleteAssetMapping ¶ added in v0.2.0
func (*OwnershipRepository) DeleteMapping ¶ added in v0.2.0
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 (*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 (*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) MarkOwnershipSourceReady ¶ added in v0.2.0
func (*OwnershipRepository) OwnershipInbox ¶ added in v0.2.0
func (r *OwnershipRepository) OwnershipInbox(ctx context.Context, f ports.OwnershipInboxFilter) (out ports.OwnershipInboxPage, err error)
func (*OwnershipRepository) RemoveMember ¶ added in v0.2.0
func (*OwnershipRepository) ReserveOwnershipBulk ¶ added in v0.2.0
func (*OwnershipRepository) SaveAssetMapping ¶ added in v0.2.0
func (r *OwnershipRepository) SaveAssetMapping(ctx context.Context, m ports.OwnershipAssetMapping) error
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 (*OwnershipRepository) UpdateTeam ¶ added in v0.2.0
func (*OwnershipRepository) VisibleOwnershipEngagement ¶ added in v0.2.0
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) CurrentAnalysisHotspotSummary ¶
func (*ProjectAnalysisStore) CurrentFindingStatuses ¶ added in v0.1.8
func (*ProjectAnalysisStore) Get ¶
func (r *ProjectAnalysisStore) Get(ctx context.Context, tenantID, projectID, analysisID shared.ID) (projectanalysis.Analysis, error)
func (*ProjectAnalysisStore) GetHotspot ¶
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 (*ProjectAnalysisStore) LatestWithResult ¶
func (r *ProjectAnalysisStore) LatestWithResult(ctx context.Context, tenantID, projectID shared.ID, branch string) (projectanalysis.Analysis, []byte, error)
func (*ProjectAnalysisStore) ListAnalysisHotspots ¶
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 (*ProjectAnalysisStore) Save ¶
func (r *ProjectAnalysisStore) Save(ctx context.Context, analysis projectanalysis.Analysis) error
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 (r *ProjectAnalysisStore) TransitionHotspot(ctx context.Context, cmd hotspot.TransitionCommand) (hotspot.Hotspot, hotspot.ReviewEvent, error)
func (*ProjectAnalysisStore) TransitionIssue ¶
func (r *ProjectAnalysisStore) TransitionIssue(ctx context.Context, cmd issue.TransitionCommand) (issue.Issue, issue.ReviewEvent, error)
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 (*ProjectRepository) DeleteByKey ¶
func (*ProjectRepository) SetPullRequestDecoration ¶ added in v0.2.0
func (*ProjectRepository) UpdateGate ¶
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:
- Acquires sorted advisory transaction locks on (judgment, fingerprint) to serialize concurrent idempotency checks (prevents deadlocks).
- Checks judgment-level idempotency (tenant+judgmentID).
- Checks fingerprint-level idempotency (tenant+fingerprint).
- Locks the finding FOR UPDATE, verifies CAS (priority + version).
- Binds command metadata to CAS state.
- Validates exact reversal for corroborating_signal_loss.
- Constructs and validates the PromotionEvent.
- Mutates the finding (priority + version) for escalating/de-escalating effects.
- 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
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
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 (*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) DeleteIfUnassigned ¶
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) Get ¶
func (s *QualityProfileStore) Get(ctx context.Context, tenantID shared.ID, key string) (qualityprofile.Profile, error)
func (*QualityProfileStore) List ¶
func (s *QualityProfileStore) List(ctx context.Context, tenantID shared.ID) ([]qualityprofile.Profile, error)
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) ListByEngagement ¶
func (r *ReconRunStore) ListByEngagement(ctx context.Context, engagementID shared.ID) ([]recon.Run, error)
ListByEngagement returns an engagement's runs, newest first.
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
func (r *ResponseHaltWriterRepository) AdvanceHaltGenerationWithAudit(ctx context.Context, expected int64, intent ports.ResponseAuditIntent) (int64, ports.ResponseAuditIntent, ports.ResponseHaltDispatch, error)
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
func (r *ResponseObserverBindingRepository) SaveResponseObserverBindingWithAudit(ctx context.Context, binding fleetagent.ResponseObserverBinding, expectedVersion int, intent ports.FleetAuditIntent) (fleetagent.ResponseObserverBinding, ports.FleetAuditIntent, error)
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
func (r *ResponseRepository) AdvanceHaltGenerationWithAudit(context.Context, int64, ports.ResponseAuditIntent) (int64, ports.ResponseAuditIntent, ports.ResponseHaltDispatch, error)
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
func (r *ResponseRepository) EnqueueResponseAudit(ctx context.Context, intent ports.ResponseAuditIntent) (ports.ResponseAuditIntent, error)
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
func (r *ResponseRepository) StartAttempt(ctx context.Context, a responsesaga.ResponseAttempt) (responsesaga.ResponseAttempt, bool, error)
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
func (r *ResponseRepository) TransitionAttempt(ctx context.Context, a responsesaga.ResponseAttempt, from responsesaga.SagaState) (responsesaga.ResponseAttempt, bool, error)
TransitionAttempt atomically persists a state transition and its outcome/provenance payload.
func (*ResponseRepository) TransitionAttemptWithAudit ¶ added in v0.2.0
func (r *ResponseRepository) TransitionAttemptWithAudit(ctx context.Context, attempt responsesaga.ResponseAttempt, from responsesaga.SagaState, intent ports.ResponseAuditIntent) (responsesaga.ResponseAttempt, bool, ports.ResponseAuditIntent, error)
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 (r *ResponseVerificationRepository) AppendResponseTargetEvidenceReceipt(ctx context.Context, receipt fleetagent.ResponseTargetEvidenceReceipt) (fleetagent.ResponseTargetEvidenceReceipt, error)
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).
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 ¶
NewRunLock returns a Postgres-backed run locker.
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
Delete removes one connector; ErrNotFound when absent under the tenant.
func (*SCMConnectorRepository) Get ¶ added in v0.2.0
func (r *SCMConnectorRepository) Get(ctx context.Context, id shared.ID) (ports.SCMConnectorMeta, error)
Get returns one connector's metadata; ErrNotFound when absent under the tenant.
func (*SCMConnectorRepository) List ¶ added in v0.2.0
func (r *SCMConnectorRepository) List(ctx context.Context) ([]ports.SCMConnectorMeta, error)
List returns the tenant's connectors as metadata (never the token), ordered by host then name.
func (*SCMConnectorRepository) Put ¶ added in v0.2.0
func (r *SCMConnectorRepository) Put(ctx context.Context, c scmconnector.Connector, token []byte) error
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 (*SLAStore) ActivePolicy ¶ added in v0.1.8
func (*SLAStore) AssessmentHistory ¶ added in v0.1.8
func (*SLAStore) LifecycleEvents ¶ added in v0.1.8
func (*SLAStore) ListCurrent ¶ added in v0.1.8
func (*SLAStore) PolicyHistory ¶ added in v0.1.8
func (*SLAStore) SLAHistories ¶ added in v0.2.0
func (*SLAStore) SaveTransition ¶ added in v0.1.8
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 (*ScanJobStore) LatestForEngagement ¶
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.
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 (*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 ¶
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) 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) 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) SaveScanRun ¶ added in v0.2.0
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 (r *SensorStateRepository) ListCoverageSensorStates(ctx context.Context, q ports.CoverageSensorStateQuery) ([]sensorstate.Observation, error)
func (*SensorStateRepository) ListSensorStates ¶ added in v0.2.0
func (r *SensorStateRepository) ListSensorStates(ctx context.Context, q ports.SensorStateQuery) ([]sensorstate.Observation, error)
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 (s *SyncRunStore) Finish(ctx context.Context, id shared.ID, state vulnerabilitysync.State, counts vulnerabilitysync.Counts, samples []string) (vulnerabilitysync.Run, error)
func (*SyncRunStore) Get ¶ added in v0.1.8
func (s *SyncRunStore) Get(ctx context.Context, id shared.ID) (vulnerabilitysync.Run, error)
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 (*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 (*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 (s *SyncRunStore) Start(ctx context.Context, request ports.SyncRunStart) (vulnerabilitysync.Run, bool, 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
func (r *TelemetryRepository) Footprint(ctx context.Context) (ports.TelemetryFootprint, error)
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
func (r *TelemetryRepository) Ingest(ctx context.Context, batch ports.TelemetryBatch) error
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
func (r *TelemetryRepository) Query(ctx context.Context, q ports.HuntQuery) (ports.HuntResult, error)
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
func (r *TelemetryTransportRepository) AgentGapRevisions(ctx context.Context, gapID shared.ID) ([]ports.TelemetryAgentGapRevision, error)
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
func (r *TelemetryTransportRepository) CommitBatch(ctx context.Context, batch ports.TelemetryEventBatch) error
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 (r *TelemetryTransportRepository) CommitBatchWithAudit( ctx context.Context, batch ports.TelemetryEventBatch, intent ports.FleetAuditIntent, ) (ports.FleetAuditIntent, error)
func (*TelemetryTransportRepository) CountBatchEvents ¶ added in v0.2.0
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
func (r *TelemetryTransportRepository) ListCoverageGapFacts(ctx context.Context, q ports.CoverageGapQuery) ([]ports.CoverageGapFact, error)
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) QueryAgentGaps ¶ added in v0.2.0
func (r *TelemetryTransportRepository) QueryAgentGaps(ctx context.Context, q ports.TelemetryGapQuery) ([]ports.TelemetryGap, error)
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 (r *TelemetryTransportRepository) QueryDeliveryGaps(ctx context.Context, q ports.TelemetryGapQuery) ([]ports.TelemetryGap, error)
func (*TelemetryTransportRepository) QueryTelemetryBatchAccounting ¶ added in v0.2.0
func (r *TelemetryTransportRepository) QueryTelemetryBatchAccounting(ctx context.Context, q ports.TelemetryBatchAccountingQuery) ([]ports.TelemetryBatchAccounting, error)
func (*TelemetryTransportRepository) RecordAgentGap ¶ added in v0.2.0
func (r *TelemetryTransportRepository) RecordAgentGap(ctx context.Context, gap ports.TelemetryAgentGap) error
func (*TelemetryTransportRepository) ResolveTelemetryAsset ¶ added in v0.2.0
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 (r *TelemetryTransportRepository) SaveStreamState(ctx context.Context, state ports.TelemetryStreamState) error
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
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.
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) GetByAPIKeyHash ¶
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) 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.
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 (*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 (*VulnerabilityActionStore) CountPendingVulnerabilityActions ¶ added in v0.1.8
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 (s *VulnerabilityActionStore) ListActions(ctx context.Context, query vulnerabilityaction.ActionQuery) (vulnerabilityaction.ActionPage, error)
func (*VulnerabilityActionStore) ListVulnerabilityTransitions ¶ added in v0.1.8
func (s *VulnerabilityActionStore) ListVulnerabilityTransitions(ctx context.Context, query vulnerabilityintel.TransitionQuery) (vulnerabilityintel.TransitionPage, error)
func (*VulnerabilityActionStore) RecordChange ¶ added in v0.1.8
func (s *VulnerabilityActionStore) RecordChange(ctx context.Context, change vulnerabilityaction.Change) (bool, error)
func (*VulnerabilityActionStore) ResolveAction ¶ added in v0.1.8
func (*VulnerabilityActionStore) RetryOutbox ¶ added in v0.1.8
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 (*VulnerabilityOccurrenceStore) CountNewlyAffectedAssets ¶ added in v0.2.0
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 (*VulnerabilityOccurrenceStore) ListVulnerabilityOccurrences ¶ added in v0.1.8
func (s *VulnerabilityOccurrenceStore) ListVulnerabilityOccurrences(ctx context.Context, query vulnerabilityintel.OccurrenceQuery) (vulnerabilityintel.OccurrencePage, error)
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
func (s *VulnerabilityOccurrenceStore) Upsert(ctx context.Context, occurrence vulnerabilityoccurrence.Occurrence) (vulnerabilityoccurrence.UpsertResult, error)
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 (s *VulnerabilityReconcileRunStore) Finish(ctx context.Context, id shared.ID, state vulnerabilityreconcile.State, counts vulnerabilityreconcile.Counts, samples []string) (vulnerabilityreconcile.Run, error)
func (*VulnerabilityReconcileRunStore) Get ¶ added in v0.1.8
func (s *VulnerabilityReconcileRunStore) Get(ctx context.Context, id shared.ID) (vulnerabilityreconcile.Run, error)
func (*VulnerabilityReconcileRunStore) GetByDurableJobID ¶ added in v0.1.8
func (s *VulnerabilityReconcileRunStore) GetByDurableJobID(ctx context.Context, jobID string) (vulnerabilityreconcile.Run, error)
func (*VulnerabilityReconcileRunStore) HasReconciliationMatch ¶ added in v0.1.8
func (*VulnerabilityReconcileRunStore) ListReconciliationDiffs ¶ added in v0.1.8
func (s *VulnerabilityReconcileRunStore) ListReconciliationDiffs(ctx context.Context, query vulnerabilityreconcile.DiffQuery) (vulnerabilityreconcile.DiffPage, error)
func (*VulnerabilityReconcileRunStore) MarkRunning ¶ added in v0.1.8
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 (s *VulnerabilityReconcileRunStore) Start(ctx context.Context, request ports.VulnerabilityReconcileStart) (vulnerabilityreconcile.Run, bool, error)
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 (*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 (s *VulnerabilityRiskAssessmentStore) ListVulnerabilityAssessments(ctx context.Context, query vulnerabilityintel.AssessmentQuery) (vulnerabilityintel.AssessmentPage, error)
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
func (s *VulnerabilityRiskAssessmentStore) Upsert(ctx context.Context, assessment vulnerabilityrisk.Assessment) (vulnerabilityrisk.Result, error)
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) Create ¶ added in v0.1.8
func (s *VulnerabilitySourceStore) Create(ctx context.Context, source vulnerabilitysource.Source) error
func (*VulnerabilitySourceStore) Get ¶ added in v0.1.8
func (s *VulnerabilitySourceStore) Get(ctx context.Context, id shared.ID) (vulnerabilitysource.Source, error)
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 (*WorkOrderRepository) CancelResponsesBelowGenerationWithAudit ¶ added in v0.2.0
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 (*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 (*WorkOrderRepository) Issue ¶ added in v0.1.8
func (r *WorkOrderRepository) Issue(ctx context.Context, wo *workorder.WorkOrder) (*workorder.WorkOrder, error)
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 (r *WorkOrderRepository) IssueWithAudit(ctx context.Context, wo *workorder.WorkOrder, intent ports.FleetAuditIntent) (*workorder.WorkOrder, ports.FleetAuditIntent, error)
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 (*WorkOrderRepository) TransitionLeasedWithAudit ¶ added in v0.2.0
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 ¶
func (r *WriteupDraftRepository) Save(ctx context.Context, d writeupdraft.Draft) error
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
¶
- accuracy_store.go
- advisory_bulk_writer.go
- advisory_materializer.go
- advisory_repo.go
- agent_decision_store.go
- agent_plan_store.go
- agent_session_store.go
- agent_signing_key_repo.go
- aitriage_review_repo.go
- approval_store.go
- assessment_closure_history.go
- assessment_closure_repo.go
- assessment_closure_report_store.go
- assessment_comparison_repo.go
- assessment_cycle_backfill_repo.go
- assessment_cycle_integrity_repo.go
- assessment_cycle_repo.go
- assessment_cycle_request_repo.go
- assessment_relationship_repo.go
- assessment_snapshot_backfill_repo.go
- assessment_snapshot_repo.go
- asset_lookup.go
- asset_repo.go
- attack_path_store.go
- aup_audit.go
- baseline_repo.go
- cloud_observation_store.go
- cloud_run_store.go
- comment_repo.go
- component_inventory_store.go
- correlation_state_repo.go
- coverage_window_repo.go
- dast_run_store.go
- db.go
- detection_provenance_repo.go
- detection_record_repo.go
- emulation_run_repo.go
- endpoint_process_repo.go
- endpoint_timeline_repo.go
- engagement_repo.go
- engagement_source_repo.go
- evidence_repo.go
- exploitation_chain_repo.go
- finding_lineage_backfill_repo.go
- finding_lineage_repo.go
- finding_repo.go
- fleet_audit_repo.go
- fleetagent_repo.go
- fleetdesired_repo.go
- fleetrollout_repo.go
- identity_store.go
- imported_sbom_store.go
- importedfinding_repo.go
- importreceipt_repo.go
- incident_event_repo.go
- integration_store.go
- jobqueue.go
- judgment_repo.go
- leader_store.go
- legal_hold_repo.go
- notification_capture.go
- notification_reconcile.go
- notification_repository.go
- notification_source.go
- ownership_assignment.go
- ownership_execution.go
- ownership_notification.go
- ownership_policy.go
- ownership_read.go
- ownership_repository.go
- ownership_source.go
- ownership_work.go
- privacy_policy_repo.go
- project_analysis_source.go
- project_analysis_store.go
- project_hotspot_store.go
- project_issue_store.go
- project_repo.go
- promotion_store.go
- purple_repo.go
- quality_gate_mutator.go
- quality_gate_store.go
- quality_profile_store.go
- recon_run_repo.go
- response_audit_repo.go
- response_halt_writer.go
- response_observer_binding_repo.go
- response_repo.go
- response_verification_repo.go
- retest_repo.go
- runlock.go
- runlock_lease.go
- scan_job_repo.go
- scan_repo.go
- scan_result_repo.go
- scan_run_evidence.go
- scan_run_repo.go
- scanned_image_store.go
- scm_connector_repo.go
- sensor_state_repo.go
- sla_store.go
- sync_run_store.go
- telemetry_agent_gap.go
- telemetry_agent_gap_query.go
- telemetry_batch_commit.go
- telemetry_repo.go
- telemetry_transport_repo.go
- tenant.go
- threat_model_repo.go
- timestamp_store.go
- user_repo.go
- vex_statement_repo.go
- vulnerability_action_store.go
- vulnerability_occurrence_store.go
- vulnerability_reconcile_run_store.go
- vulnerability_retention_store.go
- vulnerability_risk_assessment_store.go
- vulnerability_source_store.go
- workorder_repo.go
- writeupdraft_repo.go