336 lines
13 KiB
PHP
336 lines
13 KiB
PHP
<?php
|
|
|
|
declare(strict_types=1);
|
|
|
|
namespace app\common\service\qywx;
|
|
|
|
use think\facade\Db;
|
|
|
|
/** 根据实际获客回调记账,并为同一个官方链接规划下一名承接成员。 */
|
|
class QywxPromotionMemberSchedulerService
|
|
{
|
|
public static function poolIdFromState(string $state): int
|
|
{
|
|
$state = trim($state);
|
|
|
|
return preg_match('/^zyt_pool:(\d+)$/', $state, $matches) === 1
|
|
? max(0, (int) $matches[1])
|
|
: 0;
|
|
}
|
|
|
|
/** @return array{status:string,pool_id:int,member_id:int,next_member_id:int} */
|
|
public static function recordFromState(
|
|
string $state,
|
|
string $userId,
|
|
string $externalUserId,
|
|
int $eventTime = 0,
|
|
string $source = 'external_contact'
|
|
): array {
|
|
$poolId = self::poolIdFromState($state);
|
|
if ($poolId <= 0) {
|
|
return ['status' => 'ignored_state', 'pool_id' => 0, 'member_id' => 0, 'next_member_id' => 0];
|
|
}
|
|
|
|
return self::record($poolId, $userId, $externalUserId, $eventTime, $source);
|
|
}
|
|
|
|
/** @return array{status:string,pool_id:int,member_id:int,next_member_id:int} */
|
|
public static function recordFromRemoteLink(
|
|
string $remoteLinkId,
|
|
string $userId,
|
|
string $externalUserId,
|
|
int $eventTime = 0,
|
|
string $source = 'customer_acquisition'
|
|
): array {
|
|
$poolId = (int) (Db::name('qywx_promotion_link')
|
|
->where('remote_link_id', trim($remoteLinkId))
|
|
->whereNull('delete_time')
|
|
->value('pool_id') ?? 0);
|
|
if ($poolId <= 0) {
|
|
return ['status' => 'ignored_link', 'pool_id' => 0, 'member_id' => 0, 'next_member_id' => 0];
|
|
}
|
|
|
|
return self::record($poolId, $userId, $externalUserId, $eventTime, $source);
|
|
}
|
|
|
|
/**
|
|
* 管理端变更成员开关、权重或上限后,确保当前远端成员仍可用;必要时立即排队切换。
|
|
*
|
|
* @return array{pool_id:int,next_member_id:int,queued:bool,blocked:bool}
|
|
*/
|
|
public static function reconcilePool(int $poolId): array
|
|
{
|
|
return Db::transaction(function () use ($poolId): array {
|
|
$now = time();
|
|
$today = date('Y-m-d', $now);
|
|
$linkId = self::promotionLinkId($poolId);
|
|
$sync = self::lockedSyncRow($poolId);
|
|
$members = self::lockedMembers($poolId);
|
|
if ($linkId <= 0 || $members === []) {
|
|
return ['pool_id' => $poolId, 'next_member_id' => 0, 'queued' => false, 'blocked' => true];
|
|
}
|
|
|
|
$desiredId = (int) ($sync['desired_member_id'] ?? 0);
|
|
$desiredEligible = false;
|
|
foreach ($members as &$member) {
|
|
if ((string) ($member['today_date'] ?? '') !== $today) {
|
|
$member['today_date'] = $today;
|
|
$member['today_count'] = 0;
|
|
}
|
|
if ((int) $member['id'] === $desiredId) {
|
|
$desiredEligible = QywxPromotionWeightedRandom::eligible($member, $now);
|
|
}
|
|
}
|
|
unset($member);
|
|
self::persistMemberCursors($members, $now);
|
|
if ($desiredEligible) {
|
|
if ((int) ($sync['status'] ?? 0) === 4) {
|
|
Db::name('qywx_promotion_range_sync')->where('pool_id', $poolId)->update([
|
|
'status' => 0,
|
|
'next_retry' => 0,
|
|
'last_error' => '',
|
|
'update_time' => $now,
|
|
]);
|
|
}
|
|
return ['pool_id' => $poolId, 'next_member_id' => $desiredId, 'queued' => false, 'blocked' => false];
|
|
}
|
|
|
|
return self::selectAndQueueLocked($poolId, $linkId, $sync, $members, $today, $now);
|
|
});
|
|
}
|
|
|
|
public static function initialisePool(int $poolId, int $promotionLinkId, string $selectedUserId): int
|
|
{
|
|
$memberId = (int) (Db::name('qywx_promotion_pool_member')
|
|
->where('pool_id', $poolId)
|
|
->where('userid', $selectedUserId)
|
|
->whereNull('delete_time')
|
|
->value('id') ?? 0);
|
|
if ($memberId <= 0 || $promotionLinkId <= 0) {
|
|
return 0;
|
|
}
|
|
$now = time();
|
|
$existing = Db::name('qywx_promotion_range_sync')->where('pool_id', $poolId)->find();
|
|
$data = [
|
|
'promotion_link_id' => $promotionLinkId,
|
|
'desired_member_id' => $memberId,
|
|
'applied_member_id' => $memberId,
|
|
'status' => 0,
|
|
'attempts' => 0,
|
|
'next_retry' => 0,
|
|
'lock_token' => '',
|
|
'lock_until' => 0,
|
|
'last_error' => '',
|
|
'update_time' => $now,
|
|
];
|
|
if ($existing) {
|
|
$version = max(1, (int) ($existing['desired_version'] ?? 0) + 1);
|
|
$data['desired_version'] = $version;
|
|
$data['applied_version'] = $version;
|
|
Db::name('qywx_promotion_range_sync')->where('pool_id', $poolId)->update($data);
|
|
} else {
|
|
$data += ['pool_id' => $poolId, 'desired_version' => 1, 'applied_version' => 1, 'create_time' => $now];
|
|
Db::name('qywx_promotion_range_sync')->insert($data);
|
|
}
|
|
|
|
return $memberId;
|
|
}
|
|
|
|
/** @return array{status:string,pool_id:int,member_id:int,next_member_id:int} */
|
|
private static function record(
|
|
int $poolId,
|
|
string $userId,
|
|
string $externalUserId,
|
|
int $eventTime,
|
|
string $source
|
|
): array {
|
|
$userId = trim($userId);
|
|
$externalUserId = trim($externalUserId);
|
|
if ($poolId <= 0 || $userId === '' || $externalUserId === '') {
|
|
return ['status' => 'ignored_identity', 'pool_id' => $poolId, 'member_id' => 0, 'next_member_id' => 0];
|
|
}
|
|
|
|
return Db::transaction(function () use ($poolId, $userId, $externalUserId, $eventTime, $source): array {
|
|
$pool = Db::name('qywx_promotion_pool')->where('id', $poolId)->whereNull('delete_time')->lock(true)->find();
|
|
if (!$pool) {
|
|
return ['status' => 'ignored_pool', 'pool_id' => $poolId, 'member_id' => 0, 'next_member_id' => 0];
|
|
}
|
|
$members = self::lockedMembers($poolId);
|
|
$actualMember = null;
|
|
foreach ($members as $member) {
|
|
if ((string) ($member['userid'] ?? '') === $userId) {
|
|
$actualMember = $member;
|
|
break;
|
|
}
|
|
}
|
|
if ($actualMember === null) {
|
|
return ['status' => 'ignored_member', 'pool_id' => $poolId, 'member_id' => 0, 'next_member_id' => 0];
|
|
}
|
|
$actualMemberId = (int) $actualMember['id'];
|
|
$eventKey = hash('sha256', $poolId . '|' . $userId . '|' . $externalUserId);
|
|
try {
|
|
Db::name('qywx_promotion_dispatch_event')->insert([
|
|
'event_key' => $eventKey,
|
|
'pool_id' => $poolId,
|
|
'member_id' => $actualMemberId,
|
|
'userid' => $userId,
|
|
'external_userid' => $externalUserId,
|
|
'source' => mb_substr($source, 0, 32),
|
|
'event_time' => max(0, $eventTime),
|
|
'create_time' => time(),
|
|
]);
|
|
} catch (\Throwable $e) {
|
|
if (!Db::name('qywx_promotion_dispatch_event')->where('event_key', $eventKey)->find()) {
|
|
throw $e;
|
|
}
|
|
|
|
return ['status' => 'duplicate', 'pool_id' => $poolId, 'member_id' => $actualMemberId, 'next_member_id' => 0];
|
|
}
|
|
|
|
$now = time();
|
|
$today = date('Y-m-d', $now);
|
|
foreach ($members as &$member) {
|
|
if ((string) ($member['today_date'] ?? '') !== $today) {
|
|
$member['today_date'] = $today;
|
|
$member['today_count'] = 0;
|
|
}
|
|
if ((int) $member['id'] === $actualMemberId) {
|
|
$member['today_count'] = (int) ($member['today_count'] ?? 0) + 1;
|
|
$member['total_count'] = (int) ($member['total_count'] ?? 0) + 1;
|
|
$member['last_assigned_time'] = max($now, max(0, $eventTime));
|
|
}
|
|
}
|
|
unset($member);
|
|
self::persistMemberCursors($members, $now);
|
|
|
|
$linkId = self::promotionLinkId($poolId);
|
|
$sync = self::lockedSyncRow($poolId);
|
|
$desiredId = (int) ($sync['desired_member_id'] ?? 0);
|
|
if ($linkId <= 0 || ($desiredId > 0 && $desiredId !== $actualMemberId)) {
|
|
return ['status' => 'counted_stale', 'pool_id' => $poolId, 'member_id' => $actualMemberId, 'next_member_id' => $desiredId];
|
|
}
|
|
|
|
$planned = self::selectAndQueueLocked($poolId, $linkId, $sync, $members, $today, $now);
|
|
|
|
return [
|
|
'status' => $planned['blocked'] ? 'counted_blocked' : 'counted',
|
|
'pool_id' => $poolId,
|
|
'member_id' => $actualMemberId,
|
|
'next_member_id' => $planned['next_member_id'],
|
|
];
|
|
});
|
|
}
|
|
|
|
/** @return array{pool_id:int,next_member_id:int,queued:bool,blocked:bool} */
|
|
private static function selectAndQueueLocked(
|
|
int $poolId,
|
|
int $linkId,
|
|
?array $sync,
|
|
array $members,
|
|
string $today,
|
|
int $now
|
|
): array {
|
|
$selection = QywxPromotionWeightedRandom::select($members, $today, $now);
|
|
self::persistMemberCursors($selection['members'], $now);
|
|
$selectedId = (int) $selection['selected_id'];
|
|
if ($selectedId <= 0) {
|
|
self::upsertSync($poolId, $linkId, (int) ($sync['desired_member_id'] ?? 0), false, $sync, $now, '所有成员均已禁用、未生效或达到今日上限');
|
|
|
|
return ['pool_id' => $poolId, 'next_member_id' => 0, 'queued' => false, 'blocked' => true];
|
|
}
|
|
$changed = $selectedId !== (int) ($sync['desired_member_id'] ?? 0);
|
|
self::upsertSync($poolId, $linkId, $selectedId, $changed, $sync, $now);
|
|
|
|
return ['pool_id' => $poolId, 'next_member_id' => $selectedId, 'queued' => $changed, 'blocked' => false];
|
|
}
|
|
|
|
private static function upsertSync(
|
|
int $poolId,
|
|
int $linkId,
|
|
int $desiredMemberId,
|
|
bool $pending,
|
|
?array $existing,
|
|
int $now,
|
|
string $error = ''
|
|
): void {
|
|
$version = max(1, (int) ($existing['desired_version'] ?? 0) + ($pending ? 1 : 0));
|
|
$data = [
|
|
'promotion_link_id' => $linkId,
|
|
'desired_member_id' => $desiredMemberId,
|
|
'desired_version' => $version,
|
|
'status' => $error !== '' ? 4 : ($pending ? 1 : (int) ($existing['status'] ?? 0)),
|
|
'next_retry' => $error !== '' ? strtotime('tomorrow', $now) : ($pending ? $now : 0),
|
|
'last_error' => mb_substr($error, 0, 500),
|
|
'update_time' => $now,
|
|
];
|
|
if ($pending) {
|
|
$data['attempts'] = 0;
|
|
}
|
|
if ($existing) {
|
|
Db::name('qywx_promotion_range_sync')->where('pool_id', $poolId)->update($data);
|
|
} else {
|
|
$data += [
|
|
'pool_id' => $poolId,
|
|
'applied_member_id' => 0,
|
|
'applied_version' => 0,
|
|
'attempts' => 0,
|
|
'lock_token' => '',
|
|
'lock_until' => 0,
|
|
'create_time' => $now,
|
|
];
|
|
Db::name('qywx_promotion_range_sync')->insert($data);
|
|
}
|
|
}
|
|
|
|
/** @return list<array<string,mixed>> */
|
|
private static function lockedMembers(int $poolId): array
|
|
{
|
|
return Db::name('qywx_promotion_pool_member')
|
|
->where('pool_id', $poolId)
|
|
->whereNull('delete_time')
|
|
->order('id', 'asc')
|
|
->lock(true)
|
|
->select()->toArray();
|
|
}
|
|
|
|
/** @return array<string,mixed>|null */
|
|
private static function lockedSyncRow(int $poolId): ?array
|
|
{
|
|
$row = Db::name('qywx_promotion_range_sync')->where('pool_id', $poolId)->lock(true)->find();
|
|
|
|
return $row ?: null;
|
|
}
|
|
|
|
private static function promotionLinkId(int $poolId): int
|
|
{
|
|
$syncLinkId = (int) (Db::name('qywx_promotion_range_sync')
|
|
->where('pool_id', $poolId)->value('promotion_link_id') ?? 0);
|
|
if ($syncLinkId > 0) {
|
|
return $syncLinkId;
|
|
}
|
|
|
|
return (int) (Db::name('qywx_promotion_link')
|
|
->where('pool_id', $poolId)
|
|
->where('remote_link_id', '<>', '')
|
|
->where('remote_status', '<>', 2)
|
|
->whereNull('delete_time')
|
|
->order('id', 'desc')
|
|
->value('id') ?? 0);
|
|
}
|
|
|
|
/** @param list<array<string,mixed>> $members */
|
|
private static function persistMemberCursors(array $members, int $now): void
|
|
{
|
|
foreach ($members as $member) {
|
|
Db::name('qywx_promotion_pool_member')->where('id', (int) $member['id'])->update([
|
|
'current_weight' => (int) ($member['current_weight'] ?? 0),
|
|
'today_count' => max(0, (int) ($member['today_count'] ?? 0)),
|
|
'today_date' => (string) ($member['today_date'] ?? '') ?: null,
|
|
'total_count' => max(0, (int) ($member['total_count'] ?? 0)),
|
|
'last_assigned_time' => max(0, (int) ($member['last_assigned_time'] ?? 0)),
|
|
'update_time' => $now,
|
|
]);
|
|
}
|
|
}
|
|
}
|