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
4 changes: 4 additions & 0 deletions development_docs/architecture/execution-and-caching.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,10 @@ unchanged.
For an owned export graph, the adapter can:

- add one frozen producer-dependency input to native cell hashing
- retry a failed native hash with a digest of the type and Arrow IPC stream of
each referenced PyArrow value, because Marimo's NumPy view rejects
object-typed columns such as strings. Cells whose native hash succeeds keep
their native keys.
- record the effective hit or miss for authored and projection cells
- run complete-cell owners and selected exporter leaves live when their output
contract includes uncached side effects or session-bound resources
Expand Down
2 changes: 1 addition & 1 deletion docs/guide/choose-states.md
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ Each output has one source kind:
| -------------- | --------------------------------------------------------------------- | ---------------------------------------------- |
| `kind: json` | Canonical portable JSON | Python, browser, agent, or custom client |
| `kind: native` | marimo scalar, JSON, NumPy, Arrow, or `BlobAsset` representation | Typed Python or browser loader |
| `kind: export` | `BlobAsset` returned by one declared exporter | Chart, table, media, or domain-specific loader |
| `kind: export` | `BlobAsset` or canonical JSON returned by one declared exporter | Chart, table, media, or domain-specific loader |
| `kind: output` | Formatted marimo output and replay resources | marimo-aware browser application |
| `kind: cell` | Cell identity, terminal output, console records, and replay resources | marimo-aware browser application or agent |

Expand Down
3 changes: 2 additions & 1 deletion docs/reference/export-spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -184,7 +184,8 @@ exporter: altair.vegalite
```

An export source passes the selected value to one declared exporter. The
exporter returns a `BlobAsset` with bytes, media type, filename, and metadata.
exporter returns a `BlobAsset` with bytes, media type, filename, and metadata,
or a JSON value that the export stores as canonical `marimo.json.v1` JSON.

### Rendered-output source

Expand Down
18 changes: 16 additions & 2 deletions docs/reference/python/produce.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,7 @@ OutputSpec.cell(name: str | None = None, *, id: str | None = None) -> OutputSpec
| ---------- | ---------------------------------------------------------------------------------------------- |
| `json()` | Canonical portable JSON selected from a notebook definition |
| `native()` | marimo cache representation for a scalar, JSON value, NumPy array, Arrow table, or `BlobAsset` |
| `export()` | `BlobAsset` returned by an explicit exporter |
| `export()` | `BlobAsset` or canonical JSON returned by an explicit exporter |
| `output()` | Formatted `marimo.output.v1` snapshot and replay resources |
| `cell()` | Complete `marimo.cell.v1` snapshot selected by authored cell name or inspected runtime ID |

Expand Down Expand Up @@ -222,7 +222,7 @@ modules, so restart it after changing custom exporter source.

### `BlobAsset`

Custom exporters return `marimo_export.outputs.BlobAsset`:
Custom exporters return `marimo_export.outputs.BlobAsset` or a JSON value:

```python
from marimo_export.outputs import BlobAsset
Expand Down Expand Up @@ -253,6 +253,20 @@ canonical encoding is limited to 256 KiB. Supply `media_type` for a value that
will enter a notebook export. Export production rejects a `BlobAsset` whose
media type is absent or invalid.

An exporter that returns a JSON value produces the canonical `marimo.json.v1`
representation of `OutputSpec.json()`, and readers return it through
`output.json()`:

```python
def summarize(value) -> dict[str, object]:
return {"rows": value.num_rows, "columns": value.column_names}
```

One output keeps one codec across states, so an exporter returns a `BlobAsset`
for every state or a JSON value for every state. Any other result raises
`OutputError` with code `output_execution_failed` and `exception_type`
`TypeError`.

## `plan()`

```python
Expand Down
2 changes: 1 addition & 1 deletion docs/reference/terminology.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ The concept pages introduce them through worked examples.
| Output | One published name and representation available for every exported state. |
| Output source | The `json`, `native`, `export`, `output`, or `cell` selection declared by an output spec. |
| Selector | A path from one Python definition through supported attribute or item steps to a selected notebook result. |
| Exporter | A producer-side converter that returns a `BlobAsset` for one selected value. |
| Exporter | A producer-side converter that returns a `BlobAsset` or a JSON value for one selected value. |
| Output plan | The complete set of authored output declarations. Its identity changes when an output source, exporter, option, or declared dependency changes. |
| Output representation | The codec and media type that define how one output is stored and decoded. One output name keeps the same representation across every state. |
| Codec | A versioned identifier for the native storage envelope, such as `marimo.json.v1` or `numpy.npy.v1`. |
Expand Down
4 changes: 1 addition & 3 deletions packages/python/src/marimo_export/_marimo/blob.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,8 @@

from __future__ import annotations

from marimo_export.outputs import BlobAsset


def to_native_blob_asset(value: BlobAsset) -> object:
def to_native_blob_asset(value: object) -> object:
from marimo_export._marimo.compat.blob import to_native_blob_asset as convert

return convert(value)
Expand Down
25 changes: 20 additions & 5 deletions packages/python/src/marimo_export/_marimo/compat/blob.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,32 @@

from __future__ import annotations

from marimo_export._json import portable_json_object
from marimo_export._json import canonical_bytes, json_value, portable_json_object
from marimo_export.descriptors import JSON_CODEC, JSON_MEDIA_TYPE
from marimo_export.outputs import BlobAsset


def to_native_blob_asset(value: BlobAsset) -> object:
"""Return the native value required by Marimo's lazy ``.bin`` codec."""
def to_native_blob_asset(value: object) -> object:
"""Return the native value required by Marimo's lazy ``.bin`` codec.
An exporter returns a ``BlobAsset`` or a JSON value. A JSON value uses the
canonical JSON representation of a JSON source.
"""

if not isinstance(value, BlobAsset):
raise TypeError("output exporter must return marimo_export.outputs.BlobAsset")
from marimo._save.stubs import BlobAsset as NativeBlobAsset

if not isinstance(value, BlobAsset):
try:
portable = json_value(value, "exporter result")
except (TypeError, ValueError) as error:
raise TypeError(
"output exporter must return marimo_export.outputs.BlobAsset or a JSON value"
) from error
return NativeBlobAsset(
data=canonical_bytes(portable),
media_type=JSON_MEDIA_TYPE,
metadata={"schema": JSON_CODEC},
)
return NativeBlobAsset(
data=value.data,
media_type=value.media_type,
Expand Down
72 changes: 64 additions & 8 deletions packages/python/src/marimo_export/_marimo/compat/cache/attempts.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

import ast
import copy
import hashlib
import threading
from collections.abc import Callable, Iterator
from contextlib import contextmanager
Expand Down Expand Up @@ -132,17 +133,17 @@ def tracked(
) -> Cache:
with _SCOPES_LOCK:
scope = _SCOPES.get(id(graph))
tracked_graph = scope is not None and scope.graph is graph
environment = scope.environment if scope is not None and scope.graph is graph else None
if environment is not None:
module, scope_values = _with_environment(module, scope_values, environment)
attempt = native(
module,
graph,
cell_id,
scope_values,
*args,
**kwargs,
)
try:
attempt = native(module, graph, cell_id, scope_values, *args, **kwargs)
except TypeError:
Comment thread
peter-gy marked this conversation as resolved.
hashable = _with_arrow_digests(graph, cell_id, scope_values) if tracked_graph else None
if hashable is None:
raise
attempt = native(module, graph, cell_id, hashable, *args, **kwargs)
with _SCOPES_LOCK:
scope = _SCOPES.get(id(graph))
if scope is None or scope.graph is not graph:
Expand Down Expand Up @@ -194,6 +195,61 @@ def _with_environment(
return ast.fix_missing_locations(lookup), {**scope, name: environment}


def _with_arrow_digests(
graph: Any,
cell_id: Any,
scope: dict[str, Any],
) -> dict[str, Any] | None:
# Marimo hashes Arrow data through NumPy, which cannot view object-typed
# columns such as strings as bytes. After that failure the native hasher
# receives a digest of each referenced Arrow value's type and IPC stream,
# and execution keeps the value. Cells whose native hash succeeds keep
# their native keys.
cell = graph.cells.get(cell_id)
if cell is None:
return None
digests = {
name: digest
for name in cell.refs
if name in scope and (digest := _arrow_digest(scope[name])) is not None
}
return {**scope, **digests} if digests else None


def _arrow_digest(value: object) -> str | None:
if not type(value).__module__.startswith("pyarrow"):
return None
import pyarrow as pa

kind = type(value).__qualname__
if isinstance(value, (pa.Array, pa.ChunkedArray)):
value = pa.Table.from_arrays([value], names=["value"])
elif not isinstance(value, (pa.Table, pa.RecordBatch)):
return None
sink = _DigestSink()
with pa.ipc.new_stream(sink, value.schema) as writer:
writer.write(value)
return f"arrow-ipc-sha256:{kind}:{sink.digest.hexdigest()}"


class _DigestSink:
"""A write-only file that hashes the Arrow IPC stream written to it."""

def __init__(self) -> None:
self.digest = hashlib.sha256()
self.closed = False

def write(self, data: bytes) -> int:
self.digest.update(data)
return len(data)

def flush(self) -> None:
pass

def close(self) -> None:
self.closed = True


def record_cache_miss(graph: Any, cell_id: Any) -> None:
"""Record a live run chosen after native restoration."""

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import time
from typing import Any

from marimo._runtime.exceptions import MarimoRuntimeException
from marimo._runtime.executor.lifecycles import Skip
from marimo._runtime.executor.lifecycles.cached import CachedLifecycle

Expand All @@ -20,7 +21,15 @@ class CompleteCachedLifecycle(CachedLifecycle):
"""Rerun a hit when its restored values cannot serve the live session."""

def setup(self, cell: Any, glbls: Any) -> Any:
decision = super().setup(cell, glbls)
try:
decision = super().setup(cell, glbls)
except Exception as error:
if not has_cache_scope(self._graph):
raise
# Marimo records an exception that escapes a lifecycle as an
# untyped error. A runtime exception keeps the exception on the
# cell, so export failures name its type.
raise MarimoRuntimeException from error
if not has_cache_scope(self._graph):
return decision
if not isinstance(decision, Skip):
Expand Down
20 changes: 15 additions & 5 deletions packages/python/src/marimo_export/_marimo/compat/projections.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ class ProjectionRecording:
"polars.series.series.Series",
}
)
_PYARROW_TYPES = frozenset({"pyarrow.lib.RecordBatch", "pyarrow.lib.Table"})


class _RecordingPipe:
Expand Down Expand Up @@ -220,16 +221,25 @@ def capture_native_value(

def _native_arrow_value(value: object) -> object | None:
python_type = f"{type(value).__module__}.{type(value).__qualname__}"
if python_type not in _POLARS_TYPES:
if python_type in _POLARS_TYPES:
frame = cast(Any, value).to_frame() if python_type.endswith(".Series") else value
buffer = BytesIO()
cast(Any, frame).write_ipc_stream(buffer, compression="uncompressed")
data = buffer.getvalue()
elif python_type in _PYARROW_TYPES:
import pyarrow as pa

sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, cast(Any, value).schema) as writer:
writer.write(value)
data = sink.getvalue().to_pybytes()
else:
return None
frame = cast(Any, value).to_frame() if python_type.endswith(".Series") else value
buffer = BytesIO()
cast(Any, frame).write_ipc_stream(buffer, compression="uncompressed")

from marimo._save.stubs import BlobAsset

return BlobAsset(
data=buffer.getvalue(),
data=data,
media_type=ARROW_MEDIA_TYPE,
metadata={
"python_type": python_type,
Expand Down
18 changes: 16 additions & 2 deletions packages/python/src/marimo_export/_marimo/compat/receipts.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
from marimo_export._marimo.compat.projections import _NATIVE_ARROW_SCHEMA
from marimo_export.descriptors import (
ARROW_MEDIA_TYPE,
JSON_CODEC,
JSON_MEDIA_TYPE,
MARIMO_CELL_MEDIA_TYPE,
MARIMO_OUTPUT_MEDIA_TYPE,
Expand All @@ -38,7 +39,13 @@
)
from marimo_export.errors import CodecError, OutputError
from marimo_export.outputs import BlobAsset
from marimo_export.spec import CellSource, JsonSource, NativeSource, RenderedOutputSource
from marimo_export.spec import (
CellSource,
ExportSource,
JsonSource,
NativeSource,
RenderedOutputSource,
)

_BLOB_ASSET_PYTHON_TYPE = f"{BlobAsset.__module__}.{BlobAsset.__qualname__}"

Expand Down Expand Up @@ -156,7 +163,14 @@ def native_receipt(
payload=payload,
disposition=disposition,
)
if isinstance(source, (JsonSource, NativeSource)) and cached.media_type == JSON_MEDIA_TYPE:
if cached.media_type == JSON_MEDIA_TYPE and (
isinstance(source, (JsonSource, NativeSource))
or (
isinstance(source, ExportSource)
Comment thread
peter-gy marked this conversation as resolved.
and cached.metadata == {"schema": JSON_CODEC}
and cached.filename is None
)
):
try:
value = decode_json(data, f"output {output!r} JSON projection")
except (TypeError, ValueError) as error:
Expand Down
48 changes: 48 additions & 0 deletions packages/python/tests/test_custom_exporter_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,54 @@ def test_capture_sideloads_an_importable_callable(
assert notebook.read_bytes() == source


def test_custom_exporter_returns_a_json_value(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
notebook = tmp_path / "notebook.py"
_write_notebook(notebook)
(tmp_path / "export_exports.py").write_text(
"def describe(value):\n return {'answer': value, 'labels': ['forty', 'one']}\n",
encoding="utf-8",
)
monkeypatch.setenv("PYTHONPATH", str(tmp_path))
spec = ExportSpec(
default_state="baseline",
states={"baseline": {}},
outputs={"summary": OutputSpec.export("answer", importable("export_exports:describe"))},
)

_capture(notebook, spec, tmp_path / "export")
output = open_export(tmp_path / "export").state("baseline").output("summary")

assert output.descriptor.codec == "marimo.json.v1"
assert output.json() == {"answer": 41, "labels": ("forty", "one")}


def test_custom_exporter_rejects_a_result_without_a_portable_form(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
notebook = tmp_path / "notebook.py"
_write_notebook(notebook)
(tmp_path / "export_exports.py").write_text(
"def describe(value):\n return object()\n",
encoding="utf-8",
)
monkeypatch.setenv("PYTHONPATH", str(tmp_path))
spec = ExportSpec(
default_state="baseline",
states={"baseline": {}},
outputs={"summary": OutputSpec.export("answer", importable("export_exports:describe"))},
)

with pytest.raises(OutputError) as raised:
_capture(notebook, spec, tmp_path / "export")

assert raised.value.code == "output_execution_failed"
assert raised.value.details["exception_type"] == "TypeError"


def test_custom_exporter_builds_are_deterministic_while_both_execute(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
Expand Down
Loading
Loading