aboutsummaryrefslogtreecommitdiffstats
path: root/app/ConcatPlugin.py
diff options
context:
space:
mode:
authorJan Tuomi <jans.tuomi@gmail.com>2023-09-03 11:20:39 +0300
committerJan Tuomi <jans.tuomi@gmail.com>2023-09-10 19:01:00 +0300
commit21052a42c84852c3964a9a29694e6d2a7fe3b9e8 (patch)
tree341c35bdc567b51f192acd8f97bd24b65d47f995 /app/ConcatPlugin.py
parent3a4cd4935dafe6d445daac3b7d6ced409f99aa23 (diff)
Check if Aggrofile has been changed at startup, add ConcatPlugin
Diffstat (limited to 'app/ConcatPlugin.py')
-rw-r--r--app/ConcatPlugin.py44
1 files changed, 44 insertions, 0 deletions
diff --git a/app/ConcatPlugin.py b/app/ConcatPlugin.py
new file mode 100644
index 0000000..14307ac
--- /dev/null
+++ b/app/ConcatPlugin.py
@@ -0,0 +1,44 @@
+from typing import Any
+
+from tinydb import Query
+from app.Item import Item
+from app.PluginInterface import PluginInterface
+from app.DatabaseManager import database_manager
+from app.utils import dict_to_item, item_to_dict
+
+# TODO
+
+
+class Plugin(PluginInterface):
+ def __init__(self, id: str, params: dict[str, Any]) -> None:
+ super().__init__(id, params)
+ print(f"[ConcatPlugin#{self.id}] initialized")
+
+ def process(self, source_id: str | None, items: list[Item]) -> list[Item]:
+ print(f"[ConcatPlugin#{self.id}] process called, n={len(items)}")
+ if source_id is None:
+ raise Exception(f"[ConcatPlugin#{self.id}] can not be scheduled")
+
+ Q = Query()
+ q_result = database_manager.plugin_states.search(Q.plugin_id == self.id)
+ if len(q_result) > 1:
+ raise Exception(
+ f"[ConcatPlugin#{self.id}] invalid database state, found more than 1 plugin state: {len(q_result)}"
+ )
+
+ d: Any = q_result[0] if len(q_result) == 1 else {}
+ data: dict[str, list[dict[str, str]]] = d
+
+ if source_id in data:
+ items_as_dicts = list(map(item_to_dict, items))
+ d[source_id] = items_as_dicts
+
+ aggregated_item_dicts: list[dict[str, str]] = []
+ for data_source_id in data:
+ item_dicts = data[data_source_id]
+ aggregated_item_dicts += item_dicts
+
+ ret = list(map(dict_to_item, aggregated_item_dicts))
+
+ print(f"[ConcatPlugin#{self.id}] process returning items, n={len(ret)}")
+ return ret