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
44 changes: 41 additions & 3 deletions loopx/cli_commands/todo.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,11 @@
from collections.abc import Callable, Sequence
from pathlib import Path

from ..control_plane.coordination.local_authority import read_canonical_todo_fields_if_promoted
from ..control_plane.coordination.local_authority import (
local_authority_is_promoted,
read_canonical_todo_fields_if_promoted,
)
from ..control_plane.effect_runtime import effect_runtime_result
from ..control_plane.agents.workspace_guard import capture_delivery_workspace
from ..control_plane.todos.contract import (
replan_successor_semantic_binding,
Expand Down Expand Up @@ -47,6 +51,7 @@
validate_todo_claim_options,
validate_todo_complete_options,
validate_todo_list_options,
validate_todo_receipt_options,
validate_todo_project_markdown_options,
validate_todo_plan_options,
validate_todo_supersede_options,
Expand Down Expand Up @@ -182,6 +187,23 @@ def _todo_path_args(args: argparse.Namespace) -> dict[str, Path | None]:
}


def _render_todo_receipt(payload: dict[str, object]) -> str:
lines = [
"# LoopX Canonical Operation Receipt",
"",
f"- status: `{payload.get('status')}`",
f"- goal_id: `{payload.get('goal_id')}`",
f"- operation_id: `{payload.get('operation_id')}`",
f"- source_authority: `{payload.get('source_authority')}`",
f"- provider_revision: `{payload.get('provider_revision')}`",
f"- cursor: `{payload.get('cursor')}`",
"- note: Historical readback only; it does not grant a current lease or a retry.",
]
if payload.get("error") or payload.get("reason"):
lines.append(f"- error: `{payload.get('error') or payload.get('reason')}`")
return "\n".join(lines)


def handle_todo_command(
args: argparse.Namespace,
*,
Expand All @@ -194,8 +216,8 @@ def handle_todo_command(
post_writeback_projection_builder: PostWritebackProjectionBuilder | None = None,
) -> int:
renderer = (
render_task_planning_packet
if args.todo_command == "plan"
render_task_planning_packet if args.todo_command == "plan"
else _render_todo_receipt if args.todo_command == "receipt"
else render_todo_markdown
)
try:
Expand Down Expand Up @@ -228,6 +250,22 @@ def handle_todo_command(
**_todo_path_args(args),
runtime_root_arg=runtime_root_arg,
)
elif args.todo_command == "receipt":
validate_todo_receipt_options(args)
runtime_root = resolve_runtime_root(load_registry(registry_path), runtime_root_arg)
if not local_authority_is_promoted(runtime_root=runtime_root, goal_id=args.goal_id):
raise ValueError("todo receipt requires promoted canonical authority; no legacy fallback")
result = effect_runtime_result(
"coordination.local_authority.operation_receipt",
{"schema_version": "loopx_local_coordination_operation_receipt_request_v0",
"runtime_root": str(runtime_root.expanduser().resolve(strict=False)),
"goal_id": args.goal_id, "operation_id": args.operation_id},
timeout=15.0,
)
if not isinstance(result, dict):
raise RuntimeError("canonical operation receipt returned an invalid result")
payload = {"ok": result.get("status") in {"found", "missing"},
"command": "receipt", **result}
elif args.todo_command == "project-markdown":
validate_todo_project_markdown_options(args)
registry = load_registry(registry_path)
Expand Down
12 changes: 12 additions & 0 deletions loopx/cli_commands/todo_argument_validation.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
("--priority", "priority"),
("--clear-priority", "clear_priority"),
("--todo-id", "todo_id"),
("--operation-id", "operation_id"),
("--claim-operation-id", "claim_operation_id"),
("--update-operation-id", "update_operation_id"),
("--update-expected-provider-revision", "update_expected_provider_revision"),
Expand Down Expand Up @@ -300,6 +301,15 @@ def validate_todo_list_options(args: argparse.Namespace) -> None:
)


def validate_todo_receipt_options(args: argparse.Namespace) -> None:
_validate_todo_option_subset(
args, {"operation_id"},
"todo receipt only accepts --goal-id, --operation-id, and --format; unsupported: ",
)
if not args.operation_id:
raise ValueError("todo receipt requires --operation-id")


def validate_todo_plan_options(args: argparse.Namespace) -> None:
_validate_todo_option_subset(
args, {"text", "agent_id"},
Expand Down Expand Up @@ -493,6 +503,8 @@ def validate_todo_archive_completed_options(args: argparse.Namespace) -> None:


def validate_shared_todo_options(args: argparse.Namespace) -> None:
if getattr(args, "operation_id", None) and args.todo_command != "receipt":
raise ValueError("--operation-id is supported only by todo receipt")
agent_id_allowed_for_user_authoring = (
args.todo_command == "add"
and args.role == "user"
Expand Down
3 changes: 2 additions & 1 deletion loopx/cli_commands/todo_event.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,8 @@ def append_todo_rollout_event(
and getattr(args, "no_follow_up", False)
)
if (
not payload.get("ok")
args.todo_command == "receipt"
or not payload.get("ok")
or payload.get("dry_run")
or (payload.get("idempotent_replay") and not turn_instance_id)
):
Expand Down
5 changes: 5 additions & 0 deletions loopx/cli_commands/todo_registration.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ def register_todo_command(
choices=[
"add",
"list",
"receipt",
"claim",
"update",
"complete",
Expand All @@ -53,6 +54,10 @@ def register_todo_command(
todo_parser.add_argument("--priority", choices=["P0", "P1", "P2", "P3", "P4"], help="For add/update, declare Todo priority independently of text; omission retains the current value.")
todo_parser.add_argument("--clear-priority", action="store_true", help="For update, explicitly remove priority; cannot be combined with --priority.")
todo_parser.add_argument("--todo-id", help="Structured todo id from status/quota, such as todo_ab12cd34ef56.")
todo_parser.add_argument(
"--operation-id",
help="For todo receipt, read the exact historical canonical operation after an ambiguous response; this does not grant a retry or lease.",
)
todo_parser.add_argument(
"--update-operation-id",
help=("For promoted text/note, planning, validator revision or User completion update, reuse this operation id after a lost response; "
Expand Down
55 changes: 48 additions & 7 deletions loopx/control_plane/coordination/file_authority_store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,29 @@ import {AuthorityJournalScan} from "./authority_journal_scan.ts";

const FILE_AUTHORITY_STORE_SCHEMA = "loopx_file_authority_store_v0";
const STORE_IDENTITY_PATTERN = /^file:[0-9a-f]{32}$/;
// File-v0 retains every projection in one envelope. A managed Effect server
// opens a new store handle for each request, so revalidating an unchanged
// journal on every read makes one Goal's history dominate the RPC budget.
// Keep only one verified document process-wide; bytes and store identity must
// both match before a later handle may reuse that validation.
const MAX_CACHED_DOCUMENT_BYTES = 128 * 1024 * 1024;
let verifiedDocument: {
path: string;
identity: string;
digest: string;
document: FileAuthorityStoreDocument;
} | null = null;

function documentDigest(raw: Uint8Array): string {
return createHash("sha256").update(raw).digest("hex");
}

function rememberVerifiedDocument(path: string, identity: string, raw: Uint8Array,
digest: string, document: FileAuthorityStoreDocument): void {
verifiedDocument = raw.byteLength <= MAX_CACHED_DOCUMENT_BYTES
? {path, identity, digest, document}
: null;
}

interface FileAuthorityStoreDocument extends JsonObject, RetainedAuthorityJournal {
schema_version: typeof FILE_AUTHORITY_STORE_SCHEMA;
Expand Down Expand Up @@ -167,6 +190,11 @@ export class FileAuthorityStore implements AuthorityStore {
/** Filesystem-only crash seam; the archive owner must still fsync both parents. */
protected async archiveRenamed(): Promise<void> {}

/** Full-history verification seam; unchanged byte-identical reads may reuse it. */
protected decodeStoredDocument(value: unknown, identity: string): FileAuthorityStoreDocument {
return decodeDocument(value, this.goalId, identity);
}

private async readStoreIdentity(createIfMissing = !this.existingOnly): Promise<string> {
try {
const identity = await readFile(this.identityPath, "utf8");
Expand Down Expand Up @@ -213,19 +241,26 @@ export class FileAuthorityStore implements AuthorityStore {
}
}

private async readDocument(): Promise<FileAuthorityStoreDocument | null> {
let raw: string;
private async readDocument(knownIdentity?: string): Promise<FileAuthorityStoreDocument | null> {
let raw: Buffer;
try {
raw = await readFile(this.path, "utf8");
raw = await readFile(this.path);
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
throw new FileStoreUnavailableError(
error instanceof Error ? error.message : "authority document unavailable",
);
}
const identity = await this.readStoreIdentity();
const identity = knownIdentity ?? await this.readStoreIdentity();
try {
return decodeDocument(JSON.parse(raw), this.goalId, identity);
const digest = documentDigest(raw);
if (verifiedDocument?.path === this.path &&
verifiedDocument.identity === identity && verifiedDocument.digest === digest) {
return verifiedDocument.document;
}
const document = this.decodeStoredDocument(JSON.parse(raw.toString("utf8")), identity);
rememberVerifiedDocument(this.path, identity, raw, digest, document);
return document;
} catch (error) {
if (error instanceof SyntaxError) {
throw new AuthorityStoreProtocolError(`file authority store JSON is invalid: ${error.message}`);
Expand Down Expand Up @@ -284,7 +319,7 @@ export class FileAuthorityStore implements AuthorityStore {
// A restored directory must not race a missing-head bootstrap and
// bind new authority bytes to an identity observed before the lock.
identity = await this.readStoreIdentity();
current = await this.readDocument();
current = await this.readDocument(identity);
} catch (error) {
return {
status: "failed",
Expand Down Expand Up @@ -319,8 +354,14 @@ export class FileAuthorityStore implements AuthorityStore {
const document: FileAuthorityStoreDocument = {schema_version: FILE_AUTHORITY_STORE_SCHEMA,
goal_id: this.goalId, store_identity: identity, ...journal};
try {
await this.replaceDurably(this.path, canonicalAuthorityBytes(document));
const bytes = canonicalAuthorityBytes(document);
await this.replaceDurably(this.path, bytes);
rememberVerifiedDocument(this.path, identity, bytes, documentDigest(bytes), document);
} catch (error) {
// A failure after rename may already have published the new bytes.
// The next read must prove the actual file rather than reuse either
// the previous or attempted document.
verifiedDocument = null;
return {
status: "ambiguous",
reason_code: "commit_outcome_unknown",
Expand Down
33 changes: 33 additions & 0 deletions loopx/control_plane/coordination/local_authority_read.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,39 @@ import {LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_LIS
LOCAL_COORDINATION_TODO_READ_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_READ_RESULT_SCHEMA} from "./coordination_state_contract.generated.ts";
import {decodeProjectionReadback, confirmProjectionReadback} from "../todos/projection_delivery.ts";

const LOCAL_COORDINATION_OPERATION_RECEIPT_REQUEST_SCHEMA = "loopx_local_coordination_operation_receipt_request_v0";
const LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA = "loopx_local_coordination_operation_receipt_result_v0";

/** Exact historical operation readback. It grants no current lease or retry authority. */
export async function readLocalCoordinationOperationReceipt(
value: unknown,
dependencies: LocalAuthorityProviderDependencies = {},
): Promise<JsonObject> {
let sourceAuthority = "canonical_unavailable";
try {
const input = requireJsonObject(value, "local coordination operation receipt request");
if (input.schema_version !== LOCAL_COORDINATION_OPERATION_RECEIPT_REQUEST_SCHEMA) {
throw new Error("local coordination operation receipt request schema mismatch");
}
const root = runtimeRoot(input.runtime_root);
const goalId = requireAuthorityStoreId(input.goal_id, "goal id");
const operationId = requireAuthorityStoreId(input.operation_id, "operation id");
const store = await openRuntimeStore(root, goalId, dependencies);
sourceAuthority = sourceAuthorityFor(store);
const receipt = await store.readReceipt(operationId);
return {schema_version: LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA,
goal_id: goalId, operation_id: operationId, ...receipt,
source_authority: sourceAuthority, decision_read_from_provider: true,
legacy_fallback_used: false};
} catch (error) {
return {schema_version: LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA,
status: "failed", reason_code: "invalid_local_coordination_operation_receipt_request",
reason: error instanceof Error ? error.message : "invalid operation receipt request",
source_authority: sourceAuthority, decision_read_from_provider: false,
legacy_fallback_used: false, ...localAuthorityOpenFailure(error)};
}
}

/** Provider-first exact Todo read. Missing/unavailable state never falls back. */
export async function readLocalCoordinationTodo(
value: unknown,
Expand Down
7 changes: 6 additions & 1 deletion loopx/control_plane/coordination/local_authority_write.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,16 @@
import {withFileMutationLock} from "../effect_runtime_io.ts";
import {requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath} from "./shadow_management.ts";

// File-v0 still verifies and republishes its retained history on each write.
// Give a competing canonical writer a bounded wait for that real critical
// section; the provider CAS and hard lease remain the decision authority.
const CANONICAL_WRITER_LOCK_TIMEOUT_MS = 30_000;

/** Canonical command writers share the maintenance guard; provider CAS owns state. */
export async function withCanonicalWriter<T>(root: string, goalId: string, dryRun: boolean, write: () => Promise<T>): Promise<T> {
if (dryRun) return await write();
return await withFileMutationLock(shadowMaintenanceLockPath(root, goalId), async () => {
await requireShadowPrimaryWriteAllowed(root, goalId);
return await write();
});
}, CANONICAL_WRITER_LOCK_TIMEOUT_MS);
}
Loading
Loading