-
-
Notifications
You must be signed in to change notification settings - Fork 29
Expand file tree
/
Copy pathplugin_manager.py
More file actions
1927 lines (1701 loc) · 90.9 KB
/
Copy pathplugin_manager.py
File metadata and controls
1927 lines (1701 loc) · 90.9 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
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""
Plugin Manager
Manages plugin discovery, loading, and lifecycle for the LEDMatrix system.
Loads plugins from the configured plugins directory
(``plugin_system.plugins_directory``, ``plugin-repos/`` by default).
API Version: 1.0.0
"""
import json
import math
import queue
import sys
import time
import threading
import types
import uuid
from pathlib import Path
from typing import Callable, Dict, List, NamedTuple, Optional, Any, Tuple, Union
import logging
from src import display_watchdog
from src.exceptions import PluginError, ConfigError
from src.logging_config import get_logger
from src.plugin_system.plugin_loader import PluginLoader
from src.plugin_system.plugin_executor import (
PluginBusyError, PluginExecutor, PluginTimeoutError,
)
from src.plugin_system.plugin_state import PluginStateManager, PluginState
from src.plugin_system.schema_manager import (
CORE_VEGAS_TUNING_KEYS, SchemaManager, normalize_legacy_booleans,
)
from src.plugin_system.plugin_dirs import (
ManifestStatus, PluginDirectoryIndex, resolve_plugin_dir,
)
from src.common.fetch_service import plugin_scope, register_plugin_directory
from src.common.permission_utils import (
ensure_directory_permissions,
get_plugin_dir_mode
)
class _DeferredConfigChange(NamedTuple):
"""Update-queue item: apply the config change parked for ``plugin_id``.
Queued by apply_config_change() when the plugin's lock was busy; the
change itself waits in ``PluginManager._deferred_config_changes`` so only
the latest one is ever applied.
"""
plugin_id: str
class PluginManager:
"""
Manages plugin discovery, loading, and lifecycle.
The PluginManager is responsible for:
- Discovering plugins in the configured plugins directory
- Loading plugin modules and instantiating plugin classes
- Managing plugin lifecycle (load, unload, reload)
- Providing access to loaded plugins
- Maintaining plugin manifests
Uses composition with specialized components:
- PluginLoader: Handles module loading and dependency installation
- PluginExecutor: Handles plugin execution with timeout and error isolation
- PluginStateManager: Manages plugin state machine
"""
# How long unload_plugin() waits for an in-flight update() to finish
# before tearing the instance down anyway.
UNLOAD_LOCK_TIMEOUT = 5.0
# How long unload_detached_plugin() (a live reload, off the render thread)
# waits for the old instance's lock. Longer than UNLOAD_LOCK_TIMEOUT
# because nothing is blocked by the wait, and a Vegas content build of the
# old instance can hold the lock for several seconds (5.9 s seen on
# ledpi). Past it the reload is refused rather than tearing down an
# instance a Vegas call may still be running in.
DETACHED_UNLOAD_LOCK_TIMEOUT = 30.0
# How long the update worker and apply_config_change() wait for a
# plugin's lock -- the same bound unload already uses for the same lock.
# A display() frame holds it for milliseconds, so this only runs out when
# the holder is hung or pathologically slow. The worker then skips that
# plugin (recorded as a hang, so repeats open its circuit breaker)
# instead of stalling every other plugin's update behind it.
PLUGIN_LOCK_TIMEOUT = UNLOAD_LOCK_TIMEOUT
# Minimum seconds between repeats of the same hang/slow-call warning for
# one plugin. A hung plugin is re-detected every interval; a slow
# display() can be re-detected every frame.
HANG_LOG_INTERVAL = 60.0
def __init__(self, plugins_dir: str = "plugins",
config_manager: Optional[Any] = None,
display_manager: Optional[Any] = None,
cache_manager: Optional[Any] = None,
font_manager: Optional[Any] = None) -> None:
"""
Initialize the Plugin Manager.
Args:
plugins_dir: Path to the plugins directory
config_manager: Configuration manager instance
display_manager: Display manager instance
cache_manager: Cache manager instance
font_manager: Font manager instance
"""
self.plugins_dir: Path = Path(plugins_dir)
self.config_manager: Optional[Any] = config_manager
self.display_manager: Optional[Any] = display_manager
self.cache_manager: Optional[Any] = cache_manager
self.font_manager: Optional[Any] = font_manager
self.logger: logging.Logger = get_logger(__name__)
# Initialize plugin system components
self.plugin_loader = PluginLoader(logger=self.logger)
self.plugin_executor = PluginExecutor(default_timeout=30.0, logger=self.logger)
self.state_manager = PluginStateManager(logger=self.logger)
self.schema_manager = SchemaManager(plugins_dir=self.plugins_dir, logger=self.logger,
config_manager=self.config_manager)
# Lock protecting plugin_manifests and plugin_directories from
# concurrent mutation (background reconciliation) and reads (requests).
self._discovery_lock = threading.RLock()
#: Directories already reported as unloadable, so the warning is
#: emitted once rather than on every discovery scan.
self._skip_reported: set = set()
# Lock protecting plugin_last_update from concurrent mutation/iteration.
# It's written from run_scheduled_updates() (main loop) and read/diffed by run_scheduled_updates_with_changes(), which
# Vegas mode calls from its own background update-tick thread.
self._plugin_last_update_lock = threading.RLock()
# Active plugins
self.plugins: Dict[str, Any] = {}
self.plugin_manifests: Dict[str, Dict[str, Any]] = {}
self.plugin_directories: Dict[str, Path] = {}
self.plugin_last_update: Dict[str, float] = {}
# Cached static data-fetch intervals per plugin_id, so the render
# loop's scheduling tick does not repeat the manifest/config lookup
# for every plugin. Cleared on load/unload.
self._update_interval_cache: Dict[str, Optional[float]] = {}
# Health tracking (optional, set by display_controller if available)
self.health_tracker = None
self.resource_monitor = None
# --- Asynchronous plugin updates -------------------------------
# Run inline in the render loop, one slow plugin HTTP fetch in
# update() freezes scrolling for the whole fetch. Scheduling happens
# on the render thread (run_scheduled_updates); execution happens on
# this single background worker. Per-plugin locks keep a plugin's
# update() and display() mutually exclusive, including across the
# post-timeout window.
# Kill switch: plugin_system.synchronous_updates: true restores the
# inline path.
#
# Which thread runs each plugin hook, and what it holds:
# __init__, on_enable the loading thread (main thread at startup,
# the render thread on a live enable, a
# plugin-reload thread for a control socket
# reload).
# update() plugin-update-worker, under the plugin lock,
# via PluginExecutor (whose daemon thread runs
# the call; if it outlives the executor's
# timeout it keeps the lock until it returns).
# Exceptions: the startup pass
# (DisplayController._run_initial_updates, main
# thread, before the display loop starts) and
# the synchronous_updates kill switch (render
# thread) run it without the lock.
# display() the render thread, under a try-lock: a busy
# lock skips the frame. The first frame of a
# screen goes through PluginExecutor. Vegas
# mode's adapter and coordinator take the lock
# with a bounded wait.
# on_config_change() ConfigService's watcher thread, under the
# plugin lock via apply_config_change(); if the
# lock stays busy it is deferred to the update
# worker, which applies it under the lock.
# cleanup(), on_disable() whoever calls unload_plugin() (or, for a
# reload, unload_detached_plugin() on its
# plugin-reload thread), under the lock with
# UNLOAD_LOCK_TIMEOUT.
# No wait on a plugin lock is unbounded, so one hung plugin can only
# cost the worker PLUGIN_LOCK_TIMEOUT per attempt.
self._update_queue: "queue.Queue[Union[None, Tuple[str, float], _DeferredConfigChange]]" = queue.Queue()
self._pending_updates: set = set()
self._pending_lock = threading.Lock()
# Serializes the "is this plugin eligible?" -> "claim it (RUNNING)"
# transition. Two schedulers run concurrently in practice — the render
# loop's _tick_plugin_updates() and Vegas mode's vegas-plugin-tick
# daemon thread, which is never joined — so without this both can
# observe ENABLED and both call update() on the same plugin. Held only
# across the check and the state transition, never across update()
# itself: that would serialize slow plugins behind each other and
# reintroduce the stall the async worker exists to avoid.
self._reservation_lock = threading.Lock()
self._plugin_locks: Dict[str, threading.Lock] = {}
self._plugin_locks_guard = threading.Lock()
self._update_worker: Optional[threading.Thread] = None
# Plugin ids whose update() has finished since the last time anyone
# asked. Updates are dispatched to a worker thread, so a caller that
# wants to know "whose data just changed" cannot learn it by diffing
# plugin_last_update around run_scheduled_updates() -- that call only
# enqueues, and the timestamp is stamped later, on the worker. See
# run_scheduled_updates_with_changes().
self._completed_updates: set = set()
self._completed_updates_lock = threading.Lock()
# Called with a plugin id the moment its data may have changed: its
# update() completed, or it called notify_vegas_data_changed(). See
# add_update_listener(). A tuple, replaced rather than mutated, so the
# worker can iterate it without a lock.
self._update_listeners: Tuple[Callable[[str], None], ...] = ()
# Where plugins' on-demand requests go: the display controller's
# submit_plugin_on_demand. See set_on_demand_handler().
self._on_demand_handler: Optional[Callable[[Dict[str, Any]], bool]] = None
# Config changes that found the plugin's lock busy, latest per plugin,
# with the instance they were meant for. See apply_config_change().
self._deferred_config_changes: Dict[str, Tuple[Any, Dict[str, Any]]] = {}
self._deferred_config_lock = threading.Lock()
# key -> (monotonic time last logged, repeats suppressed since)
self._rate_limited_warnings: Dict[str, Tuple[float, int]] = {}
self._synchronous_updates = False
if self.config_manager is not None:
try:
cfg = self.config_manager.get_config() or {}
except (OSError, ValueError) as exc:
self.logger.warning(
"Could not load config to check plugin_system.synchronous_updates "
"(%s: %s); defaulting to synchronous updates", type(exc).__name__, exc)
self._synchronous_updates = True
else:
plugin_system_cfg = cfg.get('plugin_system', {})
if not isinstance(plugin_system_cfg, dict):
self.logger.warning(
"config plugin_system must be a mapping, got %s; "
"defaulting to synchronous updates",
type(plugin_system_cfg).__name__)
self._synchronous_updates = True
else:
sync_value = plugin_system_cfg.get('synchronous_updates', False)
if not isinstance(sync_value, bool):
self.logger.warning(
"config plugin_system.synchronous_updates must be a boolean, "
"got %r; defaulting to synchronous updates", sync_value)
self._synchronous_updates = True
else:
self._synchronous_updates = sync_value
# Ensure plugins directory exists with proper permissions
try:
ensure_directory_permissions(self.plugins_dir, get_plugin_dir_mode())
except (OSError, PermissionError) as e:
self.logger.error("Could not create plugins directory %s: %s", self.plugins_dir, e, exc_info=True)
raise PluginError(f"Could not create plugins directory: {self.plugins_dir}", context={'error': str(e)}) from e
def _report_skip_once(self, key: str, message: str, *args: Any) -> None:
"""Warn about a skipped directory once per process, not per scan.
Discovery runs on every web UI page load and every config reconcile,
so warning unconditionally would put a line in the journal each time
someone opened a page -- the same log-volume problem this is meant to
help diagnose.
"""
# setdefault rather than self._skip_reported: tests build a bare
# scanner with PluginManager.__new__ and skip __init__.
reported = self.__dict__.setdefault('_skip_reported', set())
if key in reported:
return
reported.add(key)
self.logger.warning(message, *args)
def _scan_directory_for_plugins(self, directory: Path) -> List[str]:
"""
Scan a directory for plugins.
Which directories count and how an id maps to one is decided by
:class:`PluginDirectoryIndex` (``src/plugin_system/plugin_dirs.py``),
shared with the loader, the store and reconciliation. Only
``directory`` is scanned: discovery has no fallback to ``plugins/``.
Directories set aside mid-install (``BACKUP_MARKER`` in the name) are
skipped so they don't overwrite live entries.
Args:
directory: Directory to scan
Returns:
List of plugin IDs found
"""
if not directory.exists():
return []
# Build new state locally before acquiring lock
index = PluginDirectoryIndex.scan(directory)
if index.error is not None:
self.logger.error("Error scanning directory %s: %s", directory,
index.error, exc_info=index.error)
for entry in index.entries:
if entry.status == ManifestStatus.MISSING:
# A directory here that carries no manifest is not a plugin.
# Said once, because the alternative is a plugin that is
# enabled in config, enabled in plugin state, present on disk,
# and simply absent from the running process with nothing
# anywhere to say why. Working that out afterwards means
# reading cache-file mtimes.
self._report_skip_once(
entry.name, "Skipping %s: no manifest.json, so it cannot be "
"loaded as a plugin", entry.name)
elif entry.status == ManifestStatus.UNREADABLE:
self.logger.warning("Error reading manifest from %s: %s",
entry.path / "manifest.json", entry.error,
exc_info=entry.error)
elif entry.status == ManifestStatus.NOT_OBJECT:
# json.load accepts any JSON value, so a manifest holding
# null, [] or "text" parses. It once raised AttributeError on
# .get() and aborted the whole scan, so every other plugin on
# disk, however healthy, silently failed to register.
self._report_skip_once(
entry.name, "Skipping %s: its manifest.json is %s, not a "
"JSON object", entry.name, type(entry.manifest).__name__)
elif entry.status == ManifestStatus.NO_ID:
# Parsed but unusable. This was the quietest path of all: the
# manifest is read successfully and then dropped.
self._report_skip_once(
entry.name, "Skipping %s: its manifest.json has no \"id\", "
"so there is nothing to register it under", entry.name)
plugins = index.plugins()
for plugin_id, entries in index.duplicates().items():
self._report_skip_once(
"duplicate:" + plugin_id,
"Plugin id %r is declared by %d directories (%s); using %s",
plugin_id, len(entries), ", ".join(e.name for e in entries),
plugins[plugin_id].name)
new_manifests: Dict[str, Dict[str, Any]] = {
plugin_id: entry.manifest for plugin_id, entry in plugins.items()}
new_directories: Dict[str, Path] = {
plugin_id: entry.path for plugin_id, entry in plugins.items()}
# Replace shared state under lock so uninstalled plugins don't linger
with self._discovery_lock:
self.plugin_manifests.clear()
self.plugin_manifests.update(new_manifests)
self.plugin_directories.clear()
self.plugin_directories.update(new_directories)
return list(plugins)
def discover_plugins(self) -> List[str]:
"""
Discover all plugins in the plugins directory.
Also checks for potential config key collisions and logs warnings.
Returns:
List of plugin IDs
"""
self.logger.info("Discovering plugins in %s", self.plugins_dir)
plugin_ids = self._scan_directory_for_plugins(self.plugins_dir)
self.logger.info("Discovered %d plugin(s)", len(plugin_ids))
# Check for config key collisions
collisions = self.schema_manager.detect_config_key_collisions(plugin_ids)
for collision in collisions:
self.logger.warning(
"Config collision detected: %s",
collision.get('message', str(collision))
)
return plugin_ids
def load_plugin(self, plugin_id: str, force_enabled: bool = False) -> bool:
"""Load a plugin by ID; see _load_plugin.
Loading can install the plugin's dependencies with pip -- minutes,
not seconds. When that happens on the display's render thread (a
plugin enabled from the web UI, or loaded for on-demand), its
systemd watchdog gets a longer limit for the duration. Start-up
loads, on a thread pool, are covered by the start-up allowance.
"""
with display_watchdog.extended(display_watchdog.PLUGIN_LOAD_ALLOWANCE_SECONDS,
f'loading plugin {plugin_id}'):
return self._load_plugin(plugin_id, force_enabled)
def _load_plugin(self, plugin_id: str, force_enabled: bool = False) -> bool:
"""
Load a plugin by ID.
This method:
1. Checks if plugin is already loaded
2. Validates the manifest exists
3. Uses PluginLoader to import module and instantiate plugin
4. Validates the plugin configuration
5. Stores the plugin instance
6. Updates plugin state
Args:
plugin_id: Plugin identifier
force_enabled: Run the plugin enabled even though config.json has
it disabled. On-demand uses this to show a disabled plugin
(DisplayController._load_plugin_for_on_demand). Only the
instance's config says enabled; config.json is not written.
Returns:
True if loaded successfully, False otherwise
"""
if plugin_id in self.plugins:
self.logger.warning("Plugin %s already loaded", plugin_id)
return True
manifest = self.plugin_manifests.get(plugin_id)
if not manifest:
self.logger.error("No manifest found for plugin: %s", plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ERROR)
return False
try:
# Update state to LOADED
self.state_manager.set_state(plugin_id, PluginState.LOADED)
# Find plugin directory using PluginLoader
plugin_dir = self.plugin_loader.find_plugin_directory(
plugin_id,
self.plugins_dir,
self.plugin_directories
)
if plugin_dir is None:
self.logger.error("Plugin directory not found: %s", plugin_id)
self.logger.error("Searched in: %s", self.plugins_dir)
self.state_manager.set_state(plugin_id, PluginState.ERROR)
return False
# Update mapping if found via search
if plugin_id not in self.plugin_directories:
self.plugin_directories[plugin_id] = plugin_dir
# Code under this directory is this plugin's: the fetch service
# counts a request against it even from a thread the plugin
# started itself (src/common/fetch_service.py, caller identity).
register_plugin_directory(plugin_id, plugin_dir)
# Get plugin config
if self.config_manager:
full_config = self.config_manager.load_config()
config = full_config.get(plugin_id, {})
else:
config = {}
# Check if plugin has a config schema
schema = None
schema_path = self.schema_manager.get_schema_path(plugin_id)
if schema_path is None:
# Schema file doesn't exist
self.logger.warning(
f"Plugin '{plugin_id}' has no config_schema.json - configuration will not be validated. "
f"Consider adding a schema file for better error detection and user experience."
)
else:
# Schema file exists, try to load it
schema = self.schema_manager.load_schema(plugin_id)
if schema is None:
# Schema exists but couldn't be loaded (likely invalid JSON or schema)
self.logger.warning(
f"Plugin '{plugin_id}' has a config_schema.json but it could not be loaded. "
f"The schema may be invalid. Please verify the schema file at: {schema_path}"
)
# Legacy booleans read as objects, then schema defaults: the same
# preparation saves, GET /plugins/config and hot reload apply
# (prepare_plugin_config). In memory only: config.json is written
# by saves, never by loading a plugin.
config = self.prepare_plugin_config(plugin_id, config, schema=schema)
if force_enabled:
# A copy: prepare_plugin_config can hand back the section from
# config_manager's cached config, and setting the flag there
# would read as enabled to everything else in this process.
config = dict(config)
config['enabled'] = True
# Use PluginLoader to load plugin. Fetches the constructor makes
# count against the plugin.
with plugin_scope(plugin_id):
plugin_instance, _module = self.plugin_loader.load_plugin(
plugin_id=plugin_id,
manifest=manifest,
plugin_dir=plugin_dir,
config=config,
display_manager=self.display_manager,
cache_manager=self.cache_manager,
plugin_manager=self,
install_deps=True,
plugins_dir=self.plugins_dir,
)
# Register plugin-shipped fonts with the FontManager (if any).
# Plugin manifests can declare a "fonts" block that ships custom
# fonts with the plugin; FontManager.register_plugin_fonts handles
# the actual loading. Wired here so manifest declarations take
# effect without requiring plugin code changes.
font_manifest = manifest.get('fonts')
if font_manifest and self.font_manager is not None and hasattr(
self.font_manager, 'register_plugin_fonts'
):
try:
self.font_manager.register_plugin_fonts(
plugin_id, font_manifest, plugin_dir=plugin_dir)
except Exception as e:
self.logger.warning(
"Failed to register fonts for plugin %s: %s", plugin_id, e
)
# Validate configuration
if hasattr(plugin_instance, 'validate_config'):
try:
if not plugin_instance.validate_config():
self.logger.error("Plugin %s configuration validation failed", plugin_id)
self._discard_failed_load(plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ERROR)
return False
except Exception as e:
self.logger.error("Error validating plugin %s config: %s", plugin_id, e, exc_info=True)
self._discard_failed_load(plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ERROR, error=e)
return False
# Schema validation (warn/degrade only — never blocks loading).
# A config that violates the plugin's JSON schema is surfaced to the
# user (log warning + degraded flag in the health tracker) but the
# plugin still loads exactly as it does today. This deliberately does
# NOT change load_plugin()'s pass/fail behaviour for any plugin that
# loads under the current code.
self._validate_config_schema_soft(plugin_id, config)
# Store plugin instance
self.plugins[plugin_id] = plugin_instance
with self._plugin_last_update_lock:
self.plugin_last_update[plugin_id] = 0.0
# Invalidate cached interval so next tick re-derives it for this plugin
self._update_interval_cache.pop(plugin_id, None)
# Update state based on enabled status
if config.get('enabled', True):
self.state_manager.set_state(plugin_id, PluginState.ENABLED)
# Call on_enable if plugin is enabled
if hasattr(plugin_instance, 'on_enable'):
try:
with plugin_scope(plugin_id):
plugin_instance.on_enable()
except Exception:
# Undo the registration above before the outer
# handler marks it ERROR: left in self.plugins, the
# next load_plugin() would return True as "already
# loaded" for a plugin that never enabled.
self.plugins.pop(plugin_id, None)
with self._plugin_last_update_lock:
self.plugin_last_update.pop(plugin_id, None)
self._update_interval_cache.pop(plugin_id, None)
raise
else:
self.state_manager.set_state(plugin_id, PluginState.DISABLED)
# The version this instance runs, for the runtime snapshot the
# web UI reads: the manifest on disk can move on after an update.
version = manifest.get('version')
self.state_manager.record_loaded(
plugin_id, version if isinstance(version, str) else None)
self.logger.info("Loaded plugin: %s", plugin_id)
return True
except PluginError as e:
self.logger.error("Plugin error loading %s: %s", plugin_id, e, exc_info=True)
self._discard_failed_load(plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ERROR, error=e)
return False
except Exception as e:
self.logger.error("Unexpected error loading plugin %s: %s", plugin_id, e, exc_info=True)
self._discard_failed_load(plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ERROR, error=e)
return False
def _discard_failed_load(self, plugin_id: str) -> None:
"""Forget a plugin's imported module and font registrations after a
failed load.
load_module() reuses ``plugin_<id>`` from sys.modules, so a module
left behind by a load that failed after import (instantiation,
validate_config, on_enable) would keep serving the old code even
after the user fixes the plugin and reloads it. Never raises.
"""
try:
sys.modules.pop(f"plugin_{plugin_id.replace('-', '_')}", None)
self.plugin_loader.unregister_plugin_modules(plugin_id)
except Exception as e: # pragma: no cover - defensive
self.logger.debug("Could not drop modules of %s: %s", plugin_id, e)
self._forget_plugin_fonts(plugin_id)
def _forget_plugin_fonts(self, plugin_id: str) -> None:
"""Drop what the FontManager holds for a plugin: the fonts its
instance reported using (the Fonts tab's "Used by") and the fonts its
manifest registered. Never raises."""
if self.font_manager is None:
return
for name in ('forget_manager_fonts', 'forget_plugin_fonts'):
if not hasattr(self.font_manager, name):
continue
try:
getattr(self.font_manager, name)(plugin_id)
except Exception as e:
self.logger.debug("Could not forget fonts of %s (%s): %s", plugin_id, name, e)
#: Config keys the **core** reads out of a plugin's own config block. The
#: plugin never declares them, so a schema with
#: ``"additionalProperties": false`` — most published ones do — reports
#: them as violations and the plugin gets flagged degraded in the web UI for
#: using a documented core feature.
#:
#: Listed explicitly rather than matched on a ``vegas_`` prefix, because
#: ``vegas_mode`` is the opposite case: plugins *do* declare that one, and a
#: prefix rule would silently stop validating it.
#:
#: Read by: ``vegas_mode/plugin_adapter.py`` (``vegas_width_pct``,
#: ``vegas_overflow``, ``vegas_live``) and ``base_plugin.py``
#: (``vegas_max_width_screens``, ``vegas_participation``).
#:
#: The list itself lives with the other core-owned per-plugin properties in
#: ``schema_manager.CORE_PLUGIN_PROPERTIES``, which the web save path also
#: uses to keep these keys.
CORE_OWNED_CONFIG_KEYS = CORE_VEGAS_TUNING_KEYS
def prepare_plugin_config(self, plugin_id: str, config: Any,
schema: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
"""The config a plugin runs with, built from its raw config.json section.
A plugin that turned an on/off boolean into an ``{enabled, ...}``
object still finds the boolean in config.json until its settings are
next saved; it is read as the object, and schema defaults fill in the
rest (``SchemaManager.prepare_plugin_config``). Used when loading a
plugin and on hot reload (DisplayController), so ``on_config_change``
receives the same shape the plugin was constructed with.
Never raises: on failure the legacy-boolean pass alone is applied, or
failing that the section is returned as it was.
"""
if schema is None:
try:
schema = self.schema_manager.load_schema(plugin_id)
except Exception as e:
self.logger.debug("Could not load schema for %s: %s", plugin_id, e)
schema = None
upgraded: List[str] = []
try:
prepared = self.schema_manager.prepare_plugin_config(
plugin_id, config, schema=schema, changed_paths=upgraded)
self.logger.debug("Merged config with schema defaults for %s", plugin_id)
except Exception as e:
self.logger.warning("Could not apply schema defaults for %s: %s", plugin_id, e)
# Continue without defaults if they can't be applied
upgraded = []
prepared = config if isinstance(config, dict) else {}
if schema:
try:
prepared = normalize_legacy_booleans(prepared, schema, upgraded)
except Exception as legacy_error:
self.logger.warning(
"Could not read legacy boolean settings for %s: %s",
plugin_id, legacy_error)
if upgraded:
self.logger.info(
"Plugin %s: reading legacy boolean setting %s as "
"{\"enabled\": ...}; saving the plugin's settings "
"stores the new shape",
plugin_id, ", ".join(upgraded),
)
return prepared
def _strip_core_owned_keys(self, config: Dict[str, Any]) -> Dict[str, Any]:
"""A shallow copy of ``config`` without the core's own tuning keys.
Only the top level is touched, and only when such a key is present, so
the common case allocates nothing extra.
"""
if not isinstance(config, dict):
return config
if not self.CORE_OWNED_CONFIG_KEYS.intersection(config):
return config
return {k: v for k, v in config.items()
if k not in self.CORE_OWNED_CONFIG_KEYS}
def _validate_config_schema_soft(self, plugin_id: str, config: Dict[str, Any]) -> None:
"""Validate a plugin's config against its JSON schema — warn/degrade only.
On a schema violation this logs a warning and marks the plugin degraded
in the health tracker (when one is wired), so the problem is visible in
the web UI. It never raises, never changes plugin state, and never
affects whether the plugin loads. ``config`` here has already been
merged with schema defaults by the caller, so fields that ship a default
never appear "missing" — only genuinely user-supplied required fields
(e.g. an API key) can trip the required-field check.
"""
try:
schema = self.schema_manager.load_schema(plugin_id)
except Exception as e: # pragma: no cover - defensive
self.logger.debug("Could not load schema for %s: %s", plugin_id, e)
return
if not schema:
# No schema shipped — nothing to validate. Clear any stale flag.
self._set_degraded_safe(plugin_id, None)
return
try:
is_valid, errors = self.schema_manager.validate_config_against_schema(
self._strip_core_owned_keys(config), schema, plugin_id
)
except Exception as e: # pragma: no cover - defensive
# Validation machinery itself failed — do not penalise the plugin.
self.logger.debug("Schema validation raised for %s: %s", plugin_id, e)
return
if is_valid or not errors:
self._set_degraded_safe(plugin_id, None)
return
summary = "; ".join(errors[:5])
if len(errors) > 5:
summary += f" (+{len(errors) - 5} more)"
self.logger.warning(
"Plugin %s config does not match its schema (loading anyway): %s",
plugin_id, summary,
)
self._set_degraded_safe(plugin_id, f"Config schema: {summary}")
def _set_degraded_safe(self, plugin_id: str, reason: Optional[str]) -> None:
"""Best-effort ``health_tracker.set_degraded`` that never raises."""
if not self.health_tracker:
return
try:
self.health_tracker.set_degraded(plugin_id, reason)
except Exception as e: # pragma: no cover - defensive
self.logger.debug("Could not set degraded flag for %s: %s", plugin_id, e)
def unload_plugin(self, plugin_id: str) -> bool:
"""
Unload a plugin by ID.
Args:
plugin_id: Plugin identifier
Returns:
True if unloaded successfully, False otherwise
"""
if plugin_id not in self.plugins:
self.logger.warning("Plugin %s not loaded", plugin_id)
return False
# Take the plugin's lock so cleanup()/on_disable() can't run while
# the update worker is mid-update() on this instance. Bounded: an
# update() that hangs past PluginExecutor's timeout keeps holding the
# lock from its lingering thread, and unload must still go through.
lock = self.get_plugin_lock(plugin_id)
lock_acquired = lock.acquire(timeout=self.UNLOAD_LOCK_TIMEOUT)
if not lock_acquired:
self.logger.warning(
"Plugin %s still busy after %.1fs; unloading without its lock",
plugin_id, self.UNLOAD_LOCK_TIMEOUT)
try:
return self._unload_plugin_locked(plugin_id)
finally:
if lock_acquired:
lock.release()
def detach_plugin(self, plugin_id: str) -> Optional[Any]:
"""Take a loaded plugin out of ``plugins`` without tearing it down.
The first half of a reload that must not block its caller, the render
thread (DisplayController._start_plugin_reload). Every new call into a
plugin starts by looking it up in ``plugins``: the update scheduler,
the update worker (which looks again under the plugin's lock) and
Vegas's fetches. So once detached, nothing new reaches the instance.
Work already running on it under its lock -- an update(), or a Vegas
content render that can take seconds -- carries on;
unload_detached_plugin() waits for it, on another thread.
Returns the instance, or None when the plugin was not loaded.
"""
return self.plugins.pop(plugin_id, None)
def unload_detached_plugin(self, plugin_id: str, plugin: Any) -> bool:
"""Tear down an instance taken out by detach_plugin(): unload_plugin()
for an instance that is no longer in ``plugins``.
Waits for the plugin's lock, bounded by DETACHED_UNLOAD_LOCK_TIMEOUT,
so it belongs off the render thread. Call it before loading the plugin
again: it drops the plugin's modules and lifecycle state along with the
instance. Unlike unload_plugin() it never tears down without the lock:
a call that took the lock before the detach (a Vegas content build)
may still be running in this instance. Returns False then, and the
caller must not load the plugin again over it.
"""
lock = self.get_plugin_lock(plugin_id)
if not lock.acquire(timeout=self.DETACHED_UNLOAD_LOCK_TIMEOUT):
self.logger.warning(
"Plugin %s still busy after %.1fs; not unloading it while in use",
plugin_id, self.DETACHED_UNLOAD_LOCK_TIMEOUT)
return False
try:
return self._unload_plugin_locked(plugin_id, plugin)
finally:
lock.release()
def _unload_plugin_locked(self, plugin_id: str, detached: Optional[Any] = None) -> bool:
"""Body of unload_plugin(); caller holds (or gave up on) the plugin lock.
``detached`` is an instance already taken out of ``plugins``
(detach_plugin); without it, the loaded instance is unloaded.
"""
if detached is None and plugin_id not in self.plugins: # unloaded while we waited
self.logger.warning("Plugin %s not loaded", plugin_id)
return False
try:
plugin = self.plugins[plugin_id] if detached is None else detached
# Call cleanup if available
if hasattr(plugin, 'cleanup'):
try:
plugin.cleanup()
except Exception as e:
self.logger.warning("Error during plugin cleanup: %s", e)
# Call on_disable if available
if hasattr(plugin, 'on_disable'):
try:
plugin.on_disable()
except Exception as e:
self.logger.warning("Error during plugin on_disable: %s", e)
# Remove from active plugins (a detached one already is)
if detached is None:
del self.plugins[plugin_id]
with self._deferred_config_lock:
self._deferred_config_changes.pop(plugin_id, None)
with self._plugin_last_update_lock:
self.plugin_last_update.pop(plugin_id, None)
self._update_interval_cache.pop(plugin_id, None)
# Remove main module from sys.modules if present
module_name = f"plugin_{plugin_id.replace('-', '_')}"
sys.modules.pop(module_name, None)
# Delegate sub-module and cached-module cleanup to the loader
self.plugin_loader.unregister_plugin_modules(plugin_id)
# Its font registrations go with it: the fonts it reported using
# and the ones its manifest registered.
self._forget_plugin_fonts(plugin_id)
# Update state
self.state_manager.set_state(plugin_id, PluginState.UNLOADED)
self.state_manager.clear_state(plugin_id)
self.logger.info("Unloaded plugin: %s", plugin_id)
return True
except Exception as e:
self.logger.error("Error unloading plugin %s: %s", plugin_id, e, exc_info=True)
self.state_manager.set_state(plugin_id, PluginState.ERROR, error=e)
if plugin_id not in self.plugins:
# Failed after the instance was dropped: it is not loaded.
self.state_manager.record_unloaded(plugin_id)
return False
def reload_plugin(self, plugin_id: str) -> bool:
"""
Reload a plugin (unload and load).
Args:
plugin_id: Plugin identifier
Returns:
True if reloaded successfully, False otherwise
"""
self.logger.info("Reloading plugin: %s", plugin_id)
# Unload first
if plugin_id in self.plugins:
if not self.unload_plugin(plugin_id):
return False
# Re-read the manifest so an edit to it takes effect, from the
# directory discovery found the plugin in: a directory's name need not
# be the id its manifest declares.
with self._discovery_lock:
directories = dict(self.plugin_directories)
plugin_dir = self.plugin_loader.find_plugin_directory(
plugin_id, self.plugins_dir, directories)
manifest_path = plugin_dir / "manifest.json" if plugin_dir is not None else None
if manifest_path is not None and manifest_path.exists():
try:
with open(manifest_path, 'r', encoding='utf-8') as f:
manifest = json.load(f)
with self._discovery_lock:
self.plugin_manifests[plugin_id] = manifest
except Exception as e:
self.logger.error("Error reading manifest: %s", e, exc_info=True)
return False
return self.load_plugin(plugin_id)
def discovered_plugin_ids(self) -> set:
"""Snapshot of the discovered plugin ids, taken under the discovery lock.
Callers on other threads (the config watcher) must not iterate
``plugin_manifests`` directly: discovery rebuilds it entry by entry, so
an unsynchronised reader can see a half-populated mapping or raise
"dictionary changed size during iteration".
"""
with self._discovery_lock:
return set(self.plugin_manifests)
def get_plugin(self, plugin_id: str) -> Optional[Any]:
"""
Get a loaded plugin instance by ID.
Args:
plugin_id: Plugin identifier
Returns:
Plugin instance or None if not loaded
"""
return self.plugins.get(plugin_id)
def get_all_plugins(self) -> Dict[str, Any]:
"""
Get all loaded plugins.
Returns:
Dict of plugin_id: plugin_instance
"""
return self.plugins.copy()
def get_plugin_info(self, plugin_id: str) -> Optional[Dict[str, Any]]:
"""
Get information about a plugin (manifest + runtime info).
Args:
plugin_id: Plugin identifier
Returns:
Dict with plugin information or None if not found
"""
with self._discovery_lock:
manifest = self.plugin_manifests.get(plugin_id)
if not manifest:
return None
info = manifest.copy()
# Add runtime information if plugin is loaded
plugin = self.plugins.get(plugin_id)
if plugin:
info['loaded'] = True
if hasattr(plugin, 'get_info'):
# One plugin's get_info() raising must not take down the
# whole installed-plugins listing (/api/v3/plugins/installed).
try:
info['runtime_info'] = plugin.get_info()
except Exception as e:
self.logger.warning("Plugin %s get_info() failed: %s", plugin_id, e)
info['runtime_info'] = {}
else:
info['loaded'] = False
# Add state information
info['state'] = self.state_manager.get_state_info(plugin_id)
return info
def get_all_plugin_info(self) -> List[Dict[str, Any]]:
"""
Get information about all plugins.
Returns:
List of plugin info dictionaries
"""
with self._discovery_lock:
pids = list(self.plugin_manifests.keys())
return [info for info in [self.get_plugin_info(pid) for pid in pids] if info]
def get_plugin_directory(self, plugin_id: str) -> Optional[str]:
"""