Files
Enelix-EMS/libs/NetzfahrplanV4Messaufnahme.php
T

166 lines
9.4 KiB
PHP

<?php
declare(strict_types=1);
namespace Belevo\EnelixEMS;
use InvalidArgumentException;
use RuntimeException;
use Throwable;
require_once __DIR__ . '/NetzfahrplanV4Bilanzierung.php';
/** Passive, versioned RAW capture. Never grants training/control eligibility.
* Reader supplies numeric values only; no configuration, credentials or network.
*/
final class NetzfahrplanV4Messaufnahme
{
public static function configuration(array $c): array
{
if (($c['version'] ?? null) !== 1 || !is_string($c['installationId'] ?? null)
|| !preg_match('/^[0-9a-f-]{36}$/D', $c['installationId'])
|| !is_array($c['extraSources'] ?? null) || count($c['extraSources']) > 60
|| !is_string($c['reportedInventorySha256'] ?? null)
|| !preg_match('/^[0-9a-f]{64}$/D', $c['reportedInventorySha256'])) {
throw new InvalidArgumentException('Invalid capture configuration');
}
$c['accounting'] = NetzfahrplanV4Bilanzierung::konfigurieren($c['accounting'] ?? []);
$keys = []; $ids = [];
foreach ($c['accounting']['sources'] as $s) { $keys[$s['key']] = true; $ids[$s['variableId']] = true; }
foreach ($c['extraSources'] as $s) {
if (count($s) !== 5 || !is_string($s['key'] ?? null)
|| !preg_match('/^[A-Za-z][A-Za-z0-9_-]{0,63}$/D', $s['key']) || isset($keys[$s['key']])
|| !is_int($s['variableId'] ?? null) || $s['variableId'] < 1 || isset($ids[$s['variableId']])
|| !is_int($s['parentId'] ?? null) || $s['parentId'] < 1
|| !is_string($s['ident'] ?? null) || strlen($s['ident']) > 100
|| !in_array($s['unit'] ?? null, ['W','kWh','percent'], true)) {
throw new InvalidArgumentException('Invalid or duplicate raw source');
}
$keys[$s['key']] = true; $ids[$s['variableId']] = true;
}
usort($c['extraSources'], static fn(array $a, array $b): int => strcmp($a['key'], $b['key']));
return $c;
}
private static function readSource(array $s, callable $read): array
{
try {
$r = $read($s['variableId']);
if (!is_array($r)) throw new RuntimeException('Missing reading');
} catch (Throwable $e) {
return ['value' => null, 'sourceUpdatedAt' => null, 'sourceChangedAt' => null, 'issues' => ['read_failed']];
}
$issues = [];
$v = $r['value'] ?? null;
if ((!is_int($v) && !is_float($v)) || !is_finite((float) $v) || abs((float)$v) > 1.0e12) {
$v = null; $issues[] = 'invalid_numeric_value';
}
$at = $r['updated'] ?? null;
if (!is_int($at) || $at < 1) { $at = null; $issues[] = 'missing_original_timestamp'; }
$changed = $r['changed'] ?? null;
if (!is_int($changed) || $changed < 0) $changed = null;
if (($r['parentID'] ?? null) !== $s['parentId'] || ($r['ident'] ?? null) !== $s['ident']) {
// Do not record a numeric value from an object re-used for a different purpose.
$v = null; $issues[] = 'source_identity_mismatch';
}
return ['value' => $v, 'sourceUpdatedAt' => $at, 'sourceChangedAt' => $changed, 'issues' => $issues];
}
public static function capture(array $config, callable $read, callable $clock): array
{
$c = self::configuration($config);
$start = $clock();
if (!is_int($start) || $start <= 0) throw new InvalidArgumentException('Integer UTC epoch required');
$sources = array_merge($c['accounting']['sources'], $c['extraSources']);
$first = []; $second = []; $raw = []; $issues = [];
foreach ($sources as $s) $first[$s['key']] = self::readSource($s, $read);
foreach ($sources as $s) $second[$s['key']] = self::readSource($s, $read);
$end = $clock();
if (!is_int($end) || $end < $start) throw new RuntimeException('Acquisition clock moved backwards');
foreach ($sources as $s) {
$key = $s['key']; $r = $first[$key];
if ($r !== $second[$key]) {
$r['issues'][] = 'changed_during_capture';
$r['secondValue'] = $second[$key]['value'];
$r['secondUpdatedAt'] = $second[$key]['sourceUpdatedAt'];
}
if ($r['sourceUpdatedAt'] !== null && $r['sourceUpdatedAt'] > $end) $r['issues'][] = 'future_source';
$r['sourceAgeSeconds'] = $r['sourceUpdatedAt'] === null ? null : $end - $r['sourceUpdatedAt'];
$raw[$key] = ['variableId' => $s['variableId']] + $r;
foreach ($r['issues'] as $issue) $issues[] = $key . ':' . $issue;
}
$assessment = null; $accountingUsable = true; $samples = [];
foreach ($c['accounting']['sources'] as $s) {
$r = $raw[$s['key']];
if ($r['issues'] !== []) $accountingUsable = false;
$samples[$s['variableId']] = ['value' => $r['value'], 'updated' => $r['sourceUpdatedAt'], 'parentID' => $s['parentId'], 'ident' => $s['ident']];
}
if ($accountingUsable) {
try { $assessment = NetzfahrplanV4Bilanzierung::aufnehmen($c['accounting'], static fn(int $id): array => $samples[$id], $end); }
catch (Throwable $e) { $issues[] = 'accounting_input_invalid'; }
} else { $issues[] = 'accounting_snapshot_unusable'; }
if ($end - $start > 5) $issues[] = 'long_acquisition';
return ['schemaVersion' => 1, 'kind' => 'raw_accounting_capture', 'installationId' => $c['installationId'],
'captureStartedAt' => gmdate('c', $start), 'capturedAt' => gmdate('c', $end),
'mappingSha256' => hash('sha256', json_encode($c, JSON_THROW_ON_ERROR | JSON_PRESERVE_ZERO_FRACTION)),
'reportedInventorySha256' => $c['reportedInventorySha256'],
'timestampMeaning' => 'Symcon VariableUpdated; not proof of an independent device timestamp',
'raw' => $raw, 'issues' => $issues, 'assessment' => $assessment,
'trainingEligible' => false, 'controlEligible' => false, 'measurementBoundaryVerified' => false,
'futureSdl' => 'unknown_no_assumption', 'historicalDataRewritten' => false];
}
/** Append-only daily JSONL with exclusive lock, bounds and crash-tail detection.
* No retention deletion; quota stops capture visibly, never deletes old histories.
*/
public static function append(string $directory, array $record, int $quota = 536870912): string
{
if (!is_dir($directory) || is_link($directory) || realpath($directory) !== $directory || $quota < 1) {
throw new RuntimeException('Invalid dedicated data directory');
}
$json = json_encode($record, JSON_THROW_ON_ERROR | JSON_PRESERVE_ZERO_FRACTION) . "\n";
if (strlen($json) > 262144) throw new RuntimeException('Capture exceeds record limit');
$epoch = strtotime($record['capturedAt'] ?? '');
if ($epoch === false) throw new RuntimeException('Capture timestamp missing');
$file = $directory . '/raw-' . gmdate('Ymd', $epoch) . '.jsonl';
$lockPath = $directory . '/.writer.lock';
if (is_link($lockPath) || is_link($file)) throw new RuntimeException('Symbolic link refused');
$lock = fopen($lockPath, 'c+b');
if ($lock === false) throw new RuntimeException('Capture lock unavailable');
try {
if (!flock($lock, LOCK_EX | LOCK_NB)) throw new RuntimeException('Capture writer busy');
clearstatcache(); $used = 0; $files = glob($directory . '/raw-*.jsonl');
if ($files === false || count($files) > 400) throw new RuntimeException('Data directory bounds exceeded');
foreach ($files as $p) {
if (is_link($p)) throw new RuntimeException('Symbolic journal refused');
$size = filesize($p); if ($size === false) throw new RuntimeException('Cannot size journal');
$used += $size;
}
if ($used + strlen($json) > $quota) throw new RuntimeException('Capture quota reached; no data deleted');
$free = disk_free_space($directory);
if ($free === false || $free < 104857600 + strlen($json)) throw new RuntimeException('Insufficient free disk; capture paused');
$h = fopen($file, 'c+b');
if ($h === false) throw new RuntimeException('Cannot open capture journal');
try {
fseek($h, 0, SEEK_END); $old = ftell($h);
if ($old === false || $old + strlen($json) > 67108864) throw new RuntimeException('Daily capture bounds exceeded');
if ($old > 0) {
fseek($h, -1, SEEK_END);
if (fread($h, 1) !== "\n") throw new RuntimeException('Incomplete journal tail; preserve for review');
}
fseek($h, 0, SEEK_END); $n = 0;
try {
while ($n < strlen($json)) {
$written = fwrite($h, substr($json, $n));
if ($written === false || $written === 0) throw new RuntimeException('Capture write incomplete');
$n += $written;
}
if (!fflush($h)) throw new RuntimeException('Capture flush failed');
} catch (Throwable $e) { ftruncate($h, $old); fflush($h); throw $e; }
if (!chmod($file, 0640)) throw new RuntimeException('Cannot restrict journal permissions');
} finally { fclose($h); }
return $file;
} finally { flock($lock, LOCK_UN); fclose($lock); }
}
}