"""Real daemon extension selection, context loading, and live reload. No live model."""
import contextlib
import http.server
import json
import os
from pathlib import Path
import subprocess
import tempfile
import threading
import time
import urllib.error
import urllib.request
ROOT = Path(__file__).resolve().parents[2]
class Provider(http.server.BaseHTTPRequestHandler):
requests = []
entered = threading.Event()
release = threading.Event()
def log_message(self, *_):
pass
def do_POST(self):
request = json.loads(self.rfile.read(int(self.headers["Content-Length"])))
self.requests.append(request)
messages = request["input"]
prompt = next(item.get("content", "") for item in reversed(messages) if item.get("role") == "user")
if prompt == "hold this turn":
self.entered.set()
assert self.release.wait(15), "test did not release held model request"
output = [{
"type": "message", "role": "assistant", "status": "completed",
"content": [{"type": "output_text", "text": "done", "annotations": []}],
}]
if prompt == "model probe via python" and messages[-1].get("type") != "function_call_output":
code = ("selection = await commands.model()\n"
"assert selection['model'] == 'fixture', selection\n"
"try:\n"
" await commands.model('unreachable-model')\n"
" assert False, 'switch must refuse from the model'\n"
"except CommandsError as error:\n"
" assert 'user action' in str(error), error\n"
"print('MODEL_COMMAND_OK')")
output = [{"type": "function_call", "id": "fc-model", "call_id": "call-model",
"name": "python", "arguments": json.dumps({"code": code, "timeout_ms": 10000}),
"status": "completed"}]
if prompt == "activate demo via python" and messages[-1].get("type") != "function_call_output":
code = "activation = await commands.demo('python argument')\nassert 'BODY_MUST_NOT_AUTOLOAD' in activation['instructions']\nassert activation['arguments'] == 'python argument'\nprint('SKILL_PYTHON_ACTIVATION_OK')"
output = [{"type": "function_call", "id": "fc-skills", "call_id": "call-skills",
"name": "python", "arguments": json.dumps({"code": code, "timeout_ms": 10000}),
"status": "completed"}]
if prompt == "hot probe before reload" and messages[-1].get("type") != "function_call_output":
code = ("names = [c['name'] for c in await commands.catalog()]\n"
"assert '/late' not in names, names\n"
"hot_marker = 'kernel-kept-running'\n"
"print('BEFORE_RELOAD_OK')")
output = [{"type": "function_call", "id": "fc-hot-before", "call_id": "call-hot-before",
"name": "python", "arguments": json.dumps({"code": code, "timeout_ms": 10000}),
"status": "completed"}]
if prompt == "hot probe after reload" and messages[-1].get("type") != "function_call_output":
code = ("names = [c['name'] for c in await commands.catalog()]\n"
"assert '/late' in names, names\n"
"assert hot_marker == 'kernel-kept-running', 'kernel was restarted'\n"
"print('AFTER_RELOAD_OK')")
output = [{"type": "function_call", "id": "fc-hot-after", "call_id": "call-hot-after",
"name": "python", "arguments": json.dumps({"code": code, "timeout_ms": 10000}),
"status": "completed"}]
response = {"type": "response.completed", "response": {
"id": "fixture", "status": "completed", "output": output,
"usage": {"input_tokens": 20, "output_tokens": 1},
}}
body = ("data: " + json.dumps(response) + "\n\n").encode()
self.send_response(200)
self.send_header("Content-Type", "text/event-stream")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def run(endpoint):
with tempfile.TemporaryDirectory(prefix="albedo-extensions-") as directory:
root = Path(directory)
home, workspace, user_home = root/"state", root/"workspace", root/"user"
for path in (home, workspace, user_home):
path.mkdir(mode=0o700)
skill = workspace/".albedo"/"skills"/"demo"/"SKILL.md"
skill.parent.mkdir(parents=True)
skill.write_text("---\nname: demo\ndescription: catalog-only fixture description\n---\nBODY_MUST_NOT_AUTOLOAD\n")
(home/"config.json").write_text(json.dumps({"active": "fixture", "providers": {"fixture": {
"baseUrl": endpoint, "apiKey": "fixture-key", "model": "fixture", "protocol": "responses",
}}}))
# Catalog refresh is disabled so the suite never reaches the network.
(home/"extensions.json").write_text(json.dumps({"models": {"refreshHours": 0}}))
env = dict(os.environ, HOME=str(user_home), ALBEDO_HOME=str(home), ALBEDO_PARENT_PID=str(os.getpid()))
connection = None
def cli(*args):
result = subprocess.run(["node", "cli/bin/albedo.mjs", *args], cwd=ROOT, env=env,
capture_output=True, text=True, timeout=45)
assert result.returncode == 0, result.stdout + result.stderr + (home/"daemon.log").read_text()
return result.stdout
def connect():
nonlocal connection
cli("sessions")
connection = json.loads((home/"daemon.json").read_text())
def api(path, data=None):
request = urllib.request.Request(f"http://127.0.0.1:{connection['port']}" + path,
headers={"Authorization": "Bearer " + connection["token"], "Content-Type": "application/json"},
data=None if data is None else json.dumps(data).encode())
with urllib.request.urlopen(request, timeout=25) as response:
return json.load(response)
def ready(session):
deadline = time.monotonic() + 20
while time.monotonic() < deadline:
if not api(f"/sessions/{session}/status")["running"]:
return
time.sleep(.025)
raise AssertionError("session did not settle")
def stop():
if not connection:
return
with contextlib.suppress(Exception):
api("/shutdown", {})
deadline = time.monotonic() + 15
while time.monotonic() < deadline:
try:
os.kill(connection["pid"], 0)
except ProcessLookupError:
return
time.sleep(.05)
raise AssertionError("daemon did not stop cleanly")
def rejected(path, data):
try:
api(path, data)
raise AssertionError("invalid extension change succeeded")
except urllib.error.HTTPError as error:
assert error.code == 409, error.read()
def snapshot(session):
# `api` would json.load an endless event stream; read one frame instead.
# after_seq=0 replays the live log: a fresh stream gets a transcript
# snapshot, which does not carry live-only events like `compacted`.
request = urllib.request.Request(f"http://127.0.0.1:{connection['port']}/sessions/{session}/stream?after_seq=0",
headers={"Authorization": "Bearer " + connection["token"]})
with urllib.request.urlopen(request, timeout=25) as response:
while True:
line = response.readline()
if line.startswith(b"data: "):
return json.loads(line[6:])["events"]
def catalog_request(session, prompt):
before = len(Provider.requests)
api(f"/sessions/{session}/events", {"content": prompt})
ready(session)
assert len(Provider.requests) == before + 1
return Provider.requests[-1]
try:
connect()
daemon_pid = connection["pid"]
assert "session_extensions" in api("/health")["capabilities"]
session = json.loads(cli("new", str(workspace)))["session"]
route = f"/sessions/{session}/extensions"
installed = api(route)
assert {"python", "bash", "work", "files", "skills"} <= {item["name"] for item in installed}
assert all(item["enabled"] for item in installed if item["name"] in {"python", "bash", "work", "files", "skills"})
request = catalog_request(session, "first turn")
text = json.dumps(request["input"])
assert "catalog-only fixture description" in text and str(skill.resolve()) in text, request
assert "BODY_MUST_NOT_AUTOLOAD" not in json.dumps(request), request
assert "catalog-only fixture description" not in request.get("instructions", ""), request
assert text.index("catalog-only fixture description") < text.index("first turn")
assert "returning {content, next_offset, size, truncated}" in text, request
tools = {tool["name"] for tool in request["tools"]}
installed = api(route) # Managed capabilities resolve when the lazy worker opens.
skills = next(item for item in installed if item["name"] == "skills")
assert skills["requires"] == ["python", "commands"]
assert "skills" in skills["python_modules"], skills
assert skills["tools"] == []
assert not {"skills_read", "skills_list"} & tools
catalog = api(f"/sessions/{session}/commands")
demo = [c for c in catalog if c["name"] == "/demo"]
assert demo == [{
"name": "/demo",
"description": "catalog-only fixture description",
"method": "demo",
"usage": "/demo [arguments]",
"arguments": [{"name": "arguments", "description": "arguments for the skill", "required": False}],
"modelCallable": True,
"userTurn": True,
"page": False,
}], demo
assert {c["name"] for c in catalog} >= {"/model", "/context", "/compact"}
before = len(Provider.requests)
outcome = api(f"/sessions/{session}/commands", {
"name": "/demo", "arguments": "one two", "clientId": "fixture-client",
})
assert outcome == {"submitted": True}, outcome
ready(session)
assert len(Provider.requests) == before + 1
activation = json.dumps(Provider.requests[-1]["input"])
assert "BODY_MUST_NOT_AUTOLOAD" in activation
assert "one two" in activation and str(skill.resolve()) in activation
before = len(Provider.requests)
original_tree = api(f"/sessions/{session}/tree?after=0&limit=100")["items"]
started = api(f"/sessions/{session}/commands", {"name": "/compact"})
assert started["result"]["strategy"] == "rolling" and started["result"]["started"] is True, started
ready(session)
assert len(Provider.requests) == before + 1, "only the active summarizer may call the provider"
assert api(f"/sessions/{session}/tree?after=0&limit=100")["items"] == original_tree, "manual compaction rewrote the transcript"
before = len(Provider.requests)
api(f"/sessions/{session}/events", {"content": "after manual compaction"})
ready(session)
assert len(Provider.requests) == before + 1
followup = json.dumps(Provider.requests[-1]["input"])
assert "older conversation summary" in followup, "compacted projection was not reused: LOG: " + (home/"daemon.log").read_text()[-3000:] + " CONTEXT: " + json.dumps(api(f"/sessions/{session}/context"))[:1200]
assert api(f"/sessions/{session}/context")["compaction"]["status"] == "compacted"
context = api(f"/sessions/{session}/context")
assert context["state"] == "ready" and context["compaction"]["status"] == "compacted", context
events = snapshot(session)
compacted = next(e for e in events if e.get("type") == "compacted")
assert compacted["evicted"] > 0 and "older conversation summary" in compacted["summary"], compacted
prepared = api(f"/sessions/{session}/context/history/0")["content"]
assert "older conversation summary" in prepared, prepared
listed = api(f"/sessions/{session}/commands")
assert listed == catalog, listed
original_skill = skill.read_text()
skill.write_text(original_skill.replace("catalog-only fixture description", "changed on disk"))
assert api(f"/sessions/{session}/commands") == catalog, "catalog changed without reload"
skill.write_text(original_skill)
before = len(Provider.requests)
api(f"/sessions/{session}/commands", {"name": "/demo", "arguments": "slash argument"})
ready(session)
assert len(Provider.requests) == before + 1, "slash activation must submit exactly one turn"
activated = next(item["content"] for item in reversed(Provider.requests[-1]["input"]) if item.get("role") == "user")
activation = json.loads(activated.split("\n", 1)[1])
assert activation["name"] == "demo" and activation["arguments"] == "slash argument"
assert activation["source"] == str(skill.resolve())
assert "BODY_MUST_NOT_AUTOLOAD" in activation["instructions"]
# A skill written after the session opened stays invisible until a
# session reload — then reaches the menu, the prompt context, and the
# live kernel's routes in one swap, with the kernel still running.
late = workspace/".agents"/"skills"/"late"/"SKILL.md"
late.parent.mkdir(parents=True)
late.write_text("---\nname: late\ndescription: added after the session opened\n---\nLATE_BODY\n")
assert "/late" not in {c["name"] for c in api(f"/sessions/{session}/commands")}
before = len(Provider.requests)
api(f"/sessions/{session}/events", {"content": "hot probe before reload"})
ready(session)
assert len(Provider.requests) == before + 2
tool_output = next(item["output"] for item in reversed(Provider.requests[-1]["input"]) if item.get("type") == "function_call_output")
assert "BEFORE_RELOAD_OK" in tool_output, tool_output
cached = Provider.requests[-1]
reloaded = api(f"/sessions/{session}/commands", {"name": "/reload", "args": {"target": "session"}})
assert reloaded["result"]["reloaded"] == "session", reloaded
assert "/late" in {c["name"] for c in api(f"/sessions/{session}/commands")}
request = catalog_request(session, "after session reload")
assert "added after the session opened" in json.dumps(request["input"]), request
# A live session keeps the prompt prefix the provider cached; the
# new skill arrives as a durable user-role context update instead.
assert request["instructions"] == cached["instructions"], "reload changed the cached system prompt"
assert request["input"][0] == cached["input"][0], "reload changed the cached leading context"
assert "added after the session opened" not in json.dumps(request["input"][0])
updates = [item for item in request["input"] if item.get("role") == "user"
and str(item.get("content", "")).startswith("\n[albedo] This session's extensions changed")]
assert len(updates) == 1 and "added after the session opened" in updates[0]["content"], request
assert updates[0]["content"].endswith(""), updates[0]
before = len(Provider.requests)
api(f"/sessions/{session}/events", {"content": "hot probe after reload"})
ready(session)
assert len(Provider.requests) == before + 2
tool_output = next(item["output"] for item in reversed(Provider.requests[-1]["input"]) if item.get("type") == "function_call_output")
assert "AFTER_RELOAD_OK" in tool_output, tool_output
assert Provider.requests[-1]["input"][0] == cached["input"][0], "pin must hold until compaction"
# Compaction rewrites history and loses the cache anyway, so the
# prompt prefix is rebuilt from the session's current capabilities.
api(f"/sessions/{session}/commands", {"name": "/compact"})
ready(session)
request = catalog_request(session, "after compaction following reload")
assert "added after the session opened" in json.dumps(request["input"][0]), request["input"][0]
request = catalog_request(session, "pin stays released")
assert "added after the session opened" in json.dumps(request["input"][0]), request["input"][0]
python_session = json.loads(cli("new", str(workspace)))["session"]
before = len(Provider.requests)
api(f"/sessions/{python_session}/events", {"content": "activate demo via python"})
ready(python_session)
assert len(Provider.requests) == before + 2, "python activation must not submit another user turn"
tool_output = next(item["output"] for item in reversed(Provider.requests[-1]["input"]) if item.get("type") == "function_call_output")
assert "SKILL_PYTHON_ACTIVATION_OK" in tool_output, tool_output
model_session = json.loads(cli("new", str(workspace)))["session"]
before = len(Provider.requests)
api(f"/sessions/{model_session}/events", {"content": "model probe via python"})
ready(model_session)
assert len(Provider.requests) == before + 2
tool_output = next(item["output"] for item in reversed(Provider.requests[-1]["input"]) if item.get("type") == "function_call_output")
assert "MODEL_COMMAND_OK" in tool_output, tool_output
switched = api(f"/sessions/{model_session}/commands", {"name": "/model", "args": {"model": "switched-model"}})
assert switched["result"]["model"] == "switched-model", switched
assert switched["result"]["provider"] == "fixture", switched
listed_sessions = api("/sessions")
assert next(s for s in listed_sessions if s["id"] == model_session)["model"] == "switched-model"
# /work is the work extension's own page. A user change is told to
# the agent, but never starts a turn: while idle it waits and rides
# ahead of the next message.
work_command = next(c for c in api(f"/sessions/{session}/commands") if c["name"] == "/work")
assert work_command["page"] is True and not work_command["modelCallable"], work_command
page = api(f"/sessions/{session}/commands", {"name": "/work", "args": {}})["result"]["page"]
assert page["title"] == "work" and {action["key"] for action in page["actions"]} >= {"a", "d", "x"}, page
added = api(f"/sessions/{session}/commands", {"name": "/work", "args": {"action": "add", "details": "write the release notes"}})
assert "the agent will be told" in added["result"]["message"], added
before = len(Provider.requests)
time.sleep(0.3)
assert len(Provider.requests) == before, "a ledger note must not start a turn"
request = catalog_request(session, "anything new on the ledger?")
text = json.dumps(request["input"])
assert "The user added work item" in text, text
work_note = next(item["content"] for item in request["input"] if item.get("role") == "user"
and "The user added work item" in item.get("content", ""))
assert work_note.count("") == 1 and work_note.count("") == 1, work_note
assert text.index("The user added work item") < text.index("anything new on the ledger?"), text
page = api(f"/sessions/{session}/commands", {"name": "/work", "args": {}})["result"]["page"]
assert any(row["text"] == "write the release notes" for row in page["glance"]["rows"]), page
rejected(route, {"name": "not-installed", "enabled": False})
rejected(route, {"name": "python", "enabled": False})
assert api(route) == installed
disabled = api(route, {"name": "skills", "enabled": False})
assert not next(item["enabled"] for item in disabled if item["name"] == "skills")
rejected(f"/sessions/{session}/commands", {"name": "/demo", "arguments": "must not run"})
assert json.loads((home/"daemon.json").read_text())["pid"] == daemon_pid
request = catalog_request(session, "extension disabled")
# Disabling skills keeps the cached prefix (the python tool is
# unchanged); the model is told the skills context was removed.
update = next(item["content"] for item in reversed(request["input"]) if item.get("role") == "user"
and item.get("content", "").startswith("\n[albedo] This session's extensions changed"))
assert "Removed context:\n- skills" in update and "" not in update, update
assert "Current extension instructions" not in update and len(update) < 4000, update
assert "skills" not in {module for item in disabled if item["enabled"] for module in item["python_modules"]}
assert not {"skills_read", "skills_list"} & {tool["name"] for tool in request.get("tools", [])}
api(f"/sessions/{session}/events", {"content": "hold this turn"})
assert Provider.entered.wait(10)
rejected(route, {"name": "skills", "enabled": True})
rejected(f"/sessions/{session}/commands", {"name": "/model", "args": {"model": "mid-run"}})
rejected(f"/sessions/{session}/commands", {"name": "/compact"})
assert api(route) == disabled
Provider.release.set()
ready(session)
stop()
connect()
assert not next(item["enabled"] for item in api(route) if item["name"] == "skills")
skill.write_text("---\nname: demo\ndescription: refreshed catalog description\n---\nBODY_MUST_NOT_AUTOLOAD\n")
api(route, {"name": "skills", "enabled": True})
request = catalog_request(session, "extension reloaded")
assert "refreshed catalog description" in json.dumps(request)
context = request["input"][0]["content"]
assert "catalog-only fixture description" not in context
assert "BODY_MUST_NOT_AUTOLOAD" not in context
print("extension context, live toggles, busy/dependency guards, and persistence passed")
finally:
Provider.release.set()
stop()
if __name__ == "__main__":
server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), Provider)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
try:
run(f"http://127.0.0.1:{server.server_port}/v1")
finally:
server.shutdown()