-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathlib.py
More file actions
154 lines (129 loc) · 5.5 KB
/
Copy pathlib.py
File metadata and controls
154 lines (129 loc) · 5.5 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
152
153
import logging
import datetime
from threading import Thread, Event
import time
from typing import Callable
class TimedCalls(Thread):
"""Call function again every `interval` time duration after it's first run."""
def __init__(self, func: Callable, interval: datetime.timedelta, start_offset_seconds: int = 0) -> None:
super().__init__()
self.func = func
self.interval = interval
self.stopped = Event()
self.start_offset_seconds = start_offset_seconds
def cancel(self):
self.stopped.set()
def wait_until_start_time(self, start_time):
if datetime.datetime.now() >= start_time:
return
while datetime.datetime.now() < start_time:
time.sleep(1)
def run(self):
self.func()
now = datetime.datetime.now()
start_time = self.get_nearest_interval_time(now, self.interval)
if (offset_negative := start_time - datetime.timedelta(seconds=self.start_offset_seconds)) > now:
start_time = offset_negative
else:
start_time = start_time + datetime.timedelta(seconds=self.start_offset_seconds)
logging.info(f"waiting until {start_time} to start loop...")
self.wait_until_start_time(start_time)
next_call = time.time()
logging.info(f"starting loop...")
while not self.stopped.is_set():
self.func() # Target activity.
next_call = next_call + self.interval.seconds
# Block until beginning of next interval (unless canceled).
self.stopped.wait(next_call - time.time())
def get_nearest_interval_time(self, base_time, interval):
base_time = base_time.replace(microsecond=0)
seconds_since_midnight = int(time.time())
interval_seconds = interval.total_seconds()
remainder = seconds_since_midnight % interval_seconds
if remainder == 0:
return base_time
else:
delta = interval_seconds - remainder
return base_time + datetime.timedelta(seconds=delta)
def query_influx(secrets, influx_client):
query_api = influx_client.query_api()
bucket = secrets['bucket']
power_consumed_query = f'from(bucket:"{bucket}")\
|> range(start: -15m)\
|> filter(fn: (r) => r._measurement == "fusionsolarpy")\
|> filter(fn: (r) => r.component == "power")\
|> filter(fn: (r) => r._field == "power_consumed")'
power_produced_query = f'from(bucket:"{bucket}")\
|> range(start: -15m)\
|> filter(fn: (r) => r._measurement == "fusionsolarpy")\
|> filter(fn: (r) => r.component == "power")\
|> filter(fn: (r) => r._field == "power_produced")'
server_power_draw_query = f'from(bucket:"{bucket}")\
|> range(start: -1m)\
|> filter(fn: (r) => r._measurement == "tinytuya")\
|> filter(fn: (r) => r.socket == "server")\
|> filter(fn: (r) => r._field == "power")\
|> last()'
try:
power_consumed_result = query_api.query(org=secrets['org'], query=power_consumed_query)
power_produced_result = query_api.query(org=secrets['org'], query=power_produced_query)
server_power_result = query_api.query(org=secrets['org'], query=server_power_draw_query)
except Exception as e:
logging.warning(f"failed to query influx: {e}")
return None, None
try:
server_power = next(record for record in next(table for table in server_power_result))
power_produced = []
power_consumed = []
for table in power_produced_result:
for record in table.records:
power_produced.append(record)
for table in power_consumed_result:
for record in table.records:
power_consumed.append(record)
current_values = {
"server_power": server_power,
"power_produced": power_produced[-1],
"power_consumed": power_consumed[-1],
"query_time": power_consumed[-1].get_time().astimezone()
}
previous_values = {
"power_produced": power_produced[-2],
"power_consumed": power_consumed[-2],
"query_time": power_consumed[-2].get_time().astimezone()
}
return previous_values, current_values
except StopIteration as e:
logging.warning(f"No data returned from influx: {e}")
except Exception as e:
logging.warning(f"invalid data from influx: {e}")
return None, None
def query_docker(docker_client, container_name):
container = docker_client.container.inspect(container_name)
container_running = container.state.running
container_paused = container.state.paused
return container_running, container_paused
def get_tdarr_node_running_status(session, config):
r = session.get(f'{config["tdarrUrl"]}/api/v2/get-nodes')
r.raise_for_status()
json = r.json()
if len(json.keys()) > 0:
key = list(json.keys())[0]
else:
raise Exception("no nodes found")
return json[key]['nodePaused'], key
def set_tdarr_node_status(session, node_id, pause, config):
r = session.post(f'{config["tdarrUrl"]}/api/v2/update-node', json={
"data": {
"nodeID": node_id,
"nodeUpdates": {
"nodePaused": pause
}
}
})
r.raise_for_status()
def update_tdarr_node(session, node_id, config, pause):
try:
set_tdarr_node_status(session, node_id, pause, config)
except Exception as e:
logging.warning(f"failed to update tdarr node status: {e}")