Newer
Older
navi-1 / navi / api / routes / peer.py
"""Peer-to-peer swarm endpoints — other navi instances talk to this one here.

The peer channel lives on the main API port (8099). Three endpoints:

  GET  /peer/hello   OPEN  — name, uuid prefix, version, port. A caller that
                            stumbled onto this port learns only that a navi
                            lives here; nothing sensitive is exposed.
  GET  /peer/status  PSK   — identity + machine facts + hive reachability.
  POST /peer/ask     PSK   — one-shot question answered by a real agent run
                            under settings.peer_ask_profile.

Loop safety: the answering agent run is created with ``peer`` excluded from
its tools, so a remote question can never recursively spawn another peer
ask — the recursion is cut at the tool level, deterministically. As a second
belt, an ask whose ``from_instance_id`` matches our own uuid is refused.

Audit: every ask is logged on both the request shape (peer name/uuid,
question length, duration) and the outcome — structlog events
``peer.ask_received`` / ``peer.ask_answered`` / ``peer.ask_failed``.
"""

from __future__ import annotations

import asyncio
import os
import platform
import time
import uuid

import structlog
from fastapi import APIRouter, Header, HTTPException
from pydantic import BaseModel, Field

from navi import __version__
from navi.config import settings
from navi.identity import load_identity

router = APIRouter(prefix="/peer", tags=["peer"])

log = structlog.get_logger()

_PROCESS_START = time.time()

# Peer asks run real LLM turns — a burst of them would stack GPU work. One
# concurrent peer answer at a time; extra callers wait for the semaphore.
_ASK_SEMAPHORE = asyncio.Semaphore(1)

_MAX_QUESTION_CHARS = 4000
_MAX_ANSWER_CHARS = 16_000


def _identity() -> dict:
    try:
        ident = load_identity(settings.instance_file)
        return {"name": ident.name, "instance_id": ident.instance_id}
    except Exception:
        # A navi whose identity file vanished still answers hello/status;
        # ask refuses below (nobody to pin the answer to).
        return {"name": "unnamed", "instance_id": ""}


def _require_swarm_key(x_swarm_key: str | None) -> None:
    from navi.swarm import verify_swarm_key

    if not x_swarm_key:
        raise HTTPException(status_code=401, detail="X-Swarm-Key header required")
    if not verify_swarm_key(x_swarm_key):
        raise HTTPException(status_code=403, detail="invalid swarm key")


@router.get("/hello")
async def peer_hello() -> dict:
    ident = _identity()
    return {
        "navi": True,
        "name": ident["name"],
        "instance_id": ident["instance_id"][:8] if ident["instance_id"] else "",
        "version": __version__,
        "port": settings.navi_port,
    }


@router.get("/status")
async def peer_status(x_swarm_key: str | None = Header(None)) -> dict:
    _require_swarm_key(x_swarm_key)
    ident = _identity()
    status: dict = {
        "name": ident["name"],
        "instance_id": ident["instance_id"],
        "version": __version__,
        "port": settings.navi_port,
        "uptime_sec": round(time.time() - _PROCESS_START),
    }
    from navi.swarm import get_announcer

    announcer = get_announcer()
    status["hive"] = announcer.get_status() if announcer else {"configured": False}
    try:
        status["machine"] = {
            "hostname": os.uname().nodename,
            "os": f"{platform.system()} {platform.release()}",
            "cpu_cores": os.cpu_count(),
        }
    except Exception:  # machine facts are best-effort
        pass
    return status


class PeerAskPayload(BaseModel):
    from_name: str = Field(min_length=1, max_length=64)
    from_instance_id: str = Field(min_length=1, max_length=64)
    question: str = Field(min_length=1, max_length=_MAX_QUESTION_CHARS)


@router.post("/ask")
async def peer_ask(
    payload: PeerAskPayload,
    x_swarm_key: str | None = Header(None),
) -> dict:
    _require_swarm_key(x_swarm_key)
    ident = _identity()
    if not ident["instance_id"]:
        raise HTTPException(status_code=503, detail="instance identity missing")
    # Loop guard (belt): a question claiming to come FROM us has travelled
    # a full circle — refuse it instead of answering ourselves recursively.
    if payload.from_instance_id == ident["instance_id"]:
        log.warning("peer.ask_loop_detected", from_name=payload.from_name)
        raise HTTPException(status_code=409, detail="ask loop detected - this question originated here")

    request_id = uuid.uuid4().hex[:12]
    log.info(
        "peer.ask_received",
        request_id=request_id,
        from_name=payload.from_name,
        from_instance_id=payload.from_instance_id,
        question_chars=len(payload.question),
    )

    briefing = (
        f"Another navi instance '{payload.from_name}' "
        f"(instance {payload.from_instance_id[:8]}) is asking you a question "
        "over the swarm peer channel. Answer for a fellow autonomous agent: "
        "be direct and factual, include concrete values/paths/commands where "
        "relevant. You cannot ask other peers from here — if the answer needs "
        "something you cannot check, say so explicitly."
    )
    started = time.time()
    async with _ASK_SEMAPHORE:
        try:
            from navi.api.deps import get_agent

            agent = get_agent()
            answer, completed = await asyncio.wait_for(
                agent.run_ephemeral(
                    user_message=payload.question,
                    profile_id=settings.peer_ask_profile,
                    exclude_tools=["peer"],  # deterministic loop cut
                    briefing=briefing,
                    timeout_seconds=settings.peer_ask_timeout_sec,
                ),
                timeout=settings.peer_ask_timeout_sec + 30,  # LLM ceiling + slack
            )
        except asyncio.TimeoutError:
            log.warning("peer.ask_timeout", request_id=request_id, peer=payload.from_name,
                        timeout=settings.peer_ask_timeout_sec)
            raise HTTPException(status_code=504, detail="peer ask timed out")
        except Exception as e:
            log.error("peer.ask_failed", request_id=request_id, peer=payload.from_name,
                      error=f"{type(e).__name__}: {e}")
            raise HTTPException(status_code=500, detail="peer ask failed")

    answer = (answer or "")[:_MAX_ANSWER_CHARS]
    log.info(
        "peer.ask_answered",
        request_id=request_id,
        peer=payload.from_name,
        completed=completed,
        duration_sec=round(time.time() - started, 1),
        answer_chars=len(answer),
    )
    return {"answer": answer, "completed": completed, "request_id": request_id}