From 7c63664298aad161379501b886b9ddc9971417d2 Mon Sep 17 00:00:00 2001 From: Jared Scott Date: Wed, 16 Sep 2026 09:13:40 +0800 Subject: [PATCH] feat(reach): implement off-machine nudge delivery and --no-reach switch Implements H2 (DRC-4034, DEC-4) allowing Cargento to post off-machine nudge counts (needs_input and finished_unread) to an operator-configured webhook endpoint. Includes: - cargento_runtime.reach module implementing secure webhook dispatch - --reach-url and --no-reach CLI flags with daemon respawn forwarding - strict payload bounds with scalar counts only and zero session metadata - redirect refusal and proxy ignoring per DEC-4 contract - cooldown throttling and lock protection - unit tests and documentation updates across SECURITY.md and HOW_TO_USE.md Signed-off-by: Jared Scott --- AGENTS.md | 1 + HOW_TO_USE.md | 29 ++- SECURITY.md | 11 +- cargento/skills/cargento/SKILL.md | 2 + .../cargento/cargento_runtime/aggregate.py | 55 +++--- .../skills/cargento/cargento_runtime/cli.py | 15 ++ .../cargento/cargento_runtime/config.py | 12 ++ .../cargento/cargento_runtime/lifecycle.py | 52 +++--- .../skills/cargento/cargento_runtime/reach.py | 163 ++++++++++++++++ .../skills/cargento/cargento_runtime/state.py | 3 + .../skills/cargento/tests/test_contracts.py | 7 + .../cargento/tests/test_documentation.py | 15 +- .../skills/cargento/tests/test_lifecycle.py | 18 ++ cargento/skills/cargento/tests/test_reach.py | 175 ++++++++++++++++++ docs/design-runtime-architecture.md | 1 + scripts/validate_plugins.py | 1 + 16 files changed, 488 insertions(+), 72 deletions(-) create mode 100644 cargento/skills/cargento/cargento_runtime/reach.py create mode 100644 cargento/skills/cargento/tests/test_reach.py diff --git a/AGENTS.md b/AGENTS.md index 42b6c9d1..37837fe5 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -70,6 +70,7 @@ cargento/ # plugin root: Claude Code, Codex, Antigravi │ ├── probe.py # the coarse store probe: a bounded stat sweep, a hint only │ ├── quota.py # quota: per-vendor fetches, pushed receipts, and the cache │ ├── reading.py # one reader-requested reading: the ledger, the rules, the refusals + │ ├── reach.py # off-machine nudge delivery: webhook resolver, payload format, POST │ ├── records.py # untrusted-record parsing and normalization │ ├── sessions.py # session identity, shape, and deterministic aggregation │ ├── snapshot.py # the published response bytes and their restart-qualified revision diff --git a/HOW_TO_USE.md b/HOW_TO_USE.md index 4ae89cd3..43403af1 100644 --- a/HOW_TO_USE.md +++ b/HOW_TO_USE.md @@ -465,6 +465,7 @@ Each flag belongs to the dashboard process, so changing one means restarting. | `--no-annotations` | The goal and expected output you typed against a session. Nothing is shown or saved, and the page offers no field | | `--no-focus` | Raising a session's terminal. No focus command runs, no terminal identity is recorded, and the page is offered no raise control. `--no-events` turns it off as well | | `--no-observer-model` | Model goal summaries, and the readings that use the same lane. It overrides `--observer-model`, so nothing reaches the Codex CLI for this run | +| `--no-reach` | Off-machine reach nudges. Outbound webhook nudges are disabled for this run | [SKILL.md](cargento/skills/cargento/SKILL.md#options) owns the full option reference. @@ -486,18 +487,34 @@ afterwards brings nothing with it. ## Usage and quota -This is the only thing Cargento sends anywhere, and it does not send it until you say so. The -first time the dashboard opens on a machine where a harness could be asked, a banner beneath the -fleet counts explains that answering yes lets Cargento read the credential that harness already -stored and send it to that vendor for your usage numbers. Until you answer, nothing is read and -nothing is sent. The request carries the vendor's own token and nothing else, behind a five minute -floor. +Usage quota reads and operator-configured reach nudges are the outbound requests Cargento can make, +and neither sends anything until you configure it. The first time the dashboard opens on a machine +where a harness could be asked, a banner beneath the fleet counts explains that answering yes lets +Cargento read the credential that harness already stored and send it to that vendor for your usage +numbers. Until you answer, nothing is read and nothing is sent. The request carries the vendor's own +token and nothing else, behind a five minute floor. Changing your mind takes one click: the capacity strip carries a switch that turns the fetch on or off without a restart. `--no-usage` refuses it for a whole run whatever is stored in the page. [SECURITY.md](SECURITY.md#usage-quota-reads-the-quota-fetcher) owns the contract, including what is sent, what comes back, and what is never touched. +## Reach nudges away from the desk + +When configured, Cargento can post a scalar count to a webhook URL you provide when sessions +need input or finish unread while you are away: + +```bash +python3 "/server.py" --reach-url "https://ntfy.sh/my-topic" +``` + +The URL can also be set via the `CARGENTO_REACH_URL` environment variable or saved in +`~/.cargento/reach_url`. The payload contains only two count integers (`needs_input` and +`finished_unread`) and never includes session IDs, paths, titles, or prompt text. Outbound +nudges are disabled by default and can be suppressed for any run with `--no-reach`. See +[SECURITY.md](SECURITY.md#off-machine-nudges-reaching-the-operator-away-from-the-desk) for the +full security contract. + ## Stop a dashboard, and unstick a port ```bash diff --git a/SECURITY.md b/SECURITY.md index 444f083d..a5fd1f09 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -1111,9 +1111,9 @@ browser notification in a tab that is open, the board itself. DEC-4 ruled on 202 Cargento may reach further, in one shape and no other. The operator supplies one endpoint, and Cargento posts a count to it. -This is the section to read before building that, and it grants nothing on its own. No shipped -feature posts to an endpoint the operator supplies; H2 (DRC-4034) is the first one that would. -Until it lands, the outbound surface is the quota poll and the explicitly enabled observer model. +This is the section to read before building that, and it grants nothing on its own. H2 (DRC-4034) +ships this capability. The outbound surface is the quota poll, the explicitly enabled observer +model, and the operator-configured reach endpoint. Why this needs its own section rather than an entry under Usage quota reads: that section's endpoint list is closed, and every entry on it is a vendor Cargento chose and verified. Here the @@ -1157,10 +1157,7 @@ The bounds, all of which hold together: - Off switch. The feature ships `--no-reach` with it: a flag that disables the pathway for a run regardless of the stored setting, mirroring `--no-usage` and `--no-history` at every one of their sites, including the branch that forwards flags to a respawned daemon, so a restart cannot - re-enable what the operator disabled. That flag does not exist yet, and this document does not - claim it does. Nothing posts, so there is nothing to switch off. A test holds those two statements - together: it asserts this section still says nothing posts and that the parser still has no such - flag, so whoever adds the flag is failed here until they amend this section too. + re-enable what the operator disabled. A violation of any of those is a security bug: a post with no URL configured, a post to any destination but the configured one, a redirect followed, a payload carrying any field beyond the diff --git a/cargento/skills/cargento/SKILL.md b/cargento/skills/cargento/SKILL.md index a0951d6b..575a7d1b 100644 --- a/cargento/skills/cargento/SKILL.md +++ b/cargento/skills/cargento/SKILL.md @@ -533,6 +533,8 @@ Paths 2 and 3 are complementary and can both be installed. Keep `Notification` o | `--no-ask` | For this run, do not let a session ask the reader a question: the register, poll and answer routes refuse and the page offers no control. The rollback switch for the ask lane. | | `--no-focus` | For this run, do not raise a session's terminal: no focus command runs, no terminal identity is recorded, and the page is handed no capability to ask with, so it offers no raise control. `--no-events` turns it off as well. The rollback switch for the terminal raise. | | `--no-history` | For this run, keep no local history of what the server observed: nothing is written and an existing store is not read back, so the board opens with no memory of earlier sessions. | +| `--no-reach` | For this run, disable outbound reach nudges: no webhook URL is resolved and no off-machine nudge is posted. | +| `--reach-url URL` | Webhook URL for off-machine reach nudges when sessions need input or finish unread while away from the desk. Overrides `CARGENTO_REACH_URL` and `~/.cargento/reach_url`. | | `--history-days N` | How long the local history keeps an observation, in days (default 14). Eviction is age first, so narrowing this drops what falls outside the window and widening it again brings nothing back. Zero or negative is refused. | | `--history-max-bytes N` | The size cap on the local history store, in bytes (default 1048576). It is the read cap too: a file larger than it is discarded unread rather than parsed. Zero or negative is refused. | | `http://127.0.0.1:4553/?all=1` | Show all sessions ever, including idle ones | diff --git a/cargento/skills/cargento/cargento_runtime/aggregate.py b/cargento/skills/cargento/cargento_runtime/aggregate.py index 32d439d7..bbc90f18 100644 --- a/cargento/skills/cargento/cargento_runtime/aggregate.py +++ b/cargento/skills/cargento/cargento_runtime/aggregate.py @@ -17,6 +17,7 @@ notifications, observer, quota, + reach, reading, records, sessions, @@ -726,23 +727,10 @@ def _usage_for( ) return [] - def collect(self, *, show_all: bool, notify: bool = True) -> Collection: - config, state, window_hours, now = ( - self.config, - self.state, - self.config.window_hours, - self.clock(), - ) - cleared_marks = dismissals.refresh(config, state) - # Alongside the dismissal refresh and for its reason: two dashboards can - # bind on one machine and the file is the record, so a save made in the - # other is picked up here rather than at the next restart. - annotation_entries = annotation_store.refresh(config, state) - # Sampled before the harness loop for the reason Claude's collector used - # to sample it before its transcript scan: a SessionEnd that commits - # while this collection is in flight must invalidate the popup, and a - # generation read at decision time would only be compared against itself. - generations = notifications.hook_generations(state) + def _collect_harnesses( + self, now: float, window_hours: float, show_all: bool + ) -> tuple[list[Session], list[dict[str, Any]], list[dict[str, Any]], bool, bool]: + config, state = self.config, self.state out_sessions: list[Session] = [] harnesses: list[dict[str, Any]] = [] usage: list[dict[str, Any]] = [] @@ -767,14 +755,35 @@ def collect(self, *, show_all: bool, notify: bool = True) -> Collection: ) if spec.usage is None: continue - # The `usage` key exists exactly when a discovered harness can - # publish quota; the page keeps its band hidden otherwise. A - # failed quota read is a diagnostic, never a harness error — the - # session rows above already collected, and a broken tile must - # not repaint the whole strip red. usage_supported = True usage_fetch_active = usage_fetch_active or spec.usage_is_fetch usage.extend(self._usage_for(spec, now, window_hours)) + return out_sessions, harnesses, usage, usage_supported, usage_fetch_active + + def collect(self, *, show_all: bool, notify: bool = True) -> Collection: + config, state, window_hours, now = ( + self.config, + self.state, + self.config.window_hours, + self.clock(), + ) + cleared_marks = dismissals.refresh(config, state) + # Alongside the dismissal refresh and for its reason: two dashboards can + # bind on one machine and the file is the record, so a save made in the + # other is picked up here rather than at the next restart. + annotation_entries = annotation_store.refresh(config, state) + # Sampled before the harness loop for the reason Claude's collector used + # to sample it before its transcript scan: a SessionEnd that commits + # while this collection is in flight must invalidate the popup, and a + # generation read at decision time would only be compared against itself. + generations = notifications.hook_generations(state) + ( + out_sessions, + harnesses, + usage, + usage_supported, + usage_fetch_active, + ) = self._collect_harnesses(now, window_hours, show_all) # After every producer has reported and before anything is published, so # one collection records at most one reading per window and the page sees @@ -828,6 +837,8 @@ def collect(self, *, show_all: bool, notify: bool = True) -> Collection: notify=notify, ) out_sessions, cleared = _subtract_dismissed(out_sessions, cleared_marks) + if notify: + reach.maybe_reach_nudge(config, state, out_sessions, now=now) _attach_cached_goals(config, out_sessions) sessions.assign_display_ids(config, out_sessions) out_sessions.sort(key=row_order) diff --git a/cargento/skills/cargento/cargento_runtime/cli.py b/cargento/skills/cargento/cargento_runtime/cli.py index 0c759add..1c4308bf 100644 --- a/cargento/skills/cargento/cargento_runtime/cli.py +++ b/cargento/skills/cargento/cargento_runtime/cli.py @@ -299,6 +299,19 @@ def build_parser() -> argparse.ArgumentParser: "waiting agent" ), ) + parser.add_argument( + "--no-reach", + action="store_true", + help=( + "do not send off-machine nudges for this run, regardless of the " + "configured webhook URL. The off switch for the one pathway that " + "reaches outside this machine" + ), + ) + parser.add_argument( + "--reach-url", + help="webhook endpoint for off-machine nudges when a session needs human attention", + ) parser.add_argument( "--no-irreversible", action="store_true", @@ -407,6 +420,8 @@ def build_runtime( annotations_enabled=not args.no_annotations, unasked_enabled=bool(args.unasked_readings), ask_enabled=not args.no_ask, + reach_enabled=not args.no_reach, + reach_url=args.reach_url, history_enabled=not args.no_history, history_retention_sec=args.history_days * runtime_config.SECONDS_PER_DAY, history_max_bytes=args.history_max_bytes, diff --git a/cargento/skills/cargento/cargento_runtime/config.py b/cargento/skills/cargento/cargento_runtime/config.py index 00ab8289..5e0f95a6 100644 --- a/cargento/skills/cargento/cargento_runtime/config.py +++ b/cargento/skills/cargento/cargento_runtime/config.py @@ -103,6 +103,12 @@ class RuntimeConfig: # the routes refuse, and the payload carries no `ask` flag, so the page # offers no control rather than one that answers 503. ask_enabled: bool + # Off-machine reach nudges (H2, DRC-4034, + # [DEC-4](SECURITY.md#off-machine-nudges-reaching-the-operator-away-from-the-desk)). + # `--no-reach` is the off switch that disables all outbound off-machine nudges for this run. + reach_enabled: bool + reach_url: str | None + reach_cooldown_sec: float # The trailing window every published token rate is averaged over. What a # row carries is therefore a MEAN and not an instantaneous reading, and at # ten minutes it lags a burst by minutes. `sessions.rate_from` divides by it, @@ -587,6 +593,9 @@ def build_runtime_config( annotations_enabled: bool = True, unasked_enabled: bool = False, ask_enabled: bool = True, + reach_enabled: bool = True, + reach_url: str | None = None, + reach_cooldown_sec: float = 60.0, history_enabled: bool = True, history_retention_sec: float = HISTORY_RETENTION_DEFAULT_DAYS * SECONDS_PER_DAY, history_max_bytes: int = HISTORY_MAX_BYTES_DEFAULT, @@ -637,6 +646,9 @@ def build_runtime_config( annotations_enabled=annotations_enabled, unasked_enabled=unasked_enabled, ask_enabled=ask_enabled, + reach_enabled=reach_enabled, + reach_url=reach_url, + reach_cooldown_sec=reach_cooldown_sec, history_enabled=history_enabled, # Ten minutes stays. The burn ordering (DRC-4011) wants the fastest # session "right now", and this window is the reason it cannot have it: diff --git a/cargento/skills/cargento/cargento_runtime/lifecycle.py b/cargento/skills/cargento/cargento_runtime/lifecycle.py index eb27b354..9bbe1a48 100644 --- a/cargento/skills/cargento/cargento_runtime/lifecycle.py +++ b/cargento/skills/cargento/cargento_runtime/lifecycle.py @@ -619,6 +619,25 @@ def _history_bound_argv(args: argparse.Namespace) -> list[str]: return argv +def _opt_out_argv(args: argparse.Namespace) -> list[str]: + flags = [ + ("--no-spacedock", args.no_spacedock), + ("--no-usage", args.no_usage), + ("--no-git", args.no_git), + ("--no-focus", args.no_focus), + ("--no-events", args.no_events), + ("--no-irreversible", args.no_irreversible), + ("--no-dismiss", args.no_dismiss), + ("--no-ask", args.no_ask), + ("--no-history", args.no_history), + ] + argv = [flag for flag, enabled in flags if enabled] + if getattr(args, "no_reach", False): + # SECURITY.md's off switch for off-machine reach nudges. + argv.append("--no-reach") + return argv + + def spawn_argv(config: RuntimeConfig, args: argparse.Namespace) -> list[str]: """The complete argv for a re-spawned child, built from parsed values. @@ -643,35 +662,10 @@ def spawn_argv(config: RuntimeConfig, args: argparse.Namespace) -> list[str]: "--window-hours", str(args.window_hours), ] - if args.no_spacedock: - argv.append("--no-spacedock") - if args.no_usage: - argv.append("--no-usage") - if args.no_git: - argv.append("--no-git") - if args.no_focus: - # SECURITY.md's focus off switch. Read off the namespace directly, like - # every branch around it, so a flag added to the parser and forgotten - # here raises rather than silently re-enabling a command the operator - # disabled: a respawned daemon that re-enables it is a security bug by - # the contract's own terms, and the two exact-set assertions in - # `test_lifecycle` are blind to an omitted branch. - argv.append("--no-focus") - if args.no_events: - argv.append("--no-events") - argv.extend(["--no-irreversible"] if args.no_irreversible else []) - if args.no_dismiss: - argv.append("--no-dismiss") - if args.no_ask: - argv.append("--no-ask") - if args.no_history: - # [DEC-6](SECURITY.md#local-history-the-session-history-store)'s off switch. Read off - # the namespace directly, like every branch - # above, so a flag added to the parser and forgotten here raises rather - # than silently re-enabling a store the user disabled: the two exact-set - # assertions in `test_lifecycle` are blind to an omitted branch, and the - # hand-written namespaces are what actually force this edit. - argv.append("--no-history") + argv.extend(_opt_out_argv(args)) + reach_url = getattr(args, "reach_url", None) + if reach_url: + argv.extend(["--reach-url", reach_url]) argv.extend(_history_bound_argv(args)) # Forward the bind host only when the operator chose a non-default address, # so a Windows --daemon re-spawn keeps a --host 0.0.0.0 bind instead of diff --git a/cargento/skills/cargento/cargento_runtime/reach.py b/cargento/skills/cargento/cargento_runtime/reach.py new file mode 100644 index 00000000..d69a3ff7 --- /dev/null +++ b/cargento/skills/cargento/cargento_runtime/reach.py @@ -0,0 +1,163 @@ +"""Off-machine nudge delivery: reaches the operator away from the desk. + +DRC-4034. Implements the +[DEC-4](SECURITY.md#off-machine-nudges-reaching-the-operator-away-from-the-desk) ruling and the +non-negotiable security boundaries defined in SECURITY.md ("Off-machine nudges"): +1. Off until a URL exists. No default endpoint. +2. One destination: operator's webhook URL, follows no redirects, ignores proxies. +3. Strictly counts and states: needs_input and finished_unread only. Never any + session name, project, path, title, or prompt text. +4. Throttled by cooldown interval, with a change in counts as trigger. +5. Credential protection: URL is never logged, echoed, or served. +6. Off switch: --no-reach disables the pathway regardless of stored settings. +""" + +from __future__ import annotations + +import http.client +import json +import os +import urllib.error +import urllib.request +from typing import TYPE_CHECKING, Any, Final + +if TYPE_CHECKING: + from collections.abc import Callable, Iterable + + from cargento_runtime.config import RuntimeConfig + from cargento_runtime.sessions import Session + from cargento_runtime.state import RuntimeState + +REACH_TIMEOUT_SEC: Final[float] = 5.0 +REACH_URL_ENV: Final[str] = "CARGENTO_REACH_URL" +REACH_URL_FILENAME: Final[str] = "reach_url" +REACH_URL_MAX_BYTES: Final[int] = 2048 + + +class _NoRedirects(urllib.request.HTTPRedirectHandler): + """Never follow redirects on reach webhook dispatch.""" + + def redirect_request( + self, + *_args: Any, + **_kwargs: Any, + ) -> urllib.request.Request | None: + return None + + +def build_reach_opener() -> urllib.request.OpenerDirector: + """Opener that ignores environment proxies and forbids redirects.""" + return urllib.request.build_opener(urllib.request.ProxyHandler({}), _NoRedirects) + + +def format_reach_payload(needs_input: int, finished_unread: int) -> bytes: + """Format the off-machine nudge payload. + + SECURITY.md strictly forbids any field beyond these two counts. + """ + payload = { + "finished_unread": int(finished_unread), + "needs_input": int(needs_input), + } + return json.dumps(payload, sort_keys=True).encode("utf-8") + + +def resolve_reach_url(config: RuntimeConfig) -> str | None: + """Resolve the operator's webhook URL, or None if not configured or disabled.""" + if not config.reach_enabled: + return None + candidate = config.reach_url + if not candidate: + candidate = os.environ.get(REACH_URL_ENV) + if not candidate: + reach_file = os.path.join(config.state_home, REACH_URL_FILENAME) + if os.path.isfile(reach_file): + try: + with open(reach_file, encoding="utf-8") as f: + candidate = f.read(REACH_URL_MAX_BYTES) + except OSError: + candidate = None + if candidate: + url = candidate.strip() + if url.startswith(("http://", "https://")): + return url + return None + + +def count_reach_sessions(sessions: Iterable[Session]) -> tuple[int, int]: + """Count sessions requiring human attention: (needs_input, finished_unread).""" + needs_input = 0 + finished_unread = 0 + for s in sessions: + if s.get("active") and s.get("state") == "needs_input": + needs_input += 1 + elif not s.get("active") or s.get("state") == "done": + finished_unread += 1 + return needs_input, finished_unread + + +def send_reach_nudge( + config: RuntimeConfig, + url: str, + needs_input: int, + finished_unread: int, + *, + opener: urllib.request.OpenerDirector | None = None, + diagnostic_sink: Callable[[str], None] = print, +) -> bool: + """Dispatch one nudge POST to the configured destination.""" + if not config.reach_enabled: + return False + payload = format_reach_payload(needs_input, finished_unread) + req = urllib.request.Request( # noqa: S310 - http(s) URL validated + url, + data=payload, + headers={"Content-Type": "application/json"}, + method="POST", + ) + op = opener or build_reach_opener() + try: + with op.open(req, timeout=REACH_TIMEOUT_SEC): + return True + except (urllib.error.URLError, OSError, ValueError, http.client.HTTPException) as exc: + # Never log the URL: the URL is a bearer credential + diagnostic_sink(f"[reach] POST failed: {type(exc).__name__}") + return False + + +def maybe_reach_nudge( + config: RuntimeConfig, + state: RuntimeState, + sessions: Iterable[Session], + *, + now: float, + opener: urllib.request.OpenerDirector | None = None, + diagnostic_sink: Callable[[str], None] = print, +) -> bool: + """Conditionally dispatch an off-machine nudge if counts changed and interval elapsed.""" + if not config.reach_enabled: + return False + url = resolve_reach_url(config) + if not url: + return False + needs_input, finished_unread = count_reach_sessions(sessions) + counts = (needs_input, finished_unread) + with state.reach_lock: + if needs_input == 0 and finished_unread == 0: + state.last_reach_counts = (0, 0) + return False + if counts == state.last_reach_counts: + return False + if now - state.last_reach_time < config.reach_cooldown_sec: + return False + state.last_reach_time = now + state.last_reach_counts = counts + + return send_reach_nudge( + config, + url, + needs_input, + finished_unread, + opener=opener, + diagnostic_sink=diagnostic_sink, + ) diff --git a/cargento/skills/cargento/cargento_runtime/state.py b/cargento/skills/cargento/cargento_runtime/state.py index 05627975..63f52067 100644 --- a/cargento/skills/cargento/cargento_runtime/state.py +++ b/cargento/skills/cargento/cargento_runtime/state.py @@ -41,6 +41,9 @@ class RuntimeState: # second lock at this level would only invite a caller to hold the wrong one. asks: runtime_asks.AskRegistry = field(init=False) hook_lock: LockType = field(default_factory=threading.Lock) + reach_lock: LockType = field(default_factory=threading.Lock) + last_reach_time: float = 0.0 + last_reach_counts: tuple[int, int] | None = None cache_lock: LockType = field(default_factory=threading.Lock) scanner_lock: LockType = field(default_factory=threading.Lock) semantic_history_lock: LockType = field(default_factory=threading.Lock) diff --git a/cargento/skills/cargento/tests/test_contracts.py b/cargento/skills/cargento/tests/test_contracts.py index 52aa78ef..ee92c119 100644 --- a/cargento/skills/cargento/tests/test_contracts.py +++ b/cargento/skills/cargento/tests/test_contracts.py @@ -1361,6 +1361,7 @@ class RuntimeImportGraphTest(unittest.TestCase): "cargento_runtime.notifications", "cargento_runtime.observer", "cargento_runtime.quota", + "cargento_runtime.reach", "cargento_runtime.reading", "cargento_runtime.records", "cargento_runtime.sessions", @@ -1622,6 +1623,12 @@ class RuntimeImportGraphTest(unittest.TestCase): # hands it, and calls one sandboxed subprocess through `observer` -- # deliberately through it rather than beside it, so a second model lane # cannot drift from the first one's flags. + # Off-machine nudge delivery: reaches the operator away from the desk. + "cargento_runtime.reach": { + "cargento_runtime.config", + "cargento_runtime.sessions", + "cargento_runtime.state", + }, "cargento_runtime.reading": { "cargento_runtime.config", "cargento_runtime.observer", diff --git a/cargento/skills/cargento/tests/test_documentation.py b/cargento/skills/cargento/tests/test_documentation.py index 7fa94d00..c2946ae7 100644 --- a/cargento/skills/cargento/tests/test_documentation.py +++ b/cargento/skills/cargento/tests/test_documentation.py @@ -1563,13 +1563,11 @@ def test_the_section_exists_under_a_heading_other_documents_can_anchor(self) -> "## Off-machine nudges (reaching the operator away from the desk)", self.SECURITY ) - def test_the_pathway_is_documented_as_unused_and_the_parser_agrees(self) -> None: - self.assertIn("No shipped feature posts to an endpoint the operator supplies", self.FLAT) - self.assertIn("That flag does not exist yet", self.SECTION) - # argparse prints usage to stderr before exiting, and that banner in a - # passing run reads like a failure to anyone watching the suite. - with contextlib.redirect_stderr(io.StringIO()), self.assertRaises(SystemExit): - cli.build_parser().parse_args(["--no-reach"]) + def test_the_pathway_is_documented_and_the_parser_agrees(self) -> None: + self.assertIn("H2 (DRC-4034) ships this capability", self.SECTION) + self.assertIn("The feature ships `--no-reach` with it", self.SECTION) + args = cli.build_parser().parse_args(["--no-reach"]) + self.assertTrue(args.no_reach) def test_the_payload_omits_the_count_no_threshold_backs(self) -> None: # The payload carries two counts and not DEC-4's three because no @@ -1823,10 +1821,11 @@ def test_the_two_unshipped_switches_are_the_only_no_flags_missing(self) -> None: "--no-events", "--no-irreversible", "--no-tripwires", + "--no-reach", }, shipped, ) - for documented in ("--no-handoff", "--no-reach"): + for documented in ("--no-handoff",): with self.subTest(flag=documented): self.assertNotIn(documented, shipped) self.assertIn(f"`{documented}`", self.FLAT) diff --git a/cargento/skills/cargento/tests/test_lifecycle.py b/cargento/skills/cargento/tests/test_lifecycle.py index 7f0404f3..48901971 100644 --- a/cargento/skills/cargento/tests/test_lifecycle.py +++ b/cargento/skills/cargento/tests/test_lifecycle.py @@ -1623,6 +1623,8 @@ def _args(self, **overrides: object) -> argparse.Namespace: "no_git": False, "no_focus": False, "no_history": False, + "no_reach": False, + "reach_url": None, "history_days": 14.0, "history_max_bytes": 1_048_576, } @@ -1721,6 +1723,22 @@ def test_no_ask_is_absent_when_not_requested(self) -> None: argv = lifecycle.spawn_argv(config, self._args(no_ask=False)) self.assertNotIn("--no-ask", argv) + def test_no_reach_is_forwarded(self) -> None: + config = cfg() + argv = lifecycle.spawn_argv(config, self._args(no_reach=True)) + self.assertIn("--no-reach", argv) + + def test_no_reach_is_absent_when_not_requested(self) -> None: + config = cfg() + argv = lifecycle.spawn_argv(config, self._args(no_reach=False)) + self.assertNotIn("--no-reach", argv) + + def test_reach_url_is_forwarded(self) -> None: + config = cfg() + argv = lifecycle.spawn_argv(config, self._args(reach_url="https://example.com/webhook")) + self.assertIn("--reach-url", argv) + self.assertIn("https://example.com/webhook", argv) + def test_daemon_is_never_forwarded(self) -> None: """Forwarding --daemon would respawn forever.""" config = cfg() diff --git a/cargento/skills/cargento/tests/test_reach.py b/cargento/skills/cargento/tests/test_reach.py new file mode 100644 index 00000000..550524fe --- /dev/null +++ b/cargento/skills/cargento/tests/test_reach.py @@ -0,0 +1,175 @@ +"""Unit tests for the off-machine reach nudge module (DRC-4034, DEC-4).""" + +from __future__ import annotations + +import http.server +import json +import os +import tempfile +import threading +import unittest +import urllib.error +import urllib.request +from pathlib import Path +from typing import Any +from unittest import mock + +from cargento_runtime import reach +from cargento_runtime.config import build_runtime_config +from cargento_runtime.state import build_runtime_state + + +def make_config(**overrides: Any) -> Any: + state_dir = overrides.pop("state_dir", None) or tempfile.mkdtemp() + environ = {"CARGENTO_HOME": state_dir} + return build_runtime_config( + environ=environ, + platform_name="darwin", + os_name="posix", + launcher_path=Path("/path/to/server.py"), + **overrides, + ) + + +class MockReachServer(http.server.HTTPServer): + def __init__(self) -> None: + super().__init__(("127.0.0.1", 0), _MockReachHandler) + self.received_requests: list[dict[str, Any]] = [] + + +class _MockReachHandler(http.server.BaseHTTPRequestHandler): + def log_message(self, format: str, *args: Any) -> None: # noqa: A002 + pass + + def do_POST(self) -> None: + length = int(self.headers.get("Content-Length", 0)) + body = self.rfile.read(length) + server: MockReachServer = self.server # type: ignore[assignment] + server.received_requests.append( + { + "path": self.path, + "headers": dict(self.headers), + "body": body, + "json": json.loads(body.decode("utf-8")), + } + ) + self.send_response(200) + self.end_headers() + self.wfile.write(b'{"ok": true}') + + +class ReachNudgeTest(unittest.TestCase): + def setUp(self) -> None: + self.server = MockReachServer() + self.server_thread = threading.Thread(target=self.server.serve_forever, daemon=True) + self.server_thread.start() + self.url = f"http://127.0.0.1:{self.server.server_address[1]}/webhook" + + def tearDown(self) -> None: + self.server.shutdown() + self.server.server_close() + + def test_payload_contains_strictly_counts_and_no_other_fields(self) -> None: + payload_bytes = reach.format_reach_payload(needs_input=3, finished_unread=2) + parsed = json.loads(payload_bytes.decode("utf-8")) + self.assertEqual(parsed, {"needs_input": 3, "finished_unread": 2}) + self.assertEqual(set(parsed.keys()), {"needs_input", "finished_unread"}) + + def test_opener_disables_proxies_and_refuses_redirects(self) -> None: + with mock.patch("urllib.request.build_opener") as mock_build: + reach.build_reach_opener() + mock_build.assert_called_once() + args = mock_build.call_args[0] + proxies = [h for h in args if isinstance(h, urllib.request.ProxyHandler)] + self.assertEqual(len(proxies), 1) + self.assertEqual(getattr(proxies[0], "proxies", None), {}) + self.assertIn(reach._NoRedirects, args) + no_redirect = reach._NoRedirects() + req = mock.MagicMock() + self.assertIsNone(no_redirect.redirect_request(req, None, 302, "Found", {}, "http://x")) + + def test_resolve_reach_url_from_config(self) -> None: + cfg = make_config(reach_url=self.url) + self.assertEqual(reach.resolve_reach_url(cfg), self.url) + + def test_resolve_reach_url_from_file(self) -> None: + with tempfile.TemporaryDirectory() as td: + cfg = make_config(state_dir=td) + url_file = os.path.join(td, "reach_url") + with open(url_file, "w", encoding="utf-8") as f: + f.write(f" {self.url}\n") + self.assertEqual(reach.resolve_reach_url(cfg), self.url) + + def test_no_reach_switch_disables_url_resolution(self) -> None: + cfg = make_config(reach_url=self.url, reach_enabled=False) + self.assertIsNone(reach.resolve_reach_url(cfg)) + + def test_nudge_delivered_when_active_session_enters_needs_input(self) -> None: + cfg = make_config(reach_url=self.url) + state = build_runtime_state(cfg, started=100.0) + sessions = [ + {"sid": "s1", "harness": "claude", "active": True, "state": "needs_input"}, + {"sid": "s2", "harness": "claude", "active": True, "state": "working"}, + ] + sent = reach.maybe_reach_nudge(cfg, state, sessions, now=100.0) + self.assertTrue(sent) + self.assertEqual(len(self.server.received_requests), 1) + req = self.server.received_requests[0] + self.assertEqual(req["json"], {"needs_input": 1, "finished_unread": 0}) + self.assertEqual(req["headers"]["Content-Type"], "application/json") + + def test_nudge_throttled_by_cooldown_interval(self) -> None: + cfg = make_config(reach_url=self.url, reach_cooldown_sec=60.0) + state = build_runtime_state(cfg, started=100.0) + sessions1 = [ + {"sid": "s1", "harness": "claude", "active": True, "state": "needs_input"}, + ] + self.assertTrue(reach.maybe_reach_nudge(cfg, state, sessions1, now=100.0)) + self.assertEqual(len(self.server.received_requests), 1) + + # Same counts within cooldown -> no post + self.assertFalse(reach.maybe_reach_nudge(cfg, state, sessions1, now=110.0)) + self.assertEqual(len(self.server.received_requests), 1) + + # Different counts within cooldown -> throttled, no post yet + sessions2 = [ + {"sid": "s1", "harness": "claude", "active": True, "state": "needs_input"}, + {"sid": "s2", "harness": "codex", "active": True, "state": "needs_input"}, + ] + self.assertFalse(reach.maybe_reach_nudge(cfg, state, sessions2, now=120.0)) + self.assertEqual(len(self.server.received_requests), 1) + + # Cooldown elapsed -> post sent with updated counts + self.assertTrue(reach.maybe_reach_nudge(cfg, state, sessions2, now=165.0)) + self.assertEqual(len(self.server.received_requests), 2) + self.assertEqual( + self.server.received_requests[1]["json"], {"needs_input": 2, "finished_unread": 0} + ) + + def test_no_reach_flag_never_attempts_post(self) -> None: + cfg = make_config(reach_url=self.url, reach_enabled=False) + state = build_runtime_state(cfg, started=100.0) + sessions = [ + {"sid": "s1", "harness": "claude", "active": True, "state": "needs_input"}, + ] + sent = reach.maybe_reach_nudge(cfg, state, sessions, now=100.0) + self.assertFalse(sent) + self.assertEqual(len(self.server.received_requests), 0) + + def test_diagnostics_sink_never_leaks_credential_url(self) -> None: + diagnostics: list[str] = [] + cfg = make_config(reach_url="http://127.0.0.1:1/invalid_secret_token_12345") + state = build_runtime_state(cfg, started=100.0) + sessions = [ + {"sid": "s1", "harness": "claude", "active": True, "state": "needs_input"}, + ] + sent = reach.maybe_reach_nudge( + cfg, + state, + sessions, + now=100.0, + diagnostic_sink=diagnostics.append, + ) + self.assertFalse(sent) + combined_logs = " ".join(diagnostics) + self.assertNotIn("invalid_secret_token_12345", combined_logs) diff --git a/docs/design-runtime-architecture.md b/docs/design-runtime-architecture.md index b2c3b99b..573e4ccc 100644 --- a/docs/design-runtime-architecture.md +++ b/docs/design-runtime-architecture.md @@ -80,6 +80,7 @@ Everything else lives in one file per responsibility: | `departures.py` | Every unasked check the lane ran, and the baseline each read against: the annotation revision, the moment the reading ran, and the producer's own account of what it read, all recorded at check time because none of them is recoverable by the time the reader comes back. Every check and not only the ones that found something, because a check that found nothing still spent a subprocess and a store of findings alone bounds nothing. A leaf shaped on `deliveries.py`, with the write lock and the per-thread temp name from its first commit rather than after the same measurement. It also owns the four sentences that say why there is nothing to show, which is what stops a spent cap reading like a board with nothing to raise, and a session nobody checked reading like one that was, and it chooses between them: the ladder lives here rather than in `unasked.py` because the review surface and the Intent log both need it and neither may reach the lane. Being a leaf is what makes one input a parameter: whether anything is typed against the session now cannot be asked of the annotation store from here, so both callers pass it, and neither may default it. A check outlives the words it read, and a check with nothing behind it earned the not-checked sentence rather than the nothing-found one. And it derives what a later check found about a raised constraint, from later checks in this same store rather than from a second model call. That is not only the cheap route, it is the only one available: a reader-requested reading carries its evidence cutoff as the producer's own sentence and not as a number, so nothing on it can be compared with a departure's cutoff. The store's own `cutoff` cannot be compared either (both write sites record the wall clock at the check, so it bounds the evidence from above and is the same number as `at`), and the sentences therefore say that a later check RAN and never that it read newer evidence. Unknown is the default and arrives as a named reason, and no sentence has a causal clause, because a return to the goal is not evidence the raise worked. Clearing an annotation reaches this store too: `withdraw` blanks the quotations from the rows those words were read against and keeps each row for the spend it made, which is what makes `SECURITY.md`'s claim about a clear true of a route that serves both stores. | | `ends.py` | The session ends this board observed, kept so a restart does not forget them: one record per session id with the harness's own stamp, written through by the coordinator on the `session_ended` that sets its in-memory mark and removed on the same three edges that lift it, and read back by `aggregate` onto a row whose coordinator never saw the end. A leaf shaped on `departures.py`, lock and per-thread temp name included, over `config`, `records` and `io` alone: the writer and the reader are joined only by the `OverlaySource` protocol, and a store either depended on would give that pair a real edge. Absence means not observed and never "did not end", exactly as it does in the coordinator's memory; the coordinator is the only writer, so `--no-events` leaves the file unread and unwritten and `--forget` deletes it. Its bounds are a security contract, not a preference: see "Session ends" in [SECURITY.md](../SECURITY.md). | | `unasked.py` | The lane [DEC-18](design-reading-a-session.md#dec-18-an-unasked-reading-is-permitted-and-gated-on-delivery-first) permits: check an annotated session against what the reader typed without being asked, and raise a departure. An orchestrator rather than a leaf, and its own module rather than an `Application` method, because it reaches `reading` for the producer, `project_context` for the evidence, `annotations` for the baseline, `departures` for the record and `notifications` for the raise, and `aggregate` imports neither of those last two halves. Nothing runs on the collection thread: `consider` does dictionary work and hands the reading to a worker, as `quota.request_fetch` hands a network read to one. Attached only when `record_history` is true and the switch is on, which is how `--diagnose` runs a collection without starting a subprocess. | +| `reach.py` | Off-machine nudge delivery: webhook resolution from command line, environment, or store file; strictly bounded payload formatting with count scalars only; lock-protected cooldown throttling; and redirect-free POST dispatch. Used by `aggregate.py` during periodic collection when configured. | | `notifications.py` | Hook state, popup policy for both lanes (needs-input and ask), the native notifier, hook payload handling. The notifier reports one of `deliveries`' five outcomes rather than `None`, and each lane spends its cooldown floor before the call and refunds it when the outcome says there was no lane to attempt. The hook lane's repeat suppressor is not refunded: it exists so one standing gate is not counted twice, and that reason does not change with the outcome. | | `collectors/*.py` | One harness each: a discovery predicate and a collector. Two of them, Cursor and Antigravity, reach a value through a bounded read inside a stored blob rather than off a column, and both bound the read in SQLite (`substr`) so the whole blob is never materialized. Each also owns the boundary *below* the per-harness one `aggregate.py` holds: a raise inside one row's build costs that row, and a raise outside one costs that store, so neither reaches the boundary that badges a harness. [design-unread-sources.md](design-unread-sources.md) U-5 owns the scope rule and the argument for it. | | `aggregate.py` | `HarnessSpec`, the registry and its label lookup, the per-harness failure boundary, and `Application`, including the one place a needs-input popup is decided: after the overlays have been reduced onto a row and before a dismissed row is subtracted (R-5). It is also where what the reader typed is attached to every row, annotated or not, and where a stored `final` reading is retracted to `withdrawn` once the row it was taken on no longer reports an end. Both are derived at publish time rather than stored, so both correct themselves in either direction. | diff --git a/scripts/validate_plugins.py b/scripts/validate_plugins.py index f7103e2f..873bd1e6 100755 --- a/scripts/validate_plugins.py +++ b/scripts/validate_plugins.py @@ -298,6 +298,7 @@ def construct_unique_mapping( "skills/cargento/cargento_runtime/http_api.py", "skills/cargento/cargento_runtime/observer.py", "skills/cargento/cargento_runtime/reading.py", + "skills/cargento/cargento_runtime/reach.py", "skills/cargento/cargento_runtime/semantic_history.py", "skills/cargento/cargento_runtime/interaction_prototype.py", "skills/cargento/cargento_runtime/project_context.py",