Files
easyai-ai-gateway/apps/api/internal/runner/task_history_workers.go
T
wangbo 0f0998cbcf feat(storage): 统一二进制对象存储与公开错误
新增 Aliyun OSS 与 S3 协议、通道内重试和按优先级跨通道切换,保留 server-main 兼容与环境 OSS 内存通道。

将请求及结果中的 Base64、Data URI、Buffer、multipart 和内联二进制统一对象化,生产路径不再写入本机静态目录,历史本地资源仅保留只读兼容。

引入 PublicErrorV1 并统一 API、异步查询、兼容协议和失败回调的安全错误输出,同时补充迁移、管理端、指标、OpenAPI 与本地模拟验收。

验证:go test ./... -count=1;go vet ./...;pnpm lint;pnpm test;pnpm build;pnpm openapi;tests/ci/migrations-test.sh。
2026-08-04 08:14:39 +08:00

189 lines
6.0 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/publicerror"
"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) {
payload := map[string]any{
"taskId": item.TaskID,
"seq": item.Seq,
"eventType": item.EventType,
"status": item.TaskStatus,
"createdAt": item.CreatedAt.UTC().Format(time.RFC3339Nano),
}
if item.TaskStatus == "failed" || item.TaskStatus == "cancelled" {
standard := publicerror.WithIDs(publicerror.FromFields(item.TaskErrorCode, item.TaskErrorMessage, 0, false), item.TaskRequestID, item.TaskID)
publicerror.Observe(standard)
payload["error"] = standard
}
body, err := json.Marshal(payload)
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):
}
}
}