From 31d28baace2ad7dc1c5bc5eb0ca495e20a70fd50 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 24 Jun 2026 14:35:51 +0800 Subject: [PATCH] feat: add in-memory JobStore --- spark_executor/core/job_store.py | 28 +++++++++++++++++++++++++ tests/unit/test_job_store.py | 35 ++++++++++++++++++++++++++++++++ 2 files changed, 63 insertions(+) create mode 100644 spark_executor/core/job_store.py create mode 100644 tests/unit/test_job_store.py diff --git a/spark_executor/core/job_store.py b/spark_executor/core/job_store.py new file mode 100644 index 0000000..6a99c58 --- /dev/null +++ b/spark_executor/core/job_store.py @@ -0,0 +1,28 @@ +# coding=utf-8 +""" +@Time :2026/6/24 +@Author :tao.chen +""" +from threading import Lock + +from spark_executor.models import Job + + +class JobStore: + """In-memory job registry. Stage-3 will swap this for SQLite.""" + + def __init__(self) -> None: + self._lock = Lock() + self._jobs: dict[str, Job] = {} + + def put(self, job: Job) -> None: + with self._lock: + self._jobs[job.job_id] = job + + def get(self, job_id: str) -> Job | None: + with self._lock: + return self._jobs.get(job_id) + + def list(self) -> list[Job]: + with self._lock: + return list(self._jobs.values()) diff --git a/tests/unit/test_job_store.py b/tests/unit/test_job_store.py new file mode 100644 index 0000000..1c15d9e --- /dev/null +++ b/tests/unit/test_job_store.py @@ -0,0 +1,35 @@ +# coding=utf-8 +from datetime import datetime + +from spark_executor.core.job_store import JobStore +from spark_executor.models import Job + + +def _job(jid: str) -> Job: + return Job( + job_id=jid, + application_id=f"application_{jid}", + script_path="/tmp/j.py", + queue="default", + submit_time=datetime(2026, 6, 24), + connection="prod", + ) + + +def test_put_then_get_roundtrip(): + store = JobStore() + store.put(_job("a")) + assert store.get("a") is not None + assert store.get("a").application_id == "application_a" + + +def test_get_missing_returns_none(): + store = JobStore() + assert store.get("nope") is None + + +def test_list_returns_all_jobs(): + store = JobStore() + store.put(_job("a")) + store.put(_job("b")) + assert {j.job_id for j in store.list()} == {"a", "b"}