Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

[breaking]

- Cache files follow the process umask instead of being forced to `0o777`; `permissions` options were removed. For shared caches, call `exca.utils.setup_shared_folder(folder)` once on the cache root; without default ACL support, also set `umask 002`. [#324]
- `CacheDict`: deletions require a `write()` context (like writes). [#326]
- `DumpContext.shared_file`: content suffixes must start with `.`. [#326]
- `steps`: `Parallel` cannot be a `Chain` step; call it directly. [#328, #329]
Expand Down
14 changes: 0 additions & 14 deletions docs/dev/planned-work.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,20 +48,6 @@ Non-breaking behavior changes and internal cleanup can proceed without waiting.
- **Migration:** switch to `@DumpContext.register` with new-style handlers
- **When:** after confirming no external subclasses are in use

### Simplify permission handling
- Shared filesystems (NFS) need explicit chmod on created folders and files
so other users/jobs can read/write cached results.
- Old attempt on branch `set-permissions` (aborted — mixed into a large
refactor): added `PermissionSetter` utility in `utils.py`, a
`permissions: int | None = 0o777` field on `BaseInfra`/`Backend`/`CacheDict`,
and chmod calls after each mkdir/file-write.
- Next attempt should:
- Extract the permission logic cleanly (standalone PR, no other refactors)
- Also handle submitit log/job folders (currently created by submitit
itself, which doesn't set permissions — may need upstream changes in
submitit or post-creation fixup)
- Consider a umask-based approach as an alternative to post-hoc chmod

## Internal cleanup (non-breaking, can do anytime)

### Remove `_track_legacy_files` recursion
Expand Down
2 changes: 2 additions & 0 deletions docs/infra/explanation.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,8 @@ The class is initialized with these parameters:

For details on the serialization system and how to write custom handlers, see [Serialization](serialization.md).

Files follow the process umask. For a cache shared between users, call `exca.utils.setup_shared_folder(folder)` once on its root before launching jobs. If default ACLs are unavailable, the function warns and writers must run with `umask 002`.

**Example**
```python fixture:tmp_path
import numpy as np
Expand Down
5 changes: 2 additions & 3 deletions docs/internal/proposals/inflight-registry.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ Located in `exca/cachedict/inflight.py`.
class InflightRegistry:
"""Advisory SQLite registry of in-flight cache items."""

def __init__(self, folder: Path, permissions: int | None = 0o777) -> None:
def __init__(self, folder: Path) -> None:
# DB at <folder>/inflight.db
...

Expand Down Expand Up @@ -287,8 +287,7 @@ where coordination matters — which is what `docs/internal/debug/concurrent-wri
identified as the core problem.

The DB file is visible (no leading dot) for easy manual deletion if needed. File
permissions default to `0o777` (matching CacheDict's shared-access model) and are
applied after DB creation.
permissions follow the process umask.

## Same-PID Ownership

Expand Down
2 changes: 1 addition & 1 deletion docs/internal/steps/spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -247,7 +247,7 @@ required for MapInfra parity or current step semantics.
### Safety Measures (from TaskInfra/MapInfra)

- Config consistency checking (`identity.write_configs`)
- Permissions on CacheDict (`permissions=0o777`)
- Shared cache access follows the process umask
- Force/retry one-shot tracking per Backend lifetime
- Job lifecycle status — `LookupHandle.status` returns `"success"` /
`"error"` / `"running"` / `None`
Expand Down
28 changes: 6 additions & 22 deletions exca/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,10 +128,6 @@ def model_with_infra_validator_before(obj: tp.Any) -> tp.Any:

class BaseInfra(pydantic.BaseModel):
folder: Path | str | None = None
# general permission for folders and files
# use os.chmod / path.chmod compatible numbers, or None to deactivate
# eg: 0o777 for all rights to all users
permissions: int | None = 0o777
# {folder} will be replaced by the class folder
# {user} by user id and %j by job id
logs: Path | str = "{folder}/logs/{user}/%j"
Expand Down Expand Up @@ -197,6 +193,11 @@ def _exclude_from_cls_uid(self) -> list[str]:
def model_post_init(self, log__: tp.Any) -> None:
# Pydantic's private-attr hook would otherwise shadow SubmititMixin's hook.
super().model_post_init(log__)
if ".." in Path(self.version).parts:
raise ValueError(
f"version={self.version!r} must not contain '..': it is a path "
"component of the cache folder and would escape the cache root"
)

def __repr_args__(self) -> tp.Iterator[tuple[str | None, tp.Any]]:
"""Compact repr: only show fields that differ from their default value."""
Expand Down Expand Up @@ -250,15 +251,6 @@ def _check_configs(self, write: bool = True) -> None:
)
dump.check_and_write(xpfolder, write=write)
state.checked_configs = True
# Set permissions on written files
if write:
for name in ("uid", "full-uid", "config"):
fp = xpfolder / f"{name}.yaml"
if fp.exists():
try:
self._set_permissions(fp)
except (OSError, FileNotFoundError):
pass

def _factory(self) -> str:
state = _fast_state(self)
Expand Down Expand Up @@ -346,7 +338,7 @@ def uid_folder(self, create: bool = False) -> Path | None:
folder = Path(self.folder) / self.uid()
if not create:
return folder
utils.mkdir_with_permissions(folder, self.permissions, root=self.folder)
folder.mkdir(parents=True, exist_ok=True)
return folder

def iter_cached(self) -> tp.Iterable[pydantic.BaseModel]:
Expand All @@ -362,14 +354,6 @@ def iter_cached(self) -> tp.Iterable[pydantic.BaseModel]:
cfg = ConfDict.from_yaml(fp)
yield cls(**cfg)

def _set_permissions(self, path: str | Path) -> None:
if self.permissions is not None:
try:
Path(path).chmod(self.permissions)
except Exception as e:
msg = f"Failed to set permission to {self.permissions} on '{path}'\n({e})"
logger.warning(msg)

def clone_obj(self, *args: dict[str, tp.Any], **kwargs: tp.Any) -> tp.Any:
"""Create a new decorated object by applying a diff config to the underlying object"""
if args:
Expand Down
12 changes: 3 additions & 9 deletions exca/cachedict/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,10 +90,6 @@ class CacheDict(tp.Generic[X]):
If `None`, the type will be deduced automatically (Json for JSON-serializable values,
or a type-specific handler for numpy arrays, tensors, etc.).
Loading is handled using the cache_type specified in info files.
permissions: optional int
permissions for generated files
use os.chmod / path.chmod compatible numbers, or None to deactivate
eg: 0o777 for all rights to all users

Usage
-----
Expand Down Expand Up @@ -124,10 +120,8 @@ def __init__(
folder: Path | str | None,
keep_in_ram: bool = False,
cache_type: None | str = None,
permissions: int | None = 0o777,
) -> None:
self.folder = None if folder is None else Path(folder)
self.permissions = permissions
self.cache_type = cache_type
self._keep_in_ram = keep_in_ram
if self.folder is None and not keep_in_ram:
Expand All @@ -142,7 +136,7 @@ def __init__(
# DumpContext for this folder (load/delete; writes use per-thread _write_ctx)
self._dumper: DumpContext | None = None
if self.folder is not None:
self._dumper = DumpContext(self.folder, permissions=self.permissions)
self._dumper = DumpContext(self.folder)
self._local = threading.local() # per-thread write context, see _write_ctx

def __repr__(self) -> str:
Expand All @@ -155,7 +149,7 @@ def __repr__(self) -> str:
def __reduce__(self) -> tp.Any:
return (
self.__class__,
(self.folder, self._keep_in_ram, self.cache_type, self.permissions),
(self.folder, self._keep_in_ram, self.cache_type),
)

def clear(self) -> None:
Expand Down Expand Up @@ -310,7 +304,7 @@ def write(self) -> tp.Iterator["CacheDict[X]"]:
if self._write_ctx is not None:
raise RuntimeError("Cannot re-open an already open writer")
if self.folder is not None:
self._write_ctx = DumpContext(self.folder, permissions=self.permissions)
self._write_ctx = DumpContext(self.folder)
self._local.deleted_in_scope = False
try:
if self._write_ctx is not None:
Expand Down
39 changes: 3 additions & 36 deletions exca/cachedict/dumpcontext.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,13 +114,10 @@ class DumpContext:
DATA_DIR = "data"
INFO_SUFFIX = "-info.jsonl"

def __init__(
self, folder: str | Path, *, key: str = "", permissions: int | None = None
) -> None:
def __init__(self, folder: str | Path, *, key: str = "") -> None:
self.folder = Path(folder)
self.key = key
self.level: int = -1
self.permissions = permissions
self.options = DumpOptions()
# write state
self._thread_id = threading.get_native_id()
Expand Down Expand Up @@ -198,21 +195,9 @@ def __enter__(self) -> tp.Self:
self._stack = contextlib.ExitStack()
self._stack.__enter__()
self.folder.mkdir(parents=True, exist_ok=True)
self._created_files.append(self.folder) # re-chmod for shared caches
return self

def __exit__(self, *exc: tp.Any) -> None:
if self.permissions is not None:
for fp in self._created_files:
paths = [fp, *(fp.rglob("*") if fp.is_dir() else [])]
for path in paths:
try:
path.chmod(self.permissions)
except FileNotFoundError:
pass # deleted mid-walk — nothing to fix
except Exception:
msg = "Failed to set permissions on %s"
logger.warning(msg, path, exc_info=True)
if self._stack is None:
raise RuntimeError("DumpContext.__exit__ called without __enter__")
try:
Expand All @@ -221,13 +206,6 @@ def __exit__(self, *exc: tp.Any) -> None:
self._files.clear()
self._created_files.clear()

def _ensure_parent(self, path: Path) -> None:
"""Create parent directories and track them for permission setting."""
parent = path.parent
if parent != self.folder and not parent.exists():
parent.mkdir(parents=True, exist_ok=True)
self._created_files.append(parent)

def shared_file(self, suffix: str) -> tuple[tp.IO[bytes], str]:
"""Open a shared file for appending. Returns (handle, relative_name).
Content files go under DATA_DIR/; info files (-info.jsonl)
Expand All @@ -244,7 +222,7 @@ def shared_file(self, suffix: str) -> tuple[tp.IO[bytes], str]:
name = basename if is_info else f"{self.DATA_DIR}/{basename}"
if name not in self._files:
path = self.folder / name
self._ensure_parent(path)
path.parent.mkdir(parents=True, exist_ok=True)
f = path.open("ab")
self._stack.enter_context(f)
self._files[name] = f
Expand All @@ -260,7 +238,7 @@ def key_path(self, suffix: str = "") -> str:
basename = string_uid(self.key) + suffix
name = f"{self.DATA_DIR}/{basename}"
path = self.folder / name
self._ensure_parent(path)
path.parent.mkdir(parents=True, exist_ok=True)
if path in self._created_files:
# Same dump context tried to create this path twice: user error
raise RuntimeError(
Expand Down Expand Up @@ -396,21 +374,10 @@ def _dump_cls(self, cls: tp.Any, value: tp.Any) -> tuple[dict[str, tp.Any], str]
"ctx.key must be set before dumping with a legacy DumperLoader"
)
info = self._loaders[cls].dump(self.key, value)
self._track_legacy_files(info)
else:
info = cls.__dump_info__(self, value)
return info, cls.__name__

def _track_legacy_files(self, info: tp.Any) -> None:
"""Record files from legacy DumperLoader info dicts for permission setting.
New-style handlers track files at creation (keyed_filepath / shared_file)."""
if isinstance(info, dict):
if "filename" in info:
self._created_files.append(self.folder / info["filename"])
for val in info.values():
if isinstance(val, dict):
self._track_legacy_files(val)

def _resolve_type(self, info: dict[str, tp.Any]) -> tuple[tp.Any, dict[str, tp.Any]]:
"""Extract #type and #key from an info dict, return (cls, remaining_info).
Applies ``options.replace`` before handler lookup."""
Expand Down
12 changes: 1 addition & 11 deletions exca/cachedict/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,9 +73,8 @@ class AdvisoryRegistry:
_SCHEMA: tp.ClassVar[str] # passed to executescript(), multi-statement OK
_LABEL: tp.ClassVar[str] # short prefix in log messages

def __init__(self, folder: Path | str, permissions: int | None = 0o777) -> None:
def __init__(self, folder: Path | str) -> None:
self.db_path = Path(folder) / self._DB_NAME
self.permissions = permissions
self._conn: sqlite3.Connection | None = None

def _connect(self, *, create: bool = False) -> sqlite3.Connection | None:
Expand Down Expand Up @@ -111,15 +110,6 @@ def _connect(self, *, create: bool = False) -> sqlite3.Connection | None:
# WAL needs cross-host shared memory (broken on NFS) -> DELETE journal
conn.execute("PRAGMA journal_mode=DELETE")
conn.executescript(self._SCHEMA)
if self.permissions is not None:
try:
self.db_path.chmod(self.permissions)
except Exception:
logger.warning(
"Failed to set permissions on %s",
self.db_path,
exc_info=True,
)
self._conn = conn
return conn

Expand Down
6 changes: 0 additions & 6 deletions exca/cachedict/test_cachedict.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,12 +143,6 @@ def test_specialized_dump(
assert files, "Some memmaps should stay open"
del cache
gc.collect()
# check permissions
octal_permissions = oct(tmp_path.stat().st_mode)[-3:]
assert octal_permissions == "777", f"Wrong permissions for {tmp_path}"
for fp in tmp_path.rglob("*"):
octal_permissions = oct(fp.stat().st_mode)[-3:]
assert octal_permissions == "777", f"Wrong permissions for {fp}"
# after del, all files should be closed
files = proc.open_files()
assert not files, "No file should remain open after del cache"
Expand Down
8 changes: 3 additions & 5 deletions exca/cachedict/test_dumpcontext.py
Original file line number Diff line number Diff line change
Expand Up @@ -122,16 +122,14 @@ def test_shared_file_lifecycle(tmp_path: Path) -> None:
assert (tmp_path / name1).read_bytes() == b"hello"


def test_context_permissions(tmp_path: Path) -> None:
def test_context_creates_folder_lazily(tmp_path: Path) -> None:
folder = tmp_path / "fresh"
ctx = DumpContext(folder, permissions=0o755)
ctx = DumpContext(folder)
assert not folder.exists(), "construction must not materialise the folder"
with ctx:
f, name = ctx.shared_file(".data")
f.write(b"test")
assert folder.is_dir()
assert oct(folder.stat().st_mode)[-3:] == "755"
assert oct((folder / name).stat().st_mode)[-3:] == "755"
assert (folder / name).read_bytes() == b"test"


# =============================================================================
Expand Down
8 changes: 0 additions & 8 deletions exca/cachedict/test_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -194,11 +194,3 @@ def test_graceful_degradation(
reg2._plant(["recovered"])
assert reg2.get(["recovered"]) == {"recovered"}
reg2.close()


def test_permissions_applied(tmp_path: Path) -> None:
reg = errors.ErrorRegistry(tmp_path, permissions=0o600)
reg._plant(["a"])
mode = stat.S_IMODE((tmp_path / "errors.db").stat().st_mode)
assert mode == 0o600
reg.close()
2 changes: 1 addition & 1 deletion exca/map.py
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,7 @@ def _inflight_registry(self) -> inflight.InflightRegistry | None:
cache_folder = self.uid_folder()
if cache_folder is None:
return None
return inflight.InflightRegistry(cache_folder, permissions=self.permissions)
return inflight.InflightRegistry(cache_folder)

# pylint: disable=unused-argument
def apply(
Expand Down
1 change: 0 additions & 1 deletion exca/steps/backends.py
Original file line number Diff line number Diff line change
Expand Up @@ -551,7 +551,6 @@ def _cache_dict(
folder=cache_folder,
cache_type=cache_type,
keep_in_ram=self.keep_in_ram,
permissions=0o777,
)
self._cds[cache_folder] = cd
return cd
Expand Down
3 changes: 1 addition & 2 deletions exca/steps/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
import pydantic

import exca
from exca import utils as xkutils

from . import backends, identity, items, utils

Expand Down Expand Up @@ -295,7 +294,7 @@ def _make_paths(self, aligned: tp.Sequence[Step]) -> backends.StepPaths:
identity.step_uid(aligned),
cache_type=self._infer_cache_type(),
)
xkutils.mkdir_with_permissions(paths.step_folder, 0o777, root=paths.base_folder)
paths.step_folder.mkdir(parents=True, exist_ok=True)
return paths

def _exca_uid_dict_override(self) -> dict[str, tp.Any] | None:
Expand Down
11 changes: 0 additions & 11 deletions exca/steps/test_backends.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@

import contextlib
import logging
import stat
import sys
import time
import typing as tp
Expand Down Expand Up @@ -68,16 +67,6 @@ def test_backend_execution(tmp_path: Path, backend: str) -> None:
assert (job is not None) == (backend == "LocalProcess")


def test_step_permissions(tmp_path: Path) -> None:
infra: tp.Any = {"backend": "Cached", "folder": tmp_path}
chain = Chain(steps=[conftest.Mult(coeff=2), conftest.Add(value=1)], infra=infra)
intermediate = chain.lookup(1).paths.step_folder.parent
intermediate.mkdir(parents=True)
intermediate.chmod(0o700)
chain.run(1)
assert stat.S_IMODE(intermediate.stat().st_mode) == 0o777


def test_slurm_backend_param_forwarding(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
Expand Down
Loading
Loading