aboutsummaryrefslogtreecommitdiffstats
path: root/app
diff options
context:
space:
mode:
authorJan Tuomi <jans.tuomi@gmail.com>2023-09-01 19:31:47 +0300
committerJan Tuomi <jans.tuomi@gmail.com>2023-09-10 19:01:00 +0300
commit398799de0c1826c3544eb6eec681fa48096932a0 (patch)
treec2a2cf7fcf30276ffc432e7a1750efd6b7d2c8bf /app
parentb7f129b29b8b8103336bfae265b3d5a0bc5ea03a (diff)
Refactor
Diffstat (limited to 'app')
-rw-r--r--app/AggroConfig.py7
-rw-r--r--app/FeedSourcePlugin.py28
-rw-r--r--app/FilterPlugin.py22
-rw-r--r--app/Item.py10
-rw-r--r--app/PluginInterface.py18
-rw-r--r--app/PluginManager.py71
-rw-r--r--app/utils.py5
7 files changed, 161 insertions, 0 deletions
diff --git a/app/AggroConfig.py b/app/AggroConfig.py
new file mode 100644
index 0000000..e933eec
--- /dev/null
+++ b/app/AggroConfig.py
@@ -0,0 +1,7 @@
+from dataclasses import dataclass
+
+
+@dataclass
+class AggroConfig:
+ plugins: dict[str, dict[str, str]]
+ graph: dict[str, list[str]]
diff --git a/app/FeedSourcePlugin.py b/app/FeedSourcePlugin.py
new file mode 100644
index 0000000..af686b9
--- /dev/null
+++ b/app/FeedSourcePlugin.py
@@ -0,0 +1,28 @@
+import feedparser # type: ignore
+from typing import Any
+from app.Item import Item
+from app.PluginInterface import PluginInterface
+from app.utils import get_param
+
+
+class Plugin(PluginInterface):
+ def __init__(self, id: str, params: dict[str, Any]) -> None:
+ super().__init__(id, params)
+ self.feed_url: str = get_param("feed_url", params)
+
+ print(f"[FeedSourcePlugin#{self.id}] initialized")
+
+ def validate_input_n(self, n: int) -> bool:
+ return n == 0
+
+ def process(self, inputs: list[list[Item]]) -> list[Item]:
+ print(f"[FeedSourcePlugin#{self.id}] process called")
+ feed: Any = feedparser.parse(self.feed_url) # type: ignore
+ items: list[Item] = []
+ for d in feed.entries:
+ items.append(
+ Item(title=d["title"], link=d["link"], description=d["description"])
+ )
+
+ print(f"[FeedSourcePlugin#{self.id}] process returns items, n={len(items)}")
+ return items
diff --git a/app/FilterPlugin.py b/app/FilterPlugin.py
new file mode 100644
index 0000000..fd3984a
--- /dev/null
+++ b/app/FilterPlugin.py
@@ -0,0 +1,22 @@
+from typing import Any
+from app.Item import Item
+from app.PluginInterface import PluginInterface
+from app.utils import get_param
+
+
+class Plugin(PluginInterface):
+ def __init__(self, id: str, params: dict[str, Any]) -> None:
+ super().__init__(id, params)
+ self.filter_expr: str = get_param("filter_expr", params)
+ print(f"[FilterPlugin#{self.id}] initialized")
+
+ def validate_input_n(self, n: int) -> bool:
+ return n == 1
+
+ def process(self, inputs: list[list[Item]]) -> list[Item]:
+ input_feed: list[Item] = inputs[0]
+ print(f"[FilterPlugin#{self.id}] process called, n={len(input_feed)}")
+ expr = f"filter(lambda item: {self.filter_expr}, input_feed)"
+ ret = list(eval(expr, {"input_feed": input_feed}))
+ print(f"[FilterPlugin#{self.id}] process returning items, n={len(ret)}")
+ return ret
diff --git a/app/Item.py b/app/Item.py
new file mode 100644
index 0000000..e74239f
--- /dev/null
+++ b/app/Item.py
@@ -0,0 +1,10 @@
+from dataclasses import dataclass
+
+
+@dataclass
+class Item:
+ """RSS Item"""
+
+ title: str
+ link: str
+ description: str
diff --git a/app/PluginInterface.py b/app/PluginInterface.py
new file mode 100644
index 0000000..6194436
--- /dev/null
+++ b/app/PluginInterface.py
@@ -0,0 +1,18 @@
+from abc import ABC, abstractmethod
+from typing import Any
+
+from app.Item import Item
+
+
+class PluginInterface(ABC):
+ def __init__(self, id: str, params: dict[str, Any]):
+ self.id = id
+ self.params = params
+
+ @abstractmethod
+ def validate_input_n(self, n: int) -> bool:
+ raise NotImplemented("abstract method not implemented")
+
+ @abstractmethod
+ def process(self, inputs: list[list[Item]]) -> list[Item]:
+ pass
diff --git a/app/PluginManager.py b/app/PluginManager.py
new file mode 100644
index 0000000..8d60bc1
--- /dev/null
+++ b/app/PluginManager.py
@@ -0,0 +1,71 @@
+import importlib
+import schedule
+import time
+from types import ModuleType
+from typing import Any
+from app.Item import Item
+from app.PluginInterface import PluginInterface
+from app.AggroConfig import AggroConfig
+from app.utils import get_param
+
+
+class PluginManager:
+ def __init__(self, config: AggroConfig) -> None:
+ self._plugins: dict[str, Any] = {}
+ self.plugin_instances: dict[str, PluginInterface] = {}
+ self.running = False
+ self.config = config
+
+ def load_plugin(self, plugin_name: str) -> None:
+ module: ModuleType = importlib.import_module(f"app.{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, items)
+
+ def run_plugin_job(self, id: str, items: list[Item] = []):
+ if not self.running:
+ return
+
+ plugin: PluginInterface = self.plugin_instances[id]
+ ret_items: list[Item] = plugin.process([items])
+ self.propagate(id, ret_items)
+
+ def build_plugin_instances(self):
+ self.graph: dict[str, PluginInterface] = {}
+ for id in self.config.plugins:
+ params: dict[str, str] = self.config.plugins[id]
+
+ plugin_name = get_param("plugin", params)
+ trigger_type = get_param("trigger_type", params)
+
+ if plugin_name not in self._plugins:
+ self.load_plugin(plugin_name)
+
+ match trigger_type:
+ case "schedule":
+ schedule_expr = get_param("schedule_expr", params)
+ job: schedule.Job = eval(schedule_expr, {"schedule": schedule})
+ job.do(self.run_plugin_job, id) # type: ignore
+ case "input_change":
+ pass
+ case _:
+ raise Exception("unknown trigger_type: " + trigger_type)
+
+ PluginClass: Any = self._plugins[plugin_name]
+ plugin: PluginInterface = PluginClass(id=id, params=params)
+ self.plugin_instances[id] = plugin
+
+ def run(self) -> None:
+ self.running = True
+ while self.running:
+ schedule.run_pending()
+ time.sleep(1)
diff --git a/app/utils.py b/app/utils.py
new file mode 100644
index 0000000..ac07bf7
--- /dev/null
+++ b/app/utils.py
@@ -0,0 +1,5 @@
+def get_param(key: str, params: dict[str, str]) -> str:
+ if key not in params:
+ raise Exception(f"no {key} field in config entry: " + str(params))
+
+ return params[key]