Something went wrong. Try again.
small gleam coding and (not yet) persistent agent daemon with a detachable cli
Something went wrong. Try again.
Python
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106"""Independent contract checks for the attached benchmark's generated queue."""import importlibimport jsonfrom pathlib import Pathimport subprocessimport sysimport tempfileimport unittest
sys.path.insert(0, str(Path(sys.argv.pop(1)).resolve()))Queue = importlib.import_module("durable_queue").Queue
class Contract(unittest.TestCase): def setUp(self): self.directory = tempfile.TemporaryDirectory() self.path = str(Path(self.directory.name) / "queue.sqlite") self.q = Queue(self.path)
def tearDown(self): self.q.close() self.directory.cleanup()
def test_priority_availability_and_idempotency(self): low = self.q.enqueue({"nested": [1, None, True]}, key="once", priority=0) self.assertEqual(self.q.enqueue("replacement", key="once", priority=500), low) high = self.q.enqueue("high", priority=3) later = self.q.enqueue("later", priority=9, available_at=10) first = self.q.claim("worker", now=0) self.assertEqual(first["id"], high) self.assertTrue(self.q.ack(high, first["token"], now=1)) second = self.q.claim("worker", now=1) self.assertEqual(second["id"], low) self.assertEqual(second["payload"], {"nested": [1, None, True]}) self.assertTrue(self.q.ack(low, second["token"], now=2)) self.assertIsNone(self.q.claim("worker", now=9)) self.assertEqual(self.q.claim("worker", now=10)["id"], later) self.assertEqual(self.q.enqueue(123, key="once"), low)
def test_tokens_expire_and_exhaustion_is_terminal(self): job = self.q.enqueue(None, max_attempts=2) first = self.q.claim("one", now=10, lease_seconds=2) self.assertFalse(self.q.ack(job, first["token"], now=12)) second = self.q.claim("two", now=12, lease_seconds=2) self.assertEqual(second["id"], job) self.assertNotEqual(first["token"], second["token"]) self.assertEqual(second["attempts"], 2) self.assertFalse(self.q.fail(job, first["token"], now=12)) self.assertFalse(self.q.heartbeat(job, first["token"], now=12)) self.assertIsNone(self.q.claim("three", now=14)) self.assertEqual(self.q.get(job)["state"], "dead") self.assertEqual(self.q.stats(), {"ready": 0, "leased": 0, "done": 0, "dead": 1})
def test_retry_delay_heartbeat_and_reopen(self): job = self.q.enqueue([1, "payload"]) lease = self.q.claim("worker", now=0, lease_seconds=10) self.assertTrue(self.q.heartbeat(job, lease["token"], now=9, lease_seconds=20)) self.assertIsNone(self.q.claim("other", now=15)) self.assertTrue(self.q.fail(job, lease["token"], now=16, delay=10)) self.assertIsNone(self.q.claim("other", now=25)) self.q.close() self.q = Queue(self.path) retry = self.q.claim("worker", now=26) self.assertEqual(retry["id"], job) self.assertEqual(retry["attempts"], 2) self.assertTrue(self.q.ack(job, retry["token"], now=27)) self.assertFalse(self.q.ack(job, retry["token"], now=27)) self.assertEqual(self.q.stats()["done"], 1)
def test_expiry_reaps_all_exhausted_jobs(self): ids = [self.q.enqueue(index, max_attempts=1) for index in range(5)] for _ in ids: self.q.claim("worker", now=0, lease_seconds=1) self.assertIsNone(self.q.claim("another", now=1)) self.assertEqual(self.q.stats()["dead"], 5)
def test_insertion_order_breaks_ties(self): ids = [self.q.enqueue(value) for value in [0, False, "", [], {}]] for expected in ids: lease = self.q.claim("worker", now=0) self.assertEqual(lease["id"], expected) self.assertTrue(self.q.ack(expected, lease["token"], now=0))
def test_validation_does_not_insert_bad_jobs(self): for args in [{"max_attempts": 0}, {"available_at": float("nan")}, {"available_at": float("inf")}]: with self.assertRaises((ValueError, TypeError)): self.q.enqueue("invalid", **args) for args in [{"owner": "", "now": 0}, {"owner": "a", "now": 0, "lease_seconds": 0}]: with self.assertRaises((ValueError, TypeError)): self.q.claim(**args) self.assertEqual(sum(self.q.stats().values()), 0)
def test_claim_survives_abrupt_process_exit(self): job = self.q.enqueue("survives", max_attempts=2) source = "import json,os,sys;from durable_queue import Queue;q=Queue(sys.argv[1]);print(json.dumps(q.claim('dead-worker',now=0,lease_seconds=1)),flush=True);os._exit(0)" result = subprocess.run([sys.executable, "-c", source, self.path], cwd=sys.path[0], capture_output=True, text=True, check=True, timeout=20) token = json.loads(result.stdout)["token"] lease = self.q.claim("new-worker", now=1) self.assertEqual(lease["id"], job) self.assertFalse(self.q.ack(job, token, now=1)) self.assertTrue(self.q.ack(job, lease["token"], now=1))
if __name__ == "__main__": unittest.main(verbosity=2)