fix(recommendations): make embedding backfill job timeout configurable (#229)

The embedding backfill (TriggerEmbeddings + the scheduled runEmbeddings)
ran under a hardcoded 30-minute context. That is fine for a fast hosted
embedding API, but local/self-hosted embedders (e.g. Ollama on CPU) are
far slower — on a large catalog they embed only a few thousand items
before the context deadline aborts the run with "context deadline
exceeded". The job is idempotent and resumable, so progress is not lost,
but it never finishes without repeatedly re-triggering it.

Make the per-run timeout configurable via a new
`recommendations.embeddings_job_timeout` setting (default 24h), threaded
through RecommendationsConfig -> NewWorker and applied to both the manual
trigger and the cron-scheduled run. A non-positive value falls back to
24h. Default behavior is unchanged for hosted users (a full backfill
comfortably fits in 24h); local LLM users can now complete a one-shot
backfill instead of stalling.

AI-use disclosure: implemented with assistance from Claude Code.

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
fluxis
2026-07-05 00:17:22 -04:00
committed by GitHub
co-authored by Claude Opus 4.8
parent 5b0c90b4c1
commit 2155261706
4 changed files with 19 additions and 5 deletions
+1
View File
@@ -1538,6 +1538,7 @@ func main() {
cfg.Recommendations.TasteProfilesCron,
cfg.Recommendations.CowatchCron,
cfg.Recommendations.RecommendationsCron,
cfg.Recommendations.EmbeddingsJobTimeout,
)
if err != nil {
slog.Error("failed to create recommendation worker", "error", err)
+4
View File
@@ -246,6 +246,10 @@ type RecommendationsConfig struct {
TasteDecayHalfLifeDays float64 `yaml:"-"`
DiversityLambda float64 `yaml:"-"`
CowatchCron string `yaml:"-"`
// EmbeddingsJobTimeout bounds a single embedding backfill run. A local
// CPU embedder over a large catalog needs hours, so this defaults to 24h
// (replacing a hardcoded 30m that truncated large first-run backfills).
EmbeddingsJobTimeout time.Duration `yaml:"-"`
}
// AIConfig holds the shared connection settings for Silo's AI features
+4
View File
@@ -413,6 +413,10 @@ func LoadFromDB(m map[string]string) (*Config, error) {
}())
cfg.Recommendations.EmbeddingAuthToken = stringOr(m, "recommendations.embedding_auth_token", stringOr(m, "recommendations.openai_api_key", ""))
cfg.Recommendations.EmbeddingsCron = stringOr(m, "recommendations.embeddings_cron", "0 3 * * *")
cfg.Recommendations.EmbeddingsJobTimeout, err = durationOr(m, "recommendations.embeddings_job_timeout", 24*time.Hour)
if err != nil {
return nil, err
}
cfg.Recommendations.TasteProfilesCron = stringOr(m, "recommendations.taste_profiles_cron", "0 4 * * *")
cfg.Recommendations.RecommendationsCron = stringOr(m, "recommendations.recommendations_cron", "0 5 * * *")
tasteDecayHalfLife, err := floatOr(m, "recommendations.taste_decay_half_life_days", 180)
+10 -5
View File
@@ -29,6 +29,7 @@ type Worker struct {
profileRefreshCh chan profileRefreshRequest
profileRefreshPending map[string]struct{}
cancelFunc context.CancelFunc
embeddingsJobTimeout time.Duration
}
const tasteProfileRefreshSubjectsQuery = `
@@ -45,13 +46,17 @@ const tasteProfileRefreshSubjectsQuery = `
SELECT DISTINCT user_id, profile_id FROM user_watchlist`
// NewWorker creates a new recommendation Worker.
func NewWorker(engine *Engine, embeddingsCron, tasteProfilesCron, cowatchCron, recommendationsCron string) (*Worker, error) {
func NewWorker(engine *Engine, embeddingsCron, tasteProfilesCron, cowatchCron, recommendationsCron string, embeddingsJobTimeout time.Duration) (*Worker, error) {
if embeddingsJobTimeout <= 0 {
embeddingsJobTimeout = 24 * time.Hour
}
w := &Worker{
engine: engine,
cron: cron.New(),
running: make(map[JobName]bool),
profileRefreshCh: make(chan profileRefreshRequest, 256),
profileRefreshPending: make(map[string]struct{}),
embeddingsJobTimeout: embeddingsJobTimeout,
}
if _, err := w.cron.AddFunc(embeddingsCron, w.runEmbeddings); err != nil {
@@ -123,9 +128,9 @@ func (w *Worker) TriggerEmbeddings() error {
}
go func() {
defer w.setRunning(JobEmbeddings, false)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
ctx, cancel := context.WithTimeout(context.Background(), w.embeddingsJobTimeout)
defer cancel()
slog.Info("starting embedding job (manual trigger)")
slog.Info("starting embedding job (manual trigger)", "timeout", w.embeddingsJobTimeout)
count, err := w.engine.EmbedAll(ctx)
if err != nil {
slog.Error("embedding job failed", "error", err, "embedded", count)
@@ -236,9 +241,9 @@ func (w *Worker) runEmbeddings() {
}
defer w.setRunning(JobEmbeddings, false)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
ctx, cancel := context.WithTimeout(context.Background(), w.embeddingsJobTimeout)
defer cancel()
slog.Info("starting embedding job")
slog.Info("starting embedding job", "timeout", w.embeddingsJobTimeout)
count, err := w.engine.EmbedAll(ctx)
if err != nil {
slog.Error("embedding job failed", "error", err, "embedded", count)