diff options
| author | Jan Tuomi <jans.tuomi@gmail.com> | 2023-09-03 19:07:30 +0300 |
|---|---|---|
| committer | Jan Tuomi <jans.tuomi@gmail.com> | 2023-09-10 19:01:00 +0300 |
| commit | 7f2e5ee5bf59e5814d62990401714e01cd69ba07 (patch) | |
| tree | 0f3f9fcba77ccf45d81f6090df0195c7f34e98d3 /plugins/ConcatPlugin.py | |
| parent | 7c90b13a7eaf520126a40180f932125240b46364 (diff) | |
Refactor plugins to own dir
Diffstat (limited to 'plugins/ConcatPlugin.py')
| -rw-r--r-- | plugins/ConcatPlugin.py | 50 |
1 files changed, 50 insertions, 0 deletions
diff --git a/plugins/ConcatPlugin.py b/plugins/ConcatPlugin.py new file mode 100644 index 0000000..43d7db2 --- /dev/null +++ b/plugins/ConcatPlugin.py @@ -0,0 +1,50 @@ +from typing import Any +import time + +from tinydb import Query +from app.Item import Item +from app.PluginInterface import Params, PluginInterface +from app.DatabaseManager import database_manager +from app.utils import ItemDict, dict_to_item, item_to_dict + + +class Plugin(PluginInterface): + def __init__(self, id: str, params: Params) -> None: + super().__init__(id, params) + print(f"[ConcatPlugin#{self.id}] initialized") + + def item_sort_key(self, item: Item) -> time.struct_time: + if item.pub_date is None: + return time.localtime() + + return item.pub_date + + 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") + + plugin_state_q = Query().plugin_id == self.id + _d: Any = database_manager.plugin_states.get(plugin_state_q) # type: ignore + data: dict[str, Any] = ( + _d if _d is not None else {"plugin_id": self.id, "state": {}} + ) + state = data["state"] + + items_as_dicts = list(map(item_to_dict, items)) + state[source_id] = items_as_dicts + + database_manager.plugin_states.upsert( # type: ignore + {"plugin_id": self.id, "state": state}, plugin_state_q + ) + + aggregated_item_dicts: list[ItemDict] = [] + for data_source_id in state: + item_dicts = state[data_source_id] + aggregated_item_dicts += item_dicts + + ret_items = list(map(dict_to_item, aggregated_item_dicts)) + ret = sorted(ret_items, key=self.item_sort_key) + + print(f"[ConcatPlugin#{self.id}] process returning items, n={len(ret)}") + return ret |
