Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
36 commits
Select commit Hold shift + click to select a range
17590f2
fix(control-plane): restore merged contract compatibility
Duang777 Oct 3, 2026
36dc506
test: align merged replan and Lark contracts
Duang777 Oct 3, 2026
896ef72
test: align merged catalog and probe fixtures
Duang777 Oct 3, 2026
4418fe4
refactor(presentation): place goal ownership API below root
Duang777 Oct 3, 2026
033e80e
Merge origin/main into CI baseline fixes
Duang777 Oct 3, 2026
968dd0d
test: follow workspace-aware replan binding
Duang777 Oct 3, 2026
d5ef665
fix(ci): stabilize shared runtime retries
Duang777 Oct 3, 2026
ca46f15
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 3, 2026
8afa4c2
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 3, 2026
b1449ab
fix(delegation): make preview retirement deterministic
Duang777 Oct 3, 2026
e86fe1a
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
3667c00
test(delegation): isolate retirement timer fixture
Duang777 Oct 4, 2026
478145f
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
e5f1f9e
Merge upstream/main into codex/fix-main-ci-contracts-20261003
Duang777 Oct 4, 2026
0406e5c
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
fe06e7c
fix(ci): restore post-merge contract baselines
Duang777 Oct 4, 2026
8039296
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
e423df4
fix(ci): refresh next-action mutation oracle
Duang777 Oct 4, 2026
bfe7640
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
490e01e
test(collaboration): model explicit source revocation
Duang777 Oct 4, 2026
e18d166
fix(delegation): retain ownership after unconfirmed cleanup
Duang777 Oct 4, 2026
7109a4d
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
10e699a
test(turn): use qualified delivery workspaces
Duang777 Oct 4, 2026
9ac43a7
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
aa34aca
fix(goals): ignore completion receipt in acceptance digest
Duang777 Oct 4, 2026
e2da9b2
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
63cd618
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
b4065e4
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
06e6f21
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
1796ac6
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
b44c486
test(goals): provide replan history source
Duang777 Oct 4, 2026
29cb075
chore(canary): refresh Lark runtime metric ceiling
Duang777 Oct 4, 2026
e966cd9
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 4, 2026
9ebefe8
Merge upstream/main into codex/fix-main-ci-contracts-20261003
Duang777 Oct 4, 2026
7cbe359
Merge upstream/main into codex/fix-main-ci-contracts-20261003
Duang777 Oct 4, 2026
647a5e8
test(frontend): update PR review localization contract
Duang777 Oct 4, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -459,7 +459,7 @@ for (const capabilityId of [
const matches = capabilityLocalization.match(new RegExp(`${capabilityId}:`, "g")) ?? [];
assert.equal(matches.length, 2, `${capabilityId} has English and Simplified Chinese metadata`);
}
for (const fieldKey of ["allowed_domains", "coordinator_agent_id", "eligible_endpoints", "enabled", "executor_endpoint", "executor_model", "executor_reasoning_effort", "max_children", "profile", "profile_preset", "review_priority", "route_ref", "safe_fix", "selection_policy", "strict_receipt", "timezone"]) {
for (const fieldKey of ["agent_orders", "allowed_domains", "coordinator_agent_id", "eligible_endpoints", "enabled", "executor_endpoint", "executor_model", "executor_reasoning_effort", "max_children", "profile", "profile_preset", "review_order", "route_ref", "safe_fix", "selection_policy", "strict_receipt", "timezone"]) {
const matches = capabilityLocalization.match(new RegExp(`^\\s+${fieldKey}:`, "gm")) ?? [];
assert.equal(matches.length, 2, `${fieldKey} has English and Simplified Chinese field copy`);
}
Expand Down
14 changes: 9 additions & 5 deletions examples/shared-goal-authority-e2e/mutants.py
Original file line number Diff line number Diff line change
Expand Up @@ -269,10 +269,10 @@ def apply(source: str) -> str:
Case("source_binding_outside_lock", ((COORDINATION + "legacy_writer_fence.py",
move_guard_outside_lock("require_registry_source_write_allowed")),),
WRITER_TEST + "test_waiting_override_writer_rechecks_registry_binding_inside_shared_state_lock"),
Case("remove_refresh_cas", (("loopx/state_refresh.py", replacement(
"if current_state_text != expected_write_state_text:",
"if False: # DELIBERATE MUTANT: bypass stale-state rejection.")),),
WRITER_TEST + "test_concurrent_public_refresh_preserves_the_newer_owned_paragraph"),
Case("remove_refresh_source_recheck", (("loopx/state_refresh.py", replacement(
" if normalized_next_action:",
" if False and normalized_next_action: # DELIBERATE MUTANT: bypass source recheck.")),),
"tests/control_plane/test_next_action_writeback.py::test_final_commit_rechecks_relevant_source_facts[task]"),
Case("fence_unshared_state_lock", ((COORDINATION + "legacy_writer_fence.ts", replacement(
"withFileMutationLock(statePath, () =>",
'withFileMutationLock(statePath + ".mutant-unshared", () =>')),),
Expand Down Expand Up @@ -420,7 +420,11 @@ def main() -> int:
log = mutant.stdout + mutant.stderr
# Pytest assertion rewriting can render rich comparisons as
# "E assert ..." without spelling the exception class.
assertion = "AssertionError" in log or re.search(r"^E\s+assert ", log, re.MULTILINE) is not None
assertion = (
"AssertionError" in log
or "Failed: DID NOT RAISE" in log
or re.search(r"^E\s+assert ", log, re.MULTILINE) is not None
)
killed = (mutant.returncode == 1 and assertion
and any(token in log for token in ("1 failed", "fail 1"))
and not any(token in log for token in ("SyntaxError", "ImportError", "ModuleNotFoundError")))
Expand Down
2 changes: 1 addition & 1 deletion loopx/canary/module_metric_baseline.json
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@
"dict_any_count": 0
},
"loopx/extensions/lark/goal_topic_runtime.py": {
"any_count": 49,
"any_count": 56,
"dict_any_count": 0
},
"loopx/extensions/lark/presentation/explore_results.py": {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,11 @@
from pathlib import Path
from typing import Any

from .capabilities.machine_configuration.store import read_stored_machine_configuration
from .control_plane.effect_runtime import effect_runtime_result
from .control_plane.runtime.runtime_projection_route import resolve_goal_source_runtime_route
from .history import load_registry
from .registry import registry_goals
from ..control_plane.effect_runtime import effect_runtime_result
from ..control_plane.runtime.runtime_projection_route import resolve_goal_source_runtime_route
from ..history import load_registry
from ..registry import registry_goals
from .machine_configuration.store import read_stored_machine_configuration


def capture_configuration_backup(
Expand Down
21 changes: 20 additions & 1 deletion loopx/capabilities/manager_context/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
POLICY_SCHEMA as POLICY_SCHEMA,
registered_context_recipients,
source_context_authority,
source_context_target_authority,
)
from ...control_plane.collaboration.goal_instance_scope import (
collaboration_goal_scope,
Expand Down Expand Up @@ -86,6 +87,15 @@ def authority(
grant = source_context_authority(runtime_root, registry_path, session, turn)
return {**grant, "instruction": INSTRUCTION} if grant["mode"] == "context_only" else grant


def target_authority(
runtime_root: Path, *, session: dict, turn: dict, target: dict
) -> dict:
"""Authorize one target already validated by an exact Goal scope."""
grant = source_context_target_authority(runtime_root, session, turn, target)
return {**grant, "instruction": INSTRUCTION} if grant["mode"] == "context_only" else grant


def deliver(
runtime_root: Path, registry_path: Path, *, session: dict, turn: dict, request: dict
) -> dict:
Expand All @@ -101,7 +111,16 @@ def deliver(
goal_scope,
operation="request_create",
)
grant = authority(runtime_root, registry_path, session, turn)
grant = (
target_authority(
runtime_root,
session=session,
turn=turn,
target=target,
)
if goal_scope.exact
else authority(runtime_root, registry_path, session, turn)
)
if target not in grant["targets"]:
raise ValueError("context recipient is not authorized or registered")
content = str(turn.get("message") or "")
Expand Down
11 changes: 8 additions & 3 deletions loopx/capabilities/manager_context/roundtrip.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
from datetime import datetime, timezone, timedelta
from uuid import uuid4

from . import _root, _read, _write, _hash, authority
from . import _root, _read, _write, _hash, authority, target_authority
from .tracking import _entry, _now
from ...file_lock import (
LockAcquisitionPolicy,
Expand Down Expand Up @@ -380,7 +380,7 @@ def _exact_return_scope(registry, reply):
return collaboration_goal_scope(
registry,
goal_id=reply["goal_id"],
agents=(),
agents=(reply["agent_id"],),
caller_goal_ref=reply["goal_ref"],
)

Expand Down Expand Up @@ -418,8 +418,13 @@ def _exact_return_context(root, registry, store, path, state_path, now):
or not turn
):
raise ValueError("original_conversation_unavailable")
grant = authority(root, registry, session, turn)
target = {key: row[key] for key in ("goal_id", "agent_id")}
grant = target_authority(
root,
session=session,
turn=turn,
target=target,
)
if (
target not in grant["targets"]
or grant.get("source_id") != row["source_id"]
Expand Down
2 changes: 1 addition & 1 deletion loopx/chat_configuration_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,13 @@

from collections.abc import Callable

from .presentation import configuration_backup_api as backup_api
from .presentation import goal_ownership_api as ownership_api
from . import chat_usage_statistics_api as usage_api
from . import chat_goal_configuration_api as goal_api
from . import chat_machine_configuration_api as machine_api
from . import chat_operator_provider_api as operator_api
from . import chat_automation_cadence_api as cadence_api
from . import chat_configuration_backup_api as backup_api


class ChatConfigurationRequestMixin(
Expand Down
2 changes: 1 addition & 1 deletion loopx/cli_commands/configuration_backup.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
import tempfile
from pathlib import Path

from ..configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup
from ..capabilities.configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup
from ..history import load_registry
from ..paths import resolve_runtime_root

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,6 @@ async function accept(value: unknown) {
const request = decodeHostProcessRequest(v.request);
if (request.input !== "") throw new Error("preview input must be framed");
started = true;
// Stop accepting before the Host's independent lifetime deadline begins
// cleanup; otherwise a new request could be admitted into a dying worker.
lifetime = setTimeout(() => stop("lifetime"), LIFETIME_MS);
running = runHostProcess({...request, timeout_ms: LIFETIME_MS,
stdout_limit_bytes: LIMIT * MAX_REQUESTS}, async item => {
if (item.kind !== "stdout") return; // Never relay private worker diagnostics.
Expand All @@ -64,6 +61,9 @@ async function accept(value: unknown) {
else if (!pending) armIdle();
}
}, owner.signal, undefined, {openInput: input => { write = input; }});
// Start the reuse lifetime after synchronous worker startup. This still
// stops admission before Host cleanup, without charging spawn latency.
lifetime = setTimeout(() => stop("lifetime"), LIFETIME_MS);
void running.then(async result => {
const originalPending = pending;
stop(result.outcome);
Expand Down
84 changes: 68 additions & 16 deletions loopx/control_plane/collaboration/delegation_preview_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import hashlib
import json
import os
import signal
import subprocess
import time
import weakref
Expand All @@ -19,6 +20,9 @@
from ..effect_runtime import _node_executable


BRIDGE_CLOSE_TIMEOUT_SECONDS = 5.0


def _source_snapshot(release: Path) -> tuple:
"""Loaded-code identity only; authority/configuration is read per request."""
files = []
Expand Down Expand Up @@ -47,22 +51,61 @@ def _source_snapshot(release: Path) -> tuple:
return tuple(files)


def _close_bridge(process: subprocess.Popen) -> None:
def _terminate_bridge(process: subprocess.Popen) -> None:
if process.poll() is not None:
return
if os.name != "nt":
try:
process.send_signal(signal.SIGCONT)
except ProcessLookupError:
return
process.terminate()


def _kill_bridge(process: subprocess.Popen) -> None:
if process.poll() is not None:
return
try:
process.kill()
except ProcessLookupError:
pass


def _close_bridge(
process: subprocess.Popen,
*,
force: bool = False,
cleanup_confirmed: bool = False,
) -> bool:
# Parent EOF cancels the TS-owned group; give its cleanup fence time to run.
if force:
_terminate_bridge(process)
if process.stdin is not None and not process.stdin.closed:
try:
process.stdin.close()
except OSError:
pass # A crashed/retired supervisor may already have closed its pipe.
try:
process.wait(timeout=5)
process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS)
except subprocess.TimeoutExpired:
# SIGTERM asks the supervisor to clean, not to abandon its worker.
process.terminate()
process.wait(timeout=5)
if not force:
# SIGTERM asks the supervisor to clean, not to abandon its worker.
_terminate_bridge(process)
try:
process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS)
except subprocess.TimeoutExpired:
_kill_bridge(process)
process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS)
else:
_kill_bridge(process)
process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS)
finally:
if process.stdout is not None:
process.stdout.close()
# SIGTERM can win before the bridge installs its handlers, before it can
# spawn a worker. Once initialized, normal exit follows Host group cleanup.
# A SIGKILLed supervisor provides neither guarantee.
return cleanup_confirmed or process.returncode in (0, -signal.SIGTERM)


class DelegationPreviewTransport:
Expand All @@ -75,15 +118,21 @@ def __init__(self) -> None:
self._finalizer: weakref.finalize | None = None
self._sequence = 0

def _close(self) -> None:
if self._process is not None:
_close_bridge(self._process)
if self._finalizer is not None:
self._finalizer.detach()
self._process = None
self._finalizer = None
self._partition = None
def _close(
self, *, force: bool = False, cleanup_confirmed: bool = False
) -> bool:
process, finalizer = self._process, self._finalizer
if process is not None and not _close_bridge(
process,
force=force,
cleanup_confirmed=cleanup_confirmed,
):
return False
self._process = self._partition = self._finalizer = None
self._sequence = 0
if finalizer is not None:
finalizer.detach()
return True

def close(self) -> None:
with self._lock:
Expand All @@ -108,7 +157,10 @@ def preview(self, *, command: list[str], workspace: Path, release: Path,
for replacement in (False, True):
if (self._partition != partition or self._process is None
or self._process.poll() is not None or self._sequence >= 128):
self._close()
if not self._close():
raise ValueError(
"delegation preview cleanup remains unconfirmed"
)
bridge = Path(__file__).with_name("delegation_preview_bridge.ts")
self._process = subprocess.Popen(
[_node_executable(), "--no-warnings", "--experimental-strip-types", str(bridge)],
Expand Down Expand Up @@ -144,7 +196,7 @@ def preview(self, *, command: list[str], workspace: Path, release: Path,
and type(response["last_id"]) is int):
# The TS owner confirms this request was not accepted and
# its old group stopped. Reuse the original deadline/binding.
self._close()
self._close(cleanup_confirmed=True)
continue
if response.get("kind") == "failure" and response.get("outcome") == "timeout":
raise subprocess.TimeoutExpired(["delegation-preview"], timeout)
Expand All @@ -155,7 +207,7 @@ def preview(self, *, command: list[str], workspace: Path, release: Path,
return response["value"]
raise ValueError("delegation preview retirement did not complete")
except BaseException:
self._close()
self._close(force=True)
raise
finally:
self._lock.release()
Expand Down
Loading
Loading