From bfd4373136b36d663787f5fe9dade0a5b834fd90 Mon Sep 17 00:00:00 2001 From: Gareth Price Date: Wed, 7 Oct 2026 17:25:26 +0000 Subject: [PATCH 1/2] feat: publish complete local file export sets (#1062) Register phlo.export as an executable capability asset. Stage each attempt independently, retain immutable run directories, and atomically replace the current manifest pointer. Add lifecycle and Dagster discovery tests plus runnable export documentation. --- docs/guides/export-files.md | 92 +++++++++ docs/index.md | 1 + docs/reference/python-api.md | 35 ++++ src/phlo/__init__.py | 6 + src/phlo/exports.py | 213 ++++++++++++++++++++ tests/runtime/test_exports.py | 234 ++++++++++++++++++++++ tests/runtime/test_exports_integration.py | 113 +++++++++++ 7 files changed, 694 insertions(+) create mode 100644 docs/guides/export-files.md create mode 100644 src/phlo/exports.py create mode 100644 tests/runtime/test_exports.py create mode 100644 tests/runtime/test_exports_integration.py diff --git a/docs/guides/export-files.md b/docs/guides/export-files.md new file mode 100644 index 00000000000..3f567735518 --- /dev/null +++ b/docs/guides/export-files.md @@ -0,0 +1,92 @@ +# Export a complete set of files + +This example publishes `catalog.csv` and `summary.json` as one export asset, then reads both through one manifest. + +## Before you start + +- Install `phlo` and `phlo-dagster` in your Python environment. +- Run the commands from your project directory. In this repository, prefix Python commands with `uv run --locked`. +- Choose a local destination that every worker and consumer can access at the same absolute path. A container-local directory is not shared storage. + +## 1. Write the export + +Create the `workflows` directory if it does not exist. +Save this file as `workflows/sky_survey.py`: + +```python +import csv +import json + +import phlo +from phlo.exports import ExportContext + + +@phlo.export( + name="monthly_sky_survey", + destination=".phlo/exports", + outputs={"catalog": "catalog.csv", "summary": "summary.json"}, + max_retries=2, +) +def export_sky_survey(context: ExportContext) -> None: + stars = [ + {"name": "Sirius", "ra": 101.287}, + {"name": "Vega", "ra": 279.235}, + ] + with (context.staging_dir / "catalog.csv").open( + "w", encoding="utf-8", newline="" + ) as handle: + writer = csv.DictWriter(handle, fieldnames=["name", "ra"]) + writer.writeheader() + writer.writerows(stars) + (context.staging_dir / "summary.json").write_text( + json.dumps({"count": len(stars), "run_id": context.run_id}), + encoding="utf-8", + ) +``` + +The example uses fixed observations and needs no database. For a real source, declare its asset key in `depends_on` and select data using `context.runtime.routing` and `context.runtime.resources`. Record available snapshot or version references in `context.upstream_versions` under the declared dependency key at the time you read the data. Leave unavailable versions absent. Phlo does not infer versions from a later "latest" lookup. + +Close every output before the writer returns. Create parent directories for nested output paths inside `context.staging_dir`. + +## 2. Discover and run the asset + +Save the following as `run_export.py` and run `python run_export.py`: + +```python +from pathlib import Path + +import dagster as dg +from phlo_dagster.framework.discovery import discover_user_workflows + +from phlo.exports import resolve_export_manifest +from phlo.helpers.artifacts import verify_manifest_checksums + + +definitions = discover_user_workflows("workflows", clear_registries=True) +assets = [ + asset for asset in definitions.assets + if isinstance(asset, dg.AssetsDefinition) + and dg.AssetKey("monthly_sky_survey") in asset.keys +] +result = dg.materialize(assets) +assert result.success + +manifest = resolve_export_manifest(".phlo/exports", "monthly_sky_survey") +assert all(verify_manifest_checksums(manifest).values()) +print("Published", manifest.metadata["run_id"]) +for artifact in manifest.artifacts: + print(artifact.metadata["name"], artifact.size_bytes, artifact.checksum) + print(Path(artifact.uri).read_text(encoding="utf-8")) +``` + +The output includes `Published`, the run ID, both artifact names, SHA-256 checksums, and their file contents. The summary has `count` equal to `2`. + +The same workflow loader powers project discovery. The export registers an executable `AssetSpec` with a `RunSpec`, not a deprecated flow declaration. Dagster shows one materialisation with the named artifact manifest and immutable manifest path in its metadata. + +## 3. Read one complete set + +Call `resolve_export_manifest` once per consumer operation, as the script does. Read every artifact using the paths in that returned manifest. Do not resolve current separately for each file, because another run can publish between reads. + +If a writer raises or omits an output, inspect the orchestration failure and fix the writer. The previous complete set remains current. If no run has succeeded, resolving current raises `FileNotFoundError`. + +For storage guarantees, retry identity, and all decorator parameters, see the [file export API reference](../reference/python-api.md#phloexport). diff --git a/docs/index.md b/docs/index.md index d7472bcfd3f..3c04bcdf546 100644 --- a/docs/index.md +++ b/docs/index.md @@ -51,6 +51,7 @@ Each guide solves one problem and assumes you have finished the first pipeline: | Change a table's columns without breaking consumers | [Evolve a schema](guides/evolve-a-schema.md) | | Swap storage, catalog, query engine, or ingestion tool | [Choose your stack](guides/choose-your-stack.md) | | Serve tables to apps, analysts, and BI tools | [Expose data](guides/expose-data.md) | +| Publish a complete set of local files | [Export files](guides/export-files.md) | | Add authentication, per-service credentials, and audit logs | [Secure the stack](guides/secure-the-stack.md) | | Read logs, trace lineage, and fix a failing run | [Monitor and debug](guides/monitor-and-debug.md) | | Diagnose a `PHLO-` error code | [Troubleshoot Phlo errors](guides/troubleshoot-errors.md) | diff --git a/docs/reference/python-api.md b/docs/reference/python-api.md index 6bcf3544199..3f63bf7d235 100644 --- a/docs/reference/python-api.md +++ b/docs/reference/python-api.md @@ -2,6 +2,41 @@ The public APIs below are defined in `phlo`, `phlo-dlt`, `phlo-sling`, and `phlo-pandera`. Generated API pages are not emitted under the `python-reference` route by the current pymdx build, so this page does not link to that route. +## phlo.export + +`phlo.export` registers one executable capability asset for a complete set of files. The writer accepts `phlo.exports.ExportContext` and returns `None`. The [runnable export guide](../guides/export-files.md) demonstrates discovery, execution, and consumption. + +| Parameter | Type | Default | Meaning | +| --- | --- | --- | --- | +| `name` | `str` | required | Asset key and destination subdirectory. Must be a single non-empty path component. | +| `destination` | `str \| Path` | required | Local parent directory, resolved at declaration time. Remote URLs are unsupported. | +| `outputs` | `Mapping[str, str]` | required | Artifact names mapped to distinct relative file paths. At least one output is required. Absolute paths, parent traversal, and the reserved `manifest.json` path are rejected. | +| `depends_on` | `Sequence[str]` | `()` | Upstream asset keys registered as orchestration dependencies. | +| `group` | `str` | `"exports"` | Orchestration asset group. | +| `resources` | `Iterable[str]` | `()` | Required runtime resource keys. | +| `max_retries` | `int` | `0` | Retry count passed to `RunSpec`. | +| `retry_delay_seconds` | `int` | `30` | Retry delay passed to `RunSpec`. | + +`ExportContext` contains `staging_dir`, the canonical `run_id`, the orchestrator-neutral `runtime`, and a mutable `upstream_versions: dict[str, str]`. Version keys must be declared dependencies. References describe the data selected by the writer. Missing references mean that the version is unavailable, not that an unversioned source has a fabricated version. + +`resolve_export_manifest(destination, name)` reads `current.json` once, then returns the referenced `ArtifactManifest`. Each entry has an absolute local path in `uri`, SHA-256 checksum, byte size, and its artifact name in `metadata["name"]`. Manifest metadata contains `run_id`, `partition_key`, `ref`, and `upstream_versions`. + +The materialisation metadata keys are `phlo/export_manifest`, `phlo/export_manifest_path`, `phlo/export_current_path`, and `phlo/export_run_id`. Static asset metadata includes `phlo/export_outputs` and `phlo/export_destination`. + +### Local publication and retries + +Each attempt gets a fresh temporary directory on the destination filesystem. Phlo copies only declared regular files into a separate complete-set directory, computes checksums, and writes its manifest. Symlink outputs and symlink parent directories are rejected. Undeclared scratch files are discarded. + +Phlo renames the complete directory to `destination/name/runs/` without replacing an existing complete run. It then atomically replaces `destination/name/current.json` with a pointer to that run's manifest. Consumers must reuse one resolved manifest for all files in a read operation. + +Writer failures and missing outputs leave current unchanged. A failure to replace current can leave a complete but unreferenced run directory. Retrying the same identity with the same files and manifest metadata reuses that directory and retries the pointer replacement. Different content, different upstream references, or checksum corruption cause an error instead of overwriting the published run. + +The canonical runtime routing run ID is required. Dagster step retries retain that ID. Dagster run re-execution uses the root run ID when available. The `phlo/run_id` run tag explicitly overrides the logical identity. A new export set requires a new logical identity, even if a writer would produce different files under a reused ID. + +Successful concurrent runs use last-publication-wins at the atomic pointer replacement, not at start time. Prior run directories remain available without automatic deletion. Attempt directories are removed on normal success or exception, but a killed process can leave temporary directories. + +The guarantee covers complete-set visibility on a local filesystem with atomic same-filesystem rename and replacement. It is not a power-loss durability guarantee, an object-storage commit protocol, or a guarantee for network filesystems. Files remain immutable through Phlo's API, not through filesystem permissions. External edits can corrupt them, and `verify_manifest_checksums` detects those edits. Workers must close writers before returning and must not mutate files in background tasks. + ## phlo.ingest.dlt `phlo.ingest.dlt` resolves `phlo_dlt.phlo_ingestion`. diff --git a/src/phlo/__init__.py b/src/phlo/__init__.py index 8a6e067adf8..053a9b46d3e 100644 --- a/src/phlo/__init__.py +++ b/src/phlo/__init__.py @@ -144,6 +144,7 @@ def validate_users(): __all__ = [ "__version__", + "export", *_SUBMODULE_EXPORTS, *_HELPER_EXPORTS, *_CONTRACT_EXPORTS, @@ -163,6 +164,11 @@ def __getattr__(name: str) -> Any: """ + if name == "export": + from phlo.exports import export + + globals()[name] = export + return export if name in _SUBMODULE_EXPORTS: module = import_module(f"{__name__}.{name}") globals()[name] = module diff --git a/src/phlo/exports.py b/src/phlo/exports.py new file mode 100644 index 00000000000..b782a7487ef --- /dev/null +++ b/src/phlo/exports.py @@ -0,0 +1,213 @@ +"""Execute file export assets and publish complete sets on a local filesystem. + +Each attempt stages independently. Complete run directories are committed +before atomic current-pointer replacement, and published runs are never +overwritten by Phlo. Consumers resolve one manifest for the whole set. +""" + +from __future__ import annotations + +import errno +import hashlib +import json +import shutil +from collections.abc import Callable, Iterable, Mapping, Sequence +from dataclasses import asdict, dataclass, field, replace +from pathlib import Path +from tempfile import TemporaryDirectory + +from phlo.capabilities.registry import register_capability +from phlo.capabilities.runtime import RuntimeContext, routing_from_context +from phlo.capabilities.specs import AssetSpec, MaterializeResult, RunSpec +from phlo.helpers.artifacts import ( + ArtifactEntry, + ArtifactManifest, + manifest_from_paths, + verify_manifest_checksums, +) + + +@dataclass(frozen=True, slots=True) +class ExportContext: + """One writer invocation and its selected upstream versions. + + Write every declared output beneath ``staging_dir`` and close files before + returning. Record available version/snapshot references in + ``upstream_versions`` when selecting the data, not by looking up a later + latest version. ``runtime`` exposes resources, routing, and logging. + """ + + staging_dir: Path + run_id: str + runtime: RuntimeContext + upstream_versions: dict[str, str] = field(default_factory=dict) + + +def _export_root(destination: str | Path, name: str) -> Path: + if not name or name in {".", ".."} or "/" in name or "\\" in name: + raise ValueError("Export name must be a non-empty single path component") + if "://" in str(destination): + raise ValueError("Exports support local filesystem destinations only") + return Path(destination).resolve() / name + + +def _read_manifest(path: Path) -> ArtifactManifest: + data = json.loads(path.read_text(encoding="utf-8")) + return ArtifactManifest( + name=data["name"], + artifacts=[ArtifactEntry(**entry) for entry in data["artifacts"]], + metadata=data["metadata"], + ) + + +def resolve_export_manifest(destination: str | Path, name: str) -> ArtifactManifest: + """Resolve the current pointer once and read its immutable manifest. + + Reuse the returned manifest for every artifact read. Resolving separately + per artifact could mix runs if a concurrent publication changes current. + Raises FileNotFoundError when no complete set has been published. + """ + root = _export_root(destination, name) + pointer = json.loads((root / "current.json").read_text(encoding="utf-8")) + return _read_manifest(root / pointer["manifest"]) + + +def _publish( + root: Path, + context: ExportContext, + outputs: Mapping[str, Path], + candidate: Path, +) -> tuple[ArtifactManifest, Path]: + """Copy declared files, commit an immutable run, then replace current.""" + candidate.mkdir() + for relative_path in outputs.values(): + source = context.staging_dir / relative_path + if not source.is_file(): + raise FileNotFoundError(f"Missing export output: {relative_path}") + if any(part.is_symlink() for part in (source, *source.parents)): + raise ValueError(f"Export output must not use symlinks: {relative_path}") + target = candidate / relative_path + target.parent.mkdir(parents=True, exist_ok=True) + shutil.copyfile(source, target) + + # Hash the opaque identity rather than interpreting it as a filesystem path. + run_dir = root / "runs" / hashlib.sha256(context.run_id.encode()).hexdigest() + manifest_path = run_dir / "manifest.json" + staging_manifest = manifest_from_paths( + root.name, [candidate / path for path in outputs.values()] + ) + routing = routing_from_context(context.runtime) + manifest = ArtifactManifest( + name=root.name, + artifacts=[ + replace(entry, uri=str(run_dir / path), metadata={"name": name}) + for (name, path), entry in zip(outputs.items(), staging_manifest.artifacts, strict=True) + ], + metadata={ + "run_id": context.run_id, + "partition_key": routing.partition_key, + "ref": routing.ref, + "upstream_versions": dict(context.upstream_versions), + }, + ) + (candidate / "manifest.json").write_text( + json.dumps(asdict(manifest), sort_keys=True), encoding="utf-8" + ) + run_dir.parent.mkdir(parents=True, exist_ok=True) + try: + candidate.rename(run_dir) + except OSError as exc: + # Another attempt may have committed this identity concurrently. A + # complete run directory is non-empty, so rename cannot replace it. + if exc.errno not in {errno.EEXIST, errno.ENOTEMPTY}: + raise + existing = _read_manifest(manifest_path) + if existing != manifest or not all(verify_manifest_checksums(existing).values()): + raise ValueError( + f"Export run {context.run_id!r} already has different content" + ) from exc + + pointer = candidate.parent / "current.json" + pointer.write_text( + json.dumps({"manifest": str(manifest_path.relative_to(root))}), encoding="utf-8" + ) + pointer.replace(root / "current.json") + return manifest, manifest_path + + +def export( + *, + name: str, + destination: str | Path, + outputs: Mapping[str, str], + depends_on: Sequence[str] = (), + group: str = "exports", + resources: Iterable[str] = (), + max_retries: int = 0, + retry_delay_seconds: int = 30, +) -> Callable[[Callable[[ExportContext], None]], Callable[[ExportContext], None]]: + """Register one executable asset for a complete named file export set. + + ``destination/name`` owns immutable run directories and ``current.json``. + Retries use canonical runtime run identity, fresh staging, and reject + changes to an already published run. Concurrent publications of distinct + runs are last-publication-wins. Old runs are never automatically deleted. + """ + root = _export_root(destination, name) + dependencies = tuple(depends_on) + paths = {key: Path(value) for key, value in outputs.items()} + if not paths or any(not key for key in paths): + raise ValueError("Exports require at least one named output") + for path in paths.values(): + if path.is_absolute() or ".." in path.parts or path == Path(): + raise ValueError(f"Export output must be a relative file path: {path}") + if path.parts[0] == "manifest.json": + raise ValueError("manifest.json is reserved for the export manifest") + if len(set(paths.values())) != len(paths): + raise ValueError("Export outputs must have distinct paths") + + def decorate(writer: Callable[[ExportContext], None]) -> Callable[[ExportContext], None]: + def run(runtime: RuntimeContext) -> Iterable[MaterializeResult]: + routing = routing_from_context(runtime) + if not routing.run_id: + raise ValueError("Exports require a stable runtime run_id") + root.mkdir(parents=True, exist_ok=True) + with TemporaryDirectory(prefix=".attempt-", dir=root) as attempt: + staging = Path(attempt) / "staging" + staging.mkdir() + context = ExportContext(staging, routing.run_id, runtime) + writer(context) + unknown = context.upstream_versions.keys() - set(dependencies) + if unknown: + raise ValueError(f"Undeclared export upstream versions: {sorted(unknown)}") + manifest, manifest_path = _publish(root, context, paths, Path(attempt) / "set") + yield MaterializeResult( + metadata={ + "phlo/export_manifest": asdict(manifest), + "phlo/export_manifest_path": str(manifest_path), + "phlo/export_current_path": str(root / "current.json"), + "phlo/export_run_id": routing.run_id, + } + ) + + register_capability( + "asset", + AssetSpec( + key=name, + group=group, + description=writer.__doc__, + kinds={"file"}, + deps=list(dependencies), + resources=set(resources), + metadata={ + "phlo/export_outputs": {key: str(path) for key, path in paths.items()}, + "phlo/export_destination": str(root), + }, + run=RunSpec( + fn=run, max_retries=max_retries, retry_delay_seconds=retry_delay_seconds + ), + ), + ) + return writer + + return decorate diff --git a/tests/runtime/test_exports.py b/tests/runtime/test_exports.py new file mode 100644 index 00000000000..dfb8e57ca43 --- /dev/null +++ b/tests/runtime/test_exports.py @@ -0,0 +1,234 @@ +"""File export contracts through executable capabilities and local manifests.""" + +import hashlib +from concurrent.futures import ThreadPoolExecutor +from pathlib import Path +from threading import Barrier, Event +from types import SimpleNamespace + +import pytest + +import phlo +from phlo.capabilities.registry import clear_capabilities, get_capability_registry +from phlo.capabilities.runtime import RuntimeRouting +from phlo.exports import ExportContext, resolve_export_manifest +from phlo.helpers.artifacts import verify_manifest_checksums + + +@pytest.fixture(autouse=True) +def isolated_assets(): + previous = get_capability_registry().list("asset") + clear_capabilities("asset") + yield + clear_capabilities("asset") + for spec in previous: + get_capability_registry().register("asset", spec) + + +def execute(run_id="run-1", **routing): + spec = get_capability_registry().list("asset")[0] + runtime = SimpleNamespace( + routing=RuntimeRouting(run_id=run_id, **routing), + run_id=run_id, + resources={}, + tags={}, + partition_key=None, + ) + return list(spec.run.fn(runtime))[0] + + +def test_complete_export_has_named_checksummed_artifacts(tmp_path): + @phlo.export( + name="sky_survey", + destination=tmp_path, + depends_on=["observed_stars"], + outputs={"catalog": "catalog.csv", "summary": "reports/summary.json"}, + ) + def write(context: ExportContext): + assert context.run_id == "run-1" + assert context.runtime.routing.ref == "survey-branch" + context.upstream_versions["observed_stars"] = "snapshot-42" + (context.staging_dir / "catalog.csv").write_bytes(b"star,x\nSirius,17\n") + (context.staging_dir / "reports").mkdir() + (context.staging_dir / "reports/summary.json").write_bytes(b'{"count":1}') + + result = execute(ref="survey-branch", partition_key="2026-10") + spec = get_capability_registry().list("asset")[0] + assert spec.key == "sky_survey" + assert spec.deps == ["observed_stars"] + manifest = resolve_export_manifest(tmp_path, "sky_survey") + assert manifest.metadata["run_id"] == "run-1" + assert manifest.metadata["partition_key"] == "2026-10" + assert manifest.metadata["upstream_versions"] == {"observed_stars": "snapshot-42"} + catalog, summary = manifest.artifacts + assert catalog.metadata["name"] == "catalog" + assert summary.metadata["name"] == "summary" + assert catalog.size_bytes == 17 + assert summary.checksum == hashlib.sha256(b'{"count":1}').hexdigest() + assert all(verify_manifest_checksums(manifest).values()) + assert result.metadata["phlo/export_manifest"]["metadata"] == manifest.metadata + assert Path(result.metadata["phlo/export_manifest_path"]).is_file() + + +@pytest.mark.parametrize("failure", ["writer", "missing"]) +def test_incomplete_set_preserves_current_and_retry_starts_empty(tmp_path, failure): + attempts = [] + fail = False + + @phlo.export( + name="survey", + destination=tmp_path, + outputs={"catalog": "catalog.csv", "summary": "summary.json"}, + ) + def write(context): + attempts.append((context.run_id, context.staging_dir)) + assert list(context.staging_dir.iterdir()) == [] + (context.staging_dir / "catalog.csv").write_text(context.run_id, encoding="utf-8") + if fail: + (context.staging_dir / "scratch.txt").write_text("partial", encoding="utf-8") + if failure == "writer": + raise RuntimeError("summary writer failed") + return + (context.staging_dir / "summary.json").write_text(context.run_id, encoding="utf-8") + + execute("prior") + prior = resolve_export_manifest(tmp_path, "survey") + pointer = (tmp_path / "survey/current.json").read_bytes() + fail = True + with pytest.raises((RuntimeError, FileNotFoundError)): + execute("retry-run") + assert (tmp_path / "survey/current.json").read_bytes() == pointer + assert resolve_export_manifest(tmp_path, "survey") == prior + assert not list((tmp_path / "survey").glob(".attempt-*")) + fail = False + execute("retry-run", attempt=2) + assert attempts[-1][0] == attempts[-2][0] == "retry-run" + assert attempts[-1][1] != attempts[-2][1] + current = resolve_export_manifest(tmp_path, "survey") + assert current.metadata["run_id"] == "retry-run" + assert [Path(entry.uri).read_text() for entry in current.artifacts] == ["retry-run"] * 2 + assert all(verify_manifest_checksums(prior).values()) + assert len(list((tmp_path / "survey/runs").iterdir())) == 2 + + +def test_published_identity_is_idempotent_but_cannot_change(tmp_path): + content = "original" + + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + (context.staging_dir / "data.txt").write_text(content, encoding="utf-8") + + first = execute() + first_manifest = resolve_export_manifest(tmp_path, "survey") + assert execute(attempt=2).metadata == first.metadata + content = "changed" + with pytest.raises(ValueError, match="already has different content"): + execute(attempt=3) + assert resolve_export_manifest(tmp_path, "survey") == first_manifest + assert Path(first_manifest.artifacts[0].uri).read_text() == "original" + assert len(list((tmp_path / "survey/runs").iterdir())) == 1 + # Detect external corruption as well as a writer changing its output. + Path(first_manifest.artifacts[0].uri).write_text("tampered", encoding="utf-8") + assert verify_manifest_checksums(first_manifest) == {first_manifest.artifacts[0].uri: False} + content = "original" + with pytest.raises(ValueError, match="already has different content"): + execute() + + +def test_concurrent_runs_publish_in_completion_order_and_readers_keep_one_set(tmp_path): + slow_started = Event() + release_slow = Event() + + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + if context.run_id == "slow": + slow_started.set() + assert release_slow.wait(10) + (context.staging_dir / "data.txt").write_text(context.run_id, encoding="utf-8") + + with ThreadPoolExecutor() as pool: + slow = pool.submit(execute, "slow") + assert slow_started.wait(10) + pool.submit(execute, "fast").result(timeout=10) + fast_manifest = resolve_export_manifest(tmp_path, "survey") + release_slow.set() + slow.result(timeout=10) + assert resolve_export_manifest(tmp_path, "survey").metadata["run_id"] == "slow" + assert fast_manifest.metadata["run_id"] == "fast" + assert Path(fast_manifest.artifacts[0].uri).read_text() == "fast" + + +@pytest.mark.parametrize("path", ["../escape", "/absolute", "manifest.json", "", "."]) +def test_output_paths_cannot_escape_or_replace_manifest(tmp_path, path): + with pytest.raises(ValueError): + phlo.export(name="survey", destination=tmp_path, outputs={"data": path}) + + +def test_symlink_outputs_are_not_published(tmp_path): + outside = tmp_path / "outside" + outside.write_text("mutable", encoding="utf-8") + + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + (context.staging_dir / "data.txt").symlink_to(outside) + + with pytest.raises(ValueError, match="symlink"): + execute() + with pytest.raises(FileNotFoundError): + resolve_export_manifest(tmp_path, "survey") + + +def test_concurrent_attempts_cannot_overwrite_same_identity(tmp_path): + rendezvous = Barrier(2) + + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + (context.staging_dir / "data.txt").write_text( + str(context.runtime.routing.attempt), encoding="utf-8" + ) + rendezvous.wait(timeout=10) + + with ThreadPoolExecutor() as pool: + futures = [pool.submit(execute, "same-id", attempt=index) for index in (1, 2)] + errors = [future.exception(timeout=10) for future in futures] + assert sum(error is None for error in errors) == 1 + assert sum(isinstance(error, ValueError) for error in errors) == 1 + winner = futures[errors.index(None)].result() + manifest = resolve_export_manifest(tmp_path, "survey") + assert winner.metadata["phlo/export_manifest_path"] == str( + Path(manifest.artifacts[0].uri).parent / "manifest.json" + ) + assert all(verify_manifest_checksums(manifest).values()) + assert len(list((tmp_path / "survey/runs").iterdir())) == 1 + + +def test_pointer_failure_preserves_previous_set_and_can_retry(tmp_path, monkeypatch): + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + (context.staging_dir / "data.txt").write_text(context.run_id, encoding="utf-8") + + execute("prior") + prior = resolve_export_manifest(tmp_path, "survey") + with monkeypatch.context() as patch: + + def fail_replace(self, target): + raise OSError("publication unavailable") + + patch.setattr(Path, "replace", fail_replace) + with pytest.raises(OSError, match="publication unavailable"): + execute("next") + assert resolve_export_manifest(tmp_path, "survey") == prior + assert len(list((tmp_path / "survey/runs").iterdir())) == 2 + execute("next") + assert resolve_export_manifest(tmp_path, "survey").metadata["run_id"] == "next" + assert all(verify_manifest_checksums(prior).values()) + + +def test_missing_identity_fails_before_writer(tmp_path): + @phlo.export(name="survey", destination=tmp_path, outputs={"data": "data.txt"}) + def write(context): + pytest.fail("writer must not run without identity") + + with pytest.raises(ValueError, match="stable runtime run_id"): + execute(None) + assert not (tmp_path / "survey").exists() diff --git a/tests/runtime/test_exports_integration.py b/tests/runtime/test_exports_integration.py new file mode 100644 index 00000000000..5997e37df2d --- /dev/null +++ b/tests/runtime/test_exports_integration.py @@ -0,0 +1,113 @@ +"""Exercise exports through the supported workflow discovery and Dagster adapter.""" + +import json +import sys +from pathlib import Path + +import dagster as dg +import pytest +from phlo_dagster.framework.discovery import discover_user_workflows + +from phlo.capabilities.registry import CAPABILITY_FAMILIES, get_capability_registry +from phlo.exports import resolve_export_manifest +from phlo.helpers.artifacts import verify_manifest_checksums + + +@pytest.fixture(autouse=True) +def restore_discovery_state(): + registry = get_capability_registry() + capabilities = {family: registry.list(family) for family in CAPABILITY_FAMILIES} + modules = {name: module for name, module in sys.modules.items() if name.startswith("workflows")} + search_path = sys.path[:] + yield + registry.clear_all() + for family, specs in capabilities.items(): + for spec in specs: + registry.register(family, spec) + for name in list(sys.modules): + if name.startswith("workflows"): + sys.modules.pop(name) + sys.modules.update(modules) + sys.path[:] = search_path + + +@pytest.mark.parametrize("failure", ["writer", "missing"]) +def test_discovered_export_failure_and_retry(tmp_path, failure): + workflows = tmp_path / "workflows" + workflows.mkdir() + destination = tmp_path / "exports" + attempts_path = tmp_path / "attempts.jsonl" + (workflows / "survey.py").write_text( + f""" +import json +from pathlib import Path +import phlo +from dagster import asset + +@asset +def observed_stars(): + return None + +@phlo.export( + name="survey", destination={str(destination)!r}, + depends_on=["observed_stars"], + outputs={{"catalog": "catalog.csv", "summary": "summary.json"}}, + max_retries=1, retry_delay_seconds=1, +) +def survey(context): + assert list(context.staging_dir.iterdir()) == [] + with Path({str(attempts_path)!r}).open("a", encoding="utf-8") as log: + log.write(json.dumps({{"run_id": context.run_id, "staging": str(context.staging_dir)}}) + "\\n") + context.upstream_versions["observed_stars"] = "snapshot-42" + (context.staging_dir / "catalog.csv").write_text(context.run_id, encoding="utf-8") + mode = context.runtime.tags.get("test/mode", "success") + attempts = Path({str(attempts_path)!r}).read_text(encoding="utf-8").splitlines() + same_run_attempts = sum(json.loads(line)["run_id"] == context.run_id for line in attempts) + if mode == "fail" or (mode == "retry" and same_run_attempts == 1): + if {failure!r} == "writer": + raise RuntimeError("summary writer failed") + return + (context.staging_dir / "summary.json").write_text(context.run_id, encoding="utf-8") +""", + encoding="utf-8", + ) + definitions = discover_user_workflows(workflows, clear_registries=True) + graph = definitions.resolve_asset_graph() + assert dg.AssetKey("observed_stars") in graph.get(dg.AssetKey("survey")).parent_keys + assets = [ + asset + for asset in definitions.assets + if isinstance(asset, dg.AssetsDefinition) + and asset.keys & {dg.AssetKey("survey"), dg.AssetKey("observed_stars")} + ] + with dg.DagsterInstance.ephemeral() as instance: + prior_result = dg.materialize(assets, instance=instance) + assert prior_result.success + prior = resolve_export_manifest(destination, "survey") + failed = dg.materialize( + assets, instance=instance, tags={"test/mode": "fail"}, raise_on_error=False + ) + assert not failed.success + assert failed.get_asset_materialization_events() # upstream still succeeded + assert not failed.asset_materializations_for_node("survey") + assert resolve_export_manifest(destination, "survey") == prior + result = dg.materialize(assets, instance=instance, tags={"test/mode": "retry"}) + assert result.success + manifest = resolve_export_manifest(destination, "survey") + assert manifest.metadata["run_id"] == result.run_id + assert manifest.metadata["upstream_versions"] == {"observed_stars": "snapshot-42"} + assert all(verify_manifest_checksums(manifest).values()) + assert [Path(entry.uri).read_text() for entry in manifest.artifacts] == [result.run_id] * 2 + event = result.asset_materializations_for_node("survey")[0] + assert event.metadata["phlo/export_manifest"].value["metadata"] == manifest.metadata + assert Path(event.metadata["phlo/export_manifest_path"].value).is_file() + attempts = [json.loads(line) for line in attempts_path.read_text().splitlines()] + retried = [attempt for attempt in attempts if attempt["run_id"] == result.run_id] + assert len(retried) == 2 + assert retried[0]["staging"] != retried[1]["staging"] + assert all(not Path(attempt["staging"]).exists() for attempt in attempts) + assert all(verify_manifest_checksums(prior).values()) + assert len(list((destination / "survey/runs").iterdir())) == 2 + # Reloading imports registers the executable spec again after clearing. + reloaded = discover_user_workflows(workflows, clear_registries=True) + assert dg.AssetKey("survey") in reloaded.resolve_asset_graph().get_all_asset_keys() From 24fcdcbf094042e925948b703ae55568d9775324 Mon Sep 17 00:00:00 2001 From: Gareth Price Date: Thu, 8 Oct 2026 07:28:02 +0000 Subject: [PATCH 2/2] docs: include file exports in guide navigation --- docs/guides/meta.json | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/guides/meta.json b/docs/guides/meta.json index d8c3c44bf9c..4e73664157e 100644 --- a/docs/guides/meta.json +++ b/docs/guides/meta.json @@ -8,6 +8,7 @@ "evolve-a-schema", "choose-your-stack", "expose-data", + "export-files", "secure-the-stack", "monitor-and-debug", "troubleshoot-errors",