Skip to content
2 changes: 1 addition & 1 deletion src/crossagent/check.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
from dataclasses import dataclass
from typing import Optional

from .jobs import CheckResultDict
from .types import CheckResultDict

# Keep only a bounded tail of the check's output — a full test run can emit
# megabytes, and the whole Job record is loaded into memory (and the dashboard).
Expand Down
32 changes: 29 additions & 3 deletions src/crossagent/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from . import __version__
from . import advisors as advisors_mod
from . import check as check_mod
from . import credentials as credentials_mod
from . import jobs as jobs_mod
from . import parsers as parsers_mod
from . import registry as reg
Expand Down Expand Up @@ -109,10 +110,16 @@ def build_command(


def _run_advisor(
cmd: list[str], cwd: str | None, parser_name: str
cmd: list[str],
cwd: str | None,
parser_name: str,
*,
env: dict[str, str] | None = None,
) -> tuple[int, parsers_mod.ParsedResult]:
parser = parsers_mod.get_parser(parser_name)
outcome = runner_mod.run(cmd, cwd=cwd, consumer=parser, max_runtime_seconds=None)
outcome = runner_mod.run(
cmd, cwd=cwd, env=env, consumer=parser, max_runtime_seconds=None
)
parsed = (
outcome.result
if isinstance(outcome.result, parsers_mod.ParsedResult)
Expand Down Expand Up @@ -199,6 +206,18 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser.add_argument(
"--registry", default=str(reg.DEFAULT_REGISTRY), help="Session registry path."
)
parser.add_argument(
"--pass-env",
action="append",
default=[],
dest="pass_env",
help=(
"Name of an environment variable to pass through to the advisor even "
"though it matches a credential pattern (e.g. the advisor's own API "
"key). Repeatable. By default all credential-bearing env vars are "
"withheld from the advisor, matching the durable-job path."
),
)
parser.add_argument(
"--list-advisors", action="store_true", help="Print known advisors and exit."
)
Expand Down Expand Up @@ -278,7 +297,14 @@ def _dispatch(
registry: dict[str, Any],
registry_path: Path,
) -> int:
code, parsed = _run_advisor(cmd, args.cwd, advisor.result_parser)
# A foreground advisor is a delegate too: withhold the caller's ambient
# credentials (S4 policy), sharing the exact scrub the durable-job path uses
# (worker.build_advisor_env delegates to the same helper) so the two dispatch
# modes can never drift. ``--pass-env NAME`` is the caller's opt-out.
advisor_env = credentials_mod.scrub_env(
os.environ, pass_through=getattr(args, "pass_env", [])
)
code, parsed = _run_advisor(cmd, args.cwd, advisor.result_parser, env=advisor_env)

if parsed.failure:
error = parsed.error or f"{advisor.name} exited with code {code}"
Expand Down
163 changes: 132 additions & 31 deletions src/crossagent/escalate.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,11 +111,12 @@ def maybe_escalate(
"""Re-dispatch a *failed* delegation to the next ladder rung, if any.

Returns the spawned child's job id, or ``None`` when nothing was escalated
(the delegation did not fail, no rungs remain, the rung advisor is unknown,
or the depth cap was reached). Never raises: an escalation that cannot be
launched is recorded and skipped, never allowed to crash the worker.
(the delegate did not finish cleanly, no declared gate failed, no rungs
remain, the rung advisor is unknown, or the depth cap was reached). Never
raises: an escalation that cannot be launched is recorded and skipped, never
allowed to crash the worker.
"""
if delegation_verdict(failed_job) != "failed":
if not _is_escalatable_failure(failed_job):
return None

rungs = parse_rungs(escalate_to)
Expand All @@ -134,7 +135,7 @@ def maybe_escalate(

child_id = jobs_mod.generate_job_id()
try:
parent_id, trace_id, label, depth = jobs_mod.resolve_lineage(
lineage = jobs_mod.resolve_lineage(
parent_flag=failed_job.job_id,
state_root=state_root,
new_job_id=child_id,
Expand All @@ -145,10 +146,11 @@ def maybe_escalate(
_audit_skip(job_dir, reason=f"escalation halted: {exc}")
return None

child_dir = jobs_mod.create_job_dir(state_root, child_id)
_write_child_prompt(child_dir, prompt)
_write_child_command(
child_dir,
if not _create_and_save_child(
state_root=state_root,
job_dir=job_dir,
child_id=child_id,
prompt=prompt,
advisor=advisor,
model=model,
cwd=cwd,
Expand All @@ -160,27 +162,12 @@ def maybe_escalate(
verify_with=verify_with,
verify_model=verify_model,
escalate_to=remaining,
)
child = Job(
job_id=child_id,
status=jobs_mod.JobState.PENDING,
advisor=advisor.name,
name="",
cwd=cwd,
redacted_command="",
started_at=_now(),
updated_at=_now(),
last_activity_at=_now(),
last_event="escalation.created",
max_runtime_seconds=failed_job.max_runtime_seconds,
termination_grace_seconds=failed_job.termination_grace_seconds,
parent_job_id=parent_id,
trace_id=trace_id,
orchestrator_label=label,
nesting_depth=depth,
)
jobs_mod.save_state(child_dir, child)
lineage=lineage,
failed_job=failed_job,
):
return None

_, trace_id, _, depth = lineage
jobs_mod.append_event(
job_dir,
"escalation",
Expand All @@ -206,6 +193,33 @@ def maybe_escalate(
# ---------------------------------------------------------------------------


def _is_escalatable_failure(failed_job: Job) -> bool:
"""Return True only for a delegate that FINISHED but FAILED a declared gate.

Escalation re-dispatches work that failed a check, a scope violation, or a
failing verification (this module's stated scope). A job that did not finish
cleanly is deliberately out of scope, so gate on SUCCEEDED first rather than
on ``delegation_verdict != "failed"`` alone (which also returns "failed" for
CANCELLED / TIMED_OUT / a crashed delegate):

* CANCELLED is explicit user intent to stop; re-dispatching to a larger,
costlier peer is the opposite of cancelling and spends real money.
* TIMED_OUT (and any other non-success terminal status) produced no graded
artifact — there is no gate failure to escalate, and a bigger model is
generally slower, so it is at least as likely to time out again under the
same budget. A hard task that needs a bigger model is a fresh dispatch
decision, not an automatic ladder climb that silently burns budget. So
TIMED_OUT does NOT escalate.

Gating on SUCCEEDED means ``delegation_verdict == "failed"`` can only be a
declared-gate failure — mirroring cli._failed_reason's "did not finish
cleanly" vs. gate-failure distinction.
"""
if failed_job.status != jobs_mod.JobState.SUCCEEDED:
return False
return delegation_verdict(failed_job) == "failed"


def _drop_first_nonblank(rungs: Optional[list[str]]) -> list[str]:
"""Return the non-blank rungs with the first one (the spawned rung) removed."""
cleaned = [raw for raw in (rungs or []) if raw.strip()]
Expand Down Expand Up @@ -236,7 +250,7 @@ def _write_child_command(
verify_model: Optional[str],
escalate_to: list[str],
) -> None:
info = {
command_payload = {
"command": _build_child_argv(advisor, model),
"prompt_delivery": advisor.prompt_delivery,
"cwd": cwd,
Expand All @@ -258,7 +272,94 @@ def _write_child_command(
"verify_model": verify_model,
"escalate_to": escalate_to,
}
jobs_mod.atomic_json_write(info, child_dir / "command.json")
jobs_mod.atomic_json_write(command_payload, child_dir / "command.json")


def _build_child_job(
child_id: str,
advisor: Advisor,
cwd: str,
lineage: tuple[Optional[str], str, Optional[str], Optional[int]],
failed_job: Job,
) -> Job:
"""Assemble the PENDING child ``Job`` record for an escalation re-dispatch.

*lineage* is the ``(parent_id, trace_id, label, depth)`` tuple returned by
:func:`~crossagent.jobs.resolve_lineage`; runtime bounds are inherited from
*failed_job* so the larger peer runs under the same limits.
"""
parent_id, trace_id, label, depth = lineage
return Job(
job_id=child_id,
status=jobs_mod.JobState.PENDING,
advisor=advisor.name,
name="",
cwd=cwd,
redacted_command="",
started_at=_now(),
updated_at=_now(),
last_activity_at=_now(),
last_event="escalation.created",
max_runtime_seconds=failed_job.max_runtime_seconds,
termination_grace_seconds=failed_job.termination_grace_seconds,
parent_job_id=parent_id,
trace_id=trace_id,
orchestrator_label=label,
nesting_depth=depth,
)


def _create_and_save_child(
*,
state_root: Path,
job_dir: Path,
child_id: str,
prompt: str,
advisor: Advisor,
model: Optional[str],
cwd: str,
registry_path: str,
check: Optional[str],
check_timeout: float,
scope_paths: Optional[list[str]],
pass_env: list[str],
verify_with: Optional[str],
verify_model: Optional[str],
escalate_to: list[str],
lineage: tuple[Optional[str], str, Optional[str], Optional[int]],
failed_job: Job,
) -> bool:
"""Stage the escalated child on disk: prompt, command, and state record.

Staging touches the filesystem (mkdir, two file writes, a state save). A
read-only or full disk raises OSError; it is caught here so the caller's
"never raises" contract holds — a child that cannot be staged is recorded on
*job_dir* and skipped, exactly like a launch failure. Returns True on
success, False when staging was skipped.
"""
try:
child_dir = jobs_mod.create_job_dir(state_root, child_id)
_write_child_prompt(child_dir, prompt)
_write_child_command(
child_dir,
advisor=advisor,
model=model,
cwd=cwd,
registry_path=registry_path,
check=check,
check_timeout=check_timeout,
scope_paths=scope_paths,
pass_env=pass_env,
verify_with=verify_with,
verify_model=verify_model,
escalate_to=escalate_to,
)
child = _build_child_job(child_id, advisor, cwd, lineage, failed_job)
jobs_mod.save_state(child_dir, child)
except OSError as exc:
_audit_skip(job_dir, reason=f"escalation could not be staged: {exc}")
return False
return True


def _audit_skip(job_dir: Path, *, reason: str) -> None:
Expand Down
92 changes: 8 additions & 84 deletions src/crossagent/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,97 +16,21 @@
from datetime import datetime, timezone
from enum import Enum
from pathlib import Path
from typing import Any, Callable, Literal, Optional, TypedDict
from typing import Any, Callable, Literal, Optional

# The gate-result shapes now live in the dependency-free ``types`` module.
# ``jobs`` uses these three as ``Job`` field annotations, which also re-exports
# them for callers that still reference ``jobs_mod.CheckResultDict`` (e.g.
# ``worker``). ``ScopeStatus``/``VerifyVerdict`` are imported straight from
# ``types`` by ``scope``/``verify``.
from .types import CheckResultDict, ScopeResultDict, VerifyResultDict

# Where a job's cost figure came from. ``advisor`` is the vendor-declared
# estimate, ``computed`` is derived from a local price table, and ``unknown``
# means no cost was measured — never conflate that with a measured ``0.0``.
CostSource = Literal["advisor", "computed", "unknown"]


class CheckResultDict(TypedDict):
"""Persisted outcome of the independent check-gate (slice S3).

``command`` is the caller-supplied check string; ``exit_code`` is the
deterministic delegation verdict (D5) — ``0`` means the work verified, any
non-zero value means the delegation *failed* regardless of whether the
delegate process itself exited 0. ``stdout_tail``/``stderr_tail`` are
bounded tails of the check's output.

A ``Job.check_result`` of ``None`` means *no check ran* (unverified) — that
stays structurally distinct from a check that ran and failed (D7): absent is
never the same as false.
"""

command: str
exit_code: int
stdout_tail: str
stderr_tail: str


# Outcome of the diff-scope assertion (slice S4). ``ok`` — every path the
# delegate modified was inside the declared allowlist; ``violated`` — it wrote
# outside its declared scope; ``undetermined`` — crossagent could not establish
# what changed (e.g. the cwd is not a git repo). ``undetermined`` is FAIL-CLOSED,
# never a pass: a scope check that silently passes when it cannot see the changes
# grants false assurance.
ScopeStatus = Literal["ok", "violated", "undetermined"]


class ScopeResultDict(TypedDict):
"""Persisted outcome of the diff-scope assertion (slice S4).

``declared`` is the caller-supplied allowlist; ``violating_paths`` lists the
repo-relative paths the delegate modified outside it (empty unless
``status == "violated"``); ``detail`` is a human-readable summary.

A ``Job.scope_result`` of ``None`` means *no scope was declared* — scope
enforcement was off — which stays structurally distinct from a scope that
was declared and satisfied (D7: absent is never the same as ``ok``).
"""

declared: list[str]
status: ScopeStatus
violating_paths: list[str]
detail: str


# Outcome of the independent verification pass (slice S5). A FRESH peer session
# grades the delegate's artifact supplied as user-turn input (D6), which removes
# the implicit-authorship channel that weakens self-grading — it does NOT claim
# to eliminate self-preference bias, which is a separate documented effect.
# ``pass`` — the verifier returned a machine-checkable verdict of correct.
# ``fail`` — the verifier returned a machine-checkable verdict of wrong.
# ``unverified`` — the verifier ran but produced no machine-checkable verdict
# (free prose, or the advisor lacks a structured-output
# contract). Treated as inconclusive, never as a pass (D4).
# ``error`` — the verifier could not run (no artifact, launch failure).
# ``unverified``/``error`` are inconclusive: they never green a delegation and
# never hard-fail it. Only ``fail`` blocks the green path.
VerifyVerdict = Literal["pass", "fail", "unverified", "error"]


class VerifyResultDict(TypedDict):
"""Persisted outcome of the independent verification pass (slice S5).

``advisor``/``model`` identify the FRESH peer session that graded the work.
``verdict`` is the deterministic pass/fail/inconclusive outcome; ``structured``
records whether a machine-checkable contract (e.g. Claude ``--json-schema`` →
``structured_output``) produced it, versus a JSON object parsed out of a prose
answer. ``detail`` is a human-readable summary (never the raw artifact).

A ``Job.verify_result`` of ``None`` means *no verification was requested* —
structurally distinct from a verification that ran and failed, or ran and
could not decide (D7: absent is never the same as ``fail`` or ``unverified``).
"""

advisor: str
model: Optional[str]
verdict: VerifyVerdict
structured: bool
detail: str


# The delegation verdict keeps "the delegate finished" separate from "the work
# was verified" (D5). See :func:`delegation_verdict`.
DelegationVerdict = Literal["verified", "failed", "unverified", "incomplete"]
Expand Down
2 changes: 1 addition & 1 deletion src/crossagent/scope.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@
from pathlib import Path
from typing import Optional

from .jobs import ScopeResultDict, ScopeStatus
from .types import ScopeResultDict, ScopeStatus

_GIT_TIMEOUT_SECONDS = 30.0
# porcelain -z entries are "XY <path>": two status chars, a space, then the path.
Expand Down
Loading