From 9df3a2a855867478ceb7a0f3a5005fe70998d261 Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Wed, 14 May 2025 10:50:34 +0300 Subject: Ensure ordering at microsecond resolution --- femtoqueue.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) (limited to 'femtoqueue.py') diff --git a/femtoqueue.py b/femtoqueue.py index e7ffcf8..500089c 100644 --- a/femtoqueue.py +++ b/femtoqueue.py @@ -1,7 +1,7 @@ from os import makedirs, path, listdir, rename from dataclasses import dataclass from uuid import uuid4 -from time import time +import time @dataclass class FemtoTask: @@ -37,7 +37,8 @@ class FemtoQueue: makedirs(self.dir_failed, exist_ok=True) def push(self, data: bytes) -> str: - id = f"{uuid4().hex[-12:]}_{str(int(time()))}" + now_us = 1_000_000 * time.time() + id = f"{str(int(now_us))}_{uuid4().hex[-12:]}" pending_path = path.join(self.dir_pending, id) with open(pending_path, "wb") as f: @@ -46,7 +47,7 @@ class FemtoQueue: return id def _release_stale_tasks(self): - now = time() + now = time.time() # Only run this every `timeout_stale_ms` milliseconds because iterating # through all tasks is slow @@ -67,7 +68,8 @@ class FemtoQueue: # Check tasks in this node's in-progress directory for task_file in listdir(full_dir_path): task_path = path.join(full_dir_path, task_file) - modified_time = int(task_file[13:]) + modified_time_us = int(task_file.split("_")[0]) + modified_time = modified_time_us / 1_000_000.0 if now - modified_time < timeout_sec: continue @@ -82,13 +84,13 @@ class FemtoQueue: # First check assigned tasks in progress tasks = listdir(self.dir_in_progress) if len(tasks) > 0: - id = min(tasks, key=lambda x: x[13:]) + id = min(tasks) return path.join(self.dir_in_progress, id) # Then check pending tasks tasks = listdir(self.dir_pending) if len(tasks) > 0: - id = min(tasks, key=lambda x: x[13:]) + id = min(tasks) return path.join(self.dir_pending, id) return None -- cgit v1.3