Files
CloudVisu/api/cli/run_csv_exports.php
2026-06-22 15:52:25 +02:00

721 lines
22 KiB
PHP

#!/usr/bin/env php
<?php
declare(strict_types=1);
if (PHP_SAPI !== 'cli') {
fwrite(STDERR, "This script must be run via CLI.\n");
exit(1);
}
$rootDir = dirname(__DIR__, 2);
$configFile = $rootDir . '/config.php';
if (!is_file($configFile)) {
fwrite(STDERR, "config.php not found at {$configFile}\n");
exit(1);
}
require $configFile;
if (!isset($pdo) || !($pdo instanceof PDO)) {
fwrite(STDERR, "PDO connection \$pdo missing after loading config.php\n");
exit(1);
}
date_default_timezone_set(getenv('TZ') ?: 'Europe/Berlin');
if (!defined('VISU_INFLUX_URL')) {
define('VISU_INFLUX_URL', 'http://10.99.99.6:8086');
}
if (!defined('VISU_EXPORT_STORAGE_ROOT')) {
define('VISU_EXPORT_STORAGE_ROOT', rtrim((string)(getenv('VISU_EXPORT_ROOT') ?: '/var/portal_storage'), '/'));
}
function runnerLog(string $message): void
{
fwrite(STDOUT, '[' . date('Y-m-d H:i:s') . "] {$message}\n");
}
function runnerCsvEscape($value, string $delimiter = ';'): string
{
$text = (string)($value ?? '');
return preg_match('/["\r\n]/', $text) || str_contains($text, $delimiter)
? '"' . str_replace('"', '""', $text) . '"'
: $text;
}
function runnerCsvDelimiter($value): string
{
$value = (string)$value;
return in_array($value, [',', ';', "\t"], true) ? $value : ';';
}
function runnerTrendAggregate(string $aggregate): string
{
return in_array($aggregate, ['mean', 'min', 'max', 'last'], true) ? $aggregate : 'mean';
}
function runnerTrendRange(string $range): array
{
$ranges = [
'15m' => ['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.');