-
Notifications
You must be signed in to change notification settings - Fork 0
/
logger.py
147 lines (120 loc) · 4.27 KB
/
logger.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
# # '''
# # UART communication on Raspberry Pi using Pyhton
import asyncio
import os
import serial
# from AWSIoTPythonSDK.MQTTLib import AWSIoTMQTTClient
import json
from datetime import datetime
import requests
import time
TEST_ID = 'dev-pi4'
AUTH_KEY = 'id8wn6mkJf86Svo10uiJI6iqXvebPi4US0OV9S8a'
URL = 'https://6a3blrx50f.execute-api.us-east-2.amazonaws.com/Production/ingest'
headers = {
'Content-Type': 'application/json',
'x-api-key': AUTH_KEY
}
async def producer(queue):
print("Producer: Opening serial port")
ser = serial.Serial("/dev/serial0", 57600) # Open port with baud rate
ser.reset_input_buffer()
print("Producer: Running")
batch = []
BATCH_SIZE = 10
while True:
try:
rx = ser.readline() # str(time.time() * 1000) + ",0,0,0,1394,2493,0,0,0,00.00.00"
s = datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3] + "," + str(rx, encoding="utf-8")
batch.append(s)
except UnicodeDecodeError:
print("Unicode error")
if len(batch) >= BATCH_SIZE:
await queue.put(batch)
batch = []
await asyncio.sleep(1)
# coroutine to consume work
async def consumer(queue):
print("Consumer: Connecting to AWS")
# mqttClient = AWSIoTMQTTClient("bruce-telemetry-pi")
# mqttClient.enableMetricsCollection()
# mqttClient.configureEndpoint("a134m7gripxuzz-ats.iot.us-east-2.amazonaws.com", 8883)
# mqttClient.configureCredentials(
# "/home/bruce/AWSLogger/AmazonRootCA1.pem",
# "/home/bruce/AWSLogger/b4c01e4ca7d972a7ab4b9dc3f0d8456f788bc19901ca73474c7c085054b9bd63-private.pem.key",
# "/home/bruce/AWSLogger/b4c01e4ca7d972a7ab4b9dc3f0d8456f788bc19901ca73474c7c085054b9bd63-certificate.pem.crt"
# )
# mqttClient.configureOfflinePublishQueueing(-1) # infinite
# mqttClient.configureDrainingFrequency(2) # 2 Hz
# mqttClient.configureConnectDisconnectTimeout(10) # 10s
# mqttClient.configureMQTTOperationTimeout(5) # 10s
# while not mqttClient.connect():
# print("Failed to connect to AWS")
# asyncio.sleep(1)
# print("Connected to AWS!")
# consume work
while True:
item = await queue.get()
# print(item)
formatted = []
for i in item:
pieces = i.split(',')
d = {
"timestamp": {
"S": pieces[0]
},
"test_id": {
"S": TEST_ID
},
"throttle": {
"N": pieces[2]
},
"speed": {
"N": pieces[3]
},
"rpm": {
"N": pieces[4]
},
"current": {
"N": pieces[5]
},
"voltage": {
"N": pieces[6]
},
"throttleTooHigh": {
"N": pieces[7]
},
"motorInitializing": {
"N": pieces[8]
},
"clockState": {
"N": pieces[9]
},
"lastDeadman": {
"S": pieces[10]
}
}
# print(d)
formatted.append(d)
# print(json.dumps(formatted))
try:
response = requests.post(URL, headers=headers, json={"messages": formatted})
print("Uploaded to AWS with response", response.status_code)
except:
pass
print("Failed to connect to AWS")
# Create message payload
# mqttClient.publish("bruce/telemetry", i, 0)
# entry point coroutine
async def main():
# create the shared queue
queue = asyncio.Queue()
# run the producer and consumers
await asyncio.gather(producer(queue), consumer(queue))
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
# Exit application because user indicated they wish to exit.
# This will have cancelled `main()` implicitly.
print("User initiated exit. Exiting.")