feat: add in-memory JobStore
This commit is contained in:
@@ -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())
|
||||
@@ -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"}
|
||||
Reference in New Issue
Block a user