"""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}