58 lines
2.0 KiB
Go
58 lines
2.0 KiB
Go
package redisactivity
|
|
|
|
import (
|
|
"context"
|
|
|
|
"proxy-pool/internal/domain/activitypool"
|
|
)
|
|
|
|
var _ activitypool.HealthStore = (*Adapter)(nil)
|
|
|
|
func (a *Adapter) ApplyHealth(ctx context.Context, update activitypool.HealthUpdate) (activitypool.Entry, error) {
|
|
if ctx == nil {
|
|
return activitypool.Entry{}, activitypool.ErrInvalidHealthUpdate
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return activitypool.Entry{}, err
|
|
}
|
|
if a == nil || update.ProxyID == "" || update.CheckedAt.IsZero() || update.Latency < 0 ||
|
|
!validProxyState(string(update.NextState)) {
|
|
return activitypool.Entry{}, activitypool.ErrInvalidHealthUpdate
|
|
}
|
|
operationID, err := newOperationID()
|
|
if err != nil {
|
|
return activitypool.Entry{}, err
|
|
}
|
|
result, err := runScript(ctx, a.client, healthScript, []string{
|
|
a.keys.records, a.keys.unique, a.keys.idkeys, a.keys.expiry, a.keys.available,
|
|
a.keys.inventory, a.keys.owners, a.keys.ownerExpiry, a.keys.operation(operationID),
|
|
}, update.CheckedAt.UnixMilli(), string(update.NextState), int64(update.Latency),
|
|
a.options.CleanupLimit, operationTTLMillis(a.options.OperationTTL), update.ProxyID)
|
|
if err != nil {
|
|
return activitypool.Entry{}, err
|
|
}
|
|
var reply healthScriptReply
|
|
if err := decodeScriptResult(result, &reply); err != nil {
|
|
return activitypool.Entry{}, err
|
|
}
|
|
switch reply.Status {
|
|
case scriptNotFound:
|
|
return activitypool.Entry{}, activitypool.ErrActivityNotFound
|
|
case scriptStale:
|
|
return activitypool.Entry{}, activitypool.ErrStaleHealthUpdate
|
|
case scriptInvalid:
|
|
return activitypool.Entry{}, activitypool.ErrInvalidHealthUpdate
|
|
case scriptOK:
|
|
if reply.Record == "" {
|
|
return activitypool.Entry{}, invalidScriptReply("health reply omitted record")
|
|
}
|
|
record, err := decodeProxyRecord(reply.Record)
|
|
if err != nil {
|
|
return activitypool.Entry{}, invalidScriptReply("health reply contained an invalid record")
|
|
}
|
|
return proxyRecordEntry(record), nil
|
|
default:
|
|
return activitypool.Entry{}, invalidScriptReply("unexpected health status")
|
|
}
|
|
}
|