package main import ( "context" "errors" "log" "net/http" "os" "os/signal" "strings" "syscall" "time" "github.com/descrybe/descrybe-v2/apps/api/internal/aiprompts" "github.com/descrybe/descrybe-v2/apps/api/internal/aiprovider" "github.com/descrybe/descrybe-v2/apps/api/internal/billing" "github.com/descrybe/descrybe-v2/apps/api/internal/config" "github.com/descrybe/descrybe-v2/apps/api/internal/db" "github.com/descrybe/descrybe-v2/apps/api/internal/feeds" "github.com/descrybe/descrybe-v2/apps/api/internal/jobs" "github.com/descrybe/descrybe-v2/apps/api/internal/logredact" "github.com/descrybe/descrybe-v2/apps/api/internal/metrics" "github.com/descrybe/descrybe-v2/apps/api/internal/platformsettings" "github.com/descrybe/descrybe-v2/apps/api/internal/processing" "github.com/descrybe/descrybe-v2/apps/api/internal/shopify" "github.com/descrybe/descrybe-v2/apps/api/internal/support" "github.com/descrybe/descrybe-v2/apps/api/internal/woocommerce" "github.com/google/uuid" "github.com/jackc/pgx/v5" ) func main() { log.SetOutput(logredact.Writer(os.Stderr)) cfg, err := config.Load() if err != nil { log.Fatalf("config: %v", err) } ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer cancel() pool, err := db.NewPool(ctx, cfg.DatabaseURL, db.PoolOptions{ MaxConns: int32(cfg.DBMaxConns), MinConns: int32(cfg.DBMinConns), MaxConnLifetime: cfg.DBMaxConnLifetime, MaxConnLifetimeJitter: cfg.DBMaxConnLifetimeJitter, MaxConnIdleTime: cfg.DBMaxConnIdleTime, HealthCheckPeriod: cfg.DBHealthCheckPeriod, StatementTimeout: cfg.DBStatementTimeout, }) if err != nil { log.Fatalf("db: %v", err) } defer pool.Close() if addr := strings.TrimSpace(os.Getenv("METRICS_ADDR")); addr != "" { go func() { mux := http.NewServeMux() mux.Handle("/metrics", metrics.Gate(cfg.IsProduction(), cfg.MetricsPublic)(metrics.Handler())) log.Printf("metrics listening on %s", addr) if err := http.ListenAndServe(addr, mux); err != nil { log.Printf("metrics server: %v", err) } }() } platSettings := platformsettings.NewService(pool, platformsettings.EnvConfig{ AppEncryptionKey: cfg.AppEncryptionKey, CredentialsEncryptionKey: cfg.CredentialsEncryptionKey, TokenSigningSecret: cfg.TokenSigningSecret, DatabaseURL: cfg.DatabaseURL, OpenAIAPIKey: cfg.OpenAIAPIKey, OpenAIBaseURL: cfg.OpenAIBaseURL, OpenAIModel: cfg.OpenAIModel, OpenAIEmbeddingAPIKey: cfg.OpenAIEmbeddingAPIKey, OpenAIEmbeddingBaseURL: cfg.OpenAIEmbeddingBaseURL, OpenAIEmbeddingModel: cfg.OpenAIEmbeddingModel, EPRELEnabled: cfg.EPRELEnabled, EPRELBaseURL: cfg.EPRELBaseURL, EPRELTimeout: cfg.EPRELTimeout, EPRELFicheLanguage: cfg.EPRELFicheLanguage, EPRELAPIKey: cfg.EPRELAPIKey, PineconeAPIKey: cfg.PineconeAPIKey, PineconeHost: cfg.PineconeHost, PineconeNamespace: cfg.PineconeNamespace, }) aiSvc := aiprovider.NewService(pool, aiprovider.EnvConfig{ AppEncryptionKey: cfg.AppEncryptionKey, CredentialsEncryptionKey: cfg.CredentialsEncryptionKey, TokenSigningSecret: cfg.TokenSigningSecret, DatabaseURL: cfg.DatabaseURL, OpenAIAPIKey: cfg.OpenAIAPIKey, OpenAIBaseURL: cfg.OpenAIBaseURL, OpenAIModel: cfg.OpenAIModel, ProcessingRPM: cfg.ProcessingRPM, ProcessingMaxRetries: cfg.ProcessingMaxRetries, }) aiSvc.Platform = platSettings pipeline := processing.NewPipeline(pool) pipeline.BatchSize = cfg.ProcessingBatchSize pipeline.AI = aiSvc pipeline.Prompts = aiprompts.NewService(pool) // OpenAI resolved per job via aiSvc → platformsettings.ResolveOpenAI (no boot snapshot). if oi, rerr := platSettings.ResolveOpenAI(ctx); rerr != nil { log.Printf("worker: platform OpenAI resolve failed: %v (AI enhance skipped until admin settings or BYOK)", rerr) } else if strings.TrimSpace(oi.APIKey) != "" { log.Printf("worker: OpenAI configured source=%s base=%s model=%s rpm=%d (resolved per job; company BYOK preferred when set)", oi.Source, oi.BaseURL, oi.Model, cfg.ProcessingRPM) } else { log.Println("worker: platform OpenAI unset - configure in admin settings or company BYOK; AI enhance skipped until then") } vector := processing.VectorCategorizer(&platformsettings.DynamicPinecone{Settings: platSettings}) if pc, rerr := platSettings.ResolvePinecone(ctx); rerr != nil { log.Printf("worker: platform Pinecone resolve failed: %v (vector categorize skipped until admin settings)", rerr) } else if pc.Configured() { if emb, eerr := platSettings.ResolveEmbedder(ctx); eerr != nil { log.Printf("worker: vectorization AI resolve failed: %v (Pinecone text-query mode)", eerr) } else if emb != nil { log.Println("worker: Pinecone vector categorizer ready (embeddings via admin AI role vectorization / env)") } else { log.Println("worker: Pinecone vector categorizer ready (text query; set ai_configs.vectorization or OPENAI_EMBEDDING_* for explicit embeddings)") } } else { log.Println("worker: Pinecone unset - configure in /admin/settings; vector categorize skipped until then") } var eprelClient processing.EPRELEnricher = &platformsettings.DynamicEPREL{Settings: platSettings} if cfg.EPRELEnabled { log.Printf("worker: EPREL enricher ready (env enabled=%v; admin settings can override)", cfg.EPRELEnabled) } else { log.Println("worker: EPREL enricher uses platform settings / env (default disabled)") } pipeline.Engine = &processing.Engine{ Vector: vector, EPREL: eprelClient, ProviderMode: processing.AIProviderInternal, } wooKeyMaterial := cfg.CredentialsEncryptionKey if wooKeyMaterial == "" { wooKeyMaterial = cfg.TokenSigningSecret } shopKeyMaterial := cfg.AppEncryptionKey if shopKeyMaterial == "" { shopKeyMaterial = wooKeyMaterial } woo := woocommerce.NewService(pool, woocommerce.DeriveKey(wooKeyMaterial, cfg.DatabaseURL)) shop := shopify.NewService(pool, shopify.DeriveKey(shopKeyMaterial, cfg.DatabaseURL)) feedSvc := &feeds.Service{Pool: pool, UploadDir: cfg.UploadDir} billingSvc := &billing.Service{Pool: pool} _ = billingSvc.EnsureDefaultCosts(ctx) supportSvc := support.NewService(pool) supportSvc.SupportAI = support.NewCompleterSupportAI(aiSvc) supportSvc.AIRateLimiter = support.NewAIRateLimiter(0, 0) jobSlots := processing.NewJobSlots(processing.DefaultProcessingWorkers) syncSlots := jobs.NewSyncSlots(jobs.DefaultSyncWorkers) log.Printf("worker started - processing workers=%d sync workers=%d poll=%s LISTEN=%s,%s (ClaimNext SKIP LOCKED) + feed sync claim + support AI auto + woo/shopify claim + scheduled enqueue + billing cycles", jobSlots.Workers, syncSlots.Workers, cfg.ProcessingPollInterval, jobs.ChannelProcessingJobs, jobs.ChannelFeedSyncJobs) if cfg.InsecureLocalProductionActive() { log.Printf("worker WARNING: ALLOW_INSECURE_LOCAL_PRODUCTION=1 with loopback WEB_ORIGIN=%s — not for public deploy", cfg.WebOrigin) } // This process cannot resume prior in-memory work. Always reclaim running // orphans before touching heartbeat so a fast restart (deploy/rb within the // heartbeat stale window) does not leave ClaimNext-invisible status=running // jobs. Single-worker architecture only — never run two processing workers. probe := jobs.ProbeWorkerReadiness(ctx, pool, jobs.DefaultHeartbeatStaleAfter) if res, err := processing.ReclaimOrphanedRunning(ctx, pool); err != nil { log.Printf("worker orphan reclaim: %v", err) } else if res.JobsRequeued > 0 || res.ProductsReset > 0 || res.SyncJobsRequeued > 0 { log.Printf("worker orphan reclaim jobs_requeued=%d products_reset=%d sync_requeued=%d (heartbeat was %s)", res.JobsRequeued, res.ProductsReset, res.SyncJobsRequeued, probe.WorkerCheck) } if err := jobs.TouchHeartbeat(ctx, pool, jobs.ProcessingWorkerID); err != nil { log.Printf("worker heartbeat bootstrap: %v", err) } wake := make(chan struct{}, 1) go func() { if err := jobs.ListenWake(ctx, pool, wake, jobs.ChannelProcessingJobs, jobs.ChannelFeedSyncJobs); err != nil && !errors.Is(err, context.Canceled) { log.Printf("worker listen wake stopped: %v", err) } }() ticker := time.NewTicker(cfg.ProcessingPollInterval) defer ticker.Stop() opsTicker := time.NewTicker(15 * time.Minute) defer opsTicker.Stop() markJobFailed := func(jobID uuid.UUID, jobErr error) { if jobErr == nil { log.Printf("job %s finished", jobID) return } if errors.Is(jobErr, context.Canceled) || errors.Is(jobErr, context.DeadlineExceeded) || ctx.Err() != nil { log.Printf("job %s interrupted: %v", jobID, processing.TruncateError(jobErr)) return } log.Printf("job %s failed: %v", jobID, processing.TruncateError(jobErr)) markCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() _, execErr := pool.Exec(markCtx, ` UPDATE processing_jobs SET status = 'failed', error = $2, completed_at = now(), updated_at = now() WHERE id = $1 AND status = 'running'`, jobID, processing.TruncateError(jobErr)) if execErr != nil { log.Printf("job %s mark failed: %v", jobID, execErr) } } runOnce := func() { if err := jobs.TouchHeartbeat(ctx, pool, jobs.ProcessingWorkerID); err != nil { log.Printf("worker heartbeat: %v", err) } if cfg.MaintenanceMode || cfg.ReadOnlyMode { return } if n, err := supportSvc.ProcessPendingAutoJobs(ctx, 3); err != nil { log.Printf("support auto AI jobs: %v", err) } else if n > 0 { log.Printf("support auto AI jobs processed=%d", n) } if _, err := jobSlots.Fill(ctx, pipeline.ClaimNext, func(jobCtx context.Context, jobID uuid.UUID) error { log.Printf("processing: claimed job=%s", jobID) return pipeline.ProcessJob(jobCtx, jobID) }, markJobFailed); err != nil && !errors.Is(err, pgx.ErrNoRows) { log.Printf("claim error: %v", err) } if n, autoErr := supportSvc.ProcessPendingAutoJobs(ctx, 3); autoErr != nil { log.Printf("support auto AI jobs: %v", autoErr) } else if n > 0 { log.Printf("support auto AI jobs processed=%d", n) } var syncJobID, syncCompanyID, syncFeedID uuid.UUID if started, err := syncSlots.TryStart(func() error { var e error syncJobID, syncCompanyID, syncFeedID, e = feedSvc.ClaimNextPendingSyncJob(ctx) return e }, func() { log.Printf("feed sync job %s feed %s company %s", syncJobID, syncFeedID, syncCompanyID) start := time.Now() syncErr := feedSvc.ProcessSyncJob(ctx, syncCompanyID, syncFeedID, syncJobID) metrics.ObserveSync("feed", syncErr, time.Since(start)) if syncErr != nil { log.Printf("feed sync job %s failed: %v", syncJobID, syncErr) } else { log.Printf("feed sync job %s done", syncJobID) } }); err != nil && !errors.Is(err, pgx.ErrNoRows) { log.Printf("feed sync claim error: %v", err) } else if started { return } var wooCompanyID uuid.UUID var wooKind string if started, err := syncSlots.TryStart(func() error { var e error wooCompanyID, wooKind, e = woo.ClaimNextPendingJob(ctx) return e }, func() { log.Printf("woocommerce %s sync company %s", wooKind, wooCompanyID) start := time.Now() switch wooKind { case "orders": summary, err := woo.SyncOrders(ctx, wooCompanyID) metrics.ObserveSync("woocommerce_orders", err, time.Since(start)) if err != nil { log.Printf("woocommerce orders sync %s failed: %v", wooCompanyID, err) return } log.Printf("woocommerce orders sync %s done pages=%d fetched=%d upserted=%d items=%d failed=%d", wooCompanyID, summary.Pages, summary.Fetched, summary.Upserted, summary.ItemsSaved, summary.Failed) case "reviews": summary, err := woo.SyncReviews(ctx, wooCompanyID) metrics.ObserveSync("woocommerce_reviews", err, time.Since(start)) if err != nil { log.Printf("woocommerce reviews sync %s failed: %v", wooCompanyID, err) return } log.Printf("woocommerce reviews sync %s done pages=%d fetched=%d upserted=%d failed=%d", wooCompanyID, summary.Pages, summary.Fetched, summary.Upserted, summary.Failed) default: summary, err := woo.SyncCompany(ctx, wooCompanyID) metrics.ObserveSync("woocommerce", err, time.Since(start)) if err != nil { log.Printf("woocommerce sync %s failed: %v", wooCompanyID, err) return } log.Printf("woocommerce sync %s done total=%d created=%d updated=%d failed=%d", wooCompanyID, summary.Total, summary.Created, summary.Updated, summary.Failed) } }); err != nil && !errors.Is(err, pgx.ErrNoRows) { log.Printf("woo claim error: %v", err) } else if started { return } var shopCompanyID uuid.UUID var shopKind string if _, err := syncSlots.TryStart(func() error { var e error shopCompanyID, shopKind, e = shop.ClaimNextPendingJob(ctx) return e }, func() { log.Printf("shopify %s sync company %s", shopKind, shopCompanyID) start := time.Now() switch shopKind { case "orders": summary, err := shop.SyncOrders(ctx, shopCompanyID) metrics.ObserveSync("shopify_orders", err, time.Since(start)) if err != nil { log.Printf("shopify orders sync %s failed: %v", shopCompanyID, err) return } log.Printf("shopify orders sync %s done pages=%d fetched=%d upserted=%d items=%d failed=%d", shopCompanyID, summary.Pages, summary.Fetched, summary.Upserted, summary.ItemsSaved, summary.Failed) default: summary, err := shop.SyncCompany(ctx, shopCompanyID) metrics.ObserveSync("shopify", err, time.Since(start)) if err != nil { log.Printf("shopify sync %s failed: %v", shopCompanyID, err) return } log.Printf("shopify sync %s done total=%d created=%d updated=%d failed=%d dry=%v", shopCompanyID, summary.Total, summary.Created, summary.Updated, summary.Failed, summary.DryRun) } }); err != nil && !errors.Is(err, pgx.ErrNoRows) { log.Printf("shopify claim error: %v", err) } } for { select { case <-ctx.Done(): log.Println("worker shutting down") jobSlots.Wait() syncSlots.Wait() return case <-opsTicker.C: if cfg.MaintenanceMode || cfg.ReadOnlyMode { continue } if res, err := billingSvc.RunDueBillingCycles(ctx); err != nil { log.Printf("billing cycles processed=%d failed=%d: %v", res.Processed, res.Failed, err) } else if res.Processed > 0 { log.Printf("billing cycles processed=%d", res.Processed) } if n, err := woo.EnqueueDueScheduled(ctx, 6*time.Hour); err != nil { log.Printf("woo schedule enqueue: %v", err) } else if n > 0 { log.Printf("woo schedule enqueued=%d", n) } if n, err := shop.EnqueueDueScheduled(ctx, 6*time.Hour); err != nil { log.Printf("shopify schedule enqueue: %v", err) } else if n > 0 { log.Printf("shopify schedule enqueued=%d", n) } if res, err := processing.CleanupStuck(ctx, pool); err != nil { log.Printf("stuck job cleanup: %v", err) } else if res.JobsMarkedFailed > 0 || res.ProductsReset > 0 || res.SyncJobsMarkedFailed > 0 { log.Printf("stuck cleanup jobs_failed=%d products_reset=%d sync_jobs_failed=%d", res.JobsMarkedFailed, res.ProductsReset, res.SyncJobsMarkedFailed) } if res, err := processing.CleanupExpired(ctx, pool); err != nil { log.Printf("expired job cleanup: %v", err) } else if res.JobsDeleted > 0 { log.Printf("retention cleanup jobs_deleted=%d", res.JobsDeleted) } if res, err := processing.CleanupExpiredSyncJobs(ctx, pool); err != nil { log.Printf("expired sync cleanup: %v", err) } else if res.SyncJobsDeleted > 0 { log.Printf("retention cleanup sync_jobs_deleted=%d", res.SyncJobsDeleted) } case <-ticker.C: runOnce() case <-wake: runOnce() } } }