aboutsummaryrefslogtreecommitdiffstats
path: root/app/PluginManager.py
blob: cd810b0664e075006236817fd532e70ea68fc017 (plain)
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
import importlib
import schedule
import time
from types import ModuleType
from typing import Any
from app.Item import Item
from app.PluginInterface import Params, PluginInterface
from app.AggroConfig import AggroConfig
from app.utils import get_config
from app.MemoryState import memory_state


class PluginManager:
    def __init__(self, config: AggroConfig) -> None:
        self._plugins: dict[str, Any] = {}
        self.plugin_instances: dict[str, PluginInterface] = {}
        self.config = config
        self.scheduled_plugin_ids: list[str] = []

    def load_plugin(self, plugin_name: str) -> None:
        module: ModuleType = importlib.import_module(f"plugins.{plugin_name}")
        plugin_class = module.Plugin
        self._plugins[plugin_name] = plugin_class

    def propagate(self, id: str, items: list[Item]):
        next_nodes: list[str]
        if id not in self.config.graph:
            next_nodes = []
        else:
            next_nodes = self.config.graph[id]

        for next_node_id in next_nodes:
            self.run_plugin_job(next_node_id, id, items)

    def run_plugin_job(self, id: str, source_id: str | None, items: list[Item] = []):
        if not memory_state.running:
            return

        plugin: PluginInterface = self.plugin_instances[id]
        ret_items: list[Item] = plugin.process(source_id, items)
        self.propagate(id, ret_items)

    def build_plugin_instances(self):
        for id in self.config.plugins:
            params: Params = self.config.plugins[id]

            plugin_name: str = get_config(params, "plugin")
            schedule_expr: str | None = params.get("schedule_expr", None)

            if plugin_name not in self._plugins:
                self.load_plugin(plugin_name)

            if schedule_expr is not None:
                job: schedule.Job = eval(schedule_expr, {"schedule": schedule})
                job.do(self.run_plugin_job, id, None)
                self.scheduled_plugin_ids.append(id)

            PluginClass: Any = self._plugins[plugin_name]
            plugin: PluginInterface = PluginClass(id=id, params=params)
            self.plugin_instances[id] = plugin

    def initial_run_scheduled_plugins(self):
        for id in self.scheduled_plugin_ids:
            self.run_plugin_job(id, None)

    def run(self) -> None:
        while memory_state.running:
            schedule.run_pending()
            time.sleep(1)