package runner import ( "bytes" "context" "encoding/json" "fmt" "hash/fnv" "io" "net/http" "time" "github.com/easyai/easyai-ai-gateway/apps/api/internal/store" "github.com/google/uuid" ) func (s *Service) StartTaskHistoryWorkers(ctx context.Context) { if s.cfg.TaskProgressCallbackEnabled && s.cfg.TaskProgressCallbackURL != "" { go s.runTaskCallbackWorker(ctx, "callback-"+uuid.NewString()) } if s.cfg.TaskCleanupEnabled { go s.runTaskHistoryCleanupWorker(ctx) } } func (s *Service) runTaskCallbackWorker(ctx context.Context, workerID string) { ticker := time.NewTicker(time.Second) defer ticker.Stop() client := &http.Client{Timeout: time.Duration(s.cfg.TaskProgressCallbackTimeoutMS) * time.Millisecond} for { s.processTaskCallbackBatch(ctx, client, workerID) select { case <-ctx.Done(): return case <-ticker.C: } } } func (s *Service) processTaskCallbackBatch(ctx context.Context, client *http.Client, workerID string) { items, err := s.store.ClaimTaskCallbacks(ctx, workerID, store.TaskCallbackBatchSize, store.TaskCallbackLockTTL) if err != nil { if ctx.Err() == nil { s.logger.Error("claim task callbacks failed", "error_category", "task_callback_claim_failed") } return } for _, item := range items { if ctx.Err() != nil { return } statusCode, deliveryErr := deliverTaskCallback(ctx, client, item, s.cfg.ServerMainInternalToken) if deliveryErr == nil && statusCode >= 200 && statusCode < 300 { if err := s.store.MarkTaskCallbackDelivered(context.WithoutCancel(ctx), item); err != nil { s.logger.Error("mark task callback delivered failed", "taskID", item.TaskID, "seq", item.Seq, "error_category", "task_callback_state_failed") } continue } retryable := deliveryErr != nil || statusCode == http.StatusRequestTimeout || statusCode == http.StatusTooManyRequests || statusCode >= 500 retry := retryable && item.Attempts < s.cfg.TaskProgressCallbackMaxAttempts message := "callback_network_error" if deliveryErr == nil { message = fmt.Sprintf("callback_http_%d", statusCode) } nextAttemptAt := time.Now().Add(taskCallbackRetryDelay(item.ID, item.Attempts)) if err := s.store.MarkTaskCallbackFailed(context.WithoutCancel(ctx), item, retry, nextAttemptAt, message); err != nil { s.logger.Error("mark task callback failed", "taskID", item.TaskID, "seq", item.Seq, "error_category", "task_callback_state_failed") continue } if !retry { s.logger.Warn("task callback moved to dead letter", "taskID", item.TaskID, "seq", item.Seq, "statusCode", statusCode, "error_category", "task_callback_failed") } } } func deliverTaskCallback(ctx context.Context, client *http.Client, item store.TaskCallbackDelivery, bearerToken string) (int, error) { body, err := json.Marshal(map[string]any{ "taskId": item.TaskID, "seq": item.Seq, "eventType": item.EventType, "status": item.TaskStatus, "createdAt": item.CreatedAt.UTC().Format(time.RFC3339Nano), }) if err != nil { return 0, err } request, err := http.NewRequestWithContext(ctx, http.MethodPost, item.CallbackURL, bytes.NewReader(body)) if err != nil { return 0, err } request.Header.Set("Content-Type", "application/json") request.Header.Set("Idempotency-Key", fmt.Sprintf("%s:%d", item.TaskID, item.Seq)) request.Header.Set("X-EasyAI-Event-Type", item.EventType) if bearerToken != "" { request.Header.Set("Authorization", "Bearer "+bearerToken) } response, err := client.Do(request) if err != nil { return 0, err } defer response.Body.Close() _, _ = io.Copy(io.Discard, io.LimitReader(response.Body, 4096)) return response.StatusCode, nil } func taskCallbackRetryDelay(deliveryID string, attempt int) time.Duration { delays := []time.Duration{ 5 * time.Second, 30 * time.Second, 2 * time.Minute, 10 * time.Minute, 30 * time.Minute, } index := attempt - 1 if index < 0 { index = 0 } if index >= len(delays) { index = len(delays) - 1 } base := delays[index] hash := fnv.New32a() _, _ = hash.Write([]byte(deliveryID)) jitter := time.Duration(hash.Sum32()%21) * base / 100 return base + jitter } func (s *Service) runTaskHistoryCleanupWorker(ctx context.Context) { interval := time.Duration(s.cfg.TaskCleanupIntervalSeconds) * time.Second ticker := time.NewTicker(interval) defer ticker.Stop() for { s.processTaskHistoryCleanup(ctx) select { case <-ctx.Done(): return case <-ticker.C: } } } func (s *Service) processTaskHistoryCleanup(ctx context.Context) { for batch := 0; batch < 10; batch++ { now := time.Now().UTC() result, err := s.store.CleanupTaskHistory( ctx, now.AddDate(0, 0, -s.cfg.TaskAnalysisRetentionDays), now.AddDate(0, 0, -s.cfg.TaskRetentionDays), s.cfg.TaskCleanupBatchSize, ) if err != nil { if ctx.Err() == nil { s.logger.Error("cleanup task history failed", "error_category", "task_history_cleanup_failed") } return } if result.Total() == 0 { return } s.logger.Info("task history cleanup batch completed", "compactedTasks", result.CompactedTasks, "compactedCheckpoints", result.CompactedCheckpoints, "compactedAttempts", result.CompactedAttempts, "compactedEvents", result.CompactedEvents, "compactedCallbacks", result.CompactedCallbacks, "retiredLegacyCallbacks", result.RetiredCallbacks, "compactedParamLogs", result.CompactedParamLogs, "compactedClonedVoices", result.CompactedClonedVoices, "deletedCallbacks", result.DeletedCallbacks, "deletedParamLogs", result.DeletedParamLogs, "deletedEvents", result.DeletedEvents, "deletedAttempts", result.DeletedAttempts, "deletedTasks", result.DeletedTasks, ) select { case <-ctx.Done(): return case <-time.After(100 * time.Millisecond): } } }