Telegram Cron umgebaut
This commit is contained in:
@@ -0,0 +1,221 @@
|
||||
<?php
|
||||
declare(strict_types=1);
|
||||
|
||||
/**
|
||||
* Allgemeine Hilfsfunktionen fuer den Telegram-Cron.
|
||||
*
|
||||
* Enthalten sind hier Infrastruktur-Funktionen wie Cron-Lock, Laufprotokoll
|
||||
* ueber app_import_runs, Scheduling-Pruefung, State-Datei und Konsolen-Logging.
|
||||
* Diese Datei ist bewusst fachlich neutral und kennt keine Telegram- oder T-CrB-Details.
|
||||
*/
|
||||
|
||||
function ensureImportRunsTableExists(PDO $pdo): void
|
||||
{
|
||||
$pdo->exec("
|
||||
CREATE TABLE IF NOT EXISTS `app_import_runs` (
|
||||
`id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
|
||||
`source` varchar(32) NOT NULL,
|
||||
`source_url` varchar(500) DEFAULT NULL,
|
||||
`started_at` datetime NOT NULL,
|
||||
`finished_at` datetime DEFAULT NULL,
|
||||
`records_received` int(11) NOT NULL DEFAULT 0,
|
||||
`records_inserted` int(11) NOT NULL DEFAULT 0,
|
||||
`records_updated` int(11) NOT NULL DEFAULT 0,
|
||||
`status` varchar(32) NOT NULL DEFAULT 'success',
|
||||
`message` text DEFAULT NULL,
|
||||
PRIMARY KEY (`id`),
|
||||
KEY `idx_app_import_runs_source` (`source`),
|
||||
KEY `idx_app_import_runs_started_at` (`started_at`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
|
||||
");
|
||||
}
|
||||
|
||||
|
||||
function openCronLock(): mixed
|
||||
{
|
||||
$lockFile = sys_get_temp_dir() . DIRECTORY_SEPARATOR . 'skyview_telegram_cron.lock';
|
||||
$handle = fopen($lockFile, 'c+');
|
||||
if ($handle === false) {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (!flock($handle, LOCK_EX | LOCK_NB)) {
|
||||
fclose($handle);
|
||||
return null;
|
||||
}
|
||||
|
||||
return $handle;
|
||||
}
|
||||
|
||||
|
||||
function startImportRun(PDO $pdo, string $source, ?string $sourceUrl): int
|
||||
{
|
||||
$stmt = $pdo->prepare("
|
||||
INSERT INTO `app_import_runs` (
|
||||
`source`,
|
||||
`source_url`,
|
||||
`started_at`,
|
||||
`status`,
|
||||
`message`
|
||||
) VALUES (
|
||||
:source,
|
||||
:source_url,
|
||||
NOW(),
|
||||
'running',
|
||||
NULL
|
||||
)
|
||||
");
|
||||
$stmt->execute([
|
||||
':source' => $source,
|
||||
':source_url' => $sourceUrl,
|
||||
]);
|
||||
|
||||
return (int) $pdo->lastInsertId();
|
||||
}
|
||||
|
||||
|
||||
function finishImportRun(
|
||||
PDO $pdo,
|
||||
int $runId,
|
||||
int $recordsReceived,
|
||||
int $recordsInserted,
|
||||
int $recordsUpdated,
|
||||
string $status,
|
||||
?string $message
|
||||
): void {
|
||||
$stmt = $pdo->prepare("
|
||||
UPDATE `app_import_runs`
|
||||
SET
|
||||
`finished_at` = NOW(),
|
||||
`records_received` = :records_received,
|
||||
`records_inserted` = :records_inserted,
|
||||
`records_updated` = :records_updated,
|
||||
`status` = :status,
|
||||
`message` = :message
|
||||
WHERE `id` = :id
|
||||
LIMIT 1
|
||||
");
|
||||
$stmt->execute([
|
||||
':records_received' => $recordsReceived,
|
||||
':records_inserted' => $recordsInserted,
|
||||
':records_updated' => $recordsUpdated,
|
||||
':status' => $status,
|
||||
':message' => $message,
|
||||
':id' => $runId,
|
||||
]);
|
||||
}
|
||||
|
||||
|
||||
function loadCronState(string $stateFile): array
|
||||
{
|
||||
if (!is_file($stateFile)) {
|
||||
return [];
|
||||
}
|
||||
|
||||
$raw = @file_get_contents($stateFile);
|
||||
if ($raw === false || trim($raw) === '') {
|
||||
return [];
|
||||
}
|
||||
|
||||
$decoded = json_decode($raw, true);
|
||||
|
||||
return is_array($decoded) ? $decoded : [];
|
||||
}
|
||||
|
||||
|
||||
function saveCronState(string $stateFile, array $state): void
|
||||
{
|
||||
$json = json_encode($state, JSON_PRETTY_PRINT | JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES);
|
||||
if ($json === false) {
|
||||
throw new RuntimeException('Cron-Statusdatei konnte nicht als JSON gespeichert werden.');
|
||||
}
|
||||
|
||||
if (@file_put_contents($stateFile, $json . PHP_EOL, LOCK_EX) === false) {
|
||||
throw new RuntimeException('Cron-Statusdatei konnte nicht geschrieben werden: ' . $stateFile);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
function formatDecimal(mixed $value, int $decimals): string
|
||||
{
|
||||
if ($value === null || $value === '') {
|
||||
return '—';
|
||||
}
|
||||
|
||||
if (!is_numeric((string) $value)) {
|
||||
return (string) $value;
|
||||
}
|
||||
|
||||
return number_format((float) $value, $decimals, '.', '');
|
||||
}
|
||||
|
||||
|
||||
function isTaskDue(string $schedule, DateTimeImmutable $nowBerlin): bool
|
||||
{
|
||||
if ($schedule === 'every-run') {
|
||||
return true;
|
||||
}
|
||||
|
||||
if ($schedule === 'every-minute') {
|
||||
return true;
|
||||
}
|
||||
|
||||
if (preg_match('/^every-(\d+)-minutes$/', $schedule, $matches)) {
|
||||
$interval = max(1, (int) $matches[1]);
|
||||
return ((int) $nowBerlin->format('i')) % $interval === 0;
|
||||
}
|
||||
|
||||
if (preg_match('/^hourly-minute-(\d{2})$/', $schedule, $matches)) {
|
||||
return $nowBerlin->format('i') === $matches[1];
|
||||
}
|
||||
|
||||
if (preg_match('/^daily-(\d{2}):(\d{2})$/', $schedule, $matches)) {
|
||||
return $nowBerlin->format('H') === $matches[1]
|
||||
&& $nowBerlin->format('i') === $matches[2];
|
||||
}
|
||||
|
||||
if (preg_match('/^weekday-(mon|tue|wed|thu|fri|sat|sun)-(\d{2}):(\d{2})$/', $schedule, $matches)) {
|
||||
$weekdayMap = [
|
||||
'mon' => '1',
|
||||
'tue' => '2',
|
||||
'wed' => '3',
|
||||
'thu' => '4',
|
||||
'fri' => '5',
|
||||
'sat' => '6',
|
||||
'sun' => '7',
|
||||
];
|
||||
|
||||
return $nowBerlin->format('N') === ($weekdayMap[$matches[1]] ?? '')
|
||||
&& $nowBerlin->format('H') === $matches[2]
|
||||
&& $nowBerlin->format('i') === $matches[3];
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
|
||||
function writeMessage(string $message, bool $isError = false): void
|
||||
{
|
||||
$prefix = '[' . (new DateTimeImmutable('now'))->format('Y-m-d H:i:s') . '] ';
|
||||
$defaultStream = $isError ? 'php://stderr' : 'php://stdout';
|
||||
|
||||
if ($isError && defined('STDERR')) {
|
||||
fwrite(STDERR, $prefix . $message . PHP_EOL);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!$isError && defined('STDOUT')) {
|
||||
fwrite(STDOUT, $prefix . $message . PHP_EOL);
|
||||
return;
|
||||
}
|
||||
|
||||
$stream = @fopen($defaultStream, 'wb');
|
||||
if ($stream === false) {
|
||||
error_log($prefix . $message);
|
||||
return;
|
||||
}
|
||||
|
||||
fwrite($stream, $prefix . $message . PHP_EOL);
|
||||
fclose($stream);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user