Files
webman-buildadmin-mall/app/common/library/MallDailyPushLogReplay.php

496 lines
18 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
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
);
$oldLocked = $daily
? MallPointsConversion::dailyConvertedPoints($record['date'], $oldWinLossNet, $pointsConfig)
: 0;
$lockedDelta = $newLocked - $oldLocked;
$todayLimit = intval(round($record['yesterday_total_deposit'] * $unlockRatio));
$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);
$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;
}
}
}