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
2 changes: 1 addition & 1 deletion .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ jobs:
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
with:
repository: fruwehq/determa-state-conformance
ref: 263644f951f342b0eeaa3aceef4877293d2d7c67
ref: 5ba78c7ef90b8556e76de6481a18b23a3d0c2378
path: .pinned/determa-state-conformance
- name: Check out pinned specification
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
Expand Down
8 changes: 6 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@ a language-agnostic statechart engine with a shared normative conformance suite.

This implementation supports Determa State `format: 1` at specification
commit `cc4b0d734aa1c5953de75fb53b63e390a3b72761`. Correctness is determined by the
111-case core suite, persistence profiles, and 85-vector execution-checkpoint profile
at conformance commit `263644f951f342b0eeaa3aceef4877293d2d7c67`.
111-case core suite, persistence profiles, and 91-vector execution-checkpoint profile
at conformance commit `5ba78c7ef90b8556e76de6481a18b23a3d0c2378`.

Version `0.1.0` is the published synchronized release of the specification, conformance
suite, Python engine, and Rust engine.
Expand Down Expand Up @@ -132,6 +132,10 @@ the exact supplied state object.
semantic validation path. Native values must satisfy the same portable Unicode and
numeric domain as source documents.

Category-specific `StrEnum` definitions are used by production emitters.
`PORTABLE_CODE_SETS` is the immutable category-to-string mapping derived from those
definitions.

## Persist And Migrate

`serialize_aggregate` produces the canonical §16 aggregate artifact. Restoration
Expand Down
220 changes: 220 additions & 0 deletions conformance/execution_checkpoint.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,10 @@
from __future__ import annotations

import copy
import hashlib
import json
from collections.abc import Iterator
from contextlib import contextmanager
from dataclasses import dataclass
from pathlib import Path
from typing import Any
Expand All @@ -18,13 +21,15 @@
ExecutionStore,
ExecutionStoreError,
ExecutionStoreRegistry,
ExecutionStoreTransaction,
MemoryArtifactResolver,
MemoryExecutionStore,
load_bundle,
register_bundled_execution_stores,
restore_execution_checkpoint,
serialize_execution_checkpoint,
)
from determa.state.wire import hash_value

from .harness import conformance_root

Expand Down Expand Up @@ -246,6 +251,54 @@ def health(self) -> dict[str, Any]:
return {"healthy": True}


class _ObservedTransaction(ExecutionStoreTransaction):
def __init__(
self, transaction: ExecutionStoreTransaction, calls: list[str]
) -> None:
self._transaction = transaction
self._calls = calls

@property
def root_instance_id(self) -> str:
return self._transaction.root_instance_id

def load(self) -> bytes | None:
self._calls.append("load_checkpoint")
return self._transaction.load()

def insert(self, checkpoint: bytes) -> bool:
self._calls.append("insert_checkpoint")
return self._transaction.insert(checkpoint)

def replace(
self,
expected_revision: str,
expected_checkpoint_digest: str,
checkpoint: bytes,
) -> bool:
self._calls.append("compare_and_swap_checkpoint")
return self._transaction.replace(
expected_revision, expected_checkpoint_digest, checkpoint
)


class _ObservedMemoryStore(MemoryExecutionStore):
def __init__(self, initial: dict[str, bytes], calls: list[str]) -> None:
super().__init__(initial)
self._calls = calls

@property
def root_instance_ids(self) -> frozenset[str]:
return frozenset(self._records)

@contextmanager
def transaction(
self, root_instance_id: str
) -> Iterator[ExecutionStoreTransaction]:
with super().transaction(root_instance_id) as transaction:
yield _ObservedTransaction(transaction, self._calls)


def _adapter_operation(vector: dict[str, Any]) -> dict[str, Any]:
operation = vector["operation"]
if operation == "inject_execution_store":
Expand Down Expand Up @@ -387,10 +440,177 @@ def _expected_response(
raise AssertionError(f"no exact response projection for {operation}")


def _scope_state(
case: ExecutionCheckpointCase, reference: dict[str, str]
) -> dict[str, Any]:
return copy.deepcopy(
_pointer(_json(case.path / reference["file"]), reference["pointer"])
)


def _outbox_records(checkpoint: dict[str, Any]) -> list[dict[str, str]]:
records = [
("pending", record) for record in checkpoint["pending_outbox_intents"]
]
records.extend(
("terminal", record) for record in checkpoint["terminal_outbox_records"]
)
records.extend(
("tombstone", record) for record in checkpoint["outbox_effect_tombstones"]
)
return [
{
"effect_id": (
record["effect_id"]
if kind == "tombstone"
else record["intent"]["effect_id"]
),
"record_kind": kind,
"source_digest": hash_value(record),
}
for kind, record in records
]


def _expected_scope_maps(
case: ExecutionCheckpointCase,
expected: dict[str, Any],
) -> dict[str, dict[str, Any]]:
result = {}
for scope in expected["scopes"]:
checkpoints = {}
outbox_records = {}
for root_instance_id, binding in scope["checkpoints"].items():
expected_checkpoint = _json(case.path / binding["file"])
canonical = serialize_execution_checkpoint(expected_checkpoint)
assert (
f"sha256:{hashlib.sha256(canonical).hexdigest()}"
== binding["serialization_digest"]
)
assert (
expected_checkpoint["execution_checkpoint_digest"]
== binding["execution_checkpoint_digest"]
)
checkpoints[root_instance_id] = canonical
outbox_records[root_instance_id] = _outbox_records(expected_checkpoint)
assert [
record
for root_records in outbox_records.values()
for record in root_records
] == scope["outbox_records"]
result[scope["logical_scope_id"]] = {
"checkpoints": checkpoints,
"outbox_records": outbox_records,
}
return result


def _actual_scope_maps(
hosts: dict[str, ExecutionHost],
stores: dict[str, _ObservedMemoryStore],
) -> dict[str, dict[str, Any]]:
result = {}
for scope_id, store in stores.items():
checkpoints = {}
outbox_records = {}
for root_instance_id in sorted(store.root_instance_ids):
restored = hosts[scope_id].read_checkpoint(root_instance_id)
assert restored is not None
checkpoints[root_instance_id] = restored.canonical_bytes
outbox_records[root_instance_id] = _outbox_records(restored.document)
result[scope_id] = {
"checkpoints": checkpoints,
"outbox_records": outbox_records,
}
return result


def _assert_complete_scope_maps(
actual: dict[str, dict[str, Any]],
expected: dict[str, dict[str, Any]],
) -> None:
assert actual == expected


def _scope_hosts(
case: ExecutionCheckpointCase,
state: dict[str, Any],
store_calls: list[str],
) -> tuple[dict[str, ExecutionHost], dict[str, _ObservedMemoryStore]]:
hosts = {}
stores = {}
for scope in state["scopes"]:
initial = {
root_instance_id: (case.path / binding["file"]).read_bytes()
for root_instance_id, binding in scope["checkpoints"].items()
}
scope_id = scope["logical_scope_id"]
store = _ObservedMemoryStore(initial, store_calls)
stores[scope_id] = store
hosts[scope_id] = ExecutionHost(store, _resolver(case))
return hosts, stores


def _run_scope_vector(item: ExecutionCheckpointVector) -> None:
case = item.case
vector = item.vector
expected = vector["expect"]
before = _scope_state(case, vector["scope_state_before"])
after = _scope_state(case, vector["scope_state_after"])
store_calls: list[str] = []
hosts, stores = _scope_hosts(case, before, store_calls)

resolver_calls = ["resolve_execution_store_scope"]
selection = vector["scope_selection"]
requested = selection["requested_scope_id"]
candidates = selection["candidates"]
selected = (
requested
if requested is not None
and len(candidates) == 1
and candidates[0]["logical_scope_id"] == requested
and candidates[0]["authorized"]
else None
)
host_calls = []
if selected is not None:
host_calls.append("update_pending_outbox")
selected_scope = next(
scope for scope in before["scopes"] if scope["logical_scope_id"] == selected
)
root_instance_id, binding = next(iter(selected_scope["checkpoints"].items()))
checkpoint = _json(case.path / binding["file"])
hosts[selected].update_pending_outbox(
root_instance_id,
vector["effect_id"],
vector["desired_pending_state"],
expected_revision=checkpoint["revision"],
expected_checkpoint_digest=checkpoint["execution_checkpoint_digest"],
)

assert expected["selection_result"] == (
"selected" if selected is not None else "rejected"
)
assert expected["selected_scope_id"] == selected
assert expected["calls"] == {
"resolver": resolver_calls,
"execution_host": host_calls,
"store": store_calls,
"core": [],
}
_assert_complete_scope_maps(
_actual_scope_maps(hosts, stores),
_expected_scope_maps(case, after),
)


def run_execution_checkpoint_vector(item: ExecutionCheckpointVector) -> None:
case = item.case
vector = item.vector
expected = vector["expect"]
if "scope_state_before" in vector:
_run_scope_vector(item)
return
before_name = vector.get("checkpoint_before")
after_name = expected["checkpoint_after"]
if vector["operation"] in {
Expand Down
2 changes: 1 addition & 1 deletion conformance/pins.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

from pathlib import Path

CONFORMANCE_COMMIT = "263644f951f342b0eeaa3aceef4877293d2d7c67"
CONFORMANCE_COMMIT = "5ba78c7ef90b8556e76de6481a18b23a3d0c2378"
SPEC_COMMIT = "cc4b0d734aa1c5953de75fb53b63e390a3b72761"

ROOT = Path(__file__).resolve().parent.parent
Expand Down
Loading