From 87e01eec8abadcd177725ac32c809e9f3079c728 Mon Sep 17 00:00:00 2001 From: Kristina Shishkina Date: Wed, 25 Dec 2024 16:58:29 +0300 Subject: [PATCH 1/4] refactor: log message --- pyproject.toml | 2 +- src/nats_queue/nats_worker.py | 32 +++++++++++++++++--------------- tests/test_workers.py | 2 +- 3 files changed, 19 insertions(+), 17 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index f0fa7c1..63985bf 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "nats_queue" -version = "1.1.3" +version = "1.1.4" description = "" authors = ["Kristina Shishkina "] readme = "README.md" diff --git a/src/nats_queue/nats_worker.py b/src/nats_queue/nats_worker.py index 2536dc4..34ae729 100644 --- a/src/nats_queue/nats_worker.py +++ b/src/nats_queue/nats_worker.py @@ -114,9 +114,7 @@ async def _mark_parents_failed(self, job_data: dict): parent_job = await self.kv.get(parent_id) if not parent_job: - self.logger.warning( - f"Parent job with ID {parent_id} not found in KV store." - ) + self.logger.warning(f"ParentJob with id={parent_id} not found in KV store.") return parent_job_data = json.loads(parent_job.value.decode()) @@ -132,20 +130,22 @@ async def _publish_parent_job(self, parent_job_data): subject, job_bytes, headers={"Nats-Msg-Id": parent_job_data["id"]} ) self.logger.info( - f"Parent Job id={parent_job_data['id']} " - f"subject={subject} added successfully" + f"ParentJob: name={parent_job_data['name']} " + f"id={parent_job_data['id']} " + f"added to topic={subject} successfully" ) async def _process_task(self, job: Msg): try: self.processing_now += 1 job_data = json.loads(job.data.decode()) - if job_data["meta"].get("faild"): + if job_data["meta"].get("failed"): await job.term() self.logger.warning( - f"Job: {job_data['name']} id={job_data['id']} failed because " - f"child job did not complete successfully " + f"Job: name={job_data['name']} id={job_data['id']} failed " + f"because child job failed to process" ) + return job_start_time = datetime.fromisoformat(job_data["meta"]["start_time"]) if job_start_time > datetime.now(): @@ -154,7 +154,7 @@ async def _process_task(self, job: Msg): await job.nak(delay=delay) self.logger.debug( ( - f"Job:{job_data['name']} id={job_data['id']} is " + f"Job: name={job_data['name']} id={job_data['id']} is " f"scheduled later " f"Requeueing in {delay} seconds" ) @@ -164,7 +164,7 @@ async def _process_task(self, job: Msg): if job_data.get("meta").get("retry_count") > self.max_retries: await job.term() self.logger.warning( - f"Job: {job_data['name']} id={job_data['id']} " + f"Job: name={job_data['name']} id={job_data['id']} " f"failed max retries exceeded" ) @@ -173,7 +173,7 @@ async def _process_task(self, job: Msg): self.logger.info( ( - f"Job: {job_data['name']} id={job_data['id']} is started " + f"Job: name={job_data['name']} id={job_data['id']} is started " f"with data={job_data['data']} in queue={job_data['queue_name']}" ) ) @@ -200,7 +200,7 @@ async def _process_task(self, job: Msg): except Exception as e: if isinstance(e, asyncio.TimeoutError): self.logger.error( - f"Job: {job_data['name']} id={job_data['id']} " + f"Job: name={job_data['name']} id={job_data['id']} " f"TimeoutError start retry" ) else: @@ -229,7 +229,7 @@ async def fetch_messages( self.logger.debug( ( f"Consumer: name={(await sub.consumer_info()).name} " - f"fetched {len(msgs)} messages" + f"fetched {len(msgs)} messages from queue={self.name}" "" ) ) @@ -238,7 +238,8 @@ async def fetch_messages( self.logger.debug( ( f"Consumer: name={(await sub.consumer_info()).name} " - f"failed to fetch messages: TimeoutError" + f"failed to fetch messages from from queue={self.name}: " + f"TimeoutError" ) ) return [] @@ -246,7 +247,8 @@ async def fetch_messages( self.logger.error( ( f"Consumer: name={(await sub.consumer_info()).name} " - f"error while fetching messages: {e}" + f"error while fetching messages from queue=" + f"{self.name}: {e}" ) ) raise diff --git a/tests/test_workers.py b/tests/test_workers.py index ef01063..fa7dfb6 100644 --- a/tests/test_workers.py +++ b/tests/test_workers.py @@ -686,7 +686,7 @@ async def test_publish_parent_job(get_client): @pytest.mark.asyncio -async def test__mark_parents_failed(get_client): +async def test_mark_parents_failed(get_client): client = get_client queue = Queue(client, name="my_queue") await queue.setup() From 672a2b19c89039194f8c3da9b5f345d53ac70cee Mon Sep 17 00:00:00 2001 From: Kristina Shishkina Date: Thu, 9 Jan 2025 17:47:58 +0300 Subject: [PATCH 2/4] fix: add priorityQuota const and methods _reset_quotes_counter, _handle_quota and update loop with quota --- src/nats_queue/nats_worker.py | 55 +++++++++++++++++++++++++++++++---- 1 file changed, 49 insertions(+), 6 deletions(-) diff --git a/src/nats_queue/nats_worker.py b/src/nats_queue/nats_worker.py index 34ae729..887ec1a 100644 --- a/src/nats_queue/nats_worker.py +++ b/src/nats_queue/nats_worker.py @@ -31,7 +31,11 @@ def __init__( priorities: int = 1, limiter: Dict[str, int] = None, logger: Logger = logger, + priorityQuota: Optional[Dict[int, Dict[str, int]]] = None, ): + if priorityQuota and set(priorityQuota.keys()) != set([priority for priority in range(1, priorities+1)]): + raise ValueError("The priority quota must contain settings for each priority") + self.client = client self.name = name self.processor = processor @@ -57,6 +61,14 @@ def __init__( self.loop_task: Optional[asyncio.Task] = None self.logger: Logger = logger self.kv: Optional[KeyValue] = None + self.priorityQuota: Optional[Dict[int, Dict[str, int]]] = ( + { + priority: {"quota": item["quota"], "counter": 0} + for priority, item in priorityQuota.items() + } + if priorityQuota + else None + ) self.logger.info( ( @@ -89,19 +101,50 @@ async def start(self): self.running = True self.loop_task = asyncio.create_task(self.loop()) + def _reset_quotes_counter(self): + for item in self.priorityQuota.values(): + item['counter'] = 0 + + def _handle_quota(self, consumer_priority): + + current_quota = self.priorityQuota.get(consumer_priority).get("quota") + current_counter = self.priorityQuota.get(consumer_priority, {}).get("counter") + if current_quota - current_counter <= 0: + if consumer_priority < self.priorities: + logger.debug(f"Skip consumer_priority={consumer_priority} due to quota overruns") + return True + else: + logger.debug(f"Reset counters when quota is exceeded for consumer_priority={consumer_priority}") + self._reset_quotes_counter() + return False + return None + async def loop(self): while self.running: - for consumer in self.consumers: + + for consumer_priority, consumer in enumerate(self.consumers, 1): + jobs = [] + + if self.priorityQuota: + is_quota_met = self._handle_quota(consumer_priority) + if is_quota_met is not None: + if is_quota_met: + continue + else: + break + else: + pass + max_jobs = self.limiter.get(self.concurrency - self.processing_now) if max_jobs <= 0: - continue + break + jobs = await self.fetch_messages(consumer, max_jobs) if jobs: break - else: - jobs = [] - for job in jobs: + if self.priorityQuota: + self.priorityQuota[consumer_priority]['counter'] += 1 self.limiter.inc() asyncio.create_task(self._process_task(job)) @@ -183,7 +226,7 @@ async def _process_task(self, job: Msg): await job.ack_sync() self.logger.info( - f'Job: {job_data["name"]} id={job_data["id"]} is completed' + f'Job: name={job_data["name"]} id={job_data["id"]} is completed' ) parent_id = job_data["meta"].get("parent_id") From 0ed7e2e2647687274cdf2242ac37094dccb7a2f1 Mon Sep 17 00:00:00 2001 From: Kristina Shishkina Date: Thu, 9 Jan 2025 17:48:42 +0300 Subject: [PATCH 3/4] feat: add new test test_worker_error_priorityQuota, test_handle_quota, test_reset_quotes_counter, test_priority_quota --- tests/test_workers.py | 95 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 95 insertions(+) diff --git a/tests/test_workers.py b/tests/test_workers.py index fa7dfb6..3524700 100644 --- a/tests/test_workers.py +++ b/tests/test_workers.py @@ -112,6 +112,20 @@ async def test_worker_setup_unknow_queue(get_client): with pytest.raises(Exception): await worker.setup() +@pytest.mark.asyncio +async def test_worker_error_priorityQuota(get_client): + client: Client = get_client + + priorityQuota = {1: {"quota": 2}} + with pytest.raises(Exception, match="The priority quota must contain settings for each priority"): + worker = Worker( + client, + name="my_queue_1", + processor=process_job, + priorityQuota=priorityQuota, + priorities=2, + ) + @pytest.mark.asyncio async def test_worker_connect_stop_success(get_client): @@ -791,6 +805,87 @@ async def test_process_task_with_flow_job(get_client): with pytest.raises(NoKeysError): await queue.kv.keys() +@pytest.mark.asyncio +async def test_handle_quota(get_client): + client = get_client + queue = Queue(client, name="my_queue", priorities=2) + await queue.setup() + priorityQuota = {1: {"quota": 2}, 2: {"quota": 2}} + + worker = Worker( + client, "my_queue", process_job, priorities=2, priorityQuota=priorityQuota + ) + await worker.setup() + + assert worker.priorityQuota == { + priority: {"quota": item["quota"], "counter": 0} + for priority, item in priorityQuota.items() + } + result = worker._handle_quota(1) + assert result is None + + result = worker._handle_quota(2) + assert result is None + + for priority in range(1, 3): + worker.priorityQuota[priority]['counter'] = worker.priorityQuota[priority]['quota'] + + result = worker._handle_quota(1) + assert result is True + + result = worker._handle_quota(2) + assert result is False + +@pytest.mark.asyncio +async def test_reset_quotes_counter(get_client): + client = get_client + queue = Queue(client, name="my_queue", priorities=2) + await queue.setup() + priorityQuota = {1: {"quota": 2}, 2: {"quota": 2}} + + worker = Worker( + client, "my_queue", process_job, priorities=2, priorityQuota=priorityQuota + ) + await worker.setup() + + for priority in range(1, 3): + worker.priorityQuota[priority]['counter'] = worker.priorityQuota[priority]['quota'] + + worker._reset_quotes_counter() + + assert worker.priorityQuota[1]['counter'] == 0 + assert worker.priorityQuota[2]['counter'] == 0 + +@pytest.mark.asyncio +async def test_priority_quota(get_client): + client = get_client + queue = Queue(client, name="my_queue", priorities=2) + await queue.setup() + + job_1_1 = Job("my_queue", "job_1_1") + job_1_2 = Job("my_queue", "job_1_2") + job_1_3 = Job("my_queue", "job_1_3") + + job_2_1 = Job("my_queue", "job_2_1") + job_2_2 = Job("my_queue", "job_2_2") + job_2_3 = Job("my_queue", "job_2_3") + + await queue.addJobs([job_1_1, job_1_2, job_1_3], 1) + await queue.addJobs([job_2_1, job_2_2, job_2_3], 2) + + priorityQuota = {1: {"quota": 2}, 2: {"quota": 2}} + + worker = Worker( + client, "my_queue", process_job, priorities=2, priorityQuota=priorityQuota + ) + await worker.setup() + await worker.start() + await asyncio.sleep(10) + await worker.stop() + + msg_count = (await queue.manager.stream_info(queue.name)).state.messages + assert msg_count == 6 + async def process_job(job_data: Dict): await asyncio.sleep(1) From bf1ff8ebc91c3dfe955e02c29057b5d037432939 Mon Sep 17 00:00:00 2001 From: Kristina Shishkina Date: Thu, 9 Jan 2025 18:07:09 +0300 Subject: [PATCH 4/4] fix: code style --- src/nats_queue/nats_worker.py | 23 +++++++++++++++------- tests/test_workers.py | 36 ++++++++++++++++++++++------------- 2 files changed, 39 insertions(+), 20 deletions(-) diff --git a/src/nats_queue/nats_worker.py b/src/nats_queue/nats_worker.py index 887ec1a..fe58ea1 100644 --- a/src/nats_queue/nats_worker.py +++ b/src/nats_queue/nats_worker.py @@ -33,9 +33,13 @@ def __init__( logger: Logger = logger, priorityQuota: Optional[Dict[int, Dict[str, int]]] = None, ): - if priorityQuota and set(priorityQuota.keys()) != set([priority for priority in range(1, priorities+1)]): - raise ValueError("The priority quota must contain settings for each priority") - + if priorityQuota and set(priorityQuota.keys()) != set( + [priority for priority in range(1, priorities + 1)] + ): + raise ValueError( + "The priority quota must contain settings for each priority" + ) + self.client = client self.name = name self.processor = processor @@ -103,7 +107,7 @@ async def start(self): def _reset_quotes_counter(self): for item in self.priorityQuota.values(): - item['counter'] = 0 + item["counter"] = 0 def _handle_quota(self, consumer_priority): @@ -111,10 +115,15 @@ def _handle_quota(self, consumer_priority): current_counter = self.priorityQuota.get(consumer_priority, {}).get("counter") if current_quota - current_counter <= 0: if consumer_priority < self.priorities: - logger.debug(f"Skip consumer_priority={consumer_priority} due to quota overruns") + logger.debug( + f"Skip consumer_priority={consumer_priority} due to quota overruns" + ) return True else: - logger.debug(f"Reset counters when quota is exceeded for consumer_priority={consumer_priority}") + logger.debug( + f"Reset counters when quota is exceeded " + f"for consumer_priority={consumer_priority}" + ) self._reset_quotes_counter() return False return None @@ -144,7 +153,7 @@ async def loop(self): break for job in jobs: if self.priorityQuota: - self.priorityQuota[consumer_priority]['counter'] += 1 + self.priorityQuota[consumer_priority]["counter"] += 1 self.limiter.inc() asyncio.create_task(self._process_task(job)) diff --git a/tests/test_workers.py b/tests/test_workers.py index 3524700..fbecb17 100644 --- a/tests/test_workers.py +++ b/tests/test_workers.py @@ -112,13 +112,16 @@ async def test_worker_setup_unknow_queue(get_client): with pytest.raises(Exception): await worker.setup() + @pytest.mark.asyncio async def test_worker_error_priorityQuota(get_client): client: Client = get_client priorityQuota = {1: {"quota": 2}} - with pytest.raises(Exception, match="The priority quota must contain settings for each priority"): - worker = Worker( + with pytest.raises( + Exception, match="The priority quota must contain settings for each priority" + ): + Worker( client, name="my_queue_1", processor=process_job, @@ -805,22 +808,23 @@ async def test_process_task_with_flow_job(get_client): with pytest.raises(NoKeysError): await queue.kv.keys() + @pytest.mark.asyncio async def test_handle_quota(get_client): client = get_client queue = Queue(client, name="my_queue", priorities=2) await queue.setup() priorityQuota = {1: {"quota": 2}, 2: {"quota": 2}} - + worker = Worker( client, "my_queue", process_job, priorities=2, priorityQuota=priorityQuota ) await worker.setup() assert worker.priorityQuota == { - priority: {"quota": item["quota"], "counter": 0} - for priority, item in priorityQuota.items() - } + priority: {"quota": item["quota"], "counter": 0} + for priority, item in priorityQuota.items() + } result = worker._handle_quota(1) assert result is None @@ -828,33 +832,39 @@ async def test_handle_quota(get_client): assert result is None for priority in range(1, 3): - worker.priorityQuota[priority]['counter'] = worker.priorityQuota[priority]['quota'] - + worker.priorityQuota[priority]["counter"] = worker.priorityQuota[priority][ + "quota" + ] + result = worker._handle_quota(1) assert result is True result = worker._handle_quota(2) assert result is False + @pytest.mark.asyncio async def test_reset_quotes_counter(get_client): client = get_client queue = Queue(client, name="my_queue", priorities=2) await queue.setup() priorityQuota = {1: {"quota": 2}, 2: {"quota": 2}} - + worker = Worker( client, "my_queue", process_job, priorities=2, priorityQuota=priorityQuota ) await worker.setup() for priority in range(1, 3): - worker.priorityQuota[priority]['counter'] = worker.priorityQuota[priority]['quota'] - + worker.priorityQuota[priority]["counter"] = worker.priorityQuota[priority][ + "quota" + ] + worker._reset_quotes_counter() - assert worker.priorityQuota[1]['counter'] == 0 - assert worker.priorityQuota[2]['counter'] == 0 + assert worker.priorityQuota[1]["counter"] == 0 + assert worker.priorityQuota[2]["counter"] == 0 + @pytest.mark.asyncio async def test_priority_quota(get_client):