diff --git a/AGENTS.md b/AGENTS.md index 63ed258c82c..1c0a979789b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -56,7 +56,7 @@ Repository map and Reference Documentation sections below. | Translators | `open-sse/translator/` | Format conversion (OpenAI↔Claude↔Gemini) | | Transformer | `open-sse/transformer/` | Responses API ↔ Chat Completions | | Services | `open-sse/services/` | Combo routing, rate limits, caching, etc | -| Database | `src/lib/db/` | SQLite domain modules (170 migrations) | +| Database | `src/lib/db/` | SQLite domain modules (171 migrations) | | Domain/Policy | `src/domain/` | Policy engine, cost rules, fallback logic | | MCP Server | `open-sse/mcp-server/` | 110 tools (45 canonical + memory/skill/GitHub/pool/gamification/plugin/Notion/Obsidian/local-corpus/RTK modules), 3 transports (stdio / SSE / Streamable HTTP), 33 scopes | | A2A Server | `src/lib/a2a/` | JSON-RPC 2.0 agent protocol | diff --git a/README.md b/README.md index 519a64682f7..8753569c474 100644 --- a/README.md +++ b/README.md @@ -1244,7 +1244,7 @@ Métricas canônicas em 2026-08-24: **1.029 vídeos únicos** · **11.132.922 vi RuntimeNode.js 22.x / 24.x LTS — >=22.22.2 <23 || >=24.0.0 <27 LanguageTypeScript 6.0 — 100% TypeScript across src/ and open-sse/ (zero any in core since v2.0) FrameworkNext.js 16 + React 19 + Tailwind CSS 4 - Databasebetter-sqlite3 (SQLite, WAL journaling) + LowDB (JSON legacy) — 122 domain modules, 170 migrations + Databasebetter-sqlite3 (SQLite, WAL journaling) + LowDB (JSON legacy) — 122 domain modules, 171 migrations MemorySQLite FTS5 full-text + int8-quantized vector embeddings, typed decay SchemasZod 4 — MCP tool I/O validation + API contracts ProtocolsMCP (stdio / HTTP / SSE) + A2A v0.3 (JSON-RPC 2.0 + SSE) diff --git a/llm.txt b/llm.txt index e5c68d9d3a6..2ebe387b797 100644 --- a/llm.txt +++ b/llm.txt @@ -14,7 +14,7 @@ OmniRoute solves the problem of managing multiple AI provider subscriptions, quo - **Runtime:** Node.js `>=22.22.2 <23 || >=24.0.0 <27`, ES Modules (`"type": "module"`) - **Framework:** Next.js 16 (App Router) with TypeScript 6 -- **Database:** SQLite via better-sqlite3 (local, zero-config, 170 migrations) +- **Database:** SQLite via better-sqlite3 (local, zero-config, 171 migrations) - **State management:** Zustand (client), SQLite (server persistence) - **UI:** React 19, Tailwind CSS 4, Recharts for analytics, @lobehub/icons for 130+ provider SVG icons - **Auth:** OAuth 2.0 (PKCE) for providers, bcrypt for local user auth @@ -124,7 +124,7 @@ OmniRoute solves the problem of managing multiple AI provider subscriptions, quo │ │ │ ├── secrets.ts # Secrets management │ │ │ ├── stateReset.ts # State reset utilities │ │ │ ├── migrationRunner.ts # Schema migration runner -│ │ │ └── migrations/ # 170 versioned SQL migration files +│ │ │ └── migrations/ # 171 versioned SQL migration files │ │ ├── evals/ # Eval runner and scheduler │ │ ├── memory/ # Persistent conversational memory │ │ │ ├── extraction.ts # Memory extraction from conversations @@ -389,7 +389,7 @@ diagnostics) plus **memory**, **skill**, **agentSkill**, **githubSkill**, **pool 8. **ProviderIcon component:** Unified icon system using `@lobehub/icons` (130+ SVG) with PNG fallback and generic icon fallback chain. Used on providers, dashboard, and agents pages. -9. **DB architecture:** `localDb.ts` is a re-export layer only — real logic lives in 122 `src/lib/db/` modules with 170 SQL migrations. +9. **DB architecture:** `localDb.ts` is a re-export layer only — real logic lives in 122 `src/lib/db/` modules with 171 SQL migrations. 10. **Upstream headers:** Custom headers merged in executors after default auth; same header name replaces executor value. Forbidden header names in `src/shared/constants/upstreamHeaders.ts`. @@ -433,7 +433,7 @@ diagnostics) plus **memory**, **skill**, **agentSkill**, **githubSkill**, **pool 4. **Environment variables:** All configuration is in `.env` (from `.env.example`). Key vars: `PORT`, `NEXT_PUBLIC_BASE_URL`, `API_KEY`, `ADMIN_PASSWORD`. -5. **Database layer:** Operations go through `src/lib/db/` modules (122 domain-specific files, 170 migrations). `localDb.ts` is re-exports only — add new functions to the proper `db/*.ts` module. +5. **Database layer:** Operations go through `src/lib/db/` modules (122 domain-specific files, 171 migrations). `localDb.ts` is re-exports only — add new functions to the proper `db/*.ts` module. 6. **Tests** use Node.js built-in test runner + Vitest. Run `npm test`. Vitest for MCP/autoCombo (`npm run test:vitest`). Playwright for E2E (`npm run test:e2e`). Coverage gate: ratchet vs `quality-baseline.json`, absolute floor 60% statements/lines/functions/branches. diff --git a/open-sse/ag-ui/index.ts b/open-sse/ag-ui/index.ts new file mode 100644 index 00000000000..a8e7beca910 --- /dev/null +++ b/open-sse/ag-ui/index.ts @@ -0,0 +1,90 @@ +/** + * AG-UI — contrato de eventos agente→UI (Fase 7, Agent Console). Tipos + encoder SSE + + * validação de sequência. Puro (sem I/O). Permite replay/reconexão determinísticos: a UI + * reconstrói o estado a partir do fluxo de eventos ordenado. + * + * Segue o espírito do protocolo AG-UI (ciclo de run, deltas de texto, tool calls, estado). + */ + +export type AgUiEventType = + | "RUN_STARTED" + | "TEXT_MESSAGE_START" + | "TEXT_MESSAGE_CONTENT" + | "TEXT_MESSAGE_END" + | "TOOL_CALL_START" + | "TOOL_CALL_ARGS" + | "TOOL_CALL_END" + | "STATE_SNAPSHOT" + | "STATE_DELTA" + | "RUN_FINISHED" + | "RUN_ERROR"; + +interface Base { + readonly seq: number; // ordem monotônica (para replay/reconexão) + readonly runId: string; +} + +export type AgUiEvent = + | (Base & { type: "RUN_STARTED" }) + | (Base & { type: "TEXT_MESSAGE_START"; messageId: string; role: "assistant" | "tool" }) + | (Base & { type: "TEXT_MESSAGE_CONTENT"; messageId: string; delta: string }) + | (Base & { type: "TEXT_MESSAGE_END"; messageId: string }) + | (Base & { type: "TOOL_CALL_START"; toolCallId: string; name: string }) + | (Base & { type: "TOOL_CALL_ARGS"; toolCallId: string; delta: string }) + | (Base & { type: "TOOL_CALL_END"; toolCallId: string }) + | (Base & { type: "STATE_SNAPSHOT"; state: Record }) + | (Base & { type: "STATE_DELTA"; patch: Record }) + | (Base & { type: "RUN_FINISHED" }) + | (Base & { type: "RUN_ERROR"; message: string }); + +export const TERMINAL_EVENTS: ReadonlySet = new Set(["RUN_FINISHED", "RUN_ERROR"]); + +/** Codifica um evento como frame SSE (event: \ndata: \n\n). */ +export function encodeSse(event: AgUiEvent): string { + return `event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`; +} + +export interface SequenceValidation { + readonly ok: boolean; + readonly errors: string[]; +} + +/** + * Valida invariantes do fluxo (para catch de bugs de emissão, não de segurança): + * - `seq` estritamente crescente; + * - primeiro evento = RUN_STARTED e RUN_STARTED aparece SÓ no índice 0 (nunca reinicia no meio); + * - todos os eventos do MESMO runId (fluxo de um run só); + * - exatamente um terminal (RUN_FINISHED|RUN_ERROR), e ele é o último; + * - nenhum evento após o terminal. + */ +export function validateEventSequence(events: ReadonlyArray): SequenceValidation { + const errors: string[] = []; + if (events.length === 0) return { ok: false, errors: ["fluxo vazio"] }; + + if (events[0].type !== "RUN_STARTED") errors.push("primeiro evento deve ser RUN_STARTED"); + const runId = events[0].runId; + + let lastSeq = -Infinity; + let terminalAt = -1; + events.forEach((e, i) => { + if (e.seq <= lastSeq) errors.push(`seq não crescente em ${i} (${e.seq})`); + lastSeq = e.seq; + if (e.type === "RUN_STARTED" && i !== 0) errors.push(`RUN_STARTED duplicado no índice ${i}`); + if (e.runId !== runId) errors.push(`runId inconsistente no índice ${i} (${e.runId})`); + if (TERMINAL_EVENTS.has(e.type)) { + if (terminalAt !== -1) errors.push(`múltiplos eventos terminais (índice ${i})`); + terminalAt = i; + } + }); + + if (terminalAt === -1) errors.push("faltou evento terminal (RUN_FINISHED|RUN_ERROR)"); + else if (terminalAt !== events.length - 1) errors.push("há eventos após o terminal"); + + return { ok: errors.length === 0, errors }; +} + +/** Fábrica incremental de seq — ajuda emissores a produzir fluxos válidos. */ +export function createSeq(start = 0): () => number { + let n = start; + return () => n++; +} diff --git a/open-sse/browser-guard/index.ts b/open-sse/browser-guard/index.ts new file mode 100644 index 00000000000..a3a8a55611c --- /dev/null +++ b/open-sse/browser-guard/index.ts @@ -0,0 +1,108 @@ +/** + * Browser Guard — política DETERMINÍSTICA para automação de navegador (Fase 7, "Browser Use"). + * + * Desligado por padrão. Allowlist de domínio. Efeito externo (enviar form, download, upload, + * compra) exige aprovação humana. E — crucial contra prompt injection — uma ação cuja ORIGEM é a + * PÁGINA (texto lido) NUNCA escala permissão: efeitos externos originados na página são negados. + * O CÓDIGO decide (não a IA). Puro, sem I/O; o driver real (Playwright) fica fora deste módulo. + */ + +export type BrowserActionKind = + | "navigate" + | "read" + | "click" + | "type" + | "submit" // enviar formulário + | "download" + | "upload" + | "purchase"; + +/** Origem da intenção: o usuário/agente (confiável) ou conteúdo lido da página (NÃO confiável). */ +export type BrowserActionOrigin = "user" | "page"; + +export interface BrowserAction { + readonly kind: BrowserActionKind; + /** URL alvo (quando aplicável). Usada para checar a allowlist de domínio. */ + readonly url?: string; + readonly origin: BrowserActionOrigin; +} + +export interface BrowserPolicy { + /** Browser Use está ligado? Padrão do produto: OFF. */ + readonly enabled: boolean; + /** Domínios permitidos (host exato ou sufixo, ex.: "example.com"). Vazio = nada permitido. */ + readonly allowedDomains: ReadonlyArray; +} + +export type BrowserDecision = "allow" | "require_approval" | "deny"; + +export interface BrowserVerdict { + readonly decision: BrowserDecision; + readonly reason: string; +} + +/** Ações com efeito externo — nunca automáticas: exigem humano (e nunca se originam da página). */ +export const EXTERNAL_EFFECT_KINDS: ReadonlySet = new Set([ + "submit", + "download", + "upload", + "purchase", +]); + +function hostOf(url: string): string | null { + try { + return new URL(url).hostname.toLowerCase(); + } catch { + return null; + } +} + +/** host casa a allowlist se for igual ou subdomínio de algum domínio permitido. */ +function domainAllowed(host: string, allowed: ReadonlyArray): boolean { + return allowed.some((d) => { + const dom = d.toLowerCase().replace(/^\.+/, ""); + return host === dom || host.endsWith("." + dom); + }); +} + +/** + * Decide uma ação de navegador. Fail-closed: + * - browser OFF → deny; + * - URL fora da allowlist → deny; + * - efeito externo originado na PÁGINA → deny (prompt injection não escala); + * - efeito externo (origem usuário) → require_approval (humano); + * - navegação/leitura em domínio permitido → allow. + */ +export function decideBrowserAction(action: BrowserAction, policy: BrowserPolicy): BrowserVerdict { + if (!policy.enabled) { + return { decision: "deny", reason: "Browser Use está desligado (padrão)" }; + } + + if (action.url !== undefined) { + const host = hostOf(action.url); + if (!host) return { decision: "deny", reason: "URL inválida" }; + if (!domainAllowed(host, policy.allowedDomains)) { + return { decision: "deny", reason: `domínio fora da allowlist: ${host}` }; + } + } + + const isExternal = EXTERNAL_EFFECT_KINDS.has(action.kind); + + if (isExternal && action.origin === "page") { + // Instrução vinda do conteúdo da página tentando um efeito externo → bloqueio absoluto. + return { + decision: "deny", + reason: "efeito externo originado na página (possível prompt injection) — negado", + }; + } + + if (isExternal) { + return { + decision: "require_approval", + reason: `efeito externo (${action.kind}) requer aprovação humana`, + }; + } + + // Leitura/navegação/clique/digitação em domínio permitido. + return { decision: "allow", reason: "ação de leitura/navegação em domínio permitido" }; +} diff --git a/open-sse/buzz-bridge/adapter.ts b/open-sse/buzz-bridge/adapter.ts new file mode 100644 index 00000000000..f30b0f56238 --- /dev/null +++ b/open-sse/buzz-bridge/adapter.ts @@ -0,0 +1,55 @@ +/** + * Buzz Bridge — adaptador desabilitado + guarda de autorização. + * + * Enquanto a flag `buzz_hub` estiver OFF (ou o relay não estiver rodando), usamos o + * DisabledBuzzAdapter: ele NÃO conecta e NÃO publica — só reporta que está inerte. As + * entradas ficam no outbox até haver relay + flag ON e o WebSocketBuzzAdapter (futuro). + */ +import type { + BuzzAdapter, + BuzzEvent, + BuzzIdentityMapping, + BuzzSubscriptionFilter, + OutboxEntry, +} from "./types.ts"; + +export class DisabledBuzzAdapter implements BuzzAdapter { + readonly enabled = false; + async connect(): Promise { + /* inerte: sem relay, sem conexão */ + } + async publish(_entry: OutboxEntry): Promise { + return false; // nada é publicado enquanto desabilitado + } + async subscribe( + _filter: BuzzSubscriptionFilter, + _onEvent: (e: BuzzEvent) => void + ): Promise { + /* inerte */ + } + async close(): Promise { + /* inerte */ + } +} + +/** + * Regra de segurança inegociável: uma chave Nostr (buzz_pubkey) NUNCA autoriza, por si só, + * uma ação no OmniRoute. A autorização real vem SEMPRE das políticas/aprovações do OmniRoute + * para o (tenant, workspace, user/agent) mapeado — nunca do fato de o evento estar assinado. + * + * Esta função é deliberadamente fail-closed: ela apenas confirma que existe um mapeamento + * de identidade; a decisão de permitir o efeito é do Policy Engine, fora daqui. + */ +export function buzzIdentityIsMapped( + mapping: BuzzIdentityMapping | undefined, + buzzPubkey: string +): boolean { + if (!mapping) return false; + if (!mapping.buzzPubkey || mapping.buzzPubkey !== buzzPubkey) return false; + return Boolean(mapping.tenantId && mapping.workspaceId); +} + +/** Uma chave Nostr, sozinha, jamais autoriza. Documenta a regra em código executável. */ +export function nostrKeyAuthorizes(): false { + return false; +} diff --git a/open-sse/buzz-bridge/index.ts b/open-sse/buzz-bridge/index.ts new file mode 100644 index 00000000000..9886bb9c7d0 --- /dev/null +++ b/open-sse/buzz-bridge/index.ts @@ -0,0 +1,36 @@ +/** + * Buzz Bridge — API pública (ponte OmniRoute ↔ Block Buzz), atrás da flag `buzz_hub` (OFF). + * + * Buzz é serviço SEPARADO (relay Nostr). Este módulo é a ponte tipada + idempotente; nada + * conecta enquanto desabilitado. Configuração e controle ficam no PAINEL ÚNICO do OmniRoute. + */ +export * from "./types.ts"; +export { Outbox, Inbox } from "./outbox.ts"; +export { DisabledBuzzAdapter, buzzIdentityIsMapped, nostrKeyAuthorizes } from "./adapter.ts"; +export { + finalizeEvent, + verifyEvent, + getPublicKey, + generateSecretKey, + type SignedNostrEvent, + type UnsignedNostrEvent, +} from "./nostr.ts"; +export { WebSocketBuzzAdapter, type WebSocketBuzzConfig } from "./wsAdapter.ts"; + +import { DisabledBuzzAdapter } from "./adapter.ts"; +import type { BuzzAdapter } from "./types.ts"; +import { WebSocketBuzzAdapter, type WebSocketBuzzConfig } from "./wsAdapter.ts"; + +/** + * Resolve o adaptador: inerte quando `buzz_hub` OFF ou sem config; WebSocketBuzzAdapter real + * (contra o buzz-relay) quando a flag está ON e há config (relayUrl + secretKey). + */ +export function resolveBuzzAdapter( + flagEnabled: boolean, + config?: WebSocketBuzzConfig +): BuzzAdapter { + if (!flagEnabled || !config?.relayUrl || !config?.secretKeyHex) { + return new DisabledBuzzAdapter(); + } + return new WebSocketBuzzAdapter(config); +} diff --git a/open-sse/buzz-bridge/nostr.ts b/open-sse/buzz-bridge/nostr.ts new file mode 100644 index 00000000000..9e61cbaee9b --- /dev/null +++ b/open-sse/buzz-bridge/nostr.ts @@ -0,0 +1,64 @@ +/** + * Buzz Bridge — assinatura/verificação de eventos Nostr (NIP-01), pura. + * + * Usa @noble/curves (schnorr secp256k1) + @noble/hashes (sha256). O id do evento é o + * sha256 da serialização canônica [0,pubkey,created_at,kind,tags,content]; a assinatura é + * schnorr sobre esse hash. Base para o WebSocketBuzzAdapter falar com o buzz-relay. + */ +import { schnorr } from "@noble/curves/secp256k1.js"; +import { sha256 } from "@noble/hashes/sha2.js"; +import { bytesToHex, hexToBytes } from "@noble/hashes/utils.js"; + +export interface UnsignedNostrEvent { + pubkey: string; + created_at: number; + kind: number; + tags: string[][]; + content: string; +} + +export interface SignedNostrEvent extends UnsignedNostrEvent { + id: string; + sig: string; +} + +/** Gera uma secret key Nostr (32 bytes) em hex. */ +export function generateSecretKey(): string { + return bytesToHex(schnorr.utils.randomSecretKey()); +} + +/** Deriva a pubkey (x-only, 32 bytes) em hex a partir da secret key hex. */ +export function getPublicKey(secretKeyHex: string): string { + return bytesToHex(schnorr.getPublicKey(hexToBytes(secretKeyHex))); +} + +function serialize(e: UnsignedNostrEvent): Uint8Array { + return new TextEncoder().encode( + JSON.stringify([0, e.pubkey, e.created_at, e.kind, e.tags, e.content]) + ); +} + +/** Calcula id + assina, produzindo um evento Nostr completo. */ +export function finalizeEvent( + unsigned: Omit, + secretKeyHex: string +): SignedNostrEvent { + const sk = hexToBytes(secretKeyHex); + const pubkey = bytesToHex(schnorr.getPublicKey(sk)); + const base: UnsignedNostrEvent = { ...unsigned, pubkey }; + const idBytes = sha256(serialize(base)); + const id = bytesToHex(idBytes); + const sig = bytesToHex(schnorr.sign(idBytes, sk)); + return { ...base, id, sig }; +} + +/** Verifica id + assinatura schnorr de um evento. Fail-closed em qualquer erro. */ +export function verifyEvent(e: SignedNostrEvent): boolean { + try { + const idBytes = sha256(serialize(e)); + if (bytesToHex(idBytes) !== e.id) return false; + return schnorr.verify(hexToBytes(e.sig), idBytes, hexToBytes(e.pubkey)); + } catch { + return false; + } +} diff --git a/open-sse/buzz-bridge/outbox.ts b/open-sse/buzz-bridge/outbox.ts new file mode 100644 index 00000000000..0ecec4cb2f2 --- /dev/null +++ b/open-sse/buzz-bridge/outbox.ts @@ -0,0 +1,83 @@ +/** + * Buzz Bridge — outbox/inbox idempotente (puro). + * + * Garante que nenhum evento seja publicado ou processado duas vezes, mesmo com retries. + * Dedup por `event.id` (saída) e por `event.id` (entrada); ordenação por `sequenceNumber`. + * Sem dupla fonte de verdade: aqui só orquestramos entrega; o estado durável é do OmniRoute. + */ +import type { BuzzEvent, InboxEntry, OutboxEntry } from "./types.ts"; + +/** Fila outbox in-memory com dedup determinístico (a persistência real é no DB do OmniRoute). */ +export class Outbox { + private readonly byEventId = new Map(); + private seq = 0; + + /** Enfileira um evento para publicação. Idempotente: reenfileirar o mesmo id é no-op. */ + enqueue(params: { + event: BuzzEvent; + correlationId: string; + taskId?: string; + runId?: string; + }): OutboxEntry { + const existing = this.byEventId.get(params.event.id); + if (existing) return existing; + const entry: OutboxEntry = { + id: params.event.id, + correlationId: params.correlationId, + sequenceNumber: ++this.seq, + taskId: params.taskId, + runId: params.runId, + event: params.event, + status: "pending", + attempts: 0, + }; + this.byEventId.set(entry.id, entry); + return entry; + } + + /** Entradas pendentes, em ordem de sequência (entrega ordenada). */ + pending(): OutboxEntry[] { + return [...this.byEventId.values()] + .filter((e) => e.status === "pending") + .sort((a, b) => a.sequenceNumber - b.sequenceNumber); + } + + markPublished(id: string): void { + const e = this.byEventId.get(id); + if (e) e.status = "published"; + } + + markFailed(id: string): void { + const e = this.byEventId.get(id); + if (e) { + e.status = "failed"; + e.attempts += 1; + } + } + + /** Reprograma entradas falhas para nova tentativa (backoff é decidido por quem chama). */ + requeueFailed(): void { + for (const e of this.byEventId.values()) if (e.status === "failed") e.status = "pending"; + } + + size(): number { + return this.byEventId.size; + } +} + +/** Lado de entrada: dedup por id de evento; processa cada evento no máximo uma vez. */ +export class Inbox { + private readonly seen = new Set(); + private seq = 0; + + /** Registra um evento recebido. Retorna null se já visto (dedup), senão a entrada. */ + receive(event: BuzzEvent, correlationId: string): InboxEntry | null { + if (this.seen.has(event.id)) return null; + this.seen.add(event.id); + return { event, correlationId, sequenceNumber: ++this.seq, status: "received" }; + } + + seenCount(): number { + return this.seen.size; + } +} diff --git a/open-sse/buzz-bridge/types.ts b/open-sse/buzz-bridge/types.ts new file mode 100644 index 00000000000..8e6fa28c558 --- /dev/null +++ b/open-sse/buzz-bridge/types.ts @@ -0,0 +1,80 @@ +/** + * Buzz Bridge — tipos e contratos (ponte OmniRoute ↔ Block Buzz). + * + * Buzz (github.com/block/buzz, Apache-2.0, SHA 3c7f288…) é um serviço SEPARADO: um relay + * Nostr (NIP-01) onde cada ação é um evento assinado identificado por `kind`. O relay é a + * fonte única de verdade APENAS de canais, conversas, membros e eventos colaborativos. + * + * O OmniRoute permanece plano de controle e dono das políticas. Regras não-negociáveis: + * - Uma chave Nostr (`buzz_pubkey`) NÃO autoriza nenhuma ação no OmniRoute. + * - Toda ponte é idempotente (outbox/inbox, sequence_number, correlation_id, dedup). + * - Não criar dupla fonte de verdade: estado de tarefas/runs/aprovações vive no OmniRoute. + * + * Este módulo fica atrás da feature flag `buzz_hub` (OFF por padrão). Sem o relay rodando, + * o adaptador é inerte (não conecta); as entradas ficam no outbox até haver relay + flag ON. + */ + +/** Evento colaborativo (shape derivado de NIP-01: id, pubkey, kind, tags, content, sig). */ +export interface BuzzEvent { + /** id do evento (hash) — chave de deduplicação. */ + readonly id: string; + /** pubkey Nostr do autor (humano ou agente). NUNCA é autorização no OmniRoute. */ + readonly pubkey: string; + /** kind NIP-01 (inteiro) — nova feature = novo kind. */ + readonly kind: number; + readonly createdAt: number; // epoch seconds + readonly tags: ReadonlyArray>; + readonly content: string; + readonly sig?: string; +} + +/** Mapeamento de identidade — o OmniRoute é dono; buzz_pubkey é só um atributo. */ +export interface BuzzIdentityMapping { + readonly tenantId: string; + readonly workspaceId: string; + readonly userId?: string; + readonly agentId?: string; + readonly buzzPubkey: string; +} + +/** Entrada de saída (OmniRoute → Buzz), aguardando publicação idempotente no relay. */ +export interface OutboxEntry { + readonly id: string; // id lógico local (dedup) + readonly correlationId: string; + readonly sequenceNumber: number; + readonly taskId?: string; + readonly runId?: string; + readonly event: BuzzEvent; + status: "pending" | "published" | "failed"; + attempts: number; +} + +/** Entrada de entrada (Buzz → OmniRoute), aguardando processamento idempotente. */ +export interface InboxEntry { + readonly event: BuzzEvent; + readonly correlationId: string; + readonly sequenceNumber: number; + status: "received" | "processed" | "ignored"; +} + +/** + * Contrato do adaptador tipado. NÃO há chamadas ao relay espalhadas pelo código — + * tudo passa por esta interface. Implementações: DisabledBuzzAdapter (flag OFF) e, + * futuramente, WebSocketBuzzAdapter (quando o relay rodar + flag ON). + */ +export interface BuzzAdapter { + readonly enabled: boolean; + /** Conecta ao relay (WebSocket). No-op quando desabilitado. */ + connect(): Promise; + /** Publica um evento já enfileirado no outbox. Retorna true se publicado. */ + publish(entry: OutboxEntry): Promise; + /** Assina um filtro (channel/kind) e entrega eventos ao callback. */ + subscribe(filter: BuzzSubscriptionFilter, onEvent: (e: BuzzEvent) => void): Promise; + close(): Promise; +} + +export interface BuzzSubscriptionFilter { + readonly kinds?: ReadonlyArray; + readonly channelId?: string; + readonly since?: number; +} diff --git a/open-sse/buzz-bridge/wsAdapter.ts b/open-sse/buzz-bridge/wsAdapter.ts new file mode 100644 index 00000000000..4fe868816d3 --- /dev/null +++ b/open-sse/buzz-bridge/wsAdapter.ts @@ -0,0 +1,276 @@ +/** + * Buzz Bridge — WebSocketBuzzAdapter: fala com o buzz-relay (Nostr) via WebSocket. + * + * Implementa o contrato BuzzAdapter de verdade: conecta, faz auth NIP-42 (o relay do Buzz + * tem auth_required), publica eventos assinados (kind/tags/content) e assina filtros. + * O OmniRoute continua o plano de controle — este adaptador só transporta eventos. + */ +import WebSocket from "ws"; + +import { finalizeEvent, getPublicKey, verifyEvent, type SignedNostrEvent } from "./nostr.ts"; +import type { BuzzAdapter, BuzzEvent, BuzzSubscriptionFilter, OutboxEntry } from "./types.ts"; + +export interface WebSocketBuzzConfig { + /** URL do relay, ex.: ws://127.0.0.1:3000 */ + readonly relayUrl: string; + /** Secret key Nostr (hex) da identidade deste agente. */ + readonly secretKeyHex: string; + /** Timeout de operações (ms). */ + readonly timeoutMs?: number; + /** Janela (ms) sem AUTH para assumir relay sem auth_required (default 400). */ + readonly authGraceMs?: number; + /** Reconectar automaticamente após queda inesperada (default true). */ + readonly autoReconnect?: boolean; + /** Backoff base da reconexão (ms, default 500). */ + readonly reconnectBaseMs?: number; + /** Teto do backoff da reconexão (ms, default 15000). */ + readonly reconnectMaxMs?: number; +} + +function toBuzzEvent(e: SignedNostrEvent): BuzzEvent { + return { + id: e.id, + pubkey: e.pubkey, + kind: e.kind, + createdAt: e.created_at, + tags: e.tags, + content: e.content, + sig: e.sig, + }; +} + +export class WebSocketBuzzAdapter implements BuzzAdapter { + readonly enabled = true; + readonly pubkey: string; + private ws?: WebSocket; + private readonly timeout: number; + private readonly pendingOk = new Map void>(); + private readonly subs = new Map< + string, + { filter: BuzzSubscriptionFilter; onEvent: (e: BuzzEvent) => void } + >(); + /** Janela (ms) sem AUTH para assumir relay SEM auth_required e liberar o publish. */ + private readonly authGraceMs: number; + /** Resolvida quando o AUTH (NIP-42) foi confirmado OU o relay não exige auth. publish() espera nela. */ + private authReady: Promise = Promise.resolve(); + private authReadyResolve?: () => void; + private authGraceTimer?: ReturnType; + private authChallengeSeen = false; + // Reconexão automática (resiliência do consumidor de entrada após uma queda). + private readonly autoReconnect: boolean; + private readonly reconnectBaseMs: number; + private readonly reconnectMaxMs: number; + private closedByUser = false; + private reconnectAttempts = 0; + private reconnectTimer?: ReturnType; + private connectInFlight?: Promise; + + constructor(private readonly config: WebSocketBuzzConfig) { + this.pubkey = getPublicKey(config.secretKeyHex); + this.timeout = config.timeoutMs ?? 8000; + this.authGraceMs = config.authGraceMs ?? 400; + this.autoReconnect = config.autoReconnect ?? true; + this.reconnectBaseMs = config.reconnectBaseMs ?? 500; + this.reconnectMaxMs = config.reconnectMaxMs ?? 15000; + } + + connect(): Promise { + // Single-flight: uma reconexão automática e um ensureConnected concorrentes compartilham o + // MESMO connect em andamento (nunca abrimos dois sockets). + if (this.connectInFlight) return this.connectInFlight; + this.closedByUser = false; + this.connectInFlight = new Promise((resolve, reject) => { + const ws = new WebSocket(this.config.relayUrl); + this.ws = ws; + this.authChallengeSeen = false; + // authReady só resolve quando o relay confirma o AUTH (OK=true do nosso kind 22242) OU quando + // a janela de graça passa sem nenhum challenge (relay sem auth_required). publish() aguarda nela + // para NUNCA correr contra o handshake — o que antes fazia o 1º publish falhar e (sem o requeue + // no DB) estrangular a mensagem para sempre. + this.authReady = new Promise((res) => { + this.authReadyResolve = res; + }); + const openTimer = setTimeout(() => { + this.connectInFlight = undefined; + reject(new Error("buzz relay connect timeout")); + }, this.timeout); + ws.on("open", () => { + clearTimeout(openTimer); + this.reconnectAttempts = 0; // conexão saudável zera o backoff + this.connectInFlight = undefined; + this.authGraceTimer = setTimeout(() => { + if (!this.authChallengeSeen) this.markAuthReady(); + }, this.authGraceMs); + // Re-emite os REQ das assinaturas ativas — o consumidor sobrevive a uma reconexão. + for (const [subId, sub] of this.subs) this.sendReq(subId, sub.filter); + resolve(); + }); + ws.on("error", (e: Error) => { + clearTimeout(openTimer); + this.connectInFlight = undefined; + reject(e); + }); + ws.on("close", () => this.onSocketClose(ws)); + ws.on("message", (data: WebSocket.RawData) => this.onMessage(data.toString())); + }); + return this.connectInFlight; + } + + /** + * Queda do socket: libera waiters e agenda reconexão (salvo close() explícito). Ignora eventos + * de um socket OBSOLETO — se já reconectamos, o `close` tardio do socket antigo não pode zerar o + * novo (senão o publish enviaria para um socket órfão e expiraria). + */ + private onSocketClose(closed: WebSocket): void { + if (this.ws !== closed) return; // close tardio de um socket já substituído + this.ws = undefined; + this.markAuthReady(); // não deixa um publish preso esperando authReady de um socket morto + if (this.closedByUser || !this.autoReconnect) return; + this.scheduleReconnect(); + } + + /** Reconexão com backoff exponencial + jitter, com teto. Um timer por vez. */ + private scheduleReconnect(): void { + if (this.reconnectTimer || this.closedByUser) return; + const backoff = Math.min( + this.reconnectBaseMs * 2 ** this.reconnectAttempts, + this.reconnectMaxMs + ); + const delay = backoff + Math.floor(Math.random() * this.reconnectBaseMs); + this.reconnectAttempts++; + this.reconnectTimer = setTimeout(() => { + this.reconnectTimer = undefined; + this.connect().catch(() => this.scheduleReconnect()); + }, delay); + } + + /** Libera o publish (AUTH confirmado ou dispensado). Idempotente. */ + private markAuthReady(): void { + if (this.authGraceTimer) { + clearTimeout(this.authGraceTimer); + this.authGraceTimer = undefined; + } + this.authReadyResolve?.(); + this.authReadyResolve = undefined; + } + + /** Reconecta sob demanda se o socket não estiver aberto (ex.: queda entre flushes). */ + private async ensureConnected(): Promise { + if (this.ws && this.ws.readyState === WebSocket.OPEN) return; + await this.connect(); + } + + private send(msg: unknown): void { + this.ws?.send(JSON.stringify(msg)); + } + + private onMessage(raw: string): void { + let msg: unknown; + try { + msg = JSON.parse(raw); + } catch { + return; + } + if (!Array.isArray(msg) || msg.length === 0) return; + const type = msg[0]; + + if (type === "AUTH" && typeof msg[1] === "string") { + // NIP-42: relay exige auth. Cancela a graça (não vamos "assumir sem auth") e só libera o + // publish quando o relay confirmar nosso evento de auth (OK=true do id abaixo). + this.authChallengeSeen = true; + if (this.authGraceTimer) { + clearTimeout(this.authGraceTimer); + this.authGraceTimer = undefined; + } + const authEvent = finalizeEvent( + { + created_at: Math.floor(Date.now() / 1000), + kind: 22242, + tags: [ + ["relay", this.config.relayUrl], + ["challenge", msg[1]], + ], + content: "", + }, + this.config.secretKeyHex + ); + // Quando o OK do nosso auth chegar com sucesso, libera authReady. Se vier false, não libera: + // os publishes vão expirar → 'failed' → requeue no próximo flush (backstop durável). + this.pendingOk.set(authEvent.id, (ok) => { + if (ok) this.markAuthReady(); + }); + this.send(["AUTH", authEvent]); + return; + } + if (type === "OK" && typeof msg[1] === "string") { + const resolver = this.pendingOk.get(msg[1]); + if (resolver) { + resolver(Boolean(msg[2])); + this.pendingOk.delete(msg[1]); + } + return; + } + if (type === "EVENT" && typeof msg[1] === "string") { + const sub = this.subs.get(msg[1]); + const ev = msg[2] as SignedNostrEvent | undefined; + if (sub && ev && verifyEvent(ev)) sub.onEvent(toBuzzEvent(ev)); + } + } + + private sendReq(subId: string, filter: BuzzSubscriptionFilter): void { + const f: Record = {}; + if (filter.kinds) f.kinds = filter.kinds; + if (filter.since) f.since = filter.since; + if (filter.channelId) f["#e"] = [filter.channelId]; + this.send(["REQ", subId, f]); + } + + /** + * Publica o evento do outbox, RE-ASSINANDO com a chave deste agente (o pubkey do produtor é + * ignorado por design — "uma chave Nostr nunca autoriza ação"; a autoria de saída é sempre do + * agente OmniRoute). Usa o createdAt PERSISTIDO no enqueue (nunca o relógio) para que re-tentativas + * produzam o mesmo id → dedup real do relay. Retorna true se o relay aceitou. + */ + async publish(entry: OutboxEntry): Promise { + await this.ensureConnected(); // reconecta se o socket caiu entre flushes + await this.authReady; // não corre contra o handshake NIP-42 + const signed = finalizeEvent( + { + created_at: entry.event.createdAt, + kind: entry.event.kind, + tags: entry.event.tags.map((t) => [...t]), + content: entry.event.content, + }, + this.config.secretKeyHex + ); + return new Promise((resolve) => { + const timer = setTimeout(() => { + this.pendingOk.delete(signed.id); + resolve(false); + }, this.timeout); + this.pendingOk.set(signed.id, (ok) => { + clearTimeout(timer); + resolve(ok); + }); + this.send(["EVENT", signed]); + }); + } + + async subscribe(filter: BuzzSubscriptionFilter, onEvent: (e: BuzzEvent) => void): Promise { + await this.ensureConnected(); // reconecta se necessário antes de assinar + const subId = "sub-" + Math.random().toString(36).slice(2, 10); + this.subs.set(subId, { filter, onEvent }); + this.sendReq(subId, filter); + } + + async close(): Promise { + // Fecho explícito: desliga a reconexão automática, cancela timers e libera waiters. + this.closedByUser = true; + if (this.reconnectTimer) { + clearTimeout(this.reconnectTimer); + this.reconnectTimer = undefined; + } + this.markAuthReady(); + this.ws?.close(); + } +} diff --git a/open-sse/loop-engine/budget.ts b/open-sse/loop-engine/budget.ts new file mode 100644 index 00000000000..8d48d4d7e15 --- /dev/null +++ b/open-sse/loop-engine/budget.ts @@ -0,0 +1,48 @@ +/** + * Loop Engine — controle de orçamento (puro). + * + * "Toda task tem budget (tokens, tempo, tentativas). Estourou → aborta e reporta." + * Este módulo NUNCA executa efeito; apenas calcula se o run pode continuar. + */ +import type { LoopBudget, LoopBudgetUsage } from "./types.ts"; + +export function emptyUsage(): LoopBudgetUsage { + return { tokens: 0, wallClockMs: 0, attempts: 0 }; +} + +/** Soma consumo ao uso acumulado, sem mutar o original. */ +export function addUsage( + current: LoopBudgetUsage, + delta: Partial +): LoopBudgetUsage { + return { + tokens: current.tokens + Math.max(0, delta.tokens ?? 0), + wallClockMs: current.wallClockMs + Math.max(0, delta.wallClockMs ?? 0), + attempts: current.attempts + Math.max(0, delta.attempts ?? 0), + }; +} + +export interface BudgetCheck { + readonly ok: boolean; + /** Qual limite estourou primeiro (para relatório honesto). */ + readonly exceeded: Array<"tokens" | "wallClockMs" | "attempts">; +} + +/** true = ainda dentro do orçamento em TODAS as dimensões. */ +export function checkBudget(budget: LoopBudget, usage: LoopBudgetUsage): BudgetCheck { + const exceeded: BudgetCheck["exceeded"] = []; + if (usage.tokens > budget.maxTokens) exceeded.push("tokens"); + if (usage.wallClockMs > budget.maxWallClockMs) exceeded.push("wallClockMs"); + if (usage.attempts > budget.maxAttempts) exceeded.push("attempts"); + return { ok: exceeded.length === 0, exceeded }; +} + +/** Fração consumida (0..1+) da dimensão mais próxima do limite — para alertas. */ +export function budgetPressure(budget: LoopBudget, usage: LoopBudgetUsage): number { + const fractions = [ + budget.maxTokens > 0 ? usage.tokens / budget.maxTokens : 0, + budget.maxWallClockMs > 0 ? usage.wallClockMs / budget.maxWallClockMs : 0, + budget.maxAttempts > 0 ? usage.attempts / budget.maxAttempts : 0, + ]; + return Math.max(0, ...fractions); +} diff --git a/open-sse/loop-engine/index.ts b/open-sse/loop-engine/index.ts new file mode 100644 index 00000000000..1bd482a9d00 --- /dev/null +++ b/open-sse/loop-engine/index.ts @@ -0,0 +1,60 @@ +/** + * Loop Engine — API pública do núcleo (report-only). + * + * Módulo LEVE do OmniRoute (não é serviço). O estado durável (runs/steps/checkpoints/ + * aprovações) vive no DB do OmniRoute (migração 174); este núcleo é a lógica pura de + * ciclo, orçamento e política. Fica atrás da feature flag `loop_engine` (OFF por padrão). + * + * Configuração e controle ficam no PAINEL ÚNICO do OmniRoute (dashboard), não numa UI à parte. + */ +import { randomUUID } from "node:crypto"; + +import type { LoopBudget, LoopRun, LoopStep } from "./types.ts"; + +export * from "./types.ts"; +export { addUsage, checkBudget, budgetPressure, emptyUsage } from "./budget.ts"; +export { decideEffect, isAutoExecutable, type PolicyContext } from "./policyGate.ts"; +export { advance, type AdvanceInput, type AdvanceResult } from "./stateMachine.ts"; + +export const DEFAULT_LOOP_BUDGET: LoopBudget = { + maxTokens: 200_000, + maxWallClockMs: 15 * 60_000, // 15 min + maxAttempts: 3, +}; + +/** Cria um run novo em modo report-only, na primeira fase (discover). */ +export function createLoopRun(params: { + pattern: string; + budget?: Partial; + correlationId?: string; + taskId?: string; +}): LoopRun { + return { + id: randomUUID(), + pattern: params.pattern, + phase: "discover", + status: "report_only", + budget: { ...DEFAULT_LOOP_BUDGET, ...params.budget }, + usage: { tokens: 0, wallClockMs: 0, attempts: 0 }, + steps: [], + correlationId: params.correlationId ?? randomUUID(), + taskId: params.taskId, + sequenceNumber: 0, + }; +} + +/** Anexa uma etapa proposta ao run (sem executar nada). */ +export function proposeStep( + run: LoopRun, + step: Omit +): LoopStep { + const created: LoopStep = { + id: randomUUID(), + runId: run.id, + index: run.steps.length, + status: "proposed", + ...step, + }; + run.steps.push(created); + return created; +} diff --git a/open-sse/loop-engine/policyGate.ts b/open-sse/loop-engine/policyGate.ts new file mode 100644 index 00000000000..3af7fecfcf9 --- /dev/null +++ b/open-sse/loop-engine/policyGate.ts @@ -0,0 +1,84 @@ +/** + * Loop Engine — Policy Gate (determinístico; código clássico decide, não a IA). + * + * Regra central do plano: "Loop Engine começa em modo report-only. Ele PROPÕE; o Policy + * Engine e as aprovações do OmniRoute DECIDEM. Nenhum loop envia mensagem, publica, + * exclui, faz merge, deploy, compras ou muda conta sem a aprovação exigida." + * + * Este gate NUNCA executa o efeito — só classifica. A execução (quando aprovada) é feita + * pelos conectores do OmniRoute, fora deste módulo. + */ +import type { LoopProposedEffect, PolicyDecision } from "./types.ts"; + +/** Efeitos que são SEMPRE proibidos de forma autônoma (exigem humano), mesmo fora de report-only. */ +const ALWAYS_REQUIRE_APPROVAL = new Set([ + "message_send", + "git_pr", + "deploy", + "purchase", + "account_change", + "delete", +]); + +/** Efeitos irreversíveis/destrutivos que, por padrão, são NEGADos ao loop (só humano faz). */ +const AUTONOMY_DENIED = new Set([ + "purchase", + "account_change", + "delete", +]); + +export interface PolicyContext { + /** Modo do run. Em report-only nada externo é auto-aprovado. */ + readonly reportOnly: boolean; + /** Allowlist opcional de kinds que o operador liberou para auto-execução (nunca inclui os AUTONOMY_DENIED). */ + readonly autoApproveKinds?: ReadonlyArray; +} + +/** + * Decide o destino de um efeito proposto. Determinístico e fail-closed: + * - efeito "none" → allow (inócuo); + * - kinds destrutivos → deny (autonomia proibida); + * - report-only OU kind sensível → require_approval; + * - só libera "allow" se o operador incluiu explicitamente o kind na allowlist e não é destrutivo. + */ +export function decideEffect( + effect: LoopProposedEffect | undefined, + ctx: PolicyContext +): PolicyDecision { + if (!effect || effect.kind === "none") return { outcome: "allow" }; + + if (AUTONOMY_DENIED.has(effect.kind)) { + return { + outcome: "deny", + reason: `Efeito '${effect.kind}' nunca é autônomo: exige ação humana explícita.`, + }; + } + + if (ctx.reportOnly) { + return { + outcome: "require_approval", + reason: "Run em modo report-only: todo efeito externo aguarda aprovação.", + }; + } + + if (ALWAYS_REQUIRE_APPROVAL.has(effect.kind)) { + const allowed = ctx.autoApproveKinds?.includes(effect.kind) ?? false; + if (!allowed) { + return { + outcome: "require_approval", + reason: `Efeito '${effect.kind}' exige aprovação (não está na allowlist do operador).`, + }; + } + } + + // Chegou aqui: kind não-destrutivo, fora de report-only, e liberado pelo operador. + return { outcome: "allow" }; +} + +/** true se o efeito pode ser executado sem parar para humano, sob o contexto dado. */ +export function isAutoExecutable( + effect: LoopProposedEffect | undefined, + ctx: PolicyContext +): boolean { + return decideEffect(effect, ctx).outcome === "allow"; +} diff --git a/open-sse/loop-engine/stateMachine.ts b/open-sse/loop-engine/stateMachine.ts new file mode 100644 index 00000000000..de321c01f45 --- /dev/null +++ b/open-sse/loop-engine/stateMachine.ts @@ -0,0 +1,106 @@ +/** + * Loop Engine — máquina de estados do ciclo (puro, determinístico). + * + * Avança as fases discover→plan→split→execute→checkpoint→verify→budget→escalate, + * respeitando: orçamento (estourou → aborta), report-only (efeito → aguarda aprovação), + * verifier (reprovou → repete limitado por attempts, senão escala) e handoff humano. + * + * NÃO executa efeitos. Retorna sempre um NOVO estado (não muta a entrada). + */ +import { addUsage, checkBudget } from "./budget.ts"; +import { decideEffect, type PolicyContext } from "./policyGate.ts"; +import { LOOP_PHASES, type LoopPhase, type LoopRun, type LoopVerdict } from "./types.ts"; + +function nextPhase(phase: LoopPhase): LoopPhase { + const i = LOOP_PHASES.indexOf(phase); + // após 'escalate' o ciclo recomeça em 'discover' + return LOOP_PHASES[(i + 1) % LOOP_PHASES.length]; +} + +function clone(run: LoopRun): LoopRun { + return { + ...run, + usage: { ...run.usage }, + steps: run.steps.map((s) => ({ ...s })), + }; +} + +export interface AdvanceInput { + /** Consumo desta iteração (tokens/tempo/tentativa) — contabilizado no orçamento. */ + readonly consumed?: { tokens?: number; wallClockMs?: number; attempts?: number }; + /** Veredito do verifier, quando a fase for 'verify'. */ + readonly verdict?: LoopVerdict; + readonly policy: PolicyContext; +} + +export interface AdvanceResult { + readonly run: LoopRun; + /** Descrição legível da transição, para o audit log (sempre não-sensível). */ + readonly note: string; +} + +/** + * Executa UMA transição do ciclo. Fail-closed: qualquer efeito não auto-aprovado + * pela política para o run em 'awaiting_approval'. + */ +export function advance(run: LoopRun, input: AdvanceInput): AdvanceResult { + const next = clone(run); + next.sequenceNumber += 1; + next.usage = addUsage(next.usage, input.consumed ?? {}); + + // 1) Orçamento sempre primeiro: estourou → aborta e reporta. + const budget = checkBudget(next.budget, next.usage); + if (!budget.ok) { + next.status = "aborted"; + return { run: next, note: `budget excedido: ${budget.exceeded.join(",")}` }; + } + + // 2) Antes de propor efeito na fase execute, a política decide. + if (next.phase === "execute") { + const step = next.steps.find((s) => s.status === "proposed"); + if (step) { + const decision = decideEffect(step.proposedEffect, input.policy); + if (decision.outcome === "deny") { + step.status = "rejected"; + next.status = "failed"; + return { run: next, note: `efeito negado: ${decision.reason}` }; + } + if (decision.outcome === "require_approval") { + next.status = "awaiting_approval"; + return { run: next, note: `aguardando aprovação: ${decision.reason}` }; + } + step.status = "approved"; + } + } + + // 3) Verify: aplica o veredito do verifier. + if (next.phase === "verify" && input.verdict) { + const step = next.steps.find((s) => s.id === input.verdict!.stepId); + if (step) step.status = input.verdict.approved ? "verified" : "failed"; + if (!input.verdict.approved) { + // reprovou → o MOTOR conta a tentativa aqui (o teto é auto-imposto, não depende do + // chamador passar consumed.attempts) e repete limitado; se estourou, escala para humano. + next.usage.attempts += 1; + const canRetry = next.usage.attempts < next.budget.maxAttempts; + if (canRetry) { + next.phase = "plan"; + next.status = "report_only"; + return { run: next, note: `verifier reprovou; repetindo (attempt ${next.usage.attempts})` }; + } + next.status = "escalated"; + return { run: next, note: "verifier reprovou e sem tentativas: handoff humano" }; + } + } + + // 4) Conclusão: todos os passos verificados e ciclo passou por verify. + const allVerified = next.steps.length > 0 && next.steps.every((s) => s.status === "verified"); + if (allVerified && next.phase === "verify") { + next.status = "done"; + return { run: next, note: "todos os passos verificados: done" }; + } + + // 5) Avança para a próxima fase, mantendo report-only por padrão. + next.phase = nextPhase(next.phase); + if (next.status === "awaiting_approval") next.status = "report_only"; + return { run: next, note: `-> ${next.phase}` }; +} diff --git a/open-sse/loop-engine/types.ts b/open-sse/loop-engine/types.ts new file mode 100644 index 00000000000..6cbd2f7b071 --- /dev/null +++ b/open-sse/loop-engine/types.ts @@ -0,0 +1,104 @@ +/** + * Loop Engine — tipos do núcleo (report-only). + * + * Deriva da metodologia Loop Engineering (github.com/cobusgreyling/loop-engineering, + * MIT, SHA congelado 1d1af34b5d4b1af8bd5fc963ab2ff18499e706ff): um ciclo descobre + * trabalho, planeja, divide em passos, executa por agentes, faz checkpoint, verifica + * de forma independente, controla orçamento, repete de forma limitada e escala para + * humano. Aqui o Loop apenas PROPÕE; o Policy Engine/aprovações do OmniRoute decidem. + * + * NENHUM efeito externo é executado por este módulo. Ver policyGate.ts. + */ + +/** Fases canônicas do ciclo, na ordem. */ +export const LOOP_PHASES = [ + "discover", // descobrir trabalho (issues, tarefas, sinais) + "plan", // planejar a abordagem + "split", // dividir em etapas pequenas + "execute", // executar cada etapa (via agente) — PROPOSTA, não efeito + "checkpoint", // persistir estado recuperável + "verify", // verificação independente (verifier) + "budget", // reavaliar orçamento (tokens/tempo/tentativas) + "escalate", // escalar para humano (handoff) +] as const; + +export type LoopPhase = (typeof LOOP_PHASES)[number]; + +/** Estado terminal ou de espera de um run. */ +export type LoopRunStatus = + | "report_only" // rodando em modo report-only (padrão inicial) + | "awaiting_approval" // parou aguardando aprovação humana para um efeito + | "verifying" + | "done" + | "failed" + | "escalated" // handoff para humano + | "aborted"; // budget estourado / cancelado + +/** Orçamento de um run. Estourou → aborta e reporta (nunca ultrapassa). */ +export interface LoopBudget { + readonly maxTokens: number; + readonly maxWallClockMs: number; + readonly maxAttempts: number; +} + +/** Consumo acumulado, comparado ao LoopBudget. */ +export interface LoopBudgetUsage { + tokens: number; + wallClockMs: number; + attempts: number; +} + +/** Uma etapa proposta pelo ciclo. Report-only: nada é executado sem aprovação. */ +export interface LoopStep { + readonly id: string; + readonly runId: string; + readonly index: number; + readonly title: string; + /** Efeito externo que a etapa PROPÕE (ex.: abrir PR, enviar msg). Vazio = inócuo. */ + readonly proposedEffect?: LoopProposedEffect; + status: "proposed" | "approved" | "rejected" | "verified" | "failed"; +} + +/** Efeito externo proposto — sempre passa pelo Policy Engine antes de qualquer execução. */ +export interface LoopProposedEffect { + /** Categoria do efeito (para a política decidir). */ + readonly kind: + | "none" + | "git_commit" + | "git_pr" + | "message_send" + | "deploy" + | "purchase" + | "account_change" + | "delete"; + readonly summary: string; + readonly payload?: Record; +} + +/** Veredito da verificação independente (verifier). */ +export interface LoopVerdict { + readonly stepId: string; + readonly approved: boolean; + readonly reason: string; +} + +/** Um run do ciclo. Estado durável vive no DB do OmniRoute; este é o shape em memória. */ +export interface LoopRun { + readonly id: string; + readonly pattern: string; // ex.: "daily-triage", "pr-babysitter" + phase: LoopPhase; + status: LoopRunStatus; + readonly budget: LoopBudget; + usage: LoopBudgetUsage; + steps: LoopStep[]; + /** Idempotência / correlação para a ponte outbox-inbox. */ + readonly correlationId: string; + readonly taskId?: string; + sequenceNumber: number; +} + +/** Decisão do Policy Engine para um efeito proposto. Código clássico decide, não a IA. */ +export type PolicyDecision = + | { readonly outcome: "allow" } + | { readonly outcome: "deny"; readonly reason: string } + | { readonly outcome: "require_approval"; readonly reason: string }; diff --git a/open-sse/mcp-review/index.ts b/open-sse/mcp-review/index.ts new file mode 100644 index 00000000000..c7f9e313d06 --- /dev/null +++ b/open-sse/mcp-review/index.ts @@ -0,0 +1,143 @@ +/** + * MCP Review Gate — pipeline de segurança DETERMINÍSTICO para servidores/pacotes MCP. + * + * Deriva da Fase 6 do plano (MCP Marketplace): descobrir → quarentena → verificar → revisar → + * aprovar → instalar → habilitar. O CÓDIGO decide (não a IA), fail-closed: pacote malicioso é + * negado; qualquer novidade ou AMPLIAÇÃO de permissão volta a REVIEW_REQUIRED (aprovação humana). + * Puro (sem I/O). O instalador real (checksum/download) fica no marketplace do OmniRoute. + */ + +/** Estados do ciclo de revisão de um servidor MCP. */ +export type McpReviewState = + | "discovered" // encontrado no registry (read-only) + | "quarantined" // isolado até verificação + | "verified" // integridade/assinatura conferidas + | "review_required" // aguarda aprovação humana (novo ou permissão ampliada) + | "approved" // liberado para instalar/habilitar + | "denied"; // bloqueado (malicioso/forbidden) + +/** Um candidato do registry. Permissões são capabilities declaradas (ex.: "fs:read", "net:fetch"). */ +export interface McpCandidate { + readonly name: string; + readonly source: string; // origem (registry/url) + readonly version: string; + readonly permissions: ReadonlyArray; + /** Assinatura/publisher conferidos por processo externo (verify). */ + readonly publisherVerified?: boolean; + /** Marcado como malicioso por um scanner externo (hash conhecido, etc.). */ + readonly flaggedMalicious?: boolean; +} + +/** Estado prévio conhecido do MESMO servidor (para detectar ampliação de permissão em updates). */ +export interface McpPriorApproval { + readonly version: string; + readonly permissions: ReadonlyArray; + readonly approved: boolean; +} + +export interface McpReviewVerdict { + readonly state: McpReviewState; + readonly requiresHumanApproval: boolean; + readonly reasons: string[]; + /** Permissões novas em relação ao prior (vazio quando não amplia). */ + readonly newlyRequested: string[]; +} + +/** + * Permissões que, sozinhas, NUNCA se auto-aprovam — sempre exigem revisão humana explícita. + * (efeitos amplos/perigosos). Não são "proibidas", mas nunca passam sem gate humano. + */ +export const SENSITIVE_MCP_PERMISSIONS: ReadonlySet = new Set([ + "fs:write", + "fs:delete", + "shell:exec", + "process:spawn", + "net:listen", + "secrets:read", +]); + +/** + * Capabilities PROIBIDAS: presença → denied (bloqueio, como "pacote malicioso"). Um MCP não deve + * jamais pedir exfiltração de credenciais ou billing. + */ +export const FORBIDDEN_MCP_PERMISSIONS: ReadonlySet = new Set([ + "secrets:exfiltrate", + "billing:write", + "keys:read", +]); + +function broadenedPermissions( + candidate: ReadonlyArray, + prior: ReadonlyArray +): string[] { + const priorSet = new Set(prior); + return candidate.filter((p) => !priorSet.has(p)); +} + +/** + * Decide o próximo estado de revisão para um candidato MCP. Determinístico e fail-closed: + * - malicioso ou permissão proibida → denied. + * - novo (sem prior) → review_required. + * - update que AMPLIA permissões → review_required (re-revisão). + * - permissão sensível presente e ainda não aprovada para ela → review_required. + * - update sem ampliação, com prior aprovado → approved (carrega a aprovação). + */ +export function reviewMcpCandidate( + candidate: McpCandidate, + prior?: McpPriorApproval +): McpReviewVerdict { + const reasons: string[] = []; + + if (candidate.flaggedMalicious) { + return { + state: "denied", + requiresHumanApproval: false, + reasons: ["marcado como malicioso por scanner externo"], + newlyRequested: [], + }; + } + + const forbidden = candidate.permissions.filter((p) => FORBIDDEN_MCP_PERMISSIONS.has(p)); + if (forbidden.length > 0) { + return { + state: "denied", + requiresHumanApproval: false, + reasons: [`permissões proibidas: ${forbidden.join(", ")}`], + newlyRequested: [], + }; + } + + const newlyRequested = prior + ? broadenedPermissions(candidate.permissions, prior.permissions) + : [...candidate.permissions]; + + const hasSensitive = candidate.permissions.some((p) => SENSITIVE_MCP_PERMISSIONS.has(p)); + + // Novo servidor: sempre revisão humana. + if (!prior) { + reasons.push("servidor novo — requer revisão humana"); + if (hasSensitive) reasons.push("declara permissões sensíveis"); + return { state: "review_required", requiresHumanApproval: true, reasons, newlyRequested }; + } + + // Update que amplia permissões → volta para revisão (mesmo que antes aprovado). + if (newlyRequested.length > 0) { + reasons.push(`amplia permissões: ${newlyRequested.join(", ")} — re-revisão obrigatória`); + return { state: "review_required", requiresHumanApproval: true, reasons, newlyRequested }; + } + + // Sem ampliação: carrega a aprovação anterior — MAS fail-closed no publisher: se a verificação + // externa REPROVOU o publisher (publisherVerified === false), volta à revisão mesmo sem ampliar + // (assinatura/publisher divergente é sinal de comprometimento, não um simples patch). + if (prior.approved) { + if (candidate.publisherVerified === false) { + reasons.push("publisher não verificado — re-revisão obrigatória apesar de não ampliar"); + return { state: "review_required", requiresHumanApproval: true, reasons, newlyRequested: [] }; + } + reasons.push("update sem novas permissões — aprovação anterior mantida"); + return { state: "approved", requiresHumanApproval: false, reasons, newlyRequested: [] }; + } + + reasons.push("sem aprovação anterior — requer revisão humana"); + return { state: "review_required", requiresHumanApproval: true, reasons, newlyRequested: [] }; +} diff --git a/open-sse/otel/index.ts b/open-sse/otel/index.ts new file mode 100644 index 00000000000..12bf3795f8e --- /dev/null +++ b/open-sse/otel/index.ts @@ -0,0 +1,128 @@ +/** + * OTel-lite — tracing distribuído leve (Fase 8), SEM dependência pesada e SEM conteúdo sensível. + * + * Implementa W3C Trace Context (traceparent) e spans com atributos ALLOWLISTED: nenhum prompt, + * resposta, PII ou segredo entra num span — só metadados operacionais (provider, model, rota, + * latência, tokens, status). Exportador plugável (InMemory p/ teste; OTLP fica na camada externa). + * Puro, determinístico onde possível; ids aleatórios via crypto. + */ +import { randomBytes } from "node:crypto"; + +export interface SpanContext { + readonly traceId: string; // 32 hex + readonly spanId: string; // 16 hex + readonly parentSpanId?: string; +} + +export type SpanStatus = "unset" | "ok" | "error"; + +export interface Span extends SpanContext { + name: string; + startMs: number; + endMs?: number; + status: SpanStatus; + attributes: Record; +} + +export interface SpanExporter { + export(span: Span): void; +} + +export class InMemorySpanExporter implements SpanExporter { + readonly spans: Span[] = []; + export(span: Span): void { + this.spans.push(span); + } +} + +/** Chaves de atributo PERMITIDAS num span. Qualquer outra é descartada (anti-vazamento). */ +export const SAFE_ATTR_KEYS: ReadonlySet = new Set([ + "provider", + "model", + "route", + "http.method", + "http.status_code", + "latency_ms", + "tokens_in", + "tokens_out", + "tenant", + "cache", + "error.kind", +]); + +const MAX_ATTR_STR = 64; + +/** + * Mantém APENAS chaves allowlisted e valores primitivos; strings são truncadas. Isso garante que + * conteúdo sensível (prompts/respostas/segredos) nunca chegue à telemetria, mesmo por engano. + */ +export function safeAttributes( + attrs: Record +): Record { + const out: Record = {}; + for (const [k, v] of Object.entries(attrs)) { + if (!SAFE_ATTR_KEYS.has(k)) continue; + if (typeof v === "number" || typeof v === "boolean") out[k] = v; + else if (typeof v === "string") out[k] = v.slice(0, MAX_ATTR_STR); + } + return out; +} + +const hex = (bytes: number): string => randomBytes(bytes).toString("hex"); + +export function genTraceId(): string { + return hex(16); +} +export function genSpanId(): string { + return hex(8); +} + +/** Formata um traceparent W3C: 00---. */ +export function formatTraceparent(ctx: SpanContext, sampled = true): string { + return `00-${ctx.traceId}-${ctx.spanId}-${sampled ? "01" : "00"}`; +} + +/** Faz parse de um traceparent W3C. Retorna null se malformado (fail-closed). */ +export function parseTraceparent( + header: string +): { traceId: string; spanId: string; sampled: boolean } | null { + const m = /^00-([0-9a-f]{32})-([0-9a-f]{16})-([0-9a-f]{2})$/i.exec(header.trim()); + if (!m) return null; + if (/^0+$/.test(m[1]) || /^0+$/.test(m[2])) return null; // all-zero é inválido + return { + traceId: m[1].toLowerCase(), + spanId: m[2].toLowerCase(), + sampled: (parseInt(m[3], 16) & 1) === 1, + }; +} + +/** Cria um novo span. Se `parent` for dado, herda o traceId e referencia o parentSpanId. */ +export function startSpan( + name: string, + parent?: SpanContext, + attrs: Record = {}, + now: number = Date.now() +): Span { + return { + traceId: parent?.traceId ?? genTraceId(), + spanId: genSpanId(), + parentSpanId: parent?.spanId, + name, + startMs: now, + status: "unset", + attributes: safeAttributes(attrs), + }; +} + +/** Encerra o span e o exporta. Mescla atributos finais (sempre pela allowlist). */ +export function endSpan( + span: Span, + exporter: SpanExporter, + opts: { status?: SpanStatus; attributes?: Record; now?: number } = {} +): Span { + span.endMs = opts.now ?? Date.now(); + if (opts.status) span.status = opts.status; + if (opts.attributes) Object.assign(span.attributes, safeAttributes(opts.attributes)); + exporter.export(span); + return span; +} diff --git a/package-lock.json b/package-lock.json index 6f720491dfd..1100897c177 100644 --- a/package-lock.json +++ b/package-lock.json @@ -22,6 +22,8 @@ "@modelcontextprotocol/sdk": "^1.30.0", "@monaco-editor/react": "^4.7.0", "@ngrok/ngrok": "^1.7.0", + "@noble/curves": "^2.4.0", + "@noble/hashes": "^2.4.0", "@swc/helpers": "0.5.23", "@toon-format/toon": "^4.1.1", "@types/mdx": "^2.0.13", @@ -6937,6 +6939,33 @@ "node": ">= 10" } }, + "node_modules/@noble/curves": { + "version": "2.4.0", + "resolved": "https://registry.npmjs.org/@noble/curves/-/curves-2.4.0.tgz", + "integrity": "sha512-P4/62zrgfH33CneE3Dn4WhJVA22YUU0eR51wKIan4NVRvwsA0YnPTwWGpNbpuacSujmSFLvyzpyuR30+fbq2Ew==", + "license": "MIT", + "dependencies": { + "@noble/hashes": "2.4.0" + }, + "engines": { + "node": ">= 20.19.0" + }, + "funding": { + "url": "https://paulmillr.com/funding/" + } + }, + "node_modules/@noble/hashes": { + "version": "2.4.0", + "resolved": "https://registry.npmjs.org/@noble/hashes/-/hashes-2.4.0.tgz", + "integrity": "sha512-X5XaVWZIBCT7HHZGm5I7ZQXDwLG+bGXuSrMQAW+7Zvl87h1kmc1ZB1VSRJcpUfoUrGQp4Fkoxm5kZ+Ms+aW+eA==", + "license": "MIT", + "engines": { + "node": ">= 20.19.0" + }, + "funding": { + "url": "https://paulmillr.com/funding/" + } + }, "node_modules/@nodable/entities": { "version": "3.0.0", "resolved": "https://registry.npmjs.org/@nodable/entities/-/entities-3.0.0.tgz", diff --git a/package.json b/package.json index f984ef4314f..56229414193 100644 --- a/package.json +++ b/package.json @@ -290,6 +290,8 @@ "@modelcontextprotocol/sdk": "^1.30.0", "@monaco-editor/react": "^4.7.0", "@ngrok/ngrok": "^1.7.0", + "@noble/curves": "^2.4.0", + "@noble/hashes": "^2.4.0", "@swc/helpers": "0.5.23", "@toon-format/toon": "^4.1.1", "@types/mdx": "^2.0.13", diff --git a/src/app/(dashboard)/dashboard/buzz/page.tsx b/src/app/(dashboard)/dashboard/buzz/page.tsx new file mode 100644 index 00000000000..66d189a115f --- /dev/null +++ b/src/app/(dashboard)/dashboard/buzz/page.tsx @@ -0,0 +1,263 @@ +"use client"; + +import { useCallback, useEffect, useState } from "react"; + +interface BuzzCounts { + outboxPending: number; + outboxPublished: number; + outboxFailed: number; + inboxReceived: number; +} +interface BuzzStatus { + enabled: boolean; + relayUrl: string; + agentPubkey: string; + counts: BuzzCounts; +} + +function Stat({ label, value, tone }: { label: string; value: number; tone?: string }) { + return ( +
+

{value}

+

{label}

+
+ ); +} + +export default function BuzzHubPage() { + const [status, setStatus] = useState(null); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const [relayDraft, setRelayDraft] = useState(""); + const [savingRelay, setSavingRelay] = useState(false); + const [flushing, setFlushing] = useState(false); + const [flushMsg, setFlushMsg] = useState(null); + const [copied, setCopied] = useState(false); + + const load = useCallback(async () => { + setLoading(true); + setError(null); + try { + const res = await fetch("/api/buzz"); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + const data = (await res.json()) as BuzzStatus; + setStatus(data); + setRelayDraft(data.relayUrl); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao carregar"); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + // eslint-disable-next-line react-hooks/set-state-in-effect -- carga inicial do status do Buzz (a mesma seq de load é reusada pelo refresh) + void load(); + }, [load]); + + const saveRelay = useCallback(async () => { + setSavingRelay(true); + setError(null); + try { + const res = await fetch("/api/buzz", { + method: "PUT", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ relayUrl: relayDraft }), + }); + if (!res.ok) { + const d = await res.json().catch(() => ({})); + throw new Error(d.error || `HTTP ${res.status}`); + } + const data = (await res.json()) as BuzzStatus; + setStatus(data); + setRelayDraft(data.relayUrl); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao salvar"); + } finally { + setSavingRelay(false); + } + }, [relayDraft]); + + const flush = useCallback(async () => { + setFlushing(true); + setFlushMsg(null); + try { + const res = await fetch("/api/buzz/flush", { method: "POST" }); + const data = await res.json().catch(() => ({})); + if (data.skipped) { + setFlushMsg( + data.reason === "BUZZ_HUB_ENABLED is off" + ? "Flush ignorado: a flag BUZZ_HUB_ENABLED está desligada." + : "Nada pendente para publicar." + ); + } else { + setFlushMsg(`Publicados: ${data.published ?? 0} · falhas: ${data.failed ?? 0}`); + } + await load(); + } catch (e) { + setFlushMsg(e instanceof Error ? e.message : "Falha no flush"); + } finally { + setFlushing(false); + } + }, [load]); + + const copyPubkey = useCallback(() => { + if (!status) return; + void navigator.clipboard?.writeText(status.agentPubkey).then(() => { + setCopied(true); + setTimeout(() => setCopied(false), 1500); + }); + }, [status]); + + return ( +
+
+

Buzz Hub

+ {status && ( + + {status.enabled ? "ativado" : "desligado"} + + )} +
+

+ Colaboração humano+agente sobre um relay Nostr auto-hospedado (auth NIP-42, + outbox/inbox idempotente). A chave Nostr nunca autoriza uma ação no + OmniRoute. Configurado aqui, no painel único. +

+ + {loading &&
} + + {!loading && error && ( +
+ {error} + +
+ )} + + {!loading && status && ( + <> + {!status.enabled && ( +
+
+ +
+

+ Buzz Hub está desligado. +

+

+ A ponte fica durável no DB (outbox/inbox), mas nada conecta ao relay. Ative a + flag BUZZ_HUB_ENABLED para publicar. +

+ + + Abrir Feature Flags + +
+
+
+ )} + + {/* Relay URL */} +
+ +

+ Precedência: este valor → env BUZZ_RELAY_URL → + padrão. Deixe vazio e salve para voltar ao env/padrão. +

+
+ setRelayDraft(e.target.value)} + placeholder="ws://localhost:3000" + className="flex-1 rounded-lg border border-border bg-bg-subtle px-3 py-2 font-mono text-sm text-text-primary placeholder:text-text-muted focus:border-primary focus:outline-none focus:ring-2 focus:ring-primary/15" + /> + +
+
+ + {/* Identidade do agente */} +
+

Identidade Nostr do agente

+

+ Chave pública estável do OmniRoute no relay (a secreta nunca sai do servidor). +

+
+ + {status.agentPubkey} + + +
+
+ + {/* Contagens */} +
+ 0 ? "text-amber-600 dark:text-amber-300" : undefined + } + /> + + 0 ? "text-red-600 dark:text-red-300" : undefined} + /> + +
+ + {/* Flush — desabilitado com a flag OFF (paridade com o Loop; nada a publicar sem relay). */} +
+ + {flushMsg && {flushMsg}} +
+ + )} +
+ ); +} diff --git a/src/app/(dashboard)/dashboard/loop/page.tsx b/src/app/(dashboard)/dashboard/loop/page.tsx new file mode 100644 index 00000000000..5d72be0704f --- /dev/null +++ b/src/app/(dashboard)/dashboard/loop/page.tsx @@ -0,0 +1,342 @@ +"use client"; + +import { useCallback, useEffect, useState } from "react"; + +// Espelha open-sse/loop-engine/types.ts (mantido leve para o cliente). +interface LoopStep { + id: string; + runId: string; + index: number; + title: string; + proposedEffect?: { kind: string; summary: string }; + status: "proposed" | "approved" | "rejected" | "verified" | "failed"; +} +interface LoopRun { + id: string; + pattern: string; + phase: string; + status: + "report_only" | "awaiting_approval" | "verifying" | "done" | "failed" | "escalated" | "aborted"; + budget: { maxTokens: number; maxWallClockMs: number; maxAttempts: number }; + usage: { tokens: number; wallClockMs: number; attempts: number }; + steps: LoopStep[]; + correlationId: string; + taskId?: string; + sequenceNumber: number; +} + +const STATUS_STYLE: Record = { + report_only: "bg-sky-100 text-sky-800 dark:bg-sky-500/15 dark:text-sky-300", + awaiting_approval: "bg-amber-100 text-amber-800 dark:bg-amber-500/15 dark:text-amber-300", + verifying: "bg-violet-100 text-violet-800 dark:bg-violet-500/15 dark:text-violet-300", + done: "bg-emerald-100 text-emerald-800 dark:bg-emerald-500/15 dark:text-emerald-300", + failed: "bg-red-100 text-red-800 dark:bg-red-500/15 dark:text-red-300", + escalated: "bg-orange-100 text-orange-800 dark:bg-orange-500/15 dark:text-orange-300", + aborted: "bg-neutral-200 text-neutral-700 dark:bg-white/10 dark:text-neutral-300", +}; + +function DisabledNotice() { + return ( +
+
+ +
+

+ Loop Engine está desligado. +

+

+ Ative a flag LOOP_ENGINE_ENABLED no painel único para + usar os ciclos report-only e as aprovações humanas. +

+ + + Abrir Feature Flags + +
+
+
+ ); +} + +export default function LoopEnginePage() { + const [runs, setRuns] = useState([]); + const [loading, setLoading] = useState(true); + const [disabled, setDisabled] = useState(false); + const [error, setError] = useState(null); + const [pattern, setPattern] = useState(""); + const [starting, setStarting] = useState(false); + const [expanded, setExpanded] = useState>(new Set()); + const [busyRun, setBusyRun] = useState(null); + + const load = useCallback(async () => { + setLoading(true); + setError(null); + try { + const res = await fetch("/api/loop"); + if (res.status === 404) { + setDisabled(true); + setRuns([]); + return; + } + if (!res.ok) throw new Error(`HTTP ${res.status}`); + setDisabled(false); + const data = await res.json(); + setRuns(Array.isArray(data.runs) ? data.runs : []); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao carregar"); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + // eslint-disable-next-line react-hooks/set-state-in-effect -- carga inicial dos runs (a mesma seq de load é reusada pelo refresh) + void load(); + }, [load]); + + const startRun = useCallback(async () => { + const p = pattern.trim(); + if (!p) return; + setStarting(true); + setError(null); + try { + const res = await fetch("/api/loop", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ pattern: p }), + }); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + setPattern(""); + await load(); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao iniciar"); + } finally { + setStarting(false); + } + }, [pattern, load]); + + const showApiError = useCallback(async (res: Response, fallback: string) => { + const data = await res.json().catch(() => ({})); + setError(typeof data.error === "string" ? data.error : fallback); + }, []); + + const advance = useCallback( + async (id: string) => { + setBusyRun(id); + setError(null); + try { + const res = await fetch(`/api/loop/${id}/advance`, { method: "POST" }); + if (!res.ok) await showApiError(res, `Falha ao avançar (HTTP ${res.status})`); + await load(); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao avançar"); + } finally { + setBusyRun(null); + } + }, + [load, showApiError] + ); + + const decide = useCallback( + async (id: string, stepId: string, decision: "approve" | "reject") => { + setBusyRun(id); + setError(null); + try { + const res = await fetch(`/api/loop/${id}/approve`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ stepId, decision }), + }); + if (!res.ok) await showApiError(res, `Falha ao ${decision} (HTTP ${res.status})`); + await load(); + } catch (e) { + setError(e instanceof Error ? e.message : "Falha ao decidir"); + } finally { + setBusyRun(null); + } + }, + [load, showApiError] + ); + + const toggle = (id: string) => + setExpanded((prev) => { + const next = new Set(prev); + if (next.has(id)) next.delete(id); + else next.add(id); + return next; + }); + + return ( +
+
+

Loop Engine

+

+ Ciclos agênticos report-only: o Loop propõe, o Policy Engine decide e + nenhum efeito externo roda sem aprovação humana. Serve ao painel e a qualquer harness via{" "} + /api/loop. +

+
+ + {loading &&
} + + {!loading && disabled && } + + {!loading && !disabled && ( + <> + {/* Novo run */} +
+ setPattern(e.target.value)} + onKeyDown={(e) => e.key === "Enter" && void startRun()} + placeholder="Padrão do ciclo (ex.: pr-babysitter, daily-triage)" + className="flex-1 rounded-lg border border-border bg-card px-3 py-2 text-sm text-text-primary placeholder:text-text-muted focus:border-primary focus:outline-none focus:ring-2 focus:ring-primary/15" + /> + +
+ + {error && ( +
+ {error} +
+ )} + + {runs.length === 0 ? ( +
+ +

Nenhum ciclo ainda. Inicie um acima.

+
+ ) : ( +
+ {runs.map((run) => { + const isOpen = expanded.has(run.id); + const busy = busyRun === run.id; + return ( +
+ + + {isOpen && ( +
+
+ +
+ + {run.steps.length === 0 ? ( +

Sem etapas propostas ainda.

+ ) : ( +
    + {run.steps.map((step) => { + const canDecide = + run.status === "awaiting_approval" && step.status === "proposed"; + return ( +
  • +
    +

    + {step.index + 1}. {step.title} +

    + {step.proposedEffect && step.proposedEffect.kind !== "none" && ( +

    + efeito proposto:{" "} + + {step.proposedEffect.kind} + {" "} + — {step.proposedEffect.summary} +

    + )} +
    + {canDecide ? ( +
    + + +
    + ) : ( + + {step.status} + + )} +
  • + ); + })} +
+ )} +
+ )} +
+ ); + })} +
+ )} + + )} +
+ ); +} diff --git a/src/app/api/browser/check/route.ts b/src/app/api/browser/check/route.ts new file mode 100644 index 00000000000..71715569f99 --- /dev/null +++ b/src/app/api/browser/check/route.ts @@ -0,0 +1,64 @@ +/** + * POST /api/browser/check — avalia uma ação de navegador pela política determinística. + * Body: { action: BrowserAction, allowedDomains?: string[] }. + * + * Autenticado (management) e gated por BROWSER_USE_ENABLED. Se allowedDomains não vier no corpo, + * usa o override persistido no painel (key_value namespace 'browser'). Efeito externo → aprovação + * humana; efeito originado na página → deny (prompt injection não escala). Não executa nada. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { getDbInstance } from "@/lib/db/core"; +import { traceSync } from "@/lib/otel"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; +import { + decideBrowserAction, + type BrowserAction, +} from "@omniroute/open-sse/browser-guard/index.ts"; + +function storedAllowlist(): string[] { + try { + const row = getDbInstance() + .prepare( + "SELECT value FROM key_value WHERE namespace = 'browser' AND key = 'allowed_domains'" + ) + .get() as { value: string } | undefined; + return row?.value ? (JSON.parse(row.value) as string[]) : []; + } catch { + return []; + } +} + +export async function POST(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("BROWSER_USE_ENABLED")) { + return NextResponse.json( + { error: "Browser Use is disabled. Enable BROWSER_USE_ENABLED in the OmniRoute panel." }, + { status: 404 } + ); + } + + const body = (await req.json().catch(() => ({}))) as { + action?: Partial; + allowedDomains?: unknown; + }; + const a = body.action; + if (!a || typeof a.kind !== "string" || (a.origin !== "user" && a.origin !== "page")) { + return NextResponse.json( + { error: "action {kind, origin: 'user'|'page', url?} is required" }, + { status: 400 } + ); + } + const allowedDomains = Array.isArray(body.allowedDomains) + ? (body.allowedDomains.filter((d) => typeof d === "string") as string[]) + : storedAllowlist(); + + const verdict = traceSync( + "browser.check", + { route: "/api/browser/check", "http.method": "POST" }, + () => decideBrowserAction(a as BrowserAction, { enabled: true, allowedDomains }) + ); + return NextResponse.json({ verdict, allowedDomains }); +} diff --git a/src/app/api/buzz/flush/route.ts b/src/app/api/buzz/flush/route.ts new file mode 100644 index 00000000000..d691398b32d --- /dev/null +++ b/src/app/api/buzz/flush/route.ts @@ -0,0 +1,31 @@ +/** + * POST /api/buzz/flush — publica as entradas pendentes do outbox no relay real (idempotente). + * Gated por BUZZ_HUB_ENABLED: com a flag OFF retorna { skipped: true } sem conectar. + * + * Serve ao painel único (botão "Publicar pendentes") e a qualquer harness/automação. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { flushBuzzOutbox } from "@/lib/buzzService"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export async function POST(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("BUZZ_HUB_ENABLED")) { + return NextResponse.json( + { published: 0, failed: 0, skipped: true, reason: "BUZZ_HUB_ENABLED is off" }, + { status: 200 } + ); + } + try { + const result = await flushBuzzOutbox(); + return NextResponse.json(result); + } catch (e) { + return NextResponse.json( + { error: e instanceof Error ? e.message : "flush failed" }, + { status: 502 } + ); + } +} diff --git a/src/app/api/buzz/route.ts b/src/app/api/buzz/route.ts new file mode 100644 index 00000000000..e0c83e538c6 --- /dev/null +++ b/src/app/api/buzz/route.ts @@ -0,0 +1,44 @@ +/** + * GET /api/buzz — estado do Buzz Hub para o painel único: flag, relay, pubkey do agente e + * contagens do outbox/inbox. NUNCA expõe a chave secreta. Leitura local (não conecta ao relay). + * + * Autenticado (management). Ao contrário do Loop, o status É legível mesmo com a flag OFF — + * para o painel poder mostrar "desligado" e orientar a ativação. Nada conecta enquanto OFF. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { getBuzzStatus, setBuzzRelayUrl } from "@/lib/buzzService"; + +export async function GET(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + return NextResponse.json(getBuzzStatus()); +} + +/** + * PUT /api/buzz — define a URL do relay pelo painel único. Body: { relayUrl: string }. + * + * INTENCIONALMENTE não é gated por BUZZ_HUB_ENABLED (diferente de /flush): é configuração + * (não-secreta, admin-only) que o operador ajusta ANTES de ligar a flag. Nada conecta aqui — + * a conexão só ocorre no flush, que é gated. O alvo é local por design (default ws://localhost:3000); + * apontar para um host interno exige management-auth (o admin já controla o host). + */ +export async function PUT(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + const body = (await req.json().catch(() => ({}))) as { relayUrl?: unknown }; + if (typeof body.relayUrl !== "string") { + return NextResponse.json({ error: "relayUrl (string) is required" }, { status: 400 }); + } + const trimmed = body.relayUrl.trim(); + // Aceita ws:// ou wss:// (ou vazio para limpar o override). Evita URLs inválidas no relay. + if (trimmed && !/^wss?:\/\//i.test(trimmed)) { + return NextResponse.json( + { error: "relayUrl must start with ws:// or wss://" }, + { status: 400 } + ); + } + setBuzzRelayUrl(trimmed); + return NextResponse.json(getBuzzStatus()); +} diff --git a/src/app/api/loop/[id]/advance/route.ts b/src/app/api/loop/[id]/advance/route.ts new file mode 100644 index 00000000000..7e36ba8086c --- /dev/null +++ b/src/app/api/loop/[id]/advance/route.ts @@ -0,0 +1,34 @@ +/** + * POST /api/loop/[id]/advance — avança UM passo do run (report-only, persistente). + * Body opcional: { consumed?: {tokens?,wallClockMs?,attempts?} }. + * O gate segura efeitos externos em awaiting_approval; nada é executado aqui. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { advanceRun } from "@/lib/loopRunner"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export async function POST( + req: NextRequest, + { params }: { params: Promise<{ id: string }> } +): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) { + return NextResponse.json({ error: "Loop Engine is disabled." }, { status: 404 }); + } + const { id } = await params; + const body = (await req.json().catch(() => ({}))) as { + consumed?: { tokens?: number; wallClockMs?: number; attempts?: number }; + }; + try { + const result = advanceRun(id, { consumed: body.consumed, policy: { reportOnly: true } }); + return NextResponse.json(result); + } catch (e) { + return NextResponse.json( + { error: e instanceof Error ? e.message : "advance failed" }, + { status: 404 } + ); + } +} diff --git a/src/app/api/loop/[id]/approve/route.ts b/src/app/api/loop/[id]/approve/route.ts new file mode 100644 index 00000000000..dfbf6da74eb --- /dev/null +++ b/src/app/api/loop/[id]/approve/route.ts @@ -0,0 +1,39 @@ +/** + * POST /api/loop/[id]/approve — aprova ou rejeita uma etapa aguardando aprovação humana. + * Body: { stepId: string, decision: "approve" | "reject" }. + * Aprovar libera o run; rejeitar marca a etapa. Efeito externo real (quando aprovado) é + * executado pelos conectores do OmniRoute, fora daqui. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { approveStep, rejectStep } from "@/lib/loopRunner"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export async function POST( + req: NextRequest, + { params }: { params: Promise<{ id: string }> } +): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) { + return NextResponse.json({ error: "Loop Engine is disabled." }, { status: 404 }); + } + const { id } = await params; + const body = (await req.json().catch(() => ({}))) as { + stepId?: unknown; + decision?: unknown; + }; + const stepId = typeof body.stepId === "string" ? body.stepId : ""; + const decision = body.decision === "reject" ? "reject" : "approve"; + if (!stepId) return NextResponse.json({ error: "stepId is required" }, { status: 400 }); + try { + const run = decision === "reject" ? rejectStep(id, stepId) : approveStep(id, stepId); + return NextResponse.json({ run, decision }); + } catch (e) { + return NextResponse.json( + { error: e instanceof Error ? e.message : "approve failed" }, + { status: 404 } + ); + } +} diff --git a/src/app/api/loop/[id]/route.ts b/src/app/api/loop/[id]/route.ts new file mode 100644 index 00000000000..3ab8f59bfcb --- /dev/null +++ b/src/app/api/loop/[id]/route.ts @@ -0,0 +1,24 @@ +/** + * GET /api/loop/[id] — retorna um run do Loop Engine (com suas etapas). + * Autenticado + gated por LOOP_ENGINE_ENABLED. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { getLoopRun } from "@/lib/loopRunner"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export async function GET( + req: NextRequest, + { params }: { params: Promise<{ id: string }> } +): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) { + return NextResponse.json({ error: "Loop Engine is disabled." }, { status: 404 }); + } + const { id } = await params; + const run = getLoopRun(id); + if (!run) return NextResponse.json({ error: "run not found" }, { status: 404 }); + return NextResponse.json({ run }); +} diff --git a/src/app/api/loop/[id]/stream/route.ts b/src/app/api/loop/[id]/stream/route.ts new file mode 100644 index 00000000000..ae03a68f9d5 --- /dev/null +++ b/src/app/api/loop/[id]/stream/route.ts @@ -0,0 +1,69 @@ +/** + * GET /api/loop/[id]/stream — transmite o estado de um run do Loop como eventos AG-UI (SSE). + * + * Autenticado (management) e gated por LOOP_ENGINE_ENABLED. Emite RUN_STARTED, um bloco de texto + * por etapa (START/CONTENT/END), um STATE_SNAPSHOT (fase/status) e RUN_FINISHED — fluxo AG-UI válido + * (ordenado por seq), consumível por um EventSource no Agent Console. Report-only (só leitura). + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { getLoopRun } from "@/lib/loopRunner"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; +import { createSeq, encodeSse, type AgUiEvent } from "@omniroute/open-sse/ag-ui/index.ts"; + +export async function GET( + req: NextRequest, + { params }: { params: Promise<{ id: string }> } +): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) { + return NextResponse.json({ error: "Loop Engine is disabled." }, { status: 404 }); + } + const { id } = await params; + const run = getLoopRun(id); + if (!run) return NextResponse.json({ error: "run not found" }, { status: 404 }); + + const seq = createSeq(); + const events: AgUiEvent[] = [{ seq: seq(), runId: run.id, type: "RUN_STARTED" }]; + for (const step of run.steps) { + const messageId = `step-${step.index}`; + events.push({ + seq: seq(), + runId: run.id, + type: "TEXT_MESSAGE_START", + messageId, + role: "assistant", + }); + const effect = + step.proposedEffect && step.proposedEffect.kind !== "none" + ? ` [efeito proposto: ${step.proposedEffect.kind}]` + : ""; + events.push({ + seq: seq(), + runId: run.id, + type: "TEXT_MESSAGE_CONTENT", + messageId, + delta: `${step.index + 1}. ${step.title} (${step.status})${effect}`, + }); + events.push({ seq: seq(), runId: run.id, type: "TEXT_MESSAGE_END", messageId }); + } + events.push({ + seq: seq(), + runId: run.id, + type: "STATE_SNAPSHOT", + state: { phase: run.phase, status: run.status, sequenceNumber: run.sequenceNumber }, + }); + events.push({ seq: seq(), runId: run.id, type: "RUN_FINISHED" }); + + const body = events.map(encodeSse).join(""); + return new Response(body, { + status: 200, + headers: { + "Content-Type": "text/event-stream; charset=utf-8", + "Cache-Control": "no-cache, no-transform", + Connection: "keep-alive", + }, + }); +} diff --git a/src/app/api/loop/route.ts b/src/app/api/loop/route.ts new file mode 100644 index 00000000000..50b1cec54b8 --- /dev/null +++ b/src/app/api/loop/route.ts @@ -0,0 +1,67 @@ +/** + * GET /api/loop — lista runs do Loop Engine (opcional ?status=...). + * POST /api/loop — inicia um run (report-only) { pattern, budget?, taskId? }. + * + * Autenticado (management) e gated pela flag LOOP_ENGINE_ENABLED. Report-only: iniciar um + * run NÃO executa efeito externo. Serve ao painel único e a qualquer harness. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { listLoopRuns, startRun } from "@/lib/loopRunner"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; +import type { LoopRun } from "@omniroute/open-sse/loop-engine/index.ts"; + +const LOOP_STATUSES: ReadonlyArray = [ + "report_only", + "awaiting_approval", + "verifying", + "done", + "failed", + "escalated", + "aborted", +]; + +function disabled(): NextResponse { + return NextResponse.json( + { error: "Loop Engine is disabled. Enable LOOP_ENGINE_ENABLED in the OmniRoute panel." }, + { status: 404 } + ); +} + +export async function GET(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) return disabled(); + + const raw = new URL(req.url).searchParams.get("status"); + const status = + raw && LOOP_STATUSES.includes(raw as LoopRun["status"]) + ? (raw as LoopRun["status"]) + : undefined; + return NextResponse.json({ runs: listLoopRuns(status) }); +} + +export async function POST(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("LOOP_ENGINE_ENABLED")) return disabled(); + + const body = (await req.json().catch(() => ({}))) as { + pattern?: unknown; + budget?: Record; + taskId?: unknown; + correlationId?: unknown; + }; + const pattern = typeof body.pattern === "string" ? body.pattern.trim() : ""; + if (!pattern) { + return NextResponse.json({ error: "pattern is required" }, { status: 400 }); + } + const run = startRun({ + pattern, + budget: body.budget, + taskId: typeof body.taskId === "string" ? body.taskId : undefined, + correlationId: typeof body.correlationId === "string" ? body.correlationId : undefined, + }); + return NextResponse.json({ run }, { status: 201 }); +} diff --git a/src/app/api/mcp/review/route.ts b/src/app/api/mcp/review/route.ts new file mode 100644 index 00000000000..8dcb00dba3e --- /dev/null +++ b/src/app/api/mcp/review/route.ts @@ -0,0 +1,51 @@ +/** + * POST /api/mcp/review — roda o MCP Review Gate determinístico sobre um candidato. + * Body: { candidate: McpCandidate, prior?: McpPriorApproval }. + * + * Autenticado (management) e gated por MCP_REVIEW_ENABLED. Código decide (não a IA): malicioso/ + * permissão proibida → denied; novo ou permissão ampliada → review_required (aprovação humana). + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { traceSync } from "@/lib/otel"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; +import { + reviewMcpCandidate, + type McpCandidate, + type McpPriorApproval, +} from "@omniroute/open-sse/mcp-review/index.ts"; + +export async function POST(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("MCP_REVIEW_ENABLED")) { + return NextResponse.json( + { error: "MCP Review is disabled. Enable MCP_REVIEW_ENABLED in the OmniRoute panel." }, + { status: 404 } + ); + } + + const body = (await req.json().catch(() => ({}))) as { + candidate?: Partial; + prior?: McpPriorApproval; + }; + const c = body.candidate; + if ( + !c || + typeof c.name !== "string" || + typeof c.source !== "string" || + typeof c.version !== "string" || + !Array.isArray(c.permissions) + ) { + return NextResponse.json( + { error: "candidate {name, source, version, permissions[]} is required" }, + { status: 400 } + ); + } + + const verdict = traceSync("mcp.review", { route: "/api/mcp/review", provider: "mcp" }, () => + reviewMcpCandidate(c as McpCandidate, body.prior) + ); + return NextResponse.json({ verdict }); +} diff --git a/src/app/api/otel/spans/route.ts b/src/app/api/otel/spans/route.ts new file mode 100644 index 00000000000..fee6b44cdcd --- /dev/null +++ b/src/app/api/otel/spans/route.ts @@ -0,0 +1,25 @@ +/** + * GET /api/otel/spans — lê os spans recentes coletados (W3C Trace Context), sem conteúdo sensível. + * + * Autenticado (management) e gated por OTEL_TRACING_ENABLED. Os atributos já passaram pela allowlist + * do core (nada de prompt/resposta/PII/segredo). Útil para inspecionar latência/rota/provider. + */ +import { NextRequest, NextResponse } from "next/server"; + +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; +import { recentSpans } from "@/lib/otel"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export async function GET(req: NextRequest): Promise { + const auth = await requireManagementAuth(req); + if (auth) return auth; + if (!isFeatureFlagEnabled("OTEL_TRACING_ENABLED")) { + return NextResponse.json( + { error: "OTel tracing is disabled. Enable OTEL_TRACING_ENABLED in the OmniRoute panel." }, + { status: 404 } + ); + } + const raw = new URL(req.url).searchParams.get("limit"); + const limit = Math.min(Math.max(Number(raw) || 100, 1), 500); + return NextResponse.json({ spans: recentSpans(limit) }); +} diff --git a/src/app/callback/page.tsx b/src/app/callback/page.tsx index 3be2d56b293..4f5677cea2f 100644 --- a/src/app/callback/page.tsx +++ b/src/app/callback/page.tsx @@ -1,5 +1,6 @@ "use client"; +import Link from "next/link"; import { useTranslations } from "next-intl"; import { useEffect, useState } from "react"; @@ -132,6 +133,27 @@ export default function CallbackPage() { // loopback/tunnel callback. Keep the full URL visible as a manual fallback // in case the opener cannot receive the cross-origin postMessage. queueStatusUpdate("manual"); + // Retorno automático ao OmniRoute — APENAS em máquina local (loopback) e com sucesso: + // o relay já foi enviado ao painel na mesma origem; tenta fechar (popup) e, se continuar + // aberto (era a MESMA aba), volta ao painel. Caso remoto/túnel permanece no manual. + const host = window.location.hostname; + const isLoopback = host === "localhost" || /^127(?:\.\d{1,3}){3}$/.test(host); + if (code && !error && isLoopback) { + setTimeout(() => { + try { + window.close(); + } catch { + /* popup pode não fechar por política do navegador */ + } + setTimeout(() => { + try { + window.location.replace("/dashboard/providers"); + } catch { + /* mantém o estado manual como fallback */ + } + }, 1000); + }, 2500); + } } } else { // No code/error in URL or all send methods failed — show URL for manual copy. @@ -179,6 +201,16 @@ export default function CallbackPage() {
{currentUrl}
+ {/* Login abriu na MESMA aba (sem popup/opener): caminho de volta ao OmniRoute. */} + + + Voltar ao OmniRoute + )}
diff --git a/src/lib/buzzConsumer.ts b/src/lib/buzzConsumer.ts new file mode 100644 index 00000000000..344ce3f24d8 --- /dev/null +++ b/src/lib/buzzConsumer.ts @@ -0,0 +1,55 @@ +/** + * Buzz Consumer — liga o relay ao inbox do OmniRoute (fecha o gap "sem consumidor de produção"). + * + * Assina o relay (flag ON) e roteia cada evento recebido para `receiveInbox` — que apenas + * DEDUPLICA e ARMAZENA. NÃO há caminho daqui para qualquer efeito no OmniRoute: uma chave/assinatura + * Nostr nunca autoriza uma ação (as decisões ficam com o Policy Engine, fora deste módulo). Inbound + * é só sinal/colaboração para humanos verem no painel. Opt-in e best-effort. + */ +import type { BuzzSubscriptionFilter } from "@omniroute/open-sse/buzz-bridge/index.ts"; + +import { getBuzzAdapter } from "./buzzService"; +import { receiveInbox } from "./db/buzzBridge"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +export interface InboxSubscription { + started: boolean; + /** Encerra a assinatura e fecha a conexão. */ + stop: () => Promise; +} + +/** + * Inicia a assinatura do inbox. Retorna { started:false } quando a flag BUZZ_HUB_ENABLED está OFF + * (nada conecta). Cada evento verificado é persistido via receiveInbox (dedup); erros são engolidos + * para não derrubar o processo. `correlationId` de entrada = id do evento (rastreável, sem efeito). + */ +export async function startBuzzInboxSubscription( + filter: BuzzSubscriptionFilter = { kinds: [1] }, + tenantId?: string +): Promise { + if (!isFeatureFlagEnabled("BUZZ_HUB_ENABLED")) { + return { started: false, stop: async () => {} }; + } + const adapter = getBuzzAdapter(); + if (!adapter.enabled) return { started: false, stop: async () => {} }; + + await adapter.connect(); + await adapter.subscribe(filter, (event) => { + try { + // Storage-only: dedup + persiste. NUNCA dispara efeito (Nostr não autoriza). + receiveInbox(event, event.id, tenantId); + } catch { + /* best-effort: um evento malformado não derruba a assinatura */ + } + }); + return { + started: true, + stop: async () => { + try { + await adapter.close(); + } catch { + /* ignore */ + } + }, + }; +} diff --git a/src/lib/buzzProducer.ts b/src/lib/buzzProducer.ts new file mode 100644 index 00000000000..d77e9aea266 --- /dev/null +++ b/src/lib/buzzProducer.ts @@ -0,0 +1,56 @@ +/** + * Buzz Producer — liga o Loop Engine ao outbox do Buzz (fecha o gap "sem produtor de produção"). + * + * Quando um run do Loop precisa de um humano (awaiting_approval) ou escala (escalated), enfileira + * uma notificação DURÁVEL no outbox do Buzz. Idempotente (id derivado do run+status+seq), aditivo e + * best-effort: uma falha do Buzz NUNCA quebra o Loop. Nada conecta aqui — a publicação real fica no + * flush (gated por BUZZ_HUB_ENABLED). A chave Nostr nunca autoriza ação; isto é só sinalização. + */ +import { getPublicKey, type BuzzEvent } from "@omniroute/open-sse/buzz-bridge/index.ts"; + +import { enqueueOutbox } from "./db/buzzBridge"; +import { getOrCreateAgentSecretKey } from "./buzzService"; + +export interface LoopNotice { + runId: string; + status: string; + pattern: string; + sequenceNumber: number; + tenantId?: string; +} + +/** + * Enfileira uma notificação do Loop no outbox (report-only). Retorna true se enfileirou. + * Best-effort: engole erros (ex.: tabelas ausentes num DB mínimo) para não afetar o run. + */ +export function notifyLoopEvent(notice: LoopNotice): boolean { + try { + const event: BuzzEvent = { + // id estável → dedup no enqueue (mesma transição não vira dois eventos). + id: `loop:${notice.runId}:${notice.status}:${notice.sequenceNumber}`, + pubkey: getPublicKey(getOrCreateAgentSecretKey()), + kind: 1, + createdAt: 0, // o enqueue carimba um createdAt estável + tags: [ + ["t", "loop"], + ["run", notice.runId], + ["status", notice.status], + ], + content: `Loop "${notice.pattern}" (${notice.runId}) → ${notice.status}`, + }; + enqueueOutbox({ + event, + correlationId: `loop:${notice.runId}`, + runId: notice.runId, + tenantId: notice.tenantId, + }); + return true; + } catch { + return false; + } +} + +/** Estados do Loop que merecem um aviso ao humano (aprovação/handoff). */ +export function loopStatusNeedsHuman(status: string): boolean { + return status === "awaiting_approval" || status === "escalated"; +} diff --git a/src/lib/buzzService.ts b/src/lib/buzzService.ts new file mode 100644 index 00000000000..33299528f51 --- /dev/null +++ b/src/lib/buzzService.ts @@ -0,0 +1,145 @@ +/** + * Buzz Service — integra o Buzz ao OmniRoute (config + adaptador + flush do outbox). + * + * - Identidade Nostr do agente: gerada uma vez e PERSISTIDA (key_value namespace 'buzz'), + * estável entre reinícios. + * - Config do relay: BUZZ_RELAY_URL (default ws://localhost:3000). Configurável no painel único. + * - flushBuzzOutbox: publica as entradas pendentes do outbox no relay real (idempotente). + * + * Gated por BUZZ_HUB_ENABLED (OFF por padrão). Sem flag/relay, tudo fica inerte e durável no DB. + */ +import { + generateSecretKey, + getPublicKey, + resolveBuzzAdapter, + type BuzzAdapter, + type WebSocketBuzzConfig, +} from "@omniroute/open-sse/buzz-bridge/index.ts"; + +import { + buzzCounts, + markOutbox, + pendingOutbox, + requeueFailedOutbox, + type BuzzCounts, +} from "./db/buzzBridge"; +import { getDbInstance } from "./db/core"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +/** Chave secreta Nostr do agente OmniRoute, persistida e estável. */ +export function getOrCreateAgentSecretKey(): string { + const db = getDbInstance(); + const row = db + .prepare("SELECT value FROM key_value WHERE namespace = 'buzz' AND key = 'agent_sk'") + .get() as { value: string } | undefined; + if (row?.value) return row.value; + const sk = generateSecretKey(); + db.prepare( + "INSERT OR REPLACE INTO key_value (namespace, key, value) VALUES ('buzz', 'agent_sk', ?)" + ).run(sk); + return sk; +} + +/** + * URL do relay, com precedência: override do painel (key_value) → env BUZZ_RELAY_URL → default. + * Assim o PAINEL ÚNICO configura tudo, sem perder a opção de fixar por ambiente. + */ +export function getBuzzRelayUrl(): string { + const db = getDbInstance(); + const row = db + .prepare("SELECT value FROM key_value WHERE namespace = 'buzz' AND key = 'relay_url'") + .get() as { value: string } | undefined; + return row?.value?.trim() || process.env.BUZZ_RELAY_URL?.trim() || "ws://localhost:3000"; +} + +/** Persiste a URL do relay definida no painel. String vazia remove o override (volta ao env/default). */ +export function setBuzzRelayUrl(url: string): string { + const db = getDbInstance(); + const trimmed = url.trim(); + if (!trimmed) { + db.prepare("DELETE FROM key_value WHERE namespace = 'buzz' AND key = 'relay_url'").run(); + } else { + db.prepare( + "INSERT OR REPLACE INTO key_value (namespace, key, value) VALUES ('buzz', 'relay_url', ?)" + ).run(trimmed); + } + return getBuzzRelayUrl(); +} + +/** Config do relay (URL do painel/env + chave persistida). */ +export function getBuzzConfig(): WebSocketBuzzConfig { + return { + relayUrl: getBuzzRelayUrl(), + secretKeyHex: getOrCreateAgentSecretKey(), + }; +} + +/** Adaptador conforme a flag: real (WebSocket) quando BUZZ_HUB_ENABLED, senão inerte. */ +export function getBuzzAdapter(): BuzzAdapter { + const enabled = isFeatureFlagEnabled("BUZZ_HUB_ENABLED"); + return resolveBuzzAdapter(enabled, enabled ? getBuzzConfig() : undefined); +} + +export interface BuzzStatus { + enabled: boolean; + relayUrl: string; + /** Chave PÚBLICA Nostr do agente (nunca a secreta). Identidade estável do OmniRoute no relay. */ + agentPubkey: string; + counts: BuzzCounts; +} + +/** + * Estado do Buzz para o painel único: flag, URL do relay, pubkey do agente e contagens do + * outbox/inbox. NUNCA expõe a chave secreta. Não conecta ao relay (leitura local, barata). + */ +export function getBuzzStatus(): BuzzStatus { + return { + enabled: isFeatureFlagEnabled("BUZZ_HUB_ENABLED"), + relayUrl: getBuzzRelayUrl(), + agentPubkey: getPublicKey(getOrCreateAgentSecretKey()), + counts: buzzCounts(), + }; +} + +export interface FlushResult { + published: number; + failed: number; + skipped: boolean; +} + +/** + * Publica as entradas pendentes do outbox no relay. Idempotente (dedup do relay por id de + * evento). Retorna a contagem. Skipped quando a flag está OFF ou não há relay/pendentes. + */ +export async function flushBuzzOutbox(limit = 50): Promise { + if (!isFeatureFlagEnabled("BUZZ_HUB_ENABLED")) { + return { published: 0, failed: 0, skipped: true }; + } + // Reenfileira falhas transitórias (sob o teto de tentativas) antes de coletar as pendentes, para + // que uma queda passageira do relay não estrangule a mensagem em 'failed' para sempre. + requeueFailedOutbox(); + const pending = pendingOutbox(limit); + if (pending.length === 0) return { published: 0, failed: 0, skipped: false }; + + const adapter = resolveBuzzAdapter(true, getBuzzConfig()); + if (!adapter.enabled) return { published: 0, failed: 0, skipped: true }; + + await adapter.connect(); + let published = 0; + let failed = 0; + try { + for (const entry of pending) { + const ok = await adapter.publish(entry); + if (ok) { + markOutbox(entry.id, "published"); + published++; + } else { + markOutbox(entry.id, "failed"); + failed++; + } + } + } finally { + await adapter.close(); + } + return { published, failed, skipped: false }; +} diff --git a/src/lib/db/buzzBridge.ts b/src/lib/db/buzzBridge.ts new file mode 100644 index 00000000000..2fe494bd551 --- /dev/null +++ b/src/lib/db/buzzBridge.ts @@ -0,0 +1,180 @@ +/** + * Repositório do Buzz Bridge — persistência do outbox/inbox (tabelas da migração 174). + * + * Torna a ponte funcional mesmo com o relay AUSENTE: os eventos de saída ficam duráveis no + * outbox até haver relay + flag ON; os de entrada são deduplicados no inbox. Idempotente. + */ +import type { BuzzEvent, InboxEntry, OutboxEntry } from "@omniroute/open-sse/buzz-bridge/index.ts"; + +import { getDbInstance } from "./core"; +import { DEFAULT_TENANT } from "./loopEngine"; + +interface OutboxRow { + id: string; + correlation_id: string; + sequence_number: number; + task_id: string | null; + run_id: string | null; + event_json: string; + status: string; + attempts: number; +} + +/** Enfileira um evento de saída (idempotente por event.id). Retorna a entrada. Escopado por tenant. */ +export function enqueueOutbox(params: { + event: BuzzEvent; + correlationId: string; + taskId?: string; + runId?: string; + tenantId?: string; +}): OutboxEntry { + const db = getDbInstance(); + const tenantId = params.tenantId ?? DEFAULT_TENANT; + // Carimba createdAt UMA vez no enqueue (se ausente) e persiste, para a publicação re-assinar + // sempre com o MESMO timestamp → mesmo id Nostr → dedup real do relay entre re-tentativas. + const event: BuzzEvent = { + ...params.event, + createdAt: params.event.createdAt || Math.floor(Date.now() / 1000), + }; + const nextSeq = + ( + db + .prepare( + "SELECT COALESCE(MAX(sequence_number),0) AS m FROM buzz_outbox WHERE tenant_id = ?" + ) + .get(tenantId) as { m: number } + ).m + 1; + db.prepare( + `INSERT INTO buzz_outbox (id, tenant_id, correlation_id, sequence_number, task_id, run_id, event_json, status, attempts) + VALUES (@id, @tenant_id, @correlation_id, @sequence_number, @task_id, @run_id, @event_json, 'pending', 0) + ON CONFLICT(id) DO NOTHING` + ).run({ + id: event.id, + tenant_id: tenantId, + correlation_id: params.correlationId, + sequence_number: nextSeq, + task_id: params.taskId ?? null, + run_id: params.runId ?? null, + event_json: JSON.stringify(event), + }); + const row = db + .prepare("SELECT * FROM buzz_outbox WHERE id = ?") + .get(params.event.id) as OutboxRow; + return { + id: row.id, + correlationId: row.correlation_id, + sequenceNumber: row.sequence_number, + taskId: row.task_id ?? undefined, + runId: row.run_id ?? undefined, + event: JSON.parse(row.event_json), + status: row.status as OutboxEntry["status"], + attempts: row.attempts, + }; +} + +/** Entradas pendentes do outbox do tenant, em ordem de sequência (entrega ordenada). */ +export function pendingOutbox(limit = 100, tenantId: string = DEFAULT_TENANT): OutboxEntry[] { + const rows = getDbInstance() + .prepare( + "SELECT * FROM buzz_outbox WHERE tenant_id = ? AND status = 'pending' ORDER BY sequence_number LIMIT ?" + ) + .all(tenantId, limit) as OutboxRow[]; + return rows.map((row) => ({ + id: row.id, + correlationId: row.correlation_id, + sequenceNumber: row.sequence_number, + taskId: row.task_id ?? undefined, + runId: row.run_id ?? undefined, + event: JSON.parse(row.event_json), + status: row.status as OutboxEntry["status"], + attempts: row.attempts, + })); +} + +export function markOutbox(id: string, status: "published" | "failed"): void { + const db = getDbInstance(); + if (status === "failed") { + db.prepare("UPDATE buzz_outbox SET status='failed', attempts=attempts+1 WHERE id=?").run(id); + } else { + db.prepare("UPDATE buzz_outbox SET status='published' WHERE id=?").run(id); + } +} + +/** Teto de re-tentativas de publicação antes de desistir (evita loop infinito num relay quebrado). */ +export const MAX_OUTBOX_ATTEMPTS = 5; + +/** + * Reenfileira entradas 'failed' que ainda estão sob o teto de tentativas (failed -> pending), para + * o próximo flush tentar publicar de novo. Falhas costumam ser transitórias (relay fora do ar, + * corrida com o AUTH do NIP-42); sem isto a entrada ficaria presa em 'failed' para sempre. As que + * estouraram o teto permanecem 'failed' (não voltam). Retorna quantas foram reenfileiradas. + * Escopado por tenant. Idempotente por id (o dedup do relay cobre uma eventual republicação dupla). + */ +export function requeueFailedOutbox( + tenantId: string = DEFAULT_TENANT, + maxAttempts: number = MAX_OUTBOX_ATTEMPTS +): number { + const res = getDbInstance() + .prepare( + "UPDATE buzz_outbox SET status='pending' WHERE tenant_id = ? AND status='failed' AND attempts < ?" + ) + .run(tenantId, maxAttempts); + return res.changes; +} + +export interface BuzzCounts { + outboxPending: number; + outboxPublished: number; + outboxFailed: number; + inboxReceived: number; +} + +/** Contagens do outbox/inbox do tenant para o painel único. Zero se as tabelas não existem. */ +export function buzzCounts(tenantId: string = DEFAULT_TENANT): BuzzCounts { + const db = getDbInstance(); + const count = (sql: string): number => { + try { + return (db.prepare(sql).get(tenantId) as { n: number }).n; + } catch { + return 0; // tabela ausente num DB mínimo — reporta 0 em vez de quebrar o painel + } + }; + return { + outboxPending: count( + "SELECT COUNT(*) AS n FROM buzz_outbox WHERE tenant_id = ? AND status = 'pending'" + ), + outboxPublished: count( + "SELECT COUNT(*) AS n FROM buzz_outbox WHERE tenant_id = ? AND status = 'published'" + ), + outboxFailed: count( + "SELECT COUNT(*) AS n FROM buzz_outbox WHERE tenant_id = ? AND status = 'failed'" + ), + inboxReceived: count("SELECT COUNT(*) AS n FROM buzz_inbox WHERE tenant_id = ?"), + }; +} + +/** + * Registra um evento recebido (dedup por event.id). Retorna null se já visto — garante + * processamento no máximo uma vez, mesmo com reentrega do relay. Escopado por tenant. + */ +export function receiveInbox( + event: BuzzEvent, + correlationId: string, + tenantId: string = DEFAULT_TENANT +): InboxEntry | null { + const db = getDbInstance(); + const nextSeq = + ( + db + .prepare("SELECT COALESCE(MAX(sequence_number),0) AS m FROM buzz_inbox WHERE tenant_id = ?") + .get(tenantId) as { m: number } + ).m + 1; + const res = db + .prepare( + `INSERT INTO buzz_inbox (event_id, tenant_id, correlation_id, sequence_number, event_json, status) + VALUES (?, ?, ?, ?, ?, 'received') ON CONFLICT(event_id) DO NOTHING` + ) + .run(event.id, tenantId, correlationId, nextSeq, JSON.stringify(event)); + if (res.changes === 0) return null; // ja visto + return { event, correlationId, sequenceNumber: nextSeq, status: "received" }; +} diff --git a/src/lib/db/loopEngine.ts b/src/lib/db/loopEngine.ts new file mode 100644 index 00000000000..630449c2659 --- /dev/null +++ b/src/lib/db/loopEngine.ts @@ -0,0 +1,131 @@ +/** + * Repositório do Loop Engine — persistência dos runs/steps (tabelas da migração 174). + * + * Torna o núcleo puro (`open-sse/loop-engine`) FUNCIONAL: grava e recarrega o estado no DB + * do OmniRoute (fonte de verdade). Report-only; nenhum efeito externo aqui. + */ +import type { LoopRun, LoopStep } from "@omniroute/open-sse/loop-engine/index.ts"; + +import { getDbInstance } from "./core"; + +/** Tenant único por ora (CLAUDE.md §5.1). Toda leitura/escrita é escopada por ele. */ +export const DEFAULT_TENANT = "default"; + +interface LoopRunRow { + id: string; + pattern: string; + phase: string; + status: string; + budget_json: string; + usage_json: string; + correlation_id: string; + task_id: string | null; + sequence_number: number; +} + +interface LoopStepRow { + id: string; + run_id: string; + idx: number; + title: string; + proposed_effect_json: string | null; + status: string; +} + +function rowToRun(row: LoopRunRow, steps: LoopStep[]): LoopRun { + return { + id: row.id, + pattern: row.pattern, + phase: row.phase as LoopRun["phase"], + status: row.status as LoopRun["status"], + budget: JSON.parse(row.budget_json), + usage: JSON.parse(row.usage_json), + correlationId: row.correlation_id, + taskId: row.task_id ?? undefined, + sequenceNumber: row.sequence_number, + steps, + }; +} + +/** Persiste (upsert) um run e todas as suas etapas, de forma transacional. Escopado por tenant. */ +export function saveLoopRun(run: LoopRun, tenantId: string = DEFAULT_TENANT): void { + const db = getDbInstance(); + const tx = db.transaction((r: LoopRun) => { + db.prepare( + `INSERT INTO loop_runs (id, tenant_id, pattern, phase, status, budget_json, usage_json, correlation_id, task_id, sequence_number, updated_at) + VALUES (@id, @tenant_id, @pattern, @phase, @status, @budget_json, @usage_json, @correlation_id, @task_id, @sequence_number, datetime('now')) + ON CONFLICT(id) DO UPDATE SET + phase=excluded.phase, status=excluded.status, usage_json=excluded.usage_json, + sequence_number=excluded.sequence_number, updated_at=datetime('now')` + ).run({ + id: r.id, + tenant_id: tenantId, + pattern: r.pattern, + phase: r.phase, + status: r.status, + budget_json: JSON.stringify(r.budget), + usage_json: JSON.stringify(r.usage), + correlation_id: r.correlationId, + task_id: r.taskId ?? null, + sequence_number: r.sequenceNumber, + }); + const upStep = db.prepare( + `INSERT INTO loop_steps (id, tenant_id, run_id, idx, title, proposed_effect_json, status) + VALUES (@id, @tenant_id, @run_id, @idx, @title, @proposed_effect_json, @status) + ON CONFLICT(id) DO UPDATE SET status=excluded.status` + ); + for (const s of r.steps) { + upStep.run({ + id: s.id, + tenant_id: tenantId, + run_id: r.id, + idx: s.index, + title: s.title, + proposed_effect_json: s.proposedEffect ? JSON.stringify(s.proposedEffect) : null, + status: s.status, + }); + } + }); + tx(run); +} + +/** Carrega um run com suas etapas, ou null se não existir NESTE tenant (isolamento). */ +export function getLoopRun(id: string, tenantId: string = DEFAULT_TENANT): LoopRun | null { + const db = getDbInstance(); + const row = db + .prepare("SELECT * FROM loop_runs WHERE id = ? AND tenant_id = ?") + .get(id, tenantId) as LoopRunRow | undefined; + if (!row) return null; + const stepRows = db + .prepare("SELECT * FROM loop_steps WHERE run_id = ? AND tenant_id = ? ORDER BY idx") + .all(id, tenantId) as LoopStepRow[]; + const steps: LoopStep[] = stepRows.map((sr) => ({ + id: sr.id, + runId: sr.run_id, + index: sr.idx, + title: sr.title, + proposedEffect: sr.proposed_effect_json ? JSON.parse(sr.proposed_effect_json) : undefined, + status: sr.status as LoopStep["status"], + })); + return rowToRun(row, steps); +} + +/** Lista runs do tenant (opcionalmente por status), mais recentes primeiro. */ +export function listLoopRuns( + status?: LoopRun["status"], + tenantId: string = DEFAULT_TENANT +): LoopRun[] { + const db = getDbInstance(); + const rows = ( + status + ? db + .prepare( + "SELECT * FROM loop_runs WHERE tenant_id = ? AND status = ? ORDER BY updated_at DESC" + ) + .all(tenantId, status) + : db + .prepare("SELECT * FROM loop_runs WHERE tenant_id = ? ORDER BY updated_at DESC") + .all(tenantId) + ) as LoopRunRow[]; + return rows.map((row) => getLoopRun(row.id, tenantId)!).filter(Boolean); +} diff --git a/src/lib/db/migrations/174_loop_engine_and_buzz_bridge.sql b/src/lib/db/migrations/174_loop_engine_and_buzz_bridge.sql new file mode 100644 index 00000000000..09b3e52679f --- /dev/null +++ b/src/lib/db/migrations/174_loop_engine_and_buzz_bridge.sql @@ -0,0 +1,65 @@ +-- 174_loop_engine_and_buzz_bridge.sql +-- +-- Estado durável dos módulos agentic da Fase 1 (Loop Engine + Buzz Bridge). Aditiva, +-- idempotente (IF NOT EXISTS) e não-destrutiva: só cria tabelas/índices novos, não toca +-- nada existente. O OmniRoute permanece a fonte de verdade de runs/steps/aprovações; a +-- ponte Buzz (outbox/inbox) é idempotente por design. Ambos os módulos ficam atrás das +-- flags LOOP_ENGINE_ENABLED / BUZZ_HUB_ENABLED (OFF por padrão). + +-- ── Loop Engine ──────────────────────────────────────────── +-- Multi-tenant desde o nascimento (CLAUDE.md §5.1): toda tabela nasce com tenant_id + índice. +-- Valor único por ora ('default'); pronto para escopo por tenant sem migração destrutiva depois. +CREATE TABLE IF NOT EXISTS loop_runs ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL DEFAULT 'default', + pattern TEXT NOT NULL, + phase TEXT NOT NULL, + status TEXT NOT NULL, + budget_json TEXT NOT NULL, + usage_json TEXT NOT NULL, + correlation_id TEXT NOT NULL, + task_id TEXT, + sequence_number INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_at TEXT NOT NULL DEFAULT (datetime('now')) +); +CREATE INDEX IF NOT EXISTS idx_loop_runs_status ON loop_runs (tenant_id, status); +CREATE INDEX IF NOT EXISTS idx_loop_runs_correlation ON loop_runs (correlation_id); + +CREATE TABLE IF NOT EXISTS loop_steps ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL DEFAULT 'default', + run_id TEXT NOT NULL, + idx INTEGER NOT NULL, + title TEXT NOT NULL, + proposed_effect_json TEXT, + status TEXT NOT NULL DEFAULT 'proposed', + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); +CREATE INDEX IF NOT EXISTS idx_loop_steps_run ON loop_steps (run_id, idx); + +-- ── Buzz Bridge (ponte idempotente OmniRoute ↔ relay Nostr) ── +CREATE TABLE IF NOT EXISTS buzz_outbox ( + id TEXT PRIMARY KEY, -- = event.id (dedup) + tenant_id TEXT NOT NULL DEFAULT 'default', + correlation_id TEXT NOT NULL, + sequence_number INTEGER NOT NULL, + task_id TEXT, + run_id TEXT, + event_json TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + attempts INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); +CREATE INDEX IF NOT EXISTS idx_buzz_outbox_status ON buzz_outbox (tenant_id, status, sequence_number); + +CREATE TABLE IF NOT EXISTS buzz_inbox ( + event_id TEXT PRIMARY KEY, -- dedup de entrada + tenant_id TEXT NOT NULL DEFAULT 'default', + correlation_id TEXT NOT NULL, + sequence_number INTEGER NOT NULL, + event_json TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'received', + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); +CREATE INDEX IF NOT EXISTS idx_buzz_inbox_status ON buzz_inbox (tenant_id, status); diff --git a/src/lib/loopRunner.ts b/src/lib/loopRunner.ts new file mode 100644 index 00000000000..5eeddc7f9ff --- /dev/null +++ b/src/lib/loopRunner.ts @@ -0,0 +1,108 @@ +/** + * Loop Runner — orquestra o Loop Engine de ponta a ponta (report-only, persistente). + * + * Liga o núcleo puro (`open-sse/loop-engine`) ao repositório DB (`src/lib/db/loopEngine`): + * inicia runs, avança um passo por vez (persistindo), para em `awaiting_approval` quando um + * efeito externo é proposto, e retoma após aprovação humana explícita. NUNCA executa efeito + * externo — o Policy Engine/aprovações do OmniRoute decidem; a execução real (quando aprovada) + * fica com os conectores do OmniRoute, fora deste módulo. + */ +import { + advance, + createLoopRun, + proposeStep, + type AdvanceInput, + type LoopBudget, + type LoopProposedEffect, + type LoopRun, +} from "@omniroute/open-sse/loop-engine/index.ts"; + +import { getLoopRun, listLoopRuns, saveLoopRun } from "./db/loopEngine"; +import { loopStatusNeedsHuman, notifyLoopEvent } from "./buzzProducer"; + +/** Inicia um run novo (report-only) e persiste. */ +export function startRun(params: { + pattern: string; + budget?: Partial; + correlationId?: string; + taskId?: string; +}): LoopRun { + const run = createLoopRun(params); + saveLoopRun(run); + return run; +} + +/** Anexa uma etapa proposta a um run existente e persiste. */ +export function addStep( + runId: string, + step: { title: string; proposedEffect?: LoopProposedEffect } +): LoopRun { + const run = getLoopRun(runId); + if (!run) throw new Error(`loop run não encontrado: ${runId}`); + proposeStep(run, step); + saveLoopRun(run); + return run; +} + +/** + * Avança UM passo do run e persiste. Report-only por padrão: um efeito não auto-aprovado + * para o run em `awaiting_approval` (nada é executado). Retorna o novo estado + nota. + */ +export function advanceRun( + runId: string, + input?: Partial +): { run: LoopRun; note: string } { + const current = getLoopRun(runId); + if (!current) throw new Error(`loop run não encontrado: ${runId}`); + const result = advance(current, { + consumed: input?.consumed, + verdict: input?.verdict, + policy: input?.policy ?? { reportOnly: true }, + }); + saveLoopRun(result.run); + // Produtor Buzz: ao ENTRAR num estado que exige humano (aprovação/handoff), enfileira um aviso + // durável no outbox. Só na TRANSIÇÃO (evita repetir a cada advance). Best-effort, report-only. + if (loopStatusNeedsHuman(result.run.status) && result.run.status !== current.status) { + notifyLoopEvent({ + runId: result.run.id, + status: result.run.status, + pattern: result.run.pattern, + sequenceNumber: result.run.sequenceNumber, + }); + } + return result; +} + +/** + * Aprova explicitamente uma etapa que estava aguardando aprovação, liberando o run para + * prosseguir. Só um humano/operador chama isto (via painel/endpoint autenticado). + */ +export function approveStep(runId: string, stepId: string): LoopRun { + const run = getLoopRun(runId); + if (!run) throw new Error(`loop run não encontrado: ${runId}`); + const step = run.steps.find((s) => s.id === stepId); + if (!step) throw new Error(`etapa não encontrada: ${stepId}`); + step.status = "approved"; + // libera o run que estava parado para aprovação + if (run.status === "awaiting_approval") run.status = "report_only"; + saveLoopRun(run); + return run; +} + +/** Rejeita uma etapa (marca como rejeitada); o run pode então escalar/abortar no próximo advance. */ +export function rejectStep(runId: string, stepId: string): LoopRun { + const run = getLoopRun(runId); + if (!run) throw new Error(`loop run não encontrado: ${runId}`); + const step = run.steps.find((s) => s.id === stepId); + if (!step) throw new Error(`etapa não encontrada: ${stepId}`); + step.status = "rejected"; + saveLoopRun(run); + return run; +} + +/** Runs que estão parados aguardando aprovação humana (para o painel destacar). */ +export function runsAwaitingApproval(): LoopRun[] { + return listLoopRuns("awaiting_approval"); +} + +export { getLoopRun, listLoopRuns }; diff --git a/src/lib/otel.ts b/src/lib/otel.ts new file mode 100644 index 00000000000..971c6e561e1 --- /dev/null +++ b/src/lib/otel.ts @@ -0,0 +1,50 @@ +/** + * OTel server singleton — exportador em memória (ring buffer) + helper de trace por requisição. + * + * Torna o core `open-sse/otel` USÁVEL no runtime: um exportador único guarda os últimos spans + * (sem conteúdo sensível — a allowlist do core garante) para inspeção via GET /api/otel/spans. + * Só coleta quando OTEL_TRACING_ENABLED está ON. Nada bloqueia a requisição. + */ +import { + InMemorySpanExporter, + endSpan, + startSpan, + type Span, +} from "@omniroute/open-sse/otel/index.ts"; + +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; + +const MAX_SPANS = 500; + +class RingExporter extends InMemorySpanExporter { + export(span: Span): void { + super.export(span); + if (this.spans.length > MAX_SPANS) this.spans.splice(0, this.spans.length - MAX_SPANS); + } +} + +// Singleton por processo. +const g = globalThis as unknown as { __omniOtel?: RingExporter }; +export const otelExporter: RingExporter = (g.__omniOtel ??= new RingExporter()); + +/** Spans recentes (mais novos por último). Cópia rasa para leitura segura. */ +export function recentSpans(limit = 100): Span[] { + return otelExporter.spans.slice(-limit); +} + +/** + * Envolve uma operação síncrona num span (só quando a flag está ON). Atributos passam pela + * allowlist do core (nada sensível). Retorna o resultado da operação intacto. + */ +export function traceSync(name: string, attributes: Record, fn: () => T): T { + if (!isFeatureFlagEnabled("OTEL_TRACING_ENABLED")) return fn(); + const span = startSpan(name, undefined, attributes); + try { + const out = fn(); + endSpan(span, otelExporter, { status: "ok" }); + return out; + } catch (e) { + endSpan(span, otelExporter, { status: "error", attributes: { "error.kind": "exception" } }); + throw e; + } +} diff --git a/src/lib/piiSanitizer.ts b/src/lib/piiSanitizer.ts index 33f750ef61d..be398669d1c 100644 --- a/src/lib/piiSanitizer.ts +++ b/src/lib/piiSanitizer.ts @@ -80,6 +80,21 @@ const PII_PATTERNS: PIIPattern[] = [ replacement: "[CNPJ_REDACTED]", severity: "high", }, + { + // CEP brasileiro NNNNN-NNN (exige hífen; boundary de 3 dígitos descarta ZIP+4 dos EUA). + name: "cep", + regex: /(?<=^|[^A-Za-z0-9])\d{5}-\d{3}(?=$|[^A-Za-z0-9])/g, + replacement: "[CEP_REDACTED]", + severity: "medium", + }, + { + // Chave PIX aleatória (UUID v4) redigida só com a pista "pix" por perto (evita nuke de UUIDs). + name: "pix_key", + regex: + /(?<=\bpix\b[^\n]{0,30})[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-4[0-9a-fA-F]{3}-[89abAB][0-9a-fA-F]{3}-[0-9a-fA-F]{12}/gi, + replacement: "[PIX_KEY_REDACTED]", + severity: "high", + }, { name: "ip_address", regex: /(?<=^|[^A-Za-z0-9])(?:\d{1,3}\.){3}\d{1,3}(?=$|[^A-Za-z0-9])/g, diff --git a/src/shared/components/OAuthModalPanels.tsx b/src/shared/components/OAuthModalPanels.tsx index a691a20a398..ef42cc3d0df 100644 --- a/src/shared/components/OAuthModalPanels.tsx +++ b/src/shared/components/OAuthModalPanels.tsx @@ -348,6 +348,18 @@ export function OAuthManualInputPanel({

{t("step1OpenUrl")}

+ {/* Abrir em NOVA ABA (gesto do usuário → o navegador não bloqueia e não abre por cima + do painel). O painel segue aberto e recebe o retorno do login. */} +