516 lines
19 KiB
PHP
516 lines
19 KiB
PHP
<?php
|
||
|
||
declare(strict_types=1);
|
||
|
||
namespace app\common\library;
|
||
|
||
use app\common\model\MallBusinessConfig;
|
||
use app\common\model\MallDailyPush;
|
||
use app\common\model\MallUserAsset;
|
||
use ba\Random;
|
||
use RuntimeException;
|
||
use support\think\Db;
|
||
use Throwable;
|
||
|
||
/**
|
||
* 将 daily_push_*.log 中保存的第三方原始推送幂等回放到业务表。
|
||
*/
|
||
class MallDailyPushLogReplay
|
||
{
|
||
private const QUERY_CHUNK_SIZE = 1000;
|
||
|
||
/**
|
||
* @return array<string, mixed>
|
||
*/
|
||
public function replay(string $logFile, bool $dryRun = false): array
|
||
{
|
||
$stats = $this->newStats($logFile, $dryRun);
|
||
$records = $this->parseRecords($logFile, $stats);
|
||
$stats['unique_members'] = count($records);
|
||
|
||
[$dailyMap, $assetByPlayxId, $assetByUsername] = $this->preloadExistingData($records);
|
||
$configRow = MallBusinessConfig::ensureRow();
|
||
$pointsConfig = MallPointsConversion::configFromRow($configRow);
|
||
$unlockRatio = floatval($pointsConfig['unlock_ratio'] ?? 0);
|
||
[$backupHandle, $backupPath] = $this->openBackup($dryRun);
|
||
$stats['backup_file'] = $backupPath;
|
||
|
||
try {
|
||
foreach ($records as $key => $record) {
|
||
try {
|
||
$daily = $dailyMap[$key] ?? null;
|
||
$asset = $assetByPlayxId[$record['user_id']]
|
||
?? ($record['username'] !== '' ? ($assetByUsername[$record['username']] ?? null) : null);
|
||
|
||
$result = $this->replayRecord(
|
||
$record,
|
||
$daily,
|
||
$asset,
|
||
$pointsConfig,
|
||
$unlockRatio,
|
||
$dryRun,
|
||
$backupHandle,
|
||
$stats
|
||
);
|
||
if ($result['daily']) {
|
||
$dailyMap[$key] = $result['daily'];
|
||
}
|
||
if ($result['asset']) {
|
||
$assetByPlayxId[$record['user_id']] = $result['asset'];
|
||
if ($record['username'] !== '') {
|
||
$assetByUsername[$record['username']] = $result['asset'];
|
||
}
|
||
}
|
||
} catch (Throwable $e) {
|
||
$stats['failed_members']++;
|
||
$this->appendError($stats, 'user_id=' . $record['user_id'] . ':' . $e->getMessage());
|
||
}
|
||
}
|
||
} finally {
|
||
if (is_resource($backupHandle)) {
|
||
fclose($backupHandle);
|
||
}
|
||
}
|
||
|
||
return $stats;
|
||
}
|
||
|
||
/**
|
||
* @return array<string, mixed>
|
||
*/
|
||
private function newStats(string $logFile, bool $dryRun): array
|
||
{
|
||
return [
|
||
'dry_run' => $dryRun,
|
||
'log_file' => $logFile,
|
||
'backup_file' => '',
|
||
'lines' => 0,
|
||
'members' => 0,
|
||
'unique_members' => 0,
|
||
'daily_created' => 0,
|
||
'daily_updated' => 0,
|
||
'daily_unchanged' => 0,
|
||
'assets_created' => 0,
|
||
'assets_updated' => 0,
|
||
'locked_points_delta' => 0,
|
||
'invalid_lines' => 0,
|
||
'failed_members' => 0,
|
||
'errors' => [],
|
||
];
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $stats
|
||
* @return array<string, array<string, mixed>>
|
||
*/
|
||
private function parseRecords(string $logFile, array &$stats): array
|
||
{
|
||
if (!is_file($logFile) || !is_readable($logFile)) {
|
||
throw new RuntimeException('日志文件不存在或不可读:' . $logFile);
|
||
}
|
||
$handle = fopen($logFile, 'rb');
|
||
if ($handle === false) {
|
||
throw new RuntimeException('无法打开日志文件:' . $logFile);
|
||
}
|
||
|
||
$records = [];
|
||
try {
|
||
while (($line = fgets($handle)) !== false) {
|
||
$line = trim($line);
|
||
if ($line === '') {
|
||
continue;
|
||
}
|
||
$stats['lines']++;
|
||
$entry = json_decode($line, true);
|
||
$raw = is_array($entry) ? strval($entry['raw'] ?? '') : '';
|
||
$payload = $raw !== '' ? json_decode($raw, true) : null;
|
||
if (!is_array($payload) || !isset($payload['report_date'], $payload['member']) || !is_array($payload['member'])) {
|
||
$stats['invalid_lines']++;
|
||
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行格式无效');
|
||
continue;
|
||
}
|
||
|
||
$date = $this->resolveDate($payload['report_date']);
|
||
if ($date === '') {
|
||
$stats['invalid_lines']++;
|
||
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行 report_date 无效');
|
||
continue;
|
||
}
|
||
$entryTime = strtotime(strval($entry['time'] ?? ''));
|
||
$createTime = $entryTime === false ? time() : $entryTime;
|
||
|
||
foreach ($payload['member'] as $member) {
|
||
$stats['members']++;
|
||
if (!is_array($member)) {
|
||
$stats['failed_members']++;
|
||
continue;
|
||
}
|
||
$userId = trim(strval($member['member_id'] ?? ''));
|
||
if ($userId === '') {
|
||
$stats['failed_members']++;
|
||
$this->appendError($stats, '第 ' . strval($stats['lines']) . ' 行存在空 member_id');
|
||
continue;
|
||
}
|
||
$key = $date . '|' . $userId;
|
||
$records[$key] = [
|
||
'user_id' => $userId,
|
||
'date' => $date,
|
||
'username' => trim(strval($member['login'] ?? '')),
|
||
'phone' => trim(strval($member['phone'] ?? ($member['mobile'] ?? ''))),
|
||
'yesterday_win_loss_net' => $this->numericValue($member, ['yesterday_total_wl', 'yesterday_total_w']),
|
||
'yesterday_total_deposit' => $this->numericValue($member, ['yesterday_total_deposit']),
|
||
'lifetime_total_deposit' => $this->numericValue($member, ['ltv_deposit', 'lty_deposit']),
|
||
'lifetime_total_withdraw' => $this->numericValue($member, ['ltv_withdrawal', 'lty_withdrawal']),
|
||
'create_time' => $createTime,
|
||
];
|
||
}
|
||
}
|
||
} finally {
|
||
fclose($handle);
|
||
}
|
||
return $records;
|
||
}
|
||
|
||
/**
|
||
* @param array<string, array<string, mixed>> $records
|
||
* @return array{0:array<string,MallDailyPush>,1:array<string,MallUserAsset>,2:array<string,MallUserAsset>}
|
||
*/
|
||
private function preloadExistingData(array $records): array
|
||
{
|
||
$idsByDate = [];
|
||
$userIds = [];
|
||
$usernames = [];
|
||
foreach ($records as $record) {
|
||
$date = $record['date'];
|
||
$userId = $record['user_id'];
|
||
$idsByDate[$date][$userId] = true;
|
||
$userIds[$userId] = true;
|
||
if ($record['username'] !== '') {
|
||
$usernames[$record['username']] = true;
|
||
}
|
||
}
|
||
|
||
$dailyMap = [];
|
||
foreach ($idsByDate as $date => $ids) {
|
||
foreach (array_chunk(array_keys($ids), self::QUERY_CHUNK_SIZE) as $idChunk) {
|
||
$rows = MallDailyPush::where('date', $date)->whereIn('user_id', $idChunk)->select();
|
||
foreach ($rows as $row) {
|
||
/** @var MallDailyPush $row */
|
||
$dailyMap[$date . '|' . strval($row->user_id)] = $row;
|
||
}
|
||
}
|
||
}
|
||
|
||
$assetByPlayxId = [];
|
||
foreach (array_chunk(array_keys($userIds), self::QUERY_CHUNK_SIZE) as $idChunk) {
|
||
$rows = MallUserAsset::whereIn('playx_user_id', $idChunk)->select();
|
||
foreach ($rows as $row) {
|
||
$assetByPlayxId[strval($row->playx_user_id)] = $row;
|
||
}
|
||
}
|
||
|
||
$assetByUsername = [];
|
||
foreach (array_chunk(array_keys($usernames), self::QUERY_CHUNK_SIZE) as $usernameChunk) {
|
||
$rows = MallUserAsset::whereIn('username', $usernameChunk)->select();
|
||
foreach ($rows as $row) {
|
||
$assetByUsername[strval($row->username)] = $row;
|
||
}
|
||
}
|
||
|
||
return [$dailyMap, $assetByPlayxId, $assetByUsername];
|
||
}
|
||
|
||
/**
|
||
* @return array{0:mixed,1:string}
|
||
*/
|
||
private function openBackup(bool $dryRun): array
|
||
{
|
||
if ($dryRun) {
|
||
return [null, ''];
|
||
}
|
||
$backupDir = runtime_path('backup/daily_push_replay');
|
||
if (!is_dir($backupDir) && !mkdir($backupDir, 0755, true) && !is_dir($backupDir)) {
|
||
throw new RuntimeException('无法创建回放备份目录');
|
||
}
|
||
$path = $backupDir . DIRECTORY_SEPARATOR . 'before_replay_' . date('Ymd_His') . '.jsonl';
|
||
$handle = fopen($path, 'wb');
|
||
if ($handle === false) {
|
||
throw new RuntimeException('无法创建回放备份文件');
|
||
}
|
||
return [$handle, $path];
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $record
|
||
* @param resource|null $backupHandle
|
||
* @param array<string, mixed> $stats
|
||
* @return array{daily:?MallDailyPush,asset:?MallUserAsset}
|
||
*/
|
||
private function replayRecord(
|
||
array $record,
|
||
?MallDailyPush $daily,
|
||
?MallUserAsset $asset,
|
||
array $pointsConfig,
|
||
float $unlockRatio,
|
||
bool $dryRun,
|
||
$backupHandle,
|
||
array &$stats
|
||
): array {
|
||
$oldWinLossNet = $daily ? floatval($daily->yesterday_win_loss_net ?? 0) : 0.0;
|
||
$newLocked = MallPointsConversion::dailyConvertedPoints(
|
||
$record['date'],
|
||
floatval($record['yesterday_win_loss_net']),
|
||
$pointsConfig,
|
||
MallPointsConversion::nextCalendarDate($record['date'])
|
||
);
|
||
$oldLocked = $daily
|
||
? MallPointsConversion::dailyConvertedPoints(
|
||
$record['date'],
|
||
$oldWinLossNet,
|
||
$pointsConfig,
|
||
MallPointsConversion::nextCalendarDate($record['date'])
|
||
)
|
||
: 0;
|
||
$lockedDelta = $newLocked - $oldLocked;
|
||
$todayLimit = MallPointsConversion::dailyPushTodayLimit(
|
||
$record['date'],
|
||
floatval($record['yesterday_total_deposit']),
|
||
$pointsConfig,
|
||
MallPointsConversion::nextCalendarDate($record['date'])
|
||
);
|
||
$dailyData = [
|
||
'user_id' => $record['user_id'],
|
||
'date' => $record['date'],
|
||
'username' => $record['username'],
|
||
'yesterday_win_loss_net' => $record['yesterday_win_loss_net'],
|
||
'converted_points' => $newLocked,
|
||
'yesterday_total_deposit' => $record['yesterday_total_deposit'],
|
||
'lifetime_total_deposit' => $record['lifetime_total_deposit'],
|
||
'lifetime_total_withdraw' => $record['lifetime_total_withdraw'],
|
||
];
|
||
$dailyChanged = !$daily || $this->dailyHasChanges($daily, $dailyData);
|
||
$this->countDailyChange($stats, $daily, $dailyChanged);
|
||
|
||
$assetWillCreate = !$asset;
|
||
$assetWillUpdate = $assetWillCreate
|
||
|| trim(strval($asset->playx_user_id ?? '')) !== $record['user_id']
|
||
|| $lockedDelta !== 0
|
||
|| ($record['username'] !== '' && trim(strval($asset->username ?? '')) !== $record['username'])
|
||
|| ($record['phone'] !== '' && trim(strval($asset->phone ?? '')) === '')
|
||
|| $this->shouldUpdateDailyLimit($asset, $record['date'], $todayLimit);
|
||
if ($assetWillCreate) {
|
||
$stats['assets_created']++;
|
||
} elseif ($assetWillUpdate) {
|
||
$stats['assets_updated']++;
|
||
}
|
||
$stats['locked_points_delta'] += $lockedDelta;
|
||
|
||
if ($dryRun) {
|
||
return ['daily' => $daily, 'asset' => $asset];
|
||
}
|
||
if (!$dailyChanged && !$assetWillUpdate) {
|
||
return ['daily' => $daily, 'asset' => $asset];
|
||
}
|
||
|
||
$this->writeBackup($backupHandle, $record['date'] . '|' . $record['user_id'], $daily, $asset);
|
||
Db::startTrans();
|
||
try {
|
||
if ($daily) {
|
||
if ($dailyChanged) {
|
||
$daily->save($dailyData);
|
||
}
|
||
} else {
|
||
$dailyData['create_time'] = $record['create_time'];
|
||
$daily = MallDailyPush::create($dailyData);
|
||
}
|
||
|
||
if (!$asset) {
|
||
$effectiveUsername = $record['username'] !== '' ? $record['username'] : 'playx_' . $record['user_id'];
|
||
$asset = MallUserAsset::create([
|
||
'playx_user_id' => $record['user_id'],
|
||
'username' => $effectiveUsername,
|
||
'phone' => $record['phone'],
|
||
'password' => hash_password(Random::build('alnum', 16)),
|
||
'admin_id' => 0,
|
||
'locked_points' => 0,
|
||
'available_points' => 0,
|
||
'today_limit' => 0,
|
||
'today_claimed' => 0,
|
||
'today_limit_date' => null,
|
||
'points_claim_cap' => 0,
|
||
'total_points_claimed' => 0,
|
||
'lifetime_converted_points' => 0,
|
||
'daily_converted_points' => 0,
|
||
'create_time' => $record['create_time'],
|
||
'update_time' => time(),
|
||
]);
|
||
if (!$asset) {
|
||
throw new RuntimeException('创建用户资产失败');
|
||
}
|
||
/** @var MallUserAsset $asset */
|
||
} elseif ($assetWillUpdate) {
|
||
$asset->playx_user_id = $record['user_id'];
|
||
if ($record['username'] !== '') {
|
||
$asset->username = $record['username'];
|
||
}
|
||
if ($record['phone'] !== '' && trim(strval($asset->phone ?? '')) === '') {
|
||
$asset->phone = $record['phone'];
|
||
}
|
||
}
|
||
|
||
if ($assetWillUpdate) {
|
||
if ($lockedDelta !== 0) {
|
||
MallPointsConversion::applyDailyConversionDelta($asset, $lockedDelta);
|
||
}
|
||
$this->applyDailyLimit($asset, $record['date'], $todayLimit);
|
||
$settlementDay = MallPointsConversion::nextCalendarDate($record['date']);
|
||
MallPointsConversion::applyEventPeriodWinLossForDailyPush(
|
||
$asset,
|
||
$record['date'],
|
||
floatval($record['yesterday_win_loss_net']),
|
||
$pointsConfig,
|
||
$settlementDay,
|
||
$daily ? $oldWinLossNet : null
|
||
);
|
||
$asset->save();
|
||
}
|
||
Db::commit();
|
||
} catch (Throwable $e) {
|
||
Db::rollback();
|
||
throw $e;
|
||
}
|
||
|
||
return ['daily' => $daily, 'asset' => $asset];
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $stats
|
||
*/
|
||
private function countDailyChange(array &$stats, ?MallDailyPush $daily, bool $changed): void
|
||
{
|
||
if (!$daily) {
|
||
$stats['daily_created']++;
|
||
} elseif ($changed) {
|
||
$stats['daily_updated']++;
|
||
} else {
|
||
$stats['daily_unchanged']++;
|
||
}
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $member
|
||
* @param array<int, string> $keys
|
||
*/
|
||
private function numericValue(array $member, array $keys): float
|
||
{
|
||
foreach ($keys as $key) {
|
||
if (array_key_exists($key, $member) && is_numeric($member[$key])) {
|
||
return round(floatval($member[$key]), 2);
|
||
}
|
||
}
|
||
return 0.0;
|
||
}
|
||
|
||
private function resolveDate(mixed $reportDate): string
|
||
{
|
||
if (is_numeric($reportDate)) {
|
||
$timestamp = intval($reportDate);
|
||
return $timestamp > 0 ? date('Y-m-d', $timestamp) : '';
|
||
}
|
||
$date = trim(strval($reportDate));
|
||
$parsed = \DateTime::createFromFormat('Y-m-d', $date);
|
||
return $parsed && $parsed->format('Y-m-d') === $date ? $date : '';
|
||
}
|
||
|
||
private function lockedContribution(float $winLossNet, float $returnRatio): int
|
||
{
|
||
return $winLossNet < 0 ? intval(round(abs($winLossNet) * $returnRatio)) : 0;
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $dailyData
|
||
*/
|
||
private function dailyHasChanges(MallDailyPush $daily, array $dailyData): bool
|
||
{
|
||
if (trim(strval($daily->username ?? '')) !== $dailyData['username']) {
|
||
return true;
|
||
}
|
||
foreach ([
|
||
'yesterday_win_loss_net',
|
||
'converted_points',
|
||
'yesterday_total_deposit',
|
||
'lifetime_total_deposit',
|
||
'lifetime_total_withdraw',
|
||
] as $field) {
|
||
if (abs(floatval($daily->$field ?? 0) - floatval($dailyData[$field])) > 0.000001) {
|
||
return true;
|
||
}
|
||
}
|
||
return false;
|
||
}
|
||
|
||
private function shouldUpdateDailyLimit(?MallUserAsset $asset, string $date, int $todayLimit): bool
|
||
{
|
||
if (!$asset) {
|
||
return true;
|
||
}
|
||
$currentDate = trim(strval($asset->today_limit_date ?? ''));
|
||
if ($currentDate !== '' && $currentDate > $date) {
|
||
return false;
|
||
}
|
||
return $currentDate !== $date || intval($asset->today_limit ?? 0) !== $todayLimit;
|
||
}
|
||
|
||
private function applyDailyLimit(MallUserAsset $asset, string $date, int $todayLimit): void
|
||
{
|
||
$currentDate = trim(strval($asset->today_limit_date ?? ''));
|
||
if ($currentDate !== '' && $currentDate > $date) {
|
||
return;
|
||
}
|
||
if ($currentDate !== $date) {
|
||
$asset->today_claimed = 0;
|
||
$asset->today_limit_date = $date;
|
||
}
|
||
$asset->today_limit = $todayLimit;
|
||
}
|
||
|
||
/**
|
||
* @param resource|null $handle
|
||
*/
|
||
private function writeBackup($handle, string $key, ?MallDailyPush $daily, ?MallUserAsset $asset): void
|
||
{
|
||
if (!is_resource($handle)) {
|
||
return;
|
||
}
|
||
$assetData = null;
|
||
if ($asset) {
|
||
$assetData = [
|
||
'id' => $asset->id,
|
||
'playx_user_id' => $asset->playx_user_id,
|
||
'username' => $asset->username,
|
||
'phone' => $asset->phone,
|
||
'locked_points' => $asset->locked_points,
|
||
'today_limit' => $asset->today_limit,
|
||
'today_claimed' => $asset->today_claimed,
|
||
'today_limit_date' => $asset->today_limit_date,
|
||
'update_time' => $asset->update_time,
|
||
];
|
||
}
|
||
fwrite($handle, json_encode([
|
||
'key' => $key,
|
||
'daily_push' => $daily ? $daily->toArray() : null,
|
||
'user_asset' => $assetData,
|
||
], JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES) . PHP_EOL);
|
||
}
|
||
|
||
/**
|
||
* @param array<string, mixed> $stats
|
||
*/
|
||
private function appendError(array &$stats, string $message): void
|
||
{
|
||
if (count($stats['errors']) < 100) {
|
||
$stats['errors'][] = $message;
|
||
}
|
||
}
|
||
}
|