aboutsummaryrefslogtreecommitdiffstats
path: root/app/PluginManager.py
blob: b19d0380ce5a1dfb760c7884962c684f1ebe47b0 (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
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
import importlib
import schedule
import time
import sys
import traceback
from types import ModuleType
from typing import Any
from app.Item import Item
from app.PluginInterface import Params, PluginInterface
from app.AggroConfig import AggroConfig, AggroConfigSendGridAlerter
from app.utils import get_config
from app.MemoryState import memory_state
from app.SendGridAlerter import SendGridAlerter


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] = []
        if self.config.sendgrid_alerter:
            self.sendgrid_alerter = SendGridAlerter.from_config(
                self.config.sendgrid_alerter
            )

    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

        try:
            try:
                plugin: PluginInterface = self.plugin_instances[id]
                ret_items: list[Item] = plugin.process(source_id, items)
                self.propagate(id, ret_items)
            except Exception as ex:
                ex.add_note(f"context: plugin id: {id}, source id: {source_id}")
                raise

        except:
            exc = traceback.format_exc()
            print(exc, file=sys.stderr)
            if self.sendgrid_alerter:
                self.sendgrid_alerter.send_alert(exc)

    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)