164 lines
9.1 KiB
Python
164 lines
9.1 KiB
Python
|
|
"""Offline direct foreground routing, persistence and delegation; no paid models."""
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
from pathlib import Path
|
||
|
|
import socket
|
||
|
|
import subprocess
|
||
|
|
import tempfile
|
||
|
|
import threading
|
||
|
|
import time
|
||
|
|
from http.server import ThreadingHTTPServer
|
||
|
|
from cli_flow import BINARY, Provider, response, tool
|
||
|
|
|
||
|
|
|
||
|
|
class BrowserProvider(Provider):
|
||
|
|
def do_GET(self):
|
||
|
|
body = b'<title>Local foreground fixture</title><body>Browser continuity</body>'
|
||
|
|
self.send_response(200)
|
||
|
|
self.send_header('Content-Length', str(len(body)))
|
||
|
|
self.send_header('Content-Type', 'text/html')
|
||
|
|
self.end_headers()
|
||
|
|
self.wfile.write(body)
|
||
|
|
|
||
|
|
|
||
|
|
def wait(predicate, timeout=12):
|
||
|
|
deadline = time.monotonic() + timeout
|
||
|
|
while time.monotonic() < deadline:
|
||
|
|
value = predicate()
|
||
|
|
if value:
|
||
|
|
return value
|
||
|
|
time.sleep(0.03)
|
||
|
|
raise AssertionError('fixture timed out')
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
with tempfile.TemporaryDirectory(prefix='lazyboy-foreground-') as root, ThreadingHTTPServer(('127.0.0.1', 0), BrowserProvider) as provider:
|
||
|
|
root = Path(root)
|
||
|
|
workspace = root/'workspace'
|
||
|
|
workspace.mkdir()
|
||
|
|
(workspace/'source.txt').write_text('verified local evidence')
|
||
|
|
data = root/'data'
|
||
|
|
fake_bin = root/'bin'
|
||
|
|
fake_bin.mkdir()
|
||
|
|
(fake_bin/'docker').write_text('#!/bin/sh\nexit 1\n')
|
||
|
|
(fake_bin/'docker').chmod(0o755)
|
||
|
|
provider.requests, provider.errors = [], []
|
||
|
|
provider.fragment_size = 7
|
||
|
|
foreground = {}
|
||
|
|
background = {}
|
||
|
|
fixture_url = f'http://127.0.0.1:{provider.server_port}/fixture'
|
||
|
|
|
||
|
|
def callback(body):
|
||
|
|
messages = body['messages']
|
||
|
|
system = messages[0].get('content', '')
|
||
|
|
if system.startswith('Extract memory'):
|
||
|
|
return response({'role':'assistant', 'content':'{"expertise":"","memories":[]}'})
|
||
|
|
if 'persistent LazyBoy main agent' in system:
|
||
|
|
goal = next(message['content'] for message in reversed(messages) if message['role']=='user' and not message.get('content', '').startswith('<system_reminder>'))
|
||
|
|
if goal.startswith('Background task result'):
|
||
|
|
return response({'role':'assistant', 'content':'背景結果已收到。'})
|
||
|
|
foreground[goal] = foreground.get(goal, 0) + 1
|
||
|
|
assert body['tool_choice']=='auto'
|
||
|
|
names = {definition['function']['name'] for definition in body['tools']}
|
||
|
|
assert {'external_read_file', 'external_write_file', 'browser_navigate', 'call_mcp_tool'} <= names
|
||
|
|
assert 'browser_handoff' not in names and 'wait_task' not in names
|
||
|
|
assert 'Background worker tools (available through' not in system
|
||
|
|
assert not any(message.get('name', '').startswith('lazyboy_public:') for message in messages)
|
||
|
|
assert not any('opened this turn by calling tools' in message.get('content', '') for message in messages)
|
||
|
|
round_number = foreground[goal]
|
||
|
|
if goal=='DIRECT_CHAT':
|
||
|
|
return response({'role':'assistant', 'content':'你好,直接回答。'})
|
||
|
|
if goal=='DIRECT_READ':
|
||
|
|
if round_number==1:
|
||
|
|
return response(tool('external_read_file', {'path':'source.txt'}))
|
||
|
|
assert json.loads(messages[-1]['content'])['content']=='verified local evidence'
|
||
|
|
return response({'role':'assistant', 'content':'已讀取並確認內容。'})
|
||
|
|
if goal=='DIRECT_FINAL':
|
||
|
|
return response(tool('send_message', {'type':'text', 'content':'直接結束。', 'final':True}))
|
||
|
|
if goal=='DIRECT_LONG':
|
||
|
|
assert round_number==1, 'delegation must not need another foreground model request'
|
||
|
|
message = tool('spawn_agent', {'goal':'BACKGROUND_FIXTURE'})
|
||
|
|
message['content']='這份長任務已交給背景處理。'
|
||
|
|
return response(message)
|
||
|
|
if goal=='DIRECT_BROWSER':
|
||
|
|
if round_number==1:
|
||
|
|
return response(tool('browser_navigate', {'url':fixture_url}))
|
||
|
|
if round_number==2:
|
||
|
|
assert 'error' not in json.loads(messages[-1]['content']), messages[-1]
|
||
|
|
return response(tool('browser_eval', {'expression':"localStorage.setItem('continuity','same-owner'); document.cookie='continuity=same-owner; path=/'; 'saved'"}))
|
||
|
|
assert json.loads(messages[-1]['content'])['result']=='saved'
|
||
|
|
message = tool('spawn_agent', {'goal':'BACKGROUND_BROWSER'})
|
||
|
|
message['content']='瀏覽器工作已交接。'
|
||
|
|
return response(message)
|
||
|
|
raise AssertionError(goal)
|
||
|
|
assert 'LazyBoy background task worker' in system, system
|
||
|
|
goal = next(message['content'] for message in messages if message['role']=='user')
|
||
|
|
background[goal] = background.get(goal, 0) + 1
|
||
|
|
if goal=='BACKGROUND_BROWSER':
|
||
|
|
round_number = background[goal]
|
||
|
|
if round_number==1:
|
||
|
|
return response(tool('browser_navigate', {'url':fixture_url}))
|
||
|
|
if round_number==2:
|
||
|
|
return response(tool('browser_eval', {'expression':"({cookie:document.cookie, storage:localStorage.getItem('continuity')})"}))
|
||
|
|
result = json.loads(messages[-1]['content'])['result']
|
||
|
|
assert result=={'cookie':'continuity=same-owner', 'storage':'same-owner'}, result
|
||
|
|
return response(tool('send_message', {'type':'text', 'content':'背景工作已完成。', 'final':True}))
|
||
|
|
|
||
|
|
provider.callback = callback
|
||
|
|
threading.Thread(target=provider.serve_forever, daemon=True).start()
|
||
|
|
env = {**os.environ, 'PATH':str(fake_bin)+os.pathsep+os.environ['PATH'],
|
||
|
|
'LAZYBOY_TLS':'0', 'LAZYBOY_WEB_PORT':'0', 'LAZYBOY_DATA_DIR':str(data),
|
||
|
|
'LAZYBOY_API_KEY':'mock', 'LAZYBOY_MODEL':'mock',
|
||
|
|
'LAZYBOY_BROWSER_SURFACE':'local', 'LAZYBOY_BROWSER_HEADED':'0',
|
||
|
|
'LAZYBOY_MCP_CONFIG':str(root/'mcp.json'),
|
||
|
|
'LAZYBOY_BASE_URL':f'http://127.0.0.1:{provider.server_port}/v1'}
|
||
|
|
daemon = subprocess.Popen([str(BINARY), 'serve'], cwd=workspace, env=env,
|
||
|
|
stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, text=True)
|
||
|
|
logs = []
|
||
|
|
threading.Thread(target=lambda: [logs.append(line) for line in daemon.stderr], daemon=True).start()
|
||
|
|
|
||
|
|
def rpc(op, **extra):
|
||
|
|
with socket.socket(socket.AF_UNIX) as client:
|
||
|
|
client.settimeout(5)
|
||
|
|
client.connect(str(data/'service.sock'))
|
||
|
|
client.sendall((json.dumps({'op':op, **extra})+'\n').encode())
|
||
|
|
result = json.loads(client.makefile().readline())
|
||
|
|
assert 'error' not in result, result
|
||
|
|
return result
|
||
|
|
|
||
|
|
try:
|
||
|
|
wait(lambda: (data/'service.sock').exists() and daemon.poll() is None)
|
||
|
|
owner = rpc('create', name='direct', cwd=str(workspace))['id']
|
||
|
|
for goal, expected, calls in [('DIRECT_CHAT','你好,直接回答。',1),
|
||
|
|
('DIRECT_READ','已讀取並確認內容。',2),
|
||
|
|
('DIRECT_FINAL','直接結束。',1),
|
||
|
|
('DIRECT_LONG','這份長任務已交給背景處理。',1),
|
||
|
|
('DIRECT_BROWSER','瀏覽器工作已交接。',3)]:
|
||
|
|
cursor = [rpc('event_cursor', agent=owner)['id']]
|
||
|
|
queued = rpc('chat', agent=owner, message=goal)['queued']
|
||
|
|
def completed_reply():
|
||
|
|
batch = rpc('events', agent=owner, after=cursor[0])['events']
|
||
|
|
if batch:
|
||
|
|
cursor[0] = batch[-1]['id']
|
||
|
|
return [event for event in batch if event['kind']=='reply' and event['payload'].get('chat_id')==queued]
|
||
|
|
events = wait(completed_reply)
|
||
|
|
assert events[-1]['payload']['message']==expected, events
|
||
|
|
assert foreground[goal]==calls, foreground
|
||
|
|
snapshot = rpc('get', agent=owner)
|
||
|
|
assert any(line['content']==expected for line in snapshot['transcript']), snapshot
|
||
|
|
print(f'PASS {goal}: {calls} foreground model request(s), durable visible reply')
|
||
|
|
wait(lambda: all(task['state']=='terminal' for task in rpc('tasks', agent=owner)))
|
||
|
|
assert not provider.errors, provider.errors
|
||
|
|
assert background['BACKGROUND_BROWSER']==3, background
|
||
|
|
assert all(task['verdict']=='answer' for task in rpc('tasks', agent=owner))
|
||
|
|
print('PASS direct browser delegates with the same owner cookie/profile and releases its lease')
|
||
|
|
print('PASS background final message finishes; no classifier model; provider history stays paired')
|
||
|
|
finally:
|
||
|
|
daemon.terminate()
|
||
|
|
daemon.wait(timeout=8)
|
||
|
|
provider.shutdown()
|
||
|
|
|
||
|
|
|
||
|
|
if __name__=='__main__':
|
||
|
|
main()
|