forked from fenjo26/Orbitra.link
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpostback_queue_cron.php
More file actions
349 lines (319 loc) · 15.7 KB
/
Copy pathpostback_queue_cron.php
File metadata and controls
349 lines (319 loc) · 15.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
<?php
// postback_queue_cron.php
// Cron worker: delivers queued outbound S2S postbacks with exponential backoff retry.
//
// A row is enqueued by postback.php with status='pending'. This worker picks due rows
// (next_retry_at <= now, attempts < MAX_ATTEMPTS), performs the HTTP call, and on a
// non-2xx/3xx response schedules the next attempt with growing delay. Rows that exhaust
// MAX_ATTEMPTS are marked status='failed' and kept for inspection in the S2S logs UI.
//
// Example cron (every minute):
// * * * * * php /var/www/orbitra/postback_queue_cron.php >> /var/log/orbitra_postback_queue.log 2>&1
require_once __DIR__ . '/config.php';
function orbitraPqSetSetting(PDO $pdo, string $key, string $value): void
{
try {
$stmt = $pdo->prepare("
INSERT INTO settings (key, value, updated_at)
VALUES (?, ?, datetime('now'))
ON CONFLICT(key) DO UPDATE SET
value = excluded.value,
updated_at = datetime('now')
");
$stmt->execute([$key, $value]);
} catch (Throwable $e) {
// Non-fatal: worker should still attempt delivery.
}
}
function orbitraPqLog(string $msg): void
{
echo '[' . date('Y-m-d H:i:s') . '] ' . $msg . PHP_EOL;
}
// --- Configuration -----------------------------------------------------------
// A row is delivered on attempt 1, then retried after each entry of PQ_BACKOFF_SECONDS.
// So MAX_ATTEMPTS must be count(PQ_BACKOFF_SECONDS) + 1 for the whole schedule — including
// the final 24h wait — to actually be used before the row is marked dead.
const PQ_BACKOFF_SECONDS = [60, 300, 1800, 7200, 86400]; // 1m, 5m, 30m, 2h, 24h
const PQ_MAX_ATTEMPTS = 6; // 1 initial attempt + 5 retries
const PQ_BATCH_SIZE = 50; // rows per run
const PQ_HTTP_TIMEOUT = 8; // seconds per delivery attempt
// A row claimed as 'in_flight' longer than this is assumed to belong to a worker that
// died mid-delivery (fatal error, OOM, cron kill) and is returned to the queue.
const PQ_INFLIGHT_STALE_SECONDS = 600; // 10 minutes
// --- Single-instance lock ----------------------------------------------------
$lockFile = __DIR__ . '/var/locks/postback_queue.lock';
$lockTtlSeconds = 300;
if (!is_dir(__DIR__ . '/var/locks')) {
@mkdir(__DIR__ . '/var/locks', 0777, true);
}
$fp = @fopen($lockFile, 'c+');
if ($fp) {
if (!flock($fp, LOCK_EX | LOCK_NB)) {
// Another worker running.
exit(0);
}
$st = fstat($fp);
if ($st && isset($st['mtime']) && (time() - (int) $st['mtime']) > $lockTtlSeconds) {
ftruncate($fp, 0);
}
}
$processed = 0;
$delivered = 0;
$requeued = 0;
$failed = 0;
// Defined outside the try so the catch block can always reference it.
$ts = date('Y-m-d H:i:s');
try {
orbitraPqSetSetting($pdo, 'postback_queue_last_ping_at', $ts);
// Allow disabling the worker from UI while keeping cron in place.
$enabled = '1';
try {
$val = $pdo->query("SELECT value FROM settings WHERE key='postback_queue_enabled'")->fetchColumn();
if (is_string($val) && $val !== '') {
$enabled = $val;
}
} catch (Throwable $e) {
// Ignore.
}
if ($enabled === '0') {
orbitraPqLog("postback_queue: disabled via settings");
exit(0);
}
// Reclaim rows abandoned by a crashed worker. Without this a row that was claimed
// as 'in_flight' and never finished would be invisible to the due-query forever,
// because that query only looks at 'pending'.
try {
// attempts is incremented so a "poison" row that reliably kills the worker
// burns through its budget instead of looping forever.
$reap = $pdo->prepare("
UPDATE s2s_postbacks_log
SET status = 'pending',
attempts = attempts + 1,
updated_at = datetime('now'),
last_error = 'worker died mid-delivery; requeued'
WHERE status = 'in_flight'
AND COALESCE(updated_at, created_at) <= datetime('now', ?)
");
$reap->execute(['-' . PQ_INFLIGHT_STALE_SECONDS . ' seconds']);
$reaped = $reap->rowCount();
if ($reaped > 0) {
orbitraPqLog("postback_queue: requeued $reaped stale in_flight row(s)");
}
// Retire pending rows that are out of attempts. Without this they would linger
// as 'pending' forever, invisible to the due-query and misleading in the UI.
$retire = $pdo->prepare("
UPDATE s2s_postbacks_log
SET status = 'failed', updated_at = datetime('now')
WHERE status = 'pending' AND attempts >= " . PQ_MAX_ATTEMPTS . "
");
$retire->execute();
} catch (Throwable $e) {
orbitraPqLog('postback_queue: reaper failed: ' . $e->getMessage());
}
// Select due rows. Claiming is done by flipping status to 'in_flight' inside a
// transaction so a parallel worker cannot pick the same row.
$dueStmt = $pdo->prepare("
SELECT id, conversion_id, url, method, attempts, payload_json, content_type, proxy_url, headers_json
FROM s2s_postbacks_log
WHERE status = 'pending'
AND next_retry_at <= datetime('now')
AND attempts < " . PQ_MAX_ATTEMPTS . "
ORDER BY next_retry_at ASC
LIMIT " . PQ_BATCH_SIZE . "
");
$dueStmt->execute();
$rows = $dueStmt->fetchAll(PDO::FETCH_ASSOC);
if (empty($rows)) {
orbitraPqSetSetting($pdo, 'postback_queue_last_checked_at', $ts);
exit(0);
}
// Every transition stamps updated_at so the reaper above can tell a live delivery
// from an abandoned one.
$claimStmt = $pdo->prepare("UPDATE s2s_postbacks_log SET status = 'in_flight', updated_at = datetime('now') WHERE id = ? AND status = 'pending'");
$doneStmt = $pdo->prepare("UPDATE s2s_postbacks_log SET status = 'delivered', http_code = ?, status_code = ?, last_error = NULL, attempts = attempts + 1, updated_at = datetime('now') WHERE id = ?");
$retryStmt = $pdo->prepare("UPDATE s2s_postbacks_log SET status = 'pending', attempts = attempts + 1, next_retry_at = datetime('now', ?), http_code = ?, status_code = ?, last_error = ?, updated_at = datetime('now') WHERE id = ?");
$deadStmt = $pdo->prepare("UPDATE s2s_postbacks_log SET status = 'failed', attempts = attempts + 1, http_code = ?, status_code = ?, last_error = ?, updated_at = datetime('now') WHERE id = ?");
foreach ($rows as $row) {
// Claim the row.
$pdo->beginTransaction();
try {
$claimStmt->execute([$row['id']]);
$claimed = $claimStmt->rowCount() > 0;
$pdo->commit();
} catch (Throwable $e) {
$pdo->rollBack();
continue;
}
if (!$claimed) {
continue; // taken by another worker
}
$processed++;
$url = (string) $row['url'];
$method = (string) $row['method'] === 'POST' ? 'POST' : 'GET';
$attempt = (int) $row['attempts'];
// SSRF re-check right before delivery (DNS may have changed since enqueue).
$parsedUrl = parse_url($url);
$host = $parsedUrl['host'] ?? '';
$ssrfBlocked = false;
$dnsFailed = false;
if ($host) {
// gethostbyname() returns the input unchanged both for an IP literal and on
// resolution failure, so do NOT skip the check when $ip === $host — that
// would wave through http://127.0.0.1/.
$ip = @gethostbyname($host);
if (filter_var($ip, FILTER_VALIDATE_IP)) {
// Resolved (or already an IP literal): block private/reserved targets.
$ssrfBlocked = filter_var($ip, FILTER_VALIDATE_IP, FILTER_FLAG_NO_PRIV_RANGE | FILTER_FLAG_NO_RES_RANGE) === false;
} else {
// Could not resolve. Likely transient — retry rather than kill the row.
$dnsFailed = true;
}
}
if ($dnsFailed) {
$nextAttempt = (int) $row['attempts'] + 1;
$errMsg = 'DNS resolution failed for ' . $host;
if ($nextAttempt >= PQ_MAX_ATTEMPTS) {
$deadStmt->execute([0, 0, $errMsg, $row['id']]);
$failed++;
} else {
$backoff = PQ_BACKOFF_SECONDS[min($nextAttempt - 1, count(PQ_BACKOFF_SECONDS) - 1)];
$retryStmt->execute(['+' . $backoff . ' seconds', 0, 0, $errMsg, $row['id']]);
$requeued++;
}
orbitraPqLog("postback #{$row['id']} $errMsg");
continue;
}
if ($ssrfBlocked) {
// Treat as permanent failure — do not retry a blocked target.
$deadStmt->execute([0, 0, 'SSRF: target resolves to a private/reserved IP', $row['id']]);
$failed++;
orbitraPqLog("postback #{$row['id']} SSRF-blocked -> failed");
continue;
}
// Perform the HTTP call.
$ch = curl_init($url);
curl_setopt($ch, CURLOPT_RETURNTRANSFER, true);
curl_setopt($ch, CURLOPT_TIMEOUT, PQ_HTTP_TIMEOUT);
// 4s was not enough for a proxied Meta call, and graph.facebook.com's AAAA
// records are unroutable from many hosts: curl tries v6, burns the whole
// connect budget waiting, and the row fails as "Resolving timed out" on a
// perfectly healthy box. Pin v4 and give a proxy handshake room to finish.
curl_setopt($ch, CURLOPT_CONNECTTIMEOUT, 10);
curl_setopt($ch, CURLOPT_IPRESOLVE, CURL_IPRESOLVE_V4);
curl_setopt($ch, CURLOPT_PROTOCOLS, CURLPROTO_HTTP | CURLPROTO_HTTPS);
curl_setopt($ch, CURLOPT_FOLLOWLOCATION, true);
curl_setopt($ch, CURLOPT_MAXREDIRS, 3);
// SSL verification: strict in production, relaxed for local development only
$isLocalEnv = defined('ORBITRA_ENV') && ORBITRA_ENV === 'local';
if ($isLocalEnv) {
// Local development: allow self-signed certs (e.g. localhost testing)
curl_setopt($ch, CURLOPT_SSL_VERIFYPEER, false);
curl_setopt($ch, CURLOPT_SSL_VERIFYHOST, 0);
} else {
// Production: enforce strict SSL verification
curl_setopt($ch, CURLOPT_SSL_VERIFYPEER, true);
curl_setopt($ch, CURLOPT_SSL_VERIFYHOST, 2);
}
// Some rows carry their own egress proxy (Facebook CAPI, when Meta blocks the
// tracker's own IP range). Applied before the body so a transport failure is
// attributed to the proxy, not to the payload.
$rowProxy = trim((string) ($row['proxy_url'] ?? ''));
if ($rowProxy !== '') {
$proxyParts = parse_url($rowProxy);
if (is_array($proxyParts) && !empty($proxyParts['host'])) {
$proxyScheme = strtolower($proxyParts['scheme'] ?? 'http');
$proxyType = CURLPROXY_HTTP;
if ($proxyScheme === 'socks5' || $proxyScheme === 'socks5h') {
$proxyType = CURLPROXY_SOCKS5_HOSTNAME;
} elseif ($proxyScheme === 'socks4') {
$proxyType = CURLPROXY_SOCKS4;
}
curl_setopt($ch, CURLOPT_PROXY, $proxyParts['host'] . (isset($proxyParts['port']) ? ':' . $proxyParts['port'] : ''));
curl_setopt($ch, CURLOPT_PROXYTYPE, $proxyType);
if (isset($proxyParts['user'])) {
curl_setopt($ch, CURLOPT_PROXYUSERPWD, urldecode($proxyParts['user']) . ':' . urldecode($proxyParts['pass'] ?? ''));
}
}
}
if ($method === 'POST') {
curl_setopt($ch, CURLOPT_POST, true);
$rawBody = (string) ($row['payload_json'] ?? '');
if ($rawBody !== '') {
// A row with a prepared body (Facebook Conversions API) is sent
// verbatim — its JSON structure is the message, and folding it into
// form fields would make Meta reject the event.
$contentType = trim((string) ($row['content_type'] ?? '')) ?: 'application/json';
curl_setopt($ch, CURLOPT_POSTFIELDS, $rawBody);
$requestHeaders = ['Content-Type: ' . $contentType];
$extraHeaders = json_decode((string) ($row['headers_json'] ?? ''), true);
if (is_array($extraHeaders)) {
foreach ($extraHeaders as $headerName => $headerValue) {
$headerName = trim((string) $headerName);
$headerValue = trim((string) $headerValue);
// Only simple HTTP token names are accepted; CR/LF are
// rejected so a stored credential cannot inject headers.
if ($headerName !== '' && preg_match('/^[A-Za-z0-9-]+$/', $headerName)
&& $headerValue !== '' && !preg_match('/[\r\n]/', $headerValue)) {
$requestHeaders[] = $headerName . ': ' . $headerValue;
}
}
}
curl_setopt($ch, CURLOPT_HTTPHEADER, $requestHeaders);
} else {
// Move query-string fields into the POST body so partners receive them in
// the request body (the common S2S convention) instead of an empty body.
$parsedForBody = parse_url($url);
parse_str($parsedForBody['query'] ?? '', $bodyFields);
if (!empty($bodyFields)) {
curl_setopt($ch, CURLOPT_POSTFIELDS, http_build_query($bodyFields));
curl_setopt($ch, CURLOPT_HTTPHEADER, ['Content-Type: application/x-www-form-urlencoded']);
}
}
}
$response = curl_exec($ch);
$httpCode = (int) curl_getinfo($ch, CURLINFO_HTTP_CODE);
$curlErr = curl_error($ch);
curl_close($ch);
$success = ($httpCode >= 200 && $httpCode < 400) && $curlErr === '';
if ($success) {
$doneStmt->execute([$httpCode, $httpCode, $row['id']]);
$delivered++;
orbitraPqLog("postback #{$row['id']} delivered (HTTP $httpCode)");
} else {
// Determine next state: retry or give up.
$nextAttempt = $attempt + 1;
$errMsg = $curlErr !== '' ? $curlErr : "HTTP $httpCode";
if ($nextAttempt >= PQ_MAX_ATTEMPTS) {
$deadStmt->execute([$httpCode, $httpCode, $errMsg, $row['id']]);
$failed++;
orbitraPqLog("postback #{$row['id']} FAILED after $nextAttempt attempts: $errMsg");
} else {
$backoff = PQ_BACKOFF_SECONDS[min($nextAttempt - 1, count(PQ_BACKOFF_SECONDS) - 1)];
$retryStmt->execute([
'+' . $backoff . ' seconds',
$httpCode,
$httpCode,
$errMsg,
$row['id'],
]);
$requeued++;
orbitraPqLog("postback #{$row['id']} retry #$nextAttempt in {$backoff}s: $errMsg");
}
}
}
// Health/state for the UI.
orbitraPqSetSetting($pdo, 'postback_queue_last_checked_at', $ts);
orbitraPqSetSetting($pdo, 'postback_queue_last_run_processed', (string) $processed);
orbitraPqSetSetting($pdo, 'postback_queue_last_run_delivered', (string) $delivered);
orbitraPqSetSetting($pdo, 'postback_queue_last_run_requeued', (string) $requeued);
orbitraPqSetSetting($pdo, 'postback_queue_last_run_failed', (string) $failed);
orbitraPqLog("postback_queue run: processed=$processed delivered=$delivered requeued=$requeued failed=$failed");
} catch (Throwable $e) {
orbitraPqSetSetting($pdo, 'postback_queue_last_error', $e->getMessage());
echo "[$ts] postback_queue error: " . $e->getMessage() . "\n";
} finally {
if (isset($fp) && is_resource($fp)) {
flock($fp, LOCK_UN);
fclose($fp);
}
}