Documentation
¶
Index ¶
- Constants
- type AlwaysLeader
- type Clock
- type Job
- type Leadership
- type RecoveryResult
- type Registry
- func (r *Registry) AllowPrivateNetworkSources(allow bool)
- func (r *Registry) Register(provider ports.VulnerabilityIntelligenceProvider) error
- func (r *Registry) RegisterFactory(adapterType string, factory ports.VulnerabilityProviderFactory) error
- func (r *Registry) Resolve(source vulnerabilitysource.Source) (ports.VulnerabilityIntelligenceProvider, error)
- type Scheduler
- type SchedulerConfig
- type Service
- func (s *Service) DueForSync(ctx context.Context, now time.Time) ([]vulnerabilitysource.Source, error)
- func (s *Service) Execute(ctx context.Context, runID shared.ID) (vulnerabilitysync.Run, error)
- func (s *Service) ExecuteJob(ctx context.Context, jobID string) (vulnerabilitysync.Run, error)
- func (s *Service) FailJob(ctx context.Context, jobID string, cause error) error
- func (s *Service) GetRun(ctx context.Context, id shared.ID) (vulnerabilitysync.Run, error)
- func (s *Service) Health(ctx context.Context, source vulnerabilitysource.Source) (SourceHealth, error)
- func (s *Service) RecoverStale(ctx context.Context, stale vulnerabilitysync.Run, staleBefore time.Time) (RecoveryResult, error)
- func (s *Service) RecoverStaleRuns(ctx context.Context, staleBefore time.Time, limit int) ([]RecoveryResult, error)
- func (s *Service) SetReconciler(reconciler ports.AdvisoryRevisionReconciler)
- func (s *Service) SetRollout(policy *vulnerabilityrollout.Policy)
- func (s *Service) SetRunLock(lock ports.RunLocker)
- func (s *Service) Start(ctx context.Context, request StartRequest) (vulnerabilitysync.Run, bool, error)
- func (s *Service) StartAll(ctx context.Context, mode vulnerabilitysync.Mode, ...) ([]StartAllResult, error)
- type SourceHealth
- type SourceHealthState
- type StartAllResult
- type StartRequest
Constants ¶
const JobKind = "vulnerability_sync"
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AlwaysLeader ¶ added in v0.2.0
type AlwaysLeader struct{}
AlwaysLeader is the single-worker default.
func (AlwaysLeader) IsLeader ¶ added in v0.2.0
func (AlwaysLeader) IsLeader() bool
IsLeader always reports true.
type Leadership ¶ added in v0.2.0
type Leadership interface {
IsLeader() bool
}
Leadership gates the scheduler to a single worker, so a due sync is enqueued once across the fleet.
type RecoveryResult ¶
type Registry ¶
type Registry struct {
// contains filtered or unexported fields
}
Registry is an explicit code-owned provider registry. Registration is done by composition roots; runtime source configuration cannot upload executable code.
func NewRegistry ¶
func NewRegistry() *Registry
func (*Registry) AllowPrivateNetworkSources ¶ added in v0.2.0
AllowPrivateNetworkSources opts the deployment in to sources that target private address ranges. Default false: a source that asks for private-range egress does not resolve, so it neither syncs nor tests, whatever its stored adapter config says.
func (*Registry) Register ¶
func (r *Registry) Register(provider ports.VulnerabilityIntelligenceProvider) error
func (*Registry) RegisterFactory ¶
func (r *Registry) RegisterFactory(adapterType string, factory ports.VulnerabilityProviderFactory) error
RegisterFactory registers one reviewed adapter constructor. A factory is resolved for every source execution so two source rows never share mutable endpoint, checkpoint, or client state.
func (*Registry) Resolve ¶
func (r *Registry) Resolve(source vulnerabilitysource.Source) (ports.VulnerabilityIntelligenceProvider, error)
type Scheduler ¶ added in v0.2.0
type Scheduler struct {
// contains filtered or unexported fields
}
Scheduler drives cadence-based vulnerability-source syncing: on each tick it enqueues due sources and reclaims stranded runs. It is leader-gated so exactly one worker schedules, and it honours the same provider-sync rollout gate as a manual sync, so it does nothing while sync is disabled.
func NewScheduler ¶ added in v0.2.0
func NewScheduler(service *Service, clock Clock, leadership Leadership, config SchedulerConfig, log *slog.Logger) (*Scheduler, error)
NewScheduler validates and returns a scheduler.
type SchedulerConfig ¶ added in v0.2.0
type SchedulerConfig struct {
Interval time.Duration // how often to poll for due sources and stale runs
StaleAfter time.Duration // a queued/running run older than this is reclaimed
DispatchLimit int // maximum syncs enqueued (and stale runs recovered) per tick
}
SchedulerConfig tunes the cadence-driven sync scheduler.
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
func NewService ¶
func NewService(sources ports.VulnerabilitySourceStore, runs ports.SyncRunStore, materializer ports.AdvisoryMaterializer, registry ports.VulnerabilityProviderRegistry, clock ports.Clock) (*Service, error)
func (*Service) DueForSync ¶ added in v0.2.0
func (s *Service) DueForSync(ctx context.Context, now time.Time) ([]vulnerabilitysource.Source, error)
DueForSync returns the enabled, non-archived sources whose freshness has lapsed: a source is due when it has never completed a successful sync, or its last successful sync finished more than one Cadence ago. A source with no Cadence configured is never scheduled (it syncs only on an explicit trigger). The query is O(sources): one source list plus one latest-run lookup per source.
func (*Service) Execute ¶
Execute processes one queued/running run. Every materialized batch is committed before the run checkpoint advances. A canceled provider leaves the run running with the last committed checkpoint so the durable worker can retry it safely.
func (*Service) ExecuteJob ¶
ExecuteJob resolves the run from the durable queue row. The queue payload is intentionally opaque; the durable job id is the authoritative join key.
func (*Service) Health ¶
func (s *Service) Health(ctx context.Context, source vulnerabilitysource.Source) (SourceHealth, error)
func (*Service) RecoverStale ¶
func (s *Service) RecoverStale(ctx context.Context, stale vulnerabilitysync.Run, staleBefore time.Time) (RecoveryResult, error)
func (*Service) RecoverStaleRuns ¶ added in v0.2.0
func (s *Service) RecoverStaleRuns(ctx context.Context, staleBefore time.Time, limit int) ([]RecoveryResult, error)
RecoverStaleRuns reclaims sync runs stranded in queued/running past staleBefore, so a crashed or abandoned run cannot block a source forever. It lists at most limit stale runs and recovers each; a per-run recovery error is collected, not fatal, so one bad run does not stop the sweep.
func (*Service) SetReconciler ¶
func (s *Service) SetReconciler(reconciler ports.AdvisoryRevisionReconciler)
func (*Service) SetRollout ¶
func (s *Service) SetRollout(policy *vulnerabilityrollout.Policy)
func (*Service) SetRunLock ¶ added in v0.2.0
SetRunLock guards one active durable execution per sync run across workers.
func (*Service) Start ¶
func (s *Service) Start(ctx context.Context, request StartRequest) (vulnerabilitysync.Run, bool, error)
func (*Service) StartAll ¶
func (s *Service) StartAll(ctx context.Context, mode vulnerabilitysync.Mode, trigger, actor, idempotencyPrefix string) ([]StartAllResult, error)
StartAll enqueues enabled, non-archived sources in deterministic key order. One source failure does not duplicate or roll back already accepted runs.
type SourceHealth ¶
type SourceHealth struct {
State SourceHealthState `json:"state"`
Stale bool `json:"stale"`
LatestRun *vulnerabilitysync.Run `json:"latest_run,omitempty"`
LastSuccessfulAt *time.Time `json:"last_successful_at,omitempty"`
FreshUntil *time.Time `json:"fresh_until,omitempty"`
}
type SourceHealthState ¶
type SourceHealthState string
const ( SourceHealthArchived SourceHealthState = "archived" SourceHealthDisabled SourceHealthState = "disabled" SourceHealthSyncing SourceHealthState = "syncing" SourceHealthNeverSynced SourceHealthState = "never_synced" SourceHealthHealthy SourceHealthState = "healthy" SourceHealthStale SourceHealthState = "stale" SourceHealthPartial SourceHealthState = "partial" SourceHealthFailed SourceHealthState = "failed" )