aboutsummaryrefslogtreecommitdiffstats
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/AggroConfig.py7
-rw-r--r--src/FeedSourcePlugin.py28
-rw-r--r--src/FilterPlugin.py22
-rw-r--r--src/Item.py10
-rw-r--r--src/PluginInterface.py18
-rw-r--r--src/PluginManager.py71
-rw-r--r--src/aggro.py24
-rw-r--r--src/utils.py5
8 files changed, 185 insertions, 0 deletions
diff --git a/src/AggroConfig.py b/src/AggroConfig.py
new file mode 100644
index 0000000..e933eec
--- /dev/null
+++ b/src/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/src/FeedSourcePlugin.py b/src/FeedSourcePlugin.py
new file mode 100644
index 0000000..e6e4366
--- /dev/null
+++ b/src/FeedSourcePlugin.py
@@ -0,0 +1,28 @@
+import feedparser # type: ignore
+from typing import Any
+from Item import Item
+from PluginInterface import PluginInterface
+from 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/src/FilterPlugin.py b/src/FilterPlugin.py
new file mode 100644
index 0000000..6524d20
--- /dev/null
+++ b/src/FilterPlugin.py
@@ -0,0 +1,22 @@
+from typing import Any
+from Item import Item
+from PluginInterface import PluginInterface
+from 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/src/Item.py b/src/Item.py
new file mode 100644
index 0000000..e74239f
--- /dev/null
+++ b/src/Item.py
@@ -0,0 +1,10 @@
+from dataclasses import dataclass
+
+
+@dataclass
+class Item:
+ """RSS Item"""
+
+ title: str
+ link: str
+ description: str
diff --git a/src/PluginInterface.py b/src/PluginInterface.py
new file mode 100644
index 0000000..c6be538
--- /dev/null
+++ b/src/PluginInterface.py
@@ -0,0 +1,18 @@
+from abc import ABC, abstractmethod
+from typing import Any
+
+from 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/src/PluginManager.py b/src/PluginManager.py
new file mode 100644
index 0000000..4a9bd13
--- /dev/null
+++ b/src/PluginManager.py
@@ -0,0 +1,71 @@
+import importlib
+import schedule
+import time
+from types import ModuleType
+from typing import Any
+from AggroConfig import AggroConfig
+from Item import Item
+from PluginInterface import PluginInterface
+from 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(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/src/aggro.py b/src/aggro.py
new file mode 100644
index 0000000..22e8c0b
--- /dev/null
+++ b/src/aggro.py
@@ -0,0 +1,24 @@
+import json
+from AggroConfig import AggroConfig
+from PluginManager import PluginManager
+import os
+
+if __name__ == "__main__":
+ print("Starting aggro. Press CTRL-C to exit.")
+
+ aggrofile_path = os.environ.get("AGGRO_CONFIG_PATH", "Aggrofile")
+
+ with open(aggrofile_path) as f:
+ aggrofile_content = json.loads(f.read())
+
+ aggro_config = AggroConfig(
+ plugins=aggrofile_content["plugins"], graph=aggrofile_content["graph"]
+ )
+
+ manager = PluginManager(aggro_config)
+ manager.build_plugin_instances()
+
+ try:
+ manager.run()
+ except KeyboardInterrupt:
+ print("Exiting...")
diff --git a/src/utils.py b/src/utils.py
new file mode 100644
index 0000000..ac07bf7
--- /dev/null
+++ b/src/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]