symptomre

package
v0.0.0-...-5fad17d Latest Latest
Warning

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

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

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

View Source
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
)
View Source
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

View Source
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.

func (Batch) TableName

func (Batch) TableName() string

TableName returns the PostgreSQL table name for Batch.

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

func NewBatchCanceller(gormDB *gorm.DB, riverClient *river.Client[pgx.Tx]) *BatchCanceller

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.

func (*BatchCleanupProcess) Run

Run executes the periodic cleanup loop, deleting old batches every hour until the context is canceled.

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.

func (BatchItem) TableName

func (BatchItem) TableName() string

TableName returns the PostgreSQL table name for BatchItem.

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

type ProcessBatchArgs struct {
	BatchID uuid.UUID `json:"batch_id"`
}

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

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

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

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

type SubmitResult struct {
	BatchID   uuid.UUID `json:"batch_id"`
	Requested int       `json:"requested"`
}

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

func NewSubmitter(gormDB *gorm.DB, riverClient *river.Client[pgx.Tx]) *Submitter

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.

Jump to

Keyboard shortcuts

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