feat(core): implement sync outbox mechanism and refactor provider validation

- Introduce `SyncOutboxService` and model to retry failed CP-to-Redis sync operations
- Update `SyncService` to handle sync failures by enqueuing tasks to the outbox
- Centralize provider group and API key validation logic into `ProviderGroupManager`
- Refactor API handlers to utilize the new manager and robust sync methods
- Add configuration options for sync outbox (interval, batch size, retries)
This commit is contained in:
zenfun
2025-12-25 01:24:19 +08:00
parent 44a82fa252
commit 6a16712b9d
12 changed files with 750 additions and 113 deletions

View File

@@ -6,7 +6,6 @@ import (
"github.com/ez-api/ez-api/internal/dto"
"github.com/ez-api/ez-api/internal/model"
"github.com/ez-api/foundation/provider"
"github.com/gin-gonic/gin"
)
@@ -39,11 +38,6 @@ func (h *Handler) CreateAPIKey(c *gin.Context) {
}
apiKey := strings.TrimSpace(req.APIKey)
ptype := provider.NormalizeType(group.Type)
if provider.IsGoogleFamily(ptype) && !provider.IsVertexFamily(ptype) && apiKey == "" {
c.JSON(http.StatusBadRequest, gin.H{"error": "api_key required for gemini api providers"})
return
}
status := strings.TrimSpace(req.Status)
if status == "" {
@@ -66,17 +60,21 @@ func (h *Handler) CreateAPIKey(c *gin.Context) {
tu := req.BanUntil.UTC()
key.BanUntil = &tu
}
if err := h.groupManager.ValidateAPIKey(group, key); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if err := h.db.Create(&key).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to create api key", "details": err.Error()})
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}
@@ -167,12 +165,16 @@ func (h *Handler) UpdateAPIKey(c *gin.Context) {
}
update := map[string]any{}
groupID := key.GroupID
if req.GroupID != 0 {
groupID = req.GroupID
}
var group model.ProviderGroup
if err := h.db.First(&group, groupID).Error; err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "provider group not found"})
return
}
if req.GroupID != 0 {
var group model.ProviderGroup
if err := h.db.First(&group, req.GroupID).Error; err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "provider group not found"})
return
}
update["group_id"] = req.GroupID
}
if strings.TrimSpace(req.APIKey) != "" {
@@ -197,6 +199,19 @@ func (h *Handler) UpdateAPIKey(c *gin.Context) {
if req.BanUntil.IsZero() && strings.TrimSpace(req.Status) == "active" {
update["ban_until"] = nil
}
if req.GroupID != 0 || strings.TrimSpace(req.APIKey) != "" {
nextKey := key
if v, ok := update["api_key"].(string); ok {
nextKey.APIKey = v
}
if req.GroupID != 0 {
nextKey.GroupID = req.GroupID
}
if err := h.groupManager.ValidateAPIKey(group, nextKey); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
}
if err := h.db.Model(&key).Updates(update).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to update api key", "details": err.Error()})
@@ -207,11 +222,11 @@ func (h *Handler) UpdateAPIKey(c *gin.Context) {
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}
@@ -246,11 +261,11 @@ func (h *Handler) DeleteAPIKey(c *gin.Context) {
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForAPIKey(h.db, key.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}

View File

@@ -22,6 +22,7 @@ type Handler struct {
rdb *redis.Client
logWebhook *service.LogWebhookService
logPartitioner *service.LogPartitioner
groupManager *service.ProviderGroupManager
}
func NewHandler(db *gorm.DB, logDB *gorm.DB, sync *service.SyncService, logger *service.LogWriter, rdb *redis.Client, partitioner *service.LogPartitioner) *Handler {
@@ -36,6 +37,7 @@ func NewHandler(db *gorm.DB, logDB *gorm.DB, sync *service.SyncService, logger *
rdb: rdb,
logWebhook: service.NewLogWebhookService(rdb),
logPartitioner: partitioner,
groupManager: service.NewProviderGroupManager(),
}
}

View File

@@ -7,7 +7,6 @@ import (
"github.com/ez-api/ez-api/internal/dto"
"github.com/ez-api/ez-api/internal/model"
"github.com/ez-api/foundation/provider"
"github.com/gin-gonic/gin"
"gorm.io/gorm"
)
@@ -37,64 +36,35 @@ func (h *Handler) CreateProviderGroup(c *gin.Context) {
return
}
ptype := provider.NormalizeType(req.Type)
if ptype == "" {
c.JSON(http.StatusBadRequest, gin.H{"error": "type required"})
return
}
baseURL := strings.TrimSpace(req.BaseURL)
googleLocation := provider.DefaultGoogleLocation(ptype, req.GoogleLocation)
switch ptype {
case provider.TypeOpenAI:
if baseURL == "" {
baseURL = "https://api.openai.com/v1"
}
case provider.TypeAnthropic, provider.TypeClaude:
if baseURL == "" {
baseURL = "https://api.anthropic.com"
}
case provider.TypeCompatible:
if baseURL == "" {
c.JSON(http.StatusBadRequest, gin.H{"error": "base_url required for compatible providers"})
return
}
default:
if provider.IsVertexFamily(ptype) && strings.TrimSpace(googleLocation) == "" {
googleLocation = provider.DefaultGoogleLocation(ptype, "")
}
}
status := strings.TrimSpace(req.Status)
if status == "" {
status = "active"
}
group := model.ProviderGroup{
Name: name,
Type: strings.TrimSpace(req.Type),
BaseURL: baseURL,
BaseURL: strings.TrimSpace(req.BaseURL),
GoogleProject: strings.TrimSpace(req.GoogleProject),
GoogleLocation: googleLocation,
GoogleLocation: strings.TrimSpace(req.GoogleLocation),
Models: strings.Join(req.Models, ","),
Status: status,
Status: strings.TrimSpace(req.Status),
}
if err := h.db.Create(&group).Error; err != nil {
normalized, err := h.groupManager.NormalizeGroup(group)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if err := h.db.Create(&normalized).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to create provider group", "details": err.Error()})
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForGroup(h.db, normalized.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForGroup(h.db, normalized.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}
c.JSON(http.StatusCreated, group)
c.JSON(http.StatusCreated, normalized)
}
// ListProviderGroups godoc
@@ -181,71 +151,52 @@ func (h *Handler) UpdateProviderGroup(c *gin.Context) {
return
}
nextType := strings.TrimSpace(group.Type)
if t := strings.TrimSpace(req.Type); t != "" {
nextType = t
}
nextTypeLower := provider.NormalizeType(nextType)
nextBaseURL := strings.TrimSpace(group.BaseURL)
if strings.TrimSpace(req.BaseURL) != "" {
nextBaseURL = strings.TrimSpace(req.BaseURL)
}
update := map[string]any{}
next := group
if strings.TrimSpace(req.Name) != "" {
update["name"] = strings.TrimSpace(req.Name)
next.Name = strings.TrimSpace(req.Name)
}
if strings.TrimSpace(req.Type) != "" {
update["type"] = strings.TrimSpace(req.Type)
next.Type = strings.TrimSpace(req.Type)
}
if strings.TrimSpace(req.BaseURL) != "" {
update["base_url"] = strings.TrimSpace(req.BaseURL)
next.BaseURL = strings.TrimSpace(req.BaseURL)
}
if strings.TrimSpace(req.GoogleProject) != "" {
update["google_project"] = strings.TrimSpace(req.GoogleProject)
next.GoogleProject = strings.TrimSpace(req.GoogleProject)
}
if strings.TrimSpace(req.GoogleLocation) != "" {
update["google_location"] = strings.TrimSpace(req.GoogleLocation)
} else if provider.IsVertexFamily(nextTypeLower) && strings.TrimSpace(group.GoogleLocation) == "" {
update["google_location"] = provider.DefaultGoogleLocation(nextTypeLower, "")
next.GoogleLocation = strings.TrimSpace(req.GoogleLocation)
}
if req.Models != nil {
update["models"] = strings.Join(req.Models, ",")
next.Models = strings.Join(req.Models, ",")
}
if strings.TrimSpace(req.Status) != "" {
update["status"] = strings.TrimSpace(req.Status)
next.Status = strings.TrimSpace(req.Status)
}
switch nextTypeLower {
case provider.TypeOpenAI:
if nextBaseURL == "" {
update["base_url"] = "https://api.openai.com/v1"
}
case provider.TypeAnthropic, provider.TypeClaude:
if nextBaseURL == "" {
update["base_url"] = "https://api.anthropic.com"
}
case provider.TypeCompatible:
if nextBaseURL == "" {
c.JSON(http.StatusBadRequest, gin.H{"error": "base_url required for compatible providers"})
return
}
normalized, err := h.groupManager.NormalizeGroup(next)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
group.Name = normalized.Name
group.Type = normalized.Type
group.BaseURL = normalized.BaseURL
group.GoogleProject = normalized.GoogleProject
group.GoogleLocation = normalized.GoogleLocation
group.Models = normalized.Models
group.Status = normalized.Status
if err := h.db.Model(&group).Updates(update).Error; err != nil {
if err := h.db.Save(&group).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to update provider group", "details": err.Error()})
return
}
if err := h.db.First(&group, id).Error; err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to reload provider group", "details": err.Error()})
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForGroup(h.db, group.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForGroup(h.db, group.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}
@@ -289,11 +240,11 @@ func (h *Handler) DeleteProviderGroup(c *gin.Context) {
return
}
if err := h.sync.SyncProviders(h.db); err != nil {
if err := h.sync.SyncProvidersForGroup(h.db, group.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync providers", "details": err.Error()})
return
}
if err := h.sync.SyncBindings(h.db); err != nil {
if err := h.sync.SyncBindingsForGroup(h.db, group.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to sync bindings", "details": err.Error()})
return
}