PDO::ERRMODE_EXCEPTION, PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC, ] ); ensureImportRunsTableExists($pdo); $lockHandle = openCronLock(); if ($lockHandle === null) { writeMessage('Telegram-Cron laeuft bereits. Dieser Start wird beendet.', true); exit(1); } $timezoneBerlin = new DateTimeZone('Europe/Berlin'); $timezoneUtc = new DateTimeZone('UTC'); $nowBerlin = new DateTimeImmutable('now', $timezoneBerlin); $nowUtc = $nowBerlin->setTimezone($timezoneUtc); $tasks = [ [ 'key' => 'cleanup_expired_link_tokens', 'label' => 'Abgelaufene Telegram-Verknuepfungscodes aufraeumen', 'enabled' => true, 'schedule' => 'every-15-minutes', 'handler' => static function (PDO $pdo, array $telegramConfig, DateTimeImmutable $nowBerlin): array { $stmt = $pdo->prepare(" UPDATE `app_users` SET `telegram_link_token` = NULL, `telegram_link_token_expires_at` = NULL, `updated_at` = NOW() WHERE `telegram_link_token` IS NOT NULL AND `telegram_link_token` <> '' AND `telegram_link_token_expires_at` IS NOT NULL AND `telegram_link_token_expires_at` < NOW() "); $stmt->execute(); return [ 'sent' => 0, 'updated' => (int) $stmt->rowCount(), 'message' => 'Abgelaufene Verbindungscodes bereinigt.', ]; }, ], [ 'key' => 'daily_morning_report_placeholder', 'label' => 'Morgendlichen Telegram-Bericht vorbereiten', 'enabled' => false, 'schedule' => 'daily-08:00', 'handler' => static function (PDO $pdo, array $telegramConfig, DateTimeImmutable $nowBerlin): array { return [ 'sent' => 0, 'updated' => 0, 'message' => 'Platzhalter fuer den spaeteren Morgenbericht.', ]; }, ], [ 'key' => 'daily_evening_report_placeholder', 'label' => 'Abendlichen Telegram-Bericht vorbereiten', 'enabled' => false, 'schedule' => 'daily-19:30', 'handler' => static function (PDO $pdo, array $telegramConfig, DateTimeImmutable $nowBerlin): array { return [ 'sent' => 0, 'updated' => 0, 'message' => 'Platzhalter fuer den spaeteren Abendbericht.', ]; }, ], [ 'key' => 'connected_users_heartbeat_placeholder', 'label' => 'Verbundene Telegram-Nutzer pruefen', 'enabled' => false, 'schedule' => 'hourly-minute-05', 'handler' => static function (PDO $pdo, array $telegramConfig, DateTimeImmutable $nowBerlin): array { $count = (int) $pdo->query(" SELECT COUNT(*) FROM `app_users` WHERE `telegram_chat_id` IS NOT NULL AND `telegram_chat_id` <> '' ")->fetchColumn(); return [ 'sent' => 0, 'updated' => 0, 'message' => 'Aktuell verbundene Telegram-Nutzer: ' . $count, ]; }, ], ]; $runId = startImportRun($pdo, IMPORT_RUN_SOURCE, null); $recordsReceived = count($tasks); $recordsInserted = 0; $recordsUpdated = 0; $messages = []; try { foreach ($tasks as $task) { if (($task['enabled'] ?? false) !== true) { $messages[] = sprintf('[%s] deaktiviert', (string) $task['key']); continue; } $schedule = (string) ($task['schedule'] ?? ''); if (!isTaskDue($schedule, $nowBerlin)) { $messages[] = sprintf('[%s] noch nicht faellig (%s)', (string) $task['key'], $schedule); continue; } $messages[] = sprintf('[%s] gestartet (%s)', (string) $task['key'], (string) ($task['label'] ?? $task['key'])); $handler = $task['handler'] ?? null; if (!is_callable($handler)) { throw new RuntimeException('Task-Handler ist nicht aufrufbar: ' . (string) $task['key']); } $result = $handler($pdo, $telegramConfig, $nowBerlin); $recordsInserted += (int) ($result['sent'] ?? 0); $recordsUpdated += (int) ($result['updated'] ?? 0); $messages[] = sprintf( '[%s] ok - %s', (string) $task['key'], trim((string) ($result['message'] ?? 'ohne Rueckmeldung')) ); } finishImportRun( $pdo, $runId, $recordsReceived, $recordsInserted, $recordsUpdated, 'success', implode("\n", $messages) ); writeMessage('Telegram-Cron abgeschlossen.'); foreach ($messages as $message) { writeMessage($message); } } catch (Throwable $exception) { finishImportRun( $pdo, $runId, $recordsReceived, $recordsInserted, $recordsUpdated, 'error', $exception->getMessage() ); writeMessage('Telegram-Cron mit Fehler beendet: ' . $exception->getMessage(), true); exit(1); } finally { flock($lockHandle, LOCK_UN); fclose($lockHandle); } 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 isTaskDue(string $schedule, DateTimeImmutable $nowBerlin): bool { 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); }