Something went wrong. Try again.
small gleam coding and (not yet) persistent agent daemon with a detachable cli
Something went wrong. Try again.
Python
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378"""Real daemon extension selection, context loading, and live reload. No live model."""import contextlibimport http.serverimport jsonimport osfrom pathlib import Pathimport subprocessimport tempfileimport threadingimport timeimport urllib.errorimport 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("<system-note>\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("</system-note>"), 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("<system-note>") == 1 and work_note.count("</system-note>") == 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("<system-note>\n[albedo] This session's extensions changed")) assert "Removed context:\n- skills" in update and "<available_skills>" 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()