Skip to content

locks

tit.jobs.locks

Advisory, filesystem-based locks (TODO.md §2.4).

Keys are held by the runner process itself (with tit.jobs.locks.hold(...), a few lines per __main__) and predicted by the scheduler by reading the same lock directory with :func:holders — so a bare simnibs_python -m tit.sim cfg.json from a second shell and a server-submitted job share one policy and one lock directory. Direct Python calls (run_simulation() in a notebook) remain advisory: nothing stops them, they just don't register a hold.

Storage: one directory per held (resource, mode, holder) triple, jobs/.locks/<sha1(discriminator)[:16]>/lock.json, created with :func:os.mkdir (atomic on bind mounts — no : or arbitrary user text in path components; NTFS-hostile characters never appear because the directory name is a hash). Reader/writer semantics: a "write" request conflicts with any existing holder of the same resource (read or write) other than itself; a "read" request conflicts only with an existing "write" holder of the same resource — so concurrent readers each get their own directory (discriminated by job id) and never collide.

The full lock-key table (which kind holds which resource, in which mode) lives in :func:keys_for, transcribed from TODO.md §2.4.

LockConflictError

Bases: RuntimeError

Raised by :func:hold on a write-against-write conflict (always), and on any other conflict in strict mode (TIT_LOCKS=strict).

LockRequest dataclass

LockRequest(resource: str, mode: str = 'write')

One lock a job needs. resource excludes the mode (e.g. "subject:001:m2m").

key property

key: str

Human-readable key for display (LockConflict.key), e.g. subject:001:m2m:write.

holders

holders(project_dir: str, *, reconcile_stale: bool = True) -> list[dict[str, Any]]

Every currently-held lock, as descriptor dicts ({key, resource, mode, job_id, pid, create_time, ts}). Stale entries (process gone) are removed as a side effect when reconcile_stale is true (the default; boot-time reconciliation passes it explicitly, but every read benefits since a crashed holder never cleans up after itself).

Source code in tit/jobs/locks.py
def holders(project_dir: str, *, reconcile_stale: bool = True) -> list[dict[str, Any]]:
    """Every currently-held lock, as descriptor dicts (``{key, resource, mode, job_id, pid,
    create_time, ts}``). Stale entries (process gone) are removed as a side effect when
    *reconcile_stale* is true (the default; boot-time reconciliation passes it explicitly, but
    every read benefits since a crashed holder never cleans up after itself).
    """
    base = locks_dir(project_dir)
    found: list[dict[str, Any]] = []
    try:
        entries = os.listdir(base)
    except OSError:
        return found
    for name in entries:
        try:
            directory = _storage_path(project_dir, os.path.join(base, name))
            path = _storage_path(project_dir, os.path.join(directory, DESCRIPTOR_FILE))
            with open(path, encoding="utf-8") as fh:
                data = json.load(fh)
        except (OSError, json.JSONDecodeError):
            continue
        pid = data.get("pid")
        create_time = data.get("create_time")
        if reconcile_stale and pid is not None and create_time is not None:
            if not _is_alive(int(pid), float(create_time)):
                _remove_dir(project_dir, directory)
                continue
        found.append(data)
    return found

parse_key

parse_key(key: str) -> LockRequest

Inverse of :attr:LockRequest.key ("<resource>:<mode>"); unknown/missing mode suffix defaults to "write" (the conservative choice — treat it as exclusive).

Source code in tit/jobs/locks.py
def parse_key(key: str) -> LockRequest:
    """Inverse of :attr:`LockRequest.key` (``"<resource>:<mode>"``); unknown/missing mode
    suffix defaults to ``"write"`` (the conservative choice — treat it as exclusive)."""
    if key.endswith(":read"):
        return LockRequest(resource=key[: -len(":read")], mode="read")
    if key.endswith(":write"):
        return LockRequest(resource=key[: -len(":write")], mode="write")
    return LockRequest(resource=key, mode="write")

release_job

release_job(project_dir: str, job_id: str) -> int

Forcibly remove every lock directory held by job_id, regardless of liveness.

Used by :meth:tit.jobs.manager.JobManager.force — a job being forced to a terminal state isn't waited on, so its locks must be dropped immediately rather than on the next lazy :func:holders reconciliation. Returns the number of directories removed.

Source code in tit/jobs/locks.py
def release_job(project_dir: str, job_id: str) -> int:
    """Forcibly remove every lock directory held by *job_id*, regardless of liveness.

    Used by :meth:`tit.jobs.manager.JobManager.force` — a job being forced to a terminal state
    isn't waited on, so its locks must be dropped immediately rather than on the next lazy
    :func:`holders` reconciliation. Returns the number of directories removed.
    """
    base = locks_dir(project_dir)
    removed = 0
    try:
        entries = os.listdir(base)
    except OSError:
        return 0
    for name in entries:
        try:
            path = _storage_path(project_dir, os.path.join(base, name))
            descriptor = _storage_path(project_dir, os.path.join(path, DESCRIPTOR_FILE))
            with open(descriptor, encoding="utf-8") as fh:
                data = json.load(fh)
        except (OSError, json.JSONDecodeError):
            continue
        if data.get("job_id") == job_id:
            _remove_dir(project_dir, path)
            removed += 1
    return removed

reconcile

reconcile(project_dir: str) -> int

Boot-time sweep: drop lock directories whose holder process is gone.

Returns the number of stale directories removed.

Source code in tit/jobs/locks.py
def reconcile(project_dir: str) -> int:
    """Boot-time sweep: drop lock directories whose holder process is gone.

    Returns the number of stale directories removed.
    """
    before = len(holders(project_dir, reconcile_stale=False))
    after = len(holders(project_dir, reconcile_stale=True))
    return max(before - after, 0)

match_conflicts

match_conflicts(requests: list[LockRequest], current_holders: list[dict[str, Any]], *, self_job_id: str | None = None) -> list[dict[str, Any]]

Pure function: which of current_holders block requests (reader/writer rule, §2.4).

Takes an already-fetched holder list (rather than reading the lock directory itself) so the scheduler can reuse one on-disk scan across every queued job in a tick, and so this rule is unit-testable without any filesystem I/O.

Source code in tit/jobs/locks.py
def match_conflicts(
    requests: list[LockRequest],
    current_holders: list[dict[str, Any]],
    *,
    self_job_id: str | None = None,
) -> list[dict[str, Any]]:
    """Pure function: which of *current_holders* block *requests* (reader/writer rule, §2.4).

    Takes an already-fetched holder list (rather than reading the lock directory itself) so the
    scheduler can reuse one on-disk scan across every queued job in a tick, and so this rule is
    unit-testable without any filesystem I/O.
    """
    blocking: list[dict[str, Any]] = []
    seen_job_ids: set[str] = set()
    for request in requests:
        for holder in current_holders:
            if holder.get("job_id") == self_job_id:
                continue
            if holder.get("resource") != request.resource:
                continue
            holder_mode = holder.get("mode", "write")
            if request.mode == "write" or holder_mode == "write":
                if holder["job_id"] not in seen_job_ids:
                    blocking.append(holder)
                    seen_job_ids.add(holder["job_id"])
    return blocking

conflicts

conflicts(project_dir: str, requests: list[LockRequest], *, self_job_id: str | None = None) -> list[dict[str, Any]]

Holders that block requests (reader/writer rule, §2.4), excluding self_job_id.

Source code in tit/jobs/locks.py
def conflicts(
    project_dir: str,
    requests: list[LockRequest],
    *,
    self_job_id: str | None = None,
) -> list[dict[str, Any]]:
    """Holders that block *requests* (reader/writer rule, §2.4), excluding *self_job_id*."""
    return match_conflicts(requests, holders(project_dir), self_job_id=self_job_id)

hold

hold(project_dir: str, job_id: str, requests: list[LockRequest], *, pid: int | None = None, strict: bool | None = None) -> Iterator[None]

Acquire requests for the duration of the context (called by the runner process).

Policy has two tiers. A write against a live write always raises :class:LockConflictError, whatever the mode: two processes writing one subject's head model is the corruption this whole module exists to prevent, and "warn and proceed" is not a policy for it. Every other conflict (a reader against a writer, a writer against readers) keeps the historical warn-and-continue default, and strict=True / TIT_LOCKS=strict raises on those too.

A write lock's directory is named after the resource alone, so two jobs wanting the same exclusive resource want the same directory. It is never taken from a live owner: the descriptor of a different, still-running job is left exactly as it was, so holders keeps naming the real owner and that owner's own release does not free someone else's lock.

Source code in tit/jobs/locks.py
@contextlib.contextmanager
def hold(
    project_dir: str,
    job_id: str,
    requests: list[LockRequest],
    *,
    pid: int | None = None,
    strict: bool | None = None,
) -> Iterator[None]:
    """Acquire *requests* for the duration of the context (called by the runner process).

    Policy has two tiers. A **write against a live write** always raises
    :class:`LockConflictError`, whatever the mode: two processes writing one subject's head
    model is the corruption this whole module exists to prevent, and "warn and proceed" is not
    a policy for it. Every other conflict (a reader against a writer, a writer against readers)
    keeps the historical warn-and-continue default, and ``strict=True`` / ``TIT_LOCKS=strict``
    raises on those too.

    A write lock's directory is named after the resource alone, so two jobs wanting the same
    exclusive resource want the same directory. It is never taken from a live owner: the
    descriptor of a *different, still-running* job is left exactly as it was, so `holders`
    keeps naming the real owner and that owner's own release does not free someone else's lock.
    """
    if strict is None:
        strict = os.environ.get(ENV_STRICT, "").strip().lower() == "strict"
    pid = pid if pid is not None else os.getpid()
    try:
        create_time = psutil.Process(pid).create_time()
    except (psutil.NoSuchProcess, psutil.AccessDenied):
        create_time = time.time()

    held_dirs: list[str] = []
    try:
        blocking = conflicts(project_dir, requests, self_job_id=job_id)
        if blocking:
            msg = (
                f"job {job_id}: lock conflict on "
                f"{sorted({b.get('key', b.get('resource', '?')) for b in blocking})} "
                f"held by {sorted({b['job_id'] for b in blocking})}"
            )
            write_write = _write_against_write(requests, blocking)
            if strict or write_write:
                raise LockConflictError(msg)
            logger.warning(msg)
        os.makedirs(locks_dir(project_dir), exist_ok=True)
        for request in requests:
            path = _dir_for(project_dir, request, job_id)
            try:
                os.mkdir(path)
            except FileExistsError:
                if _owned_by_another_live_job(project_dir, path, job_id):
                    # Someone else's exclusive lock. Do not overwrite the descriptor: that
                    # made `holders` report the wrong owner, and made this job's release
                    # remove the other job's lock directory.
                    raise LockConflictError(
                        f"job {job_id}: {request.key} is held by another live job"
                    ) from None
                # Same (resource, mode, job) re-entered (e.g. a rerun reusing the id), or a
                # dead holder's leftovers — fine to claim.
            descriptor = {
                "key": request.key,
                "resource": request.resource,
                "mode": request.mode,
                "job_id": job_id,
                "pid": pid,
                "create_time": create_time,
                "ts": time.time(),
            }
            descriptor_path = _storage_path(
                project_dir, os.path.join(path, DESCRIPTOR_FILE)
            )
            with open(descriptor_path, "w", encoding="utf-8") as fh:
                json.dump(descriptor, fh)
            held_dirs.append(path)
        yield
    finally:
        for path in held_dirs:
            _remove_dir(project_dir, path)

keys_for

keys_for(kind: str, subject_ids: list[str], config: dict[str, Any] | None = None) -> list[LockRequest]

Lock requests one job of kind needs, given its (still-serialized) config.

Best-effort against configs whose exact dataclass shape is owned by another lane (marked below): falls back to a coarse subject-scoped lock rather than requesting nothing, so two jobs of an unrecognised shape still serialize instead of silently racing.

Source code in tit/jobs/locks.py
def keys_for(
    kind: str, subject_ids: list[str], config: dict[str, Any] | None = None
) -> list[LockRequest]:
    """Lock requests one job of *kind* needs, given its (still-serialized) *config*.

    Best-effort against configs whose exact dataclass shape is owned by another lane (marked
    below): falls back to a coarse subject-scoped lock rather than requesting nothing, so two
    jobs of an unrecognised shape still serialize instead of silently racing.
    """
    config = config or {}
    subject_ids = subject_ids or (
        [config["subject_id"]] if config.get("subject_id") else []
    )
    requests: list[LockRequest] = []

    for sid in subject_ids:
        if kind == "pre":
            requests.extend(_pre_requests(sid, config))
        elif kind == "sim":
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="read"))
            requests.append(LockRequest(f"subject:{sid}:m2m:t1mni"))
            for montage in config.get("montages", []) or []:
                name = montage.get("name") if isinstance(montage, dict) else None
                if name:
                    requests.append(LockRequest(f"subject:{sid}:sim:{name}"))
        elif kind in ("flex", "flex_adaptive", "flex_pareto"):
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="read"))
        elif kind in ("ex", "mex"):
            run_name = config.get("run_name") or "default"
            requests.append(LockRequest(f"subject:{sid}:{kind}:{run_name}"))
            requests.append(LockRequest(f"subject:{sid}:leadfields", mode="read"))
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="read"))
            requests.append(LockRequest(f"subject:{sid}:rois", mode="read"))
        elif kind == "leadfield":
            requests.append(LockRequest(f"subject:{sid}:leadfields", mode="write"))
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="write"))
        elif kind == "analyzer":
            output_dir = (
                config.get("output_dir")
                or config.get("analysis_output_dir")
                or "default"
            )
            requests.append(LockRequest(f"subject:{sid}:analysis:{output_dir}"))
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="read"))
        elif kind == "source":
            requests.append(LockRequest(f"subject:{sid}:forward"))
            requests.append(LockRequest(f"subject:{sid}:m2m", mode="read"))

    if kind == "stats":
        # GroupComparisonConfig / CorrelationConfig (tit/stats/config.py) -- neither carries an
        # "analysis_type"/"output_dir"/"name" field. tit/stats/__main__.py distinguishes them by
        # a top-level "mode" key on the request dict itself (default "group_comparison", else
        # "correlation" -- config_io's own "_type" discriminator is not used here, that
        # mechanism is for union-typed *fields*, not this top-level kind switch), and the
        # human-chosen run name is "analysis_name" on both dataclasses.
        analysis_type = config.get("mode", "group_comparison")
        name = config.get("analysis_name") or "default"
        requests.append(LockRequest(f"project:stats:{analysis_type}/{name}"))

    return requests