"""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")
"""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)