reserve wall time for artifact completion

This commit is contained in:
pewdiepie-archdaemon
2026-09-19 12:45:09 +00:00
parent efe8dbeab3
commit d7cad0621f
3 changed files with 131 additions and 8 deletions
+1 -1
View File
@@ -1 +1 @@
0.20.18
0.20.19
+43 -2
View File
@@ -1436,6 +1436,29 @@ def native_execution_limits(max_rounds):
return round_limit, NATIVE_TOOL_CALL_LIMIT
def native_artifact_completion_timing(client_runtime_context):
"""Return the research deadline and completion reserve for trusted runtimes.
The caller owns the outer wall clock. Invalid or absent timing metadata
leaves the existing call-count boundary unchanged.
"""
context = (
client_runtime_context
if isinstance(client_runtime_context, dict) else {}
)
try:
wall_seconds = float(context.get('agent_wall_time_seconds'))
reserve_seconds = float(context.get('artifact_completion_reserve_seconds'))
except (TypeError, ValueError):
return None, None
if (
not 60 <= wall_seconds <= 86400
or not 30 <= reserve_seconds < wall_seconds
):
return None, None
return wall_seconds - reserve_seconds, reserve_seconds
def standalone_social_turn(text):
"""A complete social utterance cannot authorize a tool action."""
return bool(re.fullmatch(
@@ -4346,6 +4369,10 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
native_workspace_enabled = native_workspace_runtime(
client_runtime_context, workspace,
)
artifact_research_seconds, artifact_completion_reserve_seconds = (
native_artifact_completion_timing(client_runtime_context)
if native_workspace_enabled else (None, None)
)
executable_tools = EXPLICIT_EXECUTE_TOOLS | (
NATIVE_WORKSPACE_EXECUTE_TOOLS
if native_workspace_enabled else frozenset()
@@ -4648,11 +4675,19 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
'source evidence and prepare the briefing. Do not repeat the same query.'
)})
yield event({'type': 'completion_recovery', 'reason': 'research_before_synthesis'})
artifact_research_elapsed = time.monotonic() - started
artifact_time_reserved = bool(
artifact_research_seconds is not None
and artifact_research_elapsed >= artifact_research_seconds
)
if (
required_artifacts
and not successful_artifact_write
and not artifact_write_phase
and calls >= min(NATIVE_ARTIFACT_RESEARCH_LIMIT, tool_call_limit - 1)
and (
calls >= min(NATIVE_ARTIFACT_RESEARCH_LIMIT, tool_call_limit - 1)
or artifact_time_reserved
)
):
artifact_write_phase = True
directory_artifact_guidance = ''
@@ -4682,9 +4717,15 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
})
yield event({
'type': 'completion_recovery',
'reason': 'artifact_write_budget_reserved',
'reason': (
'artifact_write_time_reserved'
if artifact_time_reserved
else 'artifact_write_budget_reserved'
),
'required_artifacts': list(required_artifacts),
'calls_used': calls,
'elapsed_seconds': round(artifact_research_elapsed, 3),
'completion_reserve_seconds': artifact_completion_reserve_seconds,
})
request_messages = provider_request_messages(
prune_multimodal_images(history, max_images=3)
+87 -5
View File
@@ -5643,18 +5643,19 @@ async def test_native_stream_reserves_remaining_budget_for_required_artifact(mon
assert request_contract["tool_choice"] == {
"type": "function", "function": {"name": "python"},
}
reserved_names = [tool["function"]["name"] for tool in requests[12]["tools"]]
artifact_request = requests[module.NATIVE_ARTIFACT_RESEARCH_LIMIT]
reserved_names = [tool["function"]["name"] for tool in artifact_request["tools"]]
assert reserved_names == ["python"]
assert requests[12]["tool_choice"] == {
assert artifact_request["tool_choice"] == {
"type": "function", "function": {"name": "python"},
}
completion_code_schema = (
requests[12]["tools"][0]["function"]["parameters"]["properties"]["code"]
artifact_request["tools"][0]["function"]["parameters"]["properties"]["code"]
)
assert "pattern" not in completion_code_schema
assert "Complete executable Python" in completion_code_schema["description"]
assert "create one or more files inside" in requests[12]["messages"][-1]["content"]
assert "Do not pass the directory itself as a file path" in requests[12]["messages"][-1]["content"]
assert "create one or more files inside" in artifact_request["messages"][-1]["content"]
assert "Do not pass the directory itself as a file path" in artifact_request["messages"][-1]["content"]
assert executed[-1] == "python"
assert any(
event.get("type") == "tool_output" and event.get("tool") == "python"
@@ -5662,6 +5663,87 @@ async def test_native_stream_reserves_remaining_budget_for_required_artifact(mon
)
@pytest.mark.asyncio
async def test_native_stream_reserves_wall_time_for_required_artifact(monkeypatch):
"""A slow research tool must not consume the artifact completion window."""
import src.clean_agent_preview as module
responses = iter([
{"choices": [{"delta": {"tool_calls": [{
"index": 0, "id": "slow-search",
"function": {"name": "web_search", "arguments": json.dumps({"query": "topic"})},
}]}}]},
{"choices": [{"delta": {"tool_calls": [{
"index": 0, "id": "artifact-write",
"function": {"name": "python", "arguments": json.dumps({
"code": "open('/tmp_workspace/results/out.md', 'w').write('evidence')",
})},
}]}}]},
{"choices": [{"delta": {"content": "Saved and verified."}}]},
])
requests = []
now = [0.0]
class Response:
def __init__(self, payload): self.payload = payload
async def __aenter__(self): return self
async def __aexit__(self, *args): pass
def raise_for_status(self): pass
async def aiter_lines(self):
yield "data: " + json.dumps(self.payload)
yield "data: [DONE]"
class Client:
def __init__(self, **kwargs): pass
async def __aenter__(self): return self
async def __aexit__(self, *args): pass
def stream(self, *args, **kwargs):
requests.append(kwargs["json"])
return Response(next(responses))
executed = []
async def execute(block, **kwargs):
executed.append(block.tool_type)
if block.tool_type == "web_search":
now[0] = 450.0
return block.tool_type, {"output": "ok", "exit_code": 0}
monkeypatch.setattr(module.time, "monotonic", lambda: now[0])
monkeypatch.setattr(module.httpx, "AsyncClient", Client)
monkeypatch.setattr(module, "execute_tool_block", execute)
schemas = [
item for item in FUNCTION_TOOL_SCHEMAS
if item["function"]["name"] in {"python", "web_search"}
]
contract = resolve_full_inventory_contract(schemas=schemas, policy=ToolPolicy())
raw = [chunk async for chunk in stream_preview(
endpoint_url="http://test", model="test",
messages=[{"role": "user", "content": "Research and create the requested output."}],
headers={}, turn_contract=contract, session_id="test", owner="test",
disabled_tools=set(), tool_policy=ToolPolicy(), workspace="/tmp/workspace",
client_runtime_context={
"surface": "odysseus-native", "terminal_agent": True,
"unattended_mode": True,
"agent_wall_time_seconds": 600,
"artifact_completion_reserve_seconds": 180,
"completion_requirements": {"required_artifacts": ["/tmp_workspace/results"]},
}, max_rounds=32,
)]
events = [json.loads(chunk[6:]) for chunk in raw if "[DONE]" not in chunk]
recovery = next(
event for event in events
if event.get("type") == "completion_recovery"
and event.get("reason") == "artifact_write_time_reserved"
)
assert recovery["calls_used"] == 1
assert recovery["elapsed_seconds"] == 450.0
assert recovery["completion_reserve_seconds"] == 180.0
assert executed == ["web_search", "python"]
assert [tool["function"]["name"] for tool in requests[1]["tools"]] == ["python"]
@pytest.mark.asyncio
async def test_detailed_video_answer_requires_second_focused_inspection(monkeypatch):
import src.clean_agent_preview as module