Files
lotteryLaravel/app/Services/Ticket/RiskPoolService.php
kang 2ea3e810a0
Some checks failed
lotterLaravel CI / test (push) Has been cancelled
lotterLaravel E2E / e2e-api (push) Has been cancelled
feat(risk): 支持按开注商隔离风险池并新增经营报表
2026-07-10 10:23:30 +08:00

539 lines
19 KiB
PHP
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
<?php
namespace App\Services\Ticket;
use App\Models\BetProvider;
use App\Models\Draw;
use App\Models\RiskPool;
use App\Lottery\ErrorCode;
use App\Models\TicketItem;
use App\Models\RiskPoolLockLog;
use Illuminate\Support\Facades\App;
use Illuminate\Support\Facades\Redis;
use App\Exceptions\TicketOperationException;
final class RiskPoolService
{
public function __construct(
private readonly PlayCatalogResolver $catalogResolver,
private readonly RiskPoolRealtimePublisher $riskRealtime,
) {}
/**
* @param list<array{number_4d:string, amount:int}> $locks
* @return list<array{number_4d:string, amount:int, warning:bool}>
*/
public function preview(int $drawId, string $providerCode, array $locks): array
{
$providerCode = $this->normalizeProviderCode($providerCode);
$rows = [];
foreach ($locks as $lock) {
$pool = $this->firstOrMakePool($drawId, $providerCode, $lock['number_4d']);
if ((int) $pool->sold_out_status === 1) {
throw new TicketOperationException('risk_sold_out', ErrorCode::RiskPoolSoldOut->value);
}
$remaining = (int) $pool->remaining_amount;
if ($remaining < (int) $lock['amount']) {
throw new TicketOperationException('risk_sold_out', ErrorCode::RiskPoolSoldOut->value);
}
$usage = (int) $pool->total_cap_amount > 0
? ((int) $pool->locked_amount + (int) $lock['amount']) / (int) $pool->total_cap_amount
: 1;
$rows[] = [
'number_4d' => $lock['number_4d'],
'amount' => (int) $lock['amount'],
'warning' => $usage >= 0.8,
];
}
return $rows;
}
/**
* @param list<array{number_4d:string, amount:int}> $locks
*/
public function acquire(int $drawId, string $providerCode, ?TicketItem $ticketItem, array $locks): int
{
$providerCode = $this->normalizeProviderCode($providerCode);
if ($this->shouldUseRedisAtomicLocks()) {
return $this->acquireWithRedisLua($drawId, $providerCode, $ticketItem, $locks);
}
return $this->acquireWithDatabaseLocks($drawId, $providerCode, $ticketItem, $locks);
}
/**
* @param list<array{number_4d:string, amount:int}> $locks
*/
private function acquireWithDatabaseLocks(int $drawId, string $providerCode, ?TicketItem $ticketItem, array $locks): int
{
$acquired = [];
$total = 0;
try {
foreach ($locks as $lock) {
$pool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $lock['number_4d'])
->lockForUpdate()
->first();
if ($pool === null) {
$pool = $this->createPool($drawId, $providerCode, $lock['number_4d']);
$pool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $lock['number_4d'])
->lockForUpdate()
->firstOrFail();
}
$amount = (int) $lock['amount'];
if ((int) $pool->sold_out_status === 1 || (int) $pool->remaining_amount < $amount) {
throw new TicketOperationException('risk_sold_out', ErrorCode::RiskPoolSoldOut->value);
}
$soldOutBefore = (int) $pool->sold_out_status;
$lockedBefore = (int) $pool->locked_amount;
$totalBefore = (int) $pool->total_cap_amount;
$pool->forceFill([
'locked_amount' => (int) $pool->locked_amount + $amount,
'remaining_amount' => (int) $pool->remaining_amount - $amount,
'sold_out_status' => ((int) $pool->remaining_amount - $amount) <= 0 ? 1 : 0,
'version' => (int) $pool->version + 1,
])->save();
$this->riskRealtime->publishAfterLock(
$drawId,
$providerCode,
$lock['number_4d'],
$soldOutBefore,
$lockedBefore,
$totalBefore,
$pool,
);
RiskPoolLockLog::query()->create([
'draw_id' => $drawId,
'provider_code' => $providerCode,
'normalized_number' => $lock['number_4d'],
'ticket_item_id' => $ticketItem?->id,
'action_type' => 'lock',
'amount' => $amount,
'source_reason' => 'ticket_place',
'created_at' => now(),
]);
$acquired[] = ['number_4d' => $lock['number_4d'], 'amount' => $amount];
$total += $amount;
}
} catch (\Throwable $e) {
$this->releaseDatabaseLocks($drawId, $providerCode, $ticketItem, $acquired, 'ticket_failed_line');
throw $e;
}
return $total;
}
/**
* Redis Lua 负责额度判断与扣减的原子性;数据库事实表随后同步,失败时释放 Redis 占用。
*
* @param list<array{number_4d:string, amount:int}> $locks
*/
private function acquireWithRedisLua(int $drawId, string $providerCode, ?TicketItem $ticketItem, array $locks): int
{
$acquired = [];
$total = 0;
try {
foreach ($locks as $lock) {
$number4d = $lock['number_4d'];
$amount = (int) $lock['amount'];
$this->acquireRedisLockForCombination($drawId, $providerCode, $number4d, $amount);
$acquired[] = ['number_4d' => $number4d, 'amount' => $amount];
$total += $amount;
$this->syncDatabaseAfterRedisAcquire($drawId, $providerCode, $ticketItem, $number4d, $amount);
}
} catch (\Throwable $e) {
$this->releaseRedisLocks($drawId, $providerCode, $acquired);
$this->releaseDatabaseLocks($drawId, $providerCode, $ticketItem, $acquired, 'ticket_failed_line');
throw $e;
}
return $total;
}
/**
* DB 事务回滚时补偿 Redis 侧已占用额度DB 锁行会随事务回滚)。
*
* @param list<array{number_4d:string, amount:int}> $locks
*/
public function compensateRedisAcquires(int $drawId, string $providerCode, array $locks): void
{
if ($locks === [] || ! $this->shouldUseRedisAtomicLocks()) {
return;
}
$this->releaseRedisLocks($drawId, $this->normalizeProviderCode($providerCode), $locks);
}
/**
* @param list<array{number_4d:string, amount:int}> $locks
*/
public function release(int $drawId, string $providerCode, ?TicketItem $ticketItem, array $locks): void
{
$providerCode = $this->normalizeProviderCode($providerCode);
if ($this->shouldUseRedisAtomicLocks()) {
$this->releaseRedisLocks($drawId, $providerCode, $locks);
}
foreach ($locks as $lock) {
$this->releaseDatabaseLocks($drawId, $providerCode, $ticketItem, [$lock], 'ticket_rollback');
}
}
private function acquireRedisLockForCombination(int $drawId, string $providerCode, string $number4d, int $amount): void
{
for ($attempt = 0; $attempt < 2; $attempt++) {
$pool = $this->firstOrMakePool($drawId, $providerCode, $number4d);
$key = $this->redisPoolKey($drawId, $providerCode, $number4d);
Redis::eval(
$this->initLua(),
1,
$key,
(int) $pool->total_cap_amount,
(int) $pool->locked_amount,
(int) $pool->version,
$this->redisPoolTtlSeconds(),
);
$result = $this->normalizeLuaResult(Redis::eval(
$this->acquireLua(),
1,
$key,
$amount,
(int) $pool->version,
$this->redisPoolTtlSeconds(),
));
if (($result['code'] ?? null) === 'OK') {
return;
}
if (($result['code'] ?? null) === 'INSUFFICIENT_CAP') {
throw new TicketOperationException('risk_sold_out', ErrorCode::RiskPoolSoldOut->value);
}
if ($attempt === 0 && in_array($result['code'] ?? '', ['VERSION_CONFLICT', 'POOL_NOT_INITIALIZED'], true)) {
$freshPool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $number4d)
->firstOrFail();
$this->syncRedisStateFromPool($freshPool);
continue;
}
$this->throwForRedisAcquireFailure($result);
}
}
/**
* @param array{code:string, remaining:int, locked:int, version:int} $result
*/
private function throwForRedisAcquireFailure(array $result): void
{
throw new TicketOperationException(
'risk_pool_unavailable',
ErrorCode::InternalError->value,
503,
['redis_code' => $result['code'] ?? 'unknown'],
);
}
public function publishManualSoldOut(Draw $draw, string $normalizedNumber, string $providerCode): void
{
$this->riskRealtime->publishManualSoldOut($draw, $normalizedNumber, $this->normalizeProviderCode($providerCode));
}
/** 后台改池或释池后,将 Redis 风控快照与 DB 对齐。 */
public function syncRedisStateFromPool(RiskPool $pool): void
{
if (! $this->shouldUseRedisAtomicLocks()) {
return;
}
$total = (int) $pool->total_cap_amount;
$locked = (int) $pool->locked_amount;
$remaining = max(0, $total - $locked);
Redis::eval(
$this->overwriteStateLua(),
1,
$this->redisPoolKey((int) $pool->draw_id, (string) ($pool->provider_code ?? BetProvider::DEFAULT_CODE), (string) $pool->normalized_number),
$total,
$locked,
$remaining,
(int) $pool->version,
$this->redisPoolTtlSeconds(),
);
}
private function shouldUseRedisAtomicLocks(): bool
{
if (App::environment('testing')) {
return false;
}
return (bool) config('lottery.risk_pool.use_redis_lua', true);
}
private function redisPoolKey(int $drawId, string $providerCode, string $number4d): string
{
$providerCode = $this->normalizeProviderCode($providerCode);
return "risk_pool:draw:{$drawId}:provider:{$providerCode}:number:{$number4d}";
}
private function redisPoolTtlSeconds(): int
{
return max(60, (int) config('lottery.risk_pool.redis_ttl_seconds', 86400));
}
private function initLua(): string
{
return <<<'LUA'
if redis.call('EXISTS', KEYS[1]) == 0 then
redis.call('HMSET', KEYS[1], 'total', ARGV[1], 'locked', ARGV[2], 'remaining', ARGV[1] - ARGV[2], 'version', ARGV[3])
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[4]))
end
return 1
LUA;
}
private function overwriteStateLua(): string
{
return <<<'LUA'
redis.call('HMSET', KEYS[1], 'total', ARGV[1], 'locked', ARGV[2], 'remaining', ARGV[3], 'version', ARGV[4])
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[5]))
return 1
LUA;
}
private function acquireLua(): string
{
return <<<'LUA'
local amount = tonumber(ARGV[1])
local expectedVersion = tonumber(ARGV[2])
if amount == nil or amount <= 0 then
return {'INVALID_ARGUMENT', 0, 0, 0}
end
if redis.call('EXISTS', KEYS[1]) == 0 then
return {'POOL_NOT_INITIALIZED', 0, 0, 0}
end
local version = tonumber(redis.call('HGET', KEYS[1], 'version') or '0')
if expectedVersion ~= nil and version ~= expectedVersion then
return {'VERSION_CONFLICT', tonumber(redis.call('HGET', KEYS[1], 'remaining') or '0'), tonumber(redis.call('HGET', KEYS[1], 'locked') or '0'), version}
end
local remaining = tonumber(redis.call('HGET', KEYS[1], 'remaining') or '0')
if remaining < amount then
return {'INSUFFICIENT_CAP', remaining, tonumber(redis.call('HGET', KEYS[1], 'locked') or '0'), version}
end
local locked = redis.call('HINCRBY', KEYS[1], 'locked', amount)
remaining = redis.call('HINCRBY', KEYS[1], 'remaining', -amount)
version = redis.call('HINCRBY', KEYS[1], 'version', 1)
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[3]))
return {'OK', remaining, locked, version}
LUA;
}
private function releaseLua(): string
{
return <<<'LUA'
local amount = tonumber(ARGV[1])
local locked = tonumber(redis.call('HGET', KEYS[1], 'locked') or '0')
local releaseAmount = amount
if locked < releaseAmount then
releaseAmount = locked
end
redis.call('HINCRBY', KEYS[1], 'locked', -releaseAmount)
redis.call('HINCRBY', KEYS[1], 'remaining', releaseAmount)
redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2]))
return releaseAmount
LUA;
}
private function syncDatabaseAfterRedisAcquire(int $drawId, string $providerCode, ?TicketItem $ticketItem, string $number4d, int $amount): void
{
$pool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $number4d)
->lockForUpdate()
->firstOrFail();
if ((int) $pool->sold_out_status === 1) {
throw new TicketOperationException('risk_sold_out', ErrorCode::RiskPoolSoldOut->value);
}
$soldOutBefore = (int) $pool->sold_out_status;
$lockedBefore = (int) $pool->locked_amount;
$totalBefore = (int) $pool->total_cap_amount;
$pool->forceFill([
'locked_amount' => (int) $pool->locked_amount + $amount,
'remaining_amount' => (int) $pool->remaining_amount - $amount,
'sold_out_status' => ((int) $pool->remaining_amount - $amount) <= 0 ? 1 : 0,
'version' => (int) $pool->version + 1,
])->save();
$this->riskRealtime->publishAfterLock(
$drawId,
$providerCode,
$number4d,
$soldOutBefore,
$lockedBefore,
$totalBefore,
$pool,
);
RiskPoolLockLog::query()->create([
'draw_id' => $drawId,
'provider_code' => $providerCode,
'normalized_number' => $number4d,
'ticket_item_id' => $ticketItem?->id,
'action_type' => 'lock',
'amount' => $amount,
'source_reason' => 'ticket_place',
'created_at' => now(),
]);
}
/**
* @param list<array{number_4d:string, amount:int}> $locks
*/
private function releaseRedisLocks(int $drawId, string $providerCode, array $locks): void
{
foreach ($locks as $lock) {
Redis::eval(
$this->releaseLua(),
1,
$this->redisPoolKey($drawId, $providerCode, $lock['number_4d']),
(int) $lock['amount'],
$this->redisPoolTtlSeconds(),
);
}
}
/**
* @return array{code:string, remaining:int, locked:int, version:int}
*/
private function normalizeLuaResult(mixed $result): array
{
if (! is_array($result)) {
return [
'code' => (int) $result === 1 ? 'OK' : 'INSUFFICIENT_CAP',
'remaining' => 0,
'locked' => 0,
'version' => 0,
];
}
return [
'code' => (string) ($result[0] ?? 'INSUFFICIENT_CAP'),
'remaining' => (int) ($result[1] ?? 0),
'locked' => (int) ($result[2] ?? 0),
'version' => (int) ($result[3] ?? 0),
];
}
/**
* @param list<array{number_4d:string, amount:int}> $locks
*/
private function releaseDatabaseLocks(int $drawId, string $providerCode, ?TicketItem $ticketItem, array $locks, string $sourceReason): void
{
foreach ($locks as $lock) {
$pool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $lock['number_4d'])
->lockForUpdate()
->first();
if ($pool === null) {
continue;
}
$amount = min((int) $lock['amount'], (int) $pool->locked_amount);
$pool->forceFill([
'locked_amount' => (int) $pool->locked_amount - $amount,
'remaining_amount' => (int) $pool->remaining_amount + $amount,
'sold_out_status' => 0,
'version' => (int) $pool->version + 1,
])->save();
RiskPoolLockLog::query()->create([
'draw_id' => $drawId,
'provider_code' => $providerCode,
'normalized_number' => $lock['number_4d'],
'ticket_item_id' => $ticketItem?->id,
'action_type' => 'release',
'amount' => $amount,
'source_reason' => $sourceReason,
'created_at' => now(),
]);
}
}
private function firstOrMakePool(int $drawId, string $providerCode, string $number4d): RiskPool
{
$pool = RiskPool::query()
->where('draw_id', $drawId)
->where('provider_code', $providerCode)
->where('normalized_number', $number4d)
->first();
if ($pool !== null) {
return $pool;
}
return $this->createPool($drawId, $providerCode, $number4d);
}
private function createPool(int $drawId, string $providerCode, string $number4d): RiskPool
{
$cap = $this->catalogResolver->resolveCapAmount($drawId, $number4d, $providerCode);
return RiskPool::query()->firstOrCreate([
'draw_id' => $drawId,
'provider_code' => $providerCode,
'normalized_number' => $number4d,
], [
'normalized_number' => $number4d,
'total_cap_amount' => $cap,
'locked_amount' => 0,
'remaining_amount' => $cap,
'sold_out_status' => 0,
'version' => 0,
]);
}
private function normalizeProviderCode(?string $providerCode): string
{
$normalized = strtoupper(trim((string) $providerCode));
return $normalized !== '' ? $normalized : BetProvider::DEFAULT_CODE;
}
}