-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdummy_loop.py
More file actions
96 lines (79 loc) · 3.58 KB
/
Copy pathdummy_loop.py
File metadata and controls
96 lines (79 loc) · 3.58 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
"""Publish one synthetic meter reading per hour, on the hour.
Goes through the project's own publish path (publish.build_live ->
publish.send), so what lands on the topic is byte-for-byte the shape the RTSP
pipeline produces. Only the reading itself is invented.
The value climbs rather than repeating: a consumer deriving consumption from
successive readings would read a flat line as a stalled meter, which is not
what a test feed should look like. It is anchored to the wall clock rather
than an in-memory counter, so a restart continues the series instead of
resetting it.
"""
import json
import logging
import os
import time
from datetime import datetime, timedelta
logging.basicConfig(level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("dummy")
from meter_ocr import config, publish
EVERY_HOURS = int(os.environ.get("EVERY_HOURS", 1))
AT_MINUTE = int(os.environ.get("AT_MINUTE", 0))
BASE = float(os.environ.get("BASE_READING", 13288.87))
RATE = float(os.environ.get("KWH_PER_HOUR", 0.35))
# Must be in the PAST: reading_now() clamps elapsed hours at zero, so an anchor
# in the future pins every reading to BASE and the series goes flat - which
# reads downstream as a stalled meter, not as a test feed.
ANCHOR = float(os.environ.get("ANCHOR_EPOCH", 1788134400)) # 2026-08-31 00:00 UTC
def reading_now() -> float:
hours = max(0.0, (time.time() - ANCHOR) / 3600.0)
return round(BASE + hours * RATE, 2)
def next_slot(now: datetime) -> datetime:
"""First :MM slot strictly after `now` that lands on the hour grid.
Same rule as schedule.next_run, reimplemented here rather than imported:
meter_ocr.schedule pulls in pipeline -> capture -> cv2, and this image
deliberately carries no OCR stack.
"""
slot = now.replace(minute=AT_MINUTE, second=0, microsecond=0)
if slot <= now:
slot += timedelta(hours=1)
while slot.hour % max(1, EVERY_HOURS) != 0:
slot += timedelta(hours=1)
return slot
def main() -> None:
cfg = config.load("config.yaml")
# Align against the configured offset, not the container clock: the image
# runs on UTC, where the top of the hour is :30 in IST - so aligning on
# local time would fire half an hour off.
tz = publish._tz(str(cfg.kafka.get("utc_offset", "+05:30")))
log.info("Dummy publisher: every %dh at :%02d (%s) -> %s topic=%s",
EVERY_HOURS, AT_MINUTE, cfg.kafka.get("utc_offset", "+05:30"),
cfg.kafka.get("bootstrap_servers"), cfg.kafka.get("topic"))
n = 0
while True:
target = next_slot(datetime.now(tz))
log.info("Next send at %s", target.isoformat(timespec="seconds"))
while True:
left = (target - datetime.now(tz)).total_seconds()
if left <= 0:
break
time.sleep(min(20, left)) # short slices keep Ctrl+C responsive
value = reading_now()
row = {
"reading": f"{value}",
"raw_text": f"{value} kWh",
"confidence": "0.94",
"status": "OK",
"snapshot": "data/annotated/dummy_annotated.jpg",
}
try:
msg = publish.build_live(cfg, row)
sent = publish.send(cfg, [msg])
n += 1
log.info("cycle %d: %s -> acked %d/1", n, json.dumps(msg["reading"]), sent)
except Exception:
# A broker outage must not kill a loop meant to run for days.
log.exception("publish failed - retrying at the next slot")
time.sleep(1) # never fire twice inside the same minute
if __name__ == "__main__":
main()