Source code for spacr.workspace

"""Save and restore the interactive workspace associated with a run.

:mod:`spacr.run_journal` records pipeline inputs, settings, environment data,
hashes, and logs. This module complements that record with interactive state:
attached databases, generated montages, figure views, and the files those
panels reference. The workspace is stored beside the run::

    ~/.spacr/runs/<run>/workspace.json
    ~/.spacr/runs/<run>/workspace/files/<digest>__<name>   # 'copy' mode only

Public API::

    from spacr import workspace

    workspace.register("volcano", lambda: the_panel)     # GUI, once
    doc = workspace.collect(workspace.providers())
    workspace.save(run_dir, doc, mode="reference")

    doc = workspace.load(run_dir)
    report = workspace.restore(workspace.providers(), doc)

Each registered provider supplies its own state through
``workspace_state()`` and restores it through ``apply_workspace_state()``.
The collector does not inspect widget internals. Existing ``plot_state`` and
``apply_plot_state`` methods are accepted for compatibility.

Every referenced file is recorded with its size, modification time, and
SHA-256 digest. ``reference`` mode stores only this inventory; ``copy`` mode
also carries files up to the configured per-file limit. A provider can mark an
individual artifact for copying regardless of mode or size. Skipped files and
copy failures remain visible in the document and restore report.

The module uses only the Python standard library so headless pipelines can
record workspace state without importing pandas or Qt.
"""
from __future__ import annotations

import hashlib
import json
import logging
import os
import re
import shutil
import threading
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable, Dict, Iterator, List, Mapping, Optional, Tuple

LOG = logging.getLogger("spacr.workspace")

[docs] SCHEMA_VERSION = 1
"""Schema version written to ``workspace.json``.""" DOC_NAME = "workspace.json" FILES_DIR = "workspace/files"
[docs] MODES = ("off", "reference", "copy")
"""Supported workspace persistence modes. ``off`` Do not write workspace metadata. ``reference`` Record referenced files and copy only artifacts explicitly marked for retention. This is the default. ``copy`` Also copy each recorded file within the configured per-file size limit. """ DEFAULT_MODE = "reference" DEFAULT_COPY_LIMIT_MB = 512 #: Reserved key for files a section retains with the workspace, rather than #: merely recording. Values use ``[{"role": str, "path": str}, ...]`` and #: normally identify generated artifacts such as figures or exported tables. CARRY_KEY = "_workspace_carry" #: Section keys that are metadata about the section rather than state to put #: back, and are therefore not handed to ``apply_workspace_state``. _RESERVED = (CARRY_KEY,) _LOCK = threading.RLock() _PROVIDERS: "Dict[str, Callable[[], Any]]" = {} _DEFAULT_MODE = DEFAULT_MODE _DEFAULT_LIMIT_MB = float(DEFAULT_COPY_LIMIT_MB)
[docs] def set_default_mode(mode: Any, copy_limit_mb: Any = None) -> str: """Set the process-wide workspace defaults. An explicit ``save_workspace`` value in run settings overrides this mode. :param mode: ``'off'``, ``'reference'``, or ``'copy'``. Common boolean and yes/no aliases are accepted by :func:`resolve_mode`. :param copy_limit_mb: optional non-negative per-file copy limit in MB. Invalid or negative values leave the current limit unchanged. :returns: the normalized default mode. """ global _DEFAULT_MODE, _DEFAULT_LIMIT_MB with _LOCK: _DEFAULT_MODE = resolve_mode(mode) if copy_limit_mb is not None: try: limit = float(copy_limit_mb) if limit >= 0: _DEFAULT_LIMIT_MB = limit except (TypeError, ValueError): pass return _DEFAULT_MODE
[docs] def default_mode() -> str: """Return the process-wide workspace mode.""" with _LOCK: return _DEFAULT_MODE
[docs] def default_copy_limit_mb() -> float: """Return the process-wide per-file copy limit in megabytes.""" with _LOCK: return float(_DEFAULT_LIMIT_MB)
[docs] def register(name: str, provider: Callable[[], Any]) -> None: """Register a named workspace-state provider. :param name: stable section name written into ``workspace.json``. :param provider: zero-argument callable returning either a mapping or an object with ``workspace_state()``. Resolving the object at collection time avoids retaining a widget that has since been rebuilt or closed. :raises ValueError: if ``name`` is empty. :raises TypeError: if ``provider`` is not callable. """ if not name: raise ValueError("a workspace contributor needs a name") if not callable(provider): raise TypeError("a workspace provider must be callable") with _LOCK: _PROVIDERS[str(name)] = provider
[docs] def unregister(name: str) -> bool: """Remove a provider and return whether it was registered. :param name: registered workspace section name to remove. """ with _LOCK: return _PROVIDERS.pop(str(name), None) is not None
[docs] def providers() -> "Dict[str, Callable[[], Any]]": """Return a snapshot of the registered workspace providers.""" with _LOCK: return dict(_PROVIDERS)
[docs] def clear_providers() -> None: """Remove all providers and restore the built-in defaults. This resets process state during application shutdown and test teardown. """ global _DEFAULT_MODE, _DEFAULT_LIMIT_MB with _LOCK: _PROVIDERS.clear() _DEFAULT_MODE = DEFAULT_MODE _DEFAULT_LIMIT_MB = float(DEFAULT_COPY_LIMIT_MB)
def _state_of(source: Any) -> Optional[Dict[str, Any]]: """The state dict a contributor offers, or ``None`` if it offers none.""" if source is None: return None if isinstance(source, Mapping): return dict(source) getter = getattr(source, "workspace_state", None) if callable(getter): state = getter() return dict(state) if isinstance(state, Mapping) else None legacy = getattr(source, "plot_state", None) if callable(legacy): state = legacy() return dict(state) if isinstance(state, Mapping) else None return None
[docs] def section_states( contributors: Mapping[str, Any], ) -> Tuple[Dict[str, Any], List[Dict[str, str]]]: """Collect one JSON-compatible state section from each provider. :param contributors: ``{name: provider-or-object}``. A callable value is called; anything else is used directly. :returns: ``(sections, problems)``. A provider that raises or returns no usable state is omitted and described in ``problems``; other sections are still collected. """ sections: Dict[str, Any] = {} problems: List[Dict[str, str]] = [] for name in sorted(contributors): source = contributors[name] try: if callable(source) and not isinstance(source, Mapping): source = source() if source is None: continue state = _state_of(source) except Exception as exc: # noqa: BLE001 problems.append({"section": name, "why": f"{type(exc).__name__}: {exc}"}) LOG.debug("workspace section %r could not be collected", name, exc_info=True) continue if state is None: problems.append({"section": name, "why": "offers no workspace state"}) continue sections[name] = _jsonable(state) return sections, problems
def _jsonable(value: Any, _depth: int = 0) -> Any: """A JSON-writable copy of ``value``. Tuples become lists, Paths and everything exotic become strings. Depth is bounded because a widget state that accidentally contains a cycle must cost a truncated section, not the process. """ if _depth > 12: return str(value) if value is None or isinstance(value, (bool, int, float, str)): return value if isinstance(value, Path): return str(value) if isinstance(value, Mapping): return {str(k): _jsonable(v, _depth + 1) for k, v in value.items()} if isinstance(value, (list, tuple, set, frozenset)): items = sorted(value, key=str) if isinstance(value, (set, frozenset)) else value return [_jsonable(v, _depth + 1) for v in items] return str(value) _PATHISH = re.compile(r"[/\\]") def _looks_like_a_path(value: Any) -> bool: """Whether a section value is worth asking the filesystem about. Cheap and deliberately conservative: a separator and a plausible length. The filesystem decides in the end — this only keeps the walk from calling ``stat`` on every gene name in a picked-cell list. """ if not isinstance(value, str) or not (2 < len(value) < 4096): return False return bool(_PATHISH.search(value)) def _walk_strings(value: Any, trail: str = "", _depth: int = 0) -> Iterator[Tuple[str, str]]: """Yield ``(dotted-key, string)`` for every string inside ``value``.""" if _depth > 12: return if isinstance(value, str): yield trail, value elif isinstance(value, Mapping): for k, v in value.items(): if k in _RESERVED: continue yield from _walk_strings(v, f"{trail}.{k}" if trail else str(k), _depth + 1) elif isinstance(value, (list, tuple)): for i, v in enumerate(value): yield from _walk_strings(v, f"{trail}[{i}]", _depth + 1)
[docs] def hash_file(path: Path, chunk_size: int = 1 << 20) -> Optional[str]: """A file's full SHA-256, or ``None`` if it cannot be read. :param path: regular file whose bytes are hashed. """ try: h = hashlib.sha256() with open(path, "rb") as fh: for chunk in iter(lambda: fh.read(chunk_size), b""): h.update(chunk) return h.hexdigest() except Exception as exc: # noqa: BLE001 LOG.debug("could not hash %s: %s", path, exc) return None
def _file_record(role: str, path: Path, *, want_hash: bool = True) -> Dict[str, Any]: """What is recorded about one file, whether or not its bytes are carried.""" record: Dict[str, Any] = {"role": role, "path": str(path)} try: stat = path.stat() except Exception: # noqa: BLE001 record["exists"] = False return record record["exists"] = True if path.is_dir(): record["kind"] = "directory" return record record["kind"] = "file" record["size"] = int(stat.st_size) record["mtime"] = round(stat.st_mtime, 3) if want_hash: digest = hash_file(path) if digest: record["sha256"] = digest return record def _carried(section: Any) -> List[Dict[str, str]]: """The files a section explicitly asks to have carried.""" if not isinstance(section, Mapping): return [] out = [] for entry in section.get(CARRY_KEY) or []: if isinstance(entry, Mapping) and entry.get("path"): out.append({"role": str(entry.get("role") or "carried"), "path": str(entry["path"])}) elif isinstance(entry, str): out.append({"role": "carried", "path": entry}) return out
[docs] def inventory(sections: Mapping[str, Any], *, hash_files: bool = True) -> List[Dict[str, Any]]: """Build metadata records for files referenced by workspace sections. String values that resolve to existing paths are discovered recursively. A section may also list files under the reserved carry key to request that their bytes be included in the bundle. :param sections: workspace state keyed by section name. :param hash_files: compute SHA-256 for regular files when true. :returns: deterministic file records ordered by source path. """ seen: Dict[str, Dict[str, Any]] = {} def add(role: str, raw: str, carry: bool) -> None: """Add or promote one path in the captured inventory. :param role: source role recorded when the path is first seen. :param raw: path text to expand and inspect. :param carry: whether the file must be included in a bundle. :returns: None. Invalid and missing paths are ignored; repeated paths reuse their first record and are promoted when any occurrence requests carrying. File hashing follows the captured ``hash_files`` policy. """ try: path = Path(raw).expanduser() except Exception: # noqa: BLE001 return key = str(path) record = seen.get(key) if record is None: record = _file_record(role, path, want_hash=hash_files) if not record.get("exists"): return seen[key] = record if carry: record["carry"] = True for name in sorted(sections): section = sections[name] for entry in _carried(section): add(f"{name}:{entry['role']}", entry["path"], True) for trail, value in _walk_strings(section, name): if _looks_like_a_path(value): add(trail, value, False) return [seen[k] for k in sorted(seen)]
[docs] def collect( contributors: Mapping[str, Any], *, app_key: str = "", saved: str = "", hash_files: bool = True, ) -> Dict[str, Any]: """Build a workspace document without writing it to disk. :param contributors: named providers or state-bearing objects. :param app_key: application key associated with the workspace. :param saved: the timestamp to stamp, ISO-8601. Injected rather than read from the clock so a caller can produce a byte-identical document twice, which is what makes the writer testable. :param hash_files: compute file digests for the inventory when true. :returns: JSON-compatible workspace document. """ sections, problems = section_states(contributors) doc: Dict[str, Any] = { "version": SCHEMA_VERSION, "saved": saved or datetime.now(timezone.utc).isoformat(timespec="seconds"), "app_key": str(app_key or ""), "sections": sections, "files": inventory(sections, hash_files=hash_files), } if problems: doc["problems"] = problems return doc
[docs] def resolve_mode(value: Any) -> str: """Normalize a workspace-saving mode. :param value: mode name, boolean, common yes/no alias, or ``None``. ``None`` and unrecognized values select :data:`DEFAULT_MODE`; boolean and common yes/no aliases are accepted. Invalid values do not stop a run. """ if value is None: return DEFAULT_MODE if isinstance(value, bool): return "reference" if value else "off" text = str(value).strip().lower() if text in MODES: return text if text in ("none", "no", "false", "0", ""): return "off" if text in ("all", "full", "yes", "true", "1"): return "copy" LOG.debug("unrecognised save_workspace=%r, using %s", value, DEFAULT_MODE) return DEFAULT_MODE
[docs] def mode_from_settings(settings: Mapping[str, Any]) -> str: """Return the workspace mode requested by a settings mapping. :param settings: run settings that may declare ``save_workspace``. """ if not isinstance(settings, Mapping) or settings.get("save_workspace") is None: return default_mode() return resolve_mode(settings.get("save_workspace"))
[docs] def copy_limit_from_settings(settings: Mapping[str, Any]) -> float: """Return the non-negative per-file copy limit in megabytes. :param settings: run settings that may declare a workspace copy limit. """ if isinstance(settings, Mapping) and settings.get("workspace_copy_limit_mb") is not None: try: limit = float(settings["workspace_copy_limit_mb"]) if limit >= 0: return limit except (TypeError, ValueError): pass return default_copy_limit_mb()
def _safe_name(path: Path) -> str: """A filename that cannot escape the bundle or collide by basename.""" stem = re.sub(r"[^A-Za-z0-9._-]+", "_", path.name)[:80] or "file" return stem
[docs] def save( run_dir: Any, doc: Mapping[str, Any], *, mode: str = DEFAULT_MODE, copy_limit_mb: float = DEFAULT_COPY_LIMIT_MB, ) -> Optional[Path]: """Write a workspace document and optionally copy referenced files. ``off`` writes nothing at all — not an empty document — so a run folder saved with the feature disabled is byte-for-byte what it was before this module existed. :param run_dir: destination run folder. :param doc: document returned by :func:`collect`. :param mode: ``'off'``, ``'reference'``, or ``'copy'``. :param copy_limit_mb: maximum size of each automatically copied file in MB. Explicitly carried artifacts are not limited. :returns: path to ``workspace.json``, or ``None`` in ``off`` mode. """ mode = resolve_mode(mode) if mode == "off": return None root = Path(run_dir) root.mkdir(parents=True, exist_ok=True) payload = dict(doc) payload["mode"] = mode limit_bytes = float(copy_limit_mb) * 1024 * 1024 files = [dict(f) for f in payload.get("files") or []] if files: files_root = root / FILES_DIR for record in files: if not record.get("exists") or record.get("kind") != "file": continue wanted = mode == "copy" or record.get("carry") if not wanted: continue size = record.get("size") or 0 if mode == "copy" and not record.get("carry") and size > limit_bytes: record["copied"] = None record["skipped"] = ( f"{size / 1024 / 1024:.1f} MB is over the " f"{copy_limit_mb:g} MB per-file limit" ) continue source = Path(record["path"]) digest = record.get("sha256") or hash_file(source) or "" target_name = f"{digest[:16]}__{_safe_name(source)}" if digest else _safe_name(source) target = files_root / target_name try: files_root.mkdir(parents=True, exist_ok=True) if not target.exists(): shutil.copy2(source, target) record["copied"] = f"{FILES_DIR}/{target_name}" except Exception as exc: # noqa: BLE001 record["copied"] = None record["skipped"] = f"could not copy: {type(exc).__name__}: {exc}" LOG.debug("workspace could not copy %s", source, exc_info=True) payload["files"] = files path = root / DOC_NAME text = json.dumps(payload, indent=2, sort_keys=False, default=str) tmp = path.with_suffix(".json.tmp") tmp.write_text(text, encoding="utf-8") os.replace(tmp, path) LOG.info("workspace saved [%s] → %s", mode, path) return path
[docs] def save_for_run( run_dir: Any, settings: Optional[Mapping[str, Any]] = None, *, app_key: str = "", contributors: Optional[Mapping[str, Any]] = None, ) -> Optional[Path]: """Collect the registered workspace and write it beside a finished run. Nothing is written when saving is disabled or no provider supplies state. Headless commands therefore do not gain an empty workspace document. :param run_dir: destination run folder. :param settings: run settings controlling mode and copy limit. :param app_key: application key associated with the saved state. :param contributors: optional providers to use instead of the registry. :returns: path to ``workspace.json``, or ``None`` when nothing was saved. """ mode = mode_from_settings(settings or {}) if mode == "off": return None sources = providers() if contributors is None else contributors if not sources: return None doc = collect(sources, app_key=app_key) if not doc.get("sections") and not doc.get("files"): return None return save(run_dir, doc, mode=mode, copy_limit_mb=copy_limit_from_settings(settings or {}))
[docs] def load(run_dir: Any) -> Optional[Dict[str, Any]]: """Read a workspace document from a run folder or document path. :param run_dir: run directory or workspace document path to read. :returns: the decoded document, or ``None`` when it is absent or invalid. """ path = Path(run_dir) if path.is_file(): path = path if path.name == DOC_NAME else path.parent / DOC_NAME else: path = path / DOC_NAME try: doc = json.loads(path.read_text(encoding="utf-8")) except FileNotFoundError: return None except Exception as exc: # noqa: BLE001 LOG.warning("could not read %s: %s", path, exc) return None return doc if isinstance(doc, dict) else None
[docs] def has_workspace(run_dir: Any) -> bool: """Return whether ``run_dir`` contains ``workspace.json``. :param run_dir: run directory to inspect. """ try: return (Path(run_dir) / DOC_NAME).is_file() except Exception: # noqa: BLE001 return False
[docs] def check_files(doc: Mapping[str, Any], *, run_dir: Any = None) -> List[Dict[str, Any]]: """Check the current state of every file in a workspace document. ``state`` is one of ``present`` (there, and the same bytes), ``changed`` (there, different digest), ``carried`` (gone from its original place but inside the bundle), or ``missing``. :param doc: decoded workspace document. :param run_dir: bundle root used to resolve carried files. :returns: one state record per referenced file. """ out: List[Dict[str, Any]] = [] root = Path(run_dir) if run_dir is not None else None for record in doc.get("files") or []: if not isinstance(record, Mapping): continue path = Path(str(record.get("path", ""))) entry = {"role": record.get("role", ""), "path": str(path)} copied = record.get("copied") bundled = (root / copied) if (root is not None and copied) else None if path.exists(): digest = record.get("sha256") if record.get("kind") == "file" and digest: entry["state"] = "present" if hash_file(path) == digest else "changed" else: entry["state"] = "present" elif bundled is not None and bundled.exists(): entry["state"] = "carried" entry["path"] = str(bundled) else: entry["state"] = "missing" if record.get("skipped"): entry["skipped"] = record["skipped"] out.append(entry) return out
[docs] def restore( contributors: Mapping[str, Any], doc: Mapping[str, Any], *, run_dir: Any = None, ) -> Dict[str, Any]: """Restore each workspace section through its registered provider. Restoration is best-effort. Missing panels, unsupported state, exceptions, and provider refusals are reported under ``skipped`` rather than silently discarded. :param contributors: providers keyed by workspace section name. :param doc: decoded workspace document. :param run_dir: optional bundle root used to resolve carried files. :returns: report containing restored sections, skipped sections, and file states. """ report: Dict[str, Any] = {"restored": [], "skipped": [], "files": []} sections = doc.get("sections") if isinstance(doc, Mapping) else None if not isinstance(sections, Mapping): report["skipped"].append({"section": "", "why": "no sections in the document"}) return report for name in sorted(sections): state = sections[name] if not isinstance(state, Mapping): report["skipped"].append({"section": name, "why": "not a state document"}) continue source = contributors.get(name) if source is None: report["skipped"].append({"section": name, "why": "nothing on screen owns it"}) continue payload = {k: v for k, v in state.items() if k not in _RESERVED} try: if callable(source) and not isinstance(source, Mapping): source = source() setter = getattr(source, "apply_workspace_state", None) if not callable(setter): setter = getattr(source, "apply_plot_state", None) if not callable(setter): report["skipped"].append( {"section": name, "why": "cannot take a workspace state back"}) continue applied = setter(payload) except Exception as exc: # noqa: BLE001 report["skipped"].append( {"section": name, "why": f"{type(exc).__name__}: {exc}"}) LOG.debug("workspace section %r could not be restored", name, exc_info=True) continue if applied is False: report["skipped"].append({"section": name, "why": "the panel declined it"}) else: report["restored"].append(name) report["files"] = check_files(doc, run_dir=run_dir) return report
def _human_size(n: Any) -> str: """Format a byte count using compact binary-scaled units. :param n: Byte count or number-like value to display. :returns: Rounded text in B, KB, MB, GB, or TB, or ``?`` when ``n`` is not numeric. """ try: size = float(n) except (TypeError, ValueError): return "?" for unit in ("B", "KB", "MB", "GB"): if size < 1024: return f"{size:.0f} {unit}" if unit == "B" else f"{size:.1f} {unit}" size /= 1024 return f"{size:.1f} TB"
[docs] def inventory_text(doc: Mapping[str, Any], *, run_dir: Any = None) -> str: """Format a workspace document as a human-readable inventory. :param doc: decoded workspace document to summarize. """ if not isinstance(doc, Mapping): return "no workspace document" lines = [ f"workspace v{doc.get('version', '?')} " f"[{doc.get('mode', DEFAULT_MODE)}] saved {doc.get('saved', '?')}", ] if doc.get("app_key"): lines.append(f" app: {doc['app_key']}") sections = doc.get("sections") or {} lines.append(f" sections ({len(sections)}):") for name in sorted(sections): state = sections[name] keys = len(state) if isinstance(state, Mapping) else 0 lines.append(f" {name:<24} {keys} key{'' if keys == 1 else 's'}") files = doc.get("files") or [] states = {e["path"]: e.get("state", "") for e in check_files(doc, run_dir=run_dir)} lines.append(f" files ({len(files)}):") for record in files: if not isinstance(record, Mapping): continue path = str(record.get("path", "")) bits = [states.get(path, "")] if record.get("copied"): bits.append("carried") if record.get("skipped"): bits.append(f"skipped: {record['skipped']}") if record.get("size") is not None: bits.append(_human_size(record.get("size"))) note = ", ".join(b for b in bits if b) lines.append(f" {path} ({note})") for problem in doc.get("problems") or []: lines.append(f" ! {problem.get('section', '?')}: {problem.get('why', '')}") return "\n".join(lines)
[docs] def report_text(report: Mapping[str, Any]) -> str: """Format a restore report as human-readable text. :param report: workspace restoration report to summarize. """ restored = report.get("restored") or [] skipped = report.get("skipped") or [] lines = [] if restored: lines.append(f"restored: {', '.join(restored)}") else: lines.append("restored: nothing") for entry in skipped: name = entry.get("section") or "(document)" lines.append(f" not restored — {name}: {entry.get('why', '')}") trouble = [f for f in (report.get("files") or []) if f.get("state") in ("missing", "changed")] for entry in trouble: lines.append(f" {entry['state']} — {entry.get('role', '')}: {entry.get('path', '')}") return "\n".join(lines)