workers

package
v0.0.0-...-775003a Latest Latest
Warning

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

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

Documentation

Overview

Package workers provides River job workers (e.g. webhook delivery, feedback embedding).

Index

Constants

View Source
const WebhookDeliveryTimeoutBuffer = 5 * time.Second

WebhookDeliveryTimeoutBuffer is added to HTTP timeout for the River job timeout.

Variables

This section is empty.

Functions

func NewRiverWorkersAndQueues

func NewRiverWorkersAndQueues(
	cfg *config.Config, deps RiverDeps,
) (*river.Workers, map[string]river.QueueConfig)

NewRiverWorkersAndQueues builds River workers and queue config from cfg and deps. Each optional worker is registered only when its client is set, and its queue is declared alongside it, so a disabled enrichment leaves no queue for jobs to pile up on unworked.

Only hub-worker calls this: the API inserts jobs through an insert-only River client and registers nothing (see cmd/api/app.go), so the queue MaxWorkers below are always the real configured concurrency.

Types

type EmbeddingReconcileSweeper

type EmbeddingReconcileSweeper interface {
	Sweep(ctx context.Context) (service.EmbeddingReconcileResult, error)
}

EmbeddingReconcileSweeper is the service boundary used by the periodic worker.

type EmbeddingReconcileWorker

type EmbeddingReconcileWorker struct {
	river.WorkerDefaults[service.EmbeddingReconcileArgs]
	// contains filtered or unexported fields
}

EmbeddingReconcileWorker runs one bounded, level-triggered taxonomy embedding sweep.

func NewEmbeddingReconcileWorker

func NewEmbeddingReconcileWorker(
	sweeper EmbeddingReconcileSweeper,
	metrics observability.EmbeddingMetrics,
) *EmbeddingReconcileWorker

NewEmbeddingReconcileWorker creates the periodic sweep worker.

func (*EmbeddingReconcileWorker) Timeout

Timeout gives the database scan and bounded insert batch its own deadline rather than inheriting River's unrelated global job timeout.

func (*EmbeddingReconcileWorker) Work

Work executes one sweep. A failure is returned for observability, but scheduled jobs use one attempt because the next interval is the retry and sees the same level-triggered backlog.

type EnrichmentReconcileWorker

type EnrichmentReconcileWorker struct {
	river.WorkerDefaults[service.EnrichmentReconcileArgs]
	// contains filtered or unexported fields
}

EnrichmentReconcileWorker runs one reconcile sweep.

The job carries no state of its own: everything it needs is the database's current contents and the deployment config, both read at run time. That is what lets the periodic schedule be simple — a missed tick costs nothing, because the next one sees the same backlog plus whatever arrived.

func NewEnrichmentReconcileWorker

func NewEnrichmentReconcileWorker(
	sweeper ReconcileSweeper, metrics observability.EnrichmentReconcileMetrics,
) *EnrichmentReconcileWorker

NewEnrichmentReconcileWorker creates the reconcile worker. metrics may be nil.

func (*EnrichmentReconcileWorker) Timeout

Timeout bounds one sweep at enrichmentReconcileTimeout.

Declaring it is what makes that constant real. River falls back to Config.JobTimeout when a worker returns zero, and RIVER_JOB_TIMEOUT_SECONDS defaults to 0, which River in turn reads as its own one-minute default -- so without this the sweep would be cancelled at a minute and the context.WithTimeout in Work could not extend it, a child context being unable to outlive its parent. With MaxAttempts 1 that cancellation is a discarded job and an error log per tick, which is exactly the "run of them" signal Work describes as meaning coverage has stopped.

func (*EnrichmentReconcileWorker) Work

Work runs the sweep.

A failed sweep returns an error so River retries it, but the retry is not what makes the system converge — the next scheduled tick is. That matters for how alarming a failure here should look: one is noise, a run of them means coverage has stopped being guaranteed.

type FailureRecorder

type FailureRecorder interface {
	RecordFailure(ctx context.Context, failure models.EnrichmentFailure) error
}

FailureRecorder persists a durable marker for an enrichment that gave up on a record. Satisfied by repository.EnrichmentFailuresRepository; declared here so the worker depends on the behaviour rather than on the repository.

type FeedbackEmbeddingWorker

type FeedbackEmbeddingWorker struct {
	river.WorkerDefaults[service.FeedbackEmbeddingArgs]
	// contains filtered or unexported fields
}

FeedbackEmbeddingWorker generates and stores embeddings for feedback records.

func NewFeedbackEmbeddingWorker

func NewFeedbackEmbeddingWorker(
	embeddingService feedbackEmbeddingService,
	embeddingClient service.EmbeddingClient,
	docPrefix string,
	metrics observability.EmbeddingMetrics,
) *FeedbackEmbeddingWorker

NewFeedbackEmbeddingWorker creates a worker that fetches the record, calls the embedding client, and stores the result. docPrefix is the prefix for document text. Can be empty for some providers. metrics may be nil when metrics are disabled.

func NewFeedbackEmbeddingWorkerWithOptions

func NewFeedbackEmbeddingWorkerWithOptions(
	embeddingService feedbackEmbeddingService,
	embeddingClient service.EmbeddingClient,
	docPrefix string,
	metrics observability.EmbeddingMetrics,
	jobTimeout time.Duration,
	failures FailureRecorder,
	failureMetrics observability.EnrichmentFailureMetrics,
) *FeedbackEmbeddingWorker

NewFeedbackEmbeddingWorkerWithOptions creates an embedding worker with an explicit job deadline and durable failure recording. The simple constructor remains for tools/tests that do not need deployment wiring.

func (*FeedbackEmbeddingWorker) Timeout

Timeout limits how long a single embedding job can run.

func (*FeedbackEmbeddingWorker) Work

Work loads the record, generates or clears the embedding, and persists it.

type FeedbackEmotionsWorker

type FeedbackEmotionsWorker = enrichmentWorker[service.FeedbackEmotionsArgs, service.EmotionsResult]

FeedbackEmotionsWorker classifies a feedback record's value_text into a set of emotion labels and stores it — a configured enrichmentWorker. It borrows the shared rate-limit snooze (it, too, calls a rate-limited LLM provider); it has no per-tenant target, but its persist is guarded against content supersession (a job that read older text skips instead of landing its labels last — stale non-NULL labels would escape the NULL-rows-only backfill forever).

func NewFeedbackEmotionsWorker

func NewFeedbackEmotionsWorker(
	svc emotionsWorkerService, resolver tenantSettingsReader,
	client service.EmotionsClient, metrics observability.EmotionsMetrics, failures FailureRecorder,
	failureMetrics observability.EnrichmentFailureMetrics,
) *FeedbackEmotionsWorker

NewFeedbackEmotionsWorker creates a worker that fetches the record, classifies its value_text, and stores the emotion labels. metrics may be nil when metrics are disabled.

type FeedbackRecordsPurgeWorker

type FeedbackRecordsPurgeWorker struct {
	river.WorkerDefaults[service.FeedbackRecordsPurgeArgs]
	// contains filtered or unexported fields
}

FeedbackRecordsPurgeWorker deletes every feedback record for one tenant, everything derived from those records, and the taxonomy built on them, keeping the tenant's configuration. It runs off the request path because the delete is unbounded and would otherwise outlive the API server's write timeout on a large tenant.

func NewFeedbackRecordsPurgeWorker

func NewFeedbackRecordsPurgeWorker(svc feedbackRecordsPurgeService) *FeedbackRecordsPurgeWorker

NewFeedbackRecordsPurgeWorker creates the worker.

func (*FeedbackRecordsPurgeWorker) Timeout

Timeout limits how long a single purge attempt can run.

func (*FeedbackRecordsPurgeWorker) Work

Work purges the tenant's feedback records. The service re-scopes the work to the tenant in the job args rather than trusting enqueue-time validation.

A failure is logged here rather than left to River. River reports retries and final discards at INFO, and hub-worker does not set river.Config.Logger, so its fallback logger sits at WARN and filters both — a failing purge would otherwise leave no trace outside the river_job.errors column while the dataset sits partly emptied. Matches webhook_dispatch and feedback_embedding in separating a retryable attempt from the last one.

type FeedbackSentimentWorker

type FeedbackSentimentWorker = enrichmentWorker[service.FeedbackSentimentArgs, service.SentimentResult]

FeedbackSentimentWorker classifies a feedback record's value_text into a sentiment label and score and stores it — a configured enrichmentWorker. It borrows the shared rate-limit snooze (it, too, calls a rate-limited LLM provider); it has no per-tenant target, but its persist is guarded against content supersession (a job that read older text skips instead of landing its label last — a stale non-NULL label would escape the NULL-rows-only backfill forever).

func NewFeedbackSentimentWorker

func NewFeedbackSentimentWorker(
	svc sentimentWorkerService, resolver tenantSettingsReader,
	client service.SentimentClient, metrics observability.SentimentMetrics, failures FailureRecorder,
	failureMetrics observability.EnrichmentFailureMetrics,
) *FeedbackSentimentWorker

NewFeedbackSentimentWorker creates a worker that fetches the record, classifies its value_text, and stores the result. metrics may be nil when metrics are disabled.

type FeedbackTranslationWorker

type FeedbackTranslationWorker = enrichmentWorker[service.FeedbackTranslationArgs, string]

FeedbackTranslationWorker translates a feedback record's value_text into the tenant's target language and stores it — a configured enrichmentWorker. It borrows the shared rate-limit snooze (it calls a rate-limited LLM provider) and uses the supersession skip: a stale-target OR stale-content write is a no-op once a newer job owns the row.

func NewFeedbackTranslationWorker

func NewFeedbackTranslationWorker(
	svc translationWorkerService, client service.TranslationClient, metrics observability.TranslationMetrics,
	failures FailureRecorder,
	failureMetrics observability.EnrichmentFailureMetrics,
) *FeedbackTranslationWorker

NewFeedbackTranslationWorker creates a worker that fetches the record, translates its value_text into the target language (or copies it when the source already matches), and stores the result. metrics may be nil when metrics are disabled.

type ReconcileSweeper

type ReconcileSweeper interface {
	Sweep(ctx context.Context) (service.ReconcileResult, error)
}

ReconcileSweeper is the service half, declared here so the worker depends on the behaviour rather than the implementation. Exported because hub-worker has to name the type to convert a nil *EnrichmentReconcileService into a nil interface — a typed nil in an interface would register the worker and panic on the first sweep.

type RiverDeps

type RiverDeps struct {
	// Webhook worker
	WebhooksRepo       webhookDispatchRepo
	WebhookSender      service.WebhookSender
	WebhookHTTPTimeout time.Duration
	WebhookMetrics     observability.WebhookMetrics

	// Embedding worker (optional; if EmbeddingClient is nil, embedding worker is not registered)
	EmbeddingService   feedbackEmbeddingService
	EmbeddingClient    service.EmbeddingClient
	EmbeddingDocPrefix string
	EmbeddingMetrics   observability.EmbeddingMetrics
	// EmbeddingReconcileSweeper is non-nil only when automatic taxonomy embedding repair is enabled.
	EmbeddingReconcileSweeper EmbeddingReconcileSweeper

	// Translation worker (optional; if TranslationClient is nil, translation worker is not registered)
	TranslationService translationWorkerService
	TranslationClient  service.TranslationClient
	TranslationMetrics observability.TranslationMetrics
	// Per-tenant translation backfill worker (registered alongside the translation worker).
	TranslationBackfillService tenantTranslationBackfillService
	TranslationMaxAttempts     int

	// Sentiment worker (optional; if SentimentClient is nil, sentiment worker is not registered)
	SentimentService  sentimentWorkerService
	SentimentResolver tenantSettingsReader
	SentimentClient   service.SentimentClient
	SentimentMetrics  observability.SentimentMetrics

	// Emotions worker (optional; if EmotionsClient is nil, emotions worker is not registered)
	EmotionsService  emotionsWorkerService
	EmotionsResolver tenantSettingsReader
	EmotionsClient   service.EmotionsClient
	EmotionsMetrics  observability.EmotionsMetrics

	// Feedback-records purge worker (always registered; the purge is a core tenant operation, not
	// an enrichment, so it has no client to gate on).
	FeedbackRecordsPurgeService feedbackRecordsPurgeService

	// ReconcileSweeper runs the level-triggered enrichment sweep. nil leaves the reconciler out
	// entirely — the kill switch, and the shape a deployment with no enrichment provider takes.
	ReconcileSweeper ReconcileSweeper

	// Failures records the durable marker a classify worker writes when it gives up on a record.
	// Shared by the three classify pipelines; nil disables recording, which leaves the API
	// under-reporting failures but changes no enrichment behaviour.
	Failures FailureRecorder
	// ReconcileMetrics reports the sweep's own outcome, duration and enqueue counts. nil disables
	// them; the sweep still runs.
	ReconcileMetrics observability.EnrichmentReconcileMetrics
	// FailureMetrics counts permanent give-ups by cause, for whoever watches the deployment
	// rather than a single tenant. nil disables it.
	FailureMetrics observability.EnrichmentFailureMetrics
}

RiverDeps holds dependencies required to build River workers and queue config. Each optional group is gated on its client: leave a client nil and neither its worker nor its queue is registered.

type TenantTranslationBackfillWorker

type TenantTranslationBackfillWorker struct {
	river.WorkerDefaults[service.TenantTranslationBackfillArgs]
	// contains filtered or unexported fields
}

TenantTranslationBackfillWorker fans out a per-tenant re-translation: it lists the tenant's stale text records and enqueues a FeedbackTranslationArgs job for each. It is enqueued when a tenant's translation settings change (see service.TenantTranslationBackfillArgs) and runs off the request path.

func NewTenantTranslationBackfillWorker

func NewTenantTranslationBackfillWorker(
	svc tenantTranslationBackfillService, maxAttempts int,
) *TenantTranslationBackfillWorker

NewTenantTranslationBackfillWorker creates the worker. maxAttempts is applied to the per-record translation jobs it enqueues.

func (*TenantTranslationBackfillWorker) Timeout

Timeout limits how long a single tenant backfill fan-out can run.

func (*TenantTranslationBackfillWorker) Work

Work lists the tenant's stale records and enqueues per-record translation jobs onto the translation RECONCILE queue, not the live one. A settings change can fan out a tenant's entire history, and on the live queue that backlog sits in front of translations for records arriving right now -- the same starvation the reconciler's separate lanes exist to prevent. Both are bulk catch-up work, so both belong in the lane sized for it. The River client is obtained from the context (the only place River sets it) and handed to the service as the inserter, preserving the injected-inserter seam.

type WebhookDispatchWorker

type WebhookDispatchWorker struct {
	river.WorkerDefaults[service.WebhookDispatchArgs]
	// contains filtered or unexported fields
}

WebhookDispatchWorker delivers one event to one webhook endpoint.

func NewWebhookDispatchWorker

func NewWebhookDispatchWorker(
	repo webhookDispatchRepo, sender service.WebhookSender, httpTimeout time.Duration,
	metrics observability.WebhookMetrics,
) *WebhookDispatchWorker

NewWebhookDispatchWorker creates a worker that uses the given repo and sender. httpTimeout is the webhook HTTP client timeout; job timeout is httpTimeout + WebhookDeliveryTimeoutBuffer. metrics may be nil when metrics are disabled.

func (*WebhookDispatchWorker) Timeout

Timeout limits how long a single delivery can run (HTTP timeout + buffer).

func (*WebhookDispatchWorker) Work

Work loads the webhook, builds the payload, and sends once.

Jump to

Keyboard shortcuts

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