diff --git a/api/main.py b/api/main.py index 2b8b985..e91ecaa 100644 --- a/api/main.py +++ b/api/main.py @@ -45,6 +45,7 @@ from sql_console import router as sql_router from agent_ops import router as agent_ops_router, agent_dml_loop, etl_agent_loop from cdc_consumer import router as cdc_router, cdc_consumer_loop from movements import router as movements_router +from movements import MOVEMENT_BY_ID, trigger_and_watch from dataflow import router as dataflow_router from pii_catalog import router as pii_router from ssh_terminal import ssh_session @@ -1084,6 +1085,16 @@ async def decide_approval(approval_id: str, body: ApprovalDecision): ) if not item: return {"error": "not found"} + + # Approval-gated executor: run the underlying action once a supervisor approves it. + if item.get("status") == "approved": + payload = item.get("payload") or {} + if isinstance(payload, dict) and payload.get("executor") == "movement": + mid = payload.get("movement_id") + if mid in MOVEMENT_BY_ID: + add_feed("etl-guardian", f"Approval {approval_id} granted — executing movement '{mid}'", "info") + asyncio.create_task(trigger_and_watch(mid, payload.get("conf"))) + return {"ok": True, "approval": item} diff --git a/docker-compose.yml b/docker-compose.yml index adb5793..514e257 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -87,6 +87,8 @@ services: rag-api: build: ../atc-data-quality/rag-api restart: unless-stopped + env_file: + - atc.env environment: CHROMA_HOST: chromadb CHROMA_PORT: 8000 @@ -95,11 +97,14 @@ services: LLM_MODEL: gpt-4o LLM_API_KEY: sk-local RAG_DATA_DIR: /data + OPENMETADATA_URL: ${OPENMETADATA_URL:-http://10.0.21.47:8585} + COMMAND_CENTER_URL: http://api:3201 volumes: - rag_data:/data depends_on: - chromadb - docling-serve + - api dq-api: build: ../atc-data-quality/dq-api diff --git a/ui/src/components/features/KnowledgeChatView.tsx b/ui/src/components/features/KnowledgeChatView.tsx index 6b8d82d..af3ae7a 100644 --- a/ui/src/components/features/KnowledgeChatView.tsx +++ b/ui/src/components/features/KnowledgeChatView.tsx @@ -1,5 +1,5 @@ import { useCallback, useEffect, useRef, useState } from 'react' -import { BookOpen, FileText, Loader2, MessageSquare, RefreshCw, RotateCcw, Send, Upload } from 'lucide-react' +import { Bot, BookOpen, Database, FileText, Loader2, MessageSquare, RefreshCw, RotateCcw, Send, Upload, Wrench } from 'lucide-react' import { cn } from '../../lib/utils' import { subTabActive, subTabIdle } from '../../lib/tabActive' @@ -14,7 +14,8 @@ type StoredDoc = { bytes?: number } type Source = { source?: string; chunk?: number; preview?: string } -type ChatMsg = { role: 'user' | 'assistant'; content: string; sources?: Source[] } +type AgentStep = { tool: string; input?: Record } +type ChatMsg = { role: 'user' | 'assistant'; content: string; sources?: Source[]; steps?: AgentStep[] } type Health = { ok: boolean; chroma: boolean; docling: boolean; llm: boolean; embed_model?: string } @@ -34,6 +35,9 @@ export function KnowledgeChatView({ onGpuActivity }: Props = {}) { const [selectedDocId, setSelectedDocId] = useState(null) const [summarizing, setSummarizing] = useState(false) const [reindexing, setReindexing] = useState(null) + const [agentMode, setAgentMode] = useState(false) + const [syncing, setSyncing] = useState(false) + const [catalogDocs, setCatalogDocs] = useState(null) const bottomRef = useRef(null) const loadMeta = useCallback(async () => { @@ -57,9 +61,19 @@ export function KnowledgeChatView({ onGpuActivity }: Props = {}) { } }, []) + const loadCatalog = useCallback(async () => { + try { + const r = await fetch('/rag/catalog/status') + if (r.ok) setCatalogDocs((await r.json()).documents ?? 0) + } catch { + setCatalogDocs(null) + } + }, []) + useEffect(() => { loadMeta() - }, [loadMeta]) + loadCatalog() + }, [loadMeta, loadCatalog]) useEffect(() => { bottomRef.current?.scrollIntoView({ behavior: 'smooth' }) @@ -146,12 +160,95 @@ export function KnowledgeChatView({ onGpuActivity }: Props = {}) { } } + const onSyncCatalog = async () => { + setSyncing(true) + setError(null) + try { + const r = await fetch('/rag/catalog/sync', { method: 'POST' }) + const j = await r.json() + if (!r.ok || !j.ok) { + setError(j.error || 'Catalog sync failed') + return + } + const c = j.counts || {} + setMessages((m) => [...m, { + role: 'assistant', + content: `Catalog synced — ${j.documents} entries (${c.tables ?? 0} tables, ${c.pii_datasets ?? 0} PII datasets, ${c.flows ?? 0} lineage edges, ${c.movements ?? 0} movements). Agent mode can now answer catalog/PII/lineage questions.`, + }]) + loadCatalog() + } catch { + setError('Catalog sync request failed') + } finally { + setSyncing(false) + } + } + + const onSendAgent = async (msg: string) => { + setLoading(true) + const steps: AgentStep[] = [] + let answer = '' + let idx = -1 + setMessages((m) => { + idx = m.length + return [...m, { role: 'assistant', content: '', steps: [] }] + }) + const patch = () => + setMessages((m) => m.map((x, i) => (i === idx ? { ...x, content: answer, steps: [...steps] } : x))) + try { + const r = await fetch('/rag/agent', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ message: msg, max_steps: 5 }), + }) + if (!r.ok || !r.body) { + setError('Agent unavailable') + return + } + const reader = r.body.getReader() + const decoder = new TextDecoder() + let buf = '' + for (;;) { + const { done, value } = await reader.read() + if (done) break + buf += decoder.decode(value, { stream: true }) + const events = buf.split('\n\n') + buf = events.pop() || '' + for (const ev of events) { + const eMatch = ev.match(/^event: (.+)$/m) + const dMatch = ev.match(/^data: (.+)$/m) + if (!eMatch || !dMatch) continue + const type = eMatch[1].trim() + let data: Record = {} + try { data = JSON.parse(dMatch[1]) } catch { continue } + if (type === 'step') { + steps.push({ tool: String(data.tool), input: data.input as Record }) + patch() + } else if (type === 'token') { + answer += String(data.t || '') + patch() + } else if (type === 'error') { + setError(String(data.error || 'Agent error')) + } + } + } + } catch { + setError('Failed to reach agent') + } finally { + patch() + setLoading(false) + } + } + const onSend = async () => { const msg = input.trim() if (!msg || loading) return setInput('') setError(null) setMessages((m) => [...m, { role: 'user', content: msg }]) + if (agentMode) { + await onSendAgent(msg) + return + } setLoading(true) try { const r = await fetch('/rag/chat', { @@ -233,6 +330,31 @@ export function KnowledgeChatView({ onGpuActivity }: Props = {}) { +

Platform agent

+ +

+ Live tools: catalog, PII, CDC, movements + approval-gated triggers. +

+ +

Ingest documents