/* Interface for MongoDB Gemeinsam genutzt von readin (Live-API) und readarchive (CSV-Archiv). Zusammengefuehrt aus den beiden zuvor getrennten Fassungen; wo sie sich unterschieden, gilt: - MONGOAUTH wird als String verglichen. Die alte readarchive-Fassung pruefte nur auf truthy, dadurch schaltete auch MONGOAUTH=false die Authentifizierung ein. - writeDataArray(client, coll, data) nimmt den Collection-Namen direkt. readarchive baut ihn mit dataCollName(styp). - getallProperties liefert das Array selbst (readin-Fassung). Die readarchive-Fassung mit {error, properties} war dort ungenutzt. */ import { MongoClient } from 'mongodb' import { logit, logerror } from './logit.js' import { statistics } from './statistics.js' import { DateTime } from 'luxon' let DEVELOP = process.env.DEVELOP || 'false' let MONGOHOST = process.env.MONGOHOST || 'localhost' let MONGOPORT = process.env.MONGOPORT || 27017 let MONGOAUTH = process.env.MONGOAUTH || 'false' let MONGOUSRP = process.env.MONGOUSRP || '' let MONGOBASE = process.env.MONGOBASE || 'sensor_data' let MONGO_URL = 'mongodb://' + MONGOHOST + ':' + MONGOPORT; // URL to mongo database if (MONGOAUTH == 'true') { MONGO_URL = 'mongodb://' + MONGOUSRP + '@' + MONGOHOST + ':' + MONGOPORT + '/?authSource=admin'; // URL to mongo database } export const property_coll = 'properties' const data_collection = 'sensors' // Datencollection eines Sensortyps, z.B. dataCollName('noise') -> 'noise_sensors' export const dataCollName = (styp) => { return styp + '_' + data_collection } const addandshowstatistics = (client, text, field, start) => { statistics[field] = DateTime.now().diff(start, ['seconds']).toObject().seconds logit(`Write ${text} to mongoDB: Time: ${statistics[field]} sec.`) } 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) } } /* *************************************************** // READ routines ******************************************************/ // Read properties from the database - oeffnet und schliesst selbst 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(property_coll).findOne({_id: query.sid}) } else { // otherwise read props corresponding to query ret.properties = await client.db(MONGOBASE).collection(property_coll).find(query).limit(limit).toArray() } } catch (e) { ret.err = e } finally { client.close() } return ret } export const getallProperties = async (client) => { return await client.db(MONGOBASE).collection(property_coll) .find().sort({ _id: 1 }).toArray() } export const checkOneproperty = async (client, sid) => { return await client.db(MONGOBASE).collection(property_coll) .findOne({ _id: sid }) } export const getOneproperty = async (client, sid) => { let ret = {error: false, errortext: '', property: null} try { ret.property = await client.db(MONGOBASE).collection(property_coll) .findOne({_id: sid}) } catch (e) { ret = {error: true, errortext: e} } return ret } 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(dataCollName(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 getLocationIDs = async (client, stype) => { let ret = {error: false, errortext: '', locations: []} try { ret.locations = await client.db(MONGOBASE).collection(property_coll) .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(property_coll) .findOne({type: 'thp', 'location.0.id': lid}, { projection: {_id: 1}}) } catch (e) { ret.error = true ret.errortext = e } return ret } /* *************************************************** // WRITE routines ******************************************************/ export const writeOneproperty = async (client, prop) => { try { let result = await client.db(MONGOBASE).collection(property_coll) .insertOne(prop) } catch (e) { if (e.code == 11000) { return false } else { throw (e) } } return true } export const writeProperties = async (client, props) => { let result let startAll = DateTime.now(); let start let coll = client.db(MONGOBASE).collection(property_coll) if (props.new.length !== 0) { start = DateTime.now(); try { result = await coll.insertMany(props.new) } catch (e) { logerror(`Write properties new ${e}`) } logit(`Write ${props.new.length} properties NEW to mongoDB: Result: ${result.acknowledged}, Time: ${start.diffNow('seconds').toObject().seconds * -1} sec.`) } if (props.loc.length !== 0) { start = DateTime.now() try { for (let item of props.loc) { result = await coll.updateOne({ _id: item._id }, { $set: { location_id: item.location_id }, $push: { location: { $each: [item.location[0]], $position: 0 } } }) } } catch (e) { logerror(`Write properties location ${e}`) } logit(`Write ${props.loc.length} properties LOC to mongoDB: Result: ${result.acknowledged}, Time: ${start.diffNow('seconds').toObject().seconds * -1} sec.`) } if (props.sname.length !== 0) { start = DateTime.now() try { for (let item of props.sname) { result = await coll.updateOne({ _id: item._id }, { $push: { name: { $each: [item.name[0]], $position: 0 } } }) } } catch (e) { logerror(`Write properties samename ${e}`) } logit(`Write ${props.sname.length} properties NAME to mongoDB: Result: ${result.acknowledged}, Time: ${start.diffNow('seconds').toObject().seconds * -1} sec.`) } addandshowstatistics(client, 'properties', 'writePropsTime', startAll) } export const writeOneSensor = async (client, data, styp) => { let ret = {error: false, errortext: '', inserted: 1} try { await client.db(MONGOBASE).collection(dataCollName(styp)) .insertOne(data) } catch (e) { ret.error = true ret.errortext = e ret.inserted = 0 } return ret } export const writeDataArray = async (client, coll, data) => { let result let start = DateTime.now(); try { result = await client.db(MONGOBASE).collection(coll) .insertMany(data, { ordered: false }) } catch (e) { if (e.code !== 11000) { console.error(e) } } addandshowstatistics(client, `${data.length} entries for ${coll}`, `writeMongoData[${coll}]Time`, start) } export const writeStatistic = async (client, stat) => { let result let start = DateTime.now(); let entry = { timestamp: new Date(), ...stat } try { result = await client.db(MONGOBASE).collection("statistics") .insertOne(entry) } catch (e) { console.error(e) } addandshowstatistics(client, `statistics`, `writeStatisticTime`, start) } 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) } addandshowstatistics(client, `Data for ${coll}`, `writeMongoProperties[${coll}]Time`, start) return result } export const bulkUpdateMapdata = async (client, data) => { let start = DateTime.now() let result try { result = await client.db(MONGOBASE).collection("mapdata") .bulkWrite(data, { ordered: false }) } catch (e) { console.error(e) } logit(`Write MapData: Result: ${result}, Time: ${start.diffNow('second').toObject().seconds * -1} sec.`) return result } /* *************************************************** // Wartung ******************************************************/ export const dropColl = async (client, coll) => { let start = DateTime.now() let result try { result = await client.db(MONGOBASE).collection(coll).drop() } catch (e) { console.error(e) } logit(`Drop collection ${coll}: Result: ${result}, Time: ${start.diffNow('second').toObject().seconds * -1} sec.`) } export const createIndex = async (client, coll) => { let result let start = DateTime.now() try { result = await client.db(MONGOBASE).collection(coll).createIndex({ "location.loc": "2dsphere" }) } catch (e) { console.error(e) } logit(`Create-Index: Result: ${result}, Time: ${start.diffNow('second').toObject().seconds * -1} sec.`) }