Repository navigation
Expand file tree
/
Copy pathplugin.py
More file actions
377 lines (329 loc) · 14.1 KB
/
Copy pathplugin.py
File metadata and controls
377 lines (329 loc) · 14.1 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
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
from __future__ import annotations
import json
import logging
import sys
import threading
import time
import uuid
from pathlib import Path
from typing import Any
try:
from .underfed import process
from .underfed import state as state_module
from .underfed.constants import (
API_KEY_ENV,
CAUSE_TIMESTAMPS,
DEFAULT_API_URL,
DEFAULT_DEMOTE_MINUTES,
DEFAULT_STORM_CONFIRM_SECONDS,
DEFAULT_STORM_PER_MINUTE,
DEFAULT_TELEMETRY_PATH,
HEARTBEAT_INTERVAL_SECONDS,
MAX_RECENT_EVENTS,
PLUGIN_DESCRIPTION,
PLUGIN_NAME,
PLUGIN_VERSION,
)
from .underfed.detector import Thresholds
from .underfed.journal import Journal
from .underfed.replay import run as replay
except ImportError:
sys.path.insert(0, str(Path(__file__).resolve().parent))
from underfed import process
from underfed import state as state_module
from underfed.constants import (
API_KEY_ENV,
CAUSE_TIMESTAMPS,
DEFAULT_API_URL,
DEFAULT_DEMOTE_MINUTES,
DEFAULT_STORM_CONFIRM_SECONDS,
DEFAULT_STORM_PER_MINUTE,
DEFAULT_TELEMETRY_PATH,
HEARTBEAT_INTERVAL_SECONDS,
MAX_RECENT_EVENTS,
PLUGIN_DESCRIPTION,
PLUGIN_NAME,
PLUGIN_VERSION,
)
from underfed.detector import Thresholds
from underfed.journal import Journal
from underfed.replay import run as replay
logger = logging.getLogger(__name__)
BASE_DIR = Path(__file__).resolve().parent
RUNTIME_DIR = BASE_DIR / ".runtime"
STATE_PATH = RUNTIME_DIR / "state.json"
JOURNAL_PATH = RUNTIME_DIR / "underfed.jsonl"
LOG_PATH = RUNTIME_DIR / "watcher.log"
DETACHING_REASONS = frozenset({"disable", "delete"})
ACTED_EVENTS = frozenset({"switched", "would_switch"})
FAILED_EVENTS = frozenset({"failed", "api_error", "error"})
_heartbeat_lock = threading.Lock()
_heartbeat_checked_at = 0.0
def _read_manifest() -> dict[str, Any]:
try:
with (BASE_DIR / "plugin.json").open(encoding="utf-8") as handle:
loaded = json.load(handle)
except (OSError, ValueError):
return {}
return loaded if isinstance(loaded, dict) else {}
_MANIFEST = _read_manifest()
def _running_inside_uwsgi() -> bool:
try:
import uwsgi # noqa: F401
except ImportError:
return False
return True
def _as_number(value: object, fallback: float) -> float:
try:
return float(str(value))
except (TypeError, ValueError):
return fallback
def _as_text(value: object, fallback: str) -> str:
text = str(value or "").strip()
return text or fallback
def _applied_arguments(state: dict) -> list[str] | None:
try:
arguments = json.loads(str(state.get("signature") or ""))
except ValueError:
return None
if isinstance(arguments, list) and arguments and all(isinstance(a, str) for a in arguments):
return arguments
return None
def _lines(value: object) -> list[str]:
if not value:
return []
if isinstance(value, (list, tuple)):
candidates = [str(item) for item in value]
else:
candidates = str(value).splitlines()
return [item.strip() for item in candidates if item.strip()]
class Plugin:
name = _MANIFEST.get("name", PLUGIN_NAME)
version = _MANIFEST.get("version", PLUGIN_VERSION)
description = _MANIFEST.get("description", PLUGIN_DESCRIPTION)
author = _MANIFEST.get("author", "")
help_url = _MANIFEST.get("help_url", "")
fields = _MANIFEST.get("fields", [])
actions = _MANIFEST.get("actions", [])
def run(self, action: str, params: dict, context: dict) -> dict:
handlers = {
"apply": self._apply,
"status": self._status,
"test": self._test,
"restart": self._restart,
"remove": self._remove,
}
handler = handlers.get((action or "").strip().lower())
if handler is None:
return {"status": "error", "message": f"Unknown action '{action}'."}
try:
return handler(dict(context or {}))
except Exception as exc:
logger.exception("Underfed action '%s' failed", action)
return {"status": "error", "message": f"{type(exc).__name__}: {exc}"}
def stop(self, context: dict | None = None) -> dict:
context = dict(context or {})
stopped = self._stop_watcher()
if str(context.get("reason") or "") in DETACHING_REASONS:
self._remember({}, clear=["pid", "token", "signature", "applied"])
else:
self._remember({}, clear=["pid", "token"])
return {
"status": "ok",
"message": "Watcher stopped." if stopped else "Watcher was not running.",
}
def _apply(self, context: dict) -> dict:
settings = dict(context.get("settings") or {})
key = _as_text(settings.get("api_key"), "")
if not key:
return {"status": "error", "message": "An API key is required."}
telemetry = _as_text(settings.get("telemetry_path"), DEFAULT_TELEMETRY_PATH)
if not Path(telemetry).is_file():
return {
"status": "error",
"message": f"{telemetry} does not exist. Is reservoarr the stream profile?",
}
arguments = self._arguments(settings)
signature = json.dumps(arguments, sort_keys=True)
state = self._state()
if process.is_running(state.get("pid"), state.get("token")):
if state.get("signature") == signature:
return {"status": "ok", "message": self._running_message(settings)}
process.terminate(state.get("pid"), state.get("token"))
self._start(arguments, key)
self._remember({"signature": signature, "applied": True})
return {"status": "ok", "message": self._running_message(settings)}
def _running_message(self, settings: dict) -> str:
percent = int(_as_number(settings.get("ratio_percent"), 70))
seconds = int(_as_number(settings.get("confirm_seconds"), 45))
storm = self._storm_per_minute(settings)
storm_seconds = int(self._storm_confirm_seconds(settings))
jumps = (
f" or at {storm} timestamp discontinuities a minute for {storm_seconds}s"
if storm
else ""
)
mode = (
"Observing only: it records what it would switch and changes nothing."
if settings.get("observe_only", True)
else "Switching is live."
)
return (
f"Watcher running below {percent}% of content rate for {seconds}s{jumps}. {mode}"
)
def _status(self, context: dict) -> dict:
settings = dict(context.get("settings") or {})
state = self._state()
running = process.is_running(state.get("pid"), state.get("token"))
records = Journal(JOURNAL_PATH, MAX_RECENT_EVENTS).read()
acted = [row for row in records if row.get("event") in ACTED_EVENTS]
skipped = [row for row in records if row.get("event") == "skipped"]
failed = [row for row in records if row.get("event") in FAILED_EVENTS]
if not state.get("applied"):
headline = "Not applied yet. Fill in the API key and press Apply."
else:
headline = "Watcher is running." if running else "Watcher is not running."
recent = "; ".join(self._describe(row) for row in acted[-5:])
detail = f" Last: {recent}." if recent else " Nothing switched yet."
if failed:
detail += f" {len(failed)} error(s) in the journal."
return {
"status": "ok" if running or not state.get("applied") else "error",
"message": f"{headline}{detail}",
"running": running,
"observe_only": bool(settings.get("observe_only", True)),
"switched": len([row for row in acted if row.get("event") == "switched"]),
"would_switch": len([row for row in acted if row.get("event") == "would_switch"]),
"skipped": len(skipped),
"events": records[-MAX_RECENT_EVENTS:],
}
def _describe(self, record: dict) -> str:
verb = "switched" if record.get("event") == "switched" else "would switch"
channel = record.get("channel") or record.get("feed")
if record.get("cause") == CAUSE_TIMESTAMPS:
return (
f"{channel} {verb} at {record.get('per_minute')} "
f"timestamp discontinuities a minute"
)
return f"{channel} {verb} at {record.get('percent')}%"
def _test(self, context: dict) -> dict:
settings = dict(context.get("settings") or {})
telemetry = Path(_as_text(settings.get("telemetry_path"), DEFAULT_TELEMETRY_PATH))
outcome = replay(telemetry, self._thresholds(settings))
return {
"status": "ok",
"message": outcome.summary(),
"samples": outcome.samples,
"sources": len(outcome.feeds),
"triggers": len(outcome.hits),
"storms": len(outcome.storms),
}
def _restart(self, context: dict) -> dict:
settings = dict(context.get("settings") or {})
state = self._state()
if not state.get("applied"):
return {"status": "ok", "message": "Nothing to do: not applied."}
if not _running_inside_uwsgi():
return {
"status": "ok",
"message": "Skipped: the watcher only starts in the Dispatcharr web process.",
}
if not self._heartbeat_due():
return {"status": "ok", "message": "Checked recently."}
if process.is_running(state.get("pid"), state.get("token")):
return {"status": "ok", "message": "Watcher is running."}
key = _as_text(settings.get("api_key"), "")
if not key:
return {"status": "error", "message": "An API key is required."}
self._start(_applied_arguments(state) or self._arguments(settings), key)
return {"status": "ok", "message": "Watcher restarted."}
def _remove(self, context: dict) -> dict:
stopped = self._stop_watcher()
self._remember({}, clear=["pid", "token", "applied", "signature"])
return {
"status": "ok",
"message": "Watcher stopped." if stopped else "Watcher was not running.",
}
def _thresholds(self, settings: dict) -> Thresholds:
percent = min(max(_as_number(settings.get("ratio_percent"), 70), 1.0), 99.0)
return Thresholds(
ratio=percent / 100,
confirm_seconds=max(_as_number(settings.get("confirm_seconds"), 45), 5.0),
warmup_seconds=max(_as_number(settings.get("warmup_seconds"), 60), 0.0),
stable_seconds=max(_as_number(settings.get("stable_seconds"), 180), 0.0),
storm_per_minute=self._storm_per_minute(settings),
storm_confirm_seconds=self._storm_confirm_seconds(settings),
)
def _storm_per_minute(self, settings: dict) -> int:
return max(int(_as_number(settings.get("storm_per_minute"), DEFAULT_STORM_PER_MINUTE)), 0)
def _storm_confirm_seconds(self, settings: dict) -> float:
return max(
_as_number(settings.get("storm_confirm_seconds"), DEFAULT_STORM_CONFIRM_SECONDS), 5.0
)
def _demote_minutes(self, settings: dict) -> float:
return max(_as_number(settings.get("demote_minutes"), DEFAULT_DEMOTE_MINUTES), 0.0)
def _arguments(self, settings: dict) -> list[str]:
arguments = [
"--telemetry",
_as_text(settings.get("telemetry_path"), DEFAULT_TELEMETRY_PATH),
"--api-url",
_as_text(settings.get("api_url"), DEFAULT_API_URL),
"--journal",
str(JOURNAL_PATH),
"--ratio-percent",
str(_as_number(settings.get("ratio_percent"), 70)),
"--confirm-seconds",
str(_as_number(settings.get("confirm_seconds"), 45)),
"--warmup-seconds",
str(_as_number(settings.get("warmup_seconds"), 60)),
"--stable-seconds",
str(_as_number(settings.get("stable_seconds"), 180)),
"--max-switches",
str(int(_as_number(settings.get("max_switches"), 2))),
"--storm-per-minute",
str(self._storm_per_minute(settings)),
"--storm-confirm-seconds",
str(self._storm_confirm_seconds(settings)),
"--demote-minutes",
str(self._demote_minutes(settings)),
]
for name in _lines(settings.get("exclude_channels")):
arguments += ["--exclude", name]
if settings.get("observe_only", True):
arguments.append("--observe-only")
return arguments
def _start(self, arguments: list[str], api_key: str) -> None:
process.terminate_strays()
token = uuid.uuid4().hex
pid = process.spawn(
BASE_DIR, arguments, token, LOG_PATH, extra_env={API_KEY_ENV: api_key}
)
self._remember({"pid": pid, "token": token, "started_at": time.time()})
if not self._settled(pid, token):
process.terminate(pid)
self._remember({}, clear=["pid", "token"])
raise RuntimeError(f"The watcher exited on startup. See {LOG_PATH}.")
def _settled(self, pid: int, token: str, seconds: float = 3.0) -> bool:
deadline = time.monotonic() + seconds
while time.monotonic() < deadline:
if not process.is_running(pid, token):
return False
time.sleep(0.25)
return True
def _stop_watcher(self) -> bool:
state = self._state()
stopped = process.terminate(state.get("pid"), state.get("token"))
strays = process.terminate_strays()
return stopped or bool(strays)
def _heartbeat_due(self) -> bool:
global _heartbeat_checked_at
now = time.monotonic()
with _heartbeat_lock:
if now - _heartbeat_checked_at < HEARTBEAT_INTERVAL_SECONDS:
return False
_heartbeat_checked_at = now
return True
def _state(self) -> dict:
return state_module.load(STATE_PATH)
def _remember(self, updates: dict, clear: list[str] | None = None) -> dict:
return state_module.remember(STATE_PATH, updates, clear or [])