diff --git a/mempalace/cli.py b/mempalace/cli.py index 608e9e7..3684a44 100644 --- a/mempalace/cli.py +++ b/mempalace/cli.py @@ -1199,6 +1199,214 @@ def cmd_status(args): status(palace_path=palace_path) +# ── Logstream (RFC 003 agent coordination) ──────────────────────────────── + + +def _open_logstream(args): + """Open the palace logstream database for a CLI command. + + Direct SQLite access is safe alongside a running hub: logstream.sqlite3 + is WAL-mode and independent of Chroma, so CLI writes are immediately + visible to hub readers without the mine-style forwarding Chroma needs. + """ + from .logstream import LOGSTREAM_DB_FILENAME, Logstream + + palace_path = os.path.expanduser(args.palace) if args.palace else MempalaceConfig().palace_path + return Logstream(db_path=os.path.join(palace_path, LOGSTREAM_DB_FILENAME)) + + +def _logstream_fail(message: str, as_json: bool): + import json + + if as_json: + print(json.dumps({"error": message})) + else: + print(f"Error: {message}", file=sys.stderr) + sys.exit(1) + + +def _read_text_arg(inline, file_arg, default=""): + """Resolve inline text vs --*-file (with '-' meaning stdin).""" + if inline is not None and file_arg is not None: + raise ValueError("pass inline text or a file, not both") + if file_arg is not None: + if file_arg == "-": + return sys.stdin.read() + return Path(os.path.expanduser(file_arg)).read_text(encoding="utf-8") + if inline is not None: + return inline + return default + + +def _parse_metadata_arg(raw): + import json + + if raw is None: + return None + try: + value = json.loads(raw) + except ValueError as exc: + raise ValueError(f"--metadata is not valid JSON: {exc}") from None + if not isinstance(value, dict): + raise ValueError("--metadata must be a JSON object") + return value + + +def _print_event_line(event): + target = event["to_agent"] or "*" + corr = f" corr={event['correlation_id']}" if event["correlation_id"] else "" + status = f" [{event['status']}]" if event["status"] else "" + arts = f" artifacts={len(event['artifact_ids'])}" if event["artifact_ids"] else "" + body = event["body"].replace("\n", " ") + if len(body) > 80: + body = body[:77] + "..." + body = f" :: {body}" if body else "" + print( + f" {event['id']} {event['created_at']} {event['type']} " + f"{event['stream']}/{event['room']} {event['from_agent']}->{target}" + f"{status}{corr}{arts}{body}" + ) + + +def cmd_logstream(args): + import json + + as_json = getattr(args, "json", False) + try: + ls = _open_logstream(args) + except Exception as exc: + _logstream_fail(str(exc), as_json) + try: + if args.logstream_action == "append": + try: + body = _read_text_arg(args.body, args.body_file) + event = ls.append_event( + type=args.type, + stream=args.stream, + room=args.room, + from_agent=args.from_agent, + to_agent=args.to_agent, + correlation_id=args.correlation_id, + branch=args.branch, + base_commit=args.base_commit, + status=args.status, + body=body, + metadata=_parse_metadata_arg(args.metadata), + artifact_ids=args.artifact_id or None, + ) + except (ValueError, OSError) as exc: + _logstream_fail(str(exc), as_json) + if as_json: + print(json.dumps(event, indent=2, ensure_ascii=False)) + else: + print("Appended:") + _print_event_line(event) + elif args.logstream_action in ("list", "wait"): + filters = { + "stream": args.stream, + "room": args.room, + "type": args.type, + "to_agent": args.to_agent, + "from_agent": args.from_agent, + "correlation_id": args.correlation_id, + "status": args.status, + "since_event_id": args.since_event_id, + "since_created_at": args.since_created_at, + } + try: + if args.logstream_action == "list": + events = ls.list_events(limit=args.limit, **filters) + result = {"events": events, "count": len(events)} + else: + result = ls.wait_events(timeout_ms=args.timeout_ms, **filters) + result["count"] = len(result["events"]) + except ValueError as exc: + _logstream_fail(str(exc), as_json) + if as_json: + print(json.dumps(result, indent=2, ensure_ascii=False)) + else: + if result.get("timed_out"): + print("Timed out; no matching events.") + elif not result["events"]: + print("No matching events.") + else: + print(f"{result['count']} event(s):") + for event in result["events"]: + _print_event_line(event) + if result.get("timed_out"): + sys.exit(2) + elif args.logstream_action == "ack": + try: + event = ls.ack_event( + args.event_id, + from_agent=args.from_agent, + status=args.status, + body=args.body or "", + ) + except ValueError as exc: + _logstream_fail(str(exc), as_json) + if as_json: + print(json.dumps(event, indent=2, ensure_ascii=False)) + else: + print("Acknowledged:") + _print_event_line(event) + finally: + ls.close() + + +def cmd_artifact(args): + import json + + as_json = getattr(args, "json", False) + try: + ls = _open_logstream(args) + except Exception as exc: + _logstream_fail(str(exc), as_json) + try: + if args.artifact_action == "put": + try: + content = _read_text_arg(args.content, args.file, default=None) + if content is None: + content = sys.stdin.read() + artifact = ls.put_artifact( + kind=args.kind, + content=content, + created_by=args.created_by, + metadata=_parse_metadata_arg(args.metadata), + ) + except (ValueError, OSError) as exc: + _logstream_fail(str(exc), as_json) + if as_json: + print(json.dumps(artifact, indent=2, ensure_ascii=False)) + else: + print(f"Stored {artifact['id']} kind={artifact['kind']}") + print(f" sha256={artifact['sha256']}") + print(f" size={artifact['size_bytes']} bytes") + elif args.artifact_action == "get": + try: + artifact = ls.get_artifact(args.artifact_id) + except ValueError as exc: + _logstream_fail(str(exc), as_json) + if artifact is None: + _logstream_fail(f"artifact {args.artifact_id!r} not found", as_json) + if args.out: + Path(os.path.expanduser(args.out)).write_text(artifact["content"], encoding="utf-8") + if as_json: + if args.out: + artifact = {**artifact, "content_written_to": args.out} + artifact.pop("content") + print(json.dumps(artifact, indent=2, ensure_ascii=False)) + elif args.out: + print(f"Wrote {artifact['size_bytes']} bytes to {args.out}") + print(f" sha256={artifact['sha256']}") + else: + # Exact content on stdout so `mempalace artifact get ID | git apply` + # works; metadata would corrupt the stream. + sys.stdout.write(artifact["content"]) + finally: + ls.close() + + def cmd_palace_set_embedder(args): """Record (or force-override) a palace's embedder identity (RFC 001). @@ -2385,6 +2593,114 @@ def main(): help="Storage backend to use for status (default: config/env/detected/chroma)", ) + # logstream (RFC 003 agent coordination) + p_logstream = sub.add_parser( + "logstream", + help="Agent coordination events — delegate work, wait for replies (RFC 003)", + ) + logstream_sub = p_logstream.add_subparsers(dest="logstream_action") + + def _add_logstream_filters(p): + p.add_argument("--stream", default=None, help="Stream, e.g. project/mempalace") + p.add_argument("--room", default=None, help="Room, e.g. delegation, patches") + p.add_argument("--type", default=None, help="Event type, e.g. task.request") + p.add_argument("--to-agent", default=None, help="Target agent (also matches '*')") + p.add_argument("--from-agent", default=None, help="Writer agent") + p.add_argument("--correlation-id", default=None, help="Task/conversation id") + p.add_argument( + "--status", + default=None, + help="open|claimed|ready|applied|blocked|failed|superseded", + ) + p.add_argument("--since-event-id", default=None, help="Only events strictly after this id") + p.add_argument( + "--since-created-at", + default=None, + help="Only events at/after this time (YYYY-MM-DD or YYYY-MM-DDTHH:MM:SSZ)", + ) + + p_ls_append = logstream_sub.add_parser("append", help="Append a coordination event") + p_ls_append.add_argument("--type", required=True, help="Event type, e.g. task.request") + p_ls_append.add_argument("--stream", required=True, help="Stream, e.g. project/mempalace") + p_ls_append.add_argument("--room", required=True, help="Room, e.g. delegation") + p_ls_append.add_argument("--from-agent", required=True, help="Writer agent identity") + p_ls_append.add_argument("--to-agent", default=None, help="Target agent or '*'") + p_ls_append.add_argument("--correlation-id", default=None, help="Task/conversation id") + p_ls_append.add_argument("--branch", default=None, help="Git branch") + p_ls_append.add_argument("--base-commit", default=None, help="Git commit work started from") + p_ls_append.add_argument( + "--status", + default=None, + help="open|claimed|ready|applied|blocked|failed|superseded", + ) + p_ls_append.add_argument("--body", default=None, help="Verbatim body text") + p_ls_append.add_argument( + "--body-file", default=None, help="Read body from file ('-' for stdin)" + ) + p_ls_append.add_argument("--metadata", default=None, help="Extra fields as a JSON object") + p_ls_append.add_argument( + "--artifact-id", + action="append", + default=None, + help="Reference an already-stored artifact (repeatable)", + ) + p_ls_append.add_argument("--json", action="store_true", help="Machine-readable output") + + p_ls_list = logstream_sub.add_parser("list", help="List events, oldest first") + _add_logstream_filters(p_ls_list) + p_ls_list.add_argument("--limit", type=int, default=50, help="Max events (default 50)") + p_ls_list.add_argument("--json", action="store_true", help="Machine-readable output") + + p_ls_wait = logstream_sub.add_parser( + "wait", help="Block until a matching event exists (exit 2 on timeout)" + ) + _add_logstream_filters(p_ls_wait) + p_ls_wait.add_argument( + "--timeout-ms", + type=int, + default=60000, + help="How long to wait in ms (default 60000, max 300000)", + ) + p_ls_wait.add_argument("--json", action="store_true", help="Machine-readable output") + + p_ls_ack = logstream_sub.add_parser( + "ack", help="Acknowledge an event (appends event.ack, never mutates)" + ) + p_ls_ack.add_argument("event_id", help="Event id to acknowledge") + p_ls_ack.add_argument("--from-agent", required=True, help="Acknowledging agent identity") + p_ls_ack.add_argument( + "--status", + default=None, + help="open|claimed|ready|applied|blocked|failed|superseded", + ) + p_ls_ack.add_argument("--body", default=None, help="Verbatim ack notes") + p_ls_ack.add_argument("--json", action="store_true", help="Machine-readable output") + + # artifact (RFC 003 exact content exchange) + p_artifact = sub.add_parser( + "artifact", help="Exact artifact exchange for agent handoffs (RFC 003)" + ) + artifact_sub = p_artifact.add_subparsers(dest="artifact_action") + + p_art_put = artifact_sub.add_parser("put", help="Store exact artifact content") + p_art_put.add_argument("--kind", required=True, help="patch|file|log|json|note") + p_art_put.add_argument("--created-by", required=True, help="Writer agent identity") + p_art_put.add_argument("--content", default=None, help="Inline content") + p_art_put.add_argument( + "--file", default=None, help="Read content from file ('-' for stdin; default stdin)" + ) + p_art_put.add_argument("--metadata", default=None, help="Extra fields as a JSON object") + p_art_put.add_argument("--json", action="store_true", help="Machine-readable output") + + p_art_get = artifact_sub.add_parser( + "get", help="Fetch exact artifact content (stdout pipes into git apply)" + ) + p_art_get.add_argument("artifact_id", help="Artifact id") + p_art_get.add_argument("--out", default=None, help="Write content to this file instead") + p_art_get.add_argument( + "--json", action="store_true", help="Metadata as JSON (content omitted with --out)" + ) + p_palace = sub.add_parser("palace", help="Palace maintenance commands") palace_sub = p_palace.add_subparsers(dest="palace_action") p_set_embedder = palace_sub.add_parser( @@ -2441,6 +2757,20 @@ def main(): p_palace.print_help() return + if args.command == "logstream": + if not getattr(args, "logstream_action", None): + p_logstream.print_help() + return + cmd_logstream(args) + return + + if args.command == "artifact": + if not getattr(args, "artifact_action", None): + p_artifact.print_help() + return + cmd_artifact(args) + return + if args.command == "daemon": if not getattr(args, "daemon_action", None): p_daemon.print_help() diff --git a/tests/test_cli_logstream.py b/tests/test_cli_logstream.py new file mode 100644 index 0000000..5e4d14d --- /dev/null +++ b/tests/test_cli_logstream.py @@ -0,0 +1,256 @@ +"""Tests for the RFC 003 logstream/artifact CLI commands. + +Covers cmd_logstream (append/list/wait/ack) and cmd_artifact (put/get): +JSON and human output, exact-content stdout piping, timeout exit code, +and error exits. Uses SimpleNamespace args like the rest of test_cli.py. +""" + +import json +import sys +from types import SimpleNamespace + +import pytest + +from mempalace.cli import cmd_artifact, cmd_logstream, main + + +def _append_args(palace, **overrides): + fields = dict( + palace=palace, + logstream_action="append", + type="task.request", + stream="project/mempalace", + room="delegation", + from_agent="mac-fable", + to_agent="windows-codex", + correlation_id="task_cli", + branch=None, + base_commit=None, + status=None, + body="Please fix the thing.", + body_file=None, + metadata=None, + artifact_id=None, + json=True, + ) + fields.update(overrides) + return SimpleNamespace(**fields) + + +def _list_args(palace, **overrides): + fields = dict( + palace=palace, + logstream_action="list", + stream=None, + room=None, + type=None, + to_agent=None, + from_agent=None, + correlation_id=None, + status=None, + since_event_id=None, + since_created_at=None, + limit=50, + json=True, + ) + fields.update(overrides) + return SimpleNamespace(**fields) + + +def _wait_args(palace, **overrides): + args = _list_args(palace, **overrides) + args.logstream_action = "wait" + del args.limit + if not hasattr(args, "timeout_ms"): + args.timeout_ms = 100 + return args + + +def _put_args(palace, content, **overrides): + fields = dict( + palace=palace, + artifact_action="put", + kind="patch", + created_by="windows-codex", + content=content, + file=None, + metadata=None, + json=True, + ) + fields.update(overrides) + return SimpleNamespace(**fields) + + +class TestLogstreamCli: + def test_append_then_list_json_round_trip(self, palace_path, capsys): + cmd_logstream(_append_args(palace_path)) + appended = json.loads(capsys.readouterr().out) + assert appended["id"].startswith("evt_") + + cmd_logstream(_list_args(palace_path, correlation_id="task_cli")) + listed = json.loads(capsys.readouterr().out) + assert listed["count"] == 1 + assert listed["events"][0]["id"] == appended["id"] + assert listed["events"][0]["body"] == "Please fix the thing." + + def test_append_human_output(self, palace_path, capsys): + cmd_logstream(_append_args(palace_path, json=False)) + out = capsys.readouterr().out + assert "Appended:" in out + assert "task.request" in out + assert "project/mempalace/delegation" in out + assert "mac-fable->windows-codex" in out + + def test_append_invalid_status_exits_1(self, palace_path, capsys): + with pytest.raises(SystemExit) as exc: + cmd_logstream(_append_args(palace_path, status="bogus")) + assert exc.value.code == 1 + assert "status" in json.loads(capsys.readouterr().out)["error"] + + def test_append_invalid_metadata_exits_1(self, palace_path, capsys): + with pytest.raises(SystemExit) as exc: + cmd_logstream(_append_args(palace_path, metadata="not json")) + assert exc.value.code == 1 + assert "metadata" in json.loads(capsys.readouterr().out)["error"] + + def test_body_file_stdin(self, palace_path, capsys, monkeypatch): + monkeypatch.setattr(sys, "stdin", __import__("io").StringIO("body from stdin\n")) + cmd_logstream(_append_args(palace_path, body=None, body_file="-")) + appended = json.loads(capsys.readouterr().out) + assert appended["body"] == "body from stdin\n" + + def test_wait_existing_event_returns_immediately(self, palace_path, capsys): + cmd_logstream(_append_args(palace_path)) + capsys.readouterr() + cmd_logstream(_wait_args(palace_path, correlation_id="task_cli", timeout_ms=5000)) + result = json.loads(capsys.readouterr().out) + assert result["timed_out"] is False + assert result["count"] == 1 + + def test_wait_timeout_exits_2(self, palace_path, capsys): + with pytest.raises(SystemExit) as exc: + cmd_logstream(_wait_args(palace_path, correlation_id="task_never", timeout_ms=100)) + assert exc.value.code == 2 + result = json.loads(capsys.readouterr().out) + assert result["timed_out"] is True + + def test_ack_round_trip(self, palace_path, capsys): + cmd_logstream(_append_args(palace_path)) + appended = json.loads(capsys.readouterr().out) + cmd_logstream( + SimpleNamespace( + palace=palace_path, + logstream_action="ack", + event_id=appended["id"], + from_agent="windows-codex", + status="applied", + body="Done.", + json=True, + ) + ) + ack = json.loads(capsys.readouterr().out) + assert ack["type"] == "event.ack" + assert ack["to_agent"] == "mac-fable" + assert ack["correlation_id"] == "task_cli" + + +class TestArtifactCli: + PATCH = "diff --git a/x b/x\n+cli\n" + + def test_put_then_get_stdout_is_exact(self, palace_path, capsys): + cmd_artifact(_put_args(palace_path, self.PATCH)) + artifact = json.loads(capsys.readouterr().out) + assert artifact["id"].startswith("art_") + + cmd_artifact( + SimpleNamespace( + palace=palace_path, + artifact_action="get", + artifact_id=artifact["id"], + out=None, + json=False, + ) + ) + # Exact content, nothing else — must survive `| git apply`. + assert capsys.readouterr().out == self.PATCH + + def test_get_out_writes_file(self, palace_path, tmp_dir, capsys): + cmd_artifact(_put_args(palace_path, self.PATCH)) + artifact = json.loads(capsys.readouterr().out) + out_path = f"{tmp_dir}/fetched.patch" + cmd_artifact( + SimpleNamespace( + palace=palace_path, + artifact_action="get", + artifact_id=artifact["id"], + out=out_path, + json=True, + ) + ) + meta = json.loads(capsys.readouterr().out) + assert "content" not in meta + assert meta["content_written_to"] == out_path + assert open(out_path, encoding="utf-8").read() == self.PATCH + + def test_get_missing_exits_1(self, palace_path, capsys): + with pytest.raises(SystemExit) as exc: + cmd_artifact( + SimpleNamespace( + palace=palace_path, + artifact_action="get", + artifact_id="art_nope", + out=None, + json=True, + ) + ) + assert exc.value.code == 1 + assert "not found" in json.loads(capsys.readouterr().out)["error"] + + def test_put_reads_stdin_by_default(self, palace_path, capsys, monkeypatch): + monkeypatch.setattr(sys, "stdin", __import__("io").StringIO(self.PATCH)) + cmd_artifact(_put_args(palace_path, None)) + artifact = json.loads(capsys.readouterr().out) + assert artifact["size_bytes"] == len(self.PATCH.encode("utf-8")) + + def test_event_can_reference_cli_artifact(self, palace_path, capsys): + cmd_artifact(_put_args(palace_path, self.PATCH)) + artifact = json.loads(capsys.readouterr().out) + cmd_logstream(_append_args(palace_path, type="patch.ready", artifact_id=[artifact["id"]])) + event = json.loads(capsys.readouterr().out) + assert event["artifact_ids"] == [artifact["id"]] + + +class TestMainDispatch: + def test_main_dispatches_logstream_list(self, palace_path, capsys, monkeypatch): + monkeypatch.setattr( + sys, + "argv", + ["mempalace", "--palace", palace_path, "logstream", "list", "--json"], + ) + main() + result = json.loads(capsys.readouterr().out) + assert result == {"events": [], "count": 0} + + def test_main_dispatches_artifact_put(self, palace_path, capsys, monkeypatch): + monkeypatch.setattr( + sys, + "argv", + [ + "mempalace", + "--palace", + palace_path, + "artifact", + "put", + "--kind", + "note", + "--created-by", + "mac-fable", + "--content", + "hello", + "--json", + ], + ) + main() + artifact = json.loads(capsys.readouterr().out) + assert artifact["kind"] == "note" + assert artifact["size_bytes"] == 5