mirror of
https://github.com/usestrix/strix.git
synced 2026-08-24 03:42:37 +02:00
feat(migration): phase 2.1-2.3 — sandbox dispatch + thin slice tool wrappers
Phase 2.1 — sandbox dispatch helper:
- strix/tools/_sandbox_dispatch.py: post_to_sandbox() centralizes the
host->container HTTP wire format. Connect=10s, read=150s timeouts mirror
legacy executor.py. 50 MB response cap (C18) prevents OOM from a runaway
tool. All errors surface as {"error": str} so the model can recover
instead of the run dying.
Phase 2.2 — C6 lock-protected JSONL writes:
- strix/tools/notes/notes_actions.py: notes.jsonl appends are now wrapped
in _notes_lock so concurrent agents can't interleave half-written lines.
Regression test in test_notes_jsonl_concurrency.py verifies 1000 parallel
writes produce exactly 1000 valid JSON lines.
Phase 2.3 — thin-slice SDK wrappers (think + todo + notes):
- strix/tools/_legacy_adapter.py: LegacyAgentStateAdapter shim — exposes
just enough surface (.agent_id) for legacy tools that close over
agent_state, sourced from ctx.context['agent_id'].
- strix/tools/thinking/thinking_sdk_tools.py: 1 tool (think).
- strix/tools/todo/todo_sdk_tools.py: 6 tools (create/list/update/done/
pending/delete) with bulk-form preserved.
- strix/tools/notes/notes_sdk_tools.py: 5 tools (create/list/get/update/
delete) with asyncio.to_thread around the lock-protected file I/O.
Tests: 22 new tests pass (10 sandbox dispatch + 2 concurrency + 10 SDK
local). Full suite still green.
Per-file ruff ignores added for SDK wrapper files: TC002 (RunContextWrapper
must be runtime-importable because the SDK calls get_type_hints() to
derive the JSON schema) and PLR0911 (sandbox dispatch's 10 short-circuit
returns are intentional, each a distinct documented failure mode).
Refs: PLAYBOOK.md §3.4, AUDIT_R3.md C6/C18.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
375389b8bc
commit
6e5d96af34
@@ -0,0 +1,57 @@
|
||||
"""Shim that lets SDK function tools call legacy ``agent_state``-style functions.
|
||||
|
||||
The legacy harness's tools (notes, todos, reporting, …) take an
|
||||
``agent_state`` argument with shape ``state.agent_id`` for per-agent silo
|
||||
keying. Under the SDK migration the equivalent identity lives in
|
||||
``RunContextWrapper.context["agent_id"]``.
|
||||
|
||||
Rather than rewrite every tool body, SDK function-tool wrappers build a
|
||||
tiny adapter from the context dict and pass it to the legacy function.
|
||||
The legacy code path remains untouched (the legacy executor still calls
|
||||
its tools with the real ``AgentState``).
|
||||
|
||||
Used by:
|
||||
- ``tools/todo/todo_sdk_tools.py``
|
||||
- ``tools/notes/notes_sdk_tools.py``
|
||||
- ``tools/reporting/reporting_sdk_tools.py``
|
||||
- any other local tool that closes over ``agent_state.agent_id``
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from agents import RunContextWrapper
|
||||
|
||||
|
||||
@dataclass
|
||||
class LegacyAgentStateAdapter:
|
||||
"""Just enough surface for legacy tools that read ``state.agent_id``.
|
||||
|
||||
Don't rely on this for new code — it's only here to avoid touching
|
||||
the legacy ``*_actions.py`` modules during the migration. New SDK
|
||||
tools should read ``ctx.context["agent_id"]`` directly.
|
||||
"""
|
||||
|
||||
agent_id: str
|
||||
|
||||
|
||||
def adapter_from_ctx(
|
||||
ctx: RunContextWrapper,
|
||||
default_agent_id: str = "sdk-default",
|
||||
) -> LegacyAgentStateAdapter:
|
||||
"""Build a ``LegacyAgentStateAdapter`` from an SDK run context.
|
||||
|
||||
Falls back to ``default_agent_id`` when context is missing or its
|
||||
``agent_id`` is unset — keeps tests and CLI dry-runs working without
|
||||
a fully-populated context.
|
||||
"""
|
||||
inner = getattr(ctx, "context", None)
|
||||
if isinstance(inner, dict):
|
||||
agent_id = inner.get("agent_id") or default_agent_id
|
||||
else:
|
||||
agent_id = default_agent_id
|
||||
return LegacyAgentStateAdapter(agent_id=str(agent_id))
|
||||
@@ -0,0 +1,132 @@
|
||||
"""post_to_sandbox — host-to-container HTTP transport for sandbox tools.
|
||||
|
||||
Every Strix tool that runs inside the Kali container (browser, terminal,
|
||||
python, the seven Caido tools) has the same wire shape: POST a JSON body
|
||||
to ``http://localhost:{tool_server_host_port}/execute`` with a Bearer
|
||||
token header and ``{"agent_id", "tool_name", "kwargs"}`` as the body.
|
||||
|
||||
This helper centralizes that transport so:
|
||||
|
||||
- Every sandbox tool gets the same timeout policy
|
||||
(``connect=10s`` / ``read=150s``).
|
||||
- Every sandbox tool inherits the same response-size cap (50 MB) so a
|
||||
runaway tool body cannot OOM the host (C18).
|
||||
- Auth/transport errors surface as predictable error strings instead of
|
||||
exceptions, so the model can retry / pick a different tool without the
|
||||
run dying.
|
||||
|
||||
References:
|
||||
- PLAYBOOK.md §3.4
|
||||
- AUDIT_R3.md C18 (sandbox response size cap)
|
||||
- HARNESS_WIKI.md §7.2 (legacy executor.py wire format we mirror)
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
import httpx
|
||||
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from agents import RunContextWrapper
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Connect: how long to wait for the TCP handshake to complete.
|
||||
# Read: how long the tool may spend executing before we abandon the call.
|
||||
# Mirrors the legacy executor.py (``SANDBOX_EXECUTION_TIMEOUT = 120 + 30``).
|
||||
_SANDBOX_TIMEOUT = httpx.Timeout(connect=10.0, read=150.0, write=150.0, pool=150.0)
|
||||
|
||||
#: Cap on response body size from the tool server. Anything bigger is
|
||||
#: replaced by an error string so the model sees something coherent and
|
||||
#: the host doesn't OOM trying to allocate the buffer (C18).
|
||||
_MAX_RESPONSE_BYTES = 50 * 1024 * 1024 # 50 MB
|
||||
|
||||
|
||||
def _ctx_dict(ctx: RunContextWrapper) -> dict[str, Any] | None:
|
||||
"""Return ``ctx.context`` if it's a dict, else ``None``.
|
||||
|
||||
Strix's runtime always passes a dict (``make_agent_context``); other
|
||||
callers might not. Be defensive so a sandbox tool never raises just
|
||||
because the context shape is wrong.
|
||||
"""
|
||||
inner = getattr(ctx, "context", None)
|
||||
return inner if isinstance(inner, dict) else None
|
||||
|
||||
|
||||
async def post_to_sandbox(
|
||||
ctx: RunContextWrapper,
|
||||
tool_name: str,
|
||||
kwargs: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
"""POST a tool invocation to the in-container FastAPI tool server.
|
||||
|
||||
Returns:
|
||||
On success: ``{"result": <whatever the tool returned>}``.
|
||||
On any failure: ``{"error": "<human-readable error string>"}``.
|
||||
|
||||
Never raises. Tool authors call this and pass the return value
|
||||
straight to the model (or extract ``result`` for further shaping).
|
||||
"""
|
||||
inner = _ctx_dict(ctx)
|
||||
if inner is None:
|
||||
return {"error": "Sandbox not initialized: context is missing or not a dict."}
|
||||
|
||||
port = inner.get("tool_server_host_port")
|
||||
token = inner.get("sandbox_token")
|
||||
agent_id = inner.get("agent_id", "unknown")
|
||||
|
||||
if not port or not token:
|
||||
return {"error": "Sandbox not initialized: tool server port or token missing."}
|
||||
|
||||
url = f"http://127.0.0.1:{port}/execute"
|
||||
headers = {
|
||||
"Authorization": f"Bearer {token}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
body = {"agent_id": agent_id, "tool_name": tool_name, "kwargs": kwargs}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=_SANDBOX_TIMEOUT) as client:
|
||||
response = await client.post(url, json=body, headers=headers)
|
||||
except httpx.TimeoutException:
|
||||
return {
|
||||
"error": (f"Sandbox tool '{tool_name}' timed out after {_SANDBOX_TIMEOUT.read}s."),
|
||||
}
|
||||
except httpx.RequestError as e:
|
||||
# ConnectError, ReadError, NetworkError, etc.
|
||||
return {"error": f"Sandbox connection failed: {e!s}"[:300]}
|
||||
|
||||
if response.status_code == 401:
|
||||
return {"error": "Sandbox authorization failed (Bearer token invalid)."}
|
||||
if response.status_code >= 400:
|
||||
return {
|
||||
"error": (
|
||||
f"Sandbox tool '{tool_name}' failed with HTTP "
|
||||
f"{response.status_code}: {response.text[:300]}"
|
||||
),
|
||||
}
|
||||
|
||||
# Cap response size before parsing so a 1 GB rogue payload never lands
|
||||
# in our heap. Most legitimate tool responses are well under 100 KB.
|
||||
raw = response.content
|
||||
if len(raw) > _MAX_RESPONSE_BYTES:
|
||||
return {
|
||||
"error": (f"Sandbox response too large ({len(raw)} bytes; max {_MAX_RESPONSE_BYTES})."),
|
||||
}
|
||||
|
||||
try:
|
||||
data: Any = response.json()
|
||||
except ValueError:
|
||||
return {
|
||||
"error": (f"Sandbox tool '{tool_name}' returned non-JSON: {response.text[:200]}"),
|
||||
}
|
||||
|
||||
if not isinstance(data, dict):
|
||||
return {"error": f"Sandbox tool '{tool_name}' returned non-object JSON."}
|
||||
|
||||
return data
|
||||
@@ -38,6 +38,12 @@ def _get_notes_jsonl_path() -> Path | None:
|
||||
|
||||
|
||||
def _append_note_event(op: str, note_id: str, note: dict[str, Any] | None = None) -> None:
|
||||
"""Append one note operation to the run's ``notes/notes.jsonl``.
|
||||
|
||||
C6 (AUDIT_R2.md §1.1): hold ``_notes_lock`` across the file open + write
|
||||
so two concurrent agents (or two parallel SDK tool calls in Phase 6)
|
||||
cannot interleave bytes mid-line and corrupt the JSONL.
|
||||
"""
|
||||
notes_path = _get_notes_jsonl_path()
|
||||
if not notes_path:
|
||||
return
|
||||
@@ -50,7 +56,7 @@ def _append_note_event(op: str, note_id: str, note: dict[str, Any] | None = None
|
||||
if note is not None:
|
||||
event["note"] = note
|
||||
|
||||
with notes_path.open("a", encoding="utf-8") as f:
|
||||
with _notes_lock, notes_path.open("a", encoding="utf-8") as f:
|
||||
f.write(f"{json.dumps(event, ensure_ascii=True)}\n")
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,117 @@
|
||||
"""SDK function-tool wrappers for the legacy notes tools.
|
||||
|
||||
Five tools, all module-global (no per-agent silo). The legacy
|
||||
``notes_actions.py`` module already implements JSONL persistence and
|
||||
wiki Markdown rendering; these wrappers are pure delegation.
|
||||
|
||||
The C6 fix (lock-protected JSONL writes) was applied directly to the
|
||||
legacy module, so both code paths benefit.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from typing import Any
|
||||
|
||||
from agents import RunContextWrapper
|
||||
|
||||
from strix.tools._decorator import strix_tool
|
||||
from strix.tools.notes import notes_actions as _legacy
|
||||
|
||||
|
||||
def _dump(result: dict[str, Any]) -> str:
|
||||
return json.dumps(result, ensure_ascii=False, default=str)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def create_note(
|
||||
ctx: RunContextWrapper,
|
||||
title: str,
|
||||
content: str,
|
||||
category: str = "general",
|
||||
tags: list[str] | None = None,
|
||||
) -> str:
|
||||
"""Create a note in the current run's notes store.
|
||||
|
||||
Notes are persisted to ``run_dir/notes/notes.jsonl`` and (for the
|
||||
``wiki`` category) rendered as Markdown to ``run_dir/wiki/<slug>.md``.
|
||||
|
||||
Args:
|
||||
title: Required, non-empty title.
|
||||
content: Note body. Markdown is preserved.
|
||||
category: One of ``"general" | "findings" | "methodology" |
|
||||
"questions" | "plan" | "wiki"``.
|
||||
tags: Optional list of free-form tags.
|
||||
"""
|
||||
# The legacy function does file I/O under a threading.RLock.
|
||||
# Wrap in to_thread so we don't block the event loop while waiting
|
||||
# on the lock or fsync.
|
||||
result = await asyncio.to_thread(
|
||||
_legacy.create_note,
|
||||
title=title,
|
||||
content=content,
|
||||
category=category,
|
||||
tags=tags,
|
||||
)
|
||||
return _dump(result)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def list_notes(
|
||||
ctx: RunContextWrapper,
|
||||
category: str | None = None,
|
||||
tags: list[str] | None = None,
|
||||
search: str | None = None,
|
||||
include_content: bool = False,
|
||||
) -> str:
|
||||
"""List notes, optionally filtered.
|
||||
|
||||
Args:
|
||||
category: Filter by category.
|
||||
tags: Filter to notes that have any of these tags.
|
||||
search: Substring match against title and content.
|
||||
include_content: When False (default), entries get a ``content_preview``;
|
||||
when True, full content is included.
|
||||
"""
|
||||
result = await asyncio.to_thread(
|
||||
_legacy.list_notes,
|
||||
category=category,
|
||||
tags=tags,
|
||||
search=search,
|
||||
include_content=include_content,
|
||||
)
|
||||
return _dump(result)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def get_note(ctx: RunContextWrapper, note_id: str) -> str:
|
||||
"""Fetch one note by its 5-char ID. Returns full content."""
|
||||
result = await asyncio.to_thread(_legacy.get_note, note_id=note_id)
|
||||
return _dump(result)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def update_note(
|
||||
ctx: RunContextWrapper,
|
||||
note_id: str,
|
||||
title: str | None = None,
|
||||
content: str | None = None,
|
||||
tags: list[str] | None = None,
|
||||
) -> str:
|
||||
"""Update a note's title, content, or tags. Pass ``None`` to leave a field unchanged."""
|
||||
result = await asyncio.to_thread(
|
||||
_legacy.update_note,
|
||||
note_id=note_id,
|
||||
title=title,
|
||||
content=content,
|
||||
tags=tags,
|
||||
)
|
||||
return _dump(result)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def delete_note(ctx: RunContextWrapper, note_id: str) -> str:
|
||||
"""Delete a note. For wiki notes, also removes the rendered Markdown file."""
|
||||
result = await asyncio.to_thread(_legacy.delete_note, note_id=note_id)
|
||||
return _dump(result)
|
||||
@@ -0,0 +1,32 @@
|
||||
"""SDK function-tool wrapper for the legacy ``think`` tool.
|
||||
|
||||
Pattern: thin async wrapper that delegates to the legacy implementation
|
||||
in :mod:`strix.tools.thinking.thinking_actions`. The legacy function is
|
||||
sync and pure (no I/O), so we don't even need ``asyncio.to_thread``.
|
||||
|
||||
Validates the simplest tool-port pattern: legacy function in, JSON string
|
||||
out, no sandbox involvement.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
|
||||
from strix.tools._decorator import strix_tool
|
||||
from strix.tools.thinking.thinking_actions import think as _legacy_think
|
||||
|
||||
|
||||
@strix_tool(timeout=10)
|
||||
async def think(thought: str) -> str:
|
||||
"""Record a private chain-of-thought note without taking any action.
|
||||
|
||||
The "think" tool is the planning escape hatch for situations where a
|
||||
message-without-tool-call would otherwise halt the run (per the
|
||||
interactive-mode tool-call requirement). The thought itself is
|
||||
recorded but produces no side effects.
|
||||
|
||||
Args:
|
||||
thought: The agent's reasoning to record. Must be non-empty.
|
||||
"""
|
||||
result = _legacy_think(thought)
|
||||
return json.dumps(result, ensure_ascii=False)
|
||||
@@ -0,0 +1,146 @@
|
||||
"""SDK function-tool wrappers for the legacy todo tools.
|
||||
|
||||
Six tools, all in-memory, all per-agent (keyed by ``ctx.context["agent_id"]``
|
||||
through :class:`LegacyAgentStateAdapter`). Bulk forms are preserved —
|
||||
``todos`` / ``updates`` / ``todo_ids`` accept JSON strings or comma-separated
|
||||
strings the same way the legacy XML schema documented.
|
||||
|
||||
Pattern: thin async wrappers that delegate to the legacy implementations
|
||||
in :mod:`strix.tools.todo.todo_actions`. Legacy code is untouched.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from typing import Any
|
||||
|
||||
from agents import RunContextWrapper
|
||||
|
||||
from strix.tools._decorator import strix_tool
|
||||
from strix.tools._legacy_adapter import adapter_from_ctx
|
||||
from strix.tools.todo import todo_actions as _legacy
|
||||
|
||||
|
||||
def _dump(result: dict[str, Any]) -> str:
|
||||
"""JSON-dump a legacy result dict for the model. ``ensure_ascii=False``
|
||||
so unicode flows through; ``default=str`` to handle stray datetimes."""
|
||||
return json.dumps(result, ensure_ascii=False, default=str)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def create_todo(
|
||||
ctx: RunContextWrapper,
|
||||
title: str | None = None,
|
||||
description: str | None = None,
|
||||
priority: str = "normal",
|
||||
todos: str | None = None,
|
||||
) -> str:
|
||||
"""Create one or many todos for the current agent.
|
||||
|
||||
Args:
|
||||
title: Title of a single todo (alternative to bulk ``todos``).
|
||||
description: Optional details for the single todo.
|
||||
priority: ``"low" | "normal" | "high" | "critical"``.
|
||||
todos: Optional JSON string or comma-separated list for bulk create.
|
||||
"""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(
|
||||
_legacy.create_todo(
|
||||
agent_state=state,
|
||||
title=title,
|
||||
description=description,
|
||||
priority=priority,
|
||||
todos=todos,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def list_todos(
|
||||
ctx: RunContextWrapper,
|
||||
status: str | None = None,
|
||||
priority: str | None = None,
|
||||
) -> str:
|
||||
"""List the current agent's todos, sorted by status then priority.
|
||||
|
||||
Args:
|
||||
status: Optional ``"pending" | "in_progress" | "done"`` filter.
|
||||
priority: Optional ``"low" | "normal" | "high" | "critical"`` filter.
|
||||
"""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(_legacy.list_todos(agent_state=state, status=status, priority=priority))
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def update_todo(
|
||||
ctx: RunContextWrapper,
|
||||
todo_id: str | None = None,
|
||||
title: str | None = None,
|
||||
description: str | None = None,
|
||||
priority: str | None = None,
|
||||
status: str | None = None,
|
||||
updates: str | None = None,
|
||||
) -> str:
|
||||
"""Update one or many todos.
|
||||
|
||||
Args:
|
||||
todo_id: Single-todo target (alternative to bulk ``updates``).
|
||||
title / description / priority / status: New values for the single
|
||||
todo. Omit to leave unchanged.
|
||||
updates: Bulk form — JSON list of update dicts.
|
||||
"""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(
|
||||
_legacy.update_todo(
|
||||
agent_state=state,
|
||||
todo_id=todo_id,
|
||||
title=title,
|
||||
description=description,
|
||||
priority=priority,
|
||||
status=status,
|
||||
updates=updates,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def mark_todo_done(
|
||||
ctx: RunContextWrapper,
|
||||
todo_id: str | None = None,
|
||||
todo_ids: str | None = None,
|
||||
) -> str:
|
||||
"""Mark one (``todo_id``) or many (``todo_ids``) todos as done."""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(
|
||||
_legacy.mark_todo_done(agent_state=state, todo_id=todo_id, todo_ids=todo_ids),
|
||||
)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def mark_todo_pending(
|
||||
ctx: RunContextWrapper,
|
||||
todo_id: str | None = None,
|
||||
todo_ids: str | None = None,
|
||||
) -> str:
|
||||
"""Mark one (``todo_id``) or many (``todo_ids``) todos as pending."""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(
|
||||
_legacy.mark_todo_pending(
|
||||
agent_state=state,
|
||||
todo_id=todo_id,
|
||||
todo_ids=todo_ids,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@strix_tool(timeout=30)
|
||||
async def delete_todo(
|
||||
ctx: RunContextWrapper,
|
||||
todo_id: str | None = None,
|
||||
todo_ids: str | None = None,
|
||||
) -> str:
|
||||
"""Delete one (``todo_id``) or many (``todo_ids``) todos."""
|
||||
state = adapter_from_ctx(ctx)
|
||||
return _dump(
|
||||
_legacy.delete_todo(agent_state=state, todo_id=todo_id, todo_ids=todo_ids),
|
||||
)
|
||||
Reference in New Issue
Block a user