64 lines
1.8 KiB
Go
64 lines
1.8 KiB
Go
package redisactivity
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"strconv"
|
|
|
|
"proxy-pool/internal/domain/activitypool"
|
|
)
|
|
|
|
var _ activitypool.UpstreamDrainPolicyWriter = (*Adapter)(nil)
|
|
|
|
type upstreamDrainPolicyRecord struct {
|
|
Version int `json:"version"`
|
|
UpstreamID string `json:"upstreamId"`
|
|
Revision string `json:"revision"`
|
|
Enabled bool `json:"enabled"`
|
|
}
|
|
|
|
func (a *Adapter) ReplaceUpstreamDrainPolicies(ctx context.Context, policies []activitypool.UpstreamDrainPolicy) error {
|
|
if ctx == nil {
|
|
return activitypool.ErrInvalidMaintenance
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return err
|
|
}
|
|
if a == nil {
|
|
return activitypool.ErrInvalidMaintenance
|
|
}
|
|
records := make([]upstreamDrainPolicyRecord, 0, len(policies))
|
|
seen := make(map[string]struct{}, len(policies))
|
|
for _, policy := range policies {
|
|
if policy.UpstreamID == "" || policy.Revision == 0 {
|
|
return activitypool.ErrInvalidMaintenance
|
|
}
|
|
if _, duplicate := seen[policy.UpstreamID]; duplicate {
|
|
return activitypool.ErrInvalidMaintenance
|
|
}
|
|
seen[policy.UpstreamID] = struct{}{}
|
|
records = append(records, upstreamDrainPolicyRecord{
|
|
Version: 1, UpstreamID: policy.UpstreamID, Revision: strconv.FormatUint(policy.Revision, 10), Enabled: policy.Enabled,
|
|
})
|
|
}
|
|
payload, err := json.Marshal(records)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
result, err := runScript(ctx, a.client, upstreamDrainPolicyScript, []string{a.keys.upstreamDrainPolicies}, string(payload))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var reply upstreamDrainPolicyScriptReply
|
|
if err := decodeScriptResult(result, &reply); err != nil {
|
|
return err
|
|
}
|
|
if reply.Status == scriptInvalid {
|
|
return activitypool.ErrInvalidMaintenance
|
|
}
|
|
if reply.Status != scriptOK {
|
|
return invalidScriptReply("unexpected upstream drain policy reply")
|
|
}
|
|
return nil
|
|
}
|