停止持久化 provider 原始响应、兼容响应快照、attempt/event/outbox 重复 JSON,并由标准任务结果动态生成 Kling/Keling/Volces 兼容响应。 增加事件去重与预算、极简 callback 投递、7/30 天分批清理、安全删除条件、并发迁移索引及可实际恢复的任务域排除备份。历史清理默认关闭,待兼容协议和异步恢复在线验证后单独启用。 验证:Go 全量测试与 go vet、PostgreSQL 18 集成与实际备份恢复、迁移安全测试、bash -n、ShellCheck、Compose 配置和人工发布脚本测试均通过。
182 lines
5.6 KiB
Go
182 lines
5.6 KiB
Go
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):
|
|
}
|
|
}
|
|
}
|