Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
45 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
7f30062
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
ebfdc34
fix(ci): align integration fixtures with current contracts
Duang777 Oct 5, 2026
d50e3b1
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
310cafe
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
fe6e548
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
f1d4c0a
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
ac892cd
test(frontstage): bind backup fixture runtime
Duang777 Oct 5, 2026
d0dbf53
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 2026
5a7b79f
Merge remote-tracking branch 'upstream/main' into codex/fix-main-ci-c…
Duang777 Oct 5, 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 @@ -428,7 +428,7 @@ assert.match(machineSettings, /periodicReportActivationDescription/, "Machine pe
assert.match(i18n, /Enabled means automatic delivery at validated stage boundaries/, "English machine settings name automatic stage delivery");
assert.match(i18n, /开启后将在已验证的阶段节点自动投递/, "Chinese machine settings name automatic stage delivery");
assert.match(machineSettings, /localizedCapabilityFieldCopy\(locale\)/, "Machine capability fields follow the selected locale");
assert.match(goalCapabilitySettings, /localizedCapabilityFieldCopy\(locale\)/, "Goal capability fields follow the selected locale");
assert.match(goalCapabilitySettings, /localizedCapabilityFieldCopy\(locale,\s*localizedSelected\.capability_id\)/, "Goal capability fields follow the selected locale and capability");
assert.match(machineSettings, /<CapabilityCatalogNavigation/, "Machine settings use the shared capability catalog navigation");
assert.match(goalCapabilitySettings, /<CapabilityCatalogNavigation/, "Goal settings use the shared capability catalog navigation");
assert.match(goalCapabilitySettings, /capability_id === "lark_event_inbox"[\s\S]*<GoalAutoNotifyToggle/, "Lark inbox capability exposes the independent human-gate notification control");
Expand Down 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_order", "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
5 changes: 2 additions & 3 deletions examples/loopx-turn-codex-cli-e2e-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,10 +45,9 @@ def _marker_value(turn_number: int) -> str:
def _write_fixture(root: Path, *, turn_count: int) -> tuple[Path, Path, Path, Path]:
project = root / "project"
runtime = root / "runtime"
workspace = root / "workspace"
workspace = project
runtime.mkdir(parents=True)
workspace.mkdir(parents=True)
(workspace / "docs").mkdir()
(workspace / "docs").mkdir(parents=True)
state = project / ".codex" / "goals" / GOAL_ID / "ACTIVE_GOAL_STATE.md"
state.parent.mkdir(parents=True)
state.write_text(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ const serverCode = `
import json,pathlib,sys,loopx
from loopx.chat_server import ChatHTTPServer,ChatRequestHandler
r=pathlib.Path(sys.argv[1]); runtime=r/'runtime'; registry=r/'registry.json'
registry.write_text(json.dumps({'goals':[{'id':'fixture','repo':str(r),'control_plane':{'optional':{'context':'complete '*10000}}}]}))
registry.write_text(json.dumps({'common_runtime_root':str(runtime),'goals':[{'id':'fixture','repo':str(r),'control_plane':{'optional':{'context':'complete '*10000}}}]}))
p=runtime/'machine/configuration.json';p.parent.mkdir(parents=True)
p.write_text(json.dumps({'schema_version':'loopx_machine_configuration_v0','namespaces':{'goal_storage':{'schema_version':'loopx_goal_storage_defaults_v0','new_goal_provider':'sqlite'}}}))
s=ChatHTTPServer(('127.0.0.1',0),ChatRequestHandler);s.registry_path=registry;s.runtime_root=runtime;s.verbose=False
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 @@ -60,7 +60,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
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
4 changes: 2 additions & 2 deletions 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 Expand Up @@ -54,7 +54,7 @@ def handle_configuration_backup(args, *, registry_path, print_payload, output_fo
os.link(staging, path)
finally:
staging.unlink(missing_ok=True)
if json.loads(path.read_text()) != backup:
if json.loads(path.read_text(encoding="utf-8")) != backup:
raise RuntimeError("configuration backup export readback mismatch")
payload.update(status="exported", written=True)
else:
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 @@ -9,6 +9,7 @@
import json
import os
import selectors
import signal
import subprocess
import time
import weakref
Expand All @@ -20,6 +21,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 @@ -48,22 +52,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 @@ -76,15 +119,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 @@ -109,7 +158,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 @@ -145,7 +197,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 @@ -156,7 +208,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
27 changes: 27 additions & 0 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
STARTUP_LOCK_TIMEOUT_SECONDS = 15.0
STARTUP_READY_TIMEOUT_SECONDS = 15.0
STARTUP_POLL_SECONDS = 0.025
RUNTIME_RETRY_SETTLE_SECONDS = 0.25
DEFAULT_REQUEST_TIMEOUT_SECONDS = 10.0
# Canonical writers may wait 30 seconds for the per-Goal maintenance lock and
# another 5 seconds for the provider lock. Keep the client connected through
Expand Down Expand Up @@ -461,6 +462,26 @@ def _read_info(path: Path, *, fingerprint: str) -> dict[str, Any] | None:
return payload


def _wait_for_runtime_locator_turnover(
path: Path,
*,
fingerprint: str,
observed: Mapping[str, Any] | None,
timeout: float,
) -> None:
"""Give a retiring runtime time to remove or replace its locator."""

if not isinstance(observed, Mapping):
return
token = observed.get("token")
deadline = time.monotonic() + min(timeout, RUNTIME_RETRY_SETTLE_SECONDS)
while time.monotonic() < deadline:
current = _read_info(path, fingerprint=fingerprint)
if current is None or current.get("token") != token:
return
time.sleep(STARTUP_POLL_SECONDS)


_RUNTIME_IDENTITY_TEXT_FIELDS = (
"node_version",
"sqlite_version",
Expand Down Expand Up @@ -972,6 +993,12 @@ def effect_runtime_request(
# Even a token check followed by unlink would race with a
# replacement server publishing its own locator.
_reap_exited_runtime_child(info)
_wait_for_runtime_locator_turnover(
info_path,
fingerprint=fingerprint,
observed=info,
timeout=timeout,
)
continue
break
if isinstance(last_error, TimeoutError):
Expand Down
3 changes: 2 additions & 1 deletion loopx/control_plane/goals/acceptance_contract.ts
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,8 @@ export function normalizeGoalAcceptanceDocument(value: unknown): AcceptanceDocum
const NON_WORK_FIELDS = new Set([
"schema_version", "source_section", "index", "title", "priority", "status", "done", "archive_state",
"claimed_by", "created_by", "last_actor_agent_id", "updated_at", "completed_at", "completion_turn_key",
"completion_validation_sha256", "completion_recovery", "completion_continuation", "no_followup", "decision_outcome",
"completion_validation_sha256", "completion_recovery", "completion_continuation", "completion_receipt_id",
"no_followup", "decision_outcome",
"completion_result",
"decision_scope_outcomes", "note", "evidence", "reason", "handoff_note", "resume_ready",
"resume_monitor_generation", "last_checked_at", "result_hash", "consecutive_no_change",
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
"""Owner-local configuration download and isolated recovery; never activation."""
from .configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup
from .control_plane.effect_runtime import MAX_LOCAL_SNAPSHOT_BYTES
from ..capabilities.configuration_backup import (
capture_configuration_backup,
restore_configuration_backup,
verify_configuration_backup,
)
from ..control_plane.effect_runtime import MAX_LOCAL_SNAPSHOT_BYTES

CONFIGURATION_BACKUP_PATH = "/api/chat/configuration-backup"

Expand Down
Loading
Loading