Skip to content
Draft
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
6 changes: 4 additions & 2 deletions loopx/control_plane/coordination/authority_journal_scan.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
* Storage effects and snapshot acquisition remain with each provider. */
import type {AuthorityStoreCommittedTransaction, AuthorityStoreHead,
AuthorityStoreReadFailure, AuthorityStoreScanResult} from "./authority_store.ts";
import {AuthorityStoreProtocolError, canonicalAuthorityBytes, parseAuthorityCursor} from "./authority_store_codec.ts";
import {AuthorityStoreProtocolError, canonicalAuthorityBytes, copyAuthorityJson, parseAuthorityCursor} from "./authority_store_codec.ts";

export class AuthorityJournalScan {
readonly after: string | null;
Expand Down Expand Up @@ -57,7 +57,9 @@ export class AuthorityJournalScan {
throw new AuthorityStoreProtocolError("committed scan head lineage is invalid");
}
}
const transactions = structuredClone(rows.slice(0, this.limit));
// Each caller owns the mutable JSON containers; immutable primitive values
// need no additional byte copy for every retained projection in the page.
const transactions = copyAuthorityJson(rows.slice(0, this.limit)) as AuthorityStoreCommittedTransaction[];
return {status: "page", transactions,
next_cursor: transactions.at(-1)?.cursor ?? this.after, has_more: rows.length > this.limit};
}
Expand Down
7 changes: 6 additions & 1 deletion loopx/control_plane/coordination/authority_state_log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
canonicalAuthorityJson,
canonicalAuthorityObject,
canonicalAuthoritySha256,
copyAuthorityJson,
isAuthorityJsonObject,
} from "./authority_store_codec.ts";

Expand Down Expand Up @@ -215,7 +216,11 @@ export class AuthorityStateReplay {
this.#state = next;
}

snapshot(): JsonObject { return structuredClone(this.#state); }
snapshot(): JsonObject {
// Copy every mutable JSON container; immutable strings need no new byte
// allocation for each historical row carrying the same large value.
return copyAuthorityJson(this.#state) as JsonObject;
}

/** Immutable canonical text; storage readers may parse it into independent rows. */
canonicalJson(): string { return this.#encode(this.#state); }
Expand Down
23 changes: 20 additions & 3 deletions loopx/control_plane/coordination/authority_store_codec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,15 @@ export function canonicalAuthorityJson(
value: unknown,
stack = new Set<object>(),
): unknown {
return cloneAuthorityJson(value, stack, true);
}

/** Own JSON containers without copying immutable primitives or reordering keys. */
export function copyAuthorityJson(value: unknown): unknown {
return cloneAuthorityJson(value, new Set<object>(), false);
}

function cloneAuthorityJson(value: unknown, stack: Set<object>, canonicalKeys: boolean): unknown {
if (value === null || typeof value === "string" || typeof value === "boolean") {
return value;
}
Expand All @@ -59,7 +68,13 @@ export function canonicalAuthorityJson(
if (stack.has(value)) throw new AuthorityStoreProtocolError("JSON value must be acyclic");
stack.add(value);
try {
return value.map((item) => canonicalAuthorityJson(item, stack));
if (canonicalKeys) return value.map((item) => cloneAuthorityJson(item, stack, true));
const copy: unknown[] = Array(value.length);
for (const key of Object.keys(value)) Object.defineProperty(copy, key, {
value: cloneAuthorityJson(Reflect.get(value, key), stack, false),
writable: true, enumerable: true, configurable: true,
});
return copy;
} finally {
stack.delete(value);
}
Expand All @@ -74,10 +89,12 @@ export function canonicalAuthorityJson(
if (stack.has(value)) throw new AuthorityStoreProtocolError("JSON value must be acyclic");
stack.add(value);
try {
const keys = Object.keys(value);
if (canonicalKeys) keys.sort(authorityUnicodeCompare);
return Object.fromEntries(
Object.keys(value).sort(authorityUnicodeCompare).map((key) => [
keys.map((key) => [
key,
canonicalAuthorityJson(value[key], stack),
cloneAuthorityJson(value[key], stack, canonicalKeys),
]),
);
} finally {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -594,7 +594,7 @@ export class SqliteAuthorityStore implements AuthorityStore {
if (rows.length) this.verifiedRange(db, BigInt(String(rows[0]!.sequence)),
BigInt(String(rows[rows.length - 1]!.sequence)), (row, replay, identity) => {
verified.push({cursor: row.cursor.toString(), provider_revision: `${identity}:${row.cursor}`,
operation_id: row.operation_id, projection: JSON.parse(replay.canonicalJson()) as JsonObject, events: row.events, receipts: row.receipts});
operation_id: row.operation_id, projection: replay.snapshot(), events: row.events, receipts: row.receipts});
});
return scan.page(verified, head === null ? null : {cursor: head.state.cursor.toString(),
provider_revision: head.provider_revision, head: head.state.projection});
Expand Down
113 changes: 106 additions & 7 deletions tests/control_plane/test_local_authority_shadow_outbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,28 +8,45 @@

from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter
from loopx.control_plane.coordination import local_authority_shadow_outbox as outbox
from loopx.control_plane.coordination.coordination_state_contract_generated import (
COORDINATION_RUNTIME_SHADOW_BOOTSTRAP_RESULT_SCHEMA,
)
from loopx.control_plane.coordination.local_authority_shadow_projection import (
ProjectionValueError,
canonical_bytes,
lease_partition_projection,
partition_digest,
sha256_digest,
source_effect_runtime_result,
text_digest,
todo_partition_projection,
)
from loopx.control_plane.coordination.runtime_shadow import (
bootstrap_coordination_runtime_shadow,
build_runtime_shadow_source_snapshot,
RuntimeInvoker,
)
from loopx.control_plane.coordination.shadow_management import (
read_shadow_bootstrap_source_path,
read_shadow_management_state,
require_shadow_primary_write_allowed,
shadow_maintenance_lock_target,
shadow_management_directory,
shadow_management_state_path,
)
from loopx.control_plane.coordination.shadow_management import require_shadow_primary_write_allowed
from loopx.control_plane.effect_runtime import EffectRuntimeResponseAmbiguous, EffectRuntimeRejected
from loopx.file_lock import exclusive_mutation_file_lock
from loopx.history import load_registry
from loopx.registry import find_registry_goal


GOAL_ID = "goal-outbox"


def _fixture(tmp_path: Path, *, bootstrap: bool = True) -> tuple[Path, Path, Path]:
def _fixture(
tmp_path: Path, *, bootstrap: bool = True,
runtime_invoker: RuntimeInvoker = source_effect_runtime_result,
) -> tuple[Path, Path, Path]:
repo = tmp_path / "repo"
repo.mkdir()
state = repo / "ACTIVE_GOAL_STATE.md"
Expand Down Expand Up @@ -75,15 +92,35 @@ def _fixture(tmp_path: Path, *, bootstrap: bool = True) -> tuple[Path, Path, Pat
"schema_version": "loopx_coordination_runtime_shadow_config_v0",
"enabled": True, "provider": "file_v0",
}}}
def bootstrap_or_receipt(method: str, request: dict) -> dict:
try:
return runtime_invoker(method, request)
except EffectRuntimeResponseAmbiguous:
# The RPC deadline still applies. Do not resend an uncertain
# write; read its exact completed receipt after the existing
# maintenance lock releases, using the normal lock budget.
with exclusive_mutation_file_lock(shadow_maintenance_lock_target(runtime_root, GOAL_ID)):
journal = read_shadow_management_state(runtime_root, GOAL_ID)
if journal is None:
raise
assert journal["status"] == "active", journal
assert journal["operation"]["operation_id"] == request["operation_id"], journal
assert journal["operation"]["request_digest"] == sha256_digest(request), journal
binding = require_shadow_primary_write_allowed(runtime_root, GOAL_ID)
assert binding is not None, journal
assert read_shadow_bootstrap_source_path(runtime_root, GOAL_ID, binding) == state.resolve()
receipt = journal["result"]
assert receipt["operation_id"] == request["operation_id"], receipt
assert {key: receipt[key] for key in binding} == binding, receipt
return {"schema_version": COORDINATION_RUNTIME_SHADOW_BOOTSTRAP_RESULT_SCHEMA, **receipt}

result = bootstrap_coordination_runtime_shadow(
goal=enabled_goal, runtime_root=runtime_root, goal_id=GOAL_ID,
operation_id="bootstrap:outbox-test", source_version="source:initial",
projection=projection, source_snapshot=snapshot,
projection=projection, source_snapshot=snapshot, runtime_invoker=bootstrap_or_receipt,
)
# The managed Effect runtime may lose the first response after the
# durable bootstrap commit and retry the same operation. Windows CI is
# slow enough to exercise that path, so the public success contract is
# applied/recovered/replayed rather than applied-only.
# Outbox assertions require a verified bootstrap, including an exact
# durable result when the managed RPC response is ambiguous.
assert result["status"] in {"applied", "recovered", "replayed"}, result
assert result["operation_id"] == "bootstrap:outbox-test", result
assert result["cursor"] == "1", result
Expand Down Expand Up @@ -161,6 +198,68 @@ def _files(directory: Path) -> dict[str, bytes]:
return {str(path.relative_to(directory)): path.read_bytes() for path in directory.rglob("*") if path.is_file()}


def test_outbox_fixture_reads_exact_bootstrap_receipt_after_lost_response(tmp_path: Path) -> None:
calls = []
committed = {}

def lose_response(method: str, request: dict) -> dict:
result = source_effect_runtime_result(method, request)
assert result["status"] in {"applied", "recovered", "replayed"}, result
calls.append(request)
path = shadow_management_state_path(Path(request["runtime_root"]), GOAL_ID)
committed["journal_bytes"] = path.read_bytes()
raise EffectRuntimeResponseAmbiguous(method, timeout=10)

registry, state, runtime_root = _fixture(tmp_path, runtime_invoker=lose_response)
assert len(calls) == 1 # A receipt readback must never dispatch another write.
assert shadow_management_state_path(runtime_root, GOAL_ID).read_bytes() == committed["journal_bytes"]
_record_change(registry, state, runtime_root, "Continue after verified bootstrap readback.")
assert _drain(registry, runtime_root).reason_code is None


@pytest.mark.parametrize("damage", ["missing", "pending", "other_operation", "changed_request", "changed_manifest"])
def test_outbox_fixture_rejects_unverified_bootstrap_after_lost_response(tmp_path: Path, damage: str) -> None:
calls = []

def lose_response(method: str, request: dict) -> dict:
calls.append(request)
if damage != "missing":
source_effect_runtime_result(method, request)
root = Path(request["runtime_root"])
path = shadow_management_state_path(root, GOAL_ID)
journal = json.loads(path.read_text(encoding="utf-8"))
if damage == "pending":
journal.update(status="bootstrapping", binding=None, result=None)
journal["operation"]["phase"] = "prepared"
elif damage == "other_operation":
journal["operation"]["operation_id"] = "bootstrap:another-operation"
elif damage == "changed_request":
journal["operation"]["request_digest"] = "sha256:" + "0" * 64
else:
manifest = next((shadow_management_directory(root, GOAL_ID) / "operations").glob("*/manifest.json"))
manifest.write_text("{}", encoding="utf-8")
path.write_text(json.dumps(journal), encoding="utf-8")
raise EffectRuntimeResponseAmbiguous(method, timeout=10)

with pytest.raises(AssertionError):
_fixture(tmp_path, runtime_invoker=lose_response)
assert len(calls) == 1
assert not _todo_dir(tmp_path / "runtime").exists()


def test_outbox_fixture_preserves_semantic_rejection_even_with_a_bootstrap_receipt(tmp_path: Path) -> None:
calls = []

def reject_response(method: str, request: dict) -> dict:
source_effect_runtime_result(method, request)
calls.append(request)
raise EffectRuntimeRejected("deliberate semantic rejection", diagnostic_code="invalid_request")

with pytest.raises(AssertionError, match="deliberate semantic rejection"):
_fixture(tmp_path, runtime_invoker=reject_response)
assert len(calls) == 1


def test_capture_records_prepared_then_committed_and_skips_prose_only_writes(tmp_path: Path) -> None:
registry, state, runtime_root = _fixture(tmp_path)
original = state.read_text(encoding="utf-8")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,8 @@ def test_real_cli_compact_and_full_detail_preserve_canonical_todos(tmp_path, mon
"role": "agent" if i < 40 else "user", "status": "open" if i % 3 else "done",
"done": i % 3 == 0, "text": f"Retained work {i}", "note": "exact metadata🙂" * 100,
"archive_state": "active", "source_section": "Agent Todo" if i < 40 else "User Todo",
"index": i + 1, "task_class": "advancement_task"} for i in range(60)]
"index": i + 1, "task_class": "advancement_task" if i < 40 else "user_action",
**({"goal_bound": True} if i >= 40 else {})} for i in range(60)]
projection = build_todo_runtime_shadow_projection(goal_id="example", todos=records, leases=[], handoff_mode="soft_claim")
initialize_canonical_authority(runtime, "example", projection, state_path=state, provider=provider)

Expand Down Expand Up @@ -97,6 +98,7 @@ def call(mode, *details):
if record["role"] == role:
assert actual[record["todo_id"]]["note"] == record["note"]
assert actual[record["todo_id"]]["status"] == record["status"]
assert actual[record["todo_id"]]["task_class"] == record["task_class"]
assert len(json.dumps(compact)) < len(json.dumps(full)) / 2
assert not list((runtime / "goals" / "example" / "runs").glob("*.json*"))
finally:
Expand Down
4 changes: 3 additions & 1 deletion tests/control_plane/test_quota_selection.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,10 +59,12 @@ def test_public_quota_keeps_gate_from_real_canonical_provider(tmp_path: Path, di
runtime, registry, state = tmp_path / "runtime", tmp_path / "registry.json", tmp_path / "state.md"
history = "\n".join(f"- [x] [P2] Completed synthetic work {i}.\n"
f" <!-- loopx:todo todo_id=todo_history_{i} status=done task_class=advancement_task -->" for i in range(12))
# A goal-wide gate cannot also bind continuation through a legacy agent claim.
claim = " claimed_by=agent-a" if scope == "blocks_agent=agent-b" else ""
state.write_text("---\nstatus: active\n---\n# Goal\n## Objective\nDeliver a checked change.\n\n"
"## Agent Todo\n" + history + "\n\n## User Todo\n"
"- [ ] [P0] Owner approval is required.\n"
f" <!-- loopx:todo todo_id=todo_gate task_class=user_gate status=open claimed_by=agent-a {scope} -->\n")
f" <!-- loopx:todo todo_id=todo_gate task_class=user_gate status=open{claim} {scope} -->\n")
write_fixture_registry(project=tmp_path, runtime_root=runtime, registry_path=registry,
goal_id="goal-scope", domain="quota-scope", adapter_kind="generic_project_goal_v0",
state_file=str(state), registered_agents=["agent-a", "agent-b"], quota_allowed_slots=None)
Expand Down
Loading
Loading