116 lines
3.7 KiB
Python
116 lines
3.7 KiB
Python
"""Crash-safe interprocess locking scoped by workspace and job type."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import fcntl
|
|
import hashlib
|
|
import os
|
|
import re
|
|
import stat
|
|
from pathlib import Path
|
|
from types import TracebackType
|
|
|
|
|
|
class JobAlreadyRunningError(RuntimeError):
|
|
pass
|
|
|
|
|
|
_JOB_KEY = re.compile(r"^[a-z][a-z0-9_-]{0,63}$")
|
|
|
|
|
|
def _lock_name(workspace_id: str, job_type: str) -> str:
|
|
if not _JOB_KEY.fullmatch(workspace_id) or not _JOB_KEY.fullmatch(job_type):
|
|
raise ValueError("lock identifiers must be lowercase filesystem-safe keys")
|
|
workspace_key = hashlib.sha256(workspace_id.encode("utf-8")).hexdigest()[:16]
|
|
return f"{workspace_key}-{job_type}.lock"
|
|
|
|
|
|
class WorkspaceJobLock:
|
|
"""Advisory kernel lock; the inode remains stable and is never deleted by PID."""
|
|
|
|
def __init__(self, workspace_root: Path, workspace_id: str, job_type: str) -> None:
|
|
self.path = workspace_root / ".tht-jobs" / ".locks" / _lock_name(
|
|
workspace_id, job_type
|
|
)
|
|
self._fd: int | None = None
|
|
|
|
def acquire(self) -> "WorkspaceJobLock":
|
|
if self._fd is not None:
|
|
raise RuntimeError("job lock is already held by this object")
|
|
root_fd = os.open(self.path.parents[2], os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW)
|
|
try:
|
|
jobs_fd = _open_owned_directory(root_fd, ".tht-jobs")
|
|
try:
|
|
locks_fd = _open_owned_directory(jobs_fd, ".locks")
|
|
try:
|
|
fd = os.open(
|
|
self.path.name,
|
|
os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC,
|
|
0o600,
|
|
dir_fd=locks_fd,
|
|
)
|
|
try:
|
|
info = os.fstat(fd)
|
|
if (
|
|
not stat.S_ISREG(info.st_mode)
|
|
or info.st_uid != os.getuid()
|
|
or info.st_nlink != 1
|
|
):
|
|
raise OSError("unsafe job lock file")
|
|
os.fchmod(fd, 0o600)
|
|
try:
|
|
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
except BlockingIOError as error:
|
|
raise JobAlreadyRunningError(
|
|
"this workspace job is already running"
|
|
) from error
|
|
except BaseException:
|
|
os.close(fd)
|
|
raise
|
|
finally:
|
|
os.close(locks_fd)
|
|
finally:
|
|
os.close(jobs_fd)
|
|
finally:
|
|
os.close(root_fd)
|
|
self._fd = fd
|
|
return self
|
|
|
|
def release(self) -> None:
|
|
if self._fd is None:
|
|
return
|
|
fd, self._fd = self._fd, None
|
|
try:
|
|
fcntl.flock(fd, fcntl.LOCK_UN)
|
|
finally:
|
|
os.close(fd)
|
|
|
|
def __enter__(self) -> "WorkspaceJobLock":
|
|
return self.acquire()
|
|
|
|
def __exit__(
|
|
self,
|
|
exc_type: type[BaseException] | None,
|
|
exc: BaseException | None,
|
|
traceback: TracebackType | None,
|
|
) -> None:
|
|
self.release()
|
|
|
|
|
|
def _open_owned_directory(parent_fd: int, name: str) -> int:
|
|
try:
|
|
os.mkdir(name, 0o700, dir_fd=parent_fd)
|
|
os.fsync(parent_fd)
|
|
except FileExistsError:
|
|
pass
|
|
fd = os.open(name, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=parent_fd)
|
|
try:
|
|
info = os.fstat(fd)
|
|
if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid():
|
|
raise OSError("unsafe job lock directory")
|
|
os.fchmod(fd, 0o700)
|
|
except BaseException:
|
|
os.close(fd)
|
|
raise
|
|
return fd
|