275 lines
8.0 KiB
Go
275 lines
8.0 KiB
Go
package dramaser
|
|
|
|
import (
|
|
"context"
|
|
cryptorand "crypto/rand"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"strconv"
|
|
"time"
|
|
|
|
redisutil "91porn-server/common/redis"
|
|
)
|
|
|
|
const (
|
|
dramaFeedScanMultiplier = 5
|
|
dramaFeedQueueTTL = 72 * time.Hour
|
|
dramaFeedLeaseTTL = 15 * time.Second
|
|
)
|
|
|
|
var (
|
|
errDramaFeedQueueMissing = errors.New("drama feed queue is missing")
|
|
errDramaFeedCursorBusy = errors.New("drama feed cursor is busy")
|
|
)
|
|
|
|
type dramaFeedReservation struct {
|
|
Version string
|
|
Length int
|
|
Offset int
|
|
Reserved int
|
|
Token string
|
|
IDs []string
|
|
}
|
|
|
|
type dramaFeedScriptClient interface {
|
|
RunScriptContext(context.Context, *redisutil.Script, []string, ...interface{}) (interface{}, error)
|
|
}
|
|
|
|
const publishDramaFeedQueueScriptSource = `
|
|
if redis.call('EXISTS', KEYS[1]) == 1 then
|
|
return 0
|
|
end
|
|
for i = 2, #ARGV do
|
|
redis.call('RPUSH', KEYS[1], ARGV[i])
|
|
end
|
|
redis.call('EXPIRE', KEYS[1], ARGV[1])
|
|
return 1
|
|
`
|
|
|
|
const reserveDramaFeedPageScriptSource = `
|
|
local length = redis.call('LLEN', KEYS[1])
|
|
if length <= 0 then
|
|
return {'QUEUE_MISSING'}
|
|
end
|
|
if redis.call('EXISTS', KEYS[3]) == 1 then
|
|
return {'BUSY'}
|
|
end
|
|
local offset = tonumber(redis.call('HGET', KEYS[2], ARGV[1]) or '0') % length
|
|
local reserved = math.min(tonumber(ARGV[2]), length)
|
|
redis.call('HSET', KEYS[3],
|
|
'token', ARGV[3], 'length', length, 'offset', offset, 'reserved', reserved)
|
|
redis.call('EXPIRE', KEYS[3], ARGV[4])
|
|
local result = {'RESERVED', ARGV[5], tostring(length), tostring(offset), tostring(reserved), ARGV[3]}
|
|
for i = 0, reserved - 1 do
|
|
table.insert(result, redis.call('LINDEX', KEYS[1], (offset + i) % length))
|
|
end
|
|
return result
|
|
`
|
|
|
|
const commitDramaFeedPageScriptSource = `
|
|
local lease = redis.call('HMGET', KEYS[3], 'token', 'length', 'offset', 'reserved')
|
|
if not lease[1] then return -1 end
|
|
if lease[1] ~= ARGV[2] or tonumber(lease[2]) ~= tonumber(ARGV[3]) or
|
|
tonumber(lease[3]) ~= tonumber(ARGV[4]) or tonumber(lease[4]) ~= tonumber(ARGV[5]) then
|
|
return -2
|
|
end
|
|
local length = redis.call('LLEN', KEYS[1])
|
|
if length <= 0 or length ~= tonumber(ARGV[3]) then return -3 end
|
|
local consumed = tonumber(ARGV[6])
|
|
if not consumed or consumed <= 0 or consumed > tonumber(ARGV[5]) then return -2 end
|
|
local currentOffset = tonumber(redis.call('HGET', KEYS[2], ARGV[1]) or '0') % length
|
|
if currentOffset ~= tonumber(ARGV[4]) then return -2 end
|
|
redis.call('HSET', KEYS[2], ARGV[1], (currentOffset + consumed) % length)
|
|
redis.call('EXPIRE', KEYS[2], ARGV[7])
|
|
redis.call('DEL', KEYS[3])
|
|
return 1
|
|
`
|
|
|
|
const abortDramaFeedPageScriptSource = `
|
|
if redis.call('HGET', KEYS[1], 'token') ~= ARGV[1] then
|
|
return 0
|
|
end
|
|
return redis.call('DEL', KEYS[1])
|
|
`
|
|
|
|
var (
|
|
publishDramaFeedQueueScript = redisutil.NewScript(publishDramaFeedQueueScriptSource)
|
|
reserveDramaFeedPageScript = redisutil.NewScript(reserveDramaFeedPageScriptSource)
|
|
commitDramaFeedPageScript = redisutil.NewScript(commitDramaFeedPageScriptSource)
|
|
abortDramaFeedPageScript = redisutil.NewScript(abortDramaFeedPageScriptSource)
|
|
)
|
|
|
|
func dramaFeedQueueKey(version string) string {
|
|
return "recommend:drama:queue:" + version
|
|
}
|
|
|
|
func dramaFeedOffsetKey(version string) string {
|
|
return "recommend:drama:offset:" + version
|
|
}
|
|
|
|
func dramaFeedLeaseKey(version string, uid uint64) string {
|
|
return "recommend:drama:reservation:" + version + ":" + strconv.FormatUint(uid, 10)
|
|
}
|
|
|
|
func reserveDramaFeedPage(
|
|
ctx context.Context,
|
|
client dramaFeedScriptClient,
|
|
uid uint64,
|
|
size int,
|
|
version string,
|
|
queueIDs []string,
|
|
) (dramaFeedReservation, error) {
|
|
if ctx == nil {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed context must not be nil")
|
|
}
|
|
if client == nil {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed Redis client must not be nil")
|
|
}
|
|
if uid == 0 || size <= 0 || version == "" || len(queueIDs) == 0 {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reservation arguments are invalid")
|
|
}
|
|
token, err := newDramaFeedToken()
|
|
if err != nil {
|
|
return dramaFeedReservation{}, err
|
|
}
|
|
reservation, err := runReserveDramaFeedPage(ctx, client, uid, size, version, token)
|
|
if !errors.Is(err, errDramaFeedQueueMissing) {
|
|
return reservation, err
|
|
}
|
|
args := make([]interface{}, 1, len(queueIDs)+1)
|
|
args[0] = int64(dramaFeedQueueTTL / time.Second)
|
|
for _, id := range queueIDs {
|
|
args = append(args, id)
|
|
}
|
|
if _, err = client.RunScriptContext(
|
|
ctx, publishDramaFeedQueueScript, []string{dramaFeedQueueKey(version)}, args...,
|
|
); err != nil {
|
|
return dramaFeedReservation{}, err
|
|
}
|
|
return runReserveDramaFeedPage(ctx, client, uid, size, version, token)
|
|
}
|
|
|
|
func runReserveDramaFeedPage(
|
|
ctx context.Context,
|
|
client dramaFeedScriptClient,
|
|
uid uint64,
|
|
size int,
|
|
version, token string,
|
|
) (dramaFeedReservation, error) {
|
|
raw, err := client.RunScriptContext(
|
|
ctx,
|
|
reserveDramaFeedPageScript,
|
|
[]string{
|
|
dramaFeedQueueKey(version),
|
|
dramaFeedOffsetKey(version),
|
|
dramaFeedLeaseKey(version, uid),
|
|
},
|
|
strconv.FormatUint(uid, 10),
|
|
size,
|
|
token,
|
|
int64(dramaFeedLeaseTTL/time.Second),
|
|
version,
|
|
)
|
|
if err != nil {
|
|
return dramaFeedReservation{}, err
|
|
}
|
|
return parseDramaFeedReservation(raw)
|
|
}
|
|
|
|
func parseDramaFeedReservation(raw interface{}) (dramaFeedReservation, error) {
|
|
items, ok := raw.([]interface{})
|
|
if !ok || len(items) == 0 {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned malformed result")
|
|
}
|
|
switch fmt.Sprint(items[0]) {
|
|
case "QUEUE_MISSING":
|
|
return dramaFeedReservation{}, errDramaFeedQueueMissing
|
|
case "BUSY":
|
|
return dramaFeedReservation{}, errDramaFeedCursorBusy
|
|
case "RESERVED":
|
|
default:
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned unknown status")
|
|
}
|
|
if len(items) < 7 {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned incomplete result")
|
|
}
|
|
result := dramaFeedReservation{Version: fmt.Sprint(items[1]), Token: fmt.Sprint(items[5])}
|
|
var err error
|
|
if result.Length, err = strconv.Atoi(fmt.Sprint(items[2])); err != nil {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned invalid length: %w", err)
|
|
}
|
|
if result.Offset, err = strconv.Atoi(fmt.Sprint(items[3])); err != nil {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned invalid offset: %w", err)
|
|
}
|
|
if result.Reserved, err = strconv.Atoi(fmt.Sprint(items[4])); err != nil {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned invalid size: %w", err)
|
|
}
|
|
if result.Version == "" || result.Token == "" || result.Length <= 0 ||
|
|
result.Offset < 0 || result.Offset >= result.Length || result.Reserved <= 0 ||
|
|
result.Reserved > result.Length || len(items)-6 != result.Reserved {
|
|
return dramaFeedReservation{}, fmt.Errorf("drama feed reserve returned invalid bounds")
|
|
}
|
|
for _, item := range items[6:] {
|
|
result.IDs = append(result.IDs, fmt.Sprint(item))
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func commitDramaFeedPage(
|
|
ctx context.Context,
|
|
client dramaFeedScriptClient,
|
|
uid uint64,
|
|
reservation dramaFeedReservation,
|
|
consumed int,
|
|
) error {
|
|
if client == nil {
|
|
return fmt.Errorf("drama feed Redis client must not be nil")
|
|
}
|
|
raw, err := client.RunScriptContext(
|
|
ctx,
|
|
commitDramaFeedPageScript,
|
|
[]string{
|
|
dramaFeedQueueKey(reservation.Version),
|
|
dramaFeedOffsetKey(reservation.Version),
|
|
dramaFeedLeaseKey(reservation.Version, uid),
|
|
},
|
|
strconv.FormatUint(uid, 10), reservation.Token, reservation.Length,
|
|
reservation.Offset, reservation.Reserved, consumed,
|
|
int64(dramaFeedQueueTTL/time.Second),
|
|
)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if fmt.Sprint(raw) != "1" {
|
|
return fmt.Errorf("drama feed commit rejected: %v", raw)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func abortDramaFeedPage(
|
|
ctx context.Context,
|
|
client dramaFeedScriptClient,
|
|
uid uint64,
|
|
reservation dramaFeedReservation,
|
|
) error {
|
|
if client == nil || reservation.Version == "" || reservation.Token == "" {
|
|
return nil
|
|
}
|
|
_, err := client.RunScriptContext(
|
|
ctx,
|
|
abortDramaFeedPageScript,
|
|
[]string{dramaFeedLeaseKey(reservation.Version, uid)},
|
|
reservation.Token,
|
|
)
|
|
return err
|
|
}
|
|
|
|
func newDramaFeedToken() (string, error) {
|
|
value := make([]byte, 16)
|
|
if _, err := cryptorand.Read(value); err != nil {
|
|
return "", fmt.Errorf("generate drama feed reservation token: %w", err)
|
|
}
|
|
return hex.EncodeToString(value), nil
|
|
}
|