service

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: 45 Imported by: 0

Documentation

Overview

Package service implements business logic for feedback records.

Index

Constants

View Source
const (
	EmbeddingProviderOpenAI       = ProviderOpenAI
	EmbeddingProviderGoogle       = ProviderGoogle
	EmbeddingProviderGoogleGemini = ProviderGoogleGemini
)

Embedding provider names for NewEmbeddingClient.

View Source
const (
	// EmbeddingReconcileQueueName carries the level-triggered sweep itself. The queue is serialized;
	// River's elected leader schedules one job per interval across all worker replicas.
	EmbeddingReconcileQueueName = "embedding_reconcile"
	// EmbeddingsReconcileQueueName carries repaired taxonomy embeddings. It is deliberately
	// separate from live embeddings so old backlog cannot delay a record that just arrived.
	EmbeddingsReconcileQueueName = "embeddings_reconcile"
)
View Source
const (
	ProviderOpenAI       = "openai"
	ProviderGoogle       = "google"
	ProviderGoogleGemini = "google-gemini"
)

Provider identifiers shared by every enrichment client — sentiment, translation, and embedding all reuse the same OpenAI and Google SDK wrappers, so the provider names are common. Per-type factories alias these (e.g. SentimentProviderOpenAI) for readability.

View Source
const (
	// EnrichmentReconcileQueueName carries the sweep itself. Its own queue with MaxWorkers 1, so
	// "one sweep at a time" is structural rather than only enforced by the job's uniqueness, and a
	// long scan cannot occupy a worker the webhook dispatcher needs.
	EnrichmentReconcileQueueName = "enrichment_reconcile"

	SentimentsReconcileQueueName   = "sentiments_reconcile"
	EmotionsReconcileQueueName     = "emotions_reconcile"
	TranslationsReconcileQueueName = "translations_reconcile"
)

The reconcile queues. One per enrichment, separate from the live queue that the event path feeds — and named `_reconcile` rather than `_backfill` deliberately, because the pre-existing `translation_backfills` (the tenant fan-out) is one character-transposition away from a `translations_backfill`, and a mix-up there compiles, passes the parity test, and silently swaps the two queues' concurrency budgets.

Separation is the rate control. The reconciler can enqueue a large sweep without a record submitted right now having to queue behind it, because the two draw from different MaxWorkers budgets — so a backlog drains in the background at whatever rate its own workers allow while live enrichment keeps its latency. Trying to achieve the same thing by pacing inserts onto the live queue means guessing a rate that matches drain speed, and guessing wrong in the slow direction is indistinguishable from the reconciler not running.

It also fixes something that predates this: the one-off backfill commands enqueue onto the LIVE queues today (see cmd/backfill-classify), so running one already starves live enrichment.

View Source
const (

	// FeedbackRecordsPurgeQueueName is the River queue for tenant feedback-record purges. It is
	// kept separate from the enrichment queues so a large purge cannot starve live ingestion
	// throughput, mirroring TranslationBackfillsQueueName.
	FeedbackRecordsPurgeQueueName = "feedback_records_purges"
	// FeedbackRecordsPurgeMaxAttempts is the retry budget for a purge. There is no config knob for
	// it: the purge is idempotent and tenant-scoped, so a retry is always safe and always makes
	// progress, and the failures worth retrying (a lock timeout while writes drain, a rescued job)
	// clear on their own within a few attempts.
	FeedbackRecordsPurgeMaxAttempts = 5
)
View Source
const (

	// EmbeddingsQueueName is the River queue used for feedback embedding jobs.
	EmbeddingsQueueName = "embeddings"
)
View Source
const (

	// EmotionsQueueName is the River queue used for feedback emotion jobs. It is distinct from the
	// sentiment queue: River uniqueness is per-kind, so emotions and sentiment dedupe independently.
	EmotionsQueueName = "emotions"
)
View Source
const EnrichmentRetryCooldown = time.Hour

EnrichmentRetryCooldown is how long a tenant must wait between clears of the same enrichment.

The number is a judgement, not a measurement: long enough that clear-and-sweep cannot be run in a loop, short enough that a genuine "the provider changed its policy, try again" is not an overnight wait. Bounds the amplification to (size of terminal set / this).

View Source
const (

	// SentimentsQueueName is the River queue used for feedback sentiment jobs.
	SentimentsQueueName = "sentiments"
)
View Source
const SigningKeySize = 32

SigningKeySize is the number of random bytes for Standard Webhooks signing keys.

View Source
const (

	// TranslationBackfillsQueueName is the River queue for per-tenant translation
	// backfill fan-out jobs. It is kept separate from TranslationsQueueName so a large
	// fan-out does not starve live per-record translation throughput.
	TranslationBackfillsQueueName = "translation_backfills"
)
View Source
const (

	// TranslationsQueueName is the River queue used for feedback translation jobs.
	TranslationsQueueName = "translations"
)

Variables

View Source
var (
	// ErrEmbeddingBatcherClosed is returned after shutdown stops accepting document calls.
	ErrEmbeddingBatcherClosed = errors.New("embedding batcher is closed")
	// ErrEmbeddingBatchResultCount is returned when a batch provider violates one-result-per-input.
	ErrEmbeddingBatchResultCount = errors.New("embedding batch result count mismatch")
)
View Source
var (
	// ErrEmbeddingConfigInvalid is returned when the embedding provider is unsupported.
	ErrEmbeddingConfigInvalid = errors.New("embedding config invalid")
	// ErrEmbeddingProviderAPIKey is returned when an API-key-based provider is configured without a key.
	ErrEmbeddingProviderAPIKey = errors.New("EMBEDDING_PROVIDER_API_KEY is required for this provider")
	// ErrEmbeddingBaseURLUnsupported is returned when a custom base URL is configured for a non-openai provider.
	ErrEmbeddingBaseURLUnsupported = errors.New("EMBEDDING_BASE_URL is only supported for openai")
	// ErrEmbeddingGoogleGeminiConfig is returned when google-gemini is configured without project or location.
	ErrEmbeddingGoogleGeminiConfig = errors.New(
		"google-gemini requires EMBEDDING_GOOGLE_CLOUD_PROJECT and EMBEDDING_GOOGLE_CLOUD_LOCATION")
)
View Source
var (
	// ErrEmotionsConfigInvalid is returned when the emotions provider is unsupported.
	ErrEmotionsConfigInvalid = errors.New("emotions config invalid")
	// ErrEmotionsProviderAPIKey is returned when an API-key-based provider is configured without a key.
	ErrEmotionsProviderAPIKey = errors.New("EMOTIONS_PROVIDER_API_KEY is required for this provider")
	// ErrEmotionsBaseURLUnsupported is returned when a custom base URL is configured for a non-openai provider.
	ErrEmotionsBaseURLUnsupported = errors.New("EMOTIONS_BASE_URL is only supported for openai")
	// ErrEmotionsGoogleGeminiConfig is returned when google-gemini is configured without project or location.
	ErrEmotionsGoogleGeminiConfig = errors.New(
		"google-gemini requires EMOTIONS_GOOGLE_CLOUD_PROJECT and EMOTIONS_GOOGLE_CLOUD_LOCATION")
)
View Source
var (
	ErrMissingTenantID   = errors.New("tenant_id is required")
	ErrEmptyQuery        = errors.New("query is required and must be non-empty")
	ErrEmbeddingNotFound = repository.ErrEmbeddingNotFound
)

Sentinel errors for search (used by handlers for status mapping).

View Source
var (
	// ErrSentimentConfigInvalid is returned when the sentiment provider is unsupported.
	ErrSentimentConfigInvalid = errors.New("sentiment config invalid")
	// ErrSentimentProviderAPIKey is returned when an API-key-based provider is configured without a key.
	ErrSentimentProviderAPIKey = errors.New("SENTIMENT_PROVIDER_API_KEY is required for this provider")
	// ErrSentimentBaseURLUnsupported is returned when a custom base URL is configured for a non-openai provider.
	ErrSentimentBaseURLUnsupported = errors.New("SENTIMENT_BASE_URL is only supported for openai")
	// ErrSentimentGoogleGeminiConfig is returned when google-gemini is configured without project or location.
	ErrSentimentGoogleGeminiConfig = errors.New(
		"google-gemini requires SENTIMENT_GOOGLE_CLOUD_PROJECT and SENTIMENT_GOOGLE_CLOUD_LOCATION")
)
View Source
var (
	// ErrTaxonomyServiceURLRequired is returned when TAXONOMY_SERVICE_URL is missing.
	ErrTaxonomyServiceURLRequired = errors.New("TAXONOMY_SERVICE_URL is required")

	// ErrTaxonomyServiceTokenRequired is returned when TAXONOMY_SERVICE_TOKEN is missing.
	ErrTaxonomyServiceTokenRequired = errors.New("TAXONOMY_SERVICE_TOKEN is required")

	// ErrTaxonomyServiceUnexpectedStatus is returned when the taxonomy service returns a non-2xx response.
	ErrTaxonomyServiceUnexpectedStatus = errors.New("taxonomy service returned non-success status")
)
View Source
var (
	// ErrTaxonomyEmbeddingsNotConfigured is returned when Hub embeddings are not configured.
	ErrTaxonomyEmbeddingsNotConfigured = errors.New("taxonomy requires EMBEDDING_MODEL to be configured")
	// ErrTaxonomyServiceNotConfigured is returned when the taxonomy compute service is unavailable.
	ErrTaxonomyServiceNotConfigured = errors.New("taxonomy service is not configured")
	// ErrTaxonomyServiceStartFailed is returned when the taxonomy compute service rejects a run.
	ErrTaxonomyServiceStartFailed = errors.New("taxonomy service failed to start run")
)
View Source
var (
	// ErrTranslationConfigInvalid is returned when the translation provider is unsupported.
	ErrTranslationConfigInvalid = errors.New("translation config invalid")
	// ErrTranslationProviderAPIKey is returned when an API-key-based provider is configured without a key.
	ErrTranslationProviderAPIKey = errors.New("TRANSLATION_PROVIDER_API_KEY is required for this provider")
	// ErrTranslationBaseURLUnsupported is returned when a custom base URL is configured for a non-openai provider.
	ErrTranslationBaseURLUnsupported = errors.New("TRANSLATION_BASE_URL is only supported for openai")
	// ErrTranslationGoogleGeminiConfig is returned when google-gemini is configured without project or location.
	ErrTranslationGoogleGeminiConfig = errors.New(
		"google-gemini requires TRANSLATION_GOOGLE_CLOUD_PROJECT and TRANSLATION_GOOGLE_CLOUD_LOCATION")
)
View Source
var (
	ErrWebhookGone   = errors.New("webhook returned 410 Gone (endpoint disabled)")
	ErrWebhookNon2xx = errors.New("webhook returned non-2xx status")
)

Sentinel errors for webhook delivery (err113).

View Source
var ErrEmbeddingBackfillNotConfigured = errors.New("embedding backfill not configured")

ErrEmbeddingBackfillNotConfigured is returned when BackfillEmbeddings is called without embedding inserter/queue.

View Source
var ErrEmbeddingReconcileInserterUnset = errors.New("embedding reconcile: inserter not set")

ErrEmbeddingReconcileInserterUnset reports a sweep invoked before the River client was attached.

View Source
var ErrEmotionsResponseInvalid = errors.New("emotions response invalid")

ErrEmotionsResponseInvalid is returned when the provider's structured output cannot be parsed. The array-of-enum schema keeps this rare; the worker treats it as a normal provider failure (retry, then fail). Unknown labels are dropped rather than failing the whole classification, so this fires only on a decode error, not on out-of-pool labels.

View Source
var ErrInvalidCursor = errors.New("invalid cursor")

ErrInvalidCursor is returned when the cursor parameter is malformed or invalid.

View Source
var ErrInvalidEmotionLabel = errors.New("invalid emotion label")

ErrInvalidEmotionLabel is returned when SetEmotions is given an unknown emotion label.

View Source
var ErrInvalidSentimentLabel = errors.New("invalid sentiment label")

ErrInvalidSentimentLabel is returned when SetSentiment is given an unknown sentiment label.

View Source
var ErrPaginationInvariantViolated = errors.New("pagination invariant violated: hasMore with empty list")

ErrPaginationInvariantViolated indicates hasMore was true with an empty list (repository invariant violation).

View Source
var ErrPurgeInserterNotConfigured = errors.New("feedback records purge inserter not configured")

ErrPurgeInserterNotConfigured is returned when a purge is requested from a process built without a job inserter (hub-worker performs purges but never enqueues them).

View Source
var ErrReconcileInserterUnset = errors.New("enrichment reconcile: inserter not set")

ErrReconcileInserterUnset marks a sweep attempted before the River client was attached.

View Source
var ErrSentimentResponseInvalid = errors.New("sentiment response invalid")

ErrSentimentResponseInvalid is returned when the provider's structured output cannot be parsed or carries an unknown sentiment label. Structured output makes this rare, but the worker still treats it as a normal provider failure (retry, then fail) rather than trusting an out-of-contract response.

View Source
var ErrSentimentScoreRequired = errors.New("sentiment score is required when a label is set")

ErrSentimentScoreRequired is returned when a sentiment label is set without a score: a label must carry its score (clearing, where sentiment is nil, nulls both columns).

View Source
var ErrTranslationLangKeyRequired = errors.New("translation lang key is required when translated text is set")

ErrTranslationLangKeyRequired is returned when a translation is set without a target locale key: a translation must record the locale it was produced in (clearing, where translated is nil, intentionally passes an empty key to null both columns).

View Source
var ErrUserIDRequired = huberrors.NewValidationError("user_id", "user_id is required")

ErrUserIDRequired is returned when deleting feedback records by user is called without user_id.

Functions

func BuildEmbeddingInput

func BuildEmbeddingInput(fieldLabel, valueText *string, prefix string) string

BuildEmbeddingInput prepares text for vector embedding. Uses a pre-allocated strings.Builder to reduce allocations on this hot path.

We feed the model raw, natural text: only TrimSpace and Unicode NFC normalization are applied. Case, diacritics, and punctuation are preserved so the model retains semantic clues (e.g. "US" vs "us", "résumé" vs "resume").

Arguments:

  • fieldLabel: The "question" or metadata key (e.g. "What is your reasoning?").
  • valueText: The "answer" or main content (e.g. "I chose option B because...").
  • prefix: Model-specific task instruction; OpenAI and Google use "".

Returns formatted string: "[prefix]Question: [label]\nAnswer: [value]" (or "[prefix][value]" when label is empty).

func BuildEmbeddingInputForKind

func BuildEmbeddingInputForKind(record *models.FeedbackRecord, kind models.EmbeddingInputKind, prefix string) string

BuildEmbeddingInputForKind prepares text for a specific embedding input kind.

func BuildEmbeddingInputFromValues

func BuildEmbeddingInputFromValues(
	fieldLabel, valueText, valueTextTranslated *string,
	kind models.EmbeddingInputKind,
	prefix string,
) string

BuildEmbeddingInputFromValues prepares text for vector embedding from raw record values. Taxonomy embeddings prefer translated text when present, falling back to original value_text.

func DecodeSearchCursor

func DecodeSearchCursor(cursor string) (distance float64, feedbackRecordID uuid.UUID, err error)

DecodeSearchCursor parses an opaque cursor and returns (distance, feedbackRecordID). Returns ErrInvalidCursor if the cursor is malformed.

func EmbeddingPrefixForProvider

func EmbeddingPrefixForProvider(provider string) string

EmbeddingPrefixForProvider returns the document prefix for the given embedding provider. Returns "" for unknown providers (and for all providers today, as none set a prefix).

func EncodeSearchCursor

func EncodeSearchCursor(distance float64, id uuid.UUID) (string, error)

EncodeSearchCursor returns an opaque cursor for the next page. distance is the cosine distance (e.embedding <=> query) of the last result row; id is that row's feedback_record_id.

func InFlightUniqueStates

func InFlightUniqueStates() []rivertype.JobState

InFlightUniqueStates is the state set every unique enqueue in this package uses: the states a job can be in while it still has work left to do.

It exists as one helper because getting it wrong is silent in both directions. Include `completed` — as River's default set does — and with no ByPeriod the uniqueness window is unbounded, so the first job for a given argument is the only one that ever runs while every later enqueue is skipped as a duplicate. Omit `retryable` and a job waiting out its backoff is enqueued a second time, doubling the work for something already being handled.

func JobQueueNames

func JobQueueNames() []string

JobQueueNames returns the distinct River queue names from JobKindSpecs, in declaration order. Several kinds may share a queue, so callers that need queues rather than kinds use this.

func NormalizeEmbeddingProvider

func NormalizeEmbeddingProvider(provider string) string

NormalizeEmbeddingProvider returns the canonical provider name (lowercase, trimmed), mapping the legacy google-vertex alias to google-gemini.

func ProviderRequiresAPIKey

func ProviderRequiresAPIKey(provider string) bool

ProviderRequiresAPIKey returns true for providers that require EMBEDDING_PROVIDER_API_KEY (from registry).

func ProviderRequiresGoogleGeminiConfig

func ProviderRequiresGoogleGeminiConfig(provider string) bool

ProviderRequiresGoogleGeminiConfig returns true for providers that require Google Cloud project and location.

func ReconcileQueueFor

func ReconcileQueueFor(enrichment string) string

ReconcileQueueFor maps an enrichment to the queue its reconciled work belongs on. Unknown names return "" — callers treat that as a wiring mistake rather than routing to the live queue by accident, which would defeat the separation above.

func SupportedEmbeddingProviders

func SupportedEmbeddingProviders() map[string]struct{}

SupportedEmbeddingProviders returns the set of supported provider names (from registry).

func TaxonomyEmbeddingModel

func TaxonomyEmbeddingModel(embeddingModel, override string) string

TaxonomyEmbeddingModel returns the embeddings.model key used for taxonomy-specific embeddings.

func TenantIDFromEventData

func TenantIDFromEventData(data any) (string, bool)

TenantIDFromEventData extracts tenant_id from known event payload shapes.

func TenantIDPointerFromEventData

func TenantIDPointerFromEventData(data any) *string

TenantIDPointerFromEventData returns a detached pointer so it can be safely stored in job args.

func ValidateEmbeddingConfig

func ValidateEmbeddingConfig(cfg EmbeddingClientConfig) error

ValidateEmbeddingConfig checks provider support and provider-specific requirements (API key, Google Cloud project/location). Use before creating a client or at startup to fail fast with a clear error.

func WebhookMatchesTenant

func WebhookMatchesTenant(webhook *models.Webhook, tenantID *string) bool

WebhookMatchesTenant reports whether a webhook may receive an event with tenantID.

Types

type BatchEmbeddingClient

type BatchEmbeddingClient interface {
	CreateEmbeddings(ctx context.Context, inputs []string) ([][]float32, error)
}

BatchEmbeddingClient is an optional provider capability for embedding multiple documents in one request. Workers use it when batching is enabled; providers that do not implement it continue to receive one request per document. Implementations must return exactly one vector per input in the same order as inputs.

type BatchingEmbeddingClient

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

BatchingEmbeddingClient coalesces concurrent document calls while leaving query calls synchronous. It is only constructed when the provider implements BatchEmbeddingClient.

func NewBatchingEmbeddingClient

func NewBatchingEmbeddingClient(
	client EmbeddingClient,
	cfg EmbeddingBatchConfig,
	metrics EmbeddingBatchMetrics,
) (*BatchingEmbeddingClient, bool)

NewBatchingEmbeddingClient wraps client when batching is enabled and supported. The bool reports whether a wrapper was created; callers should keep using client when it is false.

func (*BatchingEmbeddingClient) CreateEmbedding

func (b *BatchingEmbeddingClient) CreateEmbedding(ctx context.Context, input string) ([]float32, error)

CreateEmbedding queues a document for the next micro-batch.

func (*BatchingEmbeddingClient) CreateEmbeddingForQuery

func (b *BatchingEmbeddingClient) CreateEmbeddingForQuery(ctx context.Context, input string) ([]float32, error)

CreateEmbeddingForQuery deliberately bypasses batching; query latency and the foreground API path remain independent from background document throughput.

func (*BatchingEmbeddingClient) Shutdown

func (b *BatchingEmbeddingClient) Shutdown(ctx context.Context) error

Shutdown stops accepting new work, flushes the final partial batch, and waits for provider requests to finish. When ctx expires, in-flight provider requests are cancelled.

type CachedTenantSettings

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

CachedTenantSettings wraps a TenantSettingsReader with a per-process, size-bounded, TTL-expiring LRU cache. Feedback-record creation is high volume and the translation enqueue gate must resolve a tenant's target language per event, so caching avoids a tenant_settings read on every feedback event. Staleness is bounded by the TTL; because the worker persists the target it actually used, a changed target self-corrects on the next write. Safe for concurrent use (expirable.LRU is internally locked).

func NewCachedTenantSettings

func NewCachedTenantSettings(
	delegate TenantSettingsReader, size int, ttl time.Duration, metrics observability.CacheMetrics,
) *CachedTenantSettings

NewCachedTenantSettings wraps delegate with an LRU of at most size entries, each expiring after ttl. A non-positive size or ttl disables caching (every read hits the delegate), keeping small deployments and tests simple.

func (*CachedTenantSettings) GetSettings

func (c *CachedTenantSettings) GetSettings(
	ctx context.Context, tenantID string,
) (*models.TenantSettings, error)

GetSettings returns the tenant's settings, serving a fresh cached value when present and otherwise loading from the delegate and caching the result. Errors are never cached.

func (*CachedTenantSettings) Invalidate

func (c *CachedTenantSettings) Invalidate(tenantID string)

Invalidate evicts the tenant's cached settings so the next GetSettings reloads from the delegate. Called after a settings write so a change (e.g. a newly enabled target language) is visible immediately instead of only after TTL expiry — otherwise records created in the staleness window are skipped by the translation enqueue gate. Eviction is per-process (it refreshes the replica that handled the write); other replicas stay TTL-bounded. No-op when caching is off.

func (*CachedTenantSettings) OnSettingsChanged

func (c *CachedTenantSettings) OnSettingsChanged(_ context.Context, tenantID string, _ []string)

OnSettingsChanged implements SettingsChangeListener: any successful settings write for a tenant evicts that tenant's cached entry. Registered alongside the enrichment backfill listener so a settings change both triggers backfills and refreshes this read cache.

type EmbeddingBatchConfig

type EmbeddingBatchConfig struct {
	BatchSize   int
	MaxWait     time.Duration
	MaxInFlight int
}

EmbeddingBatchConfig controls document micro-batching. BatchSize <= 1 disables batching.

type EmbeddingBatchMetrics

type EmbeddingBatchMetrics interface {
	RecordEmbeddingBatch(ctx context.Context, inputs int64, duration time.Duration, status string)
	AddEmbeddingBatchInFlight(ctx context.Context, delta int64)
}

EmbeddingBatchMetrics is the bounded observability seam used by the batching client.

type EmbeddingClient

type EmbeddingClient interface {
	CreateEmbedding(ctx context.Context, input string) ([]float32, error)
	CreateEmbeddingForQuery(ctx context.Context, input string) ([]float32, error)
}

EmbeddingClient generates embedding vectors for text. CreateEmbedding is for embedding documents (e.g. feedback records) for storage. CreateEmbeddingForQuery is for embedding search queries; some providers (e.g. Google) use a different task type for asymmetric retrieval.

func NewEmbeddingClient

func NewEmbeddingClient(ctx context.Context, cfg EmbeddingClientConfig) (EmbeddingClient, error)

NewEmbeddingClient creates an EmbeddingClient for the given config. Validates provider-specific requirements via the registry, then calls the registry factory.

type EmbeddingClientConfig

type EmbeddingClientConfig struct {
	Provider              string
	ProviderAPIKey        string // API key for openai/google providers; not logged or serialized
	Model                 string
	BaseURL               string
	HTTPDisableKeepAlives bool
	Normalize             bool
	GoogleCloudProject    string
	GoogleCloudLocation   string
	// UsageRecorder receives each provider call's token counts and duration. nil disables it.
	UsageRecorder llm.UsageRecorder
}

EmbeddingClientConfig holds configuration for creating an embedding client.

type EmbeddingProvider

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

EmbeddingProvider implements eventPublisher by enqueueing one River job per feedback record event when the event is FeedbackRecordCreated (with non-empty value_text) or FeedbackRecordUpdated (with value_text in ChangedFields, including when value_text is now empty so the worker can clear).

func NewEmbeddingProvider

func NewEmbeddingProvider(
	inserter RiverJobInserter,
	model string,
	queueName string,
	maxAttempts int,
	docPrefix string,
	metrics observability.EmbeddingMetrics,
) *EmbeddingProvider

NewEmbeddingProvider creates a provider that enqueues feedback_embedding jobs. model is the embedding model name (e.g. text-embedding-3-small) from EMBEDDING_MODEL. docPrefix is the prefix for document text (from EmbeddingPrefixForProvider); use "" for OpenAI/Google. metrics may be nil when metrics are disabled.

func NewEmbeddingProviderForInputKind

func NewEmbeddingProviderForInputKind(
	inserter RiverJobInserter,
	model string,
	queueName string,
	maxAttempts int,
	docPrefix string,
	metrics observability.EmbeddingMetrics,
	inputKind models.EmbeddingInputKind,
) *EmbeddingProvider

NewEmbeddingProviderForInputKind creates a provider for a specific embedding input kind.

func (*EmbeddingProvider) PublishEvent

func (p *EmbeddingProvider) PublishEvent(ctx context.Context, event Event)

PublishEvent enqueues a feedback_embedding job when the event is FeedbackRecordCreated (with non-empty value_text) or FeedbackRecordUpdated (with value_text in ChangedFields). On update, the job is enqueued even when value_text is now empty so the worker can clear the embedding for text fields. API key is required for openai and google (validated at startup).

type EmbeddingReconcileArgs

type EmbeddingReconcileArgs struct{}

EmbeddingReconcileArgs is an argument-free, level-triggered taxonomy embedding sweep. Database state and deployment configuration are read when the job runs.

func (EmbeddingReconcileArgs) Kind

Kind returns the River job kind.

type EmbeddingReconcileRepository

type EmbeddingReconcileRepository interface {
	ListPendingTaxonomyEmbeddingIDs(
		ctx context.Context, model string, retryBefore time.Time, limit int,
	) ([]uuid.UUID, error)
	CountRunnableEmbeddingJobs(ctx context.Context, queue string) (int, error)
}

EmbeddingReconcileRepository is the data boundary needed by the taxonomy embedding sweep.

type EmbeddingReconcileResult

type EmbeddingReconcileResult struct {
	Found    int
	Enqueued int
	Depth    int
	AtTarget bool
}

EmbeddingReconcileResult reports one sweep's bounded queue action.

type EmbeddingReconcileService

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

EmbeddingReconcileService keeps the low-priority repair lane topped up without materializing the complete deployment backlog or crowding live embedding work.

func NewEmbeddingReconcileService

func NewEmbeddingReconcileService(
	repo EmbeddingReconcileRepository,
	model string,
	targetDepth int,
	maxAttempts int,
	retryAfter time.Duration,
) *EmbeddingReconcileService

NewEmbeddingReconcileService creates a level-triggered taxonomy embedding reconciler.

func (*EmbeddingReconcileService) SetInserter

func (s *EmbeddingReconcileService) SetInserter(inserter RiverBatchInserter)

SetInserter attaches River after its worker registry has been built.

func (*EmbeddingReconcileService) Sweep

Sweep tops the repair queue up to targetDepth with records that are still missing their current taxonomy embedding. The query excludes active jobs across all queues and terminal failures, and cools transient failures before making them eligible again.

type EmbeddingsRepository

type EmbeddingsRepository interface {
	Upsert(
		ctx context.Context, feedbackRecordID uuid.UUID, model string, embedding []float32,
		stillCurrent func(fieldLabel, valueText, valueTextTranslated *string) bool,
	) error
	DeleteByFeedbackRecordAndModel(
		ctx context.Context, feedbackRecordID uuid.UUID, model string,
		stillCurrent func(fieldLabel, valueText, valueTextTranslated *string) bool,
	) error
	ListFeedbackRecordIDsForBackfill(
		ctx context.Context, model string, afterID uuid.UUID, limit int,
	) ([]uuid.UUID, error)
	ListFeedbackRecordIDsForBackfillByInputKind(
		ctx context.Context, model string, inputKind models.EmbeddingInputKind, afterID uuid.UUID, limit int,
	) ([]uuid.UUID, error)
	ListTenantFeedbackRecordIDsForBackfillByInputKind(
		ctx context.Context, tenantID, model string, inputKind models.EmbeddingInputKind, afterID uuid.UUID, limit int,
	) ([]uuid.UUID, error)
}

EmbeddingsRepository defines the interface for embeddings table access.

type EmbeddingsRepositoryForSearch

type EmbeddingsRepositoryForSearch interface {
	GetEmbeddingAndTenantByFeedbackRecordAndModel(
		ctx context.Context, feedbackRecordID uuid.UUID, model string,
	) ([]float32, string, error)
	NearestFeedbackRecordsByEmbedding(
		ctx context.Context, model string, queryEmbedding []float32, tenantID string, limit int, excludeID *uuid.UUID, minScore float64,
	) ([]models.FeedbackRecordWithScore, bool, error)
	NearestFeedbackRecordsByEmbeddingAfterCursor(
		ctx context.Context, model string, queryEmbedding []float32, tenantID string, limit int,
		lastDistance float64, lastFeedbackRecordID uuid.UUID, excludeID *uuid.UUID, minScore float64,
	) ([]models.FeedbackRecordWithScore, bool, error)
}

EmbeddingsRepositoryForSearch provides the embedding read operations needed for semantic search. HasMore is true when there may be additional results (full page returned or full fetch limit consumed).

type EmotionsClient

type EmotionsClient interface {
	// Classify returns the emotions expressed in text. sourceLang is the record's BCP-47 language
	// and may be empty; it is passed to the model only as a hint (classification is multilingual).
	Classify(ctx context.Context, text, sourceLang string) (EmotionsResult, error)
}

EmotionsClient classifies open feedback text into zero or more emotion labels. Implementations call an LLM provider (OpenAI or Google) via structured output; the factory selects one from configuration. It mirrors the SentimentClient seam so the worker depends on the interface, not a provider.

func NewEmotionsClient

func NewEmotionsClient(ctx context.Context, cfg EmotionsClientConfig) (EmotionsClient, error)

NewEmotionsClient creates an EmotionsClient for the given config. It validates provider-specific requirements via the registry, then calls the registry factory.

type EmotionsClientConfig

type EmotionsClientConfig = EnrichmentClientConfig

EmotionsClientConfig aliases the shared classify client config (see EnrichmentClientConfig).

type EmotionsProvider

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

EmotionsProvider enqueues one emotion job per eligible feedback record event, over the shared enrichmentProvider. Eligibility is a text field with non-empty open text; re-classification is triggered by a value_text change (emotions do not depend on source language, which is only a prompt hint). Like sentiment it resolves a per-tenant setting on the enqueue path — the per-directory emotions switch (ENG-1573) — and skips tenants that have turned emotions off.

func NewEmotionsProvider

func NewEmotionsProvider(
	inserter RiverJobInserter,
	resolver TenantSettingsReader,
	queueName string,
	maxAttempts int,
	metrics observability.EmotionsMetrics,
) *EmotionsProvider

NewEmotionsProvider creates a provider that enqueues feedback_emotions jobs. resolver reads the tenant's per-directory emotions switch; metrics may be nil when disabled.

func (EmotionsProvider) PublishEvent

func (p EmotionsProvider) PublishEvent(ctx context.Context, event Event)

PublishEvent enqueues an enrichment job for an eligible create/update event.

type EmotionsResult

type EmotionsResult struct {
	Labels []models.EmotionValue
}

EmotionsResult is the parsed, validated output of one emotion classification: the distinct, known emotion labels the text expresses (possibly empty when none apply).

type EnrichmentClearMetrics

type EnrichmentClearMetrics interface {
	RecordOutputCleared(ctx context.Context, output string)
}

EnrichmentClearMetrics records enrichment outputs nulled by an edit's eager-clear, labeled by output. Optional: nil disables it; set via SetEnrichmentClearMetrics (the clear fires only on the API's UpdateFeedbackRecord path).

type EnrichmentClientConfig

type EnrichmentClientConfig struct {
	Provider            string
	ProviderAPIKey      string // API key for openai/google providers; not logged or serialized
	Model               string
	BaseURL             string
	GoogleCloudProject  string
	GoogleCloudLocation string
	// UsageRecorder receives each provider call's token counts and duration. nil disables
	// recording; the backfill commands and tests leave it unset.
	UsageRecorder llm.UsageRecorder
}

EnrichmentClientConfig is the provider/model configuration shared by the classify clients — sentiment, emotions, and translation are structurally identical, so each exposes a readable alias of this one struct (embedding keeps its own: it adds Normalize).

type EnrichmentReconcileArgs

type EnrichmentReconcileArgs struct{}

EnrichmentReconcileArgs is one sweep: find records still owed enrichments and top the reconcile queues up towards their target depth.

It carries no arguments. The sweep's inputs are the database's current state and the deployment config, both read at run time, so an argument would only be a chance for a queued job to act on a stale view.

Uniqueness is a bare marker across the in-flight states, so a tick that arrives while the previous sweep is still running collapses into it rather than running two scans concurrently. `completed` is deliberately NOT in that set: with no ByPeriod the window is unbounded, so including it would mean the first sweep is the only one that ever runs.

Deliberately an EMPTY struct: every insert encodes identically, so every insert site shares the one unique key and there is nothing a second site could get wrong. A discriminator field was tried and rejected — it created the drift channel it claimed to document, since an insert spelling the value differently would silently stop collapsing into the scheduled job.

func (EnrichmentReconcileArgs) Kind

Kind returns the River job kind.

type EnrichmentReconcileRepository

type EnrichmentReconcileRepository interface {
	ListPendingEnrichment(
		ctx context.Context, enrichment, defaultLang string, limit int,
	) ([]repository.PendingEnrichmentTarget, error)
	CountInFlightByQueue(ctx context.Context, queues []string) (map[string]int64, error)
}

EnrichmentReconcileRepository is the pending-set half: which records still owe an enrichment, and how deep each backfill queue currently is.

type EnrichmentReconcileService

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

EnrichmentReconcileService tops the backfill queues up towards a target depth with records that are still owed an enrichment.

It is the level-triggered half of enrichment coverage. The event path enqueues a job when a record changes and is what makes enrichment fast; this finds what that path missed — a record created while the provider was unconfigured, a job lost to a crash, a transient failure that used up its retries — and is what makes enrichment eventually COMPLETE.

func NewEnrichmentReconcileService

func NewEnrichmentReconcileService(params NewEnrichmentReconcileServiceParams) *EnrichmentReconcileService

NewEnrichmentReconcileService creates the reconcile service.

func (*EnrichmentReconcileService) SetInserter

func (s *EnrichmentReconcileService) SetInserter(inserter RiverBatchInserter)

SetInserter attaches the River client after construction.

The client cannot exist yet when this service is built: River needs the worker registry, the registry needs this service as its sweeper, and this service needs the client to enqueue. The same knot is untied the same way for the embedding inserter (see SetEmbeddingInserter).

A sweep with no inserter would silently enqueue nothing, so Sweep refuses rather than reporting a successful zero.

func (*EnrichmentReconcileService) Sweep

Sweep tops up every configured enrichment's backfill queue.

Topping up TO a depth rather than enqueueing AT a rate is what makes this self-regulating: each tick can only add what the workers have already drained, so river_job stays bounded whether the backlog is a thousand records or fifty million. A rate would have to be guessed against drain speed, and guessing low is indistinguishable from the reconciler not running at all.

An error on one enrichment does not abandon the others: they are independent backlogs, and a translation-specific query failure should not leave sentiment un-swept.

type EnrichmentRetryRepository

type EnrichmentRetryRepository interface {
	ClearTerminalMarkers(ctx context.Context, tenantID, enrichment string, window time.Duration) (claimed bool, cleared int64, err error)
	CooldownRemaining(ctx context.Context, tenantID, enrichment string, window time.Duration) (time.Duration, error)
}

EnrichmentRetryRepository clears terminal markers and tracks the cooldown. ClearTerminalMarkers makes the window decision itself, atomically with the delete; CooldownRemaining only reports the wait for a refused caller's response.

type EnrichmentRetryService

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

EnrichmentRetryService gives permanently-failed records another chance.

The automatic sweep will not: a terminal failure is a property of the record's own text, so re-running it costs a provider call and fails identically. This is the deliberate override for when the premise has changed — the text was edited, or the provider changed its policy.

func NewEnrichmentRetryService

func NewEnrichmentRetryService(params NewEnrichmentRetryServiceParams) *EnrichmentRetryService

NewEnrichmentRetryService creates an enrichment retry service.

func (*EnrichmentRetryService) Retry

func (s *EnrichmentRetryService) Retry(
	ctx context.Context, tenantID string, enrichments []string,
) (*models.EnrichmentRetryResponse, error)

Retry clears terminal failure markers for the named enrichments, or for all three when none are named, so the next reconcile sweep picks those records up again.

Every enrichment gets an outcome, including the ones that were refused. A caller that cannot tell "there was nothing to clear" from "you are being rate limited" will simply call again, which is the behaviour the cooldown exists to stop.

type EnrichmentSettingsListener

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

EnrichmentSettingsListener implements SettingsChangeListener by dispatching each changed setting key to the enrichment backfill it triggers. The dispatch table is built from what the deployment has enabled (e.g. translation), so a disabled enrichment registers no handler and an unknown/irrelevant key is ignored.

It deliberately does not use the webhook MessagePublisher: this is an internal, best-effort side-effect, not a customer-facing event. Enqueue errors are logged and swallowed — the global backfill command remains the guaranteed recovery path.

Adding a future setting is one map entry (e.g. a "sentiment_enabled" → sentiment backfill handler); TenantSettingsService does not change.

func NewTranslationSettingsListener

func NewTranslationSettingsListener(
	inserter RiverJobInserter, queueName string, maxAttempts int,
) *EnrichmentSettingsListener

NewTranslationSettingsListener builds a listener that, on a target_language change, enqueues a per-tenant translation backfill job (TenantTranslationBackfillArgs) via inserter. queueName/maxAttempts configure that fan-out job.

func (*EnrichmentSettingsListener) OnSettingsChanged

func (l *EnrichmentSettingsListener) OnSettingsChanged(ctx context.Context, tenantID string, changedKeys []string)

OnSettingsChanged dispatches each changed key to its enrichment backfill handler, if one is registered. Errors are logged and swallowed (the settings write already succeeded).

type EnrichmentStatusRepository

type EnrichmentStatusRepository interface {
	CountEnrichmentStatus(ctx context.Context, tenantID, defaultLang string) (repository.EnrichmentStatusCounts, error)
}

EnrichmentStatusRepository is the count surface the status service needs.

type EnrichmentStatusService

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

EnrichmentStatusService reports a tenant's enrichment progress. It resolves the per-tenant enable state from settings (via TenantSettingsReader) and overlays the deployment-level (provider/model) gate; the repo supplies the counts.

func NewEnrichmentStatusService

func NewEnrichmentStatusService(params NewEnrichmentStatusServiceParams) *EnrichmentStatusService

NewEnrichmentStatusService creates an enrichment status service.

func (*EnrichmentStatusService) GetEnrichmentStatus

func (s *EnrichmentStatusService) GetEnrichmentStatus(
	ctx context.Context, tenantID string,
) (*models.EnrichmentStatusResponse, error)

GetEnrichmentStatus returns the tenant's per-enrichment progress. tenant_id is required and validated; the query is scoped to that tenant alone.

type EnrichmentSweepSpec

type EnrichmentSweepSpec struct {
	Name string
	// MaxAttempts mirrors the live path, so a reconciled job is retried exactly as many times as
	// an event-driven one before it writes its failure marker.
	MaxAttempts int
}

EnrichmentSweepSpec is one enrichment's sweep configuration. One struct rather than parallel collections, because the parallel version had a silent failure mode: an enrichment present in the enabled list but missing from the attempts map handed River MaxAttempts 0, which it treats as its default of 25 — quietly betraying the "retried like an event-driven job" promise by a factor of eight in provider calls.

type Event

type Event struct {
	ID            uuid.UUID           // Unique event id (UUID v7, time-ordered)
	Type          datatypes.EventType // Event type enum (e.g., FeedbackRecordCreated, WebhookCreated)
	Timestamp     time.Time           // Event creation time
	Data          any                 // Event data (FeedbackRecord, Webhook, etc.)
	ChangedFields []string            // Only for updates
}

Event represents an event that can be published to message providers (webhooks, email, etc.)

type FeedbackEmbeddingArgs

type FeedbackEmbeddingArgs struct {
	FeedbackRecordID uuid.UUID `json:"feedback_record_id" river:"unique"`
	// EventID correlates this job to the domain event that enqueued it (uuid.Nil for backfill
	// jobs, which have no originating event). Not part of the dedupe key — observability only.
	EventID uuid.UUID `json:"event_id"`
	// Model is the embedding model name; stored in embeddings.model.
	Model string `json:"model" river:"unique"`
	// InputKind selects which feedback text is embedded. Empty is treated as raw for legacy queued jobs.
	InputKind models.EmbeddingInputKind `json:"input_kind,omitempty"`
	// ValueTextHash is a hash of the input (trimmed value_text, or "empty"/"backfill") for dedupe semantics.
	ValueTextHash string `json:"value_text_hash" river:"unique"`
}

FeedbackEmbeddingArgs is the job payload for generating and storing an embedding for one feedback record. Used by EmbeddingProvider and the backfill flow to enqueue, and by FeedbackEmbeddingWorker to run. Uniqueness is by (FeedbackRecordID, Model, ValueTextHash) so that edits within the uniqueness window get a new job when value_text changes; same content within 24h is deduped; one job per record+model.

func (FeedbackEmbeddingArgs) Kind

Kind returns the River job kind.

type FeedbackEmotionsArgs

type FeedbackEmotionsArgs struct {
	FeedbackRecordID uuid.UUID `json:"feedback_record_id" river:"unique"`
	// EventID correlates this job to the domain event that enqueued it (uuid.Nil for backfill
	// jobs, which have no originating event). Not part of the dedupe key — observability only.
	EventID uuid.UUID `json:"event_id"`
	// ValueTextHash is a hash of the normalized value_text, or "empty" when value_text is blank.
	ValueTextHash string `json:"value_text_hash" river:"unique"`
}

FeedbackEmotionsArgs is the job payload for classifying one feedback record's value_text into emotion labels. The river:"unique" tags define the dedupe key (FeedbackRecordID, ValueTextHash) but only take effect where the insert passes UniqueOpts — the backfill path, whose per-run ValueTextHash discriminator stops a rescued fan-out from double-enqueuing. The event-driven path deliberately inserts WITHOUT UniqueOpts (River's completed unique state would swallow legitimate re-enrichment after an edit). Like sentiment there is no target language — emotions are classified directly from the text, independent of language.

func (FeedbackEmotionsArgs) Kind

Kind returns the River job kind.

type FeedbackRecordsPurgeArgs

type FeedbackRecordsPurgeArgs struct {
	TenantID string `json:"tenant_id" river:"unique"`
}

FeedbackRecordsPurgeArgs deletes every feedback record for one tenant, everything derived from those records, and the taxonomy built on them — leaving the tenant's configuration (its webhooks and its settings) intact.

The purge runs as a job rather than on the request path because it is unbounded: feedback_records has no per-tenant cap, and the API server's write timeout would cut the response long before a large tenant finished deleting.

Uniqueness is by TenantID across the in-flight states only, spelled out explicitly at the enqueue site (see feedbackRecordsPurgeUniqueStates). A second request while a purge is running collapses into it; a purge requested after the previous one finished starts a new run.

Do not fall back to River's default state set here. It includes `completed`, and with no ByPeriod the uniqueness window is unbounded, so the first purge of a tenant would be the only one that ever runs — every later request skipped as a duplicate while still returning 202.

The worker re-reads the tenant from these args and scopes every statement by it; nothing about the purge is resolved at enqueue time.

func (FeedbackRecordsPurgeArgs) Kind

Kind returns the River job kind.

type FeedbackRecordsPurgeRepository

type FeedbackRecordsPurgeRepository interface {
	PurgeFeedbackRecordsByTenant(ctx context.Context, tenantID string) (*models.FeedbackRecordsPurgeCounts, error)
}

FeedbackRecordsPurgeRepository is the repository surface the purge needs.

type FeedbackRecordsPurgeService

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

FeedbackRecordsPurgeService owns the two halves of a tenant feedback-records purge: accepting the request (enqueue) and performing it (the worker's call).

They are split because the purge is unbounded and must not run on the request path. The API process only ever reaches Enqueue — its River client is insert-only — while hub-worker reaches Purge. Keeping both here means the tenant normalization rule is written once and applies to whichever entry point is used. The queue and retry budget are not constructor parameters: unlike the enrichment jobs they have no config knob, and both are properties of the job kind itself (see FeedbackRecordsPurgeArgs).

func NewFeedbackRecordsPurgeService

func NewFeedbackRecordsPurgeService(
	repo FeedbackRecordsPurgeRepository, inserter RiverJobInserter,
) *FeedbackRecordsPurgeService

NewFeedbackRecordsPurgeService creates the purge service. inserter may be nil in a process that only performs purges and never enqueues them (hub-worker); Enqueue then reports a clear error rather than panicking.

func (*FeedbackRecordsPurgeService) Enqueue

Enqueue accepts a purge request for a tenant and schedules the job that performs it.

Idempotent by design: the job is unique by tenant across River's in-flight states, so requesting a purge while one is queued or running collapses into it and still reports accepted. A purge requested after the previous one finished is a new run (see FeedbackRecordsPurgeArgs). One gap is deliberate: `retryable` is not in the unique set, so a request arriving while a purge sits in backoff inserts a second job — River then discards the older one when it schedules, so the tenant still ends up with exactly one purge.

func (*FeedbackRecordsPurgeService) Purge

Purge performs the purge for a tenant. Called by the worker, never from the request path.

The tenant is re-normalized here rather than trusted from the job args: enqueue-time validation is not a substitute for scoping the work at execution time, and a job's args outlive the request that created them.

type FeedbackRecordsRepository

type FeedbackRecordsRepository interface {
	Create(ctx context.Context, req *models.CreateFeedbackRecordRequest) (*models.FeedbackRecord, error)
	GetByID(ctx context.Context, id uuid.UUID) (*models.FeedbackRecord, error)
	List(ctx context.Context, filters *models.ListFeedbackRecordsFilters) ([]models.FeedbackRecord, bool, error)
	ListAfterCursor(
		ctx context.Context, filters *models.ListFeedbackRecordsFilters,
		cursorCollectedAt time.Time, cursorID uuid.UUID,
	) ([]models.FeedbackRecord, bool, error)
	Update(ctx context.Context, id uuid.UUID, req *models.UpdateFeedbackRecordRequest,
	) (updated, previous *models.FeedbackRecord, err error)
	SetTranslation(ctx context.Context, feedbackRecordID uuid.UUID, translated *string, langKey, defaultLang string,
		stillCurrent func(valueText *string) bool) error
	SetSentiment(ctx context.Context, feedbackRecordID uuid.UUID, sentiment *models.SentimentValue, score *float64,
		stillCurrent func(valueText *string) bool) error
	SetEmotions(ctx context.Context, feedbackRecordID uuid.UUID, emotions []models.EmotionValue,
		stillCurrent func(valueText *string) bool) error
	ClearEmotions(ctx context.Context, feedbackRecordID uuid.UUID,
		stillCurrent func(valueText *string) bool) error
	ListTranslationBackfillTargets(
		ctx context.Context, afterID uuid.UUID, limit int, defaultLang string,
	) ([]models.TranslationBackfillTarget, error)
	ListTranslationBackfillTargetsForTenant(
		ctx context.Context, tenantID string, afterID uuid.UUID, limit int, defaultLang string,
	) ([]models.TranslationBackfillTarget, error)
	ListSentimentBackfillTargets(ctx context.Context, afterID uuid.UUID, limit int) ([]uuid.UUID, error)
	ListEmotionsBackfillTargets(ctx context.Context, afterID uuid.UUID, limit int) ([]uuid.UUID, error)
	Count(ctx context.Context, filters *models.ListFeedbackRecordsFilters) (int, error)
	Delete(ctx context.Context, id uuid.UUID) error
	DeleteByUser(ctx context.Context, filters *models.DeleteFeedbackRecordsByUserFilters) ([]models.DeletedFeedbackRecordsByTenant, error)
}

FeedbackRecordsRepository defines the interface for feedback records data access.

type FeedbackRecordsService

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

FeedbackRecordsService handles business logic for feedback records.

func NewFeedbackRecordsService

func NewFeedbackRecordsService(
	repo FeedbackRecordsRepository,
	embeddingsRepo EmbeddingsRepository,
	embeddingModel string,
	publisher MessagePublisher,
	embeddingInserter RiverJobInserter,
	embeddingQueueName string,
	embeddingMaxAttempts int,
	translationDefaultLang string,
) *FeedbackRecordsService

NewFeedbackRecordsService creates a new feedback records service. publisher may be nil when the service is used only for backfill (BackfillEmbeddings does not use the publisher). embeddingInserter and embeddingQueueName are optional (for backfill); when nil/empty, BackfillEmbeddings returns an error. Call SetEmbeddingInserter after the River client is created to enable backfill without building the service twice. embeddingsRepo and embeddingModel are required for SetEmbedding and BackfillEmbeddings (from EMBEDDING_PROVIDER and EMBEDDING_MODEL). translationDefaultLang (TRANSLATION_DEFAULT_LANGUAGE) is the fallback target for tenants with no target_language of their own; "" disables the fallback. It governs the SetTranslation write-guard and both translation backfills, so pass it wherever translation work runs.

func (*FeedbackRecordsService) BackfillEmbeddings

func (s *FeedbackRecordsService) BackfillEmbeddings(ctx context.Context, model string) (int, error)

BackfillEmbeddings enqueues embedding jobs for the given model for all feedback records that have non-empty value_text and no embedding row for that model (existing rows are replaced by upsert when the job runs). It streams the records in keyset pages. Returns the number of jobs enqueued. Requires embeddingInserter and embeddingQueueName to be set.

func (*FeedbackRecordsService) BackfillEmbeddingsWithInputKind

func (s *FeedbackRecordsService) BackfillEmbeddingsWithInputKind(
	ctx context.Context,
	model string,
	inputKind models.EmbeddingInputKind,
) (int, error)

BackfillEmbeddingsWithInputKind enqueues embedding jobs for the given input kind.

func (*FeedbackRecordsService) BackfillEmbeddingsWithInputKindLimit

func (s *FeedbackRecordsService) BackfillEmbeddingsWithInputKindLimit(
	ctx context.Context,
	model string,
	inputKind models.EmbeddingInputKind,
	maxRecords int,
) (int, error)

BackfillEmbeddingsWithInputKindLimit enqueues global missing embedding jobs, stopping after maxRecords candidates when maxRecords is positive. Zero means unlimited.

func (*FeedbackRecordsService) BackfillEmotions

func (s *FeedbackRecordsService) BackfillEmotions(
	ctx context.Context, inserter RiverJobInserter, queueName string, maxAttempts int, runID string,
) (int, error)

BackfillEmotions enqueues an emotions job for every eligible feedback record (across all tenants) whose emotions are not yet set, streaming the targets in keyset pages. See BackfillSentiment.

func (*FeedbackRecordsService) BackfillSentiment

func (s *FeedbackRecordsService) BackfillSentiment(
	ctx context.Context, inserter RiverJobInserter, queueName string, maxAttempts int, runID string,
) (int, error)

BackfillSentiment enqueues a sentiment job for every eligible feedback record (across all tenants) whose sentiment is not yet set, streaming the targets in keyset pages. Used by the one-off backfill command. runID discriminates this run's jobs from earlier runs' (the worker re-checks the per-tenant gate, so a since-disabled tenant is skipped there rather than here). Returns the number of jobs enqueued.

func (*FeedbackRecordsService) BackfillTenantEmbeddingsWithInputKind

func (s *FeedbackRecordsService) BackfillTenantEmbeddingsWithInputKind(
	ctx context.Context,
	tenantID string,
	model string,
	inputKind models.EmbeddingInputKind,
	maxRecords int,
) (int, error)

BackfillTenantEmbeddingsWithInputKind enqueues one tenant's missing embedding jobs, stopping after maxRecords candidates when maxRecords is positive. Zero means unlimited.

func (*FeedbackRecordsService) BackfillTranslations

func (s *FeedbackRecordsService) BackfillTranslations(
	ctx context.Context, inserter RiverJobInserter, queueName string, maxAttempts int, runID string,
) (int, error)

BackfillTranslations enqueues a translation job for every feedback record (across all tenants) that needs (re)translation, streaming the targets in keyset pages. Used by the one-off global backfill command. runID discriminates this run's jobs from earlier runs' (see enqueueTranslationBackfillJobs). Returns the number of jobs enqueued.

func (*FeedbackRecordsService) BackfillTranslationsForTenant

func (s *FeedbackRecordsService) BackfillTranslationsForTenant(
	ctx context.Context, inserter RiverJobInserter, queueName string, maxAttempts int, tenantID, runID string,
) (int, error)

BackfillTranslationsForTenant enqueues a translation job for every record of a single tenant that needs (re)translation, streaming in keyset pages so a large tenant is never fully materialized. It is the bulk work behind a settings-change re-translation (TenantTranslationBackfillArgs). runID discriminates this run's jobs from earlier runs' (see enqueueTranslationBackfillJobs). Returns the number of jobs enqueued.

func (*FeedbackRecordsService) ClearEmotions

func (s *FeedbackRecordsService) ClearEmotions(
	ctx context.Context, feedbackRecordID uuid.UUID, stillCurrent func(valueText *string) bool,
) error

ClearEmotions returns a record to the "not classified" state, dropping both the labels and the completion marker. Used when the source content is gone, so the record is excluded from progress counts rather than counting as classified-with-no-emotions.

func (*FeedbackRecordsService) CountFeedbackRecords

func (s *FeedbackRecordsService) CountFeedbackRecords(
	ctx context.Context, filters *models.ListFeedbackRecordsFilters,
) (int, error)

CountFeedbackRecords returns the count of feedback records matching the given filters.

func (*FeedbackRecordsService) CreateFeedbackRecord

CreateFeedbackRecord creates a new feedback record.

func (*FeedbackRecordsService) DeleteFeedbackRecord

func (s *FeedbackRecordsService) DeleteFeedbackRecord(ctx context.Context, id uuid.UUID) error

DeleteFeedbackRecord deletes a feedback record by ID. Publishes FeedbackRecordDeleted with tenant-aware deleted IDs for webhook isolation.

func (*FeedbackRecordsService) DeleteFeedbackRecordsByUser

func (s *FeedbackRecordsService) DeleteFeedbackRecordsByUser(
	ctx context.Context, filters *models.DeleteFeedbackRecordsByUserFilters,
) (int, error)

DeleteFeedbackRecordsByUser deletes all feedback records matching user_id. When tenant_id is provided, deletion is restricted to that tenant; otherwise all user records are deleted. It publishes one tenant-aware FeedbackRecordDeleted event per tenant represented in the deleted rows.

func (*FeedbackRecordsService) GetFeedbackRecord

func (s *FeedbackRecordsService) GetFeedbackRecord(ctx context.Context, id uuid.UUID) (*models.FeedbackRecord, error)

GetFeedbackRecord retrieves a single feedback record by ID.

func (*FeedbackRecordsService) ListFeedbackRecords

ListFeedbackRecords retrieves a list of feedback records with optional filters. Uses cursor-based pagination: omit cursor for first page, use next_cursor for subsequent pages.

The provided *filters is mutated in-place to fill in the default Sort, Order and Limit when the caller omits them. This is safe because the HTTP handler (the only caller) constructs a fresh struct per request.

func (*FeedbackRecordsService) SetEmbedding

func (s *FeedbackRecordsService) SetEmbedding(
	ctx context.Context, feedbackRecordID uuid.UUID, model string, embedding []float32,
	stillCurrent func(fieldLabel, valueText, valueTextTranslated *string) bool,
) error

SetEmbedding sets or clears the embedding for a feedback record and model (internal use by embeddings worker). If embedding is nil, the row for (feedbackRecordID, model) is deleted; otherwise upserted. It does not publish an event.

func (*FeedbackRecordsService) SetEmbeddingInserter

func (s *FeedbackRecordsService) SetEmbeddingInserter(inserter RiverJobInserter)

SetEmbeddingInserter sets the River inserter for embedding jobs (e.g. after River client is created). This allows a single service instance to be used by both handlers and the embedding worker.

func (*FeedbackRecordsService) SetEmotions

func (s *FeedbackRecordsService) SetEmotions(
	ctx context.Context, feedbackRecordID uuid.UUID, emotions []models.EmotionValue,
	stillCurrent func(valueText *string) bool,
) error

SetEmotions persists or clears the emotion labels for a feedback record. It is the accessor the emotion worker uses; the write is tenant-write-locked in the repository and publishes no event (no enrichment loop). stillCurrent (optional) is the repository's content-supersession guard: it is given the record's current value_text atomically with the write, and a false return skips the write with huberrors.ErrClassificationSuperseded (nil ⇒ unconditional). Emotions are multi-label; an empty (or nil) set still records a COMPLETED classification (the labels column stays NULL, but the completion marker is stamped), which is what distinguishes "no emotion detected" from "not yet enriched" -- use ClearEmotions for the latter.

func (*FeedbackRecordsService) SetEnrichmentClearMetrics

func (s *FeedbackRecordsService) SetEnrichmentClearMetrics(m EnrichmentClearMetrics)

SetEnrichmentClearMetrics enables the eager-clear counter. Wire it on the API service instance (the eager-clear fires on UpdateFeedbackRecord); leaving it unset disables the metric.

func (*FeedbackRecordsService) SetSentiment

func (s *FeedbackRecordsService) SetSentiment(
	ctx context.Context, feedbackRecordID uuid.UUID, sentiment *models.SentimentValue, score *float64,
	stillCurrent func(valueText *string) bool,
) error

SetSentiment persists or clears the sentiment label and score for a feedback record. It is the accessor the sentiment worker uses; the write is tenant-write-locked in the repository and publishes no event (no enrichment loop). stillCurrent (optional) is the repository's content-supersession guard: it is given the record's current value_text atomically with the write, and a false return skips the write with huberrors.ErrClassificationSuperseded (nil ⇒ unconditional). Passing a nil sentiment clears both columns; a non-nil label must be valid and carry a score.

func (*FeedbackRecordsService) SetTaxonomyEmbeddingModel

func (s *FeedbackRecordsService) SetTaxonomyEmbeddingModel(model string)

SetTaxonomyEmbeddingModel sets the model key used for taxonomy-specific translated embeddings.

func (*FeedbackRecordsService) SetTranslation

func (s *FeedbackRecordsService) SetTranslation(
	ctx context.Context, feedbackRecordID uuid.UUID, translated *string, langKey string,
	stillCurrent func(valueText *string) bool,
) error

SetTranslation persists the translated value_text and the target locale key for a feedback record. It is the accessor the translation worker uses; the write is tenant-write-locked in the repository and publishes no event (no enrichment loop). stillCurrent (optional) is the repository's content-supersession guard: it is given the record's current value_text atomically with the write, and a false return skips the write with huberrors.ErrTranslationSuperseded (nil ⇒ unconditional).

func (*FeedbackRecordsService) UpdateFeedbackRecord

UpdateFeedbackRecord updates an existing feedback record.

type FeedbackSentimentArgs

type FeedbackSentimentArgs struct {
	FeedbackRecordID uuid.UUID `json:"feedback_record_id" river:"unique"`
	// EventID correlates this job to the domain event that enqueued it (uuid.Nil for backfill
	// jobs, which have no originating event). Not part of the dedupe key — observability only.
	EventID uuid.UUID `json:"event_id"`
	// ValueTextHash is a hash of the normalized value_text, or "empty" when value_text is blank.
	ValueTextHash string `json:"value_text_hash" river:"unique"`
}

FeedbackSentimentArgs is the job payload for classifying one feedback record's value_text. The river:"unique" tags define the dedupe key (FeedbackRecordID, ValueTextHash) but only take effect where the insert passes UniqueOpts — the backfill path, whose per-run ValueTextHash discriminator stops a rescued fan-out from double-enqueuing. The event-driven path deliberately inserts WITHOUT UniqueOpts (River's completed unique state would swallow legitimate re-enrichment after an edit). Unlike translation there is no target language — sentiment is classified directly from the text, independent of language.

func (FeedbackSentimentArgs) Kind

Kind returns the River job kind.

type FeedbackTranslationArgs

type FeedbackTranslationArgs struct {
	FeedbackRecordID uuid.UUID `json:"feedback_record_id" river:"unique"`
	// EventID correlates this job to the domain event that enqueued it (uuid.Nil for backfill
	// jobs, which have no originating event). Not part of the dedupe key — observability only.
	EventID uuid.UUID `json:"event_id"`
	// TargetLang is the tenant's configured target language (BCP-47) at enqueue time.
	TargetLang string `json:"target_lang" river:"unique"`
	// ValueTextHash is a hash of the inputs that determine the translation — the
	// normalized value_text and the source language — or "empty" when value_text is blank.
	ValueTextHash string `json:"value_text_hash" river:"unique"`
}

FeedbackTranslationArgs is the job payload for translating one feedback record's value_text into the tenant's target language. Uniqueness is by (FeedbackRecordID, TargetLang, ValueTextHash): a value_text edit or a target-language change yields a new job, while identical content for the same target within the window is deduped.

func (FeedbackTranslationArgs) Kind

Kind returns the River job kind.

type JobKindSpec

type JobKindSpec struct {
	// Args is a zero value of the job's argument type; Args.Kind() is the River job kind.
	Args river.JobArgs
	// Queue is the River queue the kind is inserted on by the event path.
	Queue string
	// ReconcileQueue is the second, lower-concurrency queue the same kind is inserted on when the
	// work is repaired historical data rather than event-driven, or "" for kinds with no such lane.
	//
	// A kind, not a worker, is what River registers, so the same worker serves both lanes and the
	// queue is purely a concurrency budget. That is the entire mechanism keeping a large repair
	// backlog from starving a record submitted right now: the two lanes draw from different
	// MaxWorkers.
	ReconcileQueue string
}

JobKindSpec pairs a River job kind with the queue its inserts land on.

func JobKindSpecs

func JobKindSpecs() []JobKindSpec

JobKindSpecs enumerates the Hub's River job kinds and the queues they are inserted on. It is the declaration the API's queue-depth poller derives from, and the reference the parity test in internal/workers checks hub-worker's registration against.

hub-worker's registration (workers.NewRiverWorkersAndQueues) deliberately still names its queues itself, because it attaches a per-enrichment MaxWorkers that this list has no notion of. So the two are kept consistent by assertion, not by derivation — the parity test is what fails when they drift, and it is the reason adding a kind here is safe.

The API inserts every kind below but works none of them: its River client is insert-only, so hub-worker owns registration. Adding a kind here without registering a worker for it strands its jobs, which is exactly what that test catches.

func (JobKindSpec) Kind

func (s JobKindSpec) Kind() string

Kind returns the River job kind this spec describes.

type ListPaginationMeta

type ListPaginationMeta struct {
	Limit      int
	NextCursor string
}

ListPaginationMeta holds pagination metadata for list endpoints (feedback records, webhooks).

func BuildListPaginationMeta

func BuildListPaginationMeta(
	limit int, hasMore bool, encodeLast func() (string, error),
) (ListPaginationMeta, error)

BuildListPaginationMeta builds pagination metadata for cursor-based list responses. hasMore indicates a sentinel row was fetched (limit+1 returned, trimmed to limit). encodeLast is called only when hasMore is true to produce next_cursor. Callers must ensure that when hasMore is true, the underlying list is non-empty so encodeLast can safely access the last item.

type MessagePublisher

type MessagePublisher interface {
	// PublishEvent publishes a single event with data (no changed fields)
	PublishEvent(ctx context.Context, eventType datatypes.EventType, data any)
	// PublishEventWithChangedFields publishes a single event with data and optional changed fields (for updates)
	PublishEventWithChangedFields(ctx context.Context, eventType datatypes.EventType, data any, changedFields []string)
}

MessagePublisher defines the interface for publishing events.

type MessagePublisherManager

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

MessagePublisherManager coordinates multiple message providers.

func NewMessagePublisherManager

func NewMessagePublisherManager(
	bufferSize int, perEventTimeout time.Duration, metrics observability.EventMetrics,
) *MessagePublisherManager

NewMessagePublisherManager creates a new message publisher manager. bufferSize is the event channel capacity; perEventTimeout limits how long processing one event may take. metrics may be nil when metrics are disabled.

func (*MessagePublisherManager) PublishEvent

func (m *MessagePublisherManager) PublishEvent(ctx context.Context, eventType datatypes.EventType, data any)

PublishEvent publishes an event with data to all registered providers (convenience for no changed fields).

func (*MessagePublisherManager) PublishEventWithChangedFields

func (m *MessagePublisherManager) PublishEventWithChangedFields(
	ctx context.Context, eventType datatypes.EventType, data any, changedFields []string,
)

PublishEventWithChangedFields publishes an event with data to all registered providers.

func (*MessagePublisherManager) RegisterProvider

func (m *MessagePublisherManager) RegisterProvider(provider eventPublisher)

RegisterProvider registers a message provider (webhooks, email, SMS, etc.). Must only be called during startup, before any events are published.

func (*MessagePublisherManager) Shutdown

func (m *MessagePublisherManager) Shutdown()

Shutdown stops the background worker and waits for the buffer to drain.

type NewEnrichmentReconcileServiceParams

type NewEnrichmentReconcileServiceParams struct {
	Repo        EnrichmentReconcileRepository
	Inserter    RiverBatchInserter
	DefaultLang string
	TargetDepth int
	Specs       []EnrichmentSweepSpec
}

NewEnrichmentReconcileServiceParams configures the sweep.

type NewEnrichmentRetryServiceParams

type NewEnrichmentRetryServiceParams struct {
	Repo                  EnrichmentRetryRepository
	Settings              TenantSettingsReader
	DefaultLang           string
	TranslationConfigured bool
	SentimentConfigured   bool
	EmotionsConfigured    bool
	// Cooldown overrides the default window. Zero uses EnrichmentRetryCooldown; tests set it short.
	Cooldown time.Duration
	// ReconcileEnabled is cfg.EnrichmentReconcile.Enabled — whether anything will ever act on a
	// clear. False makes Retry refuse outright rather than accept a no-op.
	ReconcileEnabled bool
	// Metrics counts each enrichment's retry outcome. nil disables them; retries still work.
	Metrics observability.EnrichmentReconcileMetrics
}

NewEnrichmentRetryServiceParams configures the retry service.

type NewEnrichmentStatusServiceParams

type NewEnrichmentStatusServiceParams struct {
	Repo                  EnrichmentStatusRepository
	Settings              TenantSettingsReader
	DefaultLang           string
	TranslationConfigured bool
	SentimentConfigured   bool
	EmotionsConfigured    bool
}

NewEnrichmentStatusServiceParams configures an EnrichmentStatusService. The *Configured flags are the deployment-level gates (provider+model set), mirroring how the enrichment providers are constructed; defaultLang is TRANSLATION_DEFAULT_LANGUAGE.

type NewTaxonomyServiceParams

type NewTaxonomyServiceParams struct {
	Repo                  TaxonomyRepository
	Starter               TaxonomyRunStarter
	EmbeddingModel        string
	MinimumEmbeddingCount int
	Metrics               observability.TaxonomyMetrics
}

NewTaxonomyServiceParams configures a TaxonomyService.

type ReconcileResult

type ReconcileResult struct {
	// Enqueued counts jobs actually inserted — after River dropped the ones already queued, so
	// this is new work rather than attempts.
	Enqueued map[string]int
	// Skipped names enrichments whose queue was already at or above the target depth. A sweep that
	// skips everything is the steady state once a backlog is draining, not a problem.
	Skipped []string
}

ReconcileResult reports what one sweep did, per enrichment.

type RiverBatchInserter

type RiverBatchInserter interface {
	InsertMany(ctx context.Context, params []river.InsertManyParams) ([]*rivertype.JobInsertResult, error)
}

RiverBatchInserter inserts many River jobs in one round trip. Shared by the webhook fan-out and the enrichment reconcile sweep, both of which enqueue in batches; the concrete River client satisfies it, and keeping the seam this small is what lets their tests run without a database-backed queue.

InsertMany, never InsertManyFast. The fast variant skips the unique-options machinery and fails the entire batch on a conflict, and callers here depend on uniqueness to avoid enqueueing work that is already in flight.

type RiverJobInserter

type RiverJobInserter interface {
	Insert(ctx context.Context, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)
}

RiverJobInserter inserts a single River job. It is the shared seam the enrichment providers (embedding, translation) and their backfill flows use to enqueue work; satisfied by the River client.

type SSRFPolicy

type SSRFPolicy struct {
	// Blacklist is a set of normalized hostnames/IPs that can never be used as webhook URLs.
	// An entry here is an explicit deny and wins over AllowedCIDRs.
	Blacklist map[string]struct{}

	// AllowedCIDRs re-permits specific private/reserved ranges that isPrivateOrReserved would
	// otherwise block, for operators whose webhook receivers legitimately live on internal
	// addresses (e.g. a tailnet in 100.64.0.0/10). It does not override Blacklist.
	AllowedCIDRs []netip.Prefix
}

SSRFPolicy holds the webhook URL restrictions applied at create/update time and again at dial time. The zero value still rejects private/reserved ranges; Blacklist and AllowedCIDRs are both optional.

func NewSSRFPolicy

func NewSSRFPolicy(blacklist map[string]struct{}, allowedCIDRs []netip.Prefix) SSRFPolicy

NewSSRFPolicy builds a policy from config primitives. Both arguments may be nil/empty; the resulting policy still rejects private/reserved ranges.

type SearchResult

type SearchResult struct {
	Results    []models.FeedbackRecordWithScore
	NextCursor string // non-empty if there may be a next page (len(Results) == requested limit)
}

SearchResult holds the results and optional next-page cursor from semantic search or similar feedback.

type SearchService

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

SearchService performs semantic search and similar-feedback lookups using embeddings.

func NewSearchService

func NewSearchService(p SearchServiceParams) *SearchService

NewSearchService creates a SearchService.

func (*SearchService) SemanticSearch

func (s *SearchService) SemanticSearch(
	ctx context.Context, query, tenantID string, limit int, minScore float64, cursor string,
) (SearchResult, error)

SemanticSearch returns feedback record IDs and similarity scores for the given query, scoped to tenantID. Requires non-empty tenantID and non-empty (after trim) query. Uses cursor-based pagination. minScore is the minimum similarity score (0..1). NextCursor is set when there may be a next page.

func (*SearchService) SimilarFeedback

func (s *SearchService) SimilarFeedback(
	ctx context.Context, feedbackRecordID uuid.UUID, limit int, minScore float64, cursor string,
) (SearchResult, error)

SimilarFeedback returns feedback record IDs and similarity scores for records similar to the given one. The tenant boundary is derived from the SOURCE RECORD (whoever owns the given UUID) — there is no caller-supplied tenant check, by design: Hub sits behind the product gateway, which owns record-level authorization (ENG-1289). If Hub ever becomes reachable without that gateway, this endpoint needs a tenant parameter checked against the source record before the search. Returns ErrEmbeddingNotFound when the record has no embedding for the current model. Uses cursor-based pagination.

type SearchServiceParams

type SearchServiceParams struct {
	EmbeddingClient EmbeddingClient
	EmbeddingsRepo  EmbeddingsRepositoryForSearch
	Model           string
	QueryCache      *lru.Cache[string, []float32]
	CacheMetrics    observability.CacheMetrics
	Logger          *slog.Logger
}

SearchServiceParams configures SearchService. QueryCache and CacheMetrics may be nil (no caching).

type SentimentClient

type SentimentClient interface {
	// Classify returns the sentiment of text. sourceLang is the record's BCP-47 language and
	// may be empty; it is passed to the model only as a hint (classification is multilingual).
	Classify(ctx context.Context, text, sourceLang string) (SentimentResult, error)
}

SentimentClient classifies the polarity of open feedback text into a sentiment label and score. Implementations call an LLM provider (OpenAI or Google) via structured output; the factory selects one from configuration. It mirrors the TranslationClient seam so the worker depends on the interface, not a provider.

func NewSentimentClient

func NewSentimentClient(ctx context.Context, cfg SentimentClientConfig) (SentimentClient, error)

NewSentimentClient creates a SentimentClient for the given config. It validates provider-specific requirements via the registry, then calls the registry factory.

type SentimentClientConfig

type SentimentClientConfig = EnrichmentClientConfig

SentimentClientConfig aliases the shared classify client config (see EnrichmentClientConfig).

type SentimentProvider

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

SentimentProvider enqueues one sentiment job per eligible feedback record event, over the shared enrichmentProvider. Eligibility is a text field with non-empty open text; re-classification is triggered by a value_text change (sentiment does not depend on source language, unlike translation). Like translation it resolves a per-tenant setting on the enqueue path — the per-directory sentiment switch (ENG-1529) — and skips tenants that have turned sentiment off.

func NewSentimentProvider

func NewSentimentProvider(
	inserter RiverJobInserter,
	resolver TenantSettingsReader,
	queueName string,
	maxAttempts int,
	metrics observability.SentimentMetrics,
) *SentimentProvider

NewSentimentProvider creates a provider that enqueues feedback_sentiment jobs. resolver reads the tenant's per-directory sentiment switch; metrics may be nil when disabled.

func (SentimentProvider) PublishEvent

func (p SentimentProvider) PublishEvent(ctx context.Context, event Event)

PublishEvent enqueues an enrichment job for an eligible create/update event.

type SentimentResult

type SentimentResult struct {
	Label models.SentimentValue
	Score float64
}

SentimentResult is the parsed, validated output of one sentiment classification: a known label and a score clamped to [models.SentimentScoreMin, SentimentScoreMax].

type SettingsChangeListener

type SettingsChangeListener interface {
	OnSettingsChanged(ctx context.Context, tenantID string, changedKeys []string)
}

SettingsChangeListener is notified after a tenant's settings are successfully written, with the setting keys the write touched. It lets enrichment side-effects (e.g. re-translation) react to a settings change without TenantSettingsService depending on any enrichment concern — the service depends only on this translation-free port.

Implementations do the reaction themselves (a fast, durable enqueue) and own their error handling: the method returns nothing because the settings write has already committed and the side-effect must never fail it.

func NewCompositeSettingsChangeListener

func NewCompositeSettingsChangeListener(listeners ...SettingsChangeListener) SettingsChangeListener

NewCompositeSettingsChangeListener combines listeners into one that notifies each in order.

type TaxonomyClient

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

TaxonomyClient calls the standalone taxonomy service.

func NewTaxonomyClient

func NewTaxonomyClient(cfg TaxonomyClientConfig, httpClient *http.Client) (*TaxonomyClient, error)

NewTaxonomyClient creates a Hub-to-taxonomy-service client.

func (*TaxonomyClient) StartRun

func (c *TaxonomyClient) StartRun(ctx context.Context, runID string) error

StartRun asks the taxonomy service to start compute for a Hub-created run.

type TaxonomyClientConfig

type TaxonomyClientConfig struct {
	ServiceURL   string
	ServiceToken string
}

TaxonomyClientConfig configures the outbound client Hub uses to call the taxonomy service.

type TaxonomyRepository

type TaxonomyRepository interface {
	ListFieldOptions(ctx context.Context, tenantID, embeddingModel string) ([]models.TaxonomyFieldOption, error)
	CountScopeInput(ctx context.Context, scope models.TaxonomyScope, embeddingModel string) (int, int, *string, error)
	CreateRunIfAvailable(ctx context.Context, params repository.CreateTaxonomyRunParams) (*models.TaxonomyRun, bool, error)
	MarkRunRunning(ctx context.Context, runID uuid.UUID, tenantID string) (*models.TaxonomyRun, error)
	MarkRunFailed(
		ctx context.Context,
		runID uuid.UUID,
		tenantID string,
		message string,
		errorCode models.TaxonomyRunFailureCode,
		metrics json.RawMessage,
	) (*models.TaxonomyRun, error)
	Heartbeat(ctx context.Context, runID uuid.UUID, tenantID string) error
	GetRunForInternalService(ctx context.Context, runID uuid.UUID) (*models.TaxonomyRun, error)
	GetRunForTenant(ctx context.Context, runID uuid.UUID, tenantID string) (*models.TaxonomyRun, error)
	GetActiveRun(ctx context.Context, scope models.TaxonomyScope) (*models.TaxonomyRun, error)
	ListRuns(ctx context.Context, filters models.ListTaxonomyRunsFilters) ([]models.TaxonomyRun, error)
	GetRunInput(
		ctx context.Context,
		runID uuid.UUID,
		tenantID string,
		embeddingModel string,
	) (*models.TaxonomyRunInputResponse, error)
	GetRunInputRecordIDs(
		ctx context.Context,
		runID uuid.UUID,
		tenantID string,
	) ([]uuid.UUID, error)
	StoreResultAndActivate(
		ctx context.Context,
		runID uuid.UUID,
		tenantID string,
		req models.TaxonomyRunResultRequest,
	) (*models.TaxonomyRun, error)
	GetTree(ctx context.Context, runID uuid.UUID, tenantID string) (*models.TaxonomyTreeResponse, error)
	RenameNode(ctx context.Context, nodeID uuid.UUID, tenantID, actorID, label string) (*models.TaxonomyNode, error)
	RemoveNode(ctx context.Context, nodeID uuid.UUID, tenantID, actorID string) (*models.TaxonomyNode, error)
	ListNodeRecords(ctx context.Context, nodeID uuid.UUID, tenantID string, limit int) ([]models.FeedbackRecord, int, error)
	CountNodeRecords(ctx context.Context, runID uuid.UUID, tenantID string) ([]models.TaxonomyNodeRecordCount, error)
}

TaxonomyRepository persists taxonomy run state and generated artifacts.

type TaxonomyRunStarter

type TaxonomyRunStarter interface {
	StartRun(ctx context.Context, runID string) error
}

TaxonomyRunStarter starts asynchronous taxonomy compute work.

type TaxonomyService

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

TaxonomyService coordinates taxonomy run lifecycle and edits.

func NewTaxonomyService

func NewTaxonomyService(params NewTaxonomyServiceParams) *TaxonomyService

NewTaxonomyService creates a taxonomy application service.

func (*TaxonomyService) CompleteRun

CompleteRun stores taxonomy output and activates the successful run.

func (*TaxonomyService) FailRun

FailRun records a taxonomy run failure.

func (*TaxonomyService) GetActiveTree

GetActiveTree returns the active taxonomy tree for a field scope.

func (*TaxonomyService) GetNodeRecordCounts

func (s *TaxonomyService) GetNodeRecordCounts(
	ctx context.Context,
	runID uuid.UUID,
	tenantID string,
) (*models.TaxonomyRecordCountsResponse, error)

GetNodeRecordCounts returns the feedback-record count for every visible node in a run, as subtree totals (a branch reports the sum of its subtopics, the root reports the run total).

func (*TaxonomyService) GetRun

func (s *TaxonomyService) GetRun(
	ctx context.Context,
	runID uuid.UUID,
	tenantID string,
) (*models.TaxonomyRun, error)

GetRun returns a taxonomy run by ID.

func (*TaxonomyService) GetRunInput

func (s *TaxonomyService) GetRunInput(
	ctx context.Context,
	runID uuid.UUID,
) (*models.TaxonomyRunInputResponse, error)

GetRunInput returns feedback text and embeddings for the taxonomy service.

func (*TaxonomyService) GetTree

func (s *TaxonomyService) GetTree(
	ctx context.Context,
	runID uuid.UUID,
	tenantID string,
) (*models.TaxonomyTreeResponse, error)

GetTree returns a taxonomy tree by run ID.

func (*TaxonomyService) Heartbeat

func (s *TaxonomyService) Heartbeat(
	ctx context.Context,
	runID uuid.UUID,
) error

Heartbeat records that a taxonomy run is still alive, keeping it out of the stuck-run reaper's reach. Resolving the run first yields its tenant and a not-found error for unknown ids.

func (*TaxonomyService) ListFieldOptions

func (s *TaxonomyService) ListFieldOptions(
	ctx context.Context,
	tenantID string,
) (*models.TaxonomyFieldsResponse, error)

ListFieldOptions returns feedback fields that can run taxonomy generation.

func (*TaxonomyService) ListNodeRecords

ListNodeRecords returns feedback records assigned to a taxonomy node.

func (*TaxonomyService) ListRuns

ListRuns returns taxonomy run history for a scoped tenant.

func (*TaxonomyService) RemoveNode

func (s *TaxonomyService) RemoveNode(
	ctx context.Context,
	nodeID uuid.UUID,
	filters models.RemoveTaxonomyNodeFilters,
) (*models.TaxonomyNode, error)

RemoveNode soft-removes a taxonomy node.

func (*TaxonomyService) RenameNode

RenameNode renames a taxonomy node.

func (*TaxonomyService) StartManualRun

StartManualRun creates and starts a manual taxonomy generation run.

type TenantDataRepository

type TenantDataRepository interface {
	DeleteByTenant(ctx context.Context, tenantID string) (*models.TenantDataDeleteCounts, error)
}

TenantDataRepository defines tenant data purge access.

type TenantDataService

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

TenantDataService handles tenant data purge business logic.

func NewTenantDataService

func NewTenantDataService(repo TenantDataRepository) *TenantDataService

NewTenantDataService creates a new tenant data service.

func (*TenantDataService) DeleteTenantData

func (s *TenantDataService) DeleteTenantData(ctx context.Context, tenantID string) (*models.TenantDataDeleteResult, error)

DeleteTenantData deletes all Hub-owned data for a tenant.

type TenantSettingsReader

type TenantSettingsReader interface {
	GetSettings(ctx context.Context, tenantID string) (*models.TenantSettings, error)
}

TenantSettingsReader is the read surface the cache wraps. *TenantSettingsService satisfies it, so the cache is a drop-in for any consumer that only needs reads (the translation enqueue gate and worker).

type TenantSettingsRepository

type TenantSettingsRepository interface {
	Get(ctx context.Context, tenantID string) (*models.TenantSettings, bool, error)
	Upsert(ctx context.Context, tenantID string, settings models.EnrichmentSettings) (*models.TenantSettings, error)
	// Patch merges set into the tenant's settings and removes removeKeys (an RFC
	// 7396 merge patch: set + delete); the two are disjoint.
	Patch(
		ctx context.Context, tenantID string, set models.EnrichmentSettings, removeKeys []string,
	) (*models.TenantSettings, error)
}

TenantSettingsRepository is the persistence surface the settings service needs.

type TenantSettingsService

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

TenantSettingsService reads and writes tenant-scoped enrichment settings. It is the accessor enrichment workflows will use to resolve a tenant's configuration.

func NewTenantSettingsService

func NewTenantSettingsService(repo TenantSettingsRepository) *TenantSettingsService

NewTenantSettingsService creates a new tenant settings service.

func (*TenantSettingsService) GetSettings

func (s *TenantSettingsService) GetSettings(ctx context.Context, tenantID string) (*models.TenantSettings, error)

GetSettings returns the tenant's enrichment settings. When the tenant has no settings yet it returns a zero-value settings bag (target language unset) rather than a not-found error: an unconfigured tenant is a valid state, and consumers decide the fallback behavior. The lookup is always scoped to the normalized tenant_id.

func (*TenantSettingsService) PatchSettings

PatchSettings applies an RFC 7396 JSON Merge Patch: a member present with a value sets that setting (validated and normalized), a member present with JSON null removes it, and an omitted member is left unchanged. It translates the typed request into the keys to set and the keys to remove, which are disjoint. The tenant_id comes from the request path and scopes the write to that tenant alone.

func (*TenantSettingsService) SetSettingsChangeListener

func (s *TenantSettingsService) SetSettingsChangeListener(listener SettingsChangeListener)

SetSettingsChangeListener registers a listener notified after a successful settings write, used to trigger enrichment side-effects (e.g. a re-translation backfill on a target_language change). Optional; mirrors the post-construction injection of SetEmbeddingInserter. Nil means no side-effects fire.

func (*TenantSettingsService) UpdateSettings

UpdateSettings validates and normalizes the requested settings, then upserts them for the tenant (full replace). The tenant_id is supplied by the caller (from the request path) and scopes the write to that tenant alone.

type TenantTranslationBackfillArgs

type TenantTranslationBackfillArgs struct {
	TenantID string `json:"tenant_id" river:"unique"`
}

TenantTranslationBackfillArgs fans out a re-translation backfill for one tenant: the worker lists the tenant's stale text records and enqueues a FeedbackTranslationArgs job for each. It is enqueued when a tenant's translation-relevant settings change (today: target_language).

Uniqueness is by TenantID across the default (in-flight) states — the enqueue site sets ByArgs without a ByPeriod — so rapid repeated settings changes collapse to a single in-flight backfill, while a change after the previous backfill has completed re-triggers. The worker resolves the tenant's current target at run time, so a coalesced backfill always targets the latest configured language.

func (TenantTranslationBackfillArgs) Kind

Kind returns the River job kind.

type TranslateRequest

type TranslateRequest struct {
	Text       string
	SourceLang string
	TargetLang string
}

TranslateRequest is the input to a single translation. SourceLang and TargetLang are BCP-47 tags: TargetLang comes from the tenant's settings; SourceLang from the feedback record's language and may be empty when the source language is unknown.

type TranslationClient

type TranslationClient interface {
	Translate(ctx context.Context, req TranslateRequest) (string, error)
}

TranslationClient translates TranslateRequest.Text from SourceLang into TargetLang and returns the translated text. Implementations call an LLM provider (OpenAI or Google); the factory selects one from configuration. It mirrors the EmbeddingClient seam so the worker depends on the interface, not a provider.

func NewTranslationClient

func NewTranslationClient(ctx context.Context, cfg TranslationClientConfig) (TranslationClient, error)

NewTranslationClient creates a TranslationClient for the given config. It validates provider-specific requirements via the registry, then calls the registry factory.

type TranslationClientConfig

type TranslationClientConfig = EnrichmentClientConfig

TranslationClientConfig aliases the shared classify client config (see EnrichmentClientConfig).

type TranslationProvider

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

TranslationProvider enqueues one translation job per eligible feedback record event, over the shared enrichmentProvider. Eligibility is a text field with non-empty open text; re-translation is triggered by a value_text OR source-language change (translation output depends on both, unlike sentiment/embedding). The per-tenant target language (settings.target_language, falling back to defaultLang) is both the gate — a record with no resolvable target is skipped — and part of the job args.

func NewTranslationProvider

func NewTranslationProvider(
	inserter RiverJobInserter,
	resolver TenantSettingsReader,
	queueName string,
	maxAttempts int,
	defaultLang string,
	metrics observability.TranslationMetrics,
) *TranslationProvider

NewTranslationProvider creates a provider that enqueues feedback_translation jobs. defaultLang is the fallback target when a tenant has none ("" keeps translation per-tenant opt-in). metrics may be nil when metrics are disabled.

func (TranslationProvider) PublishEvent

func (p TranslationProvider) PublishEvent(ctx context.Context, event Event)

PublishEvent enqueues an enrichment job for an eligible create/update event.

type WebhookDispatchArgs

type WebhookDispatchArgs struct {
	EventID       uuid.UUID `json:"event_id"                 river:"unique"`
	EventType     string    `json:"event_type"`
	Timestamp     time.Time `json:"timestamp"`
	Data          any       `json:"data"`
	ChangedFields []string  `json:"changed_fields,omitempty"`
	TenantID      *string   `json:"tenant_id,omitempty"`
	WebhookID     uuid.UUID `json:"webhook_id"               river:"unique"`
}

WebhookDispatchArgs is the job payload for one (event, webhook) delivery. Used by WebhookProvider to enqueue and by WebhookDispatchWorker to run. Only event_id and webhook_id are used for River uniqueness (river:"unique") so the hash is fast and does not include the potentially large data payload.

func (WebhookDispatchArgs) Kind

func (WebhookDispatchArgs) Kind() string

Kind returns the River job kind.

type WebhookDispatchInserter

type WebhookDispatchInserter = RiverBatchInserter

WebhookDispatchInserter inserts webhook_dispatch jobs in batch (e.g. River client).

type WebhookPayload

type WebhookPayload struct {
	ID            uuid.UUID `json:"id"`                       // Unique event id (UUID v7)
	Type          string    `json:"type"`                     // Event type as string (e.g., "feedback_record.created", "webhook.created")
	Timestamp     time.Time `json:"timestamp"`                // Event creation timestamp
	TenantID      *string   `json:"tenant_id,omitempty"`      // Tenant boundary for the event
	Data          any       `json:"data"`                     // Event data (FeedbackRecord, Webhook, etc.)
	ChangedFields []string  `json:"changed_fields,omitempty"` // Only for update events (optional)
}

WebhookPayload represents a generic webhook payload structure for all event types. The Data field can contain FeedbackRecord, Webhook, or other event data types.

func NewWebhookPayload

func NewWebhookPayload(args WebhookDispatchArgs) *WebhookPayload

NewWebhookPayload builds the public webhook payload from internal dispatch args.

type WebhookProvider

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

WebhookProvider implements eventPublisher by enqueueing one River job per (event, webhook).

func NewWebhookProvider

func NewWebhookProvider(
	inserter RiverBatchInserter, repo WebhookProviderRepository,
	maxAttempts, maxFanOut int,
	enqueueMaxRetries int, enqueueInitialBackoff, enqueueMaxBackoff time.Duration,
	metrics observability.WebhookMetrics,
) *WebhookProvider

NewWebhookProvider creates a provider that lists enabled webhooks and enqueues jobs via InsertMany. maxFanOut is the batch size for InsertMany (all matching webhooks are enqueued in batches of maxFanOut). enqueueMaxRetries, enqueueInitialBackoff, enqueueMaxBackoff configure retries when InsertMany fails (transient River/DB errors). metrics may be nil when metrics are disabled.

func (*WebhookProvider) PublishEvent

func (p *WebhookProvider) PublishEvent(ctx context.Context, event Event)

PublishEvent lists enabled webhooks for the event type and tenant, then enqueues one job per webhook. Webhooks are only eligible when the event payload has the same tenant_id.

type WebhookProviderRepository

type WebhookProviderRepository interface {
	ListEnabledForEventTypeAndTenant(ctx context.Context, eventType string, tenantID *string) ([]models.Webhook, error)
}

WebhookProviderRepository lists tenant-scoped webhooks eligible for event fan-out.

type WebhookSender

type WebhookSender interface {
	Send(ctx context.Context, webhook *models.Webhook, payload *WebhookPayload) error
}

WebhookSender sends a single webhook payload to an endpoint (Standard Webhooks: signing, headers, 410 handling).

type WebhookSenderImpl

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

WebhookSenderImpl implements WebhookSender with Standard Webhooks conformance.

func NewWebhookSenderImpl

func NewWebhookSenderImpl(
	repo WebhookSenderRepository, metrics observability.WebhookMetrics, ssrfPolicy SSRFPolicy,
	httpTimeout time.Duration, httpClient *http.Client,
) *WebhookSenderImpl

NewWebhookSenderImpl creates a sender that uses the given repo. ssrfPolicy restricts which hosts may be dialed; its zero value still rejects private/reserved ranges. It is enforced in the transport's DialContext, so it is not retained on the struct — an injected httpClient (below) is expected to carry its own dialer. httpTimeout is the HTTP client timeout; job timeout should be httpTimeout + buffer (e.g. 5s). Client does not follow redirects and validates resolved IPs at dial time (DNS rebinding protection). metrics may be nil when metrics are disabled. If httpClient is non-nil, it is used as-is (e.g. for tests that hit loopback); otherwise a secured client is built.

func (*WebhookSenderImpl) Send

func (s *WebhookSenderImpl) Send(ctx context.Context, webhook *models.Webhook, payload *WebhookPayload) error

Send signs and POSTs the payload to the webhook URL. On 410 Gone, disables the webhook and returns an error.

type WebhookSenderRepository

type WebhookSenderRepository interface {
	Update(ctx context.Context, id uuid.UUID, req *models.UpdateWebhookRequest) (*models.Webhook, error)
}

WebhookSenderRepository persists webhook state changes caused by delivery.

type WebhooksRepository

type WebhooksRepository interface {
	Create(ctx context.Context, req *models.CreateWebhookRequest) (*models.Webhook, error)
	GetByID(ctx context.Context, id uuid.UUID) (*models.Webhook, error)
	List(ctx context.Context, filters *models.ListWebhooksFilters) ([]models.Webhook, bool, error)
	ListAfterCursor(
		ctx context.Context, filters *models.ListWebhooksFilters,
		cursorCreatedAt time.Time, cursorID uuid.UUID,
	) ([]models.Webhook, bool, error)
	Count(ctx context.Context, filters *models.ListWebhooksFilters) (int64, error)
	Update(ctx context.Context, id uuid.UUID, req *models.UpdateWebhookRequest) (*models.Webhook, error)
	Delete(ctx context.Context, id uuid.UUID) (*models.DeletedWebhook, error)
}

WebhooksRepository defines the interface for webhooks data access.

type WebhooksService

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

WebhooksService handles business logic for webhooks.

func NewWebhooksService

func NewWebhooksService(
	repo WebhooksRepository, publisher MessagePublisher, maxWebhooks int, ssrfPolicy SSRFPolicy,
) *WebhooksService

NewWebhooksService creates a new webhooks service. ssrfPolicy restricts which hosts may be used as webhook URLs (SSRF mitigation); its zero value still rejects private/reserved ranges.

func (*WebhooksService) CreateWebhook

CreateWebhook creates a new webhook.

func (*WebhooksService) DeleteWebhook

func (s *WebhooksService) DeleteWebhook(ctx context.Context, id uuid.UUID) error

DeleteWebhook deletes a webhook by ID. Publishes WebhookDeleted with tenant-aware deleted IDs.

func (*WebhooksService) GetWebhook

func (s *WebhooksService) GetWebhook(ctx context.Context, id uuid.UUID) (*models.Webhook, error)

GetWebhook retrieves a single webhook by ID.

func (*WebhooksService) ListWebhooks

ListWebhooks retrieves a list of webhooks with optional filters. Uses cursor-based pagination: omit cursor for first page, use next_cursor for subsequent pages.

func (*WebhooksService) UpdateWebhook

func (s *WebhooksService) UpdateWebhook(ctx context.Context, id uuid.UUID, req *models.UpdateWebhookRequest) (*models.Webhook, error)

UpdateWebhook updates an existing webhook.

Jump to

Keyboard shortcuts

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