Documentation
¶
Overview ¶
Package workers provides River job workers (e.g. webhook delivery, feedback embedding).
Index ¶
- Constants
- func NewRiverWorkersAndQueues(cfg *config.Config, deps RiverDeps) (*river.Workers, map[string]river.QueueConfig)
- type EmbeddingReconcileSweeper
- type EmbeddingReconcileWorker
- type EnrichmentReconcileWorker
- type FailureRecorder
- type FeedbackEmbeddingWorker
- type FeedbackEmotionsWorker
- type FeedbackRecordsPurgeWorker
- type FeedbackSentimentWorker
- type FeedbackTranslationWorker
- type ReconcileSweeper
- type RiverDeps
- type TenantTranslationBackfillWorker
- type WebhookDispatchWorker
Constants ¶
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 ¶
func (w *EmbeddingReconcileWorker) Timeout(*river.Job[service.EmbeddingReconcileArgs]) time.Duration
Timeout gives the database scan and bounded insert batch its own deadline rather than inheriting River's unrelated global job timeout.
func (*EmbeddingReconcileWorker) Work ¶
func (w *EmbeddingReconcileWorker) Work( ctx context.Context, _ *river.Job[service.EmbeddingReconcileArgs], ) error
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 ¶
func (w *EnrichmentReconcileWorker) Timeout(*river.Job[service.EnrichmentReconcileArgs]) time.Duration
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 ¶
func (w *EnrichmentReconcileWorker) Work( ctx context.Context, _ *river.Job[service.EnrichmentReconcileArgs], ) error
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 ¶
func (w *FeedbackEmbeddingWorker) Timeout(*river.Job[service.FeedbackEmbeddingArgs]) time.Duration
Timeout limits how long a single embedding job can run.
func (*FeedbackEmbeddingWorker) Work ¶
func (w *FeedbackEmbeddingWorker) Work(ctx context.Context, job *river.Job[service.FeedbackEmbeddingArgs]) error
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 ¶
func (w *FeedbackRecordsPurgeWorker) Timeout(*river.Job[service.FeedbackRecordsPurgeArgs]) time.Duration
Timeout limits how long a single purge attempt can run.
func (*FeedbackRecordsPurgeWorker) Work ¶
func (w *FeedbackRecordsPurgeWorker) Work( ctx context.Context, job *river.Job[service.FeedbackRecordsPurgeArgs], ) error
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 ¶
func (w *TenantTranslationBackfillWorker) Timeout(*river.Job[service.TenantTranslationBackfillArgs]) time.Duration
Timeout limits how long a single tenant backfill fan-out can run.
func (*TenantTranslationBackfillWorker) Work ¶
func (w *TenantTranslationBackfillWorker) Work( ctx context.Context, job *river.Job[service.TenantTranslationBackfillArgs], ) error
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 ¶
func (w *WebhookDispatchWorker) Timeout(*river.Job[service.WebhookDispatchArgs]) time.Duration
Timeout limits how long a single delivery can run (HTTP timeout + buffer).
func (*WebhookDispatchWorker) Work ¶
func (w *WebhookDispatchWorker) Work(ctx context.Context, job *river.Job[service.WebhookDispatchArgs]) error
Work loads the webhook, builds the payload, and sends once.