queue

package
v0.10.1 Latest Latest
Warning

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

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

Documentation

Overview

Package queue is the Redis-backed async ingestion layer for docz-api. It enqueues per-repo ingest jobs (asynq over Redis), coalesces bursts within a debounce window so the latest HEAD wins, and runs a worker pool that drives the Phase-2 ingest pipeline. Triggers (the manual onboard flag now, webhooks in Phase 5) enqueue through the Enqueuer interface and return promptly; the worker processes jobs with at-least-once delivery and retry, relying on the store's content-hash gate to keep re-runs cheap and idempotent.

Index

Constants

View Source
const TaskTypeIngest = "ingest:repo"

TaskTypeIngest is the asynq task type for a repository ingest job. The worker registers one handler for it.

Variables

This section is empty.

Functions

This section is empty.

Types

type Client

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

Client is the Redis-backed enqueue client. One Client serves the whole process; it is safe for concurrent use. It also holds a plain go-redis client so /readyz can probe Redis reachability, and an asynq Inspector so a task id left behind by a finished run cannot block future triggers.

func NewClient

func NewClient(redisURL string, debounce time.Duration) (*Client, error)

NewClient builds a Client from a redis:// URL and a debounce window. The URL is parsed for the asynq client, the asynq inspector, and a go-redis client used by Ping.

func (*Client) Close

func (c *Client) Close() error

Close releases the asynq client, the inspector, and the go-redis client. Safe to call once at shutdown, after the worker has drained.

func (*Client) EnqueueIngest

func (c *Client) EnqueueIngest(ctx context.Context, job *IngestJob) error

EnqueueIngest schedules an ingest job for job.Owner/job.Name. The task id is "ingest:<owner>/<name>", so a second enqueue for the same repo while a run is pending returns ErrTaskIDConflict — treated as coalesced, since the pending job already covers the trigger. ProcessIn(debounce) delays execution so a burst of triggers collapses to one run at the latest HEAD.

A conflicting id does not always mean a run is pending, so the conflict is inspected rather than assumed: see resolveTaskIDConflict.

Known gap: a trigger arriving while the job is ACTIVE is dropped, because asynq holds the task id until the active run completes. The next trigger re-enqueues once the run finishes, and the content-hash gate makes any redundant re-run a cheap no-op.

func (*Client) Ping

func (c *Client) Ping(ctx context.Context) error

Ping verifies Redis is reachable; it backs the /readyz probe.

type Enqueuer

type Enqueuer interface {
	EnqueueIngest(ctx context.Context, job *IngestJob) error
}

Enqueuer is the consumer-side interface for enqueueing an ingest job. Declared here (matching ingest.Indexer / httpapi.Searcher) so callers depend on the interface, not the Redis implementation. *Client satisfies it.

type IngestJob

type IngestJob struct {
	InstallationID int64  `json:"installation_id"`
	Owner          string `json:"owner"`
	Name           string `json:"name"`
	Reason         string `json:"reason"`
	// Trace context (W3C traceparent/tracestate) captured at the enqueue site so
	// the worker's span links back to the request that triggered it. Empty when
	// the trigger ran outside a trace (e.g. the -onboard CLI); safe either way.
	TraceParent string `json:"traceparent,omitempty"`
	TraceState  string `json:"tracestate,omitempty"`
}

IngestJob is the payload for one repository ingest task. Reason is a human-readable label ("onboard", "webhook", "manual") for logging; it does not affect processing.

The job carries no HEAD SHA: the worker always fetches current HEAD at process time, so a job delayed by the debounce window sees the truly latest state. This is the "latest-HEAD wins" property — it falls out for free.

type Ingestor

type Ingestor interface {
	Run(ctx context.Context, installationID int64, owner, name string) (store.ReconcileResult, error)
}

Ingestor is the narrow surface the worker needs to run one ingest. It matches ingest.Service.Run; the production implementation (in the composition root) builds a per-installation GitHub client per job. Declared here (consumer side) so the worker is testable with a fake.

type Worker

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

Worker runs the asynq server that drains ingest jobs. Callers Start it (non-blocking) and Shutdown it (drains in-flight jobs). The pointer receiver is required: Worker holds *asynq.Server, which must not be copied.

func NewWorker

func NewWorker(redisURL string, concurrency int, ing Ingestor) (*Worker, error)

NewWorker builds a Worker that processes ingest jobs with ing. concurrency bounds the number of parallel ingests (2–4 suits a homelab single binary: each job holds a pool connection and issues GitHub API calls). The asynq server connects to Redis only on Start, so NewWorker never blocks on Redis.

func (*Worker) Shutdown

func (w *Worker) Shutdown()

Shutdown stops the asynq server gracefully: it stops accepting new tasks and blocks until in-flight handlers return. Call it after the HTTP server has drained, so no new enqueues arrive during the drain.

func (*Worker) Start

func (w *Worker) Start() error

Start registers the ingest handler and starts the asynq server (non-blocking).

Jump to

Keyboard shortcuts

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