aboutsummaryrefslogtreecommitdiffstats
path: root/femtoqueue.py
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2025-05-14 00:38:02 +0300
committerJan Tuomi <jan@jantuomi.fi>2025-05-14 01:17:16 +0300
commitcb57365b211c7178e10f8b722abca3afdb65e206 (patch)
treed9ee7c8c55e4f95e5d09acc3938cbe7d156b1333 /femtoqueue.py
parentf543c642b9ce0b0785639e1fb8fe9dbf9d41163a (diff)
Add initial impl
Diffstat (limited to 'femtoqueue.py')
-rw-r--r--femtoqueue.py134
1 files changed, 134 insertions, 0 deletions
diff --git a/femtoqueue.py b/femtoqueue.py
new file mode 100644
index 0000000..4003336
--- /dev/null
+++ b/femtoqueue.py
@@ -0,0 +1,134 @@
+from os import makedirs, path, listdir, rename
+from dataclasses import dataclass
+from uuid import uuid4
+from time import time
+
+@dataclass
+class FemtoTask:
+ id: str
+ data: bytes
+
+class FemtoQueue:
+ def __init__(
+ self,
+ data_dir: str,
+ node_name: str,
+ timeout_stale_ms: int = 30_000,
+ ):
+ assert node_name != "pending" \
+ and node_name != "done" \
+ and node_name != "failed"
+ self.node_name = node_name
+
+ assert timeout_stale_ms > 0
+ self.timeout_stale_ms = timeout_stale_ms
+ self.latest_stale_check_ts: float | None = None
+
+ self.data_dir = data_dir
+ self.dir_pending = path.join(data_dir, "pending")
+ self.dir_in_progress = path.join(data_dir, node_name)
+ self.dir_done = path.join(data_dir, "done")
+ self.dir_failed = path.join(data_dir, "failed")
+
+ makedirs(self.data_dir, exist_ok=True)
+ makedirs(self.dir_pending, exist_ok=True)
+ makedirs(self.dir_in_progress, exist_ok=True)
+ makedirs(self.dir_done, exist_ok=True)
+ makedirs(self.dir_failed, exist_ok=True)
+
+ def push(self, data: bytes) -> str:
+ id = uuid4().hex
+ pending_path = path.join(self.dir_pending, id)
+
+ with open(pending_path, "wb") as f:
+ f.write(data)
+
+ return id
+
+ def _release_stale_tasks(self):
+ now = time()
+
+ # Only run this every `timeout_stale_ms` milliseconds because iterating
+ # through all tasks is slow
+ timeout_sec = self.timeout_stale_ms / 1000.0
+ if self.latest_stale_check_ts is not None and now - self.latest_stale_check_ts < timeout_sec:
+ return
+
+ for dir_name in listdir(self.data_dir):
+ full_dir_path = path.join(self.data_dir, dir_name)
+
+ # Skip non-directories and reserved names
+ if not path.isdir(full_dir_path):
+ continue
+ if dir_name in ("pending", "done", "failed", self.node_name):
+ continue
+
+ # 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)
+ try:
+ modified_time = path.getmtime(task_path)
+ except FileNotFoundError:
+ continue # Task may have been moved concurrently
+
+ if now - modified_time < timeout_sec:
+ continue
+
+ try:
+ pending_path = path.join(self.dir_pending, task_file)
+ rename(task_path, pending_path)
+ except FileNotFoundError:
+ continue # Task may have been moved by another node
+
+ def _pop_task_path(self) -> str | None:
+ # First check assigned tasks in progress
+ tasks = [path.join(self.dir_in_progress, p) for p in listdir(self.dir_in_progress)]
+ #tasks.sort(key=lambda x: path.getmtime(x)) # oldest first
+ if len(tasks) > 0:
+ return tasks[0]
+
+ # Then check pending tasks
+ tasks = [path.join(self.dir_pending, p) for p in listdir(self.dir_pending)]
+ #tasks.sort(key=lambda x: path.getmtime(x)) # oldest first
+ if len(tasks) > 0:
+ return tasks[0]
+
+ return None
+
+ def pop(self) -> FemtoTask | None:
+ #self._release_stale_tasks()
+
+ while True:
+ task = self._pop_task_path()
+ if task is None: return None
+
+ id = path.basename(task)
+ in_progress_path = path.join(self.dir_in_progress, id)
+
+ try:
+ rename(task, in_progress_path)
+ except FileNotFoundError:
+ # If another node grabbed the task, just get another one
+ continue
+
+ with open(in_progress_path, "rb") as f:
+ content = f.read()
+ return FemtoTask(id = id, data = content)
+
+ def done(self, task: FemtoTask):
+ in_progress_path = path.join(self.dir_in_progress, task.id)
+ done_path = path.join(self.dir_done, task.id)
+
+ try:
+ rename(in_progress_path, done_path)
+ except FileNotFoundError as e:
+ raise Exception(f"Tried to complete a task that is not in progress, id={task.id}") from e
+
+ def fail(self, task: FemtoTask):
+ in_progress_path = path.join(self.dir_in_progress, task.id)
+ failed_path = path.join(self.dir_failed, task.id)
+
+ try:
+ rename(in_progress_path, failed_path)
+ except FileNotFoundError as e:
+ raise Exception(f"Tried to fail a task that is not in progress, id={task.id}") from e