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
92 changes: 92 additions & 0 deletions docs/guides/export-files.md
Original file line number Diff line number Diff line change
@@ -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).
1 change: 1 addition & 0 deletions docs/guides/meta.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
"evolve-a-schema",
"choose-your-stack",
"expose-data",
"export-files",
"secure-the-stack",
"monitor-and-debug",
"troubleshoot-errors",
Expand Down
1 change: 1 addition & 0 deletions docs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -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) |
Expand Down
35 changes: 35 additions & 0 deletions docs/reference/python-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<sha256-of-run-id>` 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`.
Expand Down
6 changes: 6 additions & 0 deletions src/phlo/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,7 @@ def validate_users():

__all__ = [
"__version__",
"export",
*_SUBMODULE_EXPORTS,
*_HELPER_EXPORTS,
*_CONTRACT_EXPORTS,
Expand All @@ -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
Expand Down
213 changes: 213 additions & 0 deletions src/phlo/exports.py
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading