95 lines
2.9 KiB
Python
95 lines
2.9 KiB
Python
import multiprocessing
|
|
import os
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
from tht.jobs.locking import JobAlreadyRunningError, WorkspaceJobLock
|
|
|
|
|
|
def _hold_lock(root: str, ready, release):
|
|
with WorkspaceJobLock(Path(root), "demo", "evidence"):
|
|
ready.set()
|
|
release.wait(10)
|
|
|
|
|
|
def _crash_with_lock(root: str, ready):
|
|
lock = WorkspaceJobLock(Path(root), "demo", "evidence")
|
|
lock.acquire()
|
|
ready.set()
|
|
raise SystemExit(7)
|
|
|
|
|
|
def test_same_workspace_and_job_are_exclusive_across_processes(tmp_path):
|
|
context = multiprocessing.get_context("spawn")
|
|
ready = context.Event()
|
|
release = context.Event()
|
|
process = context.Process(target=_hold_lock, args=(str(tmp_path), ready, release))
|
|
process.start()
|
|
assert ready.wait(10)
|
|
try:
|
|
with pytest.raises(JobAlreadyRunningError):
|
|
WorkspaceJobLock(tmp_path, "demo", "evidence").acquire()
|
|
finally:
|
|
release.set()
|
|
process.join(10)
|
|
assert process.exitcode == 0
|
|
|
|
|
|
def test_evidence_and_dwh_jobs_have_distinct_locks(tmp_path):
|
|
with (
|
|
WorkspaceJobLock(tmp_path, "demo", "evidence"),
|
|
WorkspaceJobLock(tmp_path, "demo", "dwh"),
|
|
):
|
|
pass
|
|
|
|
|
|
def test_lock_keys_cannot_escape_lock_directory(tmp_path):
|
|
with pytest.raises(ValueError, match="filesystem-safe"):
|
|
WorkspaceJobLock(tmp_path, "demo", "../evidence")
|
|
|
|
|
|
def test_preexisting_lock_symlink_is_rejected(tmp_path):
|
|
lock = WorkspaceJobLock(tmp_path, "demo", "evidence")
|
|
lock.path.parent.mkdir(parents=True)
|
|
target = tmp_path / "target"
|
|
target.write_text("do not modify")
|
|
lock.path.symlink_to(target)
|
|
with pytest.raises(OSError):
|
|
lock.acquire()
|
|
assert target.read_text() == "do not modify"
|
|
|
|
|
|
def test_preexisting_locks_directory_symlink_is_rejected(tmp_path):
|
|
jobs = tmp_path / ".tht-jobs"
|
|
jobs.mkdir()
|
|
outside = tmp_path / "outside"
|
|
outside.mkdir()
|
|
(jobs / ".locks").symlink_to(outside, target_is_directory=True)
|
|
with pytest.raises(OSError):
|
|
WorkspaceJobLock(tmp_path, "demo", "evidence").acquire()
|
|
assert list(outside.iterdir()) == []
|
|
|
|
|
|
def test_lock_file_is_owner_only_regular_single_link(tmp_path):
|
|
with WorkspaceJobLock(tmp_path, "demo", "evidence") as lock:
|
|
stat = os.stat(lock.path, follow_symlinks=False)
|
|
assert stat.st_uid == os.getuid()
|
|
assert stat.st_nlink == 1
|
|
assert stat.st_mode & 0o777 == 0o600
|
|
|
|
|
|
def test_lock_is_recoverable_after_process_crash_without_stale_deletion(tmp_path):
|
|
context = multiprocessing.get_context("spawn")
|
|
ready = context.Event()
|
|
process = context.Process(target=_crash_with_lock, args=(str(tmp_path), ready))
|
|
process.start()
|
|
assert ready.wait(10)
|
|
process.join(10)
|
|
assert process.exitcode == 7
|
|
|
|
lock_path = WorkspaceJobLock(tmp_path, "demo", "evidence").path
|
|
assert lock_path.exists()
|
|
with WorkspaceJobLock(tmp_path, "demo", "evidence"):
|
|
assert lock_path.exists()
|