forked from cptskippy/speed-trap
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtask_report_publisher.py
More file actions
151 lines (124 loc) · 4.78 KB
/
Copy pathtask_report_publisher.py
File metadata and controls
151 lines (124 loc) · 4.78 KB
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
148
149
150
151
"""
task_report_publisher.py
Subscribes to an MQTT topic and exports video clips
from the NVR based on timestamps in published messages
TODO::
* Pull metadata into report
"""
import json
from datetime import datetime
import shutil
from shared import MqttClientWrapper, load_config
VERBOSE = False
DEBUG = False
LOGGING = False
# Load configuration from yaml
config = load_config()
mqtt_config = config["servers"]["mqtt"]
task_config = config["task"]["report_publisher"]
media_config = config["media"]
# Populate configuration variables
MQTT_URI = mqtt_config["uri"]
MQTT_USER = mqtt_config["username"]
MQTT_PASSWORD = mqtt_config["password"]
MQTT_QOS = mqtt_config["qos"]
MQTT_CLIENT_ID = task_config["mqtt"]["client_id"]
MQTT_SUBSCRIBE_TOPIC = task_config["mqtt"]["topics"]["subscribe"]
MQTT_PUBLISH_TOPIC = task_config["mqtt"]["topics"]["publish"]
MQTT_ERROR_TOPIC = task_config["mqtt"]["topics"]["error"]
FOLDER_PATH = media_config["output_path"]
FOLDER_FORMAT = media_config["output_folder_format"]
PUBLISH_PATH = task_config["publish_path"]
PUBLISH_URL_TEMPLATE = task_config["publish_url_template"]
PUBLISH_HTML_FILE_NAME = task_config["html_file_name"]
PUBLISH_HTML_FILE_CONTENTS = task_config["html_file_contents"]
def on_connect(client, userdata, flags, reason_code, properties):
"""Subscribe to topic on successful connection."""
if reason_code == 0:
client.subscribe(MQTT_SUBSCRIBE_TOPIC, MQTT_QOS)
print(f"Subscribed to topic: {MQTT_SUBSCRIBE_TOPIC}")
def on_message(client, userdata, message):
""""Processes the message from MQTT Broker"""
try:
payload = message.payload.decode('utf-8')
data = json.loads(payload)
print("\nPublishing Event Received:")
print(f" Timestamp: {data.get('timestamp')}")
print(f" Sensor ID: {data.get('sensor_id')}")
print(f" Speed: {data.get('speed')} {data.get('uom')}")
print(f" Folder: {data.get('folder')}")
print(f" Data File: {data.get('data_file')}")
print(f" Payload: {payload}")
handle_event(data)
except json.JSONDecodeError:
print("Received invalid JSON: {message}")
except Exception as e:
print(f"Error processing message: {e}")
def generate_html_page(folder_name):
"""Generates a static HTML Page in the folder"""
html_page = f"{folder_name}/{PUBLISH_HTML_FILE_NAME}"
# Hardcoded page...
html_page_output = PUBLISH_HTML_FILE_CONTENTS
# Save data to disk
with open(html_page, "w", encoding="utf-8") as f:
f.write(html_page_output)
return html_page
def handle_event(data):
"""Creates a folder based on the timestamp of the event"""
# {
# "timestamp": "2025-06-22T22:46:56.524103+00:00",
# "speed": 27.96170365068,
# "uom": "mph",
# "sensor_id": "sensor.speedometer_speed",
# "folder": "./media/20250622154656",
# "data_file": "./media/20250622154656/data.json",
# "videos": [
# "./media/20250622154656/globalshutter.mpg",
# "./media/20250622154656/street.mpg",
# "./media/20250622154656/driveway.mpg"
# ],
# "images": [
# "./media/20250622154656/street.png",
# "./media/20250622154656/globalshutter.png",
# "./media/20250622154656/driveway.png"
# ],
# "thumbnails": [
# "./media/20250622154656/street_thumb.png",
# "./media/20250622154656/globalshutter_thumb.png",
# "./media/20250622154656/driveway_thumb.png"
# ]
# }
# Convert timestamp to string
timestamp = data.get("timestamp")
occurred = datetime.fromisoformat(timestamp)
local = occurred.astimezone()
dts = local.strftime(FOLDER_FORMAT)
# Create HTML
print("Saving page...")
folder = data.get("folder")
page = generate_html_page(folder)
print(" File saved")
# Publish Files
print("Publishing files...")
publish_folder = PUBLISH_PATH + dts
print(f" Source: {folder}")
print(f" Target: {publish_folder}")
shutil.copytree(src=folder, dst=publish_folder, dirs_exist_ok=True)
print(" Files published")
# Generate URLs
url = PUBLISH_URL_TEMPLATE.format(dts) + "/"
first_thumb = url
if "thumbnails" in data and isinstance(data["thumbnails"], list) and data["thumbnails"]:
first_thumb = PUBLISH_URL_TEMPLATE.format(data["thumbnails"][0].replace(FOLDER_PATH,""))
# Update Payload
data["page"] = page
data["url"] = url
data["url_thumb"] = first_thumb
payload = json.dumps(data)
print(f" New Payload: {payload}")
# Publish message for next task
client.publish(MQTT_PUBLISH_TOPIC, payload, MQTT_QOS)
print(f" Message Published: {MQTT_PUBLISH_TOPIC}")
# Configure MQTT and wait...
client = MqttClientWrapper(MQTT_URI, MQTT_CLIENT_ID, MQTT_USER, MQTT_PASSWORD)
client.connect(on_connect, on_message)