"""Initial tamq command-line surface.""" from __future__ import annotations import argparse import asyncio import json import os import signal import sqlite3 import fcntl from uuid import uuid4 import subprocess import sys import time from datetime import datetime, timedelta, timezone from pathlib import Path from . import __version__ from .config import db_path, pid_path, lock_path, config_path, socket_path, state_dir, setting from .service import PUSHY_FRAMING_CAPABILITY, Service, ping, request from .store import Store from .tmux import TmuxError, TmuxManager from .broker import BrokerIdentity, InputBroker from .ptytap import PtyTap from .control import ControlModeClient from .composer import ComposerError, compose_message, recipient_order from .cleanup import cleanup_runtime from .diagnostics import configure from .policy import load_profile from .registry import RegistryError, validate_targets from .terminal import format_comment SUBCOMMANDS = frozenset( "start attach serve stop cleanup status ping inbox history inspect ack send reply export replay purge completion db-version config tap".split() ) START_OPTIONS = frozenset({"--command", "--cmd", "--tap", "--detach", "--no-service", "--no-display", "--mode"}) GLOBAL_FLAGS = frozenset({"--orwell", "--verbose"}) def parse_size(value: str) -> int: units = {"B": 1, "KB": 1000, "MB": 1000**2, "GB": 1000**3} text = value.upper() for unit, multiplier in sorted(units.items(), key=lambda item: -len(item[0])): if text.endswith(unit): return int(float(text[:-len(unit)]) * multiplier) return int(text) def history_advisory(store: Store) -> str | None: size, oldest = store.history_stats() limit_text = setting("history_max_size", "100MB") limit = parse_size(limit_text) if size <= limit: return None since = time.strftime("%y%m%d", time.localtime(oldest)) if oldest else "unknown" return f"Size of history is {size / 1000**2:.0f} MB since {since}, exceeding {limit_text}. Consider using 'tamq purge'." def ensure_service() -> bool: if asyncio.run(ping()): return True path = lock_path() path.parent.mkdir(parents=True, exist_ok=True) with path.open("a+") as lock: fcntl.flock(lock.fileno(), fcntl.LOCK_EX) if asyncio.run(ping()): return True subprocess.Popen( [sys.executable, "-m", "tamq.cli", "serve"], stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True, env=os.environ.copy(), ) for _ in range(20): if asyncio.run(ping()): return True time.sleep(0.05) return False def stop_service_process() -> int | None: try: pid = int(pid_path().read_text(encoding="utf-8")) except (FileNotFoundError, ValueError): return None try: os.kill(pid, signal.SIGTERM) except ProcessLookupError: pid_path().unlink(missing_ok=True) return None deadline = time.monotonic() + 2.0 while time.monotonic() < deadline: try: os.kill(pid, 0) except ProcessLookupError: break time.sleep(0.05) return pid def ensure_service_capabilities(required: set[str]) -> bool: """Restart an older broker once before registering a newer endpoint mode.""" if not ensure_service(): return False try: capabilities = asyncio.run(request({"op": "ping"})).get("capabilities", []) except (OSError, json.JSONDecodeError): return False if required.issubset(capabilities): return True stop_service_process() if not ensure_service(): return False try: capabilities = asyncio.run(request({"op": "ping"})).get("capabilities", []) except (OSError, json.JSONDecodeError): return False return required.issubset(capabilities) def ensure_manual_service() -> bool: return ensure_service_capabilities({"manual_delivery"}) def ensure_output_service() -> bool: return ensure_service_capabilities({"manual_delivery", "terminal_output"}) def ensure_pushy_service() -> bool: return ensure_service_capabilities( {"manual_delivery", "pushy_input", PUSHY_FRAMING_CAPABILITY} ) def preflight_runtime_paths() -> None: """Ensure local state/runtime parents are available before tmux mutation.""" parents = {db_path().parent, socket_path().parent, pid_path().parent, lock_path().parent} for parent in parents: parent.mkdir(parents=True, exist_ok=True) if not os.access(parent, os.W_OK | os.X_OK): raise OSError(f"runtime directory is not writable: {parent}") def attach_session(session: str) -> int: command = ["tmux", "switch-client" if os.environ.get("TMUX") else "attach-session", "-t", session] return subprocess.run(command, check=False).returncode def format_comment_message(row: sqlite3.Row) -> str: """Render a durable message so every displayed line remains a shell comment.""" return format_comment(row["sender_repo"], row["body"], row["message_id"]) def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser( prog="tamq", description="Repository-aware tmux sessions with durable local messaging.", formatter_class=argparse.RawDescriptionHelpFormatter, epilog="""session shorthand: tamq [--detach] [--command COMMAND] [--mode MODE] REPO [REPO ...] With no --command, tamq opens ordinary repository shells. --command is an explicit initial command for newly created windows. Messages appear as sanitized terminal output by default. --mode inbox keeps messages durable without display; experimental --mode pushy submits them to target input. manual messaging from a managed shell: @TARGET: MESSAGE... @ compose with a visible, Tab-selectable recipient @ MESSAGE... fast reply (ordinary shell quoting applies) #TARGET: MESSAGE... agent-input alias in a tapped pushy session tamq inbox [--filter COMMAND] """, ) parser.add_argument("--version", "-V", action="version", version=__version__) parser.add_argument("--orwell", action="store_true", help="enable unsafe local diagnostics") parser.add_argument("--verbose", action="store_true", help="enable debug diagnostics") parser.add_argument("--policy-profile", default=None, help="select configured policy profile") subparsers = parser.add_subparsers(dest="command") start = subparsers.add_parser("start", help="start or reuse a tamq endpoint") start.add_argument("repos", nargs="*", help="gita-registered repository slugs") start.add_argument("--command", "--cmd", dest="initial_command", default=None, help="explicit initial command for newly created windows (default: ordinary shell)") start.add_argument("--tap", action="store_true", help="explicitly opt in to PTY input observation and pane message injection") start.add_argument("--mode", choices=("output", "inbox", "pushy"), help="message delivery mode (default: output; pushy is experimental)") start.add_argument("--no-display", action="store_true", help="compatibility alias for --mode inbox") start.add_argument("--detach", action="store_true", help="detach after startup") start.add_argument("--no-service", action="store_true", help="open tmux windows without service registration or messaging") attach = subparsers.add_parser("attach", help="attach to the managed tmux session") attach.add_argument("--session", default="tamq", help="tmux session name") subparsers.add_parser("serve", help="run the local service in the foreground") subparsers.add_parser("stop", help="stop the local service") cleanup = subparsers.add_parser( "cleanup", help="dry-run emergency shutdown and tamq-owned runtime cleanup", ) cleanup.add_argument( "--yes", action="store_true", help="stop the broker, close the managed session, and apply cleanup", ) cleanup.add_argument("--session", default="tamq", help="managed tmux session") subparsers.add_parser("status", help="show endpoint health") subparsers.add_parser("ping", help="check local service liveness") history = subparsers.add_parser("history", help="inspect local message history") history.add_argument("--repo", dest="target_repo") history.add_argument("--state") inbox = subparsers.add_parser("inbox", help="show durable messages for a repository without injecting terminal input") inbox.add_argument("--repo", dest="target_repo", help="target repository (default: TAMQ_REPO in a managed window)") inbox.add_argument("--all", action="store_true", help="include non-pending messages") inbox.add_argument("--json", action="store_true", help="emit one JSON object per message") inbox.add_argument( "--filter", dest="filter_command", help="consume each pending comment through COMMAND and acknowledge it after a zero exit", ) inspect = subparsers.add_parser("inspect", help="inspect one message") inspect.add_argument("message_id") ack = subparsers.add_parser("ack", help="acknowledge one durable message") ack.add_argument("message_id") send = subparsers.add_parser("send", help="queue a direct message") send.add_argument("address", help="@repo: message") send.add_argument("body", nargs="*", help="message body when address is a repo slug") send.add_argument("--endpoint-id") send.add_argument("--from", dest="sender_repo", help="sender repository (default: TAMQ_REPO or local)") reply = subparsers.add_parser("reply", help="reply to the latest sender for the current repository") reply.add_argument("body", nargs="*", help="message body; omit to open the interactive composer") export = subparsers.add_parser("export", help="export history as JSONL") export.add_argument("--output", required=True) export.add_argument("--repo", dest="target_repo") export.add_argument("--state") replay = subparsers.add_parser("replay", help="replay a JSONL message file") replay.add_argument("file") replay.add_argument("--endpoint-id", help="target a registered tmux-amq endpoint") purge = subparsers.add_parser("purge", help="dry-run or remove history") purge.add_argument("--before", default=None) purge.add_argument("--max-size", default=None) purge.add_argument( "--feedback-chain", metavar="MESSAGE_ID", help="select one reflected-delivery chain by its durable root message", ) purge.add_argument("--yes", action="store_true") purge.add_argument("--orwell", action="store_true") completion = subparsers.add_parser("completion", help="emit shell completion script") completion.add_argument("shell", choices=("bash", "zsh", "fish")) subparsers.add_parser("db-version", help="show SQLite schema version") subparsers.add_parser("config", help="show effective configuration paths and policy") tap = subparsers.add_parser("tap", help="explicitly run a command behind the full-duplex PTY tap") tap.add_argument("--repo", required=True, help="source gita repository slug") tap.add_argument("--endpoint", required=True, help="tmux-amq endpoint identity") tap.add_argument("wrapped_command", nargs=argparse.REMAINDER, help="command after --") return parser def normalize_argv(argv: list[str]) -> list[str]: """Allow repository-first startup while retaining explicit subcommands.""" if not argv: return argv index = 0 while index < len(argv): token = argv[index] if token in GLOBAL_FLAGS or token.startswith("--policy-profile="): index += 1 continue if token == "--policy-profile" and index + 1 < len(argv): index += 2 continue break if index == len(argv): return argv candidate = argv[index] if candidate in START_OPTIONS or candidate.startswith("--mode=") or ( not candidate.startswith("-") and candidate not in SUBCOMMANDS ): return [*argv[:index], "start", *argv[index:]] return argv def main(argv: list[str] | None = None) -> int: parser = build_parser() raw_argv = list(sys.argv[1:] if argv is None else argv) args = parser.parse_args(normalize_argv(raw_argv)) configure(args.orwell, args.verbose) try: profile = load_profile(selected=args.policy_profile or os.environ.get("TAMQ_POLICY_PROFILE", "default")) except ValueError as exc: print(f"tamq: {exc}", file=sys.stderr) return 2 if args.command is None: parser.print_help() return 0 if args.command == "db-version": store = Store(db_path()); print(store.db.execute("SELECT value FROM metadata WHERE key='schema_version'").fetchone()[0]); store.close(); return 0 if args.command == "ping": return 0 if asyncio.run(ping()) else 1 if args.command == "stop": pid = stop_service_process() if pid is not None: print(f"stopped tamq service {pid}") return 0 print("tamq service is not running", file=sys.stderr) return 1 if args.command == "cleanup": report = cleanup_runtime( apply=args.yes, session=args.session, stop_service=stop_service_process, ) print(json.dumps(report, sort_keys=True)) return 1 if report["errors"] else 0 if args.command == "serve": try: asyncio.run(Service().run()) except KeyboardInterrupt: return 0 return 0 if args.command == "attach": return attach_session(args.session) if args.command == "completion": print(completion_script(args.shell), end="") return 0 if args.command == "config": print(json.dumps({"config": str(config_path()), "state_dir": str(state_dir()), "database": str(db_path()), "socket": str(socket_path()), "pidfile": str(pid_path()), "lockfile": str(lock_path()), "policy_profile": profile.name, "delivery_ack_mode": profile.delivery_ack_mode}, sort_keys=True)) return 0 if args.command == "tap": command = list(args.wrapped_command) if command and command[0] == "--": command = command[1:] if not command: print("tamq tap requires a command after --", file=sys.stderr) return 2 store = Store(db_path()) try: return PtyTap(command, InputBroker(store, BrokerIdentity(args.endpoint, args.repo))).run() finally: store.close() if args.command == "start": if args.no_display and args.mode not in (None, "inbox"): print("tamq: --no-display conflicts with the selected --mode", file=sys.stderr) return 2 selected_mode = args.mode or ("inbox" if args.no_display else "output") endpoint_delivery_mode = { "inbox": "manual", "output": "output", "pushy": "pushy", }[selected_mode] if args.tap and args.initial_command is None: print("tamq: --tap requires an explicit --command", file=sys.stderr) return 2 if args.tap and args.no_service: print("tamq: --tap cannot be combined with --no-service", file=sys.stderr) return 2 if args.tap and args.no_display: print("tamq: --tap cannot be combined with --no-display", file=sys.stderr) return 2 if args.tap and args.mode is not None: print("tamq: --tap cannot be combined with --mode", file=sys.stderr) return 2 if args.no_service and args.mode is not None: print("tamq: --mode cannot be combined with --no-service", file=sys.stderr) return 2 tap_enabled = args.tap or ( selected_mode == "pushy" and args.initial_command is not None ) manager = TmuxManager() try: if not args.no_service: preflight_runtime_paths() launch_plan = manager.preflight(args.repos, args.initial_command) except (TmuxError, OSError) as exc: print(f"tamq: {exc}", file=sys.stderr) return 2 if args.no_service: service_ready = True elif selected_mode == "inbox": service_ready = ensure_manual_service() elif selected_mode == "pushy": service_ready = ensure_pushy_service() else: service_ready = ensure_output_service() if not service_ready: print("tamq service failed to start with the required delivery capability; use --no-service to open repos without messaging", file=sys.stderr) return 1 try: endpoint = manager.ensure_plan(launch_plan, tap=tap_enabled) except (TmuxError, OSError) as exc: print(f"tamq: {exc}", file=sys.stderr) return 2 delivered = 0 registered = False if not args.no_service: try: delivery_mode = "pane" if args.tap else endpoint_delivery_mode registration = asyncio.run(request({"op": "register", "endpoint_id": endpoint.endpoint_id, "instance_id": endpoint.instance_key, "pid": endpoint.pid, "session": endpoint.session, "repos": endpoint.repos, "delivery_mode": delivery_mode})) if not registration.get("ok"): raise RuntimeError(f"endpoint registration failed: {registration.get('error', 'unknown error')}") registered = True if args.tap: control = ControlModeClient(endpoint.session) queue_store = Store(db_path()) try: control.start() delivered = InputBroker(queue_store, BrokerIdentity(endpoint.instance_key, endpoint.repos[0])).deliver_pending( control, window_for_repo=lambda repo: f"{endpoint.session}:{repo}" ) finally: control.close() queue_store.close() except (OSError, RuntimeError, ValueError, json.JSONDecodeError) as exc: manager.rollback(endpoint) print(f"tamq: {exc}", file=sys.stderr) return 1 summary = { "endpoint_id": endpoint.endpoint_id, "instance_id": endpoint.instance_key, "session": endpoint.session, "repos": endpoint.repos, "delivered": delivered, "service": not args.no_service, "messaging": registered, "mode": "none" if args.no_service else ("tap" if args.tap else selected_mode), "input_observation_requested": tap_enabled, "delivery_mode": "none" if args.no_service else ( "pane" if args.tap else endpoint_delivery_mode ), } print(json.dumps(summary), flush=True) if args.detach: return 0 return attach_session(endpoint.session) store = Store(db_path()) try: advisory = history_advisory(store) if advisory and args.command in {"start", "serve", "status", "history", "inbox", "send", "reply"}: print(advisory, file=sys.stderr) if args.command in {"send", "reply"}: endpoint_id = None if args.command == "reply": sender = os.environ.get("TAMQ_REPO") if not sender: print("tamq reply requires a managed window with TAMQ_REPO", file=sys.stderr) return 2 latest = store.latest_counterparty(sender) if args.body: body = " ".join(args.body).strip() target = latest if target is None: print(f"tamq: no counterparty has sent a message to {sender}", file=sys.stderr) return 2 else: try: encoded_recipients = json.loads(os.environ.get("TAMQ_RECIPIENTS", "[]")) except json.JSONDecodeError: encoded_recipients = [] session_repos = ( encoded_recipients if isinstance(encoded_recipients, list) and all(isinstance(repo, str) for repo in encoded_recipients) else [] ) recipients = recipient_order(sender, latest, session_repos) try: composed = compose_message(recipients) except ComposerError as exc: print(f"tamq: {exc}", file=sys.stderr) return 2 if composed is None: return 130 target, body = composed else: text = " ".join([args.address, *args.body]).strip() if not text.startswith("@") or ":" not in text: print("tamq send expects @repo: message", file=sys.stderr); return 2 target, body = text[1:].split(":", 1); body = body.strip() if not target or not body: print("tamq send requires a target and body", file=sys.stderr); return 2 sender = args.sender_repo or os.environ.get("TAMQ_REPO") or "local" endpoint_id = args.endpoint_id if not body: print("tamq reply requires a message body", file=sys.stderr) return 2 try: validate_targets([target]) if sender != "local": validate_targets([sender]) except RegistryError as exc: print(f"tamq: {exc}", file=sys.stderr); return 2 if asyncio.run(ping()): payload = {"op": "send", "sender_repo": sender, "target_repo": target, "body": body} if endpoint_id: payload["endpoint_id"] = endpoint_id response = asyncio.run(request(payload)) if not response.get("ok"): print(f"tamq: {response.get('error', 'send failed')}", file=sys.stderr); return 1 print(response["message_id"]); return 0 if endpoint_id: print("tamq: service is not running", file=sys.stderr); return 1 try: print(store.add(sender, target, body)) except ValueError as exc: print(f"tamq: {exc}", file=sys.stderr); return 2 return 0 if args.command == "inbox": target = args.target_repo or os.environ.get("TAMQ_REPO") if not target: print("tamq inbox requires --repo outside a managed tamq window", file=sys.stderr) return 2 try: validate_targets([target]) except RegistryError as exc: print(f"tamq: {exc}", file=sys.stderr); return 2 if args.filter_command is not None and not args.filter_command.strip(): print("tamq: --filter command must not be empty", file=sys.stderr) return 2 if args.filter_command and (args.all or args.json): print("tamq: --filter cannot be combined with --all or --json", file=sys.stderr) return 2 rows = store.list(target, None if args.all else "pending") for row in rows: if args.json: print(json.dumps(dict(row), sort_keys=True)) elif args.filter_command: environment = os.environ.copy() environment.update( { "TAMQ_MESSAGE_ID": row["message_id"], "TAMQ_SENDER_REPO": row["sender_repo"], "TAMQ_TARGET_REPO": row["target_repo"], } ) result = subprocess.run( args.filter_command, shell=True, input=format_comment_message(row) + "\n", text=True, env=environment, check=False, ) if result.returncode: print( f"tamq: filter failed for {row['message_id']} with exit status {result.returncode}; message remains pending", file=sys.stderr, ) return 1 store.acknowledge(row["message_id"]) else: print(format_comment_message(row)) return 0 if args.command == "history": for row in store.list(args.target_repo, args.state): print(json.dumps(dict(row), sort_keys=True)) return 0 if args.command == "inspect": rows = [row for row in store.list() if row["message_id"] == args.message_id] if not rows: print("tamq: message not found", file=sys.stderr) return 1 print(json.dumps(dict(rows[0]), sort_keys=True)); return 0 if args.command == "ack": if not store.acknowledge(args.message_id): print("tamq: message not found", file=sys.stderr); return 1 print(args.message_id); return 0 if args.command == "export": store.export(store.list(args.target_repo, args.state), Path(args.output)); return 0 if args.command == "replay": batch_id = f"replay-{uuid4()}" count = 0 if args.endpoint_id and store.endpoint(args.endpoint_id) is None: print(f"tamq: endpoint is not registered: {args.endpoint_id}", file=sys.stderr) return 1 try: lines = Path(args.file).read_text(encoding="utf-8").splitlines() for line in lines: if not line.strip(): continue item = json.loads(line) provenance = json.dumps({"batch_id": batch_id, "original_message_id": item.get("message_id")}, sort_keys=True) endpoint = args.endpoint_id or item.get("endpoint_id") store.add(item.get("sender_repo", "replay"), item["target_repo"], item["body"], endpoint=endpoint, provenance=provenance) count += 1 except (OSError, KeyError, TypeError, ValueError, json.JSONDecodeError) as exc: print(f"tamq: cannot replay {args.file}: {exc}", file=sys.stderr) return 2 print(json.dumps({"batch_id": batch_id, "count": count})); return 0 if args.command == "purge": if args.feedback_chain: if args.before or args.max_size: print( "tamq: --feedback-chain cannot be combined with --before or --max-size", file=sys.stderr, ) return 2 matching = store.feedback_chain(args.feedback_chain) if not matching: print("tamq: feedback-chain root not found", file=sys.stderr) return 1 if not args.yes: print( f"Dry run: {len(matching)} message(s) match feedback chain rooted at {args.feedback_chain}; re-run with --yes to purge." ) return 0 count = store.delete_messages( row["message_id"] for row in matching ) print(f"Purged {count} feedback-chain message(s).") return 0 before_text = args.before or setting("purge_before", "365d") max_size_text = args.max_size or setting("purge_max_size", "100MB") before = time.time() - int(before_text[:-1]) * 86400 if before_text.endswith("d") else None if not args.yes: matching = [row for row in store.list() if before is None or row["created_at"] < before] print(f"Dry run: {len(matching)} message(s) match; re-run with --yes to purge.") return 0 count = store.purge(before=before, max_bytes=parse_size(max_size_text)) print(f"Purged {count} message(s)."); return 0 if args.command == "status": live = asyncio.run(ping()) endpoints = asyncio.run(request({"op": "endpoints"})) if live else {"endpoints": []} print(json.dumps({"version": __version__, "db": str(db_path()), "messages": len(store.list()), "pending": len(store.list(state="pending")), "leases": store.db.execute("SELECT COUNT(*) FROM leases").fetchone()[0], "service": live, "policy_profile": profile.name, "delivery_ack_mode": profile.delivery_ack_mode, "endpoints": endpoints.get("endpoints", [])}, sort_keys=True)); return 0 print(f"tamq {__version__}: command '{args.command}' is not implemented yet", file=sys.stderr) return 2 finally: store.close() def completion_script(shell: str) -> str: commands = " ".join(sorted(SUBCOMMANDS)) if shell == "bash": return f"""_tamq_complete() {{ local commands=\"{commands}\" if [[ ${{COMP_WORDS[1]}} == start || ${{COMP_WORDS[1]}} == send ]]; then COMPREPLY=( $(compgen -W \"$(gita ls 2>/dev/null)\" -- \"${{COMP_WORDS[COMP_CWORD]}}\") ) else COMPREPLY=( $(compgen -W \"$commands\" -- \"${{COMP_WORDS[COMP_CWORD]}}\") ) fi }} complete -F _tamq_complete tamq """ if shell == "zsh": return f"#compdef tamq\n_arguments '1:command:(({commands}))' '*:repo:->repos'\nif [[ $state == repos ]]; then _values 'gita repo' $(gita ls 2>/dev/null); fi\n" return f"complete -c tamq -f -n '__fish_seen_subcommand_from start send' -a '(gita ls 2>/dev/null)'\ncomplete -c tamq -f -a '{commands}'\n" if __name__ == "__main__": raise SystemExit(main())