Update data calibration and recover stale TCP queue
This commit is contained in:
+75
@@ -284,6 +284,7 @@ function tcp_queue_get(string $requestId): ?array
|
||||
function tcp_queue_claim_next(PDO $pdo): ?array
|
||||
{
|
||||
db_ensure_tcp_queue_schema($pdo);
|
||||
tcp_queue_recover_stale($pdo);
|
||||
|
||||
$pdo->beginTransaction();
|
||||
try {
|
||||
@@ -313,6 +314,80 @@ function tcp_queue_claim_next(PDO $pdo): ?array
|
||||
}
|
||||
}
|
||||
|
||||
function tcp_queue_has_usage(PDO $pdo, string $requestId): bool
|
||||
{
|
||||
$stmt = $pdo->prepare("SELECT 1 FROM car_data_usage WHERE request_id = :request_id LIMIT 1");
|
||||
$stmt->execute([':request_id' => $requestId]);
|
||||
return (bool)$stmt->fetchColumn();
|
||||
}
|
||||
|
||||
function tcp_queue_recover_stale(PDO $pdo): void
|
||||
{
|
||||
$now = microtime(true);
|
||||
$stmt = $pdo->prepare("
|
||||
SELECT *
|
||||
FROM car_tcp_request_queue
|
||||
WHERE status IN ('processing', 'ui_timeout_pending')
|
||||
AND usage_deadline_ts > 0
|
||||
AND usage_deadline_ts < :now
|
||||
ORDER BY id ASC
|
||||
LIMIT 20
|
||||
");
|
||||
$stmt->execute([':now' => $now]);
|
||||
$rows = $stmt->fetchAll();
|
||||
|
||||
foreach ($rows as $row) {
|
||||
$requestId = (string)($row['request_id'] ?? '');
|
||||
if ($requestId === '') {
|
||||
continue;
|
||||
}
|
||||
|
||||
$sentBytes = max(0, (int)($row['sent_bytes'] ?? 0));
|
||||
$receivedBytes = max(0, (int)($row['received_bytes'] ?? 0));
|
||||
$connectMs = max(0, (int)($row['connect_ms'] ?? 0));
|
||||
$readMs = max(0, (int)($row['read_ms'] ?? 0));
|
||||
$hasUsage = tcp_queue_has_usage($pdo, $requestId);
|
||||
|
||||
if ($hasUsage) {
|
||||
tcp_queue_update($pdo, $requestId, [
|
||||
'status' => 'worker_lost_usage_recorded',
|
||||
'tcp_ok' => 0,
|
||||
'tcp_error' => 'worker_lost_after_usage_recorded',
|
||||
]);
|
||||
continue;
|
||||
}
|
||||
|
||||
if ($sentBytes <= 0) {
|
||||
tcp_queue_record_usage_and_status(
|
||||
$pdo,
|
||||
$row,
|
||||
'expired_before_send',
|
||||
0,
|
||||
0,
|
||||
false,
|
||||
'worker_lost_before_send',
|
||||
'',
|
||||
$connectMs,
|
||||
$readMs
|
||||
);
|
||||
continue;
|
||||
}
|
||||
|
||||
tcp_queue_record_usage_and_status(
|
||||
$pdo,
|
||||
$row,
|
||||
'settled_sent_only',
|
||||
$sentBytes,
|
||||
$receivedBytes,
|
||||
false,
|
||||
'worker_lost_late_settlement_sent_only',
|
||||
'',
|
||||
$connectMs,
|
||||
$readMs
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
function tcp_queue_record_usage_and_status(
|
||||
PDO $pdo,
|
||||
array $job,
|
||||
|
||||
Reference in New Issue
Block a user