from datetime import datetime from typing import Any, Dict, List from delphi_es_utils.entities.base_struct import BaseStruct class TaskResponse(BaseStruct): """Models a Task API response JSON ``response`` property""" batches = 0 created = 0 deleted = 0 failures: List[Any] = [] noops = 0 requests_per_second = 0 retries: Dict[str, Any] = {} throttled_millis = 0 throttled_until_millis = 0 timed_out = False took = 0 total = 0 updated = 0 version_conflicts = 0 class TaskStatus(BaseStruct): """Models a :class:`Task` object's ``status`` property""" batches = 0 created = 0 deleted = 0 noops = 0 retries: Dict[str, Any] = {} throttled_millis = 0 total = 0 updated = 0 version_conflicts = 0 class Task(BaseStruct): action = '' cancellable = False description = '' node = '' running_time_in_nanos = 0 start_time_in_millis = 0 def __init__(self, data: dict): """Models the ``task`` property from the Elasticsearch `Task API`_ response .. _Task API: https://www.elastic.co/guide/en/elasticsearch/reference/current/tasks.html Args: data: dictionary of a serialized task object from Elasticsearch """ # store our raw data object for reference self.data = data or {} # set allowed attributes for key, val in self.data.items(): self[key] = val # set these special cases self.status: TaskStatus = TaskStatus.from_dict(self.data.get('status', {})) # set id and type to not conflict with python builtin names self.id_: int = self.data.get('id', 0) self.type_ = self.data.get('type', '') def _ignore_fields(self): return ['data', 'status'] class TaskResult(BaseStruct): def __init__(self, data: dict): """Models a Task API response JSON (root) with a few added helper properties. Args: data: serialized JSON of the ``GET _tasks/node:task_id`` response """ self.data = data task_data = self.data.get('task', {}) response_data = self.data.get('response', {}) self.completed = self.data.get('completed') self.task = Task(task_data) self.response = TaskResponse.from_dict(response_data) @property def processed(self) -> int: """ Returns: Number of documents processed in task """ status = self.task.status return (status.updated or 0) + \ (status.created or 0) + \ (status.deleted or 0) @property def progress(self) -> int: """ Returns: Percentage of task completion as integer (0-100) """ if self.completed: return 100 return int(self.processed // (self.total or 0) * 100) @property def started(self) -> datetime: """ Returns: UTC datetime task started""" start_seconds = self.task.start_time_in_millis / 1000 return datetime.utcfromtimestamp(start_seconds) @property def time_remaining(self) -> float: """ Returns: Estimated time in seconds remaining for task """ running_seconds = self.task.running_time_in_nanos / 10**9 remaining_docs = (self.total or 0) - self.processed return (running_seconds / (self.processed or 1)) * remaining_docs @property def total(self) -> int: """ Returns: Total number of documents to be processed """ return self.task.status.total or 0 def _ignore_fields(self): return ['data', 'response']