Documentation
¶
Overview ¶
Package symptomre implements the asynchronous symptom re-evaluation workflow. It defines domain-specific database models, batch submission, status querying, and River job types for re-evaluating job run symptoms via a work queue.
Index ¶
- Constants
- Variables
- type Batch
- type BatchCanceller
- type BatchCleanupProcess
- type BatchItem
- type BatchStatusCounts
- type BatchStatusResponse
- type ItemState
- type ItemStatus
- type ProcessBatchArgs
- type ProcessBatchWorker
- type ReevaluateFunc
- type ReevaluateJobRunArgs
- type ReevaluateWorker
- type StatusQuerier
- type SubmitResult
- type Submitter
Constants ¶
const ( // CompletedBatchRetention is how long completed (or failed) batches are // kept before periodic cleanup deletes them. The ON DELETE CASCADE // foreign key on batch items handles child-row removal automatically. CompletedBatchRetention = 7 * 24 * time.Hour // StaleBatchTimeout is the maximum age for a batch in a non-terminal // status (pending, processing, running) before it is considered stuck // and removed by the cleanup process. StaleBatchTimeout = 24 * time.Hour )
const ( // BatchQueue is the River queue for batch fan-out jobs. Using a separate // queue prevents a backlog of individual re-evaluation items from blocking // prompt processing of new batches. BatchQueue = "symptom_re_batch" // ItemQueue is the River queue for individual job run re-evaluation jobs. ItemQueue = "symptom_re_item" // DedupPeriod prevents re-evaluating the same job run within a two-hour // window, reducing the chance of BigQuery streaming buffer conflicts // (rows are not deletable within 90 minutes of insertion). DedupPeriod = 120 * time.Minute // MaxAttemptsPerItem is the number of times an individual re-evaluation // job is attempted before being discarded as a permanent failure. MaxAttemptsPerItem = 3 )
Variables ¶
var ErrBatchTerminal = errors.New("batch is already in a terminal status")
ErrBatchTerminal indicates the batch is already in a terminal status and cannot be canceled. The handler maps this to 409 Conflict.
Functions ¶
This section is empty.
Types ¶
type Batch ¶
type Batch struct {
ID uuid.UUID `gorm:"type:uuid;primaryKey" json:"id"`
RequestedCount int `gorm:"not null" json:"requested_count"`
EnqueuedCount int `gorm:"not null" json:"enqueued_count"`
DedupedCount int `gorm:"not null" json:"deduped_count"`
DryRun bool `gorm:"not null;default:false" json:"dry_run"`
Status workqueue.BatchStatus `gorm:"not null;default:'pending'" json:"status"`
CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"`
CompletedAt *time.Time ` json:"completed_at,omitempty"`
}
Batch represents a user-initiated batch of symptom re-evaluation work. The API creates batches; the daemon processes them via River jobs.
type BatchCanceller ¶
type BatchCanceller struct {
// contains filtered or unexported fields
}
BatchCanceller cancels in-flight symptom re-evaluation batches by requesting cancellation of their River jobs and marking the batch as canceled. Jobs that have already completed are left alone.
func NewBatchCanceller ¶
NewBatchCanceller creates a BatchCanceller.
func (*BatchCanceller) Cancel ¶
func (c *BatchCanceller) Cancel(ctx context.Context, batchID uuid.UUID) (*BatchStatusResponse, error)
Cancel attempts to cancel all non-completed River jobs in a batch and marks the batch as canceled. Returns the batch status response after cancellation. Returns nil if the batch does not exist.
type BatchCleanupProcess ¶
type BatchCleanupProcess struct {
// contains filtered or unexported fields
}
BatchCleanupProcess periodically removes old completed batches and cancels stale non-terminal batches. It implements the DaemonProcess interface so the daemon server can manage its lifecycle.
func NewBatchCleanupProcess ¶
func NewBatchCleanupProcess(db *gorm.DB, canceller *BatchCanceller) *BatchCleanupProcess
NewBatchCleanupProcess creates a cleanup process with the default retention periods for completed and stale batches. The canceller is used to cancel stale batches (including their in-flight River jobs).
func (*BatchCleanupProcess) CancelStaleBatches ¶
func (p *BatchCleanupProcess) CancelStaleBatches(ctx context.Context) (int, error)
CancelStaleBatches finds batches stuck in non-terminal statuses (pending, processing, running) past the stale timeout and cancels them, including any in-flight River jobs. This preserves history for the frontend (which would otherwise see a 404) and lets DeleteCompletedBatches remove them after the normal retention period.
func (*BatchCleanupProcess) DeleteCompletedBatches ¶
func (p *BatchCleanupProcess) DeleteCompletedBatches(ctx context.Context) (int64, error)
DeleteCompletedBatches removes batches that have a non-null completed_at timestamp older than the configured retention period.
type BatchItem ¶
type BatchItem struct {
ID uint64 `gorm:"primaryKey;autoIncrement" json:"id"`
BatchID uuid.UUID `gorm:"type:uuid;not null;index:idx_symptom_re_batch_items" json:"batch_id"`
RiverJobID *int64 `gorm:"index:idx_symptom_re_batch_items" json:"river_job_id"`
ItemKey string `gorm:"not null" json:"item_key"`
}
BatchItem associates a batch with an individual job run to re-evaluate. Before the daemon processes the batch, RiverJobID is nil (the item is a specification). After processing, RiverJobID is populated with the River job that performs the work.
type BatchStatusCounts ¶
type BatchStatusCounts struct {
Requested int `json:"requested"`
Enqueued int `json:"enqueued"`
Deduped int `json:"deduped"`
Completed int `json:"completed"`
Failed int `json:"failed"`
Running int `json:"running"`
Pending int `json:"pending"`
}
BatchStatusCounts holds the numeric counters for a batch status response.
type BatchStatusResponse ¶
type BatchStatusResponse struct {
BatchID uuid.UUID `json:"batch_id"`
Status workqueue.BatchStatus `json:"status"`
BatchStatusCounts
Items []ItemStatus `json:"items"`
}
BatchStatusResponse is the API response for querying a batch's status.
type ItemState ¶
type ItemState = string
ItemState represents the resolved state of a batch item. Most values come directly from River's job state enum; the two synthetic states cover items that have not yet been enqueued or whose River job row is missing.
const ( // ItemStateNotEnqueued means the batch item has no river_job_id yet // (the batch fan-out hasn't reached it). ItemStateNotEnqueued ItemState = "not_enqueued" // ItemStateOrphaned means the batch item references a river_job_id that // no longer exists (e.g. cleaned up by River's job cleaner). ItemStateOrphaned ItemState = "orphaned" ItemStateAvailable ItemState = string(rivertype.JobStateAvailable) ItemStateCancelled ItemState = string(rivertype.JobStateCancelled) ItemStateCompleted ItemState = string(rivertype.JobStateCompleted) ItemStateDiscarded ItemState = string(rivertype.JobStateDiscarded) ItemStatePending ItemState = string(rivertype.JobStatePending) ItemStateRetryable ItemState = string(rivertype.JobStateRetryable) ItemStateRunning ItemState = string(rivertype.JobStateRunning) ItemStateScheduled ItemState = string(rivertype.JobStateScheduled) )
type ItemStatus ¶
type ItemStatus struct {
Result json.RawMessage `gorm:"column:result;type:jsonb" json:"result,omitempty"`
ItemKey string `gorm:"column:item_key" json:"item_key"`
State string `gorm:"column:state" json:"state"`
}
ItemStatus reports the state of a single item within a batch.
type ProcessBatchArgs ¶
ProcessBatchArgs is a River job that the API enqueues when a new batch is created. The daemon picks this up, refreshes the symptom cache, and fans out into individual ReevaluateJobRunArgs jobs.
func (ProcessBatchArgs) InsertOpts ¶
func (ProcessBatchArgs) InsertOpts() river.InsertOpts
InsertOpts configures River insertion behavior. Batch processing is not retried automatically; individual items have their own retry policy.
func (ProcessBatchArgs) Kind ¶
func (ProcessBatchArgs) Kind() string
Kind returns the River job kind identifier.
type ProcessBatchWorker ¶
type ProcessBatchWorker struct {
river.WorkerDefaults[ProcessBatchArgs]
// contains filtered or unexported fields
}
ProcessBatchWorker handles ProcessBatchArgs River jobs. When the daemon picks up a new batch, this worker refreshes the symptom cache, fans out individual re-evaluation River jobs, and tracks which items were enqueued vs deduplicated.
func NewProcessBatchWorker ¶
func NewProcessBatchWorker(reEvaluator *jobrunscan.ReEvaluator, gormDB *gorm.DB) *ProcessBatchWorker
NewProcessBatchWorker creates a ProcessBatchWorker. The riverClient field must be set via SetRiverClient before the worker processes any jobs, because the River client cannot be created until after the worker is registered (circular dependency resolved by deferred wiring).
func (*ProcessBatchWorker) SetRiverClient ¶
func (w *ProcessBatchWorker) SetRiverClient(client *river.Client[pgx.Tx])
SetRiverClient wires the River client after construction. Must be called before the River client is started.
func (*ProcessBatchWorker) Work ¶
func (w *ProcessBatchWorker) Work(ctx context.Context, job *river.Job[ProcessBatchArgs]) error
Work processes a batch: refreshes symptoms, fans out individual River jobs, and updates batch/item state. MaxAttempts is 1 (no automatic retry); on failure the batch stays in "processing" and must be resubmitted.
type ReevaluateFunc ¶
type ReevaluateFunc func(ctx context.Context, prowJobBuildID string, dryRun bool) (*jobrunscan.ReEvaluationResult, error)
ReevaluateFunc evaluates a job run against cached symptom definitions. A dry run returns results without persisting labels. This seam allows worker tests to exercise River behavior without GCS or BigQuery credentials.
type ReevaluateJobRunArgs ¶
type ReevaluateJobRunArgs struct {
ProwJobBuildID string `json:"prow_job_build_id" river:"unique"`
SymptomHash string `json:"symptom_hash" river:"unique"`
DryRun bool `json:"dry_run" river:"unique"`
}
ReevaluateJobRunArgs is a River job for re-evaluating a single job run's symptoms. ProwJobBuildID, SymptomHash, and DryRun all participate in River's uniqueness check, so the same job run is re-evaluated if symptoms change, dry runs do not deduplicate against real runs, and duplicate requests within a DedupPeriod are skipped.
func (ReevaluateJobRunArgs) InsertOpts ¶
func (ReevaluateJobRunArgs) InsertOpts() river.InsertOpts
InsertOpts configures River insertion behavior with deduplication scoped to the combination of job run identity and symptom state hash.
func (ReevaluateJobRunArgs) Kind ¶
func (ReevaluateJobRunArgs) Kind() string
Kind returns the River job kind identifier.
type ReevaluateWorker ¶
type ReevaluateWorker struct {
river.WorkerDefaults[ReevaluateJobRunArgs]
// Reevaluate supplies the single-run evaluation and must be set before Work.
Reevaluate ReevaluateFunc
}
ReevaluateWorker handles individual ReevaluateJobRunArgs River jobs by delegating to the ReEvaluator's cached symptom evaluation.
func NewReevaluateWorker ¶
func NewReevaluateWorker(reEvaluator *jobrunscan.ReEvaluator) *ReevaluateWorker
NewReevaluateWorker creates a ReevaluateWorker that delegates to the ReEvaluator's cached evaluation method.
func (*ReevaluateWorker) Work ¶
func (w *ReevaluateWorker) Work(ctx context.Context, job *river.Job[ReevaluateJobRunArgs]) error
Work re-evaluates symptoms for a single job run. Transient errors trigger River's retry logic (up to MaxAttemptsPerItem attempts with exponential backoff). Results are recorded by River when the attempt finishes. Permanent errors (e.g. missing job run) cancel the job immediately.
type StatusQuerier ¶
type StatusQuerier struct {
// contains filtered or unexported fields
}
StatusQuerier queries the current status of a symptom re-evaluation batch by joining batch items with River job state.
func NewStatusQuerier ¶
func NewStatusQuerier(gormDB *gorm.DB) *StatusQuerier
NewStatusQuerier creates a StatusQuerier.
func (*StatusQuerier) GetUpdated ¶
func (q *StatusQuerier) GetUpdated(ctx context.Context, batchID uuid.UUID) (*BatchStatusResponse, error)
GetUpdated loads a batch and its items, joining with river_job to get current states. It performs lazy completion detection: when all items have reached a terminal state, the batch is marked complete (or failed if all items failed) and completed_at is set. This is idempotent.
type SubmitResult ¶
SubmitResult is returned by Submitter.Submit with the batch ID and requested count so the caller can return it to the user.
type Submitter ¶
type Submitter struct {
// contains filtered or unexported fields
}
Submitter creates batch specifications for symptom re-evaluation. It writes the batch and item rows via GORM and enqueues a ProcessBatchArgs River job so the daemon discovers the new batch.
func NewSubmitter ¶
NewSubmitter creates a Submitter with the given GORM DB and River client. The River client may be insert-only (API server) since Submitter only inserts jobs, never processes them.
func (*Submitter) Submit ¶
func (s *Submitter) Submit(ctx context.Context, prowJobBuildIDs []string, dryRun bool) (*SubmitResult, error)
Submit creates a batch specification and enqueues a River job for daemon processing. The batch and item rows are written via GORM (pgx/v4); the River job is inserted via the River client (pgx/v5). These are separate transactions, so the data can end up inconsistent: if the GORM write succeeds but the River insert fails, attempt to change batch state to "cancelled" so clients know to retry, and if that also fails, the batch remains in "pending" status until cleaned up.