diff options
| author | Jan Tuomi <jans.tuomi@gmail.com> | 2023-09-03 11:20:39 +0300 |
|---|---|---|
| committer | Jan Tuomi <jans.tuomi@gmail.com> | 2023-09-10 19:01:00 +0300 |
| commit | 21052a42c84852c3964a9a29694e6d2a7fe3b9e8 (patch) | |
| tree | 341c35bdc567b51f192acd8f97bd24b65d47f995 /app/ConcatPlugin.py | |
| parent | 3a4cd4935dafe6d445daac3b7d6ced409f99aa23 (diff) | |
Check if Aggrofile has been changed at startup, add ConcatPlugin
Diffstat (limited to 'app/ConcatPlugin.py')
| -rw-r--r-- | app/ConcatPlugin.py | 44 |
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 |
