#!/usr/bin/env php ['sql' => '15m', 'seconds' => 900], '1h' => ['sql' => '1h', 'seconds' => 3600], '6h' => ['sql' => '6h', 'seconds' => 21600], '12h' => ['sql' => '12h', 'seconds' => 43200], '24h' => ['sql' => '24h', 'seconds' => 86400], '7d' => ['sql' => '7d', 'seconds' => 604800], '30d' => ['sql' => '30d', 'seconds' => 2592000], '90d' => ['sql' => '90d', 'seconds' => 7776000], '180d' => ['sql' => '180d', 'seconds' => 15552000], ]; return $ranges[$range] ?? $ranges['24h']; } function runnerTrendInterval(string $interval, int $rangeSeconds): string { $allowed = ['10s', '30s', '1m', '5m', '10m', '15m', '30m', '1h', '6h', '12h', '1d']; if ($interval !== 'auto' && in_array($interval, $allowed, true)) { return $interval; } if ($rangeSeconds <= 900) return '10s'; if ($rangeSeconds <= 3600) return '30s'; if ($rangeSeconds <= 21600) return '5m'; if ($rangeSeconds <= 86400) return '10m'; if ($rangeSeconds <= 604800) return '1h'; return '6h'; } function runnerInfluxIdentifier(string $value): string { return '"' . str_replace(['\\', '"'], ['\\\\', '\\"'], $value) . '"'; } function runnerInfluxString(string $value): string { return "'" . str_replace(['\\', "'"], ['\\\\', "\\'"], $value) . "'"; } function runnerInfluxQuery(string $database, string $query): array { $url = rtrim(VISU_INFLUX_URL, '/') . '/query?' . http_build_query([ 'db' => $database, 'q' => $query, ]); $ch = curl_init($url); curl_setopt($ch, CURLOPT_RETURNTRANSFER, true); curl_setopt($ch, CURLOPT_TIMEOUT, 30); $resp = curl_exec($ch); $err = curl_error($ch); $status = (int)curl_getinfo($ch, CURLINFO_RESPONSE_CODE); curl_close($ch); if ($err) { return ['ok' => false, 'error' => $err]; } $decoded = json_decode((string)$resp, true); if (!is_array($decoded)) { return ['ok' => false, 'error' => 'Influx returned no JSON', 'raw' => (string)$resp]; } if ($status >= 400) { return ['ok' => false, 'error' => $decoded['error'] ?? ('Influx HTTP ' . $status)]; } $result = $decoded['results'][0] ?? []; if (!empty($result['error'])) { return ['ok' => false, 'error' => (string)$result['error']]; } return ['ok' => true, 'result' => $result]; } function runnerInfluxSeriesValues(array $response): array { $series = $response['result']['series'][0] ?? null; if (!$series || empty($series['values']) || !is_array($series['values'])) { return []; } return $series['values']; } function runnerSelectionPointKey(array $point): string { return implode('|', [ 'influx', (string)($point['database'] ?? ''), (string)($point['measurement'] ?? ''), (string)($point['tag_key'] ?? ''), (string)($point['tag_value'] ?? ''), ]); } function runnerNormalizePoints(array $points): array { $out = []; foreach ($points as $point) { if (!is_array($point)) { continue; } $database = (string)($point['database'] ?? ''); $measurement = (string)($point['measurement'] ?? ''); $tagKey = (string)($point['tag_key'] ?? ''); $tagValue = (string)($point['tag_value'] ?? ''); if ($database === '' || $measurement === '' || $tagKey === '' || $tagValue === '') { continue; } $point['label'] = (string)($point['tag_value'] ?? $point['label'] ?? ''); $out[] = $point; } return array_values($out); } function runnerLoadTrendSeries(array $selectionPoints, string $range, string $interval, string $aggregate): array { $rangeInfo = runnerTrendRange($range); $bucket = runnerTrendInterval($interval, (int)$rangeInfo['seconds']); $agg = runnerTrendAggregate($aggregate); $result = []; foreach (array_chunk($selectionPoints, 8) as $group) { foreach ($group as $index => $point) { $query = sprintf( 'SELECT %s("value") FROM %s WHERE %s = %s AND time >= now() - %s GROUP BY time(%s) fill(null)', $agg, runnerInfluxIdentifier((string)$point['measurement']), runnerInfluxIdentifier((string)$point['tag_key']), runnerInfluxString((string)$point['tag_value']), $rangeInfo['sql'], $bucket ); $response = runnerInfluxQuery((string)$point['database'], $query); $values = []; if (!empty($response['ok'])) { foreach (runnerInfluxSeriesValues($response) as $row) { $values[] = [ 'time' => (string)($row[0] ?? ''), 'value' => $row[1] ?? '', ]; } } $result[] = [ 'id' => runnerSelectionPointKey($point) ?: ('series-' . $index), 'label' => (string)($point['tag_value'] ?? $point['label'] ?? ''), 'values' => $values, 'error' => empty($response['ok']) ? (string)($response['error'] ?? 'Influx error') : '', ]; } } return $result; } function runnerLoadLastValues(array $selectionPoints): array { $result = []; foreach (array_chunk($selectionPoints, 64) as $group) { foreach ($group as $index => $point) { $query = sprintf( 'SELECT last("value") FROM %s WHERE %s = %s', runnerInfluxIdentifier((string)$point['measurement']), runnerInfluxIdentifier((string)$point['tag_key']), runnerInfluxString((string)$point['tag_value']) ); $response = runnerInfluxQuery((string)$point['database'], $query); $rows = !empty($response['ok']) ? runnerInfluxSeriesValues($response) : []; $time = (string)($rows[0][0] ?? ''); $value = $rows[0][1] ?? ''; $result[runnerSelectionPointKey($point) ?: ('series-' . $index)] = [ 'time' => $time, 'value' => $value, 'error' => empty($response['ok']) ? (string)($response['error'] ?? 'Influx error') : '', ]; } } return $result; } function runnerBuildTrendCsv(array $selectionPoints, array $series, string $delimiter): array { $rows = []; $headers = ['Zeit']; foreach ($selectionPoints as $point) { $headers[] = (string)($point['tag_value'] ?? $point['label'] ?? ''); } foreach ($series as $serie) { foreach (($serie['values'] ?? []) as $row) { $time = (string)($row['time'] ?? ''); if ($time === '') { continue; } if (!isset($rows[$time])) { $rows[$time] = ['Zeit' => $time]; } $rows[$time][(string)($serie['label'] ?? '')] = $row['value'] ?? ''; } } ksort($rows, SORT_STRING); $lines = [ implode($delimiter, array_map(static fn($cell) => runnerCsvEscape($cell, $delimiter), $headers)), ]; foreach ($rows as $row) { $lines[] = implode($delimiter, array_map( static fn($cell) => runnerCsvEscape($row[$cell] ?? '', $delimiter), $headers )); } return [ 'content' => "\xEF\xBB\xBF" . implode("\r\n", $lines), 'row_count' => count($rows), ]; } function runnerBuildLastValuesCsv(array $selectionPoints, array $lastValues, string $delimiter): array { $header = ['Datenpunkt', 'Measurement', 'Tag', 'Wert', 'Zeit', 'Datenbank']; $lines = [ implode($delimiter, array_map(static fn($cell) => runnerCsvEscape($cell, $delimiter), $header)), ]; foreach ($selectionPoints as $point) { $last = $lastValues[runnerSelectionPointKey($point)] ?? []; $lines[] = implode($delimiter, array_map( static fn($cell) => runnerCsvEscape($cell, $delimiter), [ (string)($point['tag_value'] ?? ''), (string)($point['measurement'] ?? ''), (string)($point['tag_key'] ?? ''), $last['value'] ?? '', $last['time'] ?? '', (string)($point['database'] ?? ''), ] )); } return [ 'content' => "\xEF\xBB\xBF" . implode("\r\n", $lines), 'row_count' => count($selectionPoints), ]; } function runnerScheduleKind(string $value): string { return in_array($value, ['interval', 'daily', 'weekly', 'monthly'], true) ? $value : 'daily'; } function runnerTimeOfDay(string $value): string { $value = trim($value); if (!preg_match('/^(\d{2}):(\d{2})$/', $value, $m)) { return '00:00'; } return sprintf('%02d:%02d', max(0, min(23, (int)$m[1])), max(0, min(59, (int)$m[2]))); } function runnerNextRunAt(string $scheduleKind, int $intervalMinutes, string $timeOfDay, string $weekdays, int $monthDay, ?DateTimeImmutable $now = null): string { $now = $now ?: new DateTimeImmutable('now'); [$hour, $minute] = array_map('intval', explode(':', runnerTimeOfDay($timeOfDay))); $scheduleKind = runnerScheduleKind($scheduleKind); if ($scheduleKind === 'interval') { return $now->modify('+' . max(1, $intervalMinutes) . ' minutes')->format('Y-m-d H:i:s'); } if ($scheduleKind === 'daily') { $candidate = $now->setTime($hour, $minute, 0); if ($candidate <= $now) { $candidate = $candidate->modify('+1 day'); } return $candidate->format('Y-m-d H:i:s'); } if ($scheduleKind === 'weekly') { $days = array_filter(array_map('intval', explode(',', (string)$weekdays)), static fn($day) => $day >= 1 && $day <= 7); if (empty($days)) { $days = [(int)$now->format('N')]; } for ($offset = 0; $offset <= 14; $offset++) { $candidateDay = $now->modify('+' . $offset . ' days'); if (!in_array((int)$candidateDay->format('N'), $days, true)) { continue; } $candidate = $candidateDay->setTime($hour, $minute, 0); if ($candidate > $now) { return $candidate->format('Y-m-d H:i:s'); } } } if ($scheduleKind === 'monthly') { $monthDay = max(1, min(31, $monthDay)); $base = $now; for ($i = 0; $i < 14; $i++) { $monthStart = $base->modify('first day of this month')->setTime($hour, $minute, 0); $daysInMonth = (int)$monthStart->format('t'); $day = min($monthDay, $daysInMonth); $candidate = $monthStart->setDate( (int)$monthStart->format('Y'), (int)$monthStart->format('m'), $day ); if ($candidate > $now) { return $candidate->format('Y-m-d H:i:s'); } $base = $base->modify('first day of next month'); } } return $now->modify('+1 day')->setTime($hour, $minute, 0)->format('Y-m-d H:i:s'); } function runnerEnsureDir(string $dir): void { if (!is_dir($dir) && !@mkdir($dir, 0775, true) && !is_dir($dir)) { throw new RuntimeException('Could not create export directory: ' . $dir); } if (!is_writable($dir)) { throw new RuntimeException('Export directory is not writable: ' . $dir); } } function runnerSanitizeFilePart(string $value): string { $value = trim($value); $value = preg_replace('/[^a-zA-Z0-9._-]+/', '_', $value); $value = trim((string)$value, '._-'); return $value !== '' ? $value : 'export'; } function runnerSanitizeTargetDir(string $value): string { $value = str_replace('\\', '/', trim($value)); $parts = array_filter(explode('/', $value), static fn($part) => $part !== '' && $part !== '.' && $part !== '..'); $clean = []; foreach ($parts as $part) { $part = runnerSanitizeFilePart($part); if ($part !== '') { $clean[] = $part; } } return !empty($clean) ? implode('/', $clean) : 'csv_exports'; } function runnerBuildFilename(string $pattern, string $selectionName, string $jobName, int $jobId, ?DateTimeImmutable $now = null): string { $now = $now ?: new DateTimeImmutable('now'); $selection = runnerSanitizeFilePart($selectionName); $job = runnerSanitizeFilePart($jobName); $replacements = [ '{selection}' => $selection, '{job}' => $job, '{job_id}' => (string)$jobId, '{yyyy}' => $now->format('Y'), '{mm}' => $now->format('m'), '{dd}' => $now->format('d'), '{hh}' => $now->format('H'), '{ii}' => $now->format('i'), '{ss}' => $now->format('s'), ]; $name = strtr($pattern !== '' ? $pattern : '{selection}_{yyyy}-{mm}-{dd}_{hh}-{ii}-{ss}.csv', $replacements); $name = runnerSanitizeFilePart($name); if (!str_ends_with(strtolower($name), '.csv')) { $name .= '.csv'; } return $name; } function runnerCreateFileRecord(PDO $pdo, array $job, string $fileName, string $relativePath): int { $stmt = $pdo->prepare(" INSERT INTO visu_csv_export_files ( job_id, project_id, selection_id, status, file_name, file_path, export_started_at ) VALUES (?, ?, ?, 'running', ?, ?, NOW()) "); $stmt->execute([ (int)$job['id'], (int)$job['project_id'], (int)$job['selection_id'], $fileName, $relativePath, ]); return (int)$pdo->lastInsertId(); } function runnerMarkFileSuccess(PDO $pdo, int $fileId, int $fileSize, int $rowCount): void { $stmt = $pdo->prepare(" UPDATE visu_csv_export_files SET status = 'success', file_size = ?, row_count = ?, export_finished_at = NOW(), error_message = NULL WHERE id = ? "); $stmt->execute([$fileSize, $rowCount, $fileId]); } function runnerMarkFileError(PDO $pdo, int $fileId, string $error): void { $stmt = $pdo->prepare(" UPDATE visu_csv_export_files SET status = 'error', export_finished_at = NOW(), error_message = ? WHERE id = ? "); $stmt->execute([mb_substr($error, 0, 65535), $fileId]); } function runnerMarkJobSuccess(PDO $pdo, array $job, string $nextRunAt): void { $stmt = $pdo->prepare(" UPDATE visu_csv_export_jobs SET last_run_at = NOW(), next_run_at = ?, last_status = 'success', last_error = NULL WHERE id = ? "); $stmt->execute([$nextRunAt, (int)$job['id']]); } function runnerMarkJobError(PDO $pdo, array $job, string $error, string $nextRunAt): void { $stmt = $pdo->prepare(" UPDATE visu_csv_export_jobs SET last_run_at = NOW(), next_run_at = ?, last_status = 'error', last_error = ? WHERE id = ? "); $stmt->execute([$nextRunAt, mb_substr($error, 0, 65535), (int)$job['id']]); } function runnerClaimJob(PDO $pdo, int $jobId): bool { $stmt = $pdo->prepare(" UPDATE visu_csv_export_jobs SET last_status = 'running', last_error = NULL WHERE id = ? AND enabled = 1 AND (last_status <> 'running' OR last_status IS NULL) "); $stmt->execute([$jobId]); return $stmt->rowCount() > 0; } function runnerCleanupOldFiles(PDO $pdo, array $job): void { $retentionDays = max(1, (int)($job['retention_days'] ?? 30)); $stmt = $pdo->prepare(" SELECT id, file_path FROM visu_csv_export_files WHERE job_id = ? AND status = 'success' AND created_at < DATE_SUB(NOW(), INTERVAL ? DAY) "); $stmt->execute([(int)$job['id'], $retentionDays]); $rows = $stmt->fetchAll(PDO::FETCH_ASSOC); foreach ($rows as $row) { $relativePath = (string)($row['file_path'] ?? ''); $absolutePath = rtrim(VISU_EXPORT_STORAGE_ROOT, '/') . '/' . ltrim($relativePath, '/'); if (is_file($absolutePath)) { @unlink($absolutePath); } $upd = $pdo->prepare(" UPDATE visu_csv_export_files SET status = 'deleted' WHERE id = ? "); $upd->execute([(int)$row['id']]); } } function runnerLoadDueJobs(PDO $pdo): array { $stmt = $pdo->query(" SELECT * FROM visu_csv_export_jobs WHERE enabled = 1 AND next_run_at IS NOT NULL AND next_run_at <= NOW() AND (last_status <> 'running' OR last_status IS NULL) ORDER BY next_run_at ASC, id ASC LIMIT 25 "); return $stmt->fetchAll(PDO::FETCH_ASSOC) ?: []; } function runnerLoadSelection(PDO $pdo, int $selectionId): ?array { $stmt = $pdo->prepare(" SELECT id, name, database_name, range_value, interval_value, aggregate_value, delimiter_value, points_json FROM visu_monitoring_selections WHERE id = ? AND is_active = 1 LIMIT 1 "); $stmt->execute([$selectionId]); $row = $stmt->fetch(PDO::FETCH_ASSOC); return $row ?: null; } try { $pdo->query('SELECT 1 FROM visu_csv_export_jobs LIMIT 1'); $pdo->query('SELECT 1 FROM visu_csv_export_files LIMIT 1'); } catch (Throwable $e) { runnerLog('CSV export tables missing, skipping run.'); exit(0); } $jobs = runnerLoadDueJobs($pdo); if (empty($jobs)) { runnerLog('No due CSV export jobs.'); exit(0); } runnerLog('Found ' . count($jobs) . ' due CSV export job(s).'); foreach ($jobs as $job) { $jobId = (int)($job['id'] ?? 0); if ($jobId <= 0) { continue; } if (!runnerClaimJob($pdo, $jobId)) { continue; } $selection = runnerLoadSelection($pdo, (int)($job['selection_id'] ?? 0)); $nextRunAt = runnerNextRunAt( (string)($job['schedule_kind'] ?? 'daily'), (int)($job['interval_minutes'] ?? 60), (string)($job['time_of_day'] ?? '00:00'), (string)($job['weekdays'] ?? ''), (int)($job['month_day'] ?? 1) ); try { if (!$selection) { throw new RuntimeException('Selection not found or inactive.'); } $points = json_decode((string)($selection['points_json'] ?? '[]'), true); if (!is_array($points)) { throw new RuntimeException('Selection points JSON is invalid.'); } $points = runnerNormalizePoints($points); if (empty($points)) { throw new RuntimeException('Selection contains no valid Influx points.'); } $delimiter = runnerCsvDelimiter($job['csv_delimiter'] ?? ';'); $exportKind = (string)($job['export_kind'] ?? 'trend_csv'); $selectionName = (string)($selection['name'] ?? ('selection_' . (int)$job['selection_id'])); $jobName = (string)($job['name'] ?? ('job_' . $jobId)); $targetDir = runnerSanitizeTargetDir((string)($job['target_dir'] ?? 'csv_exports')); $fileName = runnerBuildFilename((string)($job['filename_pattern'] ?? ''), $selectionName, $jobName, $jobId); $relativePath = trim($targetDir . '/' . $fileName, '/'); $absoluteDir = rtrim(VISU_EXPORT_STORAGE_ROOT, '/') . '/' . $targetDir; $absolutePath = rtrim(VISU_EXPORT_STORAGE_ROOT, '/') . '/' . $relativePath; runnerEnsureDir($absoluteDir); $fileId = runnerCreateFileRecord($pdo, $job, $fileName, $relativePath); if ($exportKind === 'lastvalues_csv') { $lastValues = runnerLoadLastValues($points); $csv = runnerBuildLastValuesCsv($points, $lastValues, $delimiter); } else { $series = runnerLoadTrendSeries( $points, (string)($selection['range_value'] ?? '24h'), (string)($selection['interval_value'] ?? 'auto'), (string)($selection['aggregate_value'] ?? 'last') ); $csv = runnerBuildTrendCsv($points, $series, $delimiter); } if (file_put_contents($absolutePath, $csv['content']) === false) { throw new RuntimeException('CSV file could not be written: ' . $absolutePath); } @chmod($absolutePath, 0664); runnerMarkFileSuccess($pdo, $fileId, (int)filesize($absolutePath), (int)$csv['row_count']); runnerMarkJobSuccess($pdo, $job, $nextRunAt); runnerCleanupOldFiles($pdo, $job); runnerLog("Job {$jobId} finished successfully: {$relativePath}"); } catch (Throwable $e) { $error = $e->getMessage(); if (!empty($fileId ?? 0)) { runnerMarkFileError($pdo, (int)$fileId, $error); } runnerMarkJobError($pdo, $job, $error, $nextRunAt); runnerLog("Job {$jobId} failed: {$error}"); } } runnerLog('CSV export runner finished.');