import { InfluxDB } from '@influxdata/influxdb-client'; import { ResolvedRange, windowStartBefore } from '@/lib/ranges'; import { VerbrauchPoint } from '@/types/strom'; const URL = process.env.INFLUX_URL || 'http://nuccy:8086'; const TOKEN = process.env.INFLUX_TOKEN || ''; const ORG = process.env.INFLUX_ORG || 'citysensor'; const BUCKET = process.env.INFLUX_BUCKET || 'strom'; const FIELD = process.env.INFLUX_FIELD || 'arbeit'; // Rohdaten (1-Sekunden-Zaehlerstand) — fuer 24h, immer aktuell const MEAS_RAW = process.env.INFLUX_MEASUREMENT || 'vzlogger'; // Stuendliche Snapshots (Downsampling-Task) — fuer 7d/31d/365d, schnell const MEAS_ROLLUP = process.env.INFLUX_ROLLUP_MEASUREMENT || 'arbeit_hourly'; const UNIT_FACTOR = parseFloat(process.env.INFLUX_UNIT_FACTOR || '0.001'); // *Faktor -> kWh (Wh) const TZ = process.env.INFLUX_TZ || 'Europe/Berlin'; const influx = new InfluxDB({ url: URL, token: TOKEN, timeout: 90_000 }); /** * Liefert den Verbrauch je Fenster aus dem kumulativen Zaehlerstand. * * Quelle 'raw' (24h): hoechstaufgeloeste Daten, aggregateWindow(last) holt den * Zaehlerstand am Fensterende, difference() bildet den Zuwachs = Verbrauch. * * Quelle 'rollup' (7d/31d/365d): stuendliche Snapshots des Zaehlerstands. * difference() ergibt den stuendlichen Verbrauch, der dann tages-/stundenweise * summiert wird. Tagesgrenzen sind ueber `location` zeitzonen-korrekt. * * timeShift(-every) verschiebt den Zeitstempel auf den Fenster-ANFANG, sodass ein * Balken den Zeitraum [T, T+every) repraesentiert. Die Abfrage startet ein Fenster * frueher, damit das erste sichtbare Fenster bereits einen Differenzwert hat. */ export async function fetchVerbrauch(r: ResolvedRange): Promise { const queryStart = windowStartBefore(r.start, r.every); // ein Fenster Vorlauf const measurement = r.source === 'raw' ? MEAS_RAW : MEAS_ROLLUP; const aggregation = r.source === 'raw' ? // Rohdaten: letzter Zaehlerstand je Fenster, dann Differenz `|> aggregateWindow(every: ${r.every}, fn: last, createEmpty: false) |> difference(nonNegative: true)` : // Rollup: stuendliche Differenz, dann je Fenster summieren `|> difference(nonNegative: true) |> aggregateWindow(every: ${r.every}, fn: sum, createEmpty: false)`; const flux = ` import "timezone" option location = timezone.location(name: "${TZ}") from(bucket: "${BUCKET}") |> range(start: ${queryStart.toISOString()}, stop: ${r.stop.toISOString()}) |> filter(fn: (r) => r._measurement == "${measurement}") |> filter(fn: (r) => r._field == "${FIELD}") ${aggregation} |> timeShift(duration: -${r.every}) |> sort(columns: ["_time"], desc: false) |> keep(columns: ["_time", "_value"]) `; const queryApi = influx.getQueryApi(ORG); const points: VerbrauchPoint[] = []; const minTs = r.start.getTime(); const maxTs = r.stop.getTime(); return new Promise((resolve, reject) => { queryApi.queryRows(flux, { next(row, tableMeta) { const o = tableMeta.toObject(row); const ts = new Date(o._time as string).getTime(); if (ts < minTs || ts >= maxTs) return; // Vorlauf-/Ueberhangfenster verwerfen const value = Number(o._value); if (!Number.isFinite(value)) return; points.push([ts, value * UNIT_FACTOR]); }, error(err) { reject(err); }, complete() { resolve(points); }, }); }); }