Files
descrybe/apps/api/cmd/worker/main.go
T

395 lines
16 KiB
Go
Raw Normal View History

package main
import (
"context"
"errors"
"log"
"net/http"
"os"
"os/signal"
"strings"
"syscall"
"time"
2026-08-23 22:03:57 +02:00
"github.com/descrybe/descrybe-v2/apps/api/internal/aiaudit"
"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
2026-08-23 22:03:57 +02:00
// Full prompt/response capture for /admin/ai-calls (pruned after
// aiaudit.RetentionDays by the retention tick below).
aiSvc.WithAudit(aiaudit.NewRecorder(pool))
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)
2026-08-16 17:38:15 +02:00
if cfg.InsecureLocalProductionActive() {
log.Printf("worker WARNING: ALLOW_INSECURE_LOCAL_PRODUCTION=1 with loopback WEB_ORIGIN=%s — not for public deploy", cfg.WebOrigin)
}
2026-08-16 18:59:29 +02:00
// 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.
2026-08-16 17:38:15 +02:00
probe := jobs.ProbeWorkerReadiness(ctx, pool, jobs.DefaultHeartbeatStaleAfter)
2026-08-16 18:59:29 +02:00
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)
2026-08-16 17:38:15 +02:00
}
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 {
2026-08-16 18:59:29 +02:00
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)
}
2026-08-23 22:03:57 +02:00
if n, err := aiaudit.CleanupExpired(ctx, pool); err != nil {
log.Printf("worker: ai call log cleanup: %v", err)
} else if n > 0 {
log.Printf("worker: ai call log cleanup removed=%d older_than_days=%d", n, aiaudit.RetentionDays)
}
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()
}
}
}