vulnerabilitymonitor

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: 14 Imported by: 0

Documentation

Index

Constants

View Source
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 Clock added in v0.2.0

type Clock interface {
	Now() time.Time
}

Clock is the narrow time source the scheduler needs.

type Job

type Job struct {
	RunID    shared.ID `json:"run_id"`
	SourceID shared.ID `json:"source_id"`
}

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 RecoveryResult struct {
	StaleRunID shared.ID
	Run        vulnerabilitysync.Run
	Created    bool
	Disabled   bool
	// Err is set only by the batch RecoverStaleRuns sweep, where one run's recovery failure must not abort
	// the others; the single RecoverStale returns its error directly instead.
	Err error
}

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

func (r *Registry) AllowPrivateNetworkSources(allow bool)

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.

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.

func (*Scheduler) Run added in v0.2.0

func (sc *Scheduler) Run(ctx context.Context)

Run ticks immediately, then on the configured interval until the context is cancelled.

func (*Scheduler) Tick added in v0.2.0

func (sc *Scheduler) Tick(ctx context.Context) (enqueued, recovered int, err error)

Tick enqueues due sources and recovers stale runs once. It is a no-op unless this worker is the leader. Returns the number of syncs enqueued and stale runs recovered.

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 (*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

func (s *Service) Execute(ctx context.Context, runID shared.ID) (vulnerabilitysync.Run, error)

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

func (s *Service) ExecuteJob(ctx context.Context, jobID string) (vulnerabilitysync.Run, error)

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) FailJob

func (s *Service) FailJob(ctx context.Context, jobID string, cause error) error

func (*Service) GetRun

func (s *Service) GetRun(ctx context.Context, id shared.ID) (vulnerabilitysync.Run, error)

func (*Service) Health

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

func (s *Service) SetRunLock(lock ports.RunLocker)

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"
)

type StartAllResult

type StartAllResult struct {
	SourceID shared.ID
	Run      vulnerabilitysync.Run
	Created  bool
	Err      error
}

type StartRequest

type StartRequest struct {
	SourceID             shared.ID
	Mode                 vulnerabilitysync.Mode
	Trigger              string
	Actor                string
	ClientIdempotencyKey string
}

Jump to

Keyboard shortcuts

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