From 7c6692c0402f06680597ada6b3204c742ad6c6a0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Reinhard=20X=2E=20F=C3=BCrst?= Date: Fri, 31 Jul 2026 06:36:01 +0000 Subject: [PATCH] readarchive: Import aus dem CSV-Archiv dazu Gegenstueck zu readin/: waehrend readin alle 5 Minuten die Live-API abfragt, liest readarchive rueckwirkend die Tagesdateien von archive.sensor.community in dieselben Datenbanken ein. Ohne History uebernommen (das Projekt lag bisher in einem eigenen Repository unter Sensors/Laerm/laerm_readfromcsv_to_database). Stand entspricht dort e33a7b7: - Zeitstempel werden in UTC gespeichert, unabhaengig von der Zeitzone der Maschine. readin/parse.js macht das seit jeher richtig, readFromcsv.js lag zwei Stunden daneben. - Fehlgeschlagene Influx-Writes werden gemeldet und setzen den Exit-Code, statt still verloren zu gehen. - Der Container nimmt Parameter entgegen (ENTRYPOINT in Exec-Form). Hinweis: mongo.js, influx_post.js und logit.js gibt es auch unter readin/, mit abweichendem Stand. Zusammenfuehren waere der naechste Schritt. Co-Authored-By: Claude Opus 5 (1M context) --- readarchive/.gitignore | 8 + readarchive/Dockerfile_rfcsv | 27 ++ readarchive/build_and_copy.sh | 63 ++++ readarchive/crontab.tmp | 2 + readarchive/docker-compose.yml | 56 +++ readarchive/influx_post.js | 107 ++++++ readarchive/logit.js | 12 + readarchive/mongo.js | 197 ++++++++++ readarchive/mqtt.js | 33 ++ readarchive/package-lock.json | 356 ++++++++++++++++++ readarchive/package.json | 23 ++ readarchive/readFromcsv.js | 539 +++++++++++++++++++++++++++ readarchive/utilities/checkprops.js | 134 +++++++ readarchive/utilities/logit.js | 16 + readarchive/utilities/reporterror.js | 20 + 15 files changed, 1593 insertions(+) create mode 100644 readarchive/.gitignore create mode 100644 readarchive/Dockerfile_rfcsv create mode 100755 readarchive/build_and_copy.sh create mode 100644 readarchive/crontab.tmp create mode 100644 readarchive/docker-compose.yml create mode 100644 readarchive/influx_post.js create mode 100644 readarchive/logit.js create mode 100644 readarchive/mongo.js create mode 100644 readarchive/mqtt.js create mode 100644 readarchive/package-lock.json create mode 100644 readarchive/package.json create mode 100644 readarchive/readFromcsv.js create mode 100644 readarchive/utilities/checkprops.js create mode 100644 readarchive/utilities/logit.js create mode 100644 readarchive/utilities/reporterror.js diff --git a/readarchive/.gitignore b/readarchive/.gitignore new file mode 100644 index 0000000..122238e --- /dev/null +++ b/readarchive/.gitignore @@ -0,0 +1,8 @@ +node_modules +.DS* +.idea +.env +certs +data +.vscode + diff --git a/readarchive/Dockerfile_rfcsv b/readarchive/Dockerfile_rfcsv new file mode 100644 index 0000000..ee1be95 --- /dev/null +++ b/readarchive/Dockerfile_rfcsv @@ -0,0 +1,27 @@ +FROM node:22-alpine + +ADD package.json package-lock.json /tmp/ +RUN cd /tmp && npm ci --omit=dev +RUN mkdir -p /opt/app && cp -a /tmp/node_modules /tmp/package.json /opt/app/ +WORKDIR /opt/app +ADD *.js /opt/app/ +RUN mkdir /opt/app/utilities +ADD utilities/*.js /opt/app/utilities + +# readFromcsv.js legt hier die Sensorliste des Tages ab (data/list.json) +RUN mkdir -p /opt/app/data + +# betrifft nur noch die Uhrzeit in den Log-Zeilen - die Messwerte selbst +# werden seit dem UTC-Fix unabhaengig von der Zeitzone der Maschine gespeichert +RUN apk add --no-cache tzdata +ENV TZ=Europe/Berlin +RUN ln -snf /usr/share/zoneinfo/$TZ /etc/localtime && echo $TZ > /etc/timezone + +# Exec-Form: nur so werden Argumente aus "docker run ..." an das +# Programm durchgereicht. CMD ist der Default, den eigene Argumente ersetzen. +# +# docker run --rm --env-file .env -e INFLUXHOST=influxdb -e MONGOHOST=mongo \ +# --network laerm_readfromcsv_to_database_default rfcsv -t noise -s 2026-07-25 +# +ENTRYPOINT ["node", "./readFromcsv.js"] +CMD ["-t", "noise"] diff --git a/readarchive/build_and_copy.sh b/readarchive/build_and_copy.sh new file mode 100755 index 0000000..bb34212 --- /dev/null +++ b/readarchive/build_and_copy.sh @@ -0,0 +1,63 @@ +#!/bin/bash +# Build Docker-Container +# +# Call: buildit.sh name [target] +# +# The Dockerfile must be named like Dockerfile_name +# +# 2018-09-20 rxf +# - before sending docker image to remote, tag actual remote image +# +# 2018-09-14 rxf +# - first Version +# + +set -x +port="" +orgName=rfcsv +name=rfcsv + +usage() +{ + echo "Usage build_and_copy.sh [-p port] [-n name] target" + echo " Build docker container $name and copy to target" + echo "Params:" + echo " target: Where to copy the container to " + echo " -p port: ssh port (default 22)" + echo " -n name: new name for container (default: $orgName)" +} + +while getopts n:p:h? o +do + case "$o" in + n) name="$OPTARG";; + p) port="-p $OPTARG";; + h) usage; exit 0;; + *) usage; exit 1;; + esac +done +shift $((OPTIND-1)) + +while [ $# -gt 0 ]; do + if [[ -z "$target" ]]; then + target=$1 + shift + else + echo "bad option $1" + # exit 1 + shift + fi +done + +docker build -f Dockerfile_$orgName -t $name . + +dat=`date +%Y%m%d%H%M` + +if [ "$target" == "localhost" ] +then + docker tag $name $name:V_$dat + exit +fi + +ssh $port $target "docker tag $name $name:V_$dat" +docker save $name | bzip2 | pv | ssh $port $target 'bunzip2 | docker load' diff --git a/readarchive/crontab.tmp b/readarchive/crontab.tmp new file mode 100644 index 0000000..9aa5e15 --- /dev/null +++ b/readarchive/crontab.tmp @@ -0,0 +1,2 @@ +#33 10 * * * cd /opt/app && node ./readFromcsv.js -t laerm >>/var/log/readFrom.log 2>&1 +#53 10 * * * cd /opt/app && node ./readFromcsv.js -t radia >>/var/log/readFrom.log 2>&1 diff --git a/readarchive/docker-compose.yml b/readarchive/docker-compose.yml new file mode 100644 index 0000000..46f2f3a --- /dev/null +++ b/readarchive/docker-compose.yml @@ -0,0 +1,56 @@ +version: "3.9" +volumes: + mongo_vol: + influx_vol: + +services: + # rfcsv: + # image: rfcsv + # volumes: + # - /var/log/rfcsv:/var/log + # - /home/rxf/CERTS:certs + # environment: + # - MONGOHOST=207.180.224.98 + # - MONGOPORT=20019 + # #- MONGOAUTH=true + # #- "MONGOUSRP=rxf:5C5dB|m" + # - MONGOBASE=allsensors + # - TZ=Europe/Berlin + # container_name: rfcsv + # restart: always + + mongodb: + image: mongo + volumes: + - mongo_vol:/data/db + ports: + - "27017:27017" + container_name: mongo + #environment: + # - MONGO_INITDB_ROOT_USERNAME=${MONGO_ROOT_USER} + # - MONGO_INITDB_ROOT_PASSWORD=${MONGO_ROOT_PASSWD} + #command: '--auth' + restart: unless-stopped + + influx: + image: influxdb:2.7 + ports: + - '8086:8086' + volumes: + # /etc/influxdb2 NICHT mounten: das Entrypoint legt dort beim Setup das + # CLI-Profil "default" an. Ein vorhandenes Profil laesst den Setup + # scheitern - und danach loescht das Entrypoint bolt+engine im Volume. + - influx_vol:/var/lib/influxdb2 + environment: + # Werte kommen aus .env, damit App und Container dieselben benutzen + - DOCKER_INFLUXDB_INIT_MODE=setup + - DOCKER_INFLUXDB_INIT_USERNAME=${DOCKER_INFLUXDB_INIT_USERNAME} + - DOCKER_INFLUXDB_INIT_PASSWORD=${DOCKER_INFLUXDB_INIT_PASSWORD} + - DOCKER_INFLUXDB_INIT_ORG=${DOCKER_INFLUXDB_INIT_ORG} + - DOCKER_INFLUXDB_INIT_BUCKET=${DOCKER_INFLUXDB_INIT_BUCKET} + - DOCKER_INFLUXDB_INIT_ADMIN_TOKEN=${INFLUXTOKEN} + restart: + unless-stopped + container_name: influxdb + + diff --git a/readarchive/influx_post.js b/readarchive/influx_post.js new file mode 100644 index 0000000..555d137 --- /dev/null +++ b/readarchive/influx_post.js @@ -0,0 +1,107 @@ +// Access to influxDB vie HTTP + +import axios from 'axios' +import { logit, logerror } from './logit.js' +import { DateTime } from 'luxon' + +let INFLUXHOST = process.env.INFLUXHOST || "localhost" +let INFLUXPORT = process.env.INFLUXPORT || 8086 +let INFLUXTOKEN = process.env.INFLUXTOKEN || 'empty' +let INFLUXDATABUCKET = process.env.INFLUXDATABUCKET || "sensor_data" +let INFLUXORG = process.env.INFLUXORG || "citysensor" + +const INFLUXURL_READ = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/query?org=${INFLUXORG}` +const INFLUXURL_WRITE = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/write?org=${INFLUXORG}&bucket=${INFLUXDATABUCKET}&precision=ms` + +// `${e}` liefert bei einem AggregateError nur "AggregateError" - die einzelnen +// Verbindungsfehler (je einer pro aufgelöster 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}` +} + +export const influxRead = async (query) => { + let start = DateTime.now() + let data = [] + try { + let ret = await axios({ + method: 'post', + url: INFLUXURL_READ, + data: query, + headers: { + Authorization: `Token ${INFLUXTOKEN}`, + Accept: 'application/csv', + 'Content-type': 'application/vnd.flux' + }, + timeout: 10000, + }) + if (ret.status != 200) { + logerror(`doReadfromAPI Status: ${ret.status}`) + } + data = ret.data + } catch (e) { + logerror(`doReadfromAPI ${INFLUXURL_READ} ${describeError(e)}`) + } + logit(`ReadIn-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`) + return data +} + + +// liefert true, wenn Influx die Daten uebernommen hat, sonst false +export const influxWrite = async (data) => { + let start = DateTime.now() + let ok = false + try { + const ret = await axios({ + method: 'post', + url: INFLUXURL_WRITE, + data: data, + headers: { + Authorization: `Token ${INFLUXTOKEN}`, + 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 ${INFLUXURL_WRITE} ${describeError(e)}`) + } + logit(`Influx-Write-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`) + return ok +} + +/* +async function main() { + let data = ` + pm,sid=140 P1=12,P2=13 + pm,sid=142 P1=42,P2=13 + pm,sid=143 P1=43,P2=13 + pm,sid=144 P1=44,P2=13 + thp,sid=141 temperature=23.5,humidity=48,pressure=998 + ` + let ret = await influxWrite(data) + process.exit() + + let query = `from(bucket:"sensor_data") +|> range(start: -1mo) +|> filter(fn: (r) => r._measurement == "pm") +|> filter(fn: (r) => r.sid == "140") +` + let erg = await influxRead(query) + console.log(erg) +} + + +main().catch(console.error) +*/ diff --git a/readarchive/logit.js b/readarchive/logit.js new file mode 100644 index 0000000..19d8b28 --- /dev/null +++ b/readarchive/logit.js @@ -0,0 +1,12 @@ +import { DateTime} from 'luxon' + +export function logit(str) { + let s = `${DateTime.now().toISO()} => ${str}`; + console.log(s); +} + +export function logerror(str) { + let s = `${DateTime.utc().toISO()} => *** ERROR *** ${str}`; + console.log(s); +} + diff --git a/readarchive/mongo.js b/readarchive/mongo.js new file mode 100644 index 0000000..311f435 --- /dev/null +++ b/readarchive/mongo.js @@ -0,0 +1,197 @@ +/* Interface for MongoDB +*/ +import { MongoClient } from 'mongodb' +import { logit, logerror } from './utilities/logit.js' +import { DateTime } from 'luxon' +import {returnOnError} from "./utilities/reporterror.js"; + +const DEVELOP = process.env.DEVELOP || false; + +const MONGOHOST = process.env.MONGOHOST || 'localhost'; +const MONGOPORT = process.env.MONGOPORT || 27017; +const MONGOAUTH = process.env.MONGOAUTH || false; +const MONGOUSRP = process.env.MONGOUSRP || ""; +const MONGOBASE = process.env.MONGOBASE || 'sensor_data'; + +const MONGO_URL = MONGOAUTH ? 'mongodb://'+MONGOUSRP+'@' + MONGOHOST + ':' + MONGOPORT + '/?authSource=admin' : 'mongodb://'+MONGOHOST+':'+MONGOPORT; // URL to mongo database +if (DEVELOP) { + console.log(`MongoURL = "${MONGO_URL}" and Database = ${MONGOBASE}`); +} +const statistics = {} + +export const properties_collection = 'properties' +const data_collection = 'sensors' + +export const connectMongo = async () => { + try { + if(DEVELOP === 'true') { + logit(`Try to connect to ${MONGO_URL}`) + } else { + logit(`Try to connect to ${'mongodb://' + MONGOHOST + ':' + MONGOPORT}`) + } + let client = await MongoClient.connect(MONGO_URL) + if ( DEVELOP === 'true') { + logit(`Mongodbase connected to ${MONGO_URL}`) + } else { + logit('Mongodbase connected') + } + return client + } + catch (error) { + throw (error) + } +} + +export const closeMongo = async (client) => { + try { + await client.close() + } + catch(error){ + throw(error) + } +} + +const listDatabases = async (client) => { + let databasesList = await client.db().admin().listDatabases(); + + console.log("Databases:"); + databasesList.databases.forEach(db => console.log(` - ${db.name}`)); +} + +/* *************************************************** +// READ routines +******************************************************/ + +// Read properties from the database +export const readProperties = async (query, limit = 0) => { + let ret = {err: null, properties: null} + let client = await connectMongo() + try { + if ("sid" in query) { // if sid is given, read property for sid + ret.properties = await client.db(MONGOBASE).collection(properties_collection).findOne({_id: query.sid}) + } else { // otherwise read props corresponding to query + ret.properties = await client.db(MONGOBASE).collection(properties_collection).find(query).limit(limit).toArray() + } + } catch (e) { + ret.err = e + } + finally { + client.close() + } + return ret +} + + +export const getallProperties = async (client) => { + let ret = {error: false, errortext: '', properties: []} + try + { + ret.properties = await client.db(MONGOBASE).collection(properties_collection) + .find().sort({_id: 1}).toArray() + } + catch(e) { + ret = {error: true, errortext: e} + } + return ret +} + + +export const getOneproperty = async (client, sid) => { + let ret = {error: false, errortext: '', property: null} + try { + ret.property = await client.db(MONGOBASE).collection(properties_collection) + .findOne({_id: sid}) + } catch (e) { + ret = {error: true, errortext: e} + } + return ret +} + +const getCollName = (coll, st) => { + return st + '_' + coll +} + +export const getOneSensorOneday = async (client, sid, day, styp) => { + let ret = {error: false, errortext: '', date: []} + // muss dieselbe Zone benutzen wie die gespeicherten Zeitstempel (UTC), + // sonst passt das Suchfenster nicht auf die Daten des Tages + let d = DateTime.fromFormat(day, "yyyy-LL-dd", { zone: 'utc' }) + let start = d.startOf('day').toJSDate() + let end = d.startOf('day').plus({day:1}).toJSDate() + try { + let erg = await client.db(MONGOBASE).collection(getCollName(data_collection, styp)) + .find({sensorid: sid, datetime: {$gte: start, $lt: end}},{sort: {datetime: 1}}).toArray() + ret.data = erg + } catch(e) { + ret.error = true + ret.errortext = e + } + return ret +} + +export const writeOneSensor = async (client, data, styp) => { + let ret = {error: false, errortext: '', inserted: 1} + let erg + try { + erg = await client.db(MONGOBASE).collection(getCollName(data_collection, styp)) + .insertOne(data) + } catch (e) { + ret.error = true + ret.errortext = e + ret.inserted = 0 + } +return ret +} + +export const writeDataArray = async (client, data, styp) => { + let result + let start = DateTime.now(); + try { + result = await client.db(MONGOBASE).collection(getCollName(data_collection, styp)) + .insertMany(data, {ordered: false}) + } catch (e) { + if(e.code !== 11000) { + console.error(e) + } + } + let statname = `writeData${getCollName(data_collection, styp)}Time` + statistics[statname] = DateTime.now().diff(start, ['seconds']).toObject().seconds + logit(`Write Data for ${getCollName(data_collection, styp)} to mongoDB: Time: ${statistics[statname]} sec.`) +} + +export const bulkWrite = async (client, coll, data) => { + let start = DateTime.now() + let result + try { + result = await client.db(MONGOBASE).collection(coll) + .bulkWrite(data, {ordered: false}) + } catch (e) { + console.error(e) + } + logit(`Write Properties: Result: ${result}, Time: ${start.diffNow('second').toObject().seconds * -1} sec.`) + return result +} + +export const getLocationIDs = async (client, stype) => { + let ret = {error: false, errortext: '', locations: []} + try { + ret.locations = await client.db(MONGOBASE).collection(properties_collection) + .distinct('location.0.id', {type: stype}) + } catch (e) { + ret.error = true + ret.errortext = e + } + return ret +} + +export const getTHPSensors = async (client, lid) => { + let ret = {error: false, errortext: '', sensors: []} + try { + ret.sensors = await client.db(MONGOBASE).collection(properties_collection) + .findOne({type: 'thp', 'location.0.id': lid}, { projection: {_id: 1}}) + } catch (e) { + ret.error = true + ret.errortext = e + } + return ret +} \ No newline at end of file diff --git a/readarchive/mqtt.js b/readarchive/mqtt.js new file mode 100644 index 0000000..0624e1a --- /dev/null +++ b/readarchive/mqtt.js @@ -0,0 +1,33 @@ +// MQTT interface + +import * as mqtt from 'mqtt' +import { logit, logerror } from './logit.js' + +const MQTTHOST = process.env.MQTTHOST || 'rexfue.de' +const MQTTPORT = process.env.MQTTPORT || 1883 +const MQTTUSER = process.env.MQTTUSER || 'stzuhr' +const MQTTPASSWD = process.env.MQTTPASSWD || '74chQCYb' + +const client = mqtt.connect(`mqtt://${MQTTHOST}:${MQTTPORT}`, {username: MQTTUSER, password: MQTTPASSWD } ) +export let connected = false + +export const doConnect = () => { +// client = mqtt.connect(`mqtt://${MQTTHOST}:${MQTTPORT}`, {user: MQTTUSER, password: MQTTPASSWD } ) +} + +client.on('error', (e) => { + logerror(`connect to MQTT ${e}`) +}) + +client.on('connect', () => { + connected = true + logit(`Connected to MQTT Broker ${MQTTHOST}`) +}) + +export const doPublish = (topic, msg) => { + client.publish(topic, msg) +} + +export const doClose = () => { + client.end() +} diff --git a/readarchive/package-lock.json b/readarchive/package-lock.json new file mode 100644 index 0000000..1ef4b31 --- /dev/null +++ b/readarchive/package-lock.json @@ -0,0 +1,356 @@ +{ + "name": "sensors_readfromcsv", + "version": "3.1.2", + "lockfileVersion": 3, + "requires": true, + "packages": { + "": { + "name": "sensors_readfromcsv", + "version": "3.1.2", + "license": "ISC", + "dependencies": { + "axios": "^1.4.0", + "csvtojson": "^2.0.10", + "dotenv": "^17.4.2", + "luxon": "^3.3.0", + "moment": "^2.24.0", + "mongodb": "^5.5.0", + "posix-getopt": "^1.2.1" + } + }, + "node_modules/@mongodb-js/saslprep": { + "version": "1.1.1", + "resolved": "https://registry.npmjs.org/@mongodb-js/saslprep/-/saslprep-1.1.1.tgz", + "integrity": "sha512-t7c5K033joZZMspnHg/gWPE4kandgc2OxE74aYOtGKfgB9VPuVJPix0H6fhmm2erj5PBJ21mqcx34lpIGtUCsQ==", + "optional": true, + "dependencies": { + "sparse-bitfield": "^3.0.3" + } + }, + "node_modules/@types/node": { + "version": "20.2.1", + "resolved": "https://registry.npmjs.org/@types/node/-/node-20.2.1.tgz", + "integrity": "sha512-DqJociPbZP1lbZ5SQPk4oag6W7AyaGMO6gSfRwq3PWl4PXTwJpRQJhDq4W0kzrg3w6tJ1SwlvGZ5uKFHY13LIg==" + }, + "node_modules/@types/webidl-conversions": { + "version": "7.0.0", + "resolved": "https://registry.npmjs.org/@types/webidl-conversions/-/webidl-conversions-7.0.0.tgz", + "integrity": "sha512-xTE1E+YF4aWPJJeUzaZI5DRntlkY3+BCVJi0axFptnjGmAoWxkyREIh/XMrfxVLejwQxMCfDXdICo0VLxThrog==" + }, + "node_modules/@types/whatwg-url": { + "version": "8.2.2", + "resolved": "https://registry.npmjs.org/@types/whatwg-url/-/whatwg-url-8.2.2.tgz", + "integrity": "sha512-FtQu10RWgn3D9U4aazdwIE2yzphmTJREDqNdODHrbrZmmMqI0vMheC/6NE/J1Yveaj8H+ela+YwWTjq5PGmuhA==", + "dependencies": { + "@types/node": "*", + "@types/webidl-conversions": "*" + } + }, + "node_modules/asynckit": { + "version": "0.4.0", + "resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz", + "integrity": "sha512-Oei9OH4tRh0YqU3GxhX79dM/mwVgvbZJaSNaRk+bshkj0S5cfHcgYakreBjrHwatXKbz+IoIdYLxrKim2MjW0Q==" + }, + "node_modules/axios": { + "version": "1.4.0", + "resolved": "https://registry.npmjs.org/axios/-/axios-1.4.0.tgz", + "integrity": "sha512-S4XCWMEmzvo64T9GfvQDOXgYRDJ/wsSZc7Jvdgx5u1sd0JwsuPLqb3SYmusag+edF6ziyMensPVqLTSc1PiSEA==", + "dependencies": { + "follow-redirects": "^1.15.0", + "form-data": "^4.0.0", + "proxy-from-env": "^1.1.0" + } + }, + "node_modules/axios/node_modules/form-data": { + "version": "4.0.0", + "resolved": "https://registry.npmjs.org/form-data/-/form-data-4.0.0.tgz", + "integrity": "sha512-ETEklSGi5t0QMZuiXoA/Q6vcnxcLQP5vdugSpuAyi6SVGi2clPPp+xgEhuMaHC+zGgn31Kd235W35f7Hykkaww==", + "dependencies": { + "asynckit": "^0.4.0", + "combined-stream": "^1.0.8", + "mime-types": "^2.1.12" + }, + "engines": { + "node": ">= 6" + } + }, + "node_modules/bluebird": { + "version": "3.7.2", + "resolved": "https://registry.npmjs.org/bluebird/-/bluebird-3.7.2.tgz", + "integrity": "sha512-XpNj6GDQzdfW+r2Wnn7xiSAd7TM3jzkxGXBGTtWKuSXv1xUV+azxAm8jdWZN06QTQk+2N2XB9jRDkvbmQmcRtg==" + }, + "node_modules/bson": { + "version": "5.5.1", + "resolved": "https://registry.npmjs.org/bson/-/bson-5.5.1.tgz", + "integrity": "sha512-ix0EwukN2EpC0SRWIj/7B5+A6uQMQy6KMREI9qQqvgpkV2frH63T0UDVd1SYedL6dNCmDBYB3QtXi4ISk9YT+g==", + "engines": { + "node": ">=14.20.1" + } + }, + "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==", + "dependencies": { + "delayed-stream": "~1.0.0" + }, + "engines": { + "node": ">= 0.8" + } + }, + "node_modules/csvtojson": { + "version": "2.0.10", + "resolved": "https://registry.npmjs.org/csvtojson/-/csvtojson-2.0.10.tgz", + "integrity": "sha512-lUWFxGKyhraKCW8Qghz6Z0f2l/PqB1W3AO0HKJzGIQ5JRSlR651ekJDiGJbBT4sRNNv5ddnSGVEnsxP9XRCVpQ==", + "dependencies": { + "bluebird": "^3.5.1", + "lodash": "^4.17.3", + "strip-bom": "^2.0.0" + }, + "bin": { + "csvtojson": "bin/csvtojson" + }, + "engines": { + "node": ">=4.0.0" + } + }, + "node_modules/delayed-stream": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/delayed-stream/-/delayed-stream-1.0.0.tgz", + "integrity": "sha512-ZySD7Nf91aLB0RxL4KGrKHBXl7Eds1DAmEdcoVawXnLD7SDhpNgtuII2aAkg7a7QS41jxPSZ17p4VdGnMHk3MQ==", + "engines": { + "node": ">=0.4.0" + } + }, + "node_modules/dotenv": { + "version": "17.4.2", + "resolved": "https://registry.npmjs.org/dotenv/-/dotenv-17.4.2.tgz", + "integrity": "sha512-nI4U3TottKAcAD9LLud4Cb7b2QztQMUEfHbvhTH09bqXTxnSie8WnjPALV/WMCrJZ6UV/qHJ6L03OqO3LcdYZw==", + "license": "BSD-2-Clause", + "engines": { + "node": ">=12" + }, + "funding": { + "url": "https://dotenvx.com" + } + }, + "node_modules/follow-redirects": { + "version": "1.15.2", + "resolved": "https://registry.npmjs.org/follow-redirects/-/follow-redirects-1.15.2.tgz", + "integrity": "sha512-VQLG33o04KaQ8uYi2tVNbdrWp1QWxNNea+nmIB4EVM28v0hmP17z7aG1+wAkNzVq4KeXTq3221ye5qTJP91JwA==", + "funding": [ + { + "type": "individual", + "url": "https://github.com/sponsors/RubenVerborgh" + } + ], + "engines": { + "node": ">=4.0" + }, + "peerDependenciesMeta": { + "debug": { + "optional": true + } + } + }, + "node_modules/ip": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/ip/-/ip-2.0.0.tgz", + "integrity": "sha512-WKa+XuLG1A1R0UWhl2+1XQSi+fZWMsYKffMZTTYsiZaUD8k2yDAj5atimTUD2TZkyCkNEeYE5NhFZmupOGtjYQ==" + }, + "node_modules/is-utf8": { + "version": "0.2.1", + "resolved": "https://registry.npmjs.org/is-utf8/-/is-utf8-0.2.1.tgz", + "integrity": "sha512-rMYPYvCzsXywIsldgLaSoPlw5PfoB/ssr7hY4pLfcodrA5M/eArza1a9VmTiNIBNMjOGr1Ow9mTyU2o69U6U9Q==" + }, + "node_modules/lodash": { + "version": "4.17.21", + "resolved": "https://registry.npmjs.org/lodash/-/lodash-4.17.21.tgz", + "integrity": "sha512-v2kDEe57lecTulaDIuNTPy3Ry4gLGJ6Z1O3vE1krgXZNrsQ+LFTGHVxVjcXPs17LhbZVGedAJv8XZ1tvj5FvSg==" + }, + "node_modules/luxon": { + "version": "3.3.0", + "resolved": "https://registry.npmjs.org/luxon/-/luxon-3.3.0.tgz", + "integrity": "sha512-An0UCfG/rSiqtAIiBPO0Y9/zAnHUZxAMiCpTd5h2smgsj7GGmcenvrvww2cqNA8/4A5ZrD1gJpHN2mIHZQF+Mg==", + "engines": { + "node": ">=12" + } + }, + "node_modules/memory-pager": { + "version": "1.5.0", + "resolved": "https://registry.npmjs.org/memory-pager/-/memory-pager-1.5.0.tgz", + "integrity": "sha512-ZS4Bp4r/Zoeq6+NLJpP+0Zzm0pR8whtGPf1XExKLJBAczGMnSi3It14OiNCStjQjM6NU1okjQGSxgEZN8eBYKg==", + "optional": true + }, + "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==", + "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==", + "dependencies": { + "mime-db": "1.52.0" + }, + "engines": { + "node": ">= 0.6" + } + }, + "node_modules/moment": { + "version": "2.29.4", + "resolved": "https://registry.npmjs.org/moment/-/moment-2.29.4.tgz", + "integrity": "sha512-5LC9SOxjSc2HF6vO2CyuTDNivEdoz2IvyJJGj6X8DJ0eFyfszE0QiEd+iXmBvUP3WHxSjFH/vIsA0EN00cgr8w==", + "engines": { + "node": "*" + } + }, + "node_modules/mongodb": { + "version": "5.9.1", + "resolved": "https://registry.npmjs.org/mongodb/-/mongodb-5.9.1.tgz", + "integrity": "sha512-NBGA8AfJxGPeB12F73xXwozt8ZpeIPmCUeWRwl9xejozTXFes/3zaep9zhzs1B/nKKsw4P3I4iPfXl3K7s6g+Q==", + "dependencies": { + "bson": "^5.5.0", + "mongodb-connection-string-url": "^2.6.0", + "socks": "^2.7.1" + }, + "engines": { + "node": ">=14.20.1" + }, + "optionalDependencies": { + "@mongodb-js/saslprep": "^1.1.0" + }, + "peerDependencies": { + "@aws-sdk/credential-providers": "^3.188.0", + "@mongodb-js/zstd": "^1.0.0", + "kerberos": "^1.0.0 || ^2.0.0", + "mongodb-client-encryption": ">=2.3.0 <3", + "snappy": "^7.2.2" + }, + "peerDependenciesMeta": { + "@aws-sdk/credential-providers": { + "optional": true + }, + "@mongodb-js/zstd": { + "optional": true + }, + "kerberos": { + "optional": true + }, + "mongodb-client-encryption": { + "optional": true + }, + "snappy": { + "optional": true + } + } + }, + "node_modules/mongodb-connection-string-url": { + "version": "2.6.0", + "resolved": "https://registry.npmjs.org/mongodb-connection-string-url/-/mongodb-connection-string-url-2.6.0.tgz", + "integrity": "sha512-WvTZlI9ab0QYtTYnuMLgobULWhokRjtC7db9LtcVfJ+Hsnyr5eo6ZtNAt3Ly24XZScGMelOcGtm7lSn0332tPQ==", + "dependencies": { + "@types/whatwg-url": "^8.2.1", + "whatwg-url": "^11.0.0" + } + }, + "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==", + "engines": { + "node": "*" + } + }, + "node_modules/proxy-from-env": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/proxy-from-env/-/proxy-from-env-1.1.0.tgz", + "integrity": "sha512-D+zkORCbA9f1tdWRK0RaCR3GPv50cMxcrz4X8k5LTSUD1Dkw47mKJEZQNunItRTkWwgtaUSo1RVFRIG9ZXiFYg==" + }, + "node_modules/punycode": { + "version": "2.3.0", + "resolved": "https://registry.npmjs.org/punycode/-/punycode-2.3.0.tgz", + "integrity": "sha512-rRV+zQD8tVFys26lAGR9WUuS4iUAngJScM+ZRSKtvl5tKeZ2t5bvdNFdNHBW9FWR4guGHlgmsZ1G7BSm2wTbuA==", + "engines": { + "node": ">=6" + } + }, + "node_modules/smart-buffer": { + "version": "4.2.0", + "resolved": "https://registry.npmjs.org/smart-buffer/-/smart-buffer-4.2.0.tgz", + "integrity": "sha512-94hK0Hh8rPqQl2xXc3HsaBoOXKV20MToPkcXvwbISWLEs+64sBq5kFgn2kJDHb1Pry9yrP0dxrCI9RRci7RXKg==", + "engines": { + "node": ">= 6.0.0", + "npm": ">= 3.0.0" + } + }, + "node_modules/socks": { + "version": "2.7.1", + "resolved": "https://registry.npmjs.org/socks/-/socks-2.7.1.tgz", + "integrity": "sha512-7maUZy1N7uo6+WVEX6psASxtNlKaNVMlGQKkG/63nEDdLOWNbiUMoLK7X4uYoLhQstau72mLgfEWcXcwsaHbYQ==", + "dependencies": { + "ip": "^2.0.0", + "smart-buffer": "^4.2.0" + }, + "engines": { + "node": ">= 10.13.0", + "npm": ">= 3.0.0" + } + }, + "node_modules/sparse-bitfield": { + "version": "3.0.3", + "resolved": "https://registry.npmjs.org/sparse-bitfield/-/sparse-bitfield-3.0.3.tgz", + "integrity": "sha512-kvzhi7vqKTfkh0PZU+2D2PIllw2ymqJKujUcyPMd9Y75Nv4nPbGJZXNhxsgdQab2BmlDct1YnfQCguEvHr7VsQ==", + "optional": true, + "dependencies": { + "memory-pager": "^1.0.2" + } + }, + "node_modules/strip-bom": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/strip-bom/-/strip-bom-2.0.0.tgz", + "integrity": "sha512-kwrX1y7czp1E69n2ajbG65mIo9dqvJ+8aBQXOGVxqwvNbsXdFM6Lq37dLAY3mknUwru8CfcCbfOLL/gMo+fi3g==", + "dependencies": { + "is-utf8": "^0.2.0" + }, + "engines": { + "node": ">=0.10.0" + } + }, + "node_modules/tr46": { + "version": "3.0.0", + "resolved": "https://registry.npmjs.org/tr46/-/tr46-3.0.0.tgz", + "integrity": "sha512-l7FvfAHlcmulp8kr+flpQZmVwtu7nfRV7NZujtN0OqES8EL4O4e0qqzL0DC5gAvx/ZC/9lk6rhcUwYvkBnBnYA==", + "dependencies": { + "punycode": "^2.1.1" + }, + "engines": { + "node": ">=12" + } + }, + "node_modules/webidl-conversions": { + "version": "7.0.0", + "resolved": "https://registry.npmjs.org/webidl-conversions/-/webidl-conversions-7.0.0.tgz", + "integrity": "sha512-VwddBukDzu71offAQR975unBIGqfKZpM+8ZX6ySk8nYhVoo5CYaZyzt3YBvYtRtO+aoGlqxPg/B87NGVZ/fu6g==", + "engines": { + "node": ">=12" + } + }, + "node_modules/whatwg-url": { + "version": "11.0.0", + "resolved": "https://registry.npmjs.org/whatwg-url/-/whatwg-url-11.0.0.tgz", + "integrity": "sha512-RKT8HExMpoYx4igMiVMY83lN6UeITKJlBQ+vR/8ZJ8OCdSiN3RwCq+9gH0+Xzj0+5IrM6i4j/6LuvzbZIQgEcQ==", + "dependencies": { + "tr46": "^3.0.0", + "webidl-conversions": "^7.0.0" + }, + "engines": { + "node": ">=12" + } + } + } +} diff --git a/readarchive/package.json b/readarchive/package.json new file mode 100644 index 0000000..d4ca0a8 --- /dev/null +++ b/readarchive/package.json @@ -0,0 +1,23 @@ +{ + "name": "sensors_readfromcsv", + "version": "3.1.2", + "date": "2023-12-23", + "description": "", + "main": "readfromcsv.js", + "scripts": { + "test": "echo \"Error: no test specified\" && exit 1" + }, + "type": "module", + "keywords": [], + "author": "Reinhard X. Fürst", + "license": "ISC", + "dependencies": { + "axios": "^1.4.0", + "csvtojson": "^2.0.10", + "dotenv": "^17.4.2", + "luxon": "^3.3.0", + "moment": "^2.24.0", + "mongodb": "^5.5.0", + "posix-getopt": "^1.2.1" + } +} diff --git a/readarchive/readFromcsv.js b/readarchive/readFromcsv.js new file mode 100644 index 0000000..bb8ae1c --- /dev/null +++ b/readarchive/readFromcsv.js @@ -0,0 +1,539 @@ +// readFromcsv - Alte Daten vom Luftsdaten per CSV einlesen +// rxf 2017-12-12 +// +// Ausgehend von Version 2.0.0 werden hier nun die Daten nach den Sensortypen getrennt gespeichert +// und zwar einstellbar in einer Mongo-DB oder in einer Influx-DB (oder auch in beide). + +// Version: +// +// V 3.1.1 2023-11-14 rxf +// - Enddatum eingeführt +// +// V 3.1.0 2023-11-01 rxf +// - Anpassung, so dass es als Docker-Container mit dem timeseries-Container läuft +// +// V 3.0.1 2023-10-18 rxf +// - mehrer Typen gleichzeitig auswählbar + +// V 3.0.0 2023-10-07 rxf +// - Erste Version mit der Trennung nach Typen +// +// TYP: Ist der Sensortyp, der eingelesen wird. Default ist 'pm', d.h. alle Feinstaubsensoren. Mögliche Werte sind +// ['radioactivity', 'noise']. Die jeweils dazugehörenden THP-Sensoren werden mit eingelesen. +const TYP = process.env.TYP || 'noise' +const START = process.env.START +const END = process.env.END +const DBASE = process.env.DBASE || 'both' + +const DATABASE = ['mongo', 'influx', 'both'] + +const DEVELOP = process.env.DEVELOP || false; +const LIVE = (process.env.LIVE == "true") || true + +const MAXTRIES = 10 + +// MUSS der erste Import bleiben: mongo.js und influx_post.js lesen ihre +// Konfiguration beim Laden aus process.env - .env muss vorher drin sein +import 'dotenv/config' +import axios from 'axios' +import https from 'https' +import * as mongo from './mongo.js' +import * as influx from './influx_post.js' +import fs from 'fs' +import csv from 'csvtojson' +import pkg from './package.json' with { type: "json" } +import mod_getopt from 'posix-getopt' +import { logit, logerror } from './utilities/logit.js' +import { DateTime, Duration } from 'luxon' +import { checkProperties} from "./utilities/checkprops.js"; +import { version } from "os" + +const DEFAULTPORT = 3008; +const PORT = process.env.SERVERPORT || DEFAULTPORT; // Port for server + +const API_URL = 'https://archive.sensor.community/'; // URL to API on 'luftdaten.info' + +const valueTypes = ["P1","P2","P0","P4","N0","N1","N4","N25","N10","TS", "co2_ppm", "temperature","humidity","pressure","pressure_at_sealevel","noise_LA_min","noise_LA_max","noise_LAeq","counts_per_minute","hv_pulses","counts","sample_time_ms"] + +const newProps = [] + +// Sensor/Tag-Paare, deren Influx-Write fehlgeschlagen ist - Auswertung am Ende von main() +const failedInfluxWrites = [] + +const saveList = (list) => { + fs.writeFileSync('data/list.json', JSON.stringify(list)) +} + +const getListFromDisk = () => { + let buffer = fs.readFileSync('data/list.json') + return JSON.parse(buffer) +} + +const types = { + P1: 'pm', P2: 'pm', P0: 'pm', + temperature: 'thp', humidity: 'thp', pressure: 'thp', + noise_LAeq: 'noise', noise_LA_max: 'noise', noise_LA_min: 'noise', + counts_per_minute: 'radioactivity', + lat: 'gps' +}; +const sensorTypes = { + "bme280": "thp", + "bmp180": "thp", + "bmp280": "thp", + "dht22": "thp", + "ds18b20": "thp", + "hpm": "pm", + "htu21d": "thp", + "laerm": "noise", + "nextpm": "pm", + "pms1003": "pm", + "pms3003": "pm", + "pms5003": "pm", + "pms6003": "pm", + "pms7003": "pm", + "ppd42ns": "pm", + "radiation": "radioactivity", + "scd30": "co2", + "sds011": "pm", + "sht11": "thp", + "sht15": "thp", + "sht30": "thp", + "sht31": "thp", + "sht35": "thp", + "sht85": "thp", + "sps30": "pm", +} + +const getSensorType = (item) => { + let p = item.split('_') + let t = sensorTypes[p[1]] + if (t === undefined) { + console.log("Type not known: " + p[1]) + t = 'unknown' + } + return t +} + +const checkType = (styp, args) => { + if ((args.sensorType === styp) || (styp === 'thp')) { + return true + } + return false +} + +function getType(typ) { + if(typ in types) { + return types[typ]; + } else { + return 'unknown' + } +} + +// get the sensorids of all THP sensors belonging to the sensors of the give type +const getTHPList = async (client, stype) => { + let locids = await mongo.getLocationIDs(client, stype) // { sensorid: 12345, locid: 34567} + let thplist = [] + for (let lid of locids.locations) { + let x = await mongo.getTHPSensors(client, lid) + if (x.sensors) { + thplist.push(x.sensors._id) + } + } + return thplist +} + +// get list of all saved sensors for this day +const getdirlistOfOneDay = async (day) => { + let tries = 5 + let ret = { err: false, list : []} + + while(true) { + try { + const erg = await axios.get(API_URL + day); + // logit(`getdirListOfOneDay: Status = ${erg.status}`); + let a = erg.data.split('"'); // parse the list + for (let i = 0; i < a.length; i++) { + if (a[i].startsWith(day.substr(0, 4))) { // extract the sensor names + ret.list.push(a[i]); + } + } + logit(`Read ${ret.list.length} entries for ${day}`) + break + } catch (e) { + if(tries-- === 0) { + ret.err = e + return ret + } + } + } + return ret; +} + +const parseDayList = (list, thplist, args) => { + let sl = [] + for (let elem of list) { + let oneEntry = {url: elem, indoor: false} + let name = elem + if (elem.indexOf('indoor') != -1) { // there is 'indoor' + name = elem.replace('_indoor', ''); + oneEntry.indoor = true; + } + let n = name.lastIndexOf('_'); + oneEntry.sensorType = getSensorType(elem) + oneEntry.sensorid = parseInt(name.substr(n + 1).replace('.csv', '')) + if(!((oneEntry.sensorType === args.sensorType) || (oneEntry.sensorType === 'thp'))){ + continue + } + if ((oneEntry.sensorType === 'thp') && (thplist.indexOf(oneEntry.sensorid) === -1)) { + continue + } + sl.push(oneEntry) + } + return sl +} + + +// do the whole work +async function readSensorsperDay(client, args) { + let list + let insertedCount = 0 + let sensorList = [] + let st = DateTime.fromFormat(args.startDate, 'yyyy-LL-dd') // startdate + let end = DateTime.fromFormat(args.endDate, 'yyyy-LL-dd') + let thplist = await getTHPList(client, args.sensorType) // find all THP sensors belonging to the sensors of the given type + for (let d = st; d < end; d = d.plus({days: 1})) { // loop thru days + let dstr = d.toFormat('yyyy-LL-dd') + logit(`*************** ${dstr}`) // log every day + let mist = false; + // fetch sensors list of current day + if(LIVE) { + let erg = await getdirlistOfOneDay(dstr) + if (erg.err) { + logerror(erg.err) + mist = true + } else { + list = erg.list + }; + } else { + list = getListFromDisk() + logit(`Read ${list.length} entries from Disk for ${dstr}`) + } + if (mist) continue; // if day doesn't exist, continue + if (LIVE) { + saveList(list) + } + sensorList = parseDayList(list, thplist, args) + insertedCount += await enterSensors(client, sensorList, dstr, args); // fetch and enter sensor data + } + return insertedCount +} + +const scanOneSensor = async (client, sensorList, dt, args) => { + let inserted = 0 + let missed = [] + let listcnt = sensorList.length + for (let item of sensorList) { // iterate the list + let sid = item.sensorid + if(args.sensornbr !== undefined) { + if(args.sensornbr !== sid) { + continue + } + } + logit(`${item.url}`) + let icount = await putOneSensorInDb(client, item.url, sid, dt, missed, item.indoor, item.sensorType, args) // put one sensor data into DB + logit(`${listcnt} -- Sensor ${sid}: ${icount} Einträge`) + inserted += icount + listcnt-- + if (args.nbrOfEntries !== undefined) { + if (--args.nbrOfEntries === 0) { + break + } + + } + } + return {inserted: inserted, missed: missed} +} + +// Iterate thru the list and enter every sensor that matches data into db +async function enterSensors(client, sensorList, dt, args) { + let inserted = 0 + // check which types to scan + while (sensorList.length > 0) { + let missed = [] + let erg = await scanOneSensor(client, sensorList, dt, args) + inserted += erg.inserted + if(erg.missed.length > 0) { + logit(`Error reading ${erg.missed.length} missed files `) + } + logit(`enterSensors: ${inserted} entries handled`) + sensorList = erg.missed + } + return inserted +} + +// read CSV file and enter data +async function putOneSensorInDb(client, name, sid, dt, missed, indoor, styp, args) { + let z = 0 + try { + let [all, dataline] = await readOneSensorOneDay(client, name, sid, dt, missed, indoor, args); + const erg = await enterOneSensorinDB(client, sid, all, dataline, dt, styp, args) + z = erg === undefined ? 0 : erg + } + catch(err) { + console.log("Error in putOneSensorInDB()"); + } + return z +} + +// read the CSV-File and parse it into right format for DB +async function readOneSensorOneDay(client, name, sid, dt, missed, indoor, args) { + let all = []; + let url = API_URL + dt + '/' + name; // construct URL + const sname = name.split('_')[1] + let firstItem + let typ = 'unknown' + let dataline = '' + + try { + const erg = await axios(url); + if ((erg.status != 200) || (erg.data == "")) { + return {error: {status: erg.status, data: erg.data}}; + } + + const data = await csv({delimiter: ';'}).fromString(erg.data); + firstItem = data[0] + for (let item of data) { + let entry = {values: {}, sensorid: sid} + // Die Zeitstempel im CSV-Archiv sind UTC, haben aber keine Zonenangabe. + // Ohne zone:'utc' wuerde luxon sie in der Zeitzone der Maschine lesen. + let date = DateTime.fromISO(item.timestamp, { zone: 'utc' }) // extract date of entry + entry.datetime = date.toJSDate() // make date for Mongo (== ISODate) + for (let val of valueTypes) { + if (item[val] !== undefined) { + let x + let v = item[val] + try { + x = parseFloat(v) // convert value to float + if (Number.isNaN(x)) { + x = -9999.9 // default if value is invalid or unknown + } + } catch (err) { + console.log('Math parse float error on value'); + x = -9999.9 // default if value is invalid or unknown + } + if (typ === 'unknown') { // extract measurement type + typ = getType(val) + } + if(typ === 'noise') { + val = val.slice(6) + } + entry.values[val] = x + if (val == "LAeq") { + entry.values.E10tel_eq = Math.pow(10, entry.values.LAeq / 10); + } + } + } + all.push(entry) + if ((args.database === 'influx') || (args.database === 'both')) { + let influxvalue = '' + for (const [key, value] of Object.entries(entry.values)) { + influxvalue = influxvalue + `${key}=${value},` + } + influxvalue = influxvalue.slice(0,-1) + dataline += `${typ},sid=${entry.sensorid} ${influxvalue} ${entry.datetime.valueOf()}\n` + } + } + checkProperties(client, firstItem, indoor, sname, typ, dt, newProps) + } catch (e) { + console.log(e) + missed.push(name) + } + return [all, dataline] // return all the data +} + +// Neue Einträge aus dem CSV mit denen in der DB vergleichn, identische aus dem CSV raus streichen +const condenseEntry = async (csvData, dbData ) => { + for( let i = csvData.length - 1; i >= 0; i--) { + const cvt = csvData[i].datetime.valueOf() + if( dbData.findIndex( (item) => item.datetime.valueOf() == cvt) !== -1) { + csvData.splice(i,1) + } + } + return csvData +} + +// enter all data for one sensor into DB +async function enterOneSensorinDB(client, sid, erg, dataline, day, styp, args) { + const dataOneDay = await mongo.getOneSensorOneday(client, sid, day, styp) + if (dataOneDay.data.length > 0) { + await condenseEntry(erg, dataOneDay.data) + } + if (((args.database === 'both') || (args.database === 'influx')) && (dataline !== '')) { + const ok = await influx.influxWrite(dataline) + if (!ok) { + failedInfluxWrites.push({sensorid: sid, day: day}) + } + } + if((args.database === 'both') || (args.database === 'mongo')) { + if(erg.length !== 0) { + await mongo.writeDataArray(client, erg, styp) + } + } + return erg.length +} + +/* + let std = moment.utc(dt).startOf('day'); + let endd = moment.utc(dt).startOf('day').add(1,'day'); + let docs = await coll.find({datetime: {$gte: new Date(std), $lt: new Date(endd)}}, {sort: {datetime: 1}}).toArray(); + for (let i = docs.length - 1; i >= 0; i--) { + let dt = docs[i].datetime.valueOf(); + for (let a = all.length - 1; a >= 0; a--) { + let at = all[a].datetime.valueOf(); + if (dt == at) { + all.splice(a, 1); + break; + } + } + } +*/ +/* + for(let item of erg) { + // Einlesen der Einträge für diesen Sensor und diesen Tag +// {sensorid: "11999", datetime: {$gte: ISODate("2023-04-07"),$lt: ISODate("2023-04-08")}} + // check, ob es den Eintrag schon gibt + const ret = await mongo.checkOneSensor(client, sid, item.datetime) + if(ret.err) { + logit(`Èrror checking sensor ${sid} for ${item.datetime}\n${ret.errortext}`) + return + } + if (!ret.exists) { + const r = await mongo.writeOneSensor(client,item) + count += r.inserted + if(r.err) { + logit(`Error writing sensor ${sid} for ${item.datetime}\n${r.errortext}`) + return + } + } + } + return count +} +*/ + +// Parse command line options +function parse_cmdline(argv) { + let parser = new mod_getopt.BasicParser('s:(start)e:(end)t:(type)h(help)v(version)a:(entries)n:(sensorid)d:(dbase)',argv); + let option; + let std = START ? START : DateTime.now().startOf('day').minus({days: 1}).toFormat('yyyy-LL-dd'); // yesterday + let ret = {startDate: std, endDate: '', sensorType: TYP, database: DBASE}; + + while((option = parser.getopt()) !== undefined) { + switch(option.option) { + case 's': + ret.startDate = option.optarg.trim() + break; + + case 'e': + ret.endDate = option.optarg.trim() + break; + + case 'a': + ret.nbrOfEntries = parseInt(option.optarg.trim()) + break; + + case 'n': + ret.sensornbr = parseInt(option.optarg.trim()) + break; + + case 't': + let x = option.optarg.trim() + if (( x === 'noise') || ( x === 'radioactivity') || + (x === 'pm') || ( x === 'thp')) { + ret.sensorType = x + } + break; + + case 'd': + let db = option.optarg.trim() + if (DATABASE.indexOf(db) !== -1) { + ret.database = db + } + break; + + case 'v': + console.log(`Version: ${pkg.version} from ${pkg.date}`); + console.log(); + process.exit(); + break; + + case 'h': + console.log("Usage: node readFromcvs.js [-h] [-s startDate] [-e endDate] [-t sensorType] [-v version] [-h help]") + if (DEVELOP) { + console.log(" [-a entries] [-n sensorid]"); + } + console.log("Params:"); + console.log(" -s startDate: date to begin calculation (ex: 2023-10-23); default: yesterday"); + console.log(" -e endDate: date to stop calculatien (exclusive !). If not given, only the startDate") + 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(" -v version: show version"); + console.log(" -h this help text"); + if(DEVELOP) { + console.log(" -a entries: only work on 'entries' entries (only for debug)") + console.log(" -n check only this sensor number (only for debug)") + } + console.log("All parameters are optional."); + console.log(); + process.exit(); + break; + + default: + break; + } + } + if (ret.endDate === '') { + if (END) { + ret.endDate = END + } else { + let ed = DateTime.fromISO(ret.startDate) + ret.endDate = ed.plus({days: 1}).toFormat('yyyy-LL-dd') + } + } + logit(JSON.stringify(ret)); + return ret; +} + +async function main() { + logit(`Programm V: ${pkg.version} from ${pkg.date} - Start`) + let starttime = DateTime.now() + let args = parse_cmdline(process.argv) + const client = await mongo.connectMongo() + let inserted = await readSensorsperDay(client, args) + logit(`Inserted: ${inserted}`) + // enter changed properties into database + if(newProps.length > 0) { + await mongo.bulkWrite(client, mongo.properties_collection, newProps) + } + await mongo.closeMongo(client) + if (failedInfluxWrites.length > 0) { + logerror(`${failedInfluxWrites.length} Influx-Writes fehlgeschlagen - diese Daten fehlen in Influx:`) + for (const f of failedInfluxWrites) { + 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'); + logit(`Dauer: ${output1} Sekunden`) + logit("Programm - Ende") +} + +main().catch((e) => { + console.error(e) + process.exitCode = 1 +}) + + diff --git a/readarchive/utilities/checkprops.js b/readarchive/utilities/checkprops.js new file mode 100644 index 0000000..4387025 --- /dev/null +++ b/readarchive/utilities/checkprops.js @@ -0,0 +1,134 @@ +import { DateTime } from 'luxon' +import { logit, logerror } from './logit.js' +import { returnOnError } from './reporterror.js' +import * as mongo from '../mongo.js' + +// Check lat/lon and convert to float +function checkLatLon(w) { + let loc = 0.0; + if (!((w === undefined) || (w == null) || (w == ''))) { + try { + loc = parseFloat(w) + } catch (e) { + logerror(`Math error with lat/lon, ${e}`) + } + } + return loc +} + + +const buildNewEntry = (item, cdate, indoor, name, typ) => { + return { + _id: parseInt(item.sensor_id), + type: typ, + name: [{ + name: name, + since: cdate, + }], + location: [{ + loc: { + type: "Point", + coordinates: [ + checkLatLon(item.lon), + checkLatLon(item.lat) + ] + }, + id: parseInt(item.location), + altitude: 0, + since: cdate, + exact_loc: 0, + indoor: indoor, + country: '' + }] + } +} + +// Check for new properties +export const checkProperties = async (client, item, indoor, name, typ, dt, newProps) => { + if(name === "laerm") { + name = 'DNMS (Laerm)' + } + let ret + try { + ret = await mongo.getOneproperty(client, parseInt(item.sensor_id) ) // read properties for + if(ret.error) { + logerror(`Error reading properties: ${ret.errortext}`) + return null + } + } catch (e) { + logerror(`Error reading properties: ${e}`) + return null + } + let changed = false + let entry = ret.property + let cdate = DateTime.fromISO(dt + 'T00:00:00', { zone: 'utc' }).toJSDate() + // read entry from actualprops + if (entry === null) { // not in properties => new sensor + entry = buildNewEntry(item, cdate, indoor, name, typ) // so build a new entry + changed = true + } else { + // check for change of collection + let ok = false + for(let i = 0; i < entry.location.length; i++) { + if(entry.location[i].id === parseInt(item.location)) { + ok = true // same location found + } + } + if( !ok) { + // new location found: + if (newProps.findIndex((obj) => { + return (obj.updateOne.update.$set.location[0].id === item.location) + }) === -1) { // not in newProps array + const newloc = { + loc: { + type: "Point", + coordinates: [ + checkLatLon(item.lon), + checkLatLon(item.lat) + ] + }, + id: parseInt(item.location), + altitude: 0, + since: cdate, + exact_loc: 0, + indoor: indoor + } + entry.location.push(newloc) // add new lolcation to array + changed = true + logit(`New location ${item.location} for sensor ${item.sensor_id}`) + } + } else { + // same location, check the date + if (entry.location[0].since > cdate) { // new date is older than old date + entry.location[0].since = cdate // so set new date + changed = true + } + } + // Check für new name + if (entry.name[0].name.toUpperCase() !== name.toUpperCase()) { // have got a new name + if (newProps.findIndex( + (obj) => { + return (obj.updateOne.update.$set.name[0].name === name) + }) === -1) { + let newname = { + name: name, + since: cdate + } + entry.name.push(newname) + changed = true + logit(`New name ${name} for sensor ${item.sensor_id}`) + } + } else { // same name, chack the date + if (entry.name[0].since > cdate) { // new date is older than old date + entry.name[0].since = cdate // so set new date + changed = true + logit(`New date for sensor ${item.sensor_id}`) + } + } + } + // push this entry to the new properties array if somethind was changed + if(changed) { + delete entry.values // delete values array + newProps.push({updateOne: {filter: {_id: entry._id}, update: {$set: entry}, upsert: true}}) + } +} diff --git a/readarchive/utilities/logit.js b/readarchive/utilities/logit.js new file mode 100644 index 0000000..7a870e4 --- /dev/null +++ b/readarchive/utilities/logit.js @@ -0,0 +1,16 @@ +import { DateTime} from 'luxon' + +const MOCHA_TEST = process.env.MOCHA_TEST || false + +export function logit(str) { + if(MOCHA_TEST) return + let s = `${DateTime.now().toISO()} => ${str}`; + console.log(s); +} + +export function logerror(str) { + if(MOCHA_TEST) return + let s = `${DateTime.utc().toISO()} => *** ERROR *** ${str}`; + console.log(s); +} + diff --git a/readarchive/utilities/reporterror.js b/readarchive/utilities/reporterror.js new file mode 100644 index 0000000..2b3de74 --- /dev/null +++ b/readarchive/utilities/reporterror.js @@ -0,0 +1,20 @@ +import {logit} from "./logit.js"; + +export const reportError = (message, errortext) => { + message.error = true + message.errortext = errortext + return message +} + +export const returnOnError = (pr, error, name, p1='', p2='') => { + if (error.indexOf('xxx') !== -1) { + error = error.replace('xxx', p1) + } + if (error.indexOf('yyy') !== -1) { + error = error.replace('yyy', p1) + } + pr.err = error + logit(`${name}: ${error}`) + return pr +} +