memory

package
v0.2.4 Latest Latest
Warning

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

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

Documentation

Overview

Package memory provides in-memory repository implementations for the walking skeleton and tests. Replaced by the Postgres adapters.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type AITriageReviewStore added in v0.1.8

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

func NewAITriageReviewStore added in v0.1.8

func NewAITriageReviewStore() *AITriageReviewStore

func (*AITriageReviewStore) Get added in v0.1.8

func (*AITriageReviewStore) List added in v0.1.8

func (*AITriageReviewStore) SaveDecision added in v0.1.8

func (s *AITriageReviewStore) SaveDecision(_ context.Context, review aitriagereview.Review, expectedVersion int) error

func (*AITriageReviewStore) SaveOwner added in v0.1.8

func (s *AITriageReviewStore) SaveOwner(_ context.Context, review aitriagereview.Review, expectedVersion int) error

func (*AITriageReviewStore) UpsertPending added in v0.1.8

func (s *AITriageReviewStore) UpsertPending(_ context.Context, review aitriagereview.Review) error

type AccuracyRunStore added in v0.2.0

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

AccuracyRunStore is an in-memory ports.AccuracyRunStore for dev and tests. Detection-accuracy runs are deployment-global engine data (no tenant).

func NewAccuracyRunStore added in v0.2.0

func NewAccuracyRunStore() *AccuracyRunStore

NewAccuracyRunStore returns an empty in-memory accuracy-run store.

func (*AccuracyRunStore) Recent added in v0.2.0

func (s *AccuracyRunStore) Recent(_ context.Context, limit int) ([]accuracy.Run, error)

Recent returns the most recent runs, newest first, capped at limit (all when limit <= 0).

func (*AccuracyRunStore) Save added in v0.2.0

Save appends a run.

type AdvisoryMaterializer added in v0.1.8

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

AdvisoryMaterializer is an in-memory implementation of the global advisory observation history and canonical projection.

func NewAdvisoryMaterializer added in v0.1.8

func NewAdvisoryMaterializer() *AdvisoryMaterializer

func (*AdvisoryMaterializer) AdvisoryRevisionAt added in v0.1.8

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

func (*AdvisoryMaterializer) ByCPE added in v0.1.8

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

func (*AdvisoryMaterializer) ByPackage added in v0.1.8

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

func (*AdvisoryMaterializer) CountVulnerabilityAdvisoriesChangedSince added in v0.1.8

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

func (*AdvisoryMaterializer) CountVulnerabilityAdvisoryDailyImpact added in v0.2.0

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

func (*AdvisoryMaterializer) CurrentRevision added in v0.1.8

func (m *AdvisoryMaterializer) CurrentRevision(_ context.Context, id string) (int64, error)

func (*AdvisoryMaterializer) CurrentSourceRecordIDs added in v0.1.8

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

func (*AdvisoryMaterializer) CurrentSourceRecordIDsBounded added in v0.2.0

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

func (*AdvisoryMaterializer) GetCanonical added in v0.1.8

func (*AdvisoryMaterializer) GetCanonicalAtRevision added in v0.1.8

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

func (*AdvisoryMaterializer) ListAdvisoryRevisions added in v0.1.8

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

func (*AdvisoryMaterializer) ListVulnerabilityAdvisories added in v0.1.8

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

func (*AdvisoryMaterializer) ListVulnerabilityAdvisoryRevisions added in v0.1.8

func (*AdvisoryMaterializer) ListVulnerabilitySyncRunRevisions added in v0.1.8

func (m *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 (m *AdvisoryMaterializer) MarkAdvisoryEvaluated(ctx context.Context, tenantID shared.ID, advisoryID string, revision int64, evaluatedAt time.Time) error

func (*AdvisoryMaterializer) Materialize added in v0.1.8

func (*AdvisoryMaterializer) MaterializeSourceSnapshot added in v0.2.0

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

MaterializeSourceSnapshot applies a complete source view through one copy-on-write swap. Each record is materialized independently because unrelated source records do not share an advisory identity, but an error leaves every live projection unchanged.

func (*AdvisoryMaterializer) OldestUnevaluatedAdvisory added in v0.1.8

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

func (*AdvisoryMaterializer) PublishSourceSnapshot added in v0.2.0

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

func (*AdvisoryMaterializer) PublishedSourceSnapshot added in v0.2.0

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

func (*AdvisoryMaterializer) WithClock added in v0.1.8

func (m *AdvisoryMaterializer) WithClock(now func() time.Time) *AdvisoryMaterializer

WithClock overrides the revision-createdAt clock (tests only). A nil clock is ignored.

type AdvisoryStore

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

AdvisoryStore is the in-memory owned-advisory store (dev/tests, mirrors the Postgres adapter). It is GLOBAL reference data – NOT tenant-scoped. Advisories are indexed by every affected (ecosystem, package) so ByPackage is a map lookup, and Upsert is idempotent by advisory id (re-syncable reference data, replaced in place – not append-only – except the exploitation-risk enrichment carried forward raise-only via advisory.Advisory.PreserveEnrichment). The stored ecosystem+package keys are the ingester-normalized, OSV-canonical ids per the ports.AdvisoryStore KEY CONTRACT.

func NewAdvisoryStore

func NewAdvisoryStore() *AdvisoryStore

NewAdvisoryStore returns an empty in-memory advisory store.

func (*AdvisoryStore) AdvisoryAliasEdges added in v0.2.0

func (s *AdvisoryStore) AdvisoryAliasEdges(_ context.Context, ids []string) ([]advisory.AliasEdge, error)

AdvisoryAliasEdges returns the alias edges for advisories whose id is in ids or whose Aliases intersect ids, mirroring the Postgres adapter. Bounded to the given ids. Empty ids -> no edges.

func (*AdvisoryStore) AdvisoryFreshness added in v0.2.0

func (s *AdvisoryStore) AdvisoryFreshness(_ context.Context) (time.Time, int, error)

AdvisoryFreshness reports the corpus size and a freshness timestamp, so the owned detection source's provenance marker (and thus the pipeline's detection-readiness guard) is non-empty when this in-memory store is populated. The in-memory store carries no per-advisory publish date (it is dev/test-only and re-loaded each process), so a populated store reports the time of its last write: the corpus is exactly as fresh as its last load, and an empty store reports a zero time + zero count so the readiness guard correctly treats it as no coverage. Production uses the Postgres adapter, which reports real advisory dates.

func (*AdvisoryStore) ByPackage

func (s *AdvisoryStore) ByPackage(_ context.Context, ecosystem, name string) ([]advisory.Advisory, error)

ByPackage returns the advisories that list (ecosystem, name) as affected, in deterministic id order (matching the Postgres adapter's ORDER BY id COLLATE "C"); the caller runs advisory.Match to decide which actually hit the component's version.

func (*AdvisoryStore) CoveredEcosystems added in v0.2.0

func (s *AdvisoryStore) CoveredEcosystems(_ context.Context) (map[string]bool, error)

CoveredEcosystems returns the distinct ecosystems this store has any affected-package index entry for, so the readiness guard can distinguish a genuinely-covered distro from a silent gap. The index keys are "ecosystem\x00package", so the ecosystem is the prefix before the NUL separator.

func (*AdvisoryStore) Upsert

Upsert inserts or replaces an advisory by id and (re)builds its (ecosystem, package) index entries. A re-sync may change the affected set, so the prior index entries for the id are dropped first. Idempotent.

type AgentSessionStore

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

AgentSessionStore is the in-memory ports.AgentSessionStore (dev/tests). It keeps the same (session, seq) transcript fork-guard as the Postgres adapter so a duplicate seq is rejected, not silently overwritten.

func NewAgentSessionStore

func NewAgentSessionStore() *AgentSessionStore

NewAgentSessionStore returns an empty in-memory agent session store.

func (*AgentSessionStore) AppendMessage

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

func (*AgentSessionStore) GetSession

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

func (*AgentSessionStore) ListByEngagement

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

func (*AgentSessionStore) ListResumable

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

func (*AgentSessionStore) Messages

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

func (*AgentSessionStore) SaveSession

func (s *AgentSessionStore) SaveSession(_ context.Context, sess agent.Session) error

type AgentSigningKeyStore added in v0.2.0

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

AgentSigningKeyStore is the in-memory agent-signing-key registry used inline/in dev. Keys are bucketed per tenant, so a read under one tenant can never observe another's — the same isolation the Postgres store gets from RLS. It keeps deep copies so a caller mutating a returned key cannot corrupt stored state, and enforces the same immutability + monotonic-revocation invariants as the Postgres store.

func NewAgentSigningKeyStore added in v0.2.0

func NewAgentSigningKeyStore() *AgentSigningKeyStore

NewAgentSigningKeyStore constructs the store.

func (*AgentSigningKeyStore) ListByAgent added in v0.2.0

func (s *AgentSigningKeyStore) 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 (*AgentSigningKeyStore) Register added in v0.2.0

Register stores a signing key, idempotent on its identity and anti-rollback on its KeyID. Like the Postgres store it does NOT verify proof-of-possession — see the port contract; the caller must have called fleetagent.VerifyKeyPossession first.

func (*AgentSigningKeyStore) ResolveSigningKey added in v0.2.0

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

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

func (*AgentSigningKeyStore) Revoke added in v0.2.0

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

Revoke marks a key revoked at `at`, monotonic: an already-revoked key keeps its first RevokedAt.

type ApprovalStore

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

ApprovalStore is the in-memory ports.ApprovalStore (dev/tests). Decide is idempotent: the first terminal decision wins; a second returns ErrConflict (so a double-click / race cannot re-open an admitted action).

func NewApprovalStore

func NewApprovalStore() *ApprovalStore

NewApprovalStore returns an empty in-memory approval store.

func (*ApprovalStore) Consume added in v0.1.8

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

func (*ApprovalStore) Decide

func (*ApprovalStore) EngagementsWithPending

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

func (*ApprovalStore) Enqueue

func (*ApprovalStore) Get

func (*ApprovalStore) Pending

func (s *ApprovalStore) Pending(_ 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() *AssessmentComparisonRepository

func (*AssessmentComparisonRepository) CreateQueued added in v0.2.0

func (*AssessmentComparisonRepository) Get added in v0.2.0

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

func (*AssessmentComparisonRepository) GetAssessmentComparisonBacklog added in v0.2.0

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

func (*AssessmentComparisonRepository) GetByInputHash added in v0.2.0

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

func (*AssessmentComparisonRepository) GetItem added in v0.2.0

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

func (*AssessmentComparisonRepository) GetMetadata added in v0.2.0

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

func (*AssessmentComparisonRepository) ListFailedAssessmentComparisons added in v0.2.0

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

func (*AssessmentComparisonRepository) ListItems added in v0.2.0

func (*AssessmentComparisonRepository) ListMetadataByCycle added in v0.2.0

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

func (*AssessmentComparisonRepository) SummarizeItems added in v0.2.0

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

func (*AssessmentComparisonRepository) UpdateCAS added in v0.2.0

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

type AssessmentCycleBackfillRepository added in v0.2.0

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

func NewAssessmentCycleBackfillRepository added in v0.2.0

func NewAssessmentCycleBackfillRepository() *AssessmentCycleBackfillRepository

func (*AssessmentCycleBackfillRepository) AcquireAssessmentCycleBackfillRun added in v0.2.0

func (*AssessmentCycleBackfillRepository) AdvanceAssessmentCycleBackfillRun added in v0.2.0

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

func (*AssessmentCycleBackfillRepository) CommitAssessmentCycleBackfillItem added in v0.2.0

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

func (*AssessmentCycleBackfillRepository) FinishAssessmentCycleBackfillRun added in v0.2.0

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

func (*AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillItem added in v0.2.0

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

func (*AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillRun added in v0.2.0

func (repository *AssessmentCycleBackfillRepository) GetAssessmentCycleBackfillRun(_ context.Context, tenantID, runID shared.ID) (ports.AssessmentCycleBackfillRun, 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(engagements *EngagementRepository, cycles *AssessmentCycleRepository) *AssessmentCycleIntegrityRepository

func (*AssessmentCycleIntegrityRepository) AcquireAssessmentCycleIntegrityRun added in v0.2.0

func (*AssessmentCycleIntegrityRepository) AdvanceAssessmentCycleIntegrityRun added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) AssessmentCycleIntegrityGeneration added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) CountAssessmentCycleIntegritySubjects added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) FinishAssessmentCycleIntegrityRun added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegrityRun added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) GetAssessmentCycleIntegritySubject added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegrityFindings added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) ListAssessmentCycleIntegritySubjects added in v0.2.0

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

func (*AssessmentCycleIntegrityRepository) SaveAssessmentCycleIntegritySubject added in v0.2.0

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

type AssessmentCycleReaders added in v0.2.0

type AssessmentCycleReaders struct {
	Engagements ports.EngagementRepository
	Snapshots   ports.AssessmentSnapshotRepository
	Comparisons ports.AssessmentComparisonRepository
	Runs        ports.ScanRunProvenanceStore
}

AssessmentCycleReaders supplies the same read projections joined by PostgreSQL. Omit readers only for domain/repository fixtures which do not exercise reads.

type AssessmentCycleRepository added in v0.2.0

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

AssessmentCycleRepository is an in-memory implementation of ports.AssessmentCycleRepository.

func NewAssessmentCycleRepository added in v0.2.0

func NewAssessmentCycleRepository(readers ...AssessmentCycleReaders) *AssessmentCycleRepository

NewAssessmentCycleRepository creates a new in-memory AssessmentCycleRepository.

func (*AssessmentCycleRepository) CommitClosure added in v0.2.0

func (*AssessmentCycleRepository) CreateCycle added in v0.2.0

func (*AssessmentCycleRepository) CreateMember added in v0.2.0

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

func (*AssessmentCycleRepository) DeleteCycle added in v0.2.0

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

func (*AssessmentCycleRepository) DeleteMember added in v0.2.0

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

func (*AssessmentCycleRepository) GetActiveClosureManifest added in v0.2.0

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

func (*AssessmentCycleRepository) GetClosureManifest added in v0.2.0

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

func (*AssessmentCycleRepository) GetClosureReport added in v0.2.0

func (r *AssessmentCycleRepository) GetClosureReport(_ 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(_ context.Context, tenantID, cycleID shared.ID) ([]assessmentclosure.Manifest, error)

func (*AssessmentCycleRepository) ListCycles added in v0.2.0

func (*AssessmentCycleRepository) ListMembers added in v0.2.0

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

func (*AssessmentCycleRepository) ListMigrationPendingAssessments added in v0.2.0

func (*AssessmentCycleRepository) LockCycleForUpdate added in v0.2.0

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

func (*AssessmentCycleRepository) NextManifestVersion added in v0.2.0

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

func (*AssessmentCycleRepository) ReopenClosure added in v0.2.0

func (*AssessmentCycleRepository) SaveClosureReport added in v0.2.0

func (*AssessmentCycleRepository) UpdateCycleCAS added in v0.2.0

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

func (*AssessmentCycleRepository) UpdateMemberCAS added in v0.2.0

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

type AssessmentCycleRequestRepository added in v0.2.0

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

func NewAssessmentCycleRequestRepository added in v0.2.0

func NewAssessmentCycleRequestRepository() *AssessmentCycleRequestRepository

func (*AssessmentCycleRequestRepository) AbortAssessmentCycleRequest added in v0.2.0

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

func (*AssessmentCycleRequestRepository) BeginAssessmentCycleRequest added in v0.2.0

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

func (*AssessmentCycleRequestRepository) CompleteAssessmentCycleRequest added in v0.2.0

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

type AssessmentRelationshipRepository added in v0.2.0

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

func NewAssessmentRelationshipRepository added in v0.2.0

func NewAssessmentRelationshipRepository() *AssessmentRelationshipRepository

func (*AssessmentRelationshipRepository) CreateCandidate added in v0.2.0

func (*AssessmentRelationshipRepository) DecideCandidateCAS added in v0.2.0

func (*AssessmentRelationshipRepository) GetCandidate added in v0.2.0

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

func (*AssessmentRelationshipRepository) ListCandidates added in v0.2.0

type AssessmentSnapshotBackfillRepository added in v0.2.0

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

func NewAssessmentSnapshotBackfillRepository added in v0.2.0

func NewAssessmentSnapshotBackfillRepository() *AssessmentSnapshotBackfillRepository

func (*AssessmentSnapshotBackfillRepository) AcquireAssessmentSnapshotBackfillRun added in v0.2.0

func (*AssessmentSnapshotBackfillRepository) AdvanceAssessmentSnapshotBackfillRun added in v0.2.0

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

func (*AssessmentSnapshotBackfillRepository) CommitAssessmentSnapshotBackfillItem added in v0.2.0

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

func (*AssessmentSnapshotBackfillRepository) FinishAssessmentSnapshotBackfillRun added in v0.2.0

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

func (*AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillItem added in v0.2.0

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

func (*AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillRun added in v0.2.0

func (repository *AssessmentSnapshotBackfillRepository) GetAssessmentSnapshotBackfillRun(_ context.Context, tenantID, runID shared.ID) (ports.AssessmentSnapshotBackfillRun, 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() *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(_ context.Context, tenantID, snapshotID shared.ID) (*assessmentsnapshot.Snapshot, error)

func (*AssessmentSnapshotRepository) GetByRequestKey added in v0.2.0

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

func (*AssessmentSnapshotRepository) GetDefault added in v0.2.0

func (*AssessmentSnapshotRepository) ListAssessmentSnapshots added in v0.2.0

func (*AssessmentSnapshotRepository) ListByAssessment added in v0.2.0

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

type AssetStore added in v0.1.8

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

AssetStore is an in-memory ports.AssetRepository for dev and tests. It mirrors the Postgres store's tenant scoping and deterministic ordering so behaviour matches across backends.

func NewAssetStore added in v0.1.8

func NewAssetStore() *AssetStore

NewAssetStore returns an empty in-memory asset repository.

func (*AssetStore) AssignEngagementBusinessAsset added in v0.1.8

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

func (*AssetStore) CountBusinessAssetsByCriticality added in v0.2.0

func (s *AssetStore) CountBusinessAssetsByCriticality(_ context.Context, tenantID shared.ID) (map[asset.Criticality]int, error)

func (*AssetStore) CreateBusinessAsset added in v0.1.8

func (s *AssetStore) CreateBusinessAsset(_ context.Context, a *asset.BusinessAsset) error

func (*AssetStore) GetAssetByID added in v0.2.0

func (s *AssetStore) GetAssetByID(_ context.Context, tenantID, id shared.ID) (*asset.Asset, error)

GetAssetByID returns one canonical technical asset by server-issued ID. It is intentionally a narrow extension used by desired-state admission; the broad AssetRepository contract remains keyed by natural identity for normal observation/upsert flows. Invalid identifiers are validation errors; a valid but absent/cross-tenant asset is reported as ErrNotFound. The memory AssetStore is keyed by natural identity rather than ID, so duplicate canonical IDs are detected across the whole store and fail closed to mirror PostgreSQL's global fleet_assets primary-key invariant.

func (*AssetStore) GetAssetByKey added in v0.1.8

func (s *AssetStore) GetAssetByKey(_ 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 (*AssetStore) GetBusinessAssetByID added in v0.1.8

func (s *AssetStore) GetBusinessAssetByID(_ context.Context, tenantID, id shared.ID) (*asset.BusinessAsset, error)

func (*AssetStore) GetBusinessAssetByKey added in v0.1.8

func (s *AssetStore) GetBusinessAssetByKey(_ context.Context, tenantID shared.ID, key string) (*asset.BusinessAsset, error)

func (*AssetStore) ListAssets added in v0.1.8

func (s *AssetStore) ListAssets(_ context.Context, tenantID shared.ID) ([]*asset.Asset, error)

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

func (*AssetStore) ListBusinessAssetProjects added in v0.1.8

func (s *AssetStore) ListBusinessAssetProjects(_ context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)

func (*AssetStore) ListBusinessAssetTechnicalAssets added in v0.1.8

func (s *AssetStore) ListBusinessAssetTechnicalAssets(_ context.Context, tenantID, assetID shared.ID) ([]asset.ComponentMembership, error)

func (*AssetStore) ListBusinessAssets added in v0.1.8

func (s *AssetStore) ListBusinessAssets(_ context.Context, tenantID shared.ID) ([]*asset.BusinessAsset, error)

func (*AssetStore) ListBusinessAssetsPage added in v0.2.0

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

func (*AssetStore) ListEdges added in v0.1.8

func (s *AssetStore) ListEdges(_ context.Context, tenantID shared.ID) ([]*asset.Edge, error)

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

func (*AssetStore) ListEngagementsByBusinessAsset added in v0.1.8

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

func (*AssetStore) ReplaceBusinessAssetProjects added in v0.1.8

func (s *AssetStore) ReplaceBusinessAssetProjects(_ context.Context, tenantID, assetID shared.ID, links []asset.ComponentMembership) error

func (*AssetStore) ReplaceBusinessAssetTechnicalAssets added in v0.1.8

func (s *AssetStore) ReplaceBusinessAssetTechnicalAssets(_ context.Context, tenantID, assetID shared.ID, links []asset.ComponentMembership) error

func (*AssetStore) SetEngagementRepository added in v0.1.8

func (s *AssetStore) SetEngagementRepository(repo ports.EngagementRepository)

func (*AssetStore) UpdateBusinessAsset added in v0.1.8

func (s *AssetStore) UpdateBusinessAsset(_ context.Context, a *asset.BusinessAsset, expectedVersion int) error

func (*AssetStore) UpsertAsset added in v0.1.8

func (s *AssetStore) UpsertAsset(_ context.Context, a *asset.Asset) error

UpsertAsset stores a by its natural key, replacing any prior value for that key. A new host row past its reporting agent's cap is refused with shared.ErrForbidden, under the store lock, the way the fleet_assets trigger (migration 0132) refuses it in Postgres.

func (*AssetStore) UpsertEdge added in v0.1.8

func (s *AssetStore) UpsertEdge(_ context.Context, e *asset.Edge) error

UpsertEdge stores e by its natural key.

type AttackPathStore added in v0.1.8

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

AttackPathStore is the in-memory derived binding store used by dev and tests.

func NewAttackPathStore added in v0.1.8

func NewAttackPathStore() *AttackPathStore

func (*AttackPathStore) ListBindings added in v0.1.8

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

func (*AttackPathStore) ReplaceBindings added in v0.1.8

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

type BaselineStore added in v0.2.0

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

BaselineStore is the in-memory twin of the behavioral-baseline projection (Phase D / D5). It is tenant-bucketed and upholds the same upsert-by-(tenant,group) contract as the Postgres tier. Reached only through ports.BaselineStore.

func NewBaselineStore added in v0.2.0

func NewBaselineStore() *BaselineStore

NewBaselineStore creates an empty in-memory baseline store.

func (*BaselineStore) Load added in v0.2.0

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

func (*BaselineStore) Save added in v0.2.0

Save upserts a baseline record for its key.

type CloudObservationStore added in v0.1.8

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

CloudObservationStore tracks producer-owned active observations in local development.

func NewCloudObservationStore added in v0.1.8

func NewCloudObservationStore() *CloudObservationStore

func (*CloudObservationStore) FindingActive added in v0.1.8

func (s *CloudObservationStore) FindingActive(tenantID, engagementID, findingID shared.ID) bool

func (*CloudObservationStore) ReconcileCloudObservations added in v0.1.8

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

type CloudRunStore added in v0.1.8

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

CloudRunStore is the in-memory CSPM lifecycle store used in local development and tests.

func NewCloudRunStore added in v0.1.8

func NewCloudRunStore() *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(_ context.Context, tenantID, id shared.ID) (cloudposture.Run, error)

func (*CloudRunStore) SaveCloudRun added in v0.1.8

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

func (*CloudRunStore) SetQueue added in v0.1.8

func (s *CloudRunStore) SetQueue(queue ports.JobQueue)

type CommentRepository

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

CommentRepository is an in-memory per-finding comment thread (dev/tests).

func NewCommentRepository

func NewCommentRepository() *CommentRepository

NewCommentRepository returns an empty in-memory comment repository.

func (*CommentRepository) Add

func (*CommentRepository) ListByEngagementFinding

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

type ComponentInventoryStore added in v0.1.8

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

func NewComponentInventoryStore added in v0.1.8

func NewComponentInventoryStore(records ...sbom.ComponentRecord) *ComponentInventoryStore

func (*ComponentInventoryStore) ClaimInventoryWork added in v0.2.0

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

func (*ComponentInventoryStore) CompleteInventoryPublication added in v0.2.0

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

func (*ComponentInventoryStore) FinishInventoryWork added in v0.2.0

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

func (*ComponentInventoryStore) GetCurrentInventoryPublication added in v0.2.0

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

func (*ComponentInventoryStore) ListCurrentComponents added in v0.1.8

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

func (*ComponentInventoryStore) ListCurrentComponentsByEngagement added in v0.2.0

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

ListCurrentComponentsByEngagement returns the components of the engagement's LATEST SBOM (by SBOMCreatedAt, then SBOMID), deduped by ComponentID and ordered by ComponentID — the same latest-SBOM semantics as the Postgres twin. It is the by-engagement enumeration the running-vs-installed join needs to resolve a vulnerable ComponentID to a package name (the keyed ListCurrentComponents query cannot, since it requires a package/CPE key the caller does not have).

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)

func (*ComponentInventoryStore) Save added in v0.1.8

type CorrelationStateStore added in v0.2.0

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

CorrelationStateStore is the tenant-isolated in-memory twin of the durable two-phase work store. Its mutex makes each checkpoint/cursor/staging change one critical section.

func NewCorrelationStateStore added in v0.2.0

func NewCorrelationStateStore() *CorrelationStateStore

func (*CorrelationStateStore) BeginCorrelationSnapshot added in v0.2.0

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

func (*CorrelationStateStore) CommitCorrelationConsume added in v0.2.0

func (s *CorrelationStateStore) 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 (*CorrelationStateStore) ListStagedCorrelationSignals added in v0.2.0

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

func (*CorrelationStateStore) LoadCorrelationState added in v0.2.0

func (s *CorrelationStateStore) LoadCorrelationState(ctx context.Context, engagementID shared.ID, signalIDs []shared.ID, activeLimit int) (correlation.State, error)

func (*CorrelationStateStore) StageCorrelationSignals added in v0.2.0

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

type CoverageWindowStore added in v0.2.0

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

func NewCoverageWindowStore added in v0.2.0

func NewCoverageWindowStore() *CoverageWindowStore

func (*CoverageWindowStore) AppendCoverageWindow added in v0.2.0

func (*CoverageWindowStore) ListCoverageWindows added in v0.2.0

func (*CoverageWindowStore) ListCoverageWindowsBounded added in v0.2.0

func (s *CoverageWindowStore) 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 is the in-memory DAST run lifecycle store used in local development and tests.

func NewDASTRunStore added in v0.2.0

func NewDASTRunStore() *DASTRunStore

func (*DASTRunStore) EnqueueDASTRun added in v0.2.0

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

func (*DASTRunStore) FinishRun added in v0.2.0

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

FinishRun writes the terminal run only if the stored run is still running (compare-and-set).

func (*DASTRunStore) GetDASTRun added in v0.2.0

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

func (*DASTRunStore) SaveDASTRun added in v0.2.0

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

SaveDASTRun upserts the run but never moves a terminal row backward, matching the Postgres store: once a run is 'succeeded' or 'failed' its record is frozen, so a stale write cannot un-terminalize it.

func (*DASTRunStore) SetQueue added in v0.2.0

func (s *DASTRunStore) SetQueue(queue ports.JobQueue)

SetQueue binds the in-memory queue so EnqueueDASTRun can persist the run and enqueue its job together.

func (*DASTRunStore) StartRun added in v0.2.0

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

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

type DecisionStore

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

DecisionStore is the in-memory ports.DecisionStore (dev/tests). It mirrors the Postgres adapter: a monotonic per-session seq, idempotent on (session_id, action_id) for step decisions and a single stop per session (a re-record is a no-op, so a redelivered drive cannot fork the log).

func NewDecisionStore

func NewDecisionStore() *DecisionStore

NewDecisionStore returns an empty in-memory decision store.

func (*DecisionStore) AppendDecision

func (s *DecisionStore) AppendDecision(_ context.Context, d agent.AgentDecision) error

func (*DecisionStore) ListBySession

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

type DetectionProvenanceStore added in v0.2.0

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

DetectionProvenanceStore is a tenant-isolated in-memory current projection plus immutable history.

func NewDetectionProvenanceStore added in v0.2.0

func NewDetectionProvenanceStore() *DetectionProvenanceStore

func (*DetectionProvenanceStore) AdmitPending added in v0.2.0

func (*DetectionProvenanceStore) AppendTransition added in v0.2.0

func (s *DetectionProvenanceStore) AppendTransition(ctx context.Context, transition detectionprovenance.Transition) error

func (*DetectionProvenanceStore) Current added in v0.2.0

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

func (*DetectionProvenanceStore) ListCurrent added in v0.2.0

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

func (*DetectionProvenanceStore) ListPending added in v0.2.0

func (*DetectionProvenanceStore) ListReceivedTransitions added in v0.2.0

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

func (*DetectionProvenanceStore) ListTransitions added in v0.2.0

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

func (*DetectionProvenanceStore) LoadReceivedTransitions added in v0.2.0

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

type DetectionRecordStore added in v0.1.8

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

DetectionRecordStore is the in-memory detection-ledger projection used inline/in dev. Records are bucketed per tenant, so a read under one tenant can never observe another's — the same isolation the Postgres store gets from RLS. Within a tenant, a row is keyed by (engagement, id) to match the Postgres uniqueness key. It keeps deep copies so a caller mutating a returned record cannot corrupt stored state.

func NewDetectionRecordStore added in v0.1.8

func NewDetectionRecordStore() *DetectionRecordStore

NewDetectionRecordStore constructs the store.

func (*DetectionRecordStore) AppendDetection added in v0.1.8

func (s *DetectionRecordStore) AppendDetection(ctx context.Context, r detection.Record) error

AppendDetection stores one record, idempotent on (engagement, id), bucketed by the record's TenantID.

func (*DetectionRecordStore) ClassCountsByAsset added in v0.2.0

func (s *DetectionRecordStore) 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, under the ctx tenant. It feeds the behavior baseline's runtime-anomaly features (#822): the network / privilege / file per-class rates the process snapshot cannot carry.

func (*DetectionRecordStore) CorrelationHighWater added in v0.2.0

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

ListDetections returns the non-expired records for an engagement under the ctx tenant, oldest first.

func (*DetectionRecordStore) DeleteDetection added in v0.2.0

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

DeleteDetection removes one exact projection row. Repeating a successful deletion is a no-op.

func (*DetectionRecordStore) HasDetection added in v0.1.8

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

HasDetection reports whether a record with this id already exists in the given engagement under the ctx tenant, so ingest can skip an already-sealed detection on a retry (idempotent resume) rather than sealing it twice.

func (*DetectionRecordStore) LastBatchSequence added in v0.1.8

func (s *DetectionRecordStore) LastBatchSequence(ctx context.Context, agentID shared.ID) (uint64, error)

LastBatchSequence returns the highest batch sequence recorded for an agent under the ctx tenant.

func (*DetectionRecordStore) ListCorrelationSourcePage added in v0.2.0

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

func (*DetectionRecordStore) ListDetections added in v0.1.8

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

func (*DetectionRecordStore) ListExpiredDetections added in v0.2.0

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

ListExpiredDetections returns eligible record ids without mutating the projection.

func (*DetectionRecordStore) SetClock added in v0.1.8

func (s *DetectionRecordStore) SetClock(now func() time.Time)

SetClock overrides the store's clock (tests only), so expiry-on-read is deterministic.

type EmulationRunStore added in v0.1.8

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

EmulationRunStore is the in-memory emulation-run persistence used inline/in dev. It keeps a deep copy of each run's coverage so a caller mutating the returned slice cannot corrupt stored state.

func NewEmulationRunStore added in v0.1.8

func NewEmulationRunStore() *EmulationRunStore

NewEmulationRunStore constructs the store.

func (*EmulationRunStore) GetRun added in v0.1.8

func (s *EmulationRunStore) GetRun(_ context.Context, tenantID, id shared.ID) (demu.Run, error)

GetRun returns a copy of a stored run scoped to the tenant; a cross-tenant read sees ErrNotFound.

func (*EmulationRunStore) SaveRun added in v0.1.8

func (s *EmulationRunStore) SaveRun(_ context.Context, run demu.Run) error

SaveRun persists a run and its coverage records.

type EndpointProcessStore added in v0.2.0

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

EndpointProcessStore is the in-memory twin of the per-host process snapshot projection (B5). It is tenant-bucketed and upholds the same upsert-by-(tenant,asset,entity) contract as the Postgres tier. Reached only through ports.EndpointProcessStore.

func NewEndpointProcessStore added in v0.2.0

func NewEndpointProcessStore() *EndpointProcessStore

NewEndpointProcessStore creates an empty in-memory process snapshot store.

func (*EndpointProcessStore) ListRunningByAsset added in v0.2.0

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

ListRunningByAsset returns the running snapshots for an asset, ordered by EntityID.

func (*EndpointProcessStore) ReplaceRunningProcesses added in v0.2.0

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

ReplaceRunningProcesses makes the asset's running set exactly the reported snapshots: it upserts them and marks every other currently-running row for that asset not-running, so a process that exited between reports is retired instead of lingering as running=true.

func (*EndpointProcessStore) SaveProcesses added in v0.2.0

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

SaveProcesses upserts snapshots by (tenant, asset, entity).

type EndpointTimelineStore added in v0.2.0

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

EndpointTimelineStore is the in-memory twin of the durable endpoint State Timeline (Phase B / B7). It is tenant-bucketed and idempotent by (tenant, asset, EventID), upholding the same contract as the Postgres tier. Reached only through ports.EndpointTimelineStore.

func NewEndpointTimelineStore added in v0.2.0

func NewEndpointTimelineStore() *EndpointTimelineStore

NewEndpointTimelineStore creates an empty in-memory endpoint-timeline store.

func (*EndpointTimelineStore) AppendTimeline added in v0.2.0

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

AppendTimeline persists the transitions idempotently, skipping any whose EventID is already stored for its (tenant, asset). Every entry's TenantID must equal the context tenant.

func (*EndpointTimelineStore) LoadTimelineEntries added in v0.2.0

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

LoadTimelineEntries returns the requested source entries without widening correlation reads to a time range.

func (*EndpointTimelineStore) QueryTimeline added in v0.2.0

QueryTimeline returns the stored transitions matching q, ordered by (OccurredAt, EventID).

type EngagementRepository

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

EngagementRepository is a goroutine-safe in-memory engagement store.

func NewEngagementRepository

func NewEngagementRepository() *EngagementRepository

NewEngagementRepository returns an empty in-memory repository.

func (*EngagementRepository) Create

func (*EngagementRepository) Delete

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

Delete removes an engagement (idempotent). In Postgres the FK cascade removes children; in memory other stores are independent, but import rollback only needs the engagement gone so a re-import isn't blocked.

func (*EngagementRepository) GetByHostAssetID added in v0.2.0

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

func (*EngagementRepository) GetByID

func (*EngagementRepository) GetByIDInTenant

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

GetByIDInTenant loads an engagement scoped to tenantID. Empty input normalizes to the non-empty default tenant; it is never a wildcard.

func (*EngagementRepository) GetByProjectID

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

func (*EngagementRepository) List

func (*EngagementRepository) ListAssessmentCycleBackfillEngagements added in v0.2.0

func (r *EngagementRepository) ListAssessmentCycleBackfillEngagements(_ context.Context, tenantID, after shared.ID, snapshotAt time.Time, limit int) ([]*engagement.Engagement, 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(_ context.Context, tenantID shared.ID) ([]*engagement.Engagement, error)

ListHostEngagements returns the tenant's hidden host vulnerability contexts.

func (*EngagementRepository) ListProjectEngagements added in v0.1.8

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

func (*EngagementRepository) ListPromotionReconciliationScopes added in v0.1.8

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

ListPromotionReconciliationScopes returns every non-project engagement for process-local recovery. It is only wired by the API composition root.

func (*EngagementRepository) ListReconciliationEngagements added in v0.1.8

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

func (*EngagementRepository) ListTenantIDs added in v0.1.8

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

func (*EngagementRepository) ProjectContexts

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

func (*EngagementRepository) Update

type EngagementSourceRepository added in v0.2.0

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

func NewEngagementSourceRepository added in v0.2.0

func NewEngagementSourceRepository() *EngagementSourceRepository

func (*EngagementSourceRepository) Create added in v0.2.0

func (*EngagementSourceRepository) Delete added in v0.2.0

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

func (*EngagementSourceRepository) Get added in v0.2.0

func (*EngagementSourceRepository) GetByLocator added in v0.2.0

func (r *EngagementSourceRepository) GetByLocator(_ 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, engID, versionID shared.ID) (sourcepackage.Package, error)

func (*EngagementSourceRepository) ObjectUnreferenced added in v0.2.0

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

type EvidenceStore

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

EvidenceStore is an in-memory append-only evidence ledger.

func NewEvidenceStore

func NewEvidenceStore() *EvidenceStore

NewEvidenceStore returns an empty in-memory evidence store.

func (*EvidenceStore) Append

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

func (*EvidenceStore) Head

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

func (*EvidenceStore) ListByEngagement

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

func (*EvidenceStore) LookupSealedForFinding added in v0.1.8

func (s *EvidenceStore) LookupSealedForFinding(_ 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. Used for crash-recoverable evidence reservation in the promotion recorder.

type ExploitationChainStore added in v0.1.8

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

ExploitationChainStore is the in-memory chain persistence used inline/in dev. It keeps a deep copy of each chain so a caller mutating the returned pointer cannot corrupt the stored state — the same isolation the Postgres store gets for free from serialization.

func NewExploitationChainStore added in v0.1.8

func NewExploitationChainStore() *ExploitationChainStore

NewExploitationChainStore constructs the store.

func (*ExploitationChainStore) GetChain added in v0.1.8

func (s *ExploitationChainStore) GetChain(_ context.Context, tenantID, id shared.ID) (*dexploit.Chain, error)

GetChain returns a copy of the stored chain, scoped to the tenant. Cross-tenant reads see ErrNotFound rather than another tenant's chain.

func (*ExploitationChainStore) SaveChain added in v0.1.8

func (s *ExploitationChainStore) SaveChain(_ context.Context, chain *dexploit.Chain) error

SaveChain upserts a chain and its steps.

type FindingLineageBackfillRepository added in v0.2.0

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

func NewFindingLineageBackfillRepository added in v0.2.0

func NewFindingLineageBackfillRepository() *FindingLineageBackfillRepository

func (*FindingLineageBackfillRepository) AcquireFindingLineageBackfillRun added in v0.2.0

func (*FindingLineageBackfillRepository) AdvanceFindingLineageBackfillRun added in v0.2.0

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

func (*FindingLineageBackfillRepository) CommitFindingLineageBackfillItem added in v0.2.0

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

func (*FindingLineageBackfillRepository) FinishFindingLineageBackfillRun added in v0.2.0

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

func (*FindingLineageBackfillRepository) GetFindingLineageBackfillItem added in v0.2.0

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

func (*FindingLineageBackfillRepository) GetFindingLineageBackfillRun added in v0.2.0

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

func (*FindingLineageBackfillRepository) ListFindingLineageBackfillSources added in v0.2.0

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

func (*FindingLineageBackfillRepository) SetSources added in v0.2.0

func (repository *FindingLineageBackfillRepository) SetSources(tenantID shared.ID, sources []ports.FindingLineageBackfillSourceRow)

type FindingLineageRepository added in v0.2.0

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

func NewFindingLineageRepository added in v0.2.0

func NewFindingLineageRepository() *FindingLineageRepository

func (*FindingLineageRepository) AppendAlias added in v0.2.0

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

func (*FindingLineageRepository) AppendObservation added in v0.2.0

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

func (*FindingLineageRepository) AppendOverrideCAS added in v0.2.0

func (*FindingLineageRepository) AppendSkip added in v0.2.0

func (*FindingLineageRepository) CreateCandidate added in v0.2.0

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

func (*FindingLineageRepository) CreateIdentityWithObservation added in v0.2.0

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

func (*FindingLineageRepository) FindIdentitiesByAlias added in v0.2.0

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

func (*FindingLineageRepository) FindIdentitiesByFingerprint added in v0.2.0

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

func (*FindingLineageRepository) FindIdentitiesByProducerID added in v0.2.0

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

func (*FindingLineageRepository) GetActiveOverride added in v0.2.0

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

func (*FindingLineageRepository) GetCandidate added in v0.2.0

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

func (*FindingLineageRepository) GetIdentity added in v0.2.0

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

func (*FindingLineageRepository) GetObservation added in v0.2.0

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

func (*FindingLineageRepository) GetObservationBySource added in v0.2.0

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

func (*FindingLineageRepository) ListActiveOverridesBySnapshot added in v0.2.0

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

func (*FindingLineageRepository) ListCandidateResolutions added in v0.2.0

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

func (*FindingLineageRepository) ListObservationsBySnapshot added in v0.2.0

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

func (*FindingLineageRepository) ListOpenCandidatesBySnapshot added in v0.2.0

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

func (*FindingLineageRepository) ListOverrideEvents added in v0.2.0

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

func (*FindingLineageRepository) ListSkipsBySnapshot added in v0.2.0

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

func (*FindingLineageRepository) LockCorrelationNamespace added in v0.2.0

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

func (*FindingLineageRepository) ResolveCandidateCAS added in v0.2.0

type FindingRepository

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

FindingRepository is an in-memory finding store (dev/tests), deduped per engagement by dedup key. Replaced by Postgres when a DB is configured.

func NewFindingRepository

func NewFindingRepository() *FindingRepository

NewFindingRepository returns an empty in-memory finding repository.

func (*FindingRepository) CheckVulnerabilityPrimaryFinding added in v0.2.0

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

func (*FindingRepository) ClaimFindingProjection added in v0.1.8

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

func (*FindingRepository) GetByEngagementAndDedupKey added in v0.2.0

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

GetByEngagementAndDedupKey resolves the canonical repository key directly.

func (*FindingRepository) GetByEngagementAndID added in v0.1.8

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

GetByEngagementAndID loads a single finding by engagement and finding ID.

func (*FindingRepository) LinkVulnerabilityFindingOccurrence added in v0.2.0

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

func (*FindingRepository) ListByEngagement

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

func (*FindingRepository) ListPublishableByEngagement

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

ListPublishableByEngagement returns only the engagement's findings that clear the evidence gate, reusing the single domain rule finding.Publishable.

func (*FindingRepository) MapVulnerabilityPrimaryFinding added in v0.2.0

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

func (*FindingRepository) SetAssignee

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

SetAssignee sets a finding's assignee with the same optimistic-concurrency guard.

func (*FindingRepository) SetCloudObservationStore added in v0.1.8

func (r *FindingRepository) SetCloudObservationStore(store *CloudObservationStore)

ClaimFindingProjection atomically selects the CapSAST projection mode for a judgment.

func (*FindingRepository) SetEvidenceScore

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

SetEvidenceScore sets a finding's evidence score with the same optimistic-concurrency guard as UpdateStatus (the adversarial-verdict path). Returns shared.ErrConflict on a version mismatch, shared.ErrNotFound if absent.

func (*FindingRepository) SummarizeOpenFindingsByEngagements added in v0.2.0

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

SummarizeOpenFindingsByEngagements counts open findings of every kind per engagement.

func (*FindingRepository) SummarizeVulnerabilitiesByEngagements added in v0.2.0

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

ListByEngagement returns the engagement's findings, highest risk first (KEV -> EPSS x CVSS). SummarizeVulnerabilitiesByEngagements counts open SCA vulnerability findings per engagement.

func (*FindingRepository) UpdateStatus

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

UpdateStatus sets a finding's triage status with optimistic concurrency (expectedVersion must match the stored version), bumping the version. Returns shared.ErrConflict on a version mismatch, shared.ErrNotFound if absent.

func (*FindingRepository) Upsert

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

Upsert inserts or updates findings, deduped by (engagement, dedup key). On update it preserves the existing triage status + created timestamp.

type FleetAgentStore added in v0.1.8

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

FleetAgentStore is an in-memory ports.FleetAgentStore for dev and tests.

func NewFleetAgentStore added in v0.1.8

func NewFleetAgentStore() *FleetAgentStore

NewFleetAgentStore returns an empty in-memory fleet agent store.

func (*FleetAgentStore) ConsumeEnrolToken added in v0.1.8

func (s *FleetAgentStore) ConsumeEnrolToken(_ context.Context, tenantID shared.ID, hash string, now time.Time) (*fleetagent.EnrolToken, error)

func (*FleetAgentStore) CreateAgent added in v0.1.8

func (s *FleetAgentStore) CreateAgent(_ context.Context, a *fleetagent.Agent) error

func (*FleetAgentStore) CreateEnrolToken added in v0.1.8

func (s *FleetAgentStore) CreateEnrolToken(_ context.Context, t *fleetagent.EnrolToken) error

func (*FleetAgentStore) Decommission added in v0.1.8

func (s *FleetAgentStore) Decommission(_ context.Context, tenantID, id shared.ID, now time.Time) error

func (*FleetAgentStore) GetAgent added in v0.1.8

func (s *FleetAgentStore) GetAgent(_ context.Context, tenantID, id shared.ID) (*fleetagent.Agent, error)

func (*FleetAgentStore) Heartbeat added in v0.1.8

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

func (*FleetAgentStore) ListAgents added in v0.1.8

func (s *FleetAgentStore) ListAgents(_ context.Context, tenantID shared.ID) ([]*fleetagent.Agent, error)

func (*FleetAgentStore) Revoke added in v0.1.8

func (s *FleetAgentStore) Revoke(_ context.Context, tenantID, id, by shared.ID, reason string, now time.Time) error

func (*FleetAgentStore) SetFingerprint added in v0.1.8

func (s *FleetAgentStore) SetFingerprint(_ context.Context, tenantID, id shared.ID, fingerprint string, now time.Time) error

type FleetDesiredStore added in v0.2.0

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

FleetDesiredStore is the in-memory desired-state store used by development and focused tests.

func NewFleetDesiredStore added in v0.2.0

func NewFleetDesiredStore() *FleetDesiredStore

func (*FleetDesiredStore) Delete added in v0.2.0

func (s *FleetDesiredStore) Delete(_ context.Context, tenantID, assetID, expectedPolicyID shared.ID, expectedVersion int64) error

Delete clears only the lifecycle/version the caller observed. Absence is idempotent; either a newer version or a delete/recreate lifecycle returns ErrConflict and remains untouched.

func (*FleetDesiredStore) Get added in v0.2.0

func (s *FleetDesiredStore) Get(_ context.Context, tenantID, assetID shared.ID) (*fleetdesired.State, error)

func (*FleetDesiredStore) List added in v0.2.0

func (s *FleetDesiredStore) List(_ context.Context, tenantID shared.ID) ([]*fleetdesired.State, error)

func (*FleetDesiredStore) Put added in v0.2.0

Put applies a lifecycle-aware CAS. A new PolicyID may only create an absent row at version 1; updates must retain the stored PolicyID and advance version by exactly one. PolicyID is unique within a tenant so a lifecycle identifier can never alias another asset's policy.

type FleetRolloutStore added in v0.1.8

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

FleetRolloutStore is the in-memory rollout plan store.

func NewFleetRolloutStore added in v0.1.8

func NewFleetRolloutStore() *FleetRolloutStore

NewFleetRolloutStore returns an empty store.

func (*FleetRolloutStore) Get added in v0.1.8

func (s *FleetRolloutStore) Get(_ context.Context, tenantID shared.ID, channel string) (*fleetrollout.Plan, error)

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

func (*FleetRolloutStore) Put added in v0.1.8

Put stores the plan, replacing any existing one for the same (tenant, channel).

type IdentityStore added in v0.2.0

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

IdentityStore is a race-safe in-memory ports.IdentityStore for development and tests.

func NewIdentityStore added in v0.2.0

func NewIdentityStore(users ports.UserRepository) (*IdentityStore, error)

NewIdentityStore returns an empty store linked to the supplied user repository.

func (*IdentityStore) ConsumeAuthorizationTransaction added in v0.2.0

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

func (*IdentityStore) CreateAuthorizationTransaction added in v0.2.0

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

func (*IdentityStore) CreateExternalIdentity added in v0.2.0

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

func (*IdentityStore) CreateSession added in v0.2.0

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

func (*IdentityStore) GetExternalIdentity added in v0.2.0

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

func (*IdentityStore) GetSessionByTokenHash added in v0.2.0

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

func (*IdentityStore) RevokeSession added in v0.2.0

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

func (*IdentityStore) RotateSession added in v0.2.0

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

RotateSession creates replacement and revokes previous under one lock.

type ImportReceiptStore added in v0.2.0

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

ImportReceiptStore is the in-memory conformance adapter for durable import receipts. It is intentionally not wired into production; the PostgreSQL implementation is the durable adapter.

func NewImportReceiptStore added in v0.2.0

func NewImportReceiptStore() *ImportReceiptStore

NewImportReceiptStore returns an empty receipt store.

func (*ImportReceiptStore) CreateOrGet added in v0.2.0

func (s *ImportReceiptStore) CreateOrGet(_ context.Context, tenantID shared.ID, receipt importreceipt.Receipt) (importreceipt.Receipt, bool, error)

CreateOrGet atomically creates a receipt or resolves the receipt already owning the idempotency key.

func (*ImportReceiptStore) Finalize added in v0.2.0

func (s *ImportReceiptStore) Finalize(_ context.Context, tenantID, receiptID shared.ID, outcome importreceipt.Outcome, counters importreceipt.Counters, updatedAt time.Time) (importreceipt.Receipt, error)

Finalize atomically performs the sole receipt mutation.

func (*ImportReceiptStore) GetByIdentity added in v0.2.0

func (s *ImportReceiptStore) GetByIdentity(_ context.Context, tenantID, engagementID shared.ID, sourceIdentity, digest string) (importreceipt.Receipt, error)

GetByIdentity returns only a receipt visible in the requested tenant partition.

type ImportedFindingStore added in v0.1.8

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

ImportedFindingStore is the in-memory third-party finding store.

It is EPHEMERAL: everything it holds is lost on restart. It backs tests and the CLI's validate mode; the server uses the Postgres store, which is what makes the audit entry for an ingest true.

func NewImportedFindingStore added in v0.1.8

func NewImportedFindingStore() *ImportedFindingStore

NewImportedFindingStore returns an empty store.

func (*ImportedFindingStore) ExistsDigest added in v0.1.8

func (s *ImportedFindingStore) ExistsDigest(_ 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 (*ImportedFindingStore) ListByEngagement added in v0.1.8

func (s *ImportedFindingStore) ListByEngagement(_ context.Context, tenantID, engagementID shared.ID) ([]importedfinding.ImportedFinding, error)

ListByEngagement returns one engagement's imported findings, deterministically ordered.

func (*ImportedFindingStore) Save added in v0.1.8

Save persists a batch ATOMICALLY: the whole delta is built and validated first, so a finding that fails validation aborts the batch without leaving earlier findings — and without recording the document digest, which would make a retry look like a clean deduplicated ingest while the un-persisted tail was permanently lost.

type ImportedSBOMStore

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

ImportedSBOMStore keeps the active imported SBOM per engagement in memory.

func NewImportedSBOMStore

func NewImportedSBOMStore() *ImportedSBOMStore

NewImportedSBOMStore returns an empty imported-SBOM store.

func (*ImportedSBOMStore) LatestByEngagement

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

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)

LatestByEngagement returns a copy of the active imported SBOM. MetadataByEngagements returns the metadata of every listed engagement's active record.

func (*ImportedSBOMStore) SaveActive

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

SaveActive stores a copy of the active imported SBOM for its engagement.

type IncidentEventStore added in v0.2.0

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

IncidentEventStore is the in-memory twin of the append-only incident log and immutable merge index.

func NewIncidentEventStore added in v0.2.0

func NewIncidentEventStore() *IncidentEventStore

func (*IncidentEventStore) AppendEvents added in v0.2.0

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

func (*IncidentEventStore) ListIncidentIDs added in v0.2.0

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

func (*IncidentEventStore) ListMergeEdges added in v0.2.0

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

func (*IncidentEventStore) LoadEvents added in v0.2.0

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

func (*IncidentEventStore) ResolveCanonicalID added in v0.2.0

func (s *IncidentEventStore) ResolveCanonicalID(ctx context.Context, id shared.ID) (shared.ID, error)

type IntegrationStore added in v0.2.0

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

func NewIntegrationStore added in v0.2.0

func NewIntegrationStore(queue ports.JobQueue, cipher *vault.Cipher, clock ports.Clock, audit ports.AuditLogger) *IntegrationStore

func (*IntegrationStore) ArchiveIntegration added in v0.2.0

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

func (*IntegrationStore) BeginIntegrationOperation added in v0.2.0

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

func (*IntegrationStore) CancelIntegrationOperation added in v0.2.0

func (store *IntegrationStore) CancelIntegrationOperation(ctx context.Context, id shared.ID, finishedAt time.Time, audit ports.AuditEntry) (integration.Operation, 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) (integration.Operation, 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) (integration.Operation, error)

func (*IntegrationStore) GetIntegration added in v0.2.0

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

func (*IntegrationStore) GetIntegrationOperation added in v0.2.0

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

func (*IntegrationStore) IntegrationCredentialConfigured added in v0.2.0

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

func (*IntegrationStore) ListDueIntegrations added in v0.2.0

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

func (*IntegrationStore) ListIntegrationBindings added in v0.2.0

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

func (*IntegrationStore) ListIntegrationExternalRuns added in v0.2.0

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

func (*IntegrationStore) ListIntegrationOperations added in v0.2.0

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

func (*IntegrationStore) ListIntegrations added in v0.2.0

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

func (*IntegrationStore) PutIntegrationCredential added in v0.2.0

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

func (*IntegrationStore) ResolveIntegrationCredential added in v0.2.0

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

func (*IntegrationStore) SetIntegrationEnabled added in v0.2.0

func (store *IntegrationStore) SetIntegrationEnabled(ctx context.Context, id shared.ID, enabled bool, expectedVersion int, audit ports.AuditEntry) (integration.Integration, 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) (integration.Integration, error)

type JobQueue

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

JobQueue is an in-memory ports.JobQueue for dev/single-process + tests, with the same visibility-lease + reclaim semantics as the Postgres adapter (not durable across restarts). Time is injected so tests can exercise lease expiry deterministically.

func NewJobQueue

func NewJobQueue(ids ports.IDGenerator, now func() time.Time) *JobQueue

NewJobQueue returns an in-memory job queue. now may be nil (uses time.Now).

func (*JobQueue) AggregateJobQueueStats added in v0.2.0

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

AggregateJobQueueStats aggregates every in-memory job without tenant scoping. It is intended solely for the operator metrics collector and never returns tenant labels.

func (*JobQueue) Claim

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

func (*JobQueue) Complete

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

func (*JobQueue) Deadletter

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

func (*JobQueue) Depth

func (q *JobQueue) Depth(_ context.Context, kinds ...string) (int, error)

Depth counts not-yet-terminal jobs (queued or claimed); 'done'/'failed' are excluded. Optional kind filter (empty = any). Mirrors the Postgres adapter.

func (*JobQueue) Enqueue

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

func (*JobQueue) Fail

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

func (*JobQueue) Heartbeat

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

func (*JobQueue) Invalidate added in v0.2.0

func (q *JobQueue) Invalidate(_ context.Context, id string) error

Invalidate makes queued or claimed work permanently unclaimable and advances its fence so a worker holding an older claim cannot acknowledge it.

func (*JobQueue) JobStatus added in v0.1.8

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

func (*JobQueue) Retry added in v0.2.0

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

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

func (*JobQueue) Stats added in v0.1.8

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

type JudgmentStore

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

JudgmentStore is the in-memory judgment repository (dev/tests). It mirrors the Postgres adapter to come: SetScoreState is the ONLY score/state mover and is guarded by optimistic concurrency (expectedVersion → shared.ErrConflict on mismatch), the same discipline as the finding repo's SetEvidenceScore. The score mover is deliberately not exposed on a broad read port.

func NewJudgmentStore

func NewJudgmentStore() *JudgmentStore

NewJudgmentStore returns an empty in-memory judgment store.

func (*JudgmentStore) AcknowledgeJudgmentAudit added in v0.1.8

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

func (*JudgmentStore) GetByID added in v0.2.0

func (s *JudgmentStore) GetByID(_ context.Context, engagementID, id shared.ID) (judgment.Judgment, error)

GetByID provides a bounded lookup for migration/backfill consumers.

func (*JudgmentStore) ListByEngagement

func (s *JudgmentStore) ListByEngagement(_ context.Context, engagementID shared.ID) ([]judgment.Judgment, error)

ListByEngagement returns a copy of the engagement's judgments.

func (*JudgmentStore) ListBySubject

func (s *JudgmentStore) ListBySubject(_ context.Context, engagementID, subjectID shared.ID) ([]judgment.Judgment, error)

ListBySubject returns the engagement's judgments about a given subject id.

func (*JudgmentStore) ListPendingJudgmentAudits added in v0.1.8

func (s *JudgmentStore) ListPendingJudgmentAudits(ctx context.Context, engagementID shared.ID) ([]ports.PendingJudgmentAudit, error)

func (*JudgmentStore) Save

Save inserts or replaces a judgment within its engagement (idempotent by id).

func (*JudgmentStore) SaveWithProposalAudit added in v0.1.8

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

SaveWithProposalAudit inserts a proposal and its immutable pending audit payload under one lock. Exact replays do not duplicate either record.

func (*JudgmentStore) SetScoreState

func (s *JudgmentStore) SetScoreState(_ context.Context, engagementID, id shared.ID, score int, state judgment.State, expectedVersion int) (judgment.Judgment, error)

SetScoreState moves a judgment's score + state under optimistic concurrency. A version mismatch returns shared.ErrConflict (lost-update guard); an unknown id returns shared.ErrNotFound.

func (*JudgmentStore) SetVerdictState added in v0.1.8

func (s *JudgmentStore) SetVerdictState(_ 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 (*JudgmentStore) SetVerdictStateWithAudit added in v0.1.8

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

SetVerdictStateWithAudit persists a sealed verdict and its immutable pending audit payload together, so a crash can only leave a recoverable pending delivery.

type LeaderStore added in v0.1.8

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

LeaderStore is an in-memory ports.LeaderStore for dev and single-process tests. It mirrors the Postgres fenced-lease semantics under a mutex.

func NewLeaderStore added in v0.1.8

func NewLeaderStore() *LeaderStore

NewLeaderStore returns an empty in-memory leader store.

func (*LeaderStore) Acquire added in v0.1.8

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

Acquire takes the lease when it is free (expired) or already held by holder, renewing the term; a takeover from an expired holder bumps the fence. Otherwise another instance holds it.

func (*LeaderStore) Resign added in v0.1.8

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

Resign releases the lease if held by holder. It expires the lease (clears the holder, sets the term to now) rather than dropping the row, so the fence survives and stays monotonic across a graceful handover, matching the Postgres store.

type LegalHoldStore added in v0.2.0

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

LegalHoldStore is an in-memory legal-hold store (#635), tenant-scoped, keeping an append-only history per (tenant, engagement) so hold+release is auditable.

func NewLegalHoldStore added in v0.2.0

func NewLegalHoldStore() *LegalHoldStore

func (*LegalHoldStore) IsHeld added in v0.2.0

func (s *LegalHoldStore) IsHeld(ctx context.Context, engagementID shared.ID) (bool, error)

func (*LegalHoldStore) ListActive added in v0.2.0

func (s *LegalHoldStore) ListActive(ctx context.Context) ([]legalhold.Hold, error)

func (*LegalHoldStore) Place added in v0.2.0

func (*LegalHoldStore) Release added in v0.2.0

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

type MissingIntegrationAnalysisMatcher added in v0.2.0

type MissingIntegrationAnalysisMatcher struct{}

func (MissingIntegrationAnalysisMatcher) MatchIntegrationAnalysis added in v0.2.0

type PlanStore

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

PlanStore is the in-memory ports.PlanStore (dev/tests). It mirrors the Postgres adapter's contract: one plan per session, and SavePlan is an optimistic-concurrency CAS on the revision (a stale revision returns ErrConflict). Plans are deep-copied in and out so a caller mutating its own Plan value cannot retroactively change stored state.

func NewPlanStore

func NewPlanStore() *PlanStore

NewPlanStore returns an empty in-memory plan store.

func (*PlanStore) CreatePlan

func (s *PlanStore) CreatePlan(_ context.Context, p agent.Plan) error

func (*PlanStore) GetBySession

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

func (*PlanStore) SavePlan

func (s *PlanStore) SavePlan(_ context.Context, p agent.Plan) error

type PrivacyPolicyStore added in v0.2.0

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

PrivacyPolicyStore is the in-memory adapter for immutable policy history and the tenant's independently mutable active policy pointer.

func NewPrivacyPolicyStore added in v0.2.0

func NewPrivacyPolicyStore() *PrivacyPolicyStore

func (*PrivacyPolicyStore) AcknowledgeFleetAudit added in v0.2.0

func (s *PrivacyPolicyStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

AcknowledgeFleetAudit retires one of the bound tenant's obligations. The tenant comes from context for the same reason as the pending listing above.

func (*PrivacyPolicyStore) ActivatePrivacyPolicy added in v0.2.0

func (s *PrivacyPolicyStore) ActivatePrivacyPolicy(
	ctx context.Context,
	activation privacy.Activation,
) (privacy.Activation, error)

func (*PrivacyPolicyStore) ActivatePrivacyPolicyWithAudit added in v0.2.0

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

func (*PrivacyPolicyStore) ActivePrivacyPolicy added in v0.2.0

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

func (*PrivacyPolicyStore) ListPendingFleetAudits added in v0.2.0

func (s *PrivacyPolicyStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

ListPendingFleetAudits returns the bound tenant's undelivered audit obligations. It takes the tenant from context directly: there is no caller-supplied tenant to cross-check here, and routing an empty id through the request/context comparison would resolve to the DEFAULT tenant and reject every other tenant's own sweep.

func (*PrivacyPolicyStore) PrivacyPolicyActivationHistory added in v0.2.0

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

func (*PrivacyPolicyStore) PrivacyPolicyByDigest added in v0.2.0

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

func (*PrivacyPolicyStore) PrivacyPolicyHistory added in v0.2.0

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

func (*PrivacyPolicyStore) PutPrivacyPolicy added in v0.2.0

func (s *PrivacyPolicyStore) PutPrivacyPolicy(
	ctx context.Context,
	assignment privacy.Assignment,
) (bool, error)

type ProjectAnalysisStore

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

func NewProjectAnalysisStore

func NewProjectAnalysisStore() *ProjectAnalysisStore

func (*ProjectAnalysisStore) Branches added in v0.2.0

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

Branches returns the distinct branch values recorded for the project, sorted and stable.

func (*ProjectAnalysisStore) CurrentAnalysisHotspotSummary

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

func (*ProjectAnalysisStore) CurrentFindingStatuses added in v0.1.8

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

func (*ProjectAnalysisStore) Get

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

func (*ProjectAnalysisStore) GetHotspot

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

func (*ProjectAnalysisStore) GetIssue

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

func (*ProjectAnalysisStore) HotspotHistory

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

func (*ProjectAnalysisStore) IssueHistory

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

func (*ProjectAnalysisStore) LatestForProjects

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

func (*ProjectAnalysisStore) LatestWithResult

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

func (*ProjectAnalysisStore) List

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

func (*ProjectAnalysisStore) ListAnalysisHotspots

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

func (*ProjectAnalysisStore) ListHotspots

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

func (*ProjectAnalysisStore) ListIssues

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

func (*ProjectAnalysisStore) PruneBranchAnalyses added in v0.2.0

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

func (*ProjectAnalysisStore) ResolvedIssueKeys

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

func (*ProjectAnalysisStore) Save

func (*ProjectAnalysisStore) SaveWithResult

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

func (*ProjectAnalysisStore) SaveWithResultAndHotspots

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

SaveWithResultAndHotspots satisfies ports.ProjectAnalysisProjectionStore; it is the hotspot-only projection (no issues), delegating to the combined path.

func (*ProjectAnalysisStore) SaveWithResultAndProjections

func (s *ProjectAnalysisStore) SaveWithResultAndProjections(_ context.Context, analysis projectanalysis.Analysis, result []byte, candidates []hotspot.Candidate, issues []issue.Candidate) error

func (*ProjectAnalysisStore) TransitionHotspot

func (*ProjectAnalysisStore) TransitionIssue

type ProjectRepository

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

ProjectRepository is a goroutine-safe in-memory project store.

func NewProjectRepository

func NewProjectRepository() *ProjectRepository

func (*ProjectRepository) AssignProfile

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

func (*ProjectRepository) CountByGate

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

func (*ProjectRepository) Create

func (*ProjectRepository) DeleteByKey

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

func (*ProjectRepository) GetByID

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

func (*ProjectRepository) GetByKey

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

func (*ProjectRepository) List

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

func (*ProjectRepository) SetPullRequestDecoration added in v0.2.0

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

func (*ProjectRepository) UpdateGate

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

type PromotionStore added in v0.1.8

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

func NewPromotionStore added in v0.1.8

func NewPromotionStore(findingRepo *FindingRepository, engagementReader ports.EngagementOwnershipReader) (*PromotionStore, error)

NewPromotionStore returns an empty in-memory promotion store. The findingRepo must be the SAME instance used by the service layer so that CAS operates on the same data. The engagementReader verifies that the engagement belongs to the context tenant before any mutation.

func (*PromotionStore) Apply added in v0.1.8

func (s *PromotionStore) Apply(ctx context.Context, engagementID, findingID shared.ID, cmd ports.PromotionCommand) (finding.Finding, 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 (matching fingerprint with different judgment, or CAS mismatch).

func (*PromotionStore) FindByJudgment added in v0.1.8

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

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

func (*PromotionStore) LatestByFinding added in v0.1.8

func (s *PromotionStore) LatestByFinding(ctx context.Context, engagementID, findingID shared.ID) (promotion.PromotionEvent, bool, error)

LatestByFinding returns the most recent promotion event for a finding, or (zero, false) if none exist. The returned event is a defensive copy. Results are scoped to the given engagement.

func (*PromotionStore) ListByFinding added in v0.1.8

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

ListByFinding returns all promotion events for a finding, oldest first. The returned slice and each event's Inputs are defensive copies; callers may modify them safely. Results are scoped to the given engagement.

func (*PromotionStore) ListPendingAudits added in v0.1.8

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

ListPendingAudits returns applied promotion events whose required audit has not been durably acknowledged. Results are oldest first.

func (*PromotionStore) MarkAuditComplete added in v0.1.8

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

MarkAuditComplete durably acknowledges an applied event's audit record.

type PurpleStore added in v0.1.8

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

PurpleStore is the in-memory purple-coverage store, tenant-bucketed so one tenant's coverage is never visible to another. Coverage is keyed (run, technique) so a re-computation of the same run replaces its records in place and coverage across runs forms a trend.

func NewPurpleStore added in v0.1.8

func NewPurpleStore() *PurpleStore

NewPurpleStore constructs the store.

func (*PurpleStore) ListByEngagement added in v0.1.8

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

ListByEngagement returns all coverage for an engagement in the ctx tenant, oldest first (by ComputedAt, then run/technique) so a trend across runs is stable.

func (*PurpleStore) ListByRun added in v0.1.8

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

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

func (*PurpleStore) SaveCoverage added in v0.1.8

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

SaveCoverage upserts a run's coverage records under the authenticated tenant. A record claiming a different tenant is refused so a computation cannot write into another tenant's ledger.

type QualityGateMutator

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

QualityGateMutator makes managed-gate writes atomic with their audit record in memory.

func NewQualityGateMutator

func NewQualityGateMutator(gates *QualityGateStore, projects *ProjectRepository, audit ports.AuditLogger) *QualityGateMutator

func (*QualityGateMutator) AssignProjectGate

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

func (*QualityGateMutator) CreateGate

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

func (*QualityGateMutator) CreateProjectWithGate

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

func (*QualityGateMutator) DeleteGate

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

func (*QualityGateMutator) UpdateGate

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

type QualityGateStore

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

QualityGateStore is a goroutine-safe in-memory custom quality-gate store.

func NewQualityGateStore

func NewQualityGateStore() *QualityGateStore

func (*QualityGateStore) Create

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

func (*QualityGateStore) Delete

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

func (*QualityGateStore) DeleteIfUnassigned

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

func (*QualityGateStore) Get

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

func (*QualityGateStore) List

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

func (*QualityGateStore) Update

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

type QualityProfileStore

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

QualityProfileStore is a goroutine-safe in-memory custom quality-profile store.

func NewQualityProfileStore

func NewQualityProfileStore() *QualityProfileStore

func (*QualityProfileStore) Create

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

func (*QualityProfileStore) Delete

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

func (*QualityProfileStore) Get

func (*QualityProfileStore) List

func (*QualityProfileStore) Update

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

type ReconRunRepository

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

ReconRunRepository is an in-memory ports.ReconRunStore for dev/tests.

func NewReconRunRepository

func NewReconRunRepository() *ReconRunRepository

NewReconRunRepository returns an empty in-memory recon-run store.

func (*ReconRunRepository) Get

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

func (*ReconRunRepository) ListByEngagement

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

ListByEngagement returns an engagement's runs, newest first.

func (*ReconRunRepository) ListStaleRunning

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

ListStaleRunning returns runs still 'running' that started before olderThan (≤ limit), oldest first.

func (*ReconRunRepository) Save

func (r *ReconRunRepository) Save(_ context.Context, run recon.Run) error

Save upserts a run.

type ResponseObserverBindingStore added in v0.2.0

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

func NewResponseObserverBindingStore added in v0.2.0

func NewResponseObserverBindingStore() *ResponseObserverBindingStore

func (*ResponseObserverBindingStore) AcknowledgeFleetAudit added in v0.2.0

func (s *ResponseObserverBindingStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*ResponseObserverBindingStore) GetResponseObserverBinding added in v0.2.0

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

func (*ResponseObserverBindingStore) ListPendingFleetAudits added in v0.2.0

func (s *ResponseObserverBindingStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

func (*ResponseObserverBindingStore) ListResponseObserverBindings added in v0.2.0

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

func (*ResponseObserverBindingStore) SaveResponseObserverBindingWithAudit added in v0.2.0

type ResponseStore added in v0.1.8

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

ResponseStore is the in-memory response-action store, tenant-bucketed so one tenant's actions are never visible to another. It mirrors the durable response journal and halt fence used by production.

func NewResponseStore added in v0.1.8

func NewResponseStore() *ResponseStore

NewResponseStore constructs the store.

func (*ResponseStore) AcknowledgeResponseAudit added in v0.2.0

func (s *ResponseStore) AcknowledgeResponseAudit(ctx context.Context, id string) error

func (*ResponseStore) AcknowledgeResponseHaltDispatch added in v0.2.0

func (s *ResponseStore) AcknowledgeResponseHaltDispatch(ctx context.Context, generation int64) error

func (*ResponseStore) AdvanceHaltGenerationWithAudit added in v0.2.0

func (s *ResponseStore) AdvanceHaltGenerationWithAudit(ctx context.Context, expected int64, intent ports.ResponseAuditIntent) (int64, ports.ResponseAuditIntent, ports.ResponseHaltDispatch, error)

func (*ResponseStore) AttemptStillCurrent added in v0.2.0

func (s *ResponseStore) AttemptStillCurrent(ctx context.Context, key string, state responsesaga.SagaState, at time.Time) (bool, error)

func (*ResponseStore) ClaimAttempt added in v0.2.0

func (*ResponseStore) CurrentHaltGeneration added in v0.2.0

func (s *ResponseStore) CurrentHaltGeneration(ctx context.Context) (int64, error)

func (*ResponseStore) EnqueueResponseAudit added in v0.2.0

func (s *ResponseStore) EnqueueResponseAudit(ctx context.Context, intent ports.ResponseAuditIntent) (ports.ResponseAuditIntent, error)

func (*ResponseStore) Get added in v0.1.8

func (s *ResponseStore) Get(ctx context.Context, id shared.ID) (rdom.Record, bool, error)

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

func (*ResponseStore) GetAttempt added in v0.2.0

func (*ResponseStore) ListAttemptsByState added in v0.2.0

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

func (*ResponseStore) ListByState added in v0.1.8

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

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

func (*ResponseStore) ListPendingResponseAudits added in v0.2.0

func (s *ResponseStore) ListPendingResponseAudits(ctx context.Context) ([]ports.ResponseAuditIntent, error)

func (*ResponseStore) ListPendingResponseHaltDispatches added in v0.2.0

func (s *ResponseStore) ListPendingResponseHaltDispatches(ctx context.Context) ([]ports.ResponseHaltDispatch, error)

func (*ResponseStore) Put added in v0.1.8

func (s *ResponseStore) Put(ctx context.Context, r rdom.Record) error

Put upserts a record under the authenticated tenant; a record claiming a different tenant is refused.

func (*ResponseStore) StartAttempt added in v0.2.0

func (*ResponseStore) Transition added in v0.2.0

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

func (*ResponseStore) TransitionAttempt added in v0.2.0

func (*ResponseStore) TransitionWithAudit added in v0.2.0

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

type ResponseVerificationStore added in v0.2.0

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

func NewResponseVerificationStore added in v0.2.0

func NewResponseVerificationStore() *ResponseVerificationStore

func (*ResponseVerificationStore) AcknowledgeFleetAudit added in v0.2.0

func (s *ResponseVerificationStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*ResponseVerificationStore) AppendResponseTargetEvidenceReceipt added in v0.2.0

func (*ResponseVerificationStore) AppendResponseVerificationWithAudit added in v0.2.0

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

func (*ResponseVerificationStore) GetResponseTargetEvidenceReceipt added in v0.2.0

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

func (*ResponseVerificationStore) GetResponseVerification added in v0.2.0

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

func (*ResponseVerificationStore) ListPendingFleetAudits added in v0.2.0

func (s *ResponseVerificationStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

type RetestRepository

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

RetestRepository is an in-memory ports.RetestRepository for dev/tests.

func NewRetestRepository

func NewRetestRepository() *RetestRepository

NewRetestRepository returns an empty in-memory retest store.

func (*RetestRepository) Add

Add appends a retest.

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)

func (*RetestRepository) ListByEngagementFinding

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

ListByEngagementFinding returns a finding's retests oldest-first, engagement-scoped.

type RunLock

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

RunLock is the single-process ports.RunLocker (F9): an in-memory set of currently- executing run ids. It guards against a same-process redelivery (the in-memory queue's lease expiry) re-running a run that is still in flight. Cross-process guarding requires the Postgres advisory-lock implementation.

func NewRunLock

func NewRunLock() *RunLock

NewRunLock returns an in-memory run locker.

func (*RunLock) TryLock

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

type SCMConnectorStore added in v0.2.0

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

SCMConnectorStore is the in-memory twin of the tenant-scoped source-control connector store, for dev and tests. It keys by (tenant, id) and reads the tenant from ctx, mirroring the Postgres store's RLS. The token is held in memory as-is (the Postgres twin seals it); it is returned ONLY via ResolveGitCredential, never by List/Get.

func NewSCMConnectorStore added in v0.2.0

func NewSCMConnectorStore() *SCMConnectorStore

NewSCMConnectorStore constructs an empty store.

func (*SCMConnectorStore) Delete added in v0.2.0

func (s *SCMConnectorStore) Delete(ctx context.Context, id shared.ID) error

func (*SCMConnectorStore) Get added in v0.2.0

func (*SCMConnectorStore) List added in v0.2.0

func (*SCMConnectorStore) Put added in v0.2.0

Put upserts by id under the caller's tenant. A token is always required. A different connector already holding the same host is a conflict, mirroring the Postgres UNIQUE(tenant_id, host).

func (*SCMConnectorStore) ResolveGitCredential added in v0.2.0

func (s *SCMConnectorStore) ResolveGitCredential(ctx context.Context, host string) (ports.GitCredential, bool, error)

ResolveGitCredential returns the credential whose connector host matches the given normalized host, under the caller's tenant. ok=false when none matches.

type SLAStore added in v0.1.8

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

SLAStore is the development/test adapter for the complete SLA governance aggregate. It mirrors the PostgreSQL adapter's tenant checks, idempotency, immutable history, optimistic transitions, and preservation of human state during machine assessment refreshes.

func NewSLAStore added in v0.1.8

func NewSLAStore() *SLAStore

func (*SLAStore) ActivePolicy added in v0.1.8

func (s *SLAStore) ActivePolicy(ctx context.Context, tenantID shared.ID) (sla.Policy, error)

func (*SLAStore) AssessmentHistory added in v0.1.8

func (s *SLAStore) AssessmentHistory(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.Assessment, error)

func (*SLAStore) Current added in v0.1.8

func (s *SLAStore) Current(ctx context.Context, tenantID, engagementID, findingID shared.ID) (sla.Current, error)

func (*SLAStore) LifecycleEvents added in v0.1.8

func (s *SLAStore) LifecycleEvents(ctx context.Context, tenantID, engagementID, findingID shared.ID) ([]sla.LifecycleEvent, error)

func (*SLAStore) ListCurrent added in v0.1.8

func (s *SLAStore) ListCurrent(ctx context.Context, tenantID, engagementID shared.ID) ([]sla.Current, error)

func (*SLAStore) PolicyHistory added in v0.1.8

func (s *SLAStore) PolicyHistory(ctx context.Context, tenantID shared.ID) ([]sla.Policy, error)

func (*SLAStore) PutPolicy added in v0.1.8

func (s *SLAStore) PutPolicy(ctx context.Context, policy sla.Policy, activate bool) (bool, error)

func (*SLAStore) SaveTransition added in v0.1.8

func (s *SLAStore) SaveTransition(ctx context.Context, next sla.Lifecycle, event sla.LifecycleEvent) error

func (*SLAStore) UpsertAssessment added in v0.1.8

func (s *SLAStore) UpsertAssessment(ctx context.Context, assessment sla.Assessment) (sla.AssessmentUpsertResult, error)

type ScanJobStore

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

ScanJobStore is an in-memory store of asynchronous scan-job status.

func NewScanJobStore

func NewScanJobStore() *ScanJobStore

NewScanJobStore returns an empty in-memory scan-job store.

func (*ScanJobStore) CreateRunning

func (s *ScanJobStore) CreateRunning(_ context.Context, j ports.ScanJob) error

func (*ScanJobStore) GetJob

func (s *ScanJobStore) GetJob(_ context.Context, id string) (ports.ScanJob, error)

GetJob returns a job by its own id, or ErrNotFound.

func (*ScanJobStore) LatestForEngagement

func (s *ScanJobStore) LatestForEngagement(_ context.Context, engagementID shared.ID) (ports.ScanJob, error)

LatestForEngagement returns the engagement's most recent job, or ErrNotFound.

func (*ScanJobStore) LatestForEngagements

func (s *ScanJobStore) LatestForEngagements(_ context.Context, engagementIDs []shared.ID) (map[shared.ID]ports.ScanJob, error)

func (*ScanJobStore) ListStaleRunning

func (s *ScanJobStore) ListStaleRunning(_ context.Context, olderThan time.Time, limit int) ([]ports.ScanJob, error)

ListStaleRunning returns jobs still 'running' that started before olderThan (≤ limit), oldest first.

func (*ScanJobStore) Save

Save upserts a job; a newly-seen id becomes the latest for its engagement.

type ScanRepository

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

ScanRepository keeps the current process's component snapshots available to background vulnerability correlation. It is not durable across restarts.

func NewScanRepository

func NewScanRepository(inventory ...*ComponentInventoryStore) *ScanRepository

func (*ScanRepository) AdmitInventory added in v0.2.0

func (r *ScanRepository) AdmitInventory(ctx context.Context, engagementID shared.ID, scope string, admittedAt time.Time) (sbom.InventoryAdmission, error)

func (*ScanRepository) SaveScan

func (r *ScanRepository) SaveScan(ctx context.Context, engagementID shared.ID, doc *sbom.SBOM, vulns []vulnerability.Vulnerability, snap ports.ScanSnapshot) (ports.ScanSaveResult, error)

type ScanResultStore

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

ScanResultStore is an in-memory cache of the latest scan result per engagement.

func NewScanResultStore

func NewScanResultStore() *ScanResultStore

NewScanResultStore returns an empty in-memory scan-result store.

func (*ScanResultStore) LatestResult

func (s *ScanResultStore) LatestResult(_ context.Context, engagementID shared.ID) ([]byte, error)

LatestResult returns the cached scan result, or shared.ErrNotFound.

func (*ScanResultStore) SaveResult

func (s *ScanResultStore) SaveResult(_ context.Context, engagementID shared.ID, result []byte) error

SaveResult stores a copy of the engagement's latest scan result JSON.

type ScanRunStore

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

ScanRunStore is an in-memory store of scan-run manifests and sealed provenance.

func NewScanRunStore

func NewScanRunStore() *ScanRunStore

NewScanRunStore returns an empty in-memory scan-run store.

func (*ScanRunStore) Get

func (s *ScanRunStore) Get(ctx context.Context, runID string) (ports.ScanRun, error)

Get returns a legacy scan run by ID.

func (*ScanRunStore) GetScanRun added in v0.2.0

func (s *ScanRunStore) GetScanRun(_ context.Context, tenantID shared.ID, runID string) (scanrun.ScanRun, error)

GetScanRun returns a scan run with all normalized lanes scoped to tenantID.

func (*ScanRunStore) GetScanRunEvidence added in v0.2.0

func (s *ScanRunStore) GetScanRunEvidence(ctx context.Context, tenantID shared.ID, runID string) (ports.ScanRunEvidence, error)

func (*ScanRunStore) List

func (s *ScanRunStore) List(ctx context.Context, engagementID shared.ID) ([]ports.ScanRun, error)

List returns legacy scan runs for an engagement.

func (*ScanRunStore) ListScanRuns added in v0.2.0

func (s *ScanRunStore) ListScanRuns(_ context.Context, tenantID, engagementID shared.ID) ([]scanrun.ScanRun, error)

ListScanRuns lists all scan runs for a tenant's engagement, sorted newest first.

func (*ScanRunStore) Save

func (s *ScanRunStore) Save(ctx context.Context, run ports.ScanRun) error

Save records a legacy scan run for backwards compatibility.

func (*ScanRunStore) SaveScanRun added in v0.2.0

func (s *ScanRunStore) SaveScanRun(ctx context.Context, run scanrun.ScanRun) error

SaveScanRun persists a tenant-owned native or legacy scan run.

func (*ScanRunStore) SaveScanRunEvidence added in v0.2.0

func (s *ScanRunStore) SaveScanRunEvidence(ctx context.Context, item ports.ScanRunEvidence) error

func (*ScanRunStore) SealScanRun added in v0.2.0

func (s *ScanRunStore) SealScanRun(ctx context.Context, command ports.SealScanRunCommand) error

SealScanRun atomically seals a scan run and records its normalized lanes, versions, and stages.

type ScannedImageStore added in v0.1.8

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

ScannedImageStore is the in-memory scanned-image digest index (#446), tenant-scoped. It normalizes the empty tenant to the default so an image scan and the agent that observes the image correlate under one tenant, matching the Postgres store.

func NewScannedImageStore added in v0.1.8

func NewScannedImageStore() *ScannedImageStore

NewScannedImageStore constructs an empty in-memory store.

func (*ScannedImageStore) MarkScanned added in v0.1.8

func (s *ScannedImageStore) MarkScanned(_ context.Context, tenantID shared.ID, digest string, _ time.Time) error

MarkScanned records digest as scanned for the tenant (idempotent).

func (*ScannedImageStore) ScannedDigests added in v0.1.8

func (s *ScannedImageStore) ScannedDigests(_ context.Context, tenantID shared.ID) (map[string]bool, error)

ScannedDigests returns a copy of the scanned-digest set for the tenant.

type SensorStateStore added in v0.2.0

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

func NewSensorStateStore added in v0.2.0

func NewSensorStateStore() *SensorStateStore

func (*SensorStateStore) AcknowledgeFleetAudit added in v0.2.0

func (s *SensorStateStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*SensorStateStore) AppendSensorState added in v0.2.0

func (s *SensorStateStore) AppendSensorState(ctx context.Context, observation sensorstate.Observation) error

func (*SensorStateStore) AppendSensorStateWithAudit added in v0.2.0

func (s *SensorStateStore) AppendSensorStateWithAudit(
	ctx context.Context,
	observation sensorstate.Observation,
	intent ports.FleetAuditIntent,
) (ports.FleetAuditIntent, error)

func (*SensorStateStore) ListCoverageSensorStates added in v0.2.0

func (*SensorStateStore) ListPendingFleetAudits added in v0.2.0

func (s *SensorStateStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

func (*SensorStateStore) ListSensorStates added in v0.2.0

type SyncRunStore added in v0.1.8

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

func NewSyncRunStore added in v0.1.8

func NewSyncRunStore(ids ports.IDGenerator, now func() time.Time, queue ports.JobQueue) *SyncRunStore

func (*SyncRunStore) Advance added in v0.1.8

func (s *SyncRunStore) Advance(_ context.Context, id shared.ID, expectedCheckpoint, nextCheckpoint []byte, counts vulnerabilitysync.Counts, errors []string) (vulnerabilitysync.Run, error)

func (*SyncRunStore) Finish added in v0.1.8

func (*SyncRunStore) Get added in v0.1.8

func (*SyncRunStore) GetByDurableJobID added in v0.1.8

func (s *SyncRunStore) GetByDurableJobID(ctx context.Context, jobID string) (vulnerabilitysync.Run, error)

func (*SyncRunStore) GetVulnerabilitySyncRun added in v0.1.8

func (s *SyncRunStore) GetVulnerabilitySyncRun(ctx context.Context, tenantID, id shared.ID) (vulnerabilityintel.SyncRunItem, error)

func (*SyncRunStore) LatestForSource added in v0.1.8

func (s *SyncRunStore) LatestForSource(ctx context.Context, sourceID shared.ID, states []vulnerabilitysync.State) (vulnerabilitysync.Run, error)

func (*SyncRunStore) LatestSuccessfulVulnerabilitySync added in v0.1.8

func (s *SyncRunStore) LatestSuccessfulVulnerabilitySync(ctx context.Context, tenantID shared.ID) (*time.Time, error)

func (*SyncRunStore) ListStale added in v0.1.8

func (s *SyncRunStore) ListStale(ctx context.Context, olderThan time.Time, limit int) ([]vulnerabilitysync.Run, error)

func (*SyncRunStore) ListVulnerabilitySyncRuns added in v0.1.8

func (s *SyncRunStore) ListVulnerabilitySyncRuns(ctx context.Context, query vulnerabilityintel.SyncRunQuery) (vulnerabilityintel.SyncRunPage, error)

func (*SyncRunStore) MarkRunning added in v0.1.8

func (s *SyncRunStore) MarkRunning(_ context.Context, id shared.ID) error

func (*SyncRunStore) RecoverStale added in v0.1.8

func (s *SyncRunStore) RecoverStale(ctx context.Context, staleRunID shared.ID, staleBefore time.Time, request ports.SyncRunStart) (vulnerabilitysync.Run, bool, error)

func (*SyncRunStore) Start added in v0.1.8

func (*SyncRunStore) Supersede added in v0.1.8

func (s *SyncRunStore) Supersede(_ context.Context, id shared.ID) error

type TelemetryStore added in v0.1.8

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

TelemetryStore is the in-memory columnar-tier twin used inline/in dev and in tests. It upholds the same contract as the Postgres tier: tenant-bucketed, tiered retention (hot -> warm-at-reduced-resolution -> expiry), sampling recorded with the data, and sequence-gap visibility. It is reached only through ports.TelemetryStore.

func NewTelemetryStore added in v0.1.8

func NewTelemetryStore(hot, warm time.Duration) *TelemetryStore

NewTelemetryStore constructs the store with the hot/warm tier boundaries (the config in ADR 0001).

func (*TelemetryStore) Footprint added in v0.1.8

Footprint reports the GLOBAL store size (all tenants) — an operator spend metric, not a per-tenant figure — with an approximate byte estimate (the Postgres tier reports real on-disk bytes). It carries only counts/bytes, never tenant data, so a global scope leaks nothing.

func (*TelemetryStore) Ingest added in v0.1.8

func (s *TelemetryStore) Ingest(ctx context.Context, batch ports.TelemetryBatch) error

Ingest appends the batch's events, idempotent on (host, class, seq, event index). The bucket is the AUTHENTICATED ctx tenant (fail-closed), never the wire batch's self-declared tenant.

func (*TelemetryStore) LastSequence added in v0.1.8

func (s *TelemetryStore) LastSequence(ctx context.Context, hostID shared.ID, class detection.Class) (uint64, error)

LastSequence returns the highest sequence stored for a (host, class) in the ctx tenant.

func (*TelemetryStore) Query added in v0.1.8

Query runs a retro-hunt over the window and reports completeness honestly.

func (*TelemetryStore) RecordLoss added in v0.2.0

func (s *TelemetryStore) RecordLoss(ctx context.Context, loss ports.TelemetryLoss) error

func (*TelemetryStore) RetentionSweep added in v0.1.8

func (s *TelemetryStore) RetentionSweep(ctx context.Context, now time.Time) (ports.SweepReport, error)

RetentionSweep down-samples the warm window (reduced resolution) and expires past-warm rows for the ctx tenant (tenant-scoped, matching the RLS-scoped Postgres tier).

type TelemetryTransportStore added in v0.2.0

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

TelemetryTransportStore is the in-memory twin of the A3 transport-sequencing store.

func NewTelemetryTransportStore added in v0.2.0

func NewTelemetryTransportStore() *TelemetryTransportStore

func (*TelemetryTransportStore) AcceptAgentGapRevision added in v0.2.0

func (s *TelemetryTransportStore) AcceptAgentGapRevision(ctx context.Context, revision ports.TelemetryAgentGapRevision) error

AcceptAgentGapRevision atomically preserves the exact signed report and advances the current agent-gap projection under the same store lock.

func (*TelemetryTransportStore) AcceptAgentGapRevisionWithAudit added in v0.2.0

func (s *TelemetryTransportStore) AcceptAgentGapRevisionWithAudit(
	ctx context.Context,
	revision ports.TelemetryAgentGapRevision,
	intent ports.FleetAuditIntent,
) (ports.FleetAuditIntent, error)

func (*TelemetryTransportStore) AcknowledgeFleetAudit added in v0.2.0

func (s *TelemetryTransportStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*TelemetryTransportStore) AgentGapRevisions added in v0.2.0

func (s *TelemetryTransportStore) AgentGapRevisions(ctx context.Context, gapID shared.ID) ([]ports.TelemetryAgentGapRevision, error)

AgentGapRevisions returns tenant-scoped immutable revisions for focused persistence verification, ordered by acceptance in this store.

func (*TelemetryTransportStore) BindTelemetryAsset added in v0.2.0

func (s *TelemetryTransportStore) BindTelemetryAsset(ctx context.Context, binding ports.TelemetryAssetBinding) error

func (*TelemetryTransportStore) CommitBatch added in v0.2.0

func (*TelemetryTransportStore) CommitBatchWithAudit added in v0.2.0

func (*TelemetryTransportStore) CountBatchEvents added in v0.2.0

func (s *TelemetryTransportStore) CountBatchEvents(ctx context.Context, agentID, streamID shared.ID, epoch, sequence uint64) (int, error)

func (*TelemetryTransportStore) IngestBatchEvents added in v0.2.0

func (s *TelemetryTransportStore) IngestBatchEvents(ctx context.Context, batch ports.TelemetryEventBatch) (int, error)

func (*TelemetryTransportStore) ListCoverageGapFacts added in v0.2.0

ListCoverageGapFacts returns each auditable loss fact independently. The existing QueryDeliveryGaps projection intentionally remains collapsed for its older consumers; coverage revision identity needs the source and fact ID.

func (*TelemetryTransportStore) ListGapChanges added in v0.2.0

func (s *TelemetryTransportStore) ListGapChanges(
	ctx context.Context,
	agentID, streamID shared.ID,
	epoch, sequence uint64,
) ([]ports.TelemetryGap, error)

ListGapChanges returns enriched inferred-gap history affected by one sequence. Resolved entries remain available so source retries can repair coverage windows.

func (*TelemetryTransportStore) ListGaps added in v0.2.0

func (s *TelemetryTransportStore) ListGaps(ctx context.Context, agentID, streamID shared.ID) ([]ports.TelemetryGap, error)

func (*TelemetryTransportStore) ListPendingFleetAudits added in v0.2.0

func (s *TelemetryTransportStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

func (*TelemetryTransportStore) ListTelemetryAssetBindings added in v0.2.0

func (s *TelemetryTransportStore) 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 for stability.

func (*TelemetryTransportStore) MaxEpoch added in v0.2.0

func (s *TelemetryTransportStore) MaxEpoch(ctx context.Context, agentID, streamID shared.ID) (uint64, error)

func (*TelemetryTransportStore) QueryDeliveryGaps added in v0.2.0

func (*TelemetryTransportStore) QueryTelemetryBatchAccounting added in v0.2.0

func (*TelemetryTransportStore) RecordAgentGap added in v0.2.0

func (*TelemetryTransportStore) ResolveTelemetryAsset added in v0.2.0

func (s *TelemetryTransportStore) ResolveTelemetryAsset(ctx context.Context, agentID shared.ID) (shared.ID, error)

func (*TelemetryTransportStore) ResolveTelemetryReferences added in v0.2.0

func (s *TelemetryTransportStore) ResolveTelemetryReferences(ctx context.Context, agentID, assetID shared.ID, redactionPolicyDigest string, refs []fleetagent.TelemetryReference) (ports.TelemetryReferenceStatus, error)

ResolveTelemetryReferences resolves causal references from the existing accepted-event facts.

func (*TelemetryTransportStore) SaveStreamState added in v0.2.0

func (s *TelemetryTransportStore) SaveStreamState(ctx context.Context, state ports.TelemetryStreamState) error

func (*TelemetryTransportStore) StreamState added in v0.2.0

func (s *TelemetryTransportStore) StreamState(ctx context.Context, agentID, streamID shared.ID, epoch uint64) (ports.TelemetryStreamState, error)

type TenantTransactionRunner added in v0.2.0

type TenantTransactionRunner struct{}

TenantTransactionRunner provides an in-memory implementation of ports.TenantTransactionRunner.

func NewTenantTransactionRunner added in v0.2.0

func NewTenantTransactionRunner() *TenantTransactionRunner

NewTenantTransactionRunner constructs an in-memory TenantTransactionRunner.

func (*TenantTransactionRunner) Run added in v0.2.0

func (r *TenantTransactionRunner) Run(ctx context.Context, tenantID shared.ID, fn func(context.Context) error) error

Run runs the given function within a simulated thread-safe tenant context.

type ThreatModelStore

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

ThreatModelStore is the in-memory architecture-input threat-model store (dev/tests, mirrors the Postgres adapter): one model per engagement, replaced on each Save (re-syncable, not append-only). Reads are engagement-scoped (the tenant gate runs upstream at the child route).

func NewThreatModelStore

func NewThreatModelStore() *ThreatModelStore

NewThreatModelStore returns an empty in-memory threat-model store.

func (*ThreatModelStore) Get

func (s *ThreatModelStore) Get(_ context.Context, engagementID shared.ID) (threatmodel.Model, bool, error)

Get returns the engagement's model, ok=false when none has been ingested.

func (*ThreatModelStore) Save

func (s *ThreatModelStore) Save(_ context.Context, engagementID, _ shared.ID, m threatmodel.Model) error

Save upserts the engagement's model (tenant id unused in memory; the Postgres adapter persists it).

type TimestampStore

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

TimestampStore is an in-memory ports.TimestampStore for dev/tests: external RFC-3161 tokens keyed by (chain, engagement, head).

func NewTimestampStore

func NewTimestampStore() *TimestampStore

NewTimestampStore returns an empty in-memory timestamp store.

func (*TimestampStore) Get

func (s *TimestampStore) Get(_ context.Context, chain string, eng shared.ID, head string) (*ports.TimestampToken, error)

Get returns the stored token for a head, or nil if not yet anchored.

func (*TimestampStore) LatestHead

func (s *TimestampStore) LatestHead(_ context.Context, chain string, eng shared.ID) (string, bool, error)

LatestHead returns the most-recently-Put head for a chain (ok=false if none) – the retained head for out-of-band tail-truncation detection.

func (*TimestampStore) Put

func (s *TimestampStore) Put(_ context.Context, chain string, eng shared.ID, head string, token ports.TimestampToken) error

Put stores a token for a head (idempotent – first write wins, like the SQL ON CONFLICT DO NOTHING).

type UserRepository

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

UserRepository is an in-memory ports.UserRepository for dev/tests.

func NewUserRepository

func NewUserRepository() *UserRepository

NewUserRepository returns an empty in-memory user store.

func (*UserRepository) Bootstrap added in v0.2.0

func (r *UserRepository) Bootstrap(ctx context.Context, u *user.User, _ ports.AuditEntry) error

Bootstrap preserves the in-memory repository's historical seed behavior. Durable audit-chain guarantees are provided by the PostgreSQL implementation.

func (*UserRepository) Create

func (r *UserRepository) Create(_ context.Context, u *user.User) error

func (*UserRepository) GetByAPIKeyHash

func (r *UserRepository) GetByAPIKeyHash(_ context.Context, hash string) (*user.User, error)

func (*UserRepository) GetByID

func (r *UserRepository) GetByID(_ context.Context, tenantID, id shared.ID) (*user.User, error)

func (*UserRepository) List

func (r *UserRepository) List(_ context.Context, tenantID shared.ID) ([]*user.User, error)

func (*UserRepository) Update added in v0.2.0

func (r *UserRepository) Update(_ context.Context, tenantID shared.ID, u *user.User) error

Update writes the mutable fields of an existing user inside tenantID. The tenant of the stored row is preserved, so an update can never move a user between tenants.

func (*UserRepository) Upsert

func (r *UserRepository) Upsert(_ context.Context, u *user.User) error

type VEXStatementStore added in v0.2.0

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

VEXStatementStore is the in-memory VEX-statement store. It is EPHEMERAL (lost on restart); the server uses the Postgres store. It backs tests and the CLI. Statements are partitioned by tenant, then by engagement, then keyed by content digest so re-ingesting an identical assertion cannot duplicate a row.

func NewVEXStatementStore added in v0.2.0

func NewVEXStatementStore() *VEXStatementStore

NewVEXStatementStore returns an empty store.

func (*VEXStatementStore) ListByEngagement added in v0.2.0

func (s *VEXStatementStore) ListByEngagement(_ context.Context, tenantID, engagementID shared.ID) ([]vex.StoredStatement, error)

ListByEngagement returns the engagement's statements in insertion order, so a re-apply walking them applies the most-recent assertion last (newest wins per finding).

func (*VEXStatementStore) Save added in v0.2.0

func (s *VEXStatementStore) Save(_ context.Context, tenantID, engagementID shared.ID, statements []vex.StoredStatement) error

Save persists the statements idempotently by digest; an identical assertion already stored is left as-is (keeping its original insertion order). New statements take the next insertion sequence.

type VulnerabilityActionStore added in v0.1.8

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

func NewVulnerabilityActionStore added in v0.1.8

func NewVulnerabilityActionStore() *VulnerabilityActionStore

func (*VulnerabilityActionStore) AcknowledgeAction added in v0.1.8

func (s *VulnerabilityActionStore) AcknowledgeAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) ClaimOutbox added in v0.1.8

func (s *VulnerabilityActionStore) ClaimOutbox(ctx context.Context, tenantID shared.ID, now, lockedUntil time.Time, limit int) ([]vulnerabilityaction.OutboxEvent, error)

func (*VulnerabilityActionStore) CompleteOutbox added in v0.1.8

func (s *VulnerabilityActionStore) CompleteOutbox(ctx context.Context, tenantID, eventID shared.ID, at time.Time) error

func (*VulnerabilityActionStore) CountPendingVulnerabilityActions added in v0.1.8

func (s *VulnerabilityActionStore) CountPendingVulnerabilityActions(ctx context.Context, tenantID shared.ID) (int64, error)

func (*VulnerabilityActionStore) GetAction added in v0.1.8

func (s *VulnerabilityActionStore) GetAction(ctx context.Context, tenantID, actionID shared.ID) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) ListActions added in v0.1.8

func (*VulnerabilityActionStore) ListVulnerabilityTransitions added in v0.1.8

func (*VulnerabilityActionStore) RecordChange added in v0.1.8

func (*VulnerabilityActionStore) ResolveAction added in v0.1.8

func (s *VulnerabilityActionStore) ResolveAction(ctx context.Context, tenantID, actionID shared.ID, actor string, at time.Time) (vulnerabilityaction.Action, error)

func (*VulnerabilityActionStore) RetryOutbox added in v0.1.8

func (s *VulnerabilityActionStore) RetryOutbox(ctx context.Context, tenantID, eventID shared.ID, at, availableAt time.Time, lastError string, terminal bool) error

func (*VulnerabilityActionStore) SummarizeVulnerabilityActions added in v0.1.8

func (s *VulnerabilityActionStore) SummarizeVulnerabilityActions(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryActionSummary, error)

type VulnerabilityOccurrenceStore added in v0.1.8

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

func NewVulnerabilityOccurrenceStore added in v0.1.8

func NewVulnerabilityOccurrenceStore() *VulnerabilityOccurrenceStore

func (*VulnerabilityOccurrenceStore) CountActiveVulnerabilityOccurrences added in v0.1.8

func (s *VulnerabilityOccurrenceStore) CountActiveVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryID string) (int64, error)

func (*VulnerabilityOccurrenceStore) CountNewlyAffectedAssets added in v0.2.0

func (s *VulnerabilityOccurrenceStore) CountNewlyAffectedAssets(ctx context.Context, tenantID shared.ID, since time.Time) (int64, error)

func (*VulnerabilityOccurrenceStore) Get added in v0.1.8

func (s *VulnerabilityOccurrenceStore) Get(ctx context.Context, tenantID, engagementID shared.ID, advisoryID, componentFingerprint string) (vulnerabilityoccurrence.Occurrence, error)

func (*VulnerabilityOccurrenceStore) ListByEngagement added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListByEngagement(ctx context.Context, tenantID, engagementID shared.ID, states []vulnerabilityoccurrence.State) ([]vulnerabilityoccurrence.Occurrence, error)

func (*VulnerabilityOccurrenceStore) ListEvents added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListEvents(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityoccurrence.Event, error)

func (*VulnerabilityOccurrenceStore) ListUnreconciled added in v0.1.8

func (s *VulnerabilityOccurrenceStore) ListUnreconciled(ctx context.Context, tenantID, runID shared.ID, advisoryID string, after shared.ID, snapshotAt time.Time, limit int) (ports.VulnerabilityOccurrenceReconciliationPage, error)

func (*VulnerabilityOccurrenceStore) ListVulnerabilityOccurrences added in v0.1.8

func (*VulnerabilityOccurrenceStore) SummarizeVulnerabilityOccurrences added in v0.1.8

func (s *VulnerabilityOccurrenceStore) SummarizeVulnerabilityOccurrences(ctx context.Context, tenantID shared.ID, advisoryIDs []string, affectedAsset string, states []vulnerabilityoccurrence.State) (map[string]vulnerabilityintel.AdvisoryOccurrenceSummary, error)

func (*VulnerabilityOccurrenceStore) Upsert added in v0.1.8

type VulnerabilityReconcileRunStore added in v0.1.8

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

func NewVulnerabilityReconcileRunStore added in v0.1.8

func NewVulnerabilityReconcileRunStore(ids ports.IDGenerator, clock ports.Clock, queue ports.JobQueue) *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, errors []string) (vulnerabilityreconcile.Run, error)

func (*VulnerabilityReconcileRunStore) Finish added in v0.1.8

func (*VulnerabilityReconcileRunStore) Get added in v0.1.8

func (*VulnerabilityReconcileRunStore) GetByDurableJobID added in v0.1.8

func (*VulnerabilityReconcileRunStore) HasReconciliationMatch added in v0.1.8

func (s *VulnerabilityReconcileRunStore) HasReconciliationMatch(ctx context.Context, tenantID, runID, engagementID shared.ID, advisoryID, componentFingerprint string) (bool, error)

func (*VulnerabilityReconcileRunStore) ListReconciliationDiffs added in v0.1.8

func (*VulnerabilityReconcileRunStore) MarkRunning added in v0.1.8

func (s *VulnerabilityReconcileRunStore) MarkRunning(ctx context.Context, id shared.ID) error

func (*VulnerabilityReconcileRunStore) RecordReconciliationDiff added in v0.1.8

func (s *VulnerabilityReconcileRunStore) RecordReconciliationDiff(ctx context.Context, diff vulnerabilityreconcile.Diff) (bool, error)

func (*VulnerabilityReconcileRunStore) Start added in v0.1.8

func (*VulnerabilityReconcileRunStore) SummarizeReconciliationDiffs added in v0.1.8

func (s *VulnerabilityReconcileRunStore) SummarizeReconciliationDiffs(ctx context.Context, tenantID, runID shared.ID) (vulnerabilityreconcile.Counts, error)

type VulnerabilityRiskAssessmentStore added in v0.1.8

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

func NewVulnerabilityRiskAssessmentStore added in v0.1.8

func NewVulnerabilityRiskAssessmentStore() *VulnerabilityRiskAssessmentStore

func (*VulnerabilityRiskAssessmentStore) CountOpenHighCriticalVulnerabilityExposure added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) CountOpenHighCriticalVulnerabilityExposure(ctx context.Context, tenantID shared.ID) (int64, error)

func (*VulnerabilityRiskAssessmentStore) Current added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) Current(ctx context.Context, tenantID, occurrenceID shared.ID) (vulnerabilityrisk.Assessment, error)

func (*VulnerabilityRiskAssessmentStore) History added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) History(ctx context.Context, tenantID, occurrenceID shared.ID) ([]vulnerabilityrisk.Assessment, error)

func (*VulnerabilityRiskAssessmentStore) ListVulnerabilityAssessments added in v0.1.8

func (*VulnerabilityRiskAssessmentStore) SummarizeVulnerabilityRisk added in v0.1.8

func (s *VulnerabilityRiskAssessmentStore) SummarizeVulnerabilityRisk(ctx context.Context, tenantID shared.ID, advisoryIDs []string) (map[string]vulnerabilityintel.AdvisoryRiskSummary, error)

func (*VulnerabilityRiskAssessmentStore) Upsert added in v0.1.8

type VulnerabilitySourceStore added in v0.1.8

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

func NewVulnerabilitySourceStore added in v0.1.8

func NewVulnerabilitySourceStore() *VulnerabilitySourceStore

func (*VulnerabilitySourceStore) Archive added in v0.1.8

func (s *VulnerabilitySourceStore) Archive(_ context.Context, id shared.ID, expectedVersion int) error

func (*VulnerabilitySourceStore) Create added in v0.1.8

func (*VulnerabilitySourceStore) Get added in v0.1.8

func (*VulnerabilitySourceStore) List added in v0.1.8

func (s *VulnerabilitySourceStore) List(_ context.Context, includeArchived bool) ([]vulnerabilitysource.Source, error)

func (*VulnerabilitySourceStore) SetEnabled added in v0.1.8

func (s *VulnerabilitySourceStore) SetEnabled(_ context.Context, id shared.ID, enabled bool, expectedVersion int) (vulnerabilitysource.Source, error)

func (*VulnerabilitySourceStore) Update added in v0.1.8

func (s *VulnerabilitySourceStore) Update(_ context.Context, source vulnerabilitysource.Source, expectedVersion int) error

type WorkOrderStore added in v0.1.8

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

WorkOrderStore is an in-memory ports.WorkOrderStore for dev and tests. It mirrors the Postgres store's idempotency, in-flight uniqueness and CAS-transition semantics.

func NewWorkOrderStore added in v0.1.8

func NewWorkOrderStore() *WorkOrderStore

NewWorkOrderStore returns an empty in-memory work order store.

func (*WorkOrderStore) AcknowledgeFleetAudit added in v0.2.0

func (s *WorkOrderStore) AcknowledgeFleetAudit(ctx context.Context, id string) error

func (*WorkOrderStore) CancelForAgent added in v0.1.8

func (s *WorkOrderStore) CancelForAgent(_ context.Context, tenantID, agentID shared.ID, reason string, now time.Time) (int, error)

CancelForAgent cancels every live order addressed to agentID and returns the count.

func (*WorkOrderStore) CancelResponsesBelowGeneration added in v0.2.0

func (s *WorkOrderStore) CancelResponsesBelowGeneration(_ context.Context, tenantID shared.ID, generation int64, reason string, now time.Time) (int, error)

CancelResponsesBelowGeneration cancels live response commands carrying an older halt generation.

func (*WorkOrderStore) CancelResponsesBelowGenerationWithAudit added in v0.2.0

func (s *WorkOrderStore) CancelResponsesBelowGenerationWithAudit(_ context.Context, tenantID shared.ID, generation int64, reason string, now time.Time, actor string) (int, ports.FleetAuditIntent, error)

func (*WorkOrderStore) Claim added in v0.1.8

func (s *WorkOrderStore) Claim(_ 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.

func (*WorkOrderStore) ClaimWithAudit added in v0.2.0

func (s *WorkOrderStore) ClaimWithAudit(_ context.Context, tenantID, agentID shared.ID, max int, now time.Time, leaseID string, leaseUntil time.Time, actor string) ([]*workorder.WorkOrder, []ports.FleetAuditIntent, error)

func (*WorkOrderStore) CompleteResponse added in v0.2.0

func (s *WorkOrderStore) CompleteResponse(_ context.Context, tenantID, id shared.ID, result fleetagent.ResponseExecutionResult, reason string, now time.Time) (bool, error)

CompleteResponse records an exact response execution result under the current lease.

func (*WorkOrderStore) CompleteResponseWithAudit added in v0.2.0

func (s *WorkOrderStore) CompleteResponseWithAudit(_ context.Context, tenantID, id shared.ID, result fleetagent.ResponseExecutionResult, reason string, now time.Time, actor string) (bool, ports.FleetAuditIntent, error)

func (*WorkOrderStore) GetByID added in v0.1.8

func (s *WorkOrderStore) GetByID(_ context.Context, tenantID, id shared.ID) (*workorder.WorkOrder, error)

GetByID returns the order or shared.ErrNotFound.

func (*WorkOrderStore) GetByIdempotencyKey added in v0.2.0

func (s *WorkOrderStore) GetByIdempotencyKey(_ context.Context, tenantID shared.ID, idempotencyKey string) (*workorder.WorkOrder, error)

GetByIdempotencyKey returns the order for an idempotency key or shared.ErrNotFound.

func (*WorkOrderStore) Issue added in v0.1.8

Issue stores wo. It is idempotent by (tenant, idempotency key) and rejects a second live order for the same (tenant, asset, capability, time bucket) with shared.ErrConflict.

func (*WorkOrderStore) IssueWithAudit added in v0.2.0

IssueWithAudit atomically persists a work order and its exact issuance audit obligation.

func (*WorkOrderStore) ListByTenant added in v0.1.8

func (s *WorkOrderStore) ListByTenant(_ context.Context, tenantID shared.ID) ([]*workorder.WorkOrder, error)

ListByTenant returns every work order for the tenant, ordered by id for determinism.

func (*WorkOrderStore) ListPendingFleetAudits added in v0.2.0

func (s *WorkOrderStore) ListPendingFleetAudits(ctx context.Context) ([]ports.FleetAuditIntent, error)

func (*WorkOrderStore) Transition added in v0.1.8

func (s *WorkOrderStore) Transition(_ 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: shared.ErrNotFound when the order does not exist under the tenant, shared.ErrConflict when its state no longer matches expected.

func (*WorkOrderStore) TransitionLeased added in v0.2.0

func (s *WorkOrderStore) TransitionLeased(_ context.Context, tenantID, id shared.ID, leaseID string, to workorder.State, reason string, expected workorder.State, now time.Time) error

TransitionLeased applies a transition only while the presented lease still owns the order.

func (*WorkOrderStore) TransitionLeasedWithAudit added in v0.2.0

func (s *WorkOrderStore) TransitionLeasedWithAudit(_ context.Context, tenantID, id shared.ID, leaseID string, to workorder.State, reason string, expected workorder.State, now time.Time, actor string) (ports.FleetAuditIntent, error)

type WriteupDraftStore

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

WriteupDraftStore is the in-memory write-up-draft repository (dev/tests). It mirrors the Postgres adapter to come: Save is an UPSERT by draft id (a draft is mutable working data – edited then accepted/rejected), and reads are engagement-scoped (tenant isolation is enforced upstream at the route). ListByEngagement returns a deterministic (created_at, id) order to match the SQL adapter.

func NewWriteupDraftStore

func NewWriteupDraftStore() *WriteupDraftStore

NewWriteupDraftStore builds an empty in-memory draft store.

func (*WriteupDraftStore) Get

func (s *WriteupDraftStore) Get(_ context.Context, engagementID, id shared.ID) (writeupdraft.Draft, error)

Get returns the engagement's draft by id, or shared.ErrNotFound.

func (*WriteupDraftStore) ListByEngagement

func (s *WriteupDraftStore) ListByEngagement(_ context.Context, engagementID shared.ID) ([]writeupdraft.Draft, error)

ListByEngagement returns a copy of the engagement's drafts ordered by (created_at, id).

func (*WriteupDraftStore) Save

Save upserts a draft by id within its engagement (replace in place if present, else append).

Source Files

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL