Files
LEDMatrix/src/plugin_system/operation_queue.py
T
ChuckandClaude Opus 5.5 b11bcfa204 fix(plugins): store and plugin-manager bugs; tidy src/plugin_system (#635)
* fix(store): don't read a ZIP-installed plugin's remote from the LEDMatrix repo

update_plugin looked up remote.origin.url with `git -C <plugin> config
--local` for plugins that are not git checkouts. Under plugin-repos/ git
walks up to the enclosing LEDMatrix repository, so the lookup returned
LEDMatrix's own URL and a plugin missing from the registry was
"reinstalled" from the LEDMatrix repo. Only ask git when the plugin
directory has its own .git, the test _get_local_git_info already uses.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(schema): report each missing required field once, by name

validate_config_against_schema ran its own required-fields loop after
Draft7Validator.iter_errors, which already yields one `required` error
per missing field, so every missing top-level field was listed twice.
The validator's copy also printed the schema's whole `required` list
("Missing required property '['api_key', 'city']'") instead of the field.
Drop the loop and take the field name from the error itself.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(store): stop mangling repository URLs that contain ".git"

install_from_url and fetch_registry_from_url cleaned URLs with
`rstrip('/').replace('.git', '')`, which removes ".git" anywhere:
https://github.com/user/my.github.io became .../myhub.io, so installing
or browsing that repository asked GitHub for one that does not exist.

Add src/plugin_system/repo_urls.py with one anchored normalize_repo_url(),
same_repo() for comparisons, github_owner_repo() and github_api_headers(),
and use them for the five copies of the owner/repo parsing and GitHub
headers in the store and for saved repositories. GitHub URLs are now
recognised by urlparse().hostname everywhere: _get_latest_commit_info
used a substring test, and _install_from_monorepo_api parsed any host.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(store): install a repository whose only branch is not main/master

_install_via_git returned None both when every clone failed and when the
last-resort clone of the repository's default branch succeeded.
_install_plugin_impl papered over it with `and not plugin_path.exists()`;
install_from_url did not, so a repository whose only branch is e.g.
`develop` was cloned, then treated as a failure, then "downloaded" from
main/master archives that do not exist.

After a default-branch clone, return the branch the clone checked out
(read from .git/HEAD), so None means failure and nothing else, and give
both callers the same `branch_used is None` fallback.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(plugins): judge the memory limit on each call's own growth

monitor_call stores `metrics.memory_mb = max(previous, growth)`, and
_check_limits compared that high-water mark with max_memory_mb. It never
decreases, so once one update() grew the process past the limit every
later call raised ResourceLimitExceeded and the circuit breaker kept
reopening. Pass the call's own RSS growth to _check_limits; keep the
high-water mark for reporting and document what it measures.

Remove ResourceMetrics.update_average_execution_time: nothing called it,
and it overwrote the running total with the average.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(plugins): reload_plugin re-reads the manifest from the discovered directory

reload_plugin read `plugins_dir / plugin_id / "manifest.json"`, ignoring
the discovery map and the plugin_dirs rules. For a plugin whose
directory name differs from its manifest id the path did not exist, the
re-read was skipped without a word, and the reload kept the stale
manifest. Resolve the directory with find_plugin_directory, as
load_plugin does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix(plugins): drop the always-null last_display from plugin state info

PluginStateManager reported `last_display` from `_last_display`, which
nothing ever wrote, so it was null for every plugin. Recording it in
PluginExecutor.execute_display would not help: get_state_info's only
reader is the web process, whose PluginManager never calls display().
Remove the field, its dict and get_last_display() (no caller in core,
the web UI or the plugin monorepo).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(store): share the rollback and requirements helpers, drop dead code

- install_plugin and _reinstall_with_rollback set aside, discard and
  restore the old copy through _set_aside/_discard_backup/_restore_backup
  instead of two copies of the same blocks.
- The loader and the store run the same pre-pip checks through
  contained_plugin_dir() and requirements_to_install() in plugin_loader.
  They still invoke pip differently (sys.executable -m pip vs. the sudo
  wrapper). `except (BrokenPipeError, OSError)` + `isinstance(e, OSError)`
  becomes `except OSError` checking errno.EPIPE.
- load_module never returns None, so load_plugin's check is gone and the
  docstring says what it raises.
- Remove the always-true JSONSCHEMA_AVAILABLE, the inline re-imports of
  re and permission_utils, the fake status_result object nobody reads,
  hasattr(git_error, 'cmd'), a redundant "merge conflict" test and
  `import traceback` (exc_info=True does it).
- Correct comments: install_from_url names the directory for the
  caller's id when given (not always the manifest id), _get_local_git_info
  saves one git subprocess (not four), _enrich calls two helpers,
  search_plugins documents all its arguments, _find_plugin_path states
  its behaviour instead of a TODO, and history narration is gone.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(plugins): tidy base_plugin, correct plugin_manager/state comments

- base_plugin: drop the unused `import logging`; get_display_duration
  runs the instance value and the config value through one
  _positive_seconds() helper instead of two copies of the coercion; the
  'static'/'none'/fallback branches of get_vegas_display_mode, which all
  returned FIXED_SEGMENT, are one; fix the mis-indented validate_config
  example; say that get_supported_vegas_modes/get_vegas_segment_width
  are not consulted by core (kept, plugins override them).
- schema_manager: import expand_style_elements normally rather than
  swallowing an ImportError of a core module.
- plugin_manager: the plugins directory is the configured one
  (plugin-repos/ by default), not plugins/; get_config() returns the live
  dict, not a copy, so the interval cache comments say what it saves.
- state_manager: config_version and the file version are not used to
  detect corruption; say what they are.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(plugins): stop writing data/plugin_operations.json

PluginOperationQueue wrote its finished-operation history to
data/plugin_operations.json after every operation, and read it back only
into its own in-memory list, which only get_operation_history() exposes
-- and nothing calls that. The operation-history endpoint reads
OperationHistory (data/operation_history.json). No code in src/,
web_interface/, scripts/ or test/ reads the file.

Drop the history_file/lazy_load parameters and the load/save code; the
bounded in-memory history stays. web_interface/app.py and the
integration test stop passing the removed arguments. An existing
data/plugin_operations.json is left in place (data/* is gitignored).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(changelog): plugin-system

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 17:32:02 -04:00

315 lines
12 KiB
Python

"""
Plugin operation queue manager.
Serializes plugin operations to prevent conflicts and provides
status tracking and cancellation support.
"""
import threading
import queue
from typing import Dict, Optional, List, Callable, Any
from datetime import datetime
from src.plugin_system.operation_types import (
PluginOperation, OperationType, OperationStatus
)
from src.logging_config import get_logger
class PluginOperationQueue:
"""
Manages a queue of plugin operations, executing them serially
to prevent conflicts.
Features:
- Serialized execution (one operation at a time)
- Prevents concurrent operations on same plugin
- Operation status tracking
- Operation cancellation
- In-memory history of finished operations
The history is not persisted. The web UI's operation history comes from
OperationHistory (operation_history.py), which has its own file; a copy
written here was never read back by anything.
"""
def __init__(self, max_history: int = 100):
"""
Initialize operation queue.
Args:
max_history: Maximum number of operations to keep in history
"""
self.logger = get_logger(__name__)
self.max_history = max_history
# Operation tracking
self._operations: Dict[str, PluginOperation] = {}
self._operation_queue: queue.Queue = queue.Queue()
self._active_operations: Dict[str, PluginOperation] = {} # plugin_id -> operation
self._operation_history: List[PluginOperation] = []
# Threading
self._lock = threading.RLock()
self._worker_thread: Optional[threading.Thread] = None
self._stop_event = threading.Event()
self._start_worker()
def enqueue_operation(
self,
operation_type: OperationType,
plugin_id: str,
parameters: Optional[Dict] = None,
operation_callback: Optional[Callable[[PluginOperation], Dict[str, Any]]] = None
) -> str:
"""
Enqueue a plugin operation.
Args:
operation_type: Type of operation to perform
plugin_id: Plugin identifier
parameters: Optional operation parameters
operation_callback: Optional callback function to execute the operation.
If None, operation will be queued but not executed.
Returns:
Operation ID for tracking
"""
with self._lock:
# Check if plugin already has an active operation
if plugin_id in self._active_operations:
active_op = self._active_operations[plugin_id]
if active_op.status in [OperationStatus.PENDING, OperationStatus.RUNNING]:
raise ValueError(
f"Plugin {plugin_id} already has an active operation: "
f"{active_op.operation_id} ({active_op.operation_type.value})"
)
# Create operation
operation = PluginOperation(
operation_type=operation_type,
plugin_id=plugin_id,
parameters=parameters or {},
)
# Store callback if provided
if operation_callback:
operation.parameters['_callback'] = operation_callback
# Store operation
self._operations[operation.operation_id] = operation
# Enqueue
self._operation_queue.put(operation)
self.logger.info(
f"Enqueued {operation_type.value} operation for plugin {plugin_id} "
f"(operation_id: {operation.operation_id})"
)
return operation.operation_id
def get_operation_status(self, operation_id: str) -> Optional[PluginOperation]:
"""
Get status of an operation.
Args:
operation_id: Operation identifier
Returns:
PluginOperation if found, None otherwise
"""
with self._lock:
return self._operations.get(operation_id)
def cancel_operation(self, operation_id: str) -> bool:
"""
Cancel a pending operation.
Args:
operation_id: Operation identifier
Returns:
True if operation was cancelled, False if not found or already running
"""
with self._lock:
operation = self._operations.get(operation_id)
if not operation:
return False
if operation.status == OperationStatus.RUNNING:
self.logger.warning(
f"Cannot cancel running operation {operation_id}"
)
return False
if operation.status == OperationStatus.PENDING:
operation.status = OperationStatus.CANCELLED
operation.completed_at = datetime.now()
operation.message = "Operation cancelled by user"
self._add_to_history(operation)
self.logger.info(f"Cancelled operation {operation_id}")
return True
return False
def get_operation_history(self, limit: int = 50) -> List[PluginOperation]:
"""
Get operation history.
Args:
limit: Maximum number of operations to return
Returns:
List of operations, sorted by creation time (newest first)
"""
with self._lock:
# Sort by creation time (newest first)
history = sorted(
self._operation_history,
key=lambda op: op.created_at,
reverse=True
)
return history[:limit]
def get_active_operations(self) -> List[PluginOperation]:
"""
Get all currently active operations (pending or running).
Returns:
List of active operations
"""
with self._lock:
active = []
for operation in self._operations.values():
if operation.status in [OperationStatus.PENDING, OperationStatus.RUNNING]:
active.append(operation)
return active
def _start_worker(self) -> None:
"""Start the worker thread that processes operations."""
if self._worker_thread and self._worker_thread.is_alive():
return
self._stop_event.clear()
self._worker_thread = threading.Thread(
target=self._worker_loop,
daemon=True,
name="PluginOperationQueueWorker"
)
self._worker_thread.start()
self.logger.info("Started plugin operation queue worker thread")
def _worker_loop(self) -> None:
"""Worker thread loop that processes queued operations."""
while not self._stop_event.is_set():
try:
# Get next operation (with timeout to allow checking stop event)
try:
operation = self._operation_queue.get(timeout=1.0)
except queue.Empty:
continue
# Check if operation was cancelled
if operation.status == OperationStatus.CANCELLED:
self._operation_queue.task_done()
continue
# Execute operation
self._execute_operation(operation)
self._operation_queue.task_done()
except Exception as e:
self.logger.error(f"Error in operation queue worker: {e}", exc_info=True)
def _execute_operation(self, operation: PluginOperation) -> None:
"""
Execute a plugin operation.
Args:
operation: Operation to execute
"""
with self._lock:
# Check if plugin already has active operation
if operation.plugin_id in self._active_operations:
active_op = self._active_operations[operation.plugin_id]
if active_op.operation_id != operation.operation_id:
# Different operation for same plugin - mark as failed
operation.status = OperationStatus.FAILED
operation.error = f"Plugin {operation.plugin_id} has another active operation"
operation.completed_at = datetime.now()
self._add_to_history(operation)
return
# Mark as running
operation.status = OperationStatus.RUNNING
operation.started_at = datetime.now()
operation.progress = 0.0
self._active_operations[operation.plugin_id] = operation
try:
self.logger.info(
f"Executing {operation.operation_type.value} operation for "
f"plugin {operation.plugin_id} (operation_id: {operation.operation_id})"
)
# Get callback from parameters
callback = operation.parameters.pop('_callback', None)
if callback:
# Execute callback
operation.progress = 0.1
result = callback(operation)
# Update operation with result
operation.progress = 1.0
operation.status = OperationStatus.COMPLETED
operation.result = result
operation.message = result.get('message', 'Operation completed successfully')
else:
# No callback - mark as completed (operation was just queued)
operation.progress = 1.0
operation.status = OperationStatus.COMPLETED
operation.message = "Operation queued (no callback provided)"
except Exception as e:
self.logger.error(
f"Error executing operation {operation.operation_id}: {e}",
exc_info=True
)
operation.status = OperationStatus.FAILED
operation.error = str(e)
operation.message = f"Operation failed: {str(e)}"
finally:
with self._lock:
operation.completed_at = datetime.now()
# Remove from active operations
if operation.plugin_id in self._active_operations:
if self._active_operations[operation.plugin_id].operation_id == operation.operation_id:
del self._active_operations[operation.plugin_id]
self._add_to_history(operation)
def _add_to_history(self, operation: PluginOperation) -> None:
"""Add operation to history, maintaining max_history limit."""
self._operation_history.append(operation)
# Trim history if needed
if len(self._operation_history) > self.max_history:
# Remove oldest operations
self._operation_history.sort(key=lambda op: op.created_at)
self._operation_history = self._operation_history[-self.max_history:]
def shutdown(self) -> None:
"""Shutdown the operation queue and worker thread."""
self.logger.info("Shutting down plugin operation queue")
self._stop_event.set()
if self._worker_thread and self._worker_thread.is_alive():
self._worker_thread.join(timeout=5.0)