Compare commits

..

2 Commits

Author SHA1 Message Date
admin 21050f51c9 compare-laermwerte: Testprogramm zum Vergleich von Laermwerten zweier API-Instanzen
Ruft dieselben Parameter (sensorid/data/span/datetime/peak) gegen zwei
"getsensordata"-Endpunkte ab und diffed die Werte zeitstempelweise, generisch
über alle Felder (funktioniert für live genauso wie für havg/davg/daynight/lden).
Diente als Verifikation der VictoriaMetrics-Anbindung: live-Werte stimmen exakt,
havg-Werte liegen weit innerhalb der Toleranz gegenüber der Influx-Produktion.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-21 10:59:56 +02:00
admin 23fe7d2ed2 VictoriaMetrics als zusätzliche Datenbank-Option für Messwerte ergänzt (STORE/DBASE=victoria)
Dritte, zu mongo/influx exklusive Auswahl für die laufenden Messwerte. Schreibpfad
nutzt das bestehende Influx-Line-Protocol unverändert (common/victoria_post.js);
Lesepfad (sensorapi/databases/victoria.js + victoria2json.js) holt Rohdaten per
VictoriaMetrics' /api/v1/export und bucketet/aggregiert stundenweise clientseitig,
nach Mongo-Konvention (Stunden-Start als Label, kein Zeit-Shift nötig wie bei Influx).
Scope bewusst auf die schon heute per DBASE umschaltbaren Funktionen begrenzt
(getActData/getNoiseAVGData) - getAvgData/getLongAvg/getGeigerData bleiben wie bisher.

Docker-Compose um victoriametrics-Service ergänzt (Retention explizit auf 100y
gesetzt, da VictoriaMetrics sonst nach 1 Monat Daten löscht).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-20 21:31:06 +02:00
15 changed files with 932 additions and 14 deletions
+6
View File
@@ -9,3 +9,9 @@ log/
# Produktionsstand vom Server, nur Referenz
noisesensors/vonRemote/
# Laufzeitdaten der lokal getesteten Container (DB-Dateien, aktuelle Messwerte)
noisesensors/data/mongo/data/
noisesensors/data/influx/
noisesensors/data/victoria/
noisesensors/data/aktdata.json
+16 -2
View File
@@ -20,6 +20,8 @@ Die Datei **noise.tgz** enthält die folgenden Dateien und Verzeichnisse:
| | | +- <f>create.js
| + <d>influx
| | + <d>data
| + <d>victoria
| | + <d>data
<f>: file, <d>: directory
~~~
@@ -42,6 +44,10 @@ hier liegt die Datei **create.js**, mit deren Hilfe bei ersten Start der Datenba
Verzeichnis für die Influx-Datenbank
* **data/influx/data**
Hier dann die eingelesenen Daten
* **data/victoria**
Verzeichnis für die VictoriaMetrics-Datenbank (optionale Alternative zu Influx, siehe unten)
* **data/victoria/data**
Hier dann die eingelesenen Daten
### Aufrufe
@@ -60,6 +66,8 @@ STORE=influx
~~~
\<username>, \<passwort>, \<username1>, \<passwort1> und \<token> sind anzupassen.
**STORE** legt fest, in welche Datenbank(en) *readin* die laufenden Messwerte schreibt: `influx`, `mongo`, `victoria` oder `both` (= mongo **und** influx gleichzeitig; Default). `victoria` ist dabei ein zu `influx` alternativer, exklusiver Wert (nicht über `both` kombinierbar). *sensorapi* liest die Messwerte unabhängig davon über **DBASE** (dieselben Werte `influx`/`mongo`/`victoria`, siehe Container-Beschreibung unten).
Danach einloggen in das Docker-Registry auf *citysensor.de*:
~~~
@@ -91,7 +99,7 @@ Im **Portainer** als eigenes Stack anlegen (Inhalt von `readarchive-stack.yml`)
`-s` | `START` | gestern
`-e` | `END` | Startdatum + 1 Tag (**exklusiv**)
`-t` | `TYP` | `noise`
`-d` | `DBASE` | `both`, in der Stack-Datei auf `influx` vorbelegt
`-d` | `DBASE` | `mongo`, `influx`, `victoria` oder `both` (= mongo+influx); in der Stack-Datei auf `influx` vorbelegt
*Deploy the stack* startet den Lauf, der Container endet danach und bleibt als „Exited" stehen. Für den nächsten Zeitraum nur die Variablen ändern und *Update the stack*. Der Haupt-Stack wird dabei nicht angefasst.
@@ -143,13 +151,19 @@ Die Influx-Datenbank. Hier werden in dem Bucket *sensor_data* die reinen Messwe
http://\<server>:8086
Logindaten: Entsprechen der im Portainer hinterlegten.
* **victoriametrics**
Alternative zur Influx-Datenbank für die reinen Messwerte, wählbar über `STORE`/`DBASE=victoria`. Läuft als Single-Node-Instanz auf **Port 8428**, ohne Login/Auth.
**Wichtig:** VictoriaMetrics löscht standardmäßig Daten, die älter als 1 Monat sind (`-retentionPeriod=1`). Im `docker-compose.yml` ist das Flag deshalb explizit auf `--retentionPeriod=100y` gesetzt — das **nicht** entfernen, sonst verschwinden alte Messwerte nach 30 Tagen kommentarlos.
**Zugriff:**
http://\<server>:8428 (z.B. `/api/v1/export?match[]={sid="..."}` zum Rohdaten-Export)
## Sourcen für die Container
Das gesamte Projekt ist im GitHub unter **laermsensor-stack** abgelegt. Für jeden der 4 Container existiert darunter ein Verzeichnis:
**readin**, **readarchive**, **sensorapi** und **noise**. In diesen Verzeichnissen sind alle benötigten Sourcen enthalten.
Daneben gibt es **common** mit den Modulen, die sich *readin* und *readarchive* teilen: `mongo.js`, `influx_post.js`, `logit.js` und `statistics.js`. Sie lagen früher in beiden Komponenten doppelt und sind auseinandergelaufen. Wer dort etwas ändert, ändert es für beide — nach einer Änderung also **beide** Container neu bauen.
Daneben gibt es **common** mit den Modulen, die sich *readin* und *readarchive* teilen: `mongo.js`, `influx_post.js`, `victoria_post.js`, `logit.js` und `statistics.js`. Sie lagen früher in beiden Komponenten doppelt und sind auseinandergelaufen. Wer dort etwas ändert, ändert es für beide — nach einer Änderung also **beide** Container neu bauen.
Sollte was geändert werden, so muss mit **deploy.sh** der Container neu erzeugt und auf die Registry (siehe oben) gepushed werden. Der Build-Kontext von *readin* und *readarchive* ist wegen `common/` das Wurzelverzeichnis des Repositories; `deploy.sh` wechselt selbst dorthin und kann aus dem Komponentenverzeichnis aufgerufen werden.
+55
View File
@@ -0,0 +1,55 @@
/* Zugriff auf VictoriaMetrics per HTTP (InfluxDB-Line-Protocol-kompatibler Write-Endpunkt)
Gemeinsam genutzt von readin (Live-API) und readarchive (CSV-Archiv).
*/
import axios from 'axios'
import { logit, logerror } from './logit.js'
import { statistics } from './statistics.js'
import { DateTime } from 'luxon'
let VICTORIAHOST = process.env.VICTORIAHOST || "localhost"
let VICTORIAPORT = process.env.VICTORIAPORT || 8428
const VICTORIAURL_WRITE = `http://${VICTORIAHOST}:${VICTORIAPORT}/write?precision=ms`
// `${e}` liefert bei einem AggregateError nur "AggregateError" - die einzelnen
// Verbindungsfehler (je einer pro aufgeloester IP) stecken in e.errors bzw. e.cause.errors
function describeError(e) {
const sub = e.errors || e.cause?.errors
if (sub?.length) {
return `${e.message || e.name}: ` + sub.map(s => `${s.code} ${s.address}:${s.port}`).join(', ')
}
if (e.response) {
return `HTTP ${e.response.status} ${JSON.stringify(e.response.data)}`
}
return `${e.code ? e.code + ' ' : ''}${e.message || e}`
}
// liefert true, wenn VictoriaMetrics die Daten uebernommen hat, sonst false
export const victoriaWrite = async (data) => {
let start = DateTime.now()
let ok = false
try {
const ret = await axios({
method: 'post',
url: VICTORIAURL_WRITE,
data: data,
headers: {
Accept: 'application/json',
'Content-Type': 'text/plain; charset=utf-8'
},
timeout: 10000,
})
if (ret.status != 204) {
logerror(`doWrite2API Status: ${ret.status}`)
} else {
ok = true
}
} catch (e) {
logerror(`doWrite2API ${VICTORIAURL_WRITE} ${describeError(e)}`)
}
statistics['writeVictoriaData[sensor_data]Time'] = DateTime.now().diff(start, ['seconds']).toObject().seconds
logit(`Victoria-Write-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`)
return ok
}
+249
View File
@@ -0,0 +1,249 @@
// compare.js - vergleicht die Laermwerte zweier "noise"-API-Instanzen
// (z.B. lokaler VictoriaMetrics-Test-Stack vs. Produktion) fuer denselben
// Sensor/Zeitraum/Auswertungstyp. Siehe Laerm_API.md fuer die API-Parameter.
//
// Aufruf: node compare.js [-s sensorid] [-d live|havg|davg|daynight|lden] ...
// node compare.js -h fuer alle Optionen
import axios from 'axios'
import { DateTime } from 'luxon'
import mod_getopt from 'posix-getopt'
const DEFAULT_URL_A = 'http://localhost:3003/api/getsensordata'
const DEFAULT_URL_B = 'https://noise.citysensor.de/api/getsensordata'
const DEFAULT_SENSORID = '37833'
const DEFAULT_DATA = 'live'
const DEFAULT_SPAN = '1'
const DEFAULT_PEAK = '70'
const DEFAULT_TOLERANCE = 0.05
const MAX_EXAMPLES = 20
function parseArgs(argv) {
let opts = {
urlA: DEFAULT_URL_A,
urlB: DEFAULT_URL_B,
sensorid: DEFAULT_SENSORID,
data: DEFAULT_DATA,
span: DEFAULT_SPAN,
datetime: null,
peak: DEFAULT_PEAK,
tolerance: DEFAULT_TOLERANCE,
}
let parser = new mod_getopt.BasicParser(
'a:(urlA)b:(urlB)s:(sensorid)d:(data)p:(span)t:(datetime)k:(peak)e:(tolerance)h(help)v(version)',
argv
)
let option
while ((option = parser.getopt()) !== undefined) {
switch (option.option) {
case 'a': opts.urlA = option.optarg; break
case 'b': opts.urlB = option.optarg; break
case 's': opts.sensorid = option.optarg; break
case 'd': opts.data = option.optarg; break
case 'p': opts.span = option.optarg; break
case 't': opts.datetime = option.optarg; break
case 'k': opts.peak = option.optarg; break
case 'e': opts.tolerance = parseFloat(option.optarg); break
case 'v':
console.log('compare-laermwerte 1.0.0')
process.exit()
break
case 'h':
console.log('Usage: node compare.js [options]')
console.log('Options:')
console.log(` -a urlA Basis-URL Implementierung A (default: ${DEFAULT_URL_A})`)
console.log(` -b urlB Basis-URL Implementierung B (default: ${DEFAULT_URL_B})`)
console.log(` -s sensorid Sensor-ID (default: ${DEFAULT_SENSORID})`)
console.log(` -d data live|havg|davg|daynight|lden (default: ${DEFAULT_DATA})`)
console.log(` -p span Zeitspanne in Tagen (default: ${DEFAULT_SPAN})`)
console.log(' -t datetime Start-Zeitpunkt ISO8601 (default: unbelegt -> jetzt - span)')
console.log(` -k peak dB-Schwelle fuer peakcount, nur havg/davg (default: ${DEFAULT_PEAK})`)
console.log(` -e tolerance Toleranz fuer Zahlenvergleich, absolut (default: ${DEFAULT_TOLERANCE})`)
console.log(' -v Version anzeigen')
console.log(' -h diese Hilfe')
process.exit()
break
default:
break
}
}
return opts
}
function buildURL(base, opts) {
let params = {
sensorid: opts.sensorid,
data: opts.data,
span: opts.span,
peak: opts.peak,
}
if (opts.datetime) {
params.datetime = opts.datetime
}
let qs = new URLSearchParams(params).toString()
return `${base}?${qs}`
}
async function fetchSide(label, url) {
try {
let ret = await axios.get(url, { timeout: 15000 })
if (ret.data && ret.data.err) {
return { ok: false, err: `${label}: API-Fehler: ${ret.data.err}` }
}
return { ok: true, values: ret.data.values || [], options: ret.data.options }
} catch (e) {
let msg = e.response ? `HTTP ${e.response.status}` : e.message
return { ok: false, err: `${label}: Request fehlgeschlagen: ${msg}` }
}
}
function toMap(values) {
let map = new Map()
for (let row of values) {
let ms = DateTime.fromISO(row.datetime, { zone: 'utc' }).toMillis()
if (Number.isNaN(ms)) {
continue
}
map.set(ms, row)
}
return map
}
function diffRows(mapA, mapB, tolerance) {
let onlyA = []
let onlyB = []
let commonCount = 0
let fieldStats = new Map() // field -> {count, sumDiff, maxDiff, mismatchCount}
let examples = []
for (let ts of mapA.keys()) {
if (!mapB.has(ts)) {
onlyA.push(ts)
}
}
for (let ts of mapB.keys()) {
if (!mapA.has(ts)) {
onlyB.push(ts)
}
}
for (let [ts, rowA] of mapA) {
let rowB = mapB.get(ts)
if (!rowB) {
continue
}
commonCount++
let fields = new Set([...Object.keys(rowA), ...Object.keys(rowB)])
fields.delete('datetime')
for (let field of fields) {
let a = (field in rowA) ? rowA[field] : null
let b = (field in rowB) ? rowB[field] : null
if (a === null && b === null) {
continue
}
let stat = fieldStats.get(field)
if (!stat) {
stat = { count: 0, sumDiff: 0, maxDiff: 0, mismatchCount: 0 }
fieldStats.set(field, stat)
}
stat.count++
if (typeof a !== 'number' || typeof b !== 'number') {
stat.mismatchCount++
if (examples.length < MAX_EXAMPLES) {
examples.push({ ts, field, a, b, reason: 'null-mismatch' })
}
continue
}
let diff = Math.abs(a - b)
stat.sumDiff += diff
if (diff > stat.maxDiff) {
stat.maxDiff = diff
}
if (diff > tolerance) {
stat.mismatchCount++
if (examples.length < MAX_EXAMPLES) {
examples.push({ ts, field, a, b, diff })
}
}
}
}
return { onlyA, onlyB, commonCount, fieldStats, examples }
}
function fmtTs(ms) {
return DateTime.fromMillis(ms, { zone: 'utc' }).toISO()
}
function printSummary(opts, sideA, sideB, diff) {
console.log('=== Vergleich Laermwerte ===')
console.log(`Sensor: ${opts.sensorid} data: ${opts.data} span: ${opts.span}${opts.datetime ? ` datetime: ${opts.datetime}` : ''}`)
console.log(`A: ${sideA.count} Werte`)
console.log(`B: ${sideB.count} Werte`)
console.log(`Gemeinsame Zeitstempel: ${diff.commonCount}`)
console.log(`Nur in A: ${diff.onlyA.length}${diff.onlyA.length ? ` (z.B. ${diff.onlyA.slice(0, 3).map(fmtTs).join(', ')})` : ''}`)
console.log(`Nur in B: ${diff.onlyB.length}${diff.onlyB.length ? ` (z.B. ${diff.onlyB.slice(0, 3).map(fmtTs).join(', ')})` : ''}`)
console.log('')
console.log('Felder (bei gemeinsamen Zeitstempeln):')
let totalMismatches = 0
for (let [field, stat] of diff.fieldStats) {
let avg = stat.count ? stat.sumDiff / stat.count : 0
totalMismatches += stat.mismatchCount
console.log(` ${field.padEnd(14)} verglichen: ${String(stat.count).padStart(5)} max diff: ${stat.maxDiff.toFixed(4).padStart(10)} avg diff: ${avg.toFixed(4).padStart(10)} mismatches (>${opts.tolerance}): ${stat.mismatchCount}`)
}
if (diff.examples.length > 0) {
console.log('')
console.log(`Beispiel-Mismatches (max ${MAX_EXAMPLES}):`)
for (let ex of diff.examples) {
if (ex.reason === 'null-mismatch') {
console.log(` ${fmtTs(ex.ts)} ${ex.field}: A=${ex.a} B=${ex.b}`)
} else {
console.log(` ${fmtTs(ex.ts)} ${ex.field}: A=${ex.a} B=${ex.b} diff=${ex.diff.toFixed(4)}`)
}
}
}
console.log('')
if (totalMismatches > 0) {
console.log(`FEHLGESCHLAGEN: ${totalMismatches} Feld-Mismatches ausserhalb der Toleranz gefunden.`)
} else {
console.log('OK: Keine Mismatches ausserhalb der Toleranz bei gemeinsamen Zeitstempeln.')
}
return totalMismatches
}
async function main() {
let opts = parseArgs(process.argv)
let urlA = buildURL(opts.urlA, opts)
let urlB = buildURL(opts.urlB, opts)
console.log(`A: ${urlA}`)
console.log(`B: ${urlB}`)
console.log('')
let [sideA, sideB] = await Promise.all([
fetchSide('A', urlA),
fetchSide('B', urlB),
])
if (!sideA.ok || !sideB.ok) {
if (!sideA.ok) console.error(sideA.err)
if (!sideB.ok) console.error(sideB.err)
process.exitCode = 2
return
}
let mapA = toMap(sideA.values)
let mapB = toMap(sideB.values)
let diff = diffRows(mapA, mapB, opts.tolerance)
let mismatches = printSummary(
opts,
{ count: sideA.values.length },
{ count: sideB.values.length },
diff
)
process.exitCode = mismatches > 0 ? 1 : 0
}
main().catch((e) => {
console.error(e)
process.exitCode = 2
})
+366
View File
@@ -0,0 +1,366 @@
{
"name": "compare-laermwerte",
"version": "1.0.0",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "compare-laermwerte",
"version": "1.0.0",
"dependencies": {
"axios": "^1.12.0",
"luxon": "^3.3.0",
"posix-getopt": "^1.2.1"
}
},
"node_modules/agent-base": {
"version": "6.0.2",
"resolved": "https://registry.npmjs.org/agent-base/-/agent-base-6.0.2.tgz",
"integrity": "sha512-RZNwNclF7+MS/8bDg70amg32dyeZGZxiDuQmZxKLAlQjr3jGyLx+4Kkk58UO7D2QdgFIQCovuSuZESne6RG6XQ==",
"license": "MIT",
"dependencies": {
"debug": "4"
},
"engines": {
"node": ">= 6.0.0"
}
},
"node_modules/asynckit": {
"version": "0.4.0",
"resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz",
"integrity": "sha512-Oei9OH4tRh0YqU3GxhX79dM/mwVgvbZJaSNaRk+bshkj0S5cfHcgYakreBjrHwatXKbz+IoIdYLxrKim2MjW0Q==",
"license": "MIT"
},
"node_modules/axios": {
"version": "1.19.0",
"resolved": "https://registry.npmjs.org/axios/-/axios-1.19.0.tgz",
"integrity": "sha512-ht/iuYZXEjFxLH/Hkezgd7m6JKlHHXEUSneaDz8uZe1Gj5QZtCnpyDsckvAiEnT89OEbCLmnte4R4sn7P0EKFw==",
"license": "MIT",
"dependencies": {
"follow-redirects": "^1.16.0",
"form-data": "^4.0.6",
"https-proxy-agent": "^5.0.1",
"proxy-from-env": "^2.1.0"
}
},
"node_modules/call-bind-apply-helpers": {
"version": "1.0.2",
"resolved": "https://registry.npmjs.org/call-bind-apply-helpers/-/call-bind-apply-helpers-1.0.2.tgz",
"integrity": "sha512-Sp1ablJ0ivDkSzjcaJdxEunN5/XvksFJ2sMBFfq6x0ryhQV/2b/KwFe21cMpmHtPOSij8K99/wSfoEuTObmuMQ==",
"license": "MIT",
"dependencies": {
"es-errors": "^1.3.0",
"function-bind": "^1.1.2"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/combined-stream": {
"version": "1.0.8",
"resolved": "https://registry.npmjs.org/combined-stream/-/combined-stream-1.0.8.tgz",
"integrity": "sha512-FQN4MRfuJeHf7cBbBMJFXhKSDq+2kAArBlmRBvcvFE5BB1HZKXtSFASDhdlz9zOYwxh8lDdnvmMOe/+5cdoEdg==",
"license": "MIT",
"dependencies": {
"delayed-stream": "~1.0.0"
},
"engines": {
"node": ">= 0.8"
}
},
"node_modules/debug": {
"version": "4.4.3",
"resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz",
"integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==",
"license": "MIT",
"dependencies": {
"ms": "^2.1.3"
},
"engines": {
"node": ">=6.0"
},
"peerDependenciesMeta": {
"supports-color": {
"optional": true
}
}
},
"node_modules/delayed-stream": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/delayed-stream/-/delayed-stream-1.0.0.tgz",
"integrity": "sha512-ZySD7Nf91aLB0RxL4KGrKHBXl7Eds1DAmEdcoVawXnLD7SDhpNgtuII2aAkg7a7QS41jxPSZ17p4VdGnMHk3MQ==",
"license": "MIT",
"engines": {
"node": ">=0.4.0"
}
},
"node_modules/dunder-proto": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/dunder-proto/-/dunder-proto-1.0.1.tgz",
"integrity": "sha512-KIN/nDJBQRcXw0MLVhZE9iQHmG68qAVIBg9CqmUYjmQIhgij9U5MFvrqkUL5FbtyyzZuOeOt0zdeRe4UY7ct+A==",
"license": "MIT",
"dependencies": {
"call-bind-apply-helpers": "^1.0.1",
"es-errors": "^1.3.0",
"gopd": "^1.2.0"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/es-define-property": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/es-define-property/-/es-define-property-1.0.1.tgz",
"integrity": "sha512-e3nRfgfUZ4rNGL232gUgX06QNyyez04KdjFrF+LTRoOXmrOgFKDg4BCdsjW8EnT69eqdYGmRpJwiPVYNrCaW3g==",
"license": "MIT",
"engines": {
"node": ">= 0.4"
}
},
"node_modules/es-errors": {
"version": "1.3.0",
"resolved": "https://registry.npmjs.org/es-errors/-/es-errors-1.3.0.tgz",
"integrity": "sha512-Zf5H2Kxt2xjTvbJvP2ZWLEICxA6j+hAmMzIlypy4xcBg1vKVnx89Wy0GbS+kf5cwCVFFzdCFh2XSCFNULS6csw==",
"license": "MIT",
"engines": {
"node": ">= 0.4"
}
},
"node_modules/es-object-atoms": {
"version": "1.1.2",
"resolved": "https://registry.npmjs.org/es-object-atoms/-/es-object-atoms-1.1.2.tgz",
"integrity": "sha512-HWcBoN6NileqtSydK2FqHbS/LoDd2pqrnQHLyJzBj4kOp/ky2MWMN694xOfkK8/SnUsW2DH7EfyVlydKCsm1Zw==",
"license": "MIT",
"dependencies": {
"es-errors": "^1.3.0"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/es-set-tostringtag": {
"version": "2.1.0",
"resolved": "https://registry.npmjs.org/es-set-tostringtag/-/es-set-tostringtag-2.1.0.tgz",
"integrity": "sha512-j6vWzfrGVfyXxge+O0x5sh6cvxAog0a/4Rdd2K36zCMV5eJ+/+tOAngRO8cODMNWbVRdVlmGZQL2YS3yR8bIUA==",
"license": "MIT",
"dependencies": {
"es-errors": "^1.3.0",
"get-intrinsic": "^1.2.6",
"has-tostringtag": "^1.0.2",
"hasown": "^2.0.2"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/follow-redirects": {
"version": "1.16.0",
"resolved": "https://registry.npmjs.org/follow-redirects/-/follow-redirects-1.16.0.tgz",
"integrity": "sha512-y5rN/uOsadFT/JfYwhxRS5R7Qce+g3zG97+JrtFZlC9klX/W5hD7iiLzScI4nZqUS7DNUdhPgw4xI8W2LuXlUw==",
"funding": [
{
"type": "individual",
"url": "https://github.com/sponsors/RubenVerborgh"
}
],
"license": "MIT",
"engines": {
"node": ">=4.0"
},
"peerDependenciesMeta": {
"debug": {
"optional": true
}
}
},
"node_modules/form-data": {
"version": "4.0.6",
"resolved": "https://registry.npmjs.org/form-data/-/form-data-4.0.6.tgz",
"integrity": "sha512-vKatAh4SlVfgbv+YtmhiRjhEMJsYpsG1Y2rMQtR+SVSbytsSD1YGzDIcrAJmdFec88u/+VoGmxnl+80gL1tRCQ==",
"license": "MIT",
"dependencies": {
"asynckit": "^0.4.0",
"combined-stream": "^1.0.8",
"es-set-tostringtag": "^2.1.0",
"hasown": "^2.0.4",
"mime-types": "^2.1.35"
},
"engines": {
"node": ">= 6"
}
},
"node_modules/function-bind": {
"version": "1.1.2",
"resolved": "https://registry.npmjs.org/function-bind/-/function-bind-1.1.2.tgz",
"integrity": "sha512-7XHNxH7qX9xG5mIwxkhumTox/MIRNcOgDrxWsMt2pAr23WHp6MrRlN7FBSFpCpr+oVO0F744iUgR82nJMfG2SA==",
"license": "MIT",
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/get-intrinsic": {
"version": "1.3.0",
"resolved": "https://registry.npmjs.org/get-intrinsic/-/get-intrinsic-1.3.0.tgz",
"integrity": "sha512-9fSjSaos/fRIVIp+xSJlE6lfwhES7LNtKaCBIamHsjr2na1BiABJPo0mOjjz8GJDURarmCPGqaiVg5mfjb98CQ==",
"license": "MIT",
"dependencies": {
"call-bind-apply-helpers": "^1.0.2",
"es-define-property": "^1.0.1",
"es-errors": "^1.3.0",
"es-object-atoms": "^1.1.1",
"function-bind": "^1.1.2",
"get-proto": "^1.0.1",
"gopd": "^1.2.0",
"has-symbols": "^1.1.0",
"hasown": "^2.0.2",
"math-intrinsics": "^1.1.0"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/get-proto": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/get-proto/-/get-proto-1.0.1.tgz",
"integrity": "sha512-sTSfBjoXBp89JvIKIefqw7U2CCebsc74kiY6awiGogKtoSGbgjYE/G/+l9sF3MWFPNc9IcoOC4ODfKHfxFmp0g==",
"license": "MIT",
"dependencies": {
"dunder-proto": "^1.0.1",
"es-object-atoms": "^1.0.0"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/gopd": {
"version": "1.2.0",
"resolved": "https://registry.npmjs.org/gopd/-/gopd-1.2.0.tgz",
"integrity": "sha512-ZUKRh6/kUFoAiTAtTYPZJ3hw9wNxx+BIBOijnlG9PnrJsCcSjs1wyyD6vJpaYtgnzDrKYRSqf3OO6Rfa93xsRg==",
"license": "MIT",
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/has-symbols": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/has-symbols/-/has-symbols-1.1.0.tgz",
"integrity": "sha512-1cDNdwJ2Jaohmb3sg4OmKaMBwuC48sYni5HUw2DvsC8LjGTLK9h+eb1X6RyuOHe4hT0ULCW68iomhjUoKUqlPQ==",
"license": "MIT",
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/has-tostringtag": {
"version": "1.0.2",
"resolved": "https://registry.npmjs.org/has-tostringtag/-/has-tostringtag-1.0.2.tgz",
"integrity": "sha512-NqADB8VjPFLM2V0VvHUewwwsw0ZWBaIdgo+ieHtK3hasLz4qeCRjYcqfB6AQrBggRKppKF8L52/VqdVsO47Dlw==",
"license": "MIT",
"dependencies": {
"has-symbols": "^1.0.3"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/hasown": {
"version": "2.0.4",
"resolved": "https://registry.npmjs.org/hasown/-/hasown-2.0.4.tgz",
"integrity": "sha512-T2UbfbBEF32wiepXIsMlTW9+dDYC6wMh/t/vYA4tuOMKqWz/n3vr1NFSxQiyP+zk2mXsoMA/i/7qV6LKut1t1A==",
"license": "MIT",
"dependencies": {
"function-bind": "^1.1.2"
},
"engines": {
"node": ">= 0.4"
}
},
"node_modules/https-proxy-agent": {
"version": "5.0.1",
"resolved": "https://registry.npmjs.org/https-proxy-agent/-/https-proxy-agent-5.0.1.tgz",
"integrity": "sha512-dFcAjpTQFgoLMzC2VwU+C/CbS7uRL0lWmxDITmqm7C+7F0Odmj6s9l6alZc6AELXhrnggM2CeWSXHGOdX2YtwA==",
"license": "MIT",
"dependencies": {
"agent-base": "6",
"debug": "4"
},
"engines": {
"node": ">= 6"
}
},
"node_modules/luxon": {
"version": "3.7.2",
"resolved": "https://registry.npmjs.org/luxon/-/luxon-3.7.2.tgz",
"integrity": "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew==",
"license": "MIT",
"engines": {
"node": ">=12"
}
},
"node_modules/math-intrinsics": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/math-intrinsics/-/math-intrinsics-1.1.0.tgz",
"integrity": "sha512-/IXtbwEk5HTPyEwyKX6hGkYXxM9nbj64B+ilVJnC/R6B0pH5G4V3b0pVbL7DBj4tkhBAppbQUlf6F6Xl9LHu1g==",
"license": "MIT",
"engines": {
"node": ">= 0.4"
}
},
"node_modules/mime-db": {
"version": "1.52.0",
"resolved": "https://registry.npmjs.org/mime-db/-/mime-db-1.52.0.tgz",
"integrity": "sha512-sPU4uV7dYlvtWJxwwxHD0PuihVNiE7TyAbQ5SWxDCB9mUYvOgroQOwYQQOKPJ8CIbE+1ETVlOoK1UC2nU3gYvg==",
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/mime-types": {
"version": "2.1.35",
"resolved": "https://registry.npmjs.org/mime-types/-/mime-types-2.1.35.tgz",
"integrity": "sha512-ZDY+bPm5zTTF+YpCrAU9nK0UgICYPT0QtT1NZWFv4s++TNkcgVaT0g6+4R2uI4MjQjzysHB1zxuWL50hzaeXiw==",
"license": "MIT",
"dependencies": {
"mime-db": "1.52.0"
},
"engines": {
"node": ">= 0.6"
}
},
"node_modules/ms": {
"version": "2.1.3",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz",
"integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==",
"license": "MIT"
},
"node_modules/posix-getopt": {
"version": "1.2.1",
"resolved": "https://registry.npmjs.org/posix-getopt/-/posix-getopt-1.2.1.tgz",
"integrity": "sha512-BbGTiH8MOWAuc6h5yITkSn9k3HP4+QOCV9t6I5F62OrH7zqTHRo08QNsgELRreTBxcvRhbSpMoUnAx77Dz4yUA==",
"license": "MIT",
"engines": {
"node": "*"
}
},
"node_modules/proxy-from-env": {
"version": "2.1.0",
"resolved": "https://registry.npmjs.org/proxy-from-env/-/proxy-from-env-2.1.0.tgz",
"integrity": "sha512-cJ+oHTW1VAEa8cJslgmUZrc+sjRKgAKl3Zyse6+PV38hZe/V6Z14TbCuXcan9F9ghlz4QrFr2c92TNF82UkYHA==",
"license": "MIT",
"engines": {
"node": ">=10"
}
}
}
}
+15
View File
@@ -0,0 +1,15 @@
{
"name": "compare-laermwerte",
"version": "1.0.0",
"description": "Vergleicht Laermwerte zweier API-Instanzen (z.B. lokaler Victoria-Test-Stack vs. Produktion)",
"type": "module",
"main": "compare.js",
"scripts": {
"start": "node compare.js"
},
"dependencies": {
"axios": "^1.12.0",
"luxon": "^3.3.0",
"posix-getopt": "^1.2.1"
}
}
+4 -1
View File
@@ -9,9 +9,12 @@ MONGO_ROOT_PASSWORD=<Mongo_Passwort>
LOCALDIR=<Aktuelle Direcrory ab />
# Datenbank, in welche die laufenden Messwerte gespeichert werden (kannn 'influx' oder 'mongo' sein)
# Datenbank, in welche die laufenden Messwerte gespeichert werden (kann 'influx', 'mongo' oder 'victoria' sein)
STORE=influx
# Datenbank, aus welcher sensorapi die Messwerte liest (kann 'influx', 'mongo' oder 'victoria' sein)
DBASE=influx
# Secret zum Signieren der Session-Cookies von esp2sensor (langer Zufallsstring)
ESP2SENSOR_SESSION_SECRET=<zufaelliger_string>
+19 -5
View File
@@ -26,20 +26,33 @@ services:
DOCKER_INFLUXDB_INIT_ORG: citysensor
DOCKER_INFLUXDB_INIT_BUCKET: sensor_data
DOCKER_INFLUXDB_INIT_ADMIN_TOKEN: ${DOCKER_INFLUXTOKEN}
restart:
unless-stopped
container_name: influxdb
restart:
unless-stopped
container_name: influxdb
victoriametrics:
image: victoriametrics/victoria-metrics:v1.102.0
ports:
- "8428:8428"
volumes:
- ${LOCALDIR}/noisesensors/data/victoria/data:/victoria-metrics-data
command:
- "--storageDataPath=/victoria-metrics-data"
- "--retentionPeriod=100y"
restart: unless-stopped
container_name: victoriametrics
readin:
image: docker.citysensor.de/readin
environment:
MONGOHOST: mongodb
INFLUXHOST: influxdb
VICTORIAHOST: victoriametrics
TYP: "[\"noise\"]"
STORE: ${STORE:-mongo}
MONGOAUTH: true
MONGOUSRP: ${MONGO_ROOT_USERNAME}:${MONGO_ROOT_PASSWORD}
INFLUXTOKEN: ${DOCKER_INFLUXTOKEN}
INFLUXTOKEN: ${DOCKER_INFLUXTOKEN}
volumes:
- ${LOCALDIR}/noisesensors/log:/var/log
# Pfad im Image ist jetzt /opt/app/readin, weil readin und readarchive
@@ -67,7 +80,8 @@ services:
environment:
INFLUXHOST: influxdb
INFLUXTOKEN: ${DOCKER_INFLUXTOKEN}
DBASE: influx
VICTORIAHOST: victoriametrics
DBASE: ${DBASE:-influx}
MONGOHOST: mongodb
MONGOAUTH: true
MONGOUSRP: ${MONGO_ROOT_USERNAME}:${MONGO_ROOT_PASSWORD}
+3 -1
View File
@@ -16,7 +16,8 @@
#
# readFromcsv.js liest jeden Parameter auch aus einer Env-Variablen:
# -s = START (Default: gestern), -e = END (Default: Startdatum + 1 Tag,
# exklusiv), -t = TYP (Default: noise), -d = DBASE (Default: both),
# exklusiv), -t = TYP (Default: noise), -d = DBASE ('mongo', 'influx',
# 'victoria' oder 'both' = mongo+influx; Default: both),
# -p = PARALLEL (Default: 4) - so viele Sensoren werden gleichzeitig geholt.
#
# Voraussetzungen:
@@ -39,6 +40,7 @@ services:
environment:
MONGOHOST: mongodb
INFLUXHOST: influxdb
VICTORIAHOST: victoriametrics
MONGOAUTH: "true"
MONGOUSRP: ${MONGO_ROOT_USERNAME}:${MONGO_ROOT_PASSWORD}
INFLUXTOKEN: ${DOCKER_INFLUXTOKEN}
+18 -2
View File
@@ -33,7 +33,7 @@ const DBASE = process.env.DBASE || 'both'
// Mehr als 4 bringt nichts: ab da liefert das Archiv nicht mehr schneller.
const PARALLEL = parseInt(process.env.PARALLEL) || 4
const DATABASE = ['mongo', 'influx', 'both']
const DATABASE = ['mongo', 'influx', 'victoria', 'both']
const DEVELOP = process.env.DEVELOP || false;
const LIVE = (process.env.LIVE == "true") || true
@@ -47,6 +47,7 @@ import axios from 'axios'
import https from 'https'
import * as mongo from '../common/mongo.js'
import * as influx from '../common/influx_post.js'
import * as victoria from '../common/victoria_post.js'
import fs from 'fs'
import csv from 'csvtojson'
import pkg from './package.json' with { type: "json" }
@@ -67,6 +68,8 @@ const newProps = []
// Sensor/Tag-Paare, deren Influx-Write fehlgeschlagen ist - Auswertung am Ende von main()
const failedInfluxWrites = []
// Sensor/Tag-Paare, deren VictoriaMetrics-Write fehlgeschlagen ist - Auswertung am Ende von main()
const failedVictoriaWrites = []
const saveList = (list) => {
fs.writeFileSync('data/list.json', JSON.stringify(list))
@@ -397,6 +400,12 @@ async function enterOneSensorinDB(client, sid, erg, dataline, day, styp, args) {
failedInfluxWrites.push({sensorid: sid, day: day})
}
}
if ((args.database === 'victoria') && (dataline !== '')) {
const ok = await victoria.victoriaWrite(dataline)
if (!ok) {
failedVictoriaWrites.push({sensorid: sid, day: day})
}
}
if((args.database === 'both') || (args.database === 'mongo')) {
if(erg.length !== 0) {
await mongo.writeDataArray(client, mongo.dataCollName(styp), erg)
@@ -507,7 +516,7 @@ function parse_cmdline(argv) {
console.log(" will be used and endDate = startDate + 1 day")
console.log(" -t sensorType: if given, only those sensors will be used; default: 'noise'");
console.log(" allowed types: 'noise', 'radiactivity', 'pm' and 'thp'")
console.log(" -d database: 'mongo', 'influx' or 'both'; default: 'both'")
console.log(" -d database: 'mongo', 'influx', 'victoria' or 'both' ('both' = mongo + influx); default: 'both'")
console.log(` -p parallel: number of sensors fetched at the same time; default: ${PARALLEL}`)
console.log(" -v version: show version");
console.log(" -h this help text");
@@ -555,6 +564,13 @@ async function main() {
}
process.exitCode = 1 // damit cron/docker den Fehlschlag sieht
}
if (failedVictoriaWrites.length > 0) {
logerror(`${failedVictoriaWrites.length} Victoria-Writes fehlgeschlagen - diese Daten fehlen in VictoriaMetrics:`)
for (const f of failedVictoriaWrites) {
logerror(` Sensor ${f.sensorid}, Tag ${f.day}`)
}
process.exitCode = 1 // damit cron/docker den Fehlschlag sieht
}
let dauer = starttime.diffNow('seconds').toObject().seconds * -1
const duration1 = Duration.fromObject({ seconds: dauer });
const output1 = duration1.toFormat('hh:mm:ss');
+9 -2
View File
@@ -14,13 +14,14 @@
// TYP: Ist ein Array von Strings, wenn nur einzelne Typen gespeichert werden sollen, also z.B.:
// ['pm', 'noise']
const TYP = process.env.TYP || ''
const STORE = (process.env.STORE || 'both').toLowerCase() // 'mongo' | 'influx' | 'both'
const STORE = (process.env.STORE || 'both').toLowerCase() // 'mongo' | 'influx' | 'victoria' | 'both' ('both' = mongo + influx)
import { doReadfromAPI as readin } from './readdata.js'
import { constructDBaseEntries as parse} from './parse.js'
import * as mongo from '../common/mongo.js'
import * as influx from '../common/influx_post.js'
import * as victoria from '../common/victoria_post.js'
import { logit, logerror } from '../common/logit.js'
import { statistics } from '../common/statistics.js'
import { DateTime } from 'luxon'
@@ -62,6 +63,11 @@ const fetchNewData = async (args) => {
await influx.influxWrite(idata)
}
// write sensor data to VictoriaMetrics
if(args.victoria) {
await victoria.victoriaWrite(idata)
}
// write properties to mongoDB
await mongo.bulkWrite(client, mongo.property_coll, props)
}
@@ -78,7 +84,7 @@ const fetchNewData = async (args) => {
function parse_cmdline(argv) {
let parser = new mod_getopt.BasicParser('i(influx)m(mongo)t:(typ)h(help)v(version)',argv);
let option;
let ret = {influx: STORE === 'both' || STORE === 'influx', mongo: STORE === 'both' || STORE === 'mongo', typ: TYP}
let ret = {influx: STORE === 'both' || STORE === 'influx', mongo: STORE === 'both' || STORE === 'mongo', victoria: STORE === 'victoria', typ: TYP}
while((option = parser.getopt()) !== undefined) {
switch(option.option) {
case 'i':
@@ -111,6 +117,7 @@ function parse_cmdline(argv) {
console.log("Params:");
console.log(" -i use only InfluxDB to store the data (default use both, InfluxDB and MongoDB))");
console.log(" -m use only MongoDB to store the data (default use both, InfluxDB and MongoDB)");
console.log(" (VictoriaMetrics storage is selected exclusively via STORE=victoria, no CLI flag)");
console.log(" -t [ sensorType, ..]: if given, only those sensors will be used (ex: laerm) default: all");
console.log(" MUST BE AN ARRAY!; allowed types: 'pm', 'noise', 'radiactivity', 'thp', 'gps'.")
console.log(" -v version: show version");
+4 -1
View File
@@ -4,6 +4,7 @@ const DBASE = process.env.DBASE || 'mongo'
import {DateTime} from "luxon"
import * as influx from "../databases/influx.js"
import * as mongo from "../databases/mongo.js"
import * as victoria from "../databases/victoria.js"
import {returnOnError} from "../utilities/reporterror.js"
import {csv2Json} from "../utilities/csv2json.js"
import checkParams from "../utilities/checkparams.js"
@@ -155,8 +156,10 @@ export async function getActData(opts) {
return await mongo.fetchActData(opts)
} else if (DBASE === 'influx') {
return await influx.fetchActData(opts)
} else if (DBASE === 'victoria') {
return await victoria.fetchActData(opts)
}
return {err: 'DBASEUNKNOWN', values: []}
return {err: 'DBASEUNKNOWN', values: []}
}
+67
View File
@@ -0,0 +1,67 @@
// Access to VictoriaMetrics via HTTP (Prometheus-compatible /api/v1/export,
// no Flux support - see common/victoria_post.js for the write side, which
// reuses the InfluxDB line-protocol format on VictoriaMetrics' /write endpoint).
import axios from 'axios'
import { logit, logerror } from '../utilities/logit.js'
import { returnOnError } from "../utilities/reporterror.js"
import { exportToRows, bucketNoiseAVG } from "../utilities/victoria2json.js"
let VICTORIAHOST = process.env.VICTORIAHOST || "localhost"
let VICTORIAPORT = process.env.VICTORIAPORT || 8428
const VICTORIAURL_EXPORT = `http://${VICTORIAHOST}:${VICTORIAPORT}/api/v1/export`
// opts.start/opts.stop arrive as Flux range() fragments ("start: <iso>" /
// "stop: <iso>") built by calcRange() in getsensorData.js - strip the Flux
// keyword the same way sensorapi/databases/mongo.js already does.
const isoStart = (opts) => opts.start.slice(7)
const isoStop = (opts) => opts.stop.slice(6)
const victoriaExport = async (matchSelector, start, end) => {
let erg = { values: '', err: null }
try {
let ret = await axios({
method: 'get',
url: VICTORIAURL_EXPORT,
params: { 'match[]': matchSelector, start, end },
timeout: 10000,
transformResponse: [(data) => data], // response body is ndjson, not a single JSON document - keep it raw
})
if (ret.status !== 200) {
return returnOnError(erg, 'RESPSTATUS', victoriaExport.name, ret.status)
}
erg.values = ret.data
} catch (e) {
return returnOnError(erg, e, victoriaExport.name)
}
return erg
}
export const fetchActData = async (opts) => {
let ret = { err: null, values: [] }
const match = `{__name__=~"noise_(LAeq|LA_min|LA_max|E10tel_eq)", sid="${opts.sensorid}"}`
let { values, err } = await victoriaExport(match, isoStart(opts), isoStop(opts))
if (err) {
return returnOnError(ret, err, fetchActData.name)
}
ret.values = exportToRows(values, opts.sort)
if (ret.values.length === 0) {
return returnOnError(ret, 'NODATA', fetchActData.name)
}
return ret
}
export const fetchNoiseAVGData = async (opts) => {
let ret = { err: null, values: [] }
const match = `{__name__=~"noise_(E10tel_eq|LA_max)", sid="${opts.sensorid}"}`
let { values, err } = await victoriaExport(match, isoStart(opts), isoStop(opts))
if (err) {
return returnOnError(ret, err, fetchNoiseAVGData.name)
}
ret.values = bucketNoiseAVG(values, opts.peak, opts.long)
if (ret.values.length === 0) {
return returnOnError(ret, 'NODATA', fetchNoiseAVGData.name)
}
return ret
}
+4
View File
@@ -9,6 +9,7 @@ import {DateTime} from 'luxon'
import { translate as trans } from '../routes/api.js'
import * as influx from "../databases/influx.js"
import * as mongo from "../databases/mongo.js"
import * as victoria from "../databases/victoria.js"
import { setoptionfromtable } from "../utilities/chartoptions.js"
export const getNoiseData = async (params, possibles, props) => {
@@ -451,6 +452,9 @@ const getNoiseAVGData = async (opts) => {
for (let x=0; x < ret.values.length; x++) {
ret.values[x].datetime = DateTime.fromISO(ret.values[x].datetime).toUTC().minus({hours:1}).toFormat("yyyy-LL-dd'T'HH:mm:ss'Z'")
}
} else if (DBASE === 'victoria') {
// victoria.js buckets by hour-start itself (like mongo), so no shift is needed here
ret = await victoria.fetchNoiseAVGData(opts)
} else {
ret.err = 'DBASEUNKNOWN'
}
+97
View File
@@ -0,0 +1,97 @@
// Parse VictoriaMetrics' /api/v1/export ndjson responses into the same row
// shapes influx.js builds via csv2Json() / pivot(), for the noise-specific
// fields (measurement "noise", separator "_", i.e. metric names noise_LAeq,
// noise_LA_min, noise_LA_max, noise_E10tel_eq).
import { DateTime } from 'luxon'
const parseExportLines = (ndjsonBody) => {
return String(ndjsonBody)
.split('\n')
.filter((line) => line.trim() !== '')
.map((line) => {
try {
return JSON.parse(line)
} catch (e) {
return null
}
})
.filter((series) => series !== null)
}
const fieldName = (metricName) => metricName.replace(/^noise_/, '')
// Merges the per-field series (one ndjson line per field) back into rows of
// {datetime, LAeq, LA_min, LA_max, E10tel_eq}, matched by raw timestamp -
// mirrors influx's pivot(rowKey:["_time"], columnKey:["_field"]).
export const exportToRows = (ndjsonBody, sort) => {
const rowsByTime = new Map()
for (const series of parseExportLines(ndjsonBody)) {
const field = fieldName(series.metric.__name__)
const { values, timestamps } = series
for (let i = 0; i < timestamps.length; i++) {
const ts = timestamps[i]
let row = rowsByTime.get(ts)
if (!row) {
row = { datetime: DateTime.fromMillis(ts).toUTC().toISO() }
rowsByTime.set(ts, row)
}
row[field] = values[i]
}
}
const rows = Array.from(rowsByTime.values())
rows.sort((a, b) => (a.datetime < b.datetime ? -1 : a.datetime > b.datetime ? 1 : 0))
if (sort === -1) {
rows.reverse()
}
return rows
}
// Buckets raw E10tel_eq/LA_max samples by hour (truncated to the start of the
// hour of each sample, same convention as sensorapi/databases/mongo.js's
// $dateToString hour-grouping) and computes n_AVG/n_sum/count/peakcount -
// equivalent to influx.js's aggregateWindow()/reduce() Flux pipeline.
// Bucketing by hour-start (not hour-end, as Influx does) means no extra
// 1-hour shift is needed downstream, unlike the DBASE === 'influx' branch.
export const bucketNoiseAVG = (ndjsonBody, peak, long) => {
const buckets = new Map()
for (const series of parseExportLines(ndjsonBody)) {
const field = fieldName(series.metric.__name__)
if (field !== 'E10tel_eq' && field !== 'LA_max') {
continue
}
const { values, timestamps } = series
for (let i = 0; i < timestamps.length; i++) {
const key = DateTime.fromMillis(timestamps[i]).toUTC().startOf('hour').toISO()
let b = buckets.get(key)
if (!b) {
b = { sum: 0, count: 0, peakcount: 0 }
buckets.set(key, b)
}
if (field === 'E10tel_eq') {
b.sum += values[i]
b.count += 1
} else {
if (values[i] >= peak) {
b.peakcount += 1
}
}
}
}
const rows = Array.from(buckets.entries())
.filter(([, b]) => b.count > 0) // matches Influx's inner join: hours without E10tel_eq samples are dropped
.map(([datetime, b]) => {
let row = {
datetime,
n_AVG: 10 * Math.log10(b.sum / b.count),
peakcount: b.peakcount,
}
if (long) {
row.count = b.count
row.n_sum = b.sum
}
return row
})
rows.sort((a, b) => (a.datetime < b.datetime ? -1 : a.datetime > b.datetime ? 1 : 0))
return rows
}