Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@ dependencies = [
"opentelemetry-exporter-otlp>=1.27",
"ddgs>=9,<10",
"networkx>=3.6,<4",
# Windows ships no system tz database, so zoneinfo falls back to this
# package. Without it current_datetime() raises for every IANA zone except
# UTC, and the agent cannot resolve "today" in the user's own timezone.
"tzdata; platform_system == 'Windows'",
]

[dependency-groups]
Expand Down
48 changes: 43 additions & 5 deletions s16code/capabilities.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,29 @@
"""
from __future__ import annotations

import re
from dataclasses import dataclass, field
from datetime import date as _date
from typing import Any
from urllib.parse import urlparse

from .core.redaction import scrub_host_paths

# A source that describes the machine rather than the material. Every producer
# is supposed to emit a relative name or an opaque scheme, but this is the one
# funnel every capability result passes through, so it is the cheapest place to
# stop a future capability reintroducing the same disclosure.
#
# A denial rather than an allowlist of schemes: an allowlist silently discards
# legitimate provenance the day someone adds a new scheme, and losing evidence
# is a worse failure than keeping an odd-looking source.
_HOST_PATH = re.compile(r"(?i)(^file:|\b[A-Za-z]:[\\/]|^\\\\|/home/|/Users/)")


def names_a_host_path(value: str) -> bool:
"""True when this string discloses the filesystem the agent runs on."""
return bool(_HOST_PATH.search(value))

# A capability worker may report that it ran but could not obtain what was
# asked of it. These two optional result keys are the declared, registry-owned
# way to say so; the planner's evidence review reads them generically.
Expand Down Expand Up @@ -260,6 +278,11 @@ def resolve_sources() -> list[str]:
found.append(projection.source_template.format(**result))
except (KeyError, IndexError):
pass
# A declared projection is no more trusted than a generic one: the leak
# this guards against came from a source_template interpolating a
# resolved path. Dropping the offending source falls back to graph://
# provenance, which is what the answer worker already treats as internal.
found = [source for source in found if not names_a_host_path(source)]
if not found and projection.external_sources:
found = list(external)
return found or [fallback_source]
Expand Down Expand Up @@ -322,9 +345,14 @@ def generic_evidence(result: dict[str, Any], *, skill: str, fallback_source: str
"""The projection for a capability that declared none: never silently lost."""
public = {key: value for key, value in result.items() if key not in drop}
sources = [value for key, value in public.items()
if key in {"uri", "url", "endpoint", "source_uri", "path"} and isinstance(value, str)]
if key in {"uri", "url", "endpoint", "source_uri", "path"}
and isinstance(value, str) and not names_a_host_path(value)]
import json as _json
return {"text": _json.dumps(public, ensure_ascii=False, default=str)[:12_000],
# The whole result is serialised into the evidence body, so scrubbing only
# the promoted sources would leave a host path sitting in the text under any
# key this function has never heard of.
body = scrub_host_paths(_json.dumps(public, ensure_ascii=False, default=str)[:12_000])
return {"text": body,
"sources": sources or [fallback_source], "kind": f"capability:{skill}"}


Expand Down Expand Up @@ -372,8 +400,14 @@ def boolean(description: str, **kwargs: Any) -> Argument:
Capability("read_file", "Read a UTF-8 file inside the configured sandbox. Paths are relative to the sandbox.",
{"path": string("Sandbox-relative file path.", maximum=2_000)},
families=("evidence",),
# `sandbox://`, not `file://`: this string is rendered into the
# model's context as the source of the evidence, and a `file://`
# URI built from a resolved path carries the host's directory
# layout and username into anything the model then says.
# `sandbox://reminders.txt` names the same file without
# describing the machine it lives on.
evidence=EvidenceProjection(kind="local_file", text="text",
source_template="file://{path}")),
source_template="sandbox://{path}")),
Capability("write_file", "Write a UTF-8 text artifact inside the configured sandbox, then return its path and SHA-256 digest. Use only when the user explicitly requests a file or a repair.",
{"path": string("Sandbox-relative destination path.", maximum=2_000),
"content": string("Complete text to write.", maximum=60_000),
Expand Down Expand Up @@ -412,8 +446,12 @@ def boolean(description: str, **kwargs: Any) -> Argument:
{"title": string("Human-readable event title.", maximum=300),
"dates": Argument("array", "ISO-8601 dates (YYYY-MM-DD).", minimum=1, maximum=20, item_kind="string")},
side_effect=True,
evidence=EvidenceProjection(kind="calendar_artifact", join="artifacts",
prefix="Calendar events created: ", sources=("artifacts",))),
# Names, not URIs: this projection's `join` puts the values
# into the evidence *body*, which the model reads as prose to
# summarise, so a file:// URI here is quoted back verbatim.
evidence=EvidenceProjection(kind="calendar_artifact", join="artifact_names",
prefix="Calendar events created: ",
sources=("artifact_names",))),
Capability("list_channels", "Ask GLC for the currently installed channel adapters and their connection state. The list is discovered at runtime, never encoded in the planner.",
{}),
Capability("send_channel_message", "Send one text message through an installed GLC channel. Use only when the request or an authorised subscription explicitly identifies the channel and recipient.",
Expand Down
9 changes: 8 additions & 1 deletion s16code/core/live_graph/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
from enum import StrEnum
from typing import TYPE_CHECKING, Any, Awaitable, Callable, Protocol

from ..redaction import scrub_host_paths

if TYPE_CHECKING:
from .store import GraphStore

Expand Down Expand Up @@ -192,7 +194,12 @@ async def _execute(self, task: TaskSpec) -> tuple[TaskSpec, bool, dict[str, Any]
try:
return task, True, await worker(task)
except Exception as exc: # worker failures become planner-visible events
return task, False, {"error": f"{type(exc).__name__}: {exc}"}
# Scrubbed here rather than where the error is displayed, because
# this string has three destinations: the durable graph journal, the
# planner's next prompt, and the answer worker's evidence. An OSError
# carries the absolute filename it failed on, so an unscrubbed error
# writes the operator's home directory into all three at once.
return task, False, {"error": scrub_host_paths(f"{type(exc).__name__}: {exc}")}

def _cancel_graph_cancelled_tasks(
self,
Expand Down
24 changes: 23 additions & 1 deletion s16code/core/memory/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,8 @@ def _scope_where(scope: MemoryScope) -> tuple[str, list[str | None]]:
return " AND ".join(clauses), values

def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind] | None = None,
limit: int = 8, include_history: bool = False, expand_neighbors: int = 0) -> list[MemoryRecord]:
limit: int = 8, include_history: bool = False, expand_neighbors: int = 0,
include_own_answers: bool = False) -> list[MemoryRecord]:
where, values = self._scope_where(scope)
if not include_history:
where += " AND status='current' AND (valid_to IS NULL OR valid_to > ?)"
Expand All @@ -225,6 +226,27 @@ def recall(self, query: str, scope: MemoryScope, *, kinds: Iterable[MemoryKind]
# answer evidence. Callers such as `audit_events` request them
# explicitly.
where += " AND kind NOT IN ('audit', 'policy')"
if not include_own_answers:
# An answer the agent produced is output, not evidence. It is stored
# as an episode so a conversation has continuity, but returning it
# here as ambient evidence closes a loop with no damping: the answer
# is a near-perfect lexical match for the question that produced it,
# so it becomes the top hit next time that question is asked, and
# the agent ends up citing itself.
#
# Observed: asked three times about a calendar invitation held in a
# sandbox file, recall returned the agent's own three previous
# "there is no information about it" replies as the top three hits.
# Each attempt made it more certain of something it had never
# checked, and rephrasing could not recover, because every retry
# added another denial.
#
# Only self-authored answers are excluded. What the user said is
# still evidence, and so are facts, documents and everything else;
# the discriminator is the run://<id>/answer source this store
# writes on the way out. Callers that genuinely want "what did you
# tell me before" pass include_own_answers=True.
where += " AND NOT (kind='episode' AND sources_json LIKE '%run://%/answer\"%')"
rows = self.db.execute(f"SELECT * FROM records WHERE {where}", values).fetchall()
if not rows:
return []
Expand Down
62 changes: 62 additions & 0 deletions s16code/core/redaction.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
"""Keep host filesystem paths out of text the agent stores or shows a model.

Capability results are the obvious way a path reaches the model, and those are
fixed at the producer. Exceptions are the way it happens by accident: an
`OSError` carries the resolved filename it failed on, and a failed node's error
string becomes planner-visible, is written to the durable graph journal, and is
handed to the answer worker as evidence. There is a test asserting exactly that
a failed file read reaches the final answer, so this is a live route rather than
a theoretical one.

Deliberately duplicated rather than shared with the gateway's
`glc.security.redaction`. The two live in separate processes with separate
dependency trees, and each defends its own output: the gateway cannot scrub
what the agent writes into its own journal, and the agent cannot scrub an
approval message the gateway composes. A shared library would couple two
services to protect against a class of bug that is cheapest to stop locally.

Stdlib only, so it stays importable from anywhere inside `s16code.core`.
"""

from __future__ import annotations

import re

_HTTP = re.compile(r"https?://\S+")

_TERMINATOR = r"[^\s`'\")\]}>]"

_PATTERNS = (
re.compile(rf"(?i)file:/{{0,3}}[A-Za-z]:[\\/]{_TERMINATOR}*"),
re.compile(rf"(?i)file://{_TERMINATOR}+"),
re.compile(rf"\\\\[A-Za-z0-9._-]+\\{_TERMINATOR}*"),
# One letter, then a colon, then a separator. The word boundary keeps this
# off "sandbox://" and "mxc://", whose scheme also ends in letter-colon.
re.compile(rf"(?i)\b[A-Za-z]:[\\/]{_TERMINATOR}*"),
re.compile(rf"(?<![\w.])/(?:home|Users)/{_TERMINATOR}+"),
)

_SEPARATORS = re.compile(r"[\\/]")


def _label(match: re.Match[str]) -> str:
name = _SEPARATORS.split(match.group(0))[-1].strip()
return f"[local file: {name}]" if name else "[local file]"


def scrub_host_paths(text: str) -> str:
"""Reduce any host path in `text` to its basename, leaving prose alone."""
if not text:
return text
held: list[str] = []

def stash(match: re.Match[str]) -> str:
held.append(match.group(0))
return f"\x00{len(held) - 1}\x00"

out = _HTTP.sub(stash, text)
for pattern in _PATTERNS:
out = pattern.sub(_label, out)
for index, original in enumerate(held):
out = out.replace(f"\x00{index}\x00", original)
return out
29 changes: 23 additions & 6 deletions s16code/events/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,16 +26,33 @@ def _json_object(text: str) -> dict[str, Any]:


def _spend_of(reply: dict[str, Any]) -> float:
"""What the gate itself cost. Deciding not to act is not free."""
"""What the gate itself cost. Deciding not to act is not free.

A metered call record is produced by the economics controller and its cost
field is named ``cost`` (see economics.budget.Charge). Reading only
``cost_usd``/``usd`` silently returned zero for every real call, which made
``daily_triage_budget`` unenforceable. All three spellings are accepted so
the helper works against the controller's records and the gateway client's.
"""
total = 0.0
for call in reply.get("metered_calls", []) or []:
if not isinstance(call, dict):
continue
try:
total += float(call.get("cost_usd") or call.get("usd") or 0.0)
except (TypeError, ValueError):
continue
return total
for key in ("cost", "cost_usd", "usd"):
value = call.get(key)
if value is not None:
try:
total += float(value)
except (TypeError, ValueError):
pass
break
if total:
return total
# Fall back to a single-call reply that carries its price at the top level.
try:
return float(reply.get("cost_usd") or (reply.get("cost") or {}).get("total_usd") or 0.0)
except (TypeError, ValueError, AttributeError):
return 0.0


class AutonomousEventEngine:
Expand Down
65 changes: 63 additions & 2 deletions s16code/events/report.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,57 @@ def liveness_status(store: EventStore, *, now: datetime | None = None,
}


def _awaiting_a_human(store: EventStore, window_start: datetime,
graph: Any = None) -> list[dict[str, Any]]:
"""Runs parked on a question nobody has answered yet.

This is the section the module docstring promises and the only place a
parked approval surfaces on its own. A run started by an event has no
channel to reply on, so when it stops to ask something nothing is sent
anywhere: not to Telegram, not by email, not on completion. Without this
list the agent can ask a question that no surface ever shows.

``graph`` is optional so the report still works without a runtime. When it
is supplied the live node state decides, because a decision record keeps
the status it had when it was written and would otherwise keep reporting a
question that was answered hours ago.
"""
waiting: list[dict[str, Any]] = []
for record in store.events():
received = record.get("received_at")
if received:
try:
if datetime.fromisoformat(received) < window_start:
continue
except ValueError:
pass
event = record.get("event") or {}
for decision in record.get("decisions", []):
run_id = decision.get("run_id")
if not run_id or decision.get("run_status") != "waiting":
continue
question = None
if graph is not None:
try:
nodes = graph.snapshot(run_id).nodes
except (KeyError, AttributeError):
continue
parked = [node for node in nodes.values() if node.get("state") == "waiting"]
if not parked:
continue # answered since; not still awaiting anyone
question = (parked[0].get("result") or {}).get("question") \
or (parked[0].get("input") or {}).get("question")
waiting.append({
"event": f"{event.get('source')}/{event.get('id')}",
"subscription": decision.get("subscription_id"),
"run_id": run_id,
"question": str(question)[:2_000] if question else None,
})
return waiting


def morning_report(store: EventStore, *, since: datetime | None = None,
now: datetime | None = None) -> dict[str, Any]:
now: datetime | None = None, graph: Any = None) -> dict[str, Any]:
moment = now or datetime.now(UTC)
window_start = since or (moment - timedelta(hours=24))
governor = AutonomyGovernor(store)
Expand Down Expand Up @@ -108,7 +157,7 @@ def morning_report(store: EventStore, *, since: datetime | None = None,
"subscription": item.get("subscription_id")}
for item in refusals],
"budgets": budgets,
"awaiting_a_human": [],
"awaiting_a_human": _awaiting_a_human(store, window_start, graph),
}


Expand Down Expand Up @@ -136,4 +185,16 @@ def render_markdown(report: dict[str, Any]) -> str:
for item in report[key][:100]:
lines.append(f"- `{item.get('event')}` — {item.get('reason') or item.get('control')}")
lines.append("")

# The section an operator acts on. A parked run is the one outcome that
# needs a person, and it is announced nowhere else.
parked = report.get("awaiting_a_human") or []
lines.append(f"## Awaiting a human ({len(parked)})")
if not parked:
lines.append("_nothing_")
for item in parked[:100]:
lines.append(f"- `{item.get('event')}` run `{item.get('run_id')}`")
if item.get("question"):
lines.append(f" - {item['question']}")
lines.append("")
return "\n".join(lines)
5 changes: 4 additions & 1 deletion s16code/events/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,10 @@ async def report(request: Request, hours: int = Query(default=24, ge=1, le=720),
fmt: str = Query(default="json", pattern="^(json|markdown)$")):
"""The human-reviewable account of a period nobody watched."""
since = datetime.now(UTC) - timedelta(hours=hours)
document = morning_report(request.app.state.event_store, since=since)
# The graph decides whether a parked run is still parked. Without it the
# report would keep listing questions that were answered hours ago.
document = morning_report(request.app.state.event_store, since=since,
graph=getattr(request.app.state.runtime, "graph", None))
if fmt == "markdown":
return PlainTextResponse(render_markdown(document))
return document
Expand Down
Loading