annotate service/mqtt_to_rdf/rdf_from_mqtt.py @ 1532:7cc7700302c2

more service renaming; start a lot more serv.n3 job files Ignore-this: 635aaefc7bd2fa5558eefb8b3fc9ec75 darcs-hash:2c8b587cbefa4db427f9a82676abdb47e651187e
author drewp <drewp@bigasterisk.com>
date Thu, 06 Feb 2020 16:36:35 -0800
parents
children
Ignore whitespace changes - Everywhere: Within whitespace: At end of lines:
rev   line source
1532
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
1 """
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
2 Subscribe to mqtt topics; generate RDF statements.
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
3 """
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
4 import json
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
5 import sys
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
6 from docopt import docopt
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
7 from rdflib import Namespace, URIRef, Literal, Graph, RDF, XSD
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
8 from rdflib.parser import StringInputSource
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
9 from rdflib.term import Node
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
10 from twisted.internet import reactor
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
11 import cyclone.web
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
12 import rx, rx.operators, rx.scheduler.eventloop
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
13 from greplin import scales
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
14 from greplin.scales.cyclonehandler import StatsHandler
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
15
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
16 from export_to_influxdb import InfluxExporter
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
17 from mqtt_client import MqttClient
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
18
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
19 from patchablegraph import PatchableGraph, CycloneGraphHandler, CycloneGraphEventsHandler
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
20 from rdfdb.patch import Patch
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
21 from rdfdb.rdflibpatch import graphFromQuads
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
22 from standardservice.logsetup import log, verboseLogging
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
23 from standardservice.scalessetup import gatherProcessStats
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
24
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
25 ROOM = Namespace('http://projects.bigasterisk.com/room/')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
26
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
27 gatherProcessStats()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
28
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
29 def parseDurationLiteral(lit: Literal) -> float:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
30 if lit.endswith('s'):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
31 return float(lit.split('s')[0])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
32 raise NotImplementedError(f'duration literal: {lit}')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
33
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
34
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
35 class MqttStatementSource:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
36 def __init__(self, uri, config, masterGraph, mqtt, influx):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
37 self.uri = uri
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
38 self.config = config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
39 self.masterGraph = masterGraph
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
40 self.mqtt = mqtt
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
41 self.influx = influx
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
42
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
43 self.mqttTopic = self.topicFromConfig(self.config)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
44
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
45 statPath = '/subscribed_topic/' + self.mqttTopic.decode('ascii').replace('/', '|')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
46 scales.init(self, statPath)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
47 self._mqttStats = scales.collection(
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
48 statPath + '/incoming', scales.IntStat('count'),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
49 scales.RecentFpsStat('fps'))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
50
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
51
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
52 rawBytes = self.subscribeMqtt()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
53 rawBytes = rx.operators.do_action(self.countIncomingMessage)(rawBytes)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
54 parsed = self.getParser()(rawBytes)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
55
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
56 g = self.config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
57 for conv in g.items(g.value(self.uri, ROOM['conversions'])):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
58 parsed = self.conversionStep(conv)(parsed)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
59
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
60 outputQuadsSets = rx.combine_latest(
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
61 *[self.makeQuads(parsed, plan)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
62 for plan in g.objects(self.uri, ROOM['graphStatements'])])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
63
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
64 outputQuadsSets.subscribe_(self.updateQuads)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
65
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
66 def topicFromConfig(self, config) -> bytes:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
67 topicParts = list(config.items(config.value(self.uri, ROOM['mqttTopic'])))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
68 return b'/'.join(t.encode('ascii') for t in topicParts)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
69
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
70
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
71 def subscribeMqtt(self):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
72 return self.mqtt.subscribe(self.mqttTopic)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
73
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
74 def countIncomingMessage(self, _):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
75 self._mqttStats.fps.mark()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
76 self._mqttStats.count += 1
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
77
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
78 def getParser(self):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
79 g = self.config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
80 parser = g.value(self.uri, ROOM['parser'])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
81 if parser == XSD.double:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
82 return rx.operators.map(lambda v: Literal(float(v.decode('ascii'))))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
83 elif parser == ROOM['tagIdToUri']:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
84 return rx.operators.map(self.tagIdToUri)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
85 elif parser == ROOM['onOffBrightness']:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
86 return rx.operators.map(lambda v: Literal(0.0 if v == b'OFF' else 1.0))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
87 elif parser == ROOM['jsonBrightness']:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
88 return rx.operators.map(self.parseJsonBrightness)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
89 elif ROOM['ValueMap'] in g.objects(parser, RDF.type):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
90 return rx.operators.map(lambda v: self.remap(parser, v.decode('ascii')))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
91 else:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
92 raise NotImplementedError(parser)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
93
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
94 def parseJsonBrightness(self, mqttValue: bytes):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
95 msg = json.loads(mqttValue.decode('ascii'))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
96 return Literal(float(msg['brightness'] / 255) if msg['state'] == 'ON' else 0.0)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
97
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
98 def conversionStep(self, conv: Node):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
99 g = self.config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
100 if conv == ROOM['celsiusToFarenheit']:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
101 return rx.operators.map(lambda value: Literal(round(value.toPython() * 1.8 + 32, 2)))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
102 elif g.value(conv, ROOM['ignoreValueBelow'], default=None) is not None:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
103 threshold = g.value(conv, ROOM['ignoreValueBelow'])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
104 return rx.operators.filter(lambda value: value.toPython() >= threshold.toPython())
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
105 else:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
106 raise NotImplementedError(conv)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
107
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
108 def makeQuads(self, parsed, plan):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
109 g = self.config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
110 def quadsFromValue(valueNode):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
111 return set([
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
112 (self.uri,
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
113 g.value(plan, ROOM['outputPredicate']),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
114 valueNode,
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
115 self.uri)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
116 ])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
117
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
118 def emptyQuads(element):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
119 return set([])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
120
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
121 quads = rx.operators.map(quadsFromValue)(parsed)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
122
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
123 dur = g.value(plan, ROOM['statementLifetime'])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
124 if dur is not None:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
125 sec = parseDurationLiteral(dur)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
126 quads = quads.pipe(
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
127 rx.operators.debounce(sec, rx.scheduler.eventloop.TwistedScheduler(reactor)),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
128 rx.operators.map(emptyQuads),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
129 rx.operators.merge(quads),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
130 )
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
131
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
132 return quads
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
133
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
134 def updateQuads(self, newGraphs):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
135 newQuads = set.union(*newGraphs)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
136 g = graphFromQuads(newQuads)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
137 log.debug(f'{self.uri} update to {len(newQuads)} statements')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
138
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
139 self.influx.exportToInflux(newQuads)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
140
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
141 self.masterGraph.patchSubgraph(self.uri, g)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
142
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
143 def tagIdToUri(self, value: bytearray) -> URIRef:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
144 justHex = value.decode('ascii').replace('-', '').lower()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
145 int(justHex, 16) # validate
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
146 return URIRef(f'http://bigasterisk.com/rfidCard/{justHex}')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
147
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
148 def remap(self, parser, valueStr: str):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
149 g = self.config
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
150 value = Literal(valueStr)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
151 for entry in g.objects(parser, ROOM['map']):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
152 if value == g.value(entry, ROOM['from']):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
153 return g.value(entry, ROOM['to'])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
154 raise KeyError(value)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
155
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
156
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
157 if __name__ == '__main__':
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
158 arg = docopt("""
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
159 Usage: rdf_from_mqtt.py [options]
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
160
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
161 -v Verbose
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
162 --cs=STR Only process config filenames with this substring
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
163 """)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
164 verboseLogging(arg['-v'])
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
165
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
166 config = Graph()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
167 for fn in [
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
168 "config_cardreader.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
169 "config_nightlight_ari.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
170 "config_bed_bar.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
171 "config_air_quality_indoor.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
172 "config_air_quality_outdoor.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
173 "config_living_lamps.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
174 "config_kitchen.n3",
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
175 ]:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
176 if not arg['--cs'] or arg['--cs'] in fn:
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
177 config.parse(fn, format='n3')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
178
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
179 masterGraph = PatchableGraph()
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
180
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
181 mqtt = MqttClient(clientId='rdf_from_mqtt', brokerHost='bang',
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
182 brokerPort=1883)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
183 influx = InfluxExporter(config)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
184
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
185 srcs = []
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
186 for src in config.subjects(RDF.type, ROOM['MqttStatementSource']):
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
187 srcs.append(MqttStatementSource(src, config, masterGraph, mqtt, influx))
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
188 log.info(f'set up {len(srcs)} sources')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
189
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
190 port = 10018
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
191 reactor.listenTCP(port, cyclone.web.Application([
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
192 (r"/()", cyclone.web.StaticFileHandler,
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
193 {"path": ".", "default_filename": "index.html"}),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
194 (r'/stats/(.*)', StatsHandler, {'serverName': 'rdf_from_mqtt'}),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
195 (r"/graph/mqtt", CycloneGraphHandler, {'masterGraph': masterGraph}),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
196 (r"/graph/mqtt/events", CycloneGraphEventsHandler,
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
197 {'masterGraph': masterGraph}),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
198 ], mqtt=mqtt, masterGraph=masterGraph, debug=arg['-v']),
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
199 interface='::')
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
200 log.warn('serving on %s', port)
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
201
7cc7700302c2 more service renaming; start a lot more serv.n3 job files
drewp <drewp@bigasterisk.com>
parents:
diff changeset
202 reactor.run()