From 7739fd9d0866949c990d0fefafd7ceec9ffe0d71 Mon Sep 17 00:00:00 2001 From: Mourya Balabhadra Date: Wed, 2 Sep 2026 03:58:11 -0700 Subject: [PATCH 1/4] SCAL-336134: Add developer examples related to chat history for spotter mcp server MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Brings the `python-react-agent-simple-ui` example up to date with the Spotter 3 MCP toolset, and fixes the things that made it unreliable to run: embed auth, chart rendering after an answer expires, CSP-blocked iframes, and an MCP connection that failed intermittently and was slow to open. It also makes the agent faster, converts the client to TypeScript, and adds a README guide for integrating the ThoughtSpot MCP server into your own application. - `claude_agent_mcp_server_v2.py` → `claude_agent_with_spotter3_mcp_server.py` - `claude_agent_mcp_server_v2_with_chat_history.py` → `claude_agent_with_spotter3_mcp_server_and_chat_history.py` - README and `env.template` re-worded off `v1`/`v2` onto Spotter 3 terms, and the OpenAI / Azure OpenAI (v1) documentation removed. The client's `getAuthToken` returned one constant `VITE_TS_AUTH_TOKEN`. The SDK requires a *fresh* token per call: once that static token stops verifying, the SDK reports a duplicate token and the callback can never recover, because it hands back the same string. - New `GET /api/ts-token` on both servers, minting a short-lived token per request via `POST /api/rest/2.0/auth/token/full`. The server used the static `TS_AUTH_TOKEN` for its own MCP and REST calls, so every tool call failed with "Failed to validate connection" once that token expired, while the tool *list* still loaded. - `server_token()` mints the server's own token and caches it until shortly before expiry. Its lifetime is set by `TS_SERVER_TOKEN_VALIDITY_SEC` (default 3600). - Expiry is tracked with `time.time()`, not `time.monotonic()`: on macOS the monotonic clock pauses during sleep, which would keep an expired token in use after a laptop sleep. - A background task (`keep_server_token_fresh()`) mints the token at startup and renews it about 12 minutes before expiry (at half-life for short lifetimes). Requests holding a valid token never wait on a renewal in progress. Previously a request that found the token expired paid the whole mint, which took ~40s on a staging cluster. - A static `TS_AUTH_TOKEN` is now needed only when minting credentials are not set. The server exits with a clear message when neither is configured. A ThoughtSpot answer object lives ~8 hours. The chat history stored each answer's `iframe_url` and `answer_id` (which is really a `{session_id, gen_no}` pair), so reopening an older conversation rendered a row of dead embeds. - The server no longer persists those fields, and strips them on read too, so rows written before this change take the same path. - `GET /api/conversations/{id}` now returns each answer's `answer_index` plus the conversation's `analytical_session_id`; the client emits a resolver placeholder and the Visual Embed SDK resolves a live URL. **Requires the SDK change in PR 2.** - `reconcile_answers` asks `getConversation` how many answers each turn actually has, rather than trusting the stored copy. This also recovers answers from turns whose SSE stream was cut off mid-flight: the Agent finishes regardless, so the answer exists on ThoughtSpot's side even when the app recorded none of it. `TS_MCP_API_VERSION` now defaults to `latest`, and the update readers accept both the server's digested shape and the Agent's raw shape (`text-chunk` vs `text_chunk`, `metadata.type == "thinking"` vs `is_thinking`, …), so `&enable-raw-session-updates=true` can be turned on through `TS_MCP_URL` without a code change. Intermediate "thinking" answers are filtered out — one three-question chat produced eleven answer updates but a single settled one — which also keeps our answer ordering aligned with `getConversation`'s. A timed turn ("total sales in 2023") split into: MCP connection setup 9.7s, four Claude calls 7.6s, ThoughtSpot tool calls 51.6s. - **Shared MCP session.** Opening a session costs several round trips (discovery probe, `initialize`, `tools/list`), measured at 6-11s and previously paid on every chat message. One session is now shared across turns (`McpPool` / `McpTurn`). Measured: connection setup went from 8.3s on the first message to 0.0s on the second; the turn went from 66.9s to 49.2s. - The session is opened and closed by its own background task (anyio requires the task that opened a session to close it). - It is replaced on token rotation; the old session stays open until the turns using it finish. It is also replaced when it breaks. - A session replaced mid-turn stays open until that turn ends, so a parallel call still running on it is not cut off. - **Retries.** - A failed handshake is retried once with a freshly minted token. `initialize` intermittently returns a bare HTTP 500, which surfaced in the UI as `MCPError: Server returned an error response`. - A call that fails with `-32600 "Session terminated"` is retried once on a fresh session. The check matches the message as well as the code, because the MCP client also uses -32600 when a call may already have run. - **No spurious error after each answer.** The SSE stream no longer cancels the agent task after `done`, and a browser hang-up is treated as a cancel, not an error. In the history server, a failed SQLite save on an error path is logged rather than raised, so the browser always receives its `error` event. - Default model is now `claude-haiku-4-5` (overridable with `ANTHROPIC_MODEL`). The agent only orchestrates ThoughtSpot's tools; measured on one question, its four Claude calls took ~4s on Haiku 4.5, ~8s on Opus 5 and ~17s on Sonnet 5. - `model_request_options()` sends adaptive thinking only to models that support it; Haiku 4.5 runs without thinking. - The server-side refusal fallback is sent only to models with fallback targets. `/v1/models` reports none for Sonnet 5 or Haiku 4.5. - `[Timing]` logs in the chat-history server break each turn down into MCP session, Claude calls, tool calls and total. - **TypeScript.** `App.jsx` → `App.tsx` and `main.jsx` → `main.tsx`, with typed messages, answers and SSE events. Adds `tsconfig.json` (type-check only; Vite still builds), `vite-env.d.ts`, and `npm run typecheck`. - **Loading state for stored chats.** Opening a stored chat shows a spinner and marks the sidebar entry busy until it loads (2.5-5s in testing, while answers are reconciled against ThoughtSpot). - Stale responses are ignored if another chat is opened or a new one started meanwhile. - Sending is disabled while a chat loads, so a message can't go to the previous conversation. - **Status line** shows the Analytics Agent's steps ("Searching for Datasets") instead of its reasoning text streamed one word-sized fragment at a time. - **System theme.** `App.css` moved fully onto CSS custom properties with a `prefers-color-scheme: dark` block (no hardcoded colours left bypassing the tokens), and the embed itself gets matching dark `customizations` variables so the chart doesn't stay white inside a dark page. - **`AnswerFrame` removed.** Answer iframes are injected as markup so React owns only the wrapper and doesn't fight the renderer's `replaceWith()`. - **README: "Integrating the ThoughtSpot MCP server into your own application".** An eight-step guide, each step naming the function to copy: 1. ThoughtSpot and Anthropic prerequisites 2. Minting tokens on the server 3. Connecting to MCP 4. Giving the tools to the model 5. Streaming to the UI 6. Rendering charts with the Visual Embed SDK 7. Optional chat history 8. A production checklist - **README fixes.** - Ports: the frontend is `:8000` (strict, because it is the CSP-allowlisted origin) and the backend `:8001`. The `uvicorn` commands now pass `--port 8001`; previously they clashed with Vite on port 8000. - Token-minting environment variables. - Snippets updated to `App.tsx`. - How stored charts are replayed. - The shared MCP session. - New troubleshooting rows. - **`env.template`:** documents the model default and `TS_SERVER_TOKEN_VALIDITY_SEC`, and says when the static token is required. - Dependency bumps: `anthropic>=1.2.0,<2`, `mcp>=2.1.1,<3`, `httpx` → `httpx2`, FastAPI/uvicorn; `@thoughtspot/visual-embed-sdk` `1.45.3-mcp.2` → `^1.52.1`; client dev dependencies `typescript`, `tslib`, `@types/react`, `@types/react-dom`. - `.gitignore`: local `*.db` / `-wal` / `-shm` chat-history files. - Live chats through both servers: answers stream, embeds render, follow-up turns reuse the same analytical session. Timed turns back to back confirm the second one skips MCP connection setup. - `/api/ts-token` verified minting: `minted: true`, a different token per call, and the minted token authenticates against `/callosum/v1/session/isactive`. - `/api/conversations` list/open/delete exercised against the stored database; the loading indicator and disabled input verified in the browser. - Live MCP session tests: - Ending the session out-of-band triggers one reconnect, and the call succeeds. - On token rotation, the old session stays open for its turn, then closes. - Hanging up mid-turn and mid-handshake leaves no traceback, and the next request works. - Unit tests with fakes: - A -32600 with a different message is not retried. - A parallel call on a replaced session finishes normally. - Background renewal: minting at startup; a request during a slow renewal returns the current token in 0ms; failed renewals keep the old token and retry; a short token life renews at half-life. - Startup with minting only (no static token) succeeds; with neither it fails with a clear message. - `py_compile` on both servers; `tsc` on the client (also passes with `--strict`); client `vite build` passes. - A separate review agent audited the diff; its findings are fixed in this PR. - `server/agent.py` (the OpenAI/Azure v1 backend) and the `openai` entry in `requirements.txt` are left in place but are no longer documented. Say if they should be removed. - The client renders model markdown with `rehypeRaw` without sanitizing it, and does not validate iframe `src` values. Both are listed in the README's production checklist. - On Haiku 4.5 the prompt cache does not engage: the ~3.6K-token prefix is under its 4,096-token minimum. --- .gitignore | 5 + mcp/python-react-agent-simple-ui/README.md | 503 ++++-- .../client/index.html | 2 +- .../client/package-lock.json | 72 +- .../client/package.json | 9 +- .../client/src/App.css | 198 +- .../client/src/App.jsx | 275 --- .../client/src/App.tsx | 669 +++++++ .../client/src/main.jsx | 10 - .../client/src/main.tsx | 13 + .../client/src/vite-env.d.ts | 9 + .../client/tsconfig.json | 26 + .../client/vite.config.js | 7 +- mcp/python-react-agent-simple-ui/env.template | 44 +- .../server/claude_agent_mcp_server_v2.py | 263 --- .../claude_agent_with_spotter3_mcp_server.py | 1104 ++++++++++++ ...th_spotter3_mcp_server_and_chat_history.py | 1599 +++++++++++++++++ .../server/requirements.txt | 10 +- 18 files changed, 4070 insertions(+), 748 deletions(-) delete mode 100644 mcp/python-react-agent-simple-ui/client/src/App.jsx create mode 100644 mcp/python-react-agent-simple-ui/client/src/App.tsx delete mode 100644 mcp/python-react-agent-simple-ui/client/src/main.jsx create mode 100644 mcp/python-react-agent-simple-ui/client/src/main.tsx create mode 100644 mcp/python-react-agent-simple-ui/client/src/vite-env.d.ts create mode 100644 mcp/python-react-agent-simple-ui/client/tsconfig.json delete mode 100644 mcp/python-react-agent-simple-ui/server/claude_agent_mcp_server_v2.py create mode 100644 mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server.py create mode 100644 mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server_and_chat_history.py diff --git a/.gitignore b/.gitignore index e5c9c28..6dadd5c 100644 --- a/.gitignore +++ b/.gitignore @@ -137,3 +137,8 @@ dist *.egg-info __pycache__ .env*.local + +# Local chat history databases (mcp/python-react-agent-simple-ui) +*.db +*.db-wal +*.db-shm diff --git a/mcp/python-react-agent-simple-ui/README.md b/mcp/python-react-agent-simple-ui/README.md index 1e4adab..d3a43eb 100644 --- a/mcp/python-react-agent-simple-ui/README.md +++ b/mcp/python-react-agent-simple-ui/README.md @@ -1,6 +1,6 @@ # Python Agent with Simple React UI -A full-stack example that pairs a **Python (FastAPI) agent** with a **React chat UI**. Supports two backend implementations: +A full-stack example that pairs a **Python (FastAPI) agent** with a **React chat UI**, running against the ThoughtSpot MCP server's **Spotter 3** toolset. Two backends are included: -| Backend | File | AI Provider | MCP Integration | -|-----------------|----------------------------------------|-----------------------|------------------------------------------------| -| **v1 (OpenAI)** | `server/agent.py` | Azure OpenAI / OpenAI | Server-side (OpenAI manages MCP) | -| **v2 (Claude)** | `server/claude_agent_mcp_server_v2.py` | Anthropic Claude | Client-side (FastAPI connects to MCP directly) | -| **v2 (OpenAI)** | _(coming soon)_ | Azure OpenAI / OpenAI | Client-side (FastAPI connects to MCP directly) | +| Backend | File | What it gives you | +|-------------------------|--------------------------------------------------------------------|------------------------------------------------------------| +| **Spotter 3** | `server/claude_agent_with_spotter3_mcp_server.py` | The agent, with history held in memory for the process life | +| **Spotter 3 + history** | `server/claude_agent_with_spotter3_mcp_server_and_chat_history.py` | The same, plus SQLite history you can list, reopen and delete | - +Both use Anthropic Claude with a client-side MCP loop — the FastAPI process connects to the MCP server directly. The backend streams responses to the frontend using Server-Sent Events (SSE), giving users a real-time chat experience while the agent queries ThoughtSpot for data insights and displays ThoughtSpot charts in an embed. @@ -32,37 +33,112 @@ The backend streams responses to the frontend using Server-Sent Events (SSE), gi --- -## MCP Server v2: Claude + Client-side MCP (Recommended) +## Claude + the Spotter 3 MCP Server -`claude_agent_mcp_server_v2.py` uses Anthropic's Claude API with a **client-side agentic loop** — the FastAPI process connects directly to the ThoughtSpot MCP server using custom HTTP headers (`Authorization` + `x-ts-host`). This approach is required because Anthropic's server-side MCP integration does not support custom headers. +`claude_agent_with_spotter3_mcp_server.py` uses Anthropic's Claude API with a **client-side agentic loop** — the FastAPI process connects directly to the [ThoughtSpot MCP server](https://github.com/thoughtspot/mcp-server) using custom HTTP headers (`Authorization` + `x-ts-host`). This is required because Anthropic's server-side MCP connector cannot send custom headers. -### Architecture (v2) +### Architecture ``` ┌──────────────┐ SSE stream ┌──────────────────────┐ MCP (streamable-http) ┌─────────────┐ │ React Chat │ ◄────────────► │ FastAPI + Claude │ ◄──────────────────────────────► │ ThoughtSpot │ -│ (Vite) │ /api/chat │ / OpenAI │ Authorization + x-ts-host │ MCP Server │ +│ (Vite) │ /api/chat │ │ Authorization + x-ts-host │ MCP Server │ └──────────────┘ └──────────────────────┘ └─────────────┘ - :5173 :8000 agent.thoughtspot.app + :8000 :8001 agent.thoughtspot.app ``` **Request flow:** 1. User sends a message from the React UI -2. FastAPI opens a new MCP session to `agent.thoughtspot.app` with auth headers +2. FastAPI borrows the shared MCP session to `agent.thoughtspot.app` - opening it with auth headers and fetching the tool list only on the first message, or when the session has to be replaced (see [Shared MCP session](#shared-mcp-session)) 3. Claude receives the user message + ThoughtSpot tool definitions -4. Claude calls ThoughtSpot tools as needed; FastAPI executes each call via the MCP session -5. The agentic loop continues until Claude stops calling tools -6. Text deltas and status events are streamed to the UI over SSE in real time +4. Claude calls ThoughtSpot tools; FastAPI executes each call over the MCP session +5. For `get_session_updates`, FastAPI polls the Analytics Agent to completion itself (see below) +6. Text deltas, agent progress, and rendered answers stream to the UI over SSE in real time -### Prerequisites (v2) +### Integrating the ThoughtSpot MCP server into your own application -- Python 3.10+ +This example is meant to be copied from. The steps below are the path from "ThoughtSpot cluster" to "charts in my app's chat", each pointing at the code to lift. + +#### 1. Prepare ThoughtSpot and Anthropic + +- **A trusted-authentication secret key**, so your backend can mint tokens: *Develop > Customizations > Security Settings > Trusted authentication*. (A username + password also works for a demo, but MFA-enabled instances reject it.) +- **Your app's origin allowlisted for embedding.** The charts render in a ThoughtSpot iframe, and the cluster's CSP `frame-ancestors` must include the origin your UI is served from - otherwise the browser blocks the embed. This example's dev server is pinned to `http://localhost:8000` for exactly that reason (`client/vite.config.js`). See the [security settings docs](https://developers.thoughtspot.com/docs/security-settings) for CSP and CORS allowlists. +- **An Anthropic API key** for the agent. + +#### 2. Mint ThoughtSpot tokens on your server + +The secret key never reaches the browser. The backend calls `POST /api/rest/2.0/auth/token/full` and uses the result two ways: + +- `server_token()` - a cached token for the agent's own MCP and REST calls. A background task (`keep_server_token_fresh()`) mints it at startup and renews it about 12 minutes before it expires, so no request waits on a mint - which can take tens of seconds on a loaded cluster. +- `GET /api/ts-token` - a **fresh** token per call for the browser's Visual Embed SDK. The SDK requires a new token from every `getAuthToken` call; handing it the same one twice triggers its "Duplicate token" alert once that token stops verifying. + +Copy `mint_token()`, `server_token()` and the `/api/ts-token` endpoint. + +> **Production:** this example mints every token for one service user (`TS_EMBED_USERNAME`). In your app, mint for the **authenticated end user** instead, so ThoughtSpot's own permissions and row-level security apply to every question they ask. + +#### 3. Connect to the MCP server from your backend + +- Endpoint: `https://agent.thoughtspot.app/token/mcp?api-version=...`, with headers `Authorization: Bearer ` and `x-ts-host: `. +- Connect **client-side** (your process holds the MCP session), because Anthropic's server-side MCP connector cannot send the custom `x-ts-host` header. +- **Pin `api-version`** to a release date for anything you depend on - see [MCP endpoint and API version](#mcp-endpoint-and-api-version). +- **Share one session across requests.** Opening a session took 6-11 s on the clusters this was tested against. `connect_mcp()`, `McpPool` and `McpTurn` open it once, reconnect when it breaks or the token rotates, and retry a failed handshake once. + +#### 4. Hand the tools to your LLM + +- Convert `list_tools()` to your model's tool format (`build_tools()`), run the tool loop (`agent_loop()`), and return all parallel tool results in one message. +- **Poll `get_session_updates` in your code, not in the model.** `send_session_message` returns immediately and the Analytics Agent answers asynchronously; `autopoll_session_updates()` polls to completion and gives the model one consolidated result. +- **Strip `iframe_url` from what the model sees** (`strip_rendered_answers()`) - your UI renders the chart, so the model only needs to summarize it. + +#### 5. Stream progress and answers to your UI + +Stream text, progress and answers as they arrive - an Analytics Agent answer typically takes tens of seconds. This example uses Server-Sent Events; the event contract is in [SSE events](#sse-events). + +#### 6. Render the charts with the Visual Embed SDK + +```ts +init({ + thoughtSpotHost: import.meta.env.VITE_TS_HOST, + authType: AuthType.TrustedAuthTokenCookieless, + getAuthToken: async () => { + const response = await fetch("/api/ts-token", { cache: "no-store" }); + if (!response.ok) throw new Error(`Token request failed: ${response.status}`); + return (await response.json()).token; + }, + autoLogin: true, // fetch a replacement before the current token expires +}); + +// Upgrades every `; +}; + +function App() { + const [messages, setMessages] = useState([]); + const [input, setInput] = useState(""); + const [isLoading, setIsLoading] = useState(false); + const [status, setStatus] = useState(""); + const [responseId, setResponseId] = useState(null); + // ThoughtSpot analytical session for the open conversation. Only a reopened + // chat needs it - live answers arrive with a URL already attached. + const [sessionId, setSessionId] = useState(null); + // Chat history. `historyAvailable` stays false against a backend without the + // /api/conversations endpoints, and the sidebar simply isn't rendered. + const [conversations, setConversations] = useState([]); + const [historyAvailable, setHistoryAvailable] = useState(false); + // The stored chat being fetched. Opening one waits on ThoughtSpot (the server + // reconciles its answers against the live conversation), so it can take seconds. + const [openingId, setOpeningId] = useState(null); + // Bumped on every open / new chat, so only the latest request's response lands. + const openRequestRef = useRef(0); + const messagesEndRef = useRef(null); + const textareaRef = useRef(null); + + const refreshConversations = useCallback(async () => { + try { + const response = await fetch(HISTORY_URL); + if (!response.ok) throw new Error(String(response.status)); + const data = (await response.json()) as { + conversations?: ConversationSummary[]; + }; + setConversations(data.conversations || []); + setHistoryAvailable(true); + } catch { + setHistoryAvailable(false); + } + }, []); + + useEffect(() => { + refreshConversations(); + }, [refreshConversations]); + + // Reopen a stored conversation. Stored turns carry their answers, so the charts + // come back as embeds rather than as text - resolved live, since the stored + // answers deliberately carry no URL. See answerSrc. + const openConversation = useCallback( + async (id: string) => { + if (isLoading) return; + const request = ++openRequestRef.current; + setOpeningId(id); + setStatus(""); + try { + const response = await fetch(`${HISTORY_URL}/${id}`); + if (!response.ok) throw new Error(String(response.status)); + const data = (await response.json()) as ConversationDetail; + // Another chat was opened, or a new one started, while this was loading. + if (request !== openRequestRef.current) return; + setMessages( + (data.turns || []).map((turn) => ({ + role: turn.role, + content: turn.content, + answers: turn.answers || [], + })), + ); + setResponseId(data.id); + setSessionId(data.analytical_session_id || null); + setStatus(""); + } catch (error) { + if (request !== openRequestRef.current) return; + setStatus(`Could not open that conversation: ${errorMessage(error)}`); + } finally { + if (request === openRequestRef.current) setOpeningId(null); + } + }, + [isLoading], + ); + + const deleteConversation = useCallback( + async (id: string, event: MouseEvent) => { + event.stopPropagation(); + if (id === openingId) { + // Deleting the chat that is still loading: drop the load too. + openRequestRef.current++; + setOpeningId(null); + } + await fetch(`${HISTORY_URL}/${id}`, { method: "DELETE" }); + if (id === responseId) { + setMessages([]); + setResponseId(null); + } + refreshConversations(); + }, + [responseId, openingId, refreshConversations], + ); + + const scrollToBottom = useCallback(() => { + messagesEndRef.current?.scrollIntoView({ behavior: "smooth" }); + }, []); + + useEffect(() => { + scrollToBottom(); + }, [messages, status, scrollToBottom]); + + useEffect(() => { + if (textareaRef.current) { + textareaRef.current.style.height = "auto"; + textareaRef.current.style.height = + Math.min(textareaRef.current.scrollHeight, 150) + "px"; + } + }, [input]); + + const sendMessage = async () => { + // While a stored chat loads, a send would go to the previous conversation and + // its reply would land in the one being opened. + if (!input.trim() || isLoading || openingId) return; + + const userMessage = input.trim(); + setInput(""); + setIsLoading(true); + setStatus("Thinking..."); + + setMessages((prev) => [ + ...prev, + { role: "user", content: userMessage }, + { role: "assistant", content: "", answers: [] }, + ]); + + try { + const response = await fetch(API_URL, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + message: userMessage, + response_id: responseId, + }), + }); + + if (!response.ok) { + throw new Error(`Server error: ${response.status}`); + } + + if (!response.body) throw new Error("Server sent an empty response"); + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + + buffer += decoder.decode(value, { stream: true }); + const lines = buffer.split("\n"); + buffer = lines.pop() || ""; + + for (const line of lines) { + if (!line.startsWith("data: ")) continue; + + try { + const data = JSON.parse(line.slice(6)) as StreamEvent; + + switch (data.type) { + case "delta": + setStatus(""); + setMessages((prev) => { + const updated = [...prev]; + const last = updated[updated.length - 1]; + updated[updated.length - 1] = { + ...last, + content: last.content + data.text, + }; + return updated; + }); + break; + case "status": + setStatus(data.message); + break; + // An answer the Analytics Agent produced. The server streams these as + // soon as they arrive, so charts appear while the agent is still working. + case "answer": { + // Raw MCP updates have no `iframe_url`, so an answer is renderable + // as long as a src can be built from what it does carry. + const src = answerSrc(data, null); + if (!src) break; + setMessages((prev) => { + const updated = [...prev]; + const last = updated[updated.length - 1]; + const answers = last.answers || []; + if (answers.some((a) => answerSrc(a, null) === src)) return prev; + updated[updated.length - 1] = { + ...last, + answers: [...answers, data], + }; + return updated; + }); + break; + } + case "done": + setResponseId(data.response_id); + setStatus(""); + // The server has written the turn by the time it sends `done`. + refreshConversations(); + break; + case "error": + setMessages((prev) => { + const updated = [...prev]; + const last = updated[updated.length - 1]; + updated[updated.length - 1] = { + role: "assistant", + content: `Error: ${data.message}`, + answers: last.answers || [], + isError: true, + }; + return updated; + }); + break; + } + } catch { + /* skip malformed lines */ + } + } + } + } catch (error) { + setMessages((prev) => { + const updated = [...prev]; + if (updated.length > 0) { + updated[updated.length - 1] = { + role: "assistant", + content: `Failed to connect to server: ${errorMessage(error)}`, + isError: true, + }; + } + return updated; + }); + } finally { + setIsLoading(false); + setStatus(""); + } + }; + + const handleKeyDown = (e: KeyboardEvent) => { + if (e.key === "Enter" && !e.shiftKey) { + e.preventDefault(); + sendMessage(); + } + }; + + const startNewChat = () => { + // Drop any chat still loading, so it doesn't replace the new one when it lands. + openRequestRef.current++; + setOpeningId(null); + setMessages([]); + setResponseId(null); + setSessionId(null); + setStatus(""); + setInput(""); + }; + + return ( +
+ {historyAvailable && ( + + )} + +
+
+
+

ThoughtSpot Agent

+ {messages.length > 0 && ( + + )} +
+
+ +
+ {openingId ? ( +
+
+ ) : messages.length === 0 ? ( +
+
TS
+

Ask anything about your data

+

+ Powered by ThoughtSpot and Claude. Ask questions and get + insights from your connected data sources. +

+
+ ) : ( +
+ {messages.map((msg, i) => ( +
+
+ {msg.role === "user" ? "U" : "TS"} +
+
+ {msg.role !== "assistant" ? ( +

{msg.content}

+ ) : ( + <> + {(msg.answers || []).map((answer) => ( +
+ {answer.title && ( +
{answer.title}
+ )} +
+
+ ))} + {msg.content ? ( + // rehypeRaw so an " # important! - "Do not ask to create charts, as thoughtspot will already create interactive charts for you." - "Respond in an engaging markdown format, with html tags when needed." - "Keep the response short and to the point." - # "Use this datasource: cd252e5c-b552-49a8-821d-3eadaa049cca to answer all data questions." -) - -# ThoughtSpot MCP Tools (v2): -# To restrict which tools are accessible to the agent, set ALLOWED_TOOLS to a list of tool names. -# Set to None to allow all tools. -# -# The v2 MCP server uses an analytical session workflow: -# 1. check_connectivity - Test connectivity and authentication. No inputs. -# 2. create_analysis_session - Start a session. Optional: data_source_id. -# Returns: analytical_session_id. -# 3. send_session_message - Send a natural-language question to the session. -# Inputs: analytical_session_id, message, additional_context (optional). -# 4. get_session_updates - Poll for incremental updates. Inputs: analytical_session_id. -# Returns: session_updates (list), is_done (bool). -# Poll until is_done=True. Each update has type: text | text_chunk | answer. -# Answer updates include: answer_id, answer_title, answer_query, iframe_url. -# 5. create_dashboard - Create a dashboard from answer IDs. -# Inputs: title, answers (list of answer_ids), note_tile. -# Returns: link. -ALLOWED_TOOLS = None -# ALLOWED_TOOLS = ["check_connectivity", "create_analysis_session", "send_session_message", "get_session_updates", "create_dashboard"] - -# In-memory conversation store: conv_id -> full message history (including tool interactions) -conversations: dict[str, list] = {} - -# 1:1 mapping: conv_id -> analytical_session_id returned by create_analysis_session tool. -# Passed to Claude via system prompt so follow-up send_session_message / get_session_updates -# calls use the same ThoughtSpot analytical session. -analytical_sessions: dict[str, str] = {} - - -class ChatRequest(BaseModel): - message: str - response_id: str | None = None - - -def format_sse(data: dict) -> str: - return f"data: {json.dumps(data)}\n\n" - - -async def agent_loop(messages: list, queue: asyncio.Queue, conv_id: str) -> None: - """ - Client-side agentic loop. Connects to the ThoughtSpot MCP server directly - (with Authorization + x-ts-host headers), fetches tool definitions, then - runs the Claude tool-use loop until the model stops calling tools. - Puts SSE event dicts into queue for streaming to the frontend. - """ - try: - headers = dict(MCP_HEADERS) - - print(f"[MCP] Connecting to {MCP_URL}") - async with streamablehttp_client(MCP_URL, headers=headers) as (read, write, _): - async with ClientSession(read, write) as session: - print("[MCP] Initializing session...") - await session.initialize() - print("[MCP] Session initialized. Fetching tools...") - - # Fetch tool definitions from ThoughtSpot MCP server - tools_result = await session.list_tools() - print(f"[MCP] Got {len(tools_result.tools)} tools") - available_tools = tools_result.tools - # Optionally filter tools based on ALLOWED_TOOLS - if ALLOWED_TOOLS is not None: - available_tools = [t for t in available_tools if t.name in ALLOWED_TOOLS] - - # Convert MCP tool definitions to Anthropic format - anthropic_tools = [ - { - "name": t.name, - "description": t.description or "", - "input_schema": t.inputSchema, - } - for t in available_tools - ] - - current_messages = messages[:] - final_text_parts: list[str] = [] - - # Build system prompt, injecting analytical_session_id for follow-up turns - system = SYSTEM_PROMPT - existing_session_id = analytical_sessions.get(conv_id) - if existing_session_id: - system += ( - f"\n\nActive ThoughtSpot analytical session ID: {existing_session_id}. " - "Use this ID when calling send_session_message or get_session_updates " - "so follow-up questions continue in the same session." - ) - - while True: - async with claude_client.messages.stream( - model="claude-opus-4-6", - max_tokens=16000, - system=system, - messages=current_messages, - tools=anthropic_tools, - ) as stream: - async for event in stream: - t = getattr(event, "type", None) - if t == "content_block_start": - if getattr(event.content_block, "type", None) == "tool_use": - await queue.put({"type": "status", "message": "Querying ThoughtSpot..."}) - elif t == "content_block_delta": - delta = event.delta - if getattr(delta, "type", None) == "text_delta": - await queue.put({"type": "delta", "text": delta.text}) - final_text_parts.append(delta.text) - - final_message = await stream.get_final_message() - - if final_message.stop_reason != "tool_use": - break - - # Execute each tool call via MCP client (headers are set on the session) - tool_results = [] - for block in final_message.content: - if getattr(block, "type", None) == "tool_use": - try: - mcp_result = await session.call_tool(block.name, block.input) - print(f"[MCP] Tool {block.name} and input {block.input} returned: {mcp_result}") - result_text = " ".join( - getattr(c, "text", str(c)) for c in mcp_result.content - ) if mcp_result.content else "" - is_error = getattr(mcp_result, "isError", False) - - # Store analytical_session_id (1:1 with conv_id) so follow-up - # requests can reference the same ThoughtSpot session. - if block.name == "create_analysis_session" and not analytical_sessions.get(conv_id): - try: - sid = json.loads(result_text).get("analytical_session_id") - if sid: - analytical_sessions[conv_id] = sid - print(f"[MCP] Stored analytical_session_id for conv {conv_id}: {sid}") - except Exception: - pass - - except McpError as e: - print(f"[MCP] Tool {block.name} failed: {e}") - result_text = f"Tool call failed: {e}" - is_error = True - tool_results.append({ - "type": "tool_result", - "tool_use_id": block.id, - "content": result_text, - "is_error": is_error, - }) - - # Append assistant turn + tool results and continue the loop - current_messages = current_messages + [ - {"role": "assistant", "content": final_message.content}, - {"role": "user", "content": tool_results}, - ] - final_text_parts = [] # reset; next iteration may stream more text - - # Persist full conversation history (including tool interactions) so - # follow-up turns have complete context (e.g. analytical_session_id in prior results). - conversations[conv_id] = current_messages + [ - {"role": "assistant", "content": final_message.content} - ] - await queue.put({"type": "done", "response_id": conv_id}) - - except BaseException as e: - traceback.print_exc() - # Recursively unwrap ExceptionGroup to get the root cause - err = e - while hasattr(err, "exceptions") and getattr(err, "exceptions", None): - err = err.exceptions[0] - await queue.put({"type": "error", "message": f"{type(err).__name__}: {err}"}) - - -@app.post("/api/chat") -async def chat(request: ChatRequest): - conv_id = request.response_id or str(uuid.uuid4()) - print(f"[Chat] Received message for conv_id {conv_id}: {request.response_id}") - history = conversations.get(conv_id, []) - messages = history + [{"role": "user", "content": request.message}] - - queue: asyncio.Queue = asyncio.Queue() - asyncio.create_task(agent_loop(messages, queue, conv_id)) - - async def event_stream() -> AsyncGenerator[str, None]: - while True: - item = await queue.get() - yield format_sse(item) - if item.get("type") in ("done", "error"): - break - - return StreamingResponse(event_stream(), media_type="text/event-stream") - - -@app.get("/api/health") -async def health(): - return {"status": "ok"} diff --git a/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server.py b/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server.py new file mode 100644 index 0000000..cdecaeb --- /dev/null +++ b/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server.py @@ -0,0 +1,1104 @@ +""" +Python agent: Anthropic Claude + the ThoughtSpot MCP server (Spotter 3 toolset). + +FastAPI streams chat responses to the React frontend over Server-Sent Events. + +Why client-side MCP: the ThoughtSpot MCP server's static-token endpoint needs custom +HTTP headers (Authorization + x-ts-host). Anthropic's server-side MCP connector cannot +send those, so this process connects to the MCP server itself and executes tool calls. + +Two things this server does that a plain pass-through loop does not: + +1. `get_session_updates` polling happens here, not in the model. The ThoughtSpot + Analytics Agent answers asynchronously, so one `get_session_updates` call usually + returns `is_done: false` and an empty list. Letting the model poll costs a full + model round-trip per poll. Instead we poll until `is_done: true` and hand the model + one consolidated result, streaming the Agent's progress to the UI as it arrives. + +2. Answers are rendered by the client, not by the model. Each `answer` update becomes + an iframe marked `tsmcp=true`, which the React client mounts and the Visual Embed + SDK's `startAutoMCPFrameRenderer` upgrades into a real ThoughtSpot embed. We keep + the URL out of what the model sees - it is long, and the model does not need to + echo markup for a chart the UI has already drawn. + +MCP server: https://github.com/thoughtspot/mcp-server +""" + +import asyncio +import json +import os +import time +import traceback +import uuid +from collections.abc import AsyncGenerator, AsyncIterator, Callable +from contextlib import AsyncExitStack, asynccontextmanager +from pathlib import Path +from typing import Any + +import anthropic +import httpx2 +from dotenv import load_dotenv +from fastapi import FastAPI, HTTPException +from fastapi.middleware.cors import CORSMiddleware +from fastapi.responses import StreamingResponse +from mcp import Client, MCPError +from mcp.client.streamable_http import streamable_http_client +from pydantic import BaseModel + +load_dotenv(dotenv_path=Path(__file__).resolve().parent.parent / ".env") + +@asynccontextmanager +async def lifespan(_app: FastAPI) -> AsyncIterator[None]: + refresher = asyncio.create_task(keep_server_token_fresh()) if CAN_MINT_TOKENS else None + yield + if refresher: + refresher.cancel() + # Close the shared MCP session so the server side is not left holding it. + await mcp_pool.close() + + +app = FastAPI(title="ThoughtSpot Agent", lifespan=lifespan) + +app.add_middleware( + CORSMiddleware, + allow_origins=["*"], + allow_credentials=True, + allow_methods=["*"], + allow_headers=["*"], +) + +# ── Anthropic ─────────────────────────────────────────────────────────────────── +claude_client = anthropic.AsyncAnthropic(api_key=os.getenv("ANTHROPIC_API_KEY")) + +# Haiku 4.5: the agent only orchestrates ThoughtSpot's tools - pick a tool, pass the +# question on, summarise the answer - and the Analytics Agent does the analysis. Measured +# on one question, Haiku's four model calls took ~4s in total, against ~17s on Sonnet 5 +# and ~8s on Opus 5. Set ANTHROPIC_MODEL=claude-sonnet-5 for stronger reasoning on +# multi-step follow-ups, or claude-opus-5 for the most capable. +MODEL = os.getenv("ANTHROPIC_MODEL", "claude-haiku-4-5") +MAX_TOKENS = 16000 + +# Safety classifiers can decline a request (HTTP 200, stop_reason="refusal"). The +# server-side fallback beta re-routes those to a fallback model automatically - but +# only models with published fallback targets accept it. Sonnet 5 and Haiku 4.5 have +# none (`allowed_fallback_models` is empty on /v1/models), so for them a refusal is +# simply reported to the user. +REFUSAL_FALLBACK_BETA = "server-side-fallback-2026-07-01" +FALLBACK_CAPABLE_MODELS = ("claude-opus-5", "claude-fable-5-1") + + +def model_request_options(model: str) -> dict: + """The thinking and fallback options this model accepts.""" + options: dict[str, Any] = {} + # Haiku 4.5 predates adaptive thinking (it would need a fixed `budget_tokens`). + # Choosing tools does not need extended reasoning, so run it without thinking - + # that is also the fastest configuration. + if not model.startswith("claude-haiku"): + options["thinking"] = {"type": "adaptive", "display": "summarized"} + if model in FALLBACK_CAPABLE_MODELS: + options["betas"] = [REFUSAL_FALLBACK_BETA] + options["fallbacks"] = "default" + return options + +# ── ThoughtSpot MCP server ────────────────────────────────────────────────────── +TS_HOST = os.getenv("VITE_TS_HOST") or os.getenv("TS_HOST") +TS_AUTH_TOKEN = os.getenv("VITE_TS_AUTH_TOKEN") or os.getenv("TS_AUTH_TOKEN") + +if not TS_HOST: + raise RuntimeError("TS_HOST (or VITE_TS_HOST) must be set in .env") + +# The `/token/*` endpoint family is the static-bearer-token transport. +# (`/bearer/*` is the legacy path and is frozen on the older v1 toolset.) +# +# `latest` tracks the newest toolset, so a ThoughtSpot release can change the tools +# and the update shape under this app. Set TS_MCP_API_VERSION to a release date +# (e.g. 2026-05-01) to pin it instead. +MCP_API_VERSION = os.getenv("TS_MCP_API_VERSION", "latest") +# +# Adding `&enable-raw-session-updates=true` makes the server stream the Agent's own +# updates instead of its digested ones. This app reads either - see "Reading session +# updates" below - so you can turn it on through TS_MCP_URL without touching the code. +MCP_URL = os.getenv( + "TS_MCP_URL", + f"https://agent.thoughtspot.app/token/mcp?api-version={MCP_API_VERSION}", +) + +# Headers for the MCP session are built per request by `mcp_headers()` below, so a +# long-running process picks up a renewed token instead of holding an expired one. + +# ── ThoughtSpot tokens ────────────────────────────────────────────────────────── +# Two consumers, two policies: +# +# The browser. The Visual Embed SDK requires a FRESH token from every `getAuthToken` +# call. Handing it the same static token twice trips its duplicate-token check the +# moment the token stops verifying, and the callback can never recover because it +# returns the same string. So `/api/ts-token` mints one per request. +# +# This server, for its own MCP and REST calls. A static TS_AUTH_TOKEN eventually +# expires, and when it does every tool call fails with "Failed to validate +# connection" while the tool *list* still succeeds - so the app looks healthy right +# up until someone asks a question. Instead we mint one and cache it until shortly +# before it expires. TS_AUTH_TOKEN stays as the fallback for when no minting +# credentials are configured. +# +# Minting needs a username plus either the cluster secret key (Develop > Customizations +# > Security Settings > Trusted authentication) or that user's password. `secret_key` +# takes precedence when both are set. +TS_EMBED_USERNAME = os.getenv("TS_EMBED_USERNAME") +TS_SECRET_KEY = os.getenv("TS_SECRET_KEY") +TS_EMBED_PASSWORD = os.getenv("TS_EMBED_PASSWORD") +TS_TOKEN_VALIDITY_SEC = int(os.getenv("TS_TOKEN_VALIDITY_SEC", "1800")) + +CAN_MINT_TOKENS = bool(TS_EMBED_USERNAME and (TS_SECRET_KEY or TS_EMBED_PASSWORD)) + +if not CAN_MINT_TOKENS and not TS_AUTH_TOKEN: + raise RuntimeError( + "Set TS_EMBED_USERNAME plus TS_SECRET_KEY (or TS_EMBED_PASSWORD) so the server can " + "mint ThoughtSpot tokens, or a static TS_AUTH_TOKEN (or VITE_TS_AUTH_TOKEN) in .env" + ) + +if not CAN_MINT_TOKENS: + print( + "[Auth] TS_EMBED_USERNAME + TS_SECRET_KEY (or TS_EMBED_PASSWORD) are not set - " + "/api/ts-token will serve the static TS_AUTH_TOKEN. Fine for a local demo, but " + "the SDK cannot recover once that token expires." + ) + + +# Renew this long before the token actually expires, so a call already in flight does +# not race the expiry. +TOKEN_REFRESH_MARGIN_SEC = 120 +# The server's own token. Longer-lived than the browser's: it is never handed out, and +# a longer life means fewer mints on a slow cluster. +SERVER_TOKEN_VALIDITY_SEC = int(os.getenv("TS_SERVER_TOKEN_VALIDITY_SEC", "3600")) + +_server_token: str | None = None +# Wall-clock (time.time()) instant to renew at. Not time.monotonic(): on macOS that clock +# stops while the machine sleeps, so after a laptop sleep the cache would keep serving a +# token the cluster has already expired. +_server_token_expiry = 0.0 +# When the background task renews the token - well before `_server_token_expiry`, so no +# request has to wait on a mint. Minting can take 40s on a loaded cluster, and a request +# that finds the token expired pays all of it. +_server_token_refresh_at = 0.0 +# Serialises minting: an agent turn fires several tool calls at once, and without this +# they would each mint a token on a cold cache. +_server_token_lock = asyncio.Lock() + +# Renew this long before the token expires. Capped at half the token's life, so a short +# cluster-imposed lifetime still gets renewed at a sane interval. +TOKEN_BACKGROUND_LEAD_SEC = 720 +# After a failed background mint, try again this soon. The current token keeps serving. +TOKEN_BACKGROUND_RETRY_SEC = 30.0 +# The background loop re-reads the wall clock at least this often, so it notices when a +# laptop sleep has carried it past the renewal time (asyncio.sleep pauses with the host). +TOKEN_BACKGROUND_TICK_SEC = 60.0 + + +class TokenMintError(Exception): + """The cluster would not issue a token.""" + + def __init__(self, message: str, status_code: int) -> None: + super().__init__(message) + self.status_code = status_code + + +async def mint_token(validity_sec: int) -> dict: + """Ask ThoughtSpot for a bearer token. Returns the parsed response body.""" + payload: dict[str, Any] = { + "username": TS_EMBED_USERNAME, + "validity_time_in_sec": validity_sec, + } + # secret_key wins over password when both are present, per the REST API contract. + if TS_SECRET_KEY: + payload["secret_key"] = TS_SECRET_KEY + else: + payload["password"] = TS_EMBED_PASSWORD + + # This endpoint is occasionally slow to respond - 15s is normal on a loaded + # cluster - and a bare timeout would surface as an unhandled 500 with no + # explanation. + try: + async with httpx2.AsyncClient(timeout=httpx2.Timeout(10.0, read=90.0)) as client: + response = await client.post( + f"{TS_HOST.rstrip('/')}/api/rest/2.0/auth/token/full", json=payload + ) + except httpx2.HTTPError as exc: + raise TokenMintError(f"{type(exc).__name__}", 504) from exc + + if response.status_code != 200: + # Surface the cluster's own reason (bad secret key, trusted auth disabled, ...) + # rather than a blank 500 the browser cannot explain. + print(f"[Auth] Token mint failed {response.status_code}: {response.text[:300]}") + raise TokenMintError(f"HTTP {response.status_code}", 502) + + return response.json() + + +def server_token_valid() -> bool: + return bool(_server_token) and time.time() < _server_token_expiry + + +async def mint_server_token() -> str: + """Mint a server token and cache it. The caller holds `_server_token_lock`.""" + global _server_token, _server_token_expiry, _server_token_refresh_at + + body = await mint_token(SERVER_TOKEN_VALIDITY_SEC) + created = body.get("creation_time_in_millis") + expires = body.get("expiration_time_in_millis") + lifetime = ( + (expires - created) / 1000 + if isinstance(created, (int, float)) and isinstance(expires, (int, float)) + else SERVER_TOKEN_VALIDITY_SEC + ) + now = time.time() + _server_token = body["token"] + _server_token_expiry = now + max(60.0, lifetime - TOKEN_REFRESH_MARGIN_SEC) + _server_token_refresh_at = now + max(lifetime / 2, lifetime - TOKEN_BACKGROUND_LEAD_SEC) + print(f"[Auth] Minted a server token, good for {lifetime:.0f}s", flush=True) + return _server_token + + +async def server_token() -> str: + """A currently-valid bearer token for this server's own ThoughtSpot calls.""" + if not CAN_MINT_TOKENS: + return TS_AUTH_TOKEN + + # Fast path, outside the lock: a valid token never waits on a mint in progress - + # the background renewal holds the lock for as long as the cluster takes. + if server_token_valid(): + return _server_token + + async with _server_token_lock: + if server_token_valid(): + return _server_token + try: + return await mint_server_token() + except TokenMintError as exc: + # Keep serving the token we have: an expiring token still beats no token, + # and the next call retries the mint. + fallback = _server_token or TS_AUTH_TOKEN + if not fallback: + raise + print(f"[Auth] Could not mint a server token ({exc}); using the previous one.") + return fallback + + +async def keep_server_token_fresh() -> None: + """Background task: mint at startup, then renew ahead of expiry. + + Renewing on demand still works - this only moves the wait off the request path. + """ + while True: + due = _server_token_refresh_at if server_token_valid() else 0.0 + remaining = due - time.time() + if remaining > 0: + await asyncio.sleep(min(TOKEN_BACKGROUND_TICK_SEC, remaining)) + continue + try: + async with _server_token_lock: + # A request may have minted while this task waited for the lock. + if server_token_valid() and time.time() < _server_token_refresh_at: + continue + await mint_server_token() + except TokenMintError as exc: + print(f"[Auth] Background token renewal failed ({exc}); retrying soon.") + await asyncio.sleep(TOKEN_BACKGROUND_RETRY_SEC) + except Exception: # noqa: BLE001 - the loop must outlive any one failure + traceback.print_exc() + await asyncio.sleep(TOKEN_BACKGROUND_RETRY_SEC) + + +def invalidate_server_token() -> None: + """Force the next server_token() call to mint, keeping the old token as a fallback.""" + global _server_token_expiry, _server_token_refresh_at + _server_token_expiry = 0.0 + _server_token_refresh_at = 0.0 + + +async def mcp_headers() -> dict: + """Auth headers for an MCP session, with a token that is current right now.""" + return {"Authorization": f"Bearer {await server_token()}", "x-ts-host": TS_HOST} + +# Long read timeout: MCP replies stream over SSE and the Analytics Agent is slow. +MCP_TIMEOUT = httpx2.Timeout(30.0, read=300.0) + +# The MCP server's `initialize` intermittently fails with a bare HTTP 500 (surfaced as +# "MCPError: Server returned an error response") that a second attempt gets past. So a +# failed handshake is retried once. The retry also mints a fresh token - cheap, and it +# rules the token out. Only the handshake is retried: past it, the turn has streamed +# output to the browser and cannot be replayed. +MCP_CONNECT_ATTEMPTS = 2 +MCP_CONNECT_RETRY_DELAY_SEC = 2.0 + + +def root_cause(exc: BaseException) -> BaseException: + """anyio wraps transport failures in ExceptionGroups; unwrap to the first leaf.""" + while getattr(exc, "exceptions", None): + exc = exc.exceptions[0] + return exc + + +@asynccontextmanager +async def connect_mcp( + send_status: Callable[[dict], None] | None = None, +) -> AsyncIterator[tuple[Client, Any, str]]: + """An initialised MCP session, its tool listing and the bearer token it runs on. + + Retries a failed handshake once. + """ + for attempt in range(1, MCP_CONNECT_ATTEMPTS + 1): + stack = AsyncExitStack() + try: + headers = await mcp_headers() + token = headers["Authorization"].removeprefix("Bearer ") + http_client = await stack.enter_async_context( + httpx2.AsyncClient(headers=headers, timeout=MCP_TIMEOUT, follow_redirects=True) + ) + mcp = await stack.enter_async_context( + Client(streamable_http_client(MCP_URL, http_client=http_client)) + ) + listing = await mcp.list_tools() + except BaseException as exc: + await stack.aclose() + if not isinstance(exc, Exception) or attempt == MCP_CONNECT_ATTEMPTS: + raise + err = root_cause(exc) + print( + f"[MCP] Handshake attempt {attempt} failed ({type(err).__name__}: {err}); " + "retrying with a fresh token" + ) + if send_status: + send_status({"type": "status", "message": "Reconnecting to ThoughtSpot..."}) + invalidate_server_token() + await asyncio.sleep(MCP_CONNECT_RETRY_DELAY_SEC) + continue + + async with stack: + yield mcp, listing, token + return + + +# ── Shared MCP session ────────────────────────────────────────────────────────── +# +# Opening an MCP session costs ~10s on a loaded cluster (discover probe, initialize, +# tools/list - each several seconds), which used to be paid on every chat message. One +# session is now shared by every turn and reopened only when it has to be: +# +# - the server token rotated (the session's HTTP client carries the old one), or +# - the session broke: a transport error, or the server dropped it (-32600 +# "Session terminated", e.g. after an idle timeout). +# +# The session is opened and closed by its own background task, never by a request. +# anyio requires a session to be closed by the task that opened it, and closing it in +# a request task is what used to get cancelled half-way when the browser stream ended. + +# What the MCP client raises once the server no longer knows the session: -32600 with +# the message "Session terminated". Only that pairing means the call never ran and is +# safe to repeat. -32600 alone is not enough - it is JSON-RPC's generic "Invalid +# Request", which the client also uses when a call may already have been processed. +MCP_SESSION_TERMINATED = (-32600, "Session terminated") + + +def is_session_terminated(exc: MCPError) -> bool: + return (exc.code, exc.message) == MCP_SESSION_TERMINATED + + +class McpConnection: + """One open MCP session, held open by its owner task until retired. + + `leases` counts the turns using it. A retired connection stays open until the last + of them is done, so rotating the token never pulls a session out from under a turn. + """ + + def __init__( + self, mcp: Client, listing: Any, token: str, stop: asyncio.Event, task: asyncio.Task + ) -> None: + self.mcp = mcp + self.listing = listing + self.token = token + self.task = task + self._stop = stop + self.leases = 0 + self.retired = False + + @property + def usable(self) -> bool: + return not self.retired and not self.task.done() + + # Both of these are synchronous on purpose: a turn gives up its lease without + # awaiting anything, so a cancelled turn cannot be interrupted mid-teardown. The + # owner task does the actual closing. + def retire(self) -> None: + self.retired = True + if self.leases == 0: + self._stop.set() + + def release(self) -> None: + self.leases -= 1 + if self.retired and self.leases == 0: + self._stop.set() + + +class McpPool: + """Hands out the shared MCP session, opening a new one when the current can't serve.""" + + def __init__(self) -> None: + # Held across get-or-create, so turns arriving together share one handshake + # instead of each opening a session. + self._lock = asyncio.Lock() + self._conn: McpConnection | None = None + + async def acquire(self, send_status: Callable[[dict], None] | None = None) -> McpConnection: + async with self._lock: + token = await server_token() + conn = self._conn + if conn and conn.usable and conn.token == token: + conn.leases += 1 + return conn + if conn: + reason = "token rotated" if conn.usable else "session closed" + print(f"[MCP] Replacing shared session ({reason})") + conn.retire() + self._conn = None + conn = await self._open(send_status) + conn.leases += 1 + self._conn = conn + return conn + + async def _open(self, send_status: Callable[[dict], None] | None) -> McpConnection: + ready: asyncio.Future = asyncio.get_running_loop().create_future() + stop = asyncio.Event() + + async def own() -> None: + try: + async with connect_mcp(send_status) as opened: + if ready.done(): + return # the asker left mid-handshake; just let the session close + ready.set_result(opened) + await stop.wait() + except Exception as exc: # noqa: BLE001 - reported to the waiter, or logged + if not ready.done(): + ready.set_exception(exc) + else: + err = root_cause(exc) + print(f"[MCP] Shared session ended: {type(err).__name__}: {err}") + except BaseException: + # Cancellation, possibly grouped with other errors by anyio. Resolve the + # waiter so a turn never waits forever on a handshake that is gone. + if not ready.done(): + ready.cancel() + raise + + task = asyncio.create_task(own()) + try: + mcp, listing, token = await ready + except asyncio.CancelledError: + # Close the session once it opens rather than leave it with no owner. + stop.set() + if not asyncio.current_task().cancelling(): + # This turn was not cancelled - the handshake was (e.g. at shutdown). + # Report it as a failure, not as a hang-up nobody is told about. + raise RuntimeError("The MCP session closed while it was opening") from None + raise + return McpConnection(mcp, listing, token, stop, task) + + async def close(self) -> None: + """Retire the current session and wait for it to close (app shutdown).""" + async with self._lock: + conn, self._conn = self._conn, None + if conn: + conn.leases = 0 + conn.retire() + await asyncio.wait({conn.task}, timeout=10) + + +mcp_pool = McpPool() + + +class McpTurn: + """One chat turn's handle on the shared MCP session. + + Exposes `call_tool` with Client's signature, so the tool code does not know the + session is shared. A call that finds the session dropped is repeated once on a + fresh one; any other connection failure retires the session for the next turn + and surfaces the error, since the call may already have run. + """ + + def __init__(self, send_status: Callable[[dict], None] | None = None) -> None: + self.send_status = send_status + self.conn: McpConnection | None = None + # Every connection this turn has leased. A replaced one is released only when + # the turn ends, so a sibling call still in flight on it is not cut off. + self._held: list[McpConnection] = [] + # Tool calls run concurrently; only one of them should swap the session. + self._swap_lock = asyncio.Lock() + + async def __aenter__(self) -> "McpTurn": + self.conn = await mcp_pool.acquire(self.send_status) + self._held.append(self.conn) + return self + + async def __aexit__(self, *exc_info: Any) -> None: + for conn in self._held: + conn.release() + self._held.clear() + self.conn = None + + @property + def listing(self) -> Any: + return self.conn.listing + + async def call_tool(self, name: str, arguments: dict) -> Any: + for attempt in (1, 2): + conn = self.conn + try: + return await conn.mcp.call_tool(name, arguments) + except MCPError as exc: + if attempt == 2 or not is_session_terminated(exc): + raise + print(f"[MCP] {name}: session terminated by the server; reconnecting") + await self._replace(conn) + except Exception: + # Transport-level failure: the session is suspect, but the call may + # have reached the server, so it is not repeated. + conn.retire() + raise + raise AssertionError("unreachable") + + async def _replace(self, failed: McpConnection) -> None: + async with self._swap_lock: + if self.conn is not failed: + return # a concurrent call already swapped it + if self.send_status: + self.send_status({"type": "status", "message": "Reconnecting to ThoughtSpot..."}) + failed.retire() + self.conn = await mcp_pool.acquire(self.send_status) + self._held.append(self.conn) + +# ThoughtSpot MCP tools, as exposed by the Spotter 3 toolset: +# +# check_connectivity - test connectivity + auth. No inputs. +# search_objects - find existing Liveboards / Answers / Worksheets by name. +# Metadata only; never returns data. +# create_analysis_session - start a session. Optional: data_source_id. +# Returns analytical_session_id. +# send_session_message - ask the Analytics Agent a question. +# Inputs: analytical_session_id, message, additional_context. +# get_session_updates - poll for updates. Input: analytical_session_id. +# Returns session_updates[] + is_done. See AUTOPOLL below. +# create_dashboard - build a dashboard from answer_ids. +# Inputs: title, answers[], note_tile. Returns link. +# list_orgs / switch_org - OAuth-only. The server hides them on `/token/*`, so they +# never appear in the tool list for this static-token setup. +# +# Set ALLOWED_TOOLS to a list of names to restrict the agent. None allows everything +# the server exposes, which is the right default: the tool list is version-negotiated, +# so a hardcoded list silently drops tools added in later API versions. +ALLOWED_TOOLS: list[str] | None = None + +# ── Server-side polling of get_session_updates ────────────────────────────────── +POLL_TOOL = "get_session_updates" +POLL_INITIAL_DELAY = 0.75 # seconds before the first re-poll +POLL_MAX_DELAY = 4.0 # cap on the backoff +POLL_TIMEOUT = 300.0 # give up after this long without is_done + +SYSTEM_PROMPT = """You are a data analyst assistant powered by ThoughtSpot's Analytics Agent. + +Workflow: +- Create one analysis session per conversation with `create_analysis_session`, then ask + questions with `send_session_message`, then call `get_session_updates` once. +- `get_session_updates` is polled to completion for you: a single call returns the Agent's + full response, so never call it twice for the same question. +- Use `search_objects` to find existing Liveboards, Answers or Worksheets by name. It + returns metadata only, never data - to answer a data question, ask the Agent. +- Use `create_dashboard` when the user wants to save or share results, passing the + `answer_id` values from the answers you want on it. + +Presenting answers: +- Every `answer` update is ALREADY rendered in the UI as an interactive ThoughtSpot chart, + in the order it was returned. Do not emit `; }; From f314d4327348e0f2a17395e6609e9c7451836778 Mon Sep 17 00:00:00 2001 From: Mourya Balabhadra Date: Mon, 28 Sep 2026 04:02:55 -0700 Subject: [PATCH 3/4] fix duplicate calls and package lock --- mcp/python-react-agent-simple-ui/README.md | 11 +++++- .../client/.npmrc | 3 ++ .../client/package-lock.json | 16 ++++---- .../client/src/App.tsx | 39 +++++++++++++------ 4 files changed, 48 insertions(+), 21 deletions(-) create mode 100644 mcp/python-react-agent-simple-ui/client/.npmrc diff --git a/mcp/python-react-agent-simple-ui/README.md b/mcp/python-react-agent-simple-ui/README.md index b9ba111..e972a7c 100644 --- a/mcp/python-react-agent-simple-ui/README.md +++ b/mcp/python-react-agent-simple-ui/README.md @@ -468,9 +468,18 @@ This works for both paths into the DOM: charts the server streams as `answer` ev The client injects each answer's iframe as markup rather than rendering ``; }; +/** + * One MCP answer's iframe, written into the DOM exactly once. + * + * startAutoMCPFrameRenderer swaps the iframe for the real embed via replaceWith(), so + * React must not own that node - it owns only this wrapper, filled with + * `dangerouslySetInnerHTML`. React re-applies that prop whenever it receives a new + * object, and a fresh `{ __html }` literal is a new object on every render of the message + * list - each keystroke, each streamed delta. Every re-application throws the embed away + * and the renderer resolves the answer again from scratch: two ThoughtSpot calls per + * chart, and a stored answer's take tens of seconds. + * + * So the prop object is created once per mount and kept in state (which also survives + * dev Fast Refresh). Callers key this by `html`, so a different answer gets a fresh + * wrapper. The markup goes in before the wrapper is attached, so the renderer sees the + * iframe once - writing it from an effect instead would show it twice, via the attached + * container and via the iframe itself, and resolve it twice. + */ +function AnswerFrame({ html }: { html: string }) { + const [markup] = useState(() => ({ __html: html })); + return
; +} + function App() { const [messages, setMessages] = useState([]); const [input, setInput] = useState(""); @@ -618,11 +634,10 @@ function App() { {answer.title && (
{answer.title}
)} -
+ {(() => { + const html = answerHtml(answer, sessionId); + return ; + })()} ))} {msg.content ? ( From ab03d751a3172477a4ca2317a3ebcefad8e78b09 Mon Sep 17 00:00:00 2001 From: Mourya Balabhadra Date: Wed, 30 Sep 2026 01:21:45 -0700 Subject: [PATCH 4/4] Add thinking answers --- mcp/python-react-agent-simple-ui/README.md | 16 +- .../client/src/App.css | 168 ++++++++++ .../client/src/App.tsx | 314 +++++++++++++++++- ...th_spotter3_mcp_server_and_chat_history.py | 228 +++++++++++-- 4 files changed, 687 insertions(+), 39 deletions(-) diff --git a/mcp/python-react-agent-simple-ui/README.md b/mcp/python-react-agent-simple-ui/README.md index e972a7c..f88c1fb 100644 --- a/mcp/python-react-agent-simple-ui/README.md +++ b/mcp/python-react-agent-simple-ui/README.md @@ -231,6 +231,13 @@ The ThoughtSpot Analytics Agent answers asynchronously: `send_session_message` r Instead `autopoll_session_updates()` does it in-process: it polls with backoff until the Agent is done, accumulates every update, and hands Claude **one** consolidated tool result. Progress streams to the UI as `status` events while it waits - the Agent's steps ("Searching for Datasets", ...), not its reasoning prose, which arrives in word-sized chunks. The model still receives the full reasoning in the tool result. +The chat-history backend also keeps that reasoning for the UI. `emit_work_event()` sends each step, each stretch of reasoning prose, and each query the Agent tried before settling as a `work` event, and the client shows them in a collapsible **Show work** section above the answer. This is the Analytics Agent's own work as the Spotter MCP server reports it — it has nothing to do with Claude's extended thinking. Each group of tried queries is a **Visualized data** row. Its charts load only when the row is expanded: + +- **Live:** the `work` event carries the thinking answer's own `iframe_url` / `frame_params`, so the chart renders immediately. +- **Replayed:** those URLs expire with the answer (about 8 hours), so they are never stored. `GET /api/conversations/{id}` returns each turn's `work_answer_ids` — ThoughtSpot's ids for its thinking answers, in step order. Expanding a row calls `GET /api/conversations/{id}/work-answers/{answer_id}`, which loads that answer again through the conversation service and returns fresh embed ids. Loading re-runs the query, so it can take 30s+ on a busy cluster; the client starts each load once and reuses it. + +Thinking answers stay out of the turn's `answers` list: stored answers are resolved by their position among the *non-thinking* answers (by the SDK too), so counting a thinking answer there would shift every later chart onto the wrong one. + ```python POLL_INITIAL_DELAY = 0.75 # seconds before the first re-poll POLL_MAX_DELAY = 4.0 # backoff cap; resets whenever new updates arrive @@ -286,6 +293,7 @@ Other Claude API details worth noting: | `delta` | `text` | Streamed assistant text | | `status` | `message` | Thinking, tool start, Analytics Agent step, or reconnecting | | `answer` | `answer_id`, `title`, `query`, `iframe_url`, `frame_params` | A chart to render. `frame_params` carries the embed ids when there is no ready-made `iframe_url` | +| `work` | `kind`, `text`, `query`, `iframe_url`, `frame_params` | One step of the Analytics Agent's work, for **Show work** (chat-history backend only). `kind` is `step` (progress), `thought` (reasoning prose — consecutive ones join into one), or `query` (a query tried before the final answer) | | `done` | `response_id` | Turn complete; pass `response_id` back for follow-ups | | `error` | `message` | Fatal error for this turn | @@ -335,7 +343,7 @@ ALLOWED_TOOLS = ["create_analysis_session", "send_session_message", "get_session | Field | Present when | Description | |-------|--------------|-------------| | `type` | always | `text`, `text_chunk`, `answer`, or `step_notification` | -| `is_thinking` | always | Whether this update is part of the Agent's reasoning rather than its final answer. The server shows reasoning *steps* in the UI as status text, and passes the reasoning prose to the model only. | +| `is_thinking` | always | Whether this update is part of the Agent's reasoning rather than its final answer. The server shows reasoning *steps* in the UI as status text and passes the reasoning prose to the model. The chat-history backend also keeps both for the **Show work** section. | | `text` | `text`, `text_chunk`, `step_notification` | Message text. Consecutive `text_chunk` values are concatenated by the server before the model sees them. | | `answer_id` | `answer` | Identifier to pass to `create_dashboard`. | | `answer_title` | `answer` | Human-readable title. | @@ -370,7 +378,7 @@ Two tables, because the browser and the model need different things: | `conversations` | `id`, `title`, `created_at`, `updated_at` | The sidebar list. `title` is the first line of the first user message. | | | `analytical_session_id` | The ThoughtSpot session, so a reopened conversation continues in the same one. | | | `claude_messages` | The **raw** Claude message list — `tool_use` / `tool_result` / `thinking` blocks included — replayed into the next request so follow-ups keep full context after a restart. | -| `turns` | `role`, `content`, `answers` | What the UI renders. `answers` holds each chart's title and query - **not** its `iframe_url`, which expires with the ThoughtSpot answer (about 8 hours). | +| `turns` | `role`, `content`, `answers`, `work` | What the UI renders. `answers` holds each chart's title - **not** its `iframe_url`, which expires with the ThoughtSpot answer (about 8 hours). `work` holds the Analytics Agent's steps for **Show work**; `db_init()` adds the column to databases created before it existed. | `turns` cascades on delete (`PRAGMA foreign_keys=ON`), so removing a conversation removes its transcript. @@ -401,7 +409,8 @@ In-memory `conversations` / `analytical_sessions` dicts stay as a hot cache in f | Method | Path | Purpose | |--------|------|---------| | `GET` | `/api/conversations` | List conversations, newest first. Returns `id`, `title`, timestamps, `turn_count`. | -| `GET` | `/api/conversations/{id}` | One conversation with its `turns` (each `role`, `content`, `answers` with `answer_index`). | +| `GET` | `/api/conversations/{id}` | One conversation with its `turns` (each `role`, `content`, `answers` with `answer_index`, `work`, `work_answer_ids`). | +| `GET` | `/api/conversations/{id}/work-answers/{answer_id}` | Fresh `frame_params` for one thinking answer, so an expanded **Visualized data** row can render it. | | `PATCH` | `/api/conversations/{id}` | Rename. Body: `{"title": "..."}`. | | `DELETE` | `/api/conversations/{id}` | Delete the conversation and its turns. | | `GET` | `/api/ts-token` | A freshly minted ThoughtSpot token for the Visual Embed SDK (both servers). | @@ -447,6 +456,7 @@ The React client is shared across both backends: - Renders assistant text as **markdown** (tables, code blocks, and raw HTML via `rehypeRaw`) - Renders ThoughtSpot charts from `answer` events as auto-upgraded embeds - Shows **real-time status** — Claude thinking, tool calls, and Analytics Agent progress +- Shows a collapsible **Show work** section with the Analytics Agent's steps, reasoning and tried queries, when the backend sends `work` events - Tracks `response_id` across turns for multi-turn continuity - Shows a **chat history sidebar** when the backend exposes `/api/conversations` — click to reopen a chat (charts and all), `×` to delete. A loading indicator shows while a stored chat is fetched. Against the backend without history the probe fails and the sidebar is simply not rendered. diff --git a/mcp/python-react-agent-simple-ui/client/src/App.css b/mcp/python-react-agent-simple-ui/client/src/App.css index 68b5d3f..734b607 100644 --- a/mcp/python-react-agent-simple-ui/client/src/App.css +++ b/mcp/python-react-agent-simple-ui/client/src/App.css @@ -503,6 +503,174 @@ body { border: none; } +/* "Show work" - the Analytics Agent's steps as a timeline, collapsed by default */ +.show-work { + margin: 2px 0 14px; + font-size: 14px; +} + +.show-work summary { + display: inline-flex; + align-items: center; + gap: 4px; + cursor: pointer; + list-style: none; + user-select: none; + color: var(--text-secondary); +} + +.show-work summary::-webkit-details-marker { + display: none; +} + +.show-work summary:hover { + color: var(--text); +} + +.show-work summary > .work-icon { + width: 14px; + height: 14px; + transition: transform 0.15s ease; +} + +.show-work[open] > summary > .work-icon, +.work-queries[open] > summary > .work-icon { + transform: rotate(90deg); +} + +.show-work .work-timeline { + list-style: none; + margin: 12px 0 0; + padding: 0; +} + +.show-work .work-timeline > li { + position: relative; + margin: 0; + display: flex; + gap: 12px; + padding-bottom: 14px; + line-height: 1.55; + color: var(--text); +} + +/* The connector between markers. */ +.show-work .work-timeline > li:not(:last-child)::before { + content: ""; + position: absolute; + left: 7px; + top: 22px; + bottom: 2px; + width: 1px; + background: var(--border); +} + +.show-work .work-timeline > li:last-child { + padding-bottom: 0; +} + +.work-marker { + flex: 0 0 16px; + height: 22px; + display: flex; + align-items: center; + justify-content: center; + color: var(--text-secondary); +} + +.work-dot::after { + content: ""; + width: 6px; + height: 6px; + border-radius: 50%; + background: var(--text-secondary); + opacity: 0.6; +} + +.work-icon { + width: 15px; + height: 15px; + fill: none; + stroke: currentColor; + stroke-width: 1.4; + stroke-linecap: round; + stroke-linejoin: round; +} + +.work-body { + flex: 1; + min-width: 0; + overflow-wrap: anywhere; +} + +/* Reasoning is markdown; keep it as tight as the rest of the timeline. */ +.work-thought .work-body p, +.work-thought .work-body ul, +.work-thought .work-body ol { + margin: 0 0 6px; +} + +.work-thought .work-body > :last-child { + margin-bottom: 0; +} + +.work-thought .work-body ul, +.work-thought .work-body ol { + padding-left: 20px; +} + +.work-queries > summary { + color: var(--text); +} + +.work-query { + margin-top: 8px; +} + +.work-query-title { + color: var(--text-secondary); + font-size: 13px; + margin-bottom: 4px; +} + +.work-query code { + display: block; + white-space: pre-wrap; + padding: 8px 10px; + border: 1px solid var(--border); + border-radius: 6px; + font-size: 12px; +} + +.work-chart { + margin-top: 8px; + border: 1px solid var(--border); + border-radius: 10px; + overflow: hidden; + background: var(--surface); +} + +.work-chart iframe { + display: block; + width: 100%; + border: none; +} + +.work-chart-note { + margin-top: 8px; + font-size: 13px; + color: var(--text-secondary); +} + +.work-retry { + border: none; + background: none; + padding: 0; + font: inherit; + color: var(--primary); + cursor: pointer; +} + /* ---- Chat history sidebar ---- */ .shell { diff --git a/mcp/python-react-agent-simple-ui/client/src/App.tsx b/mcp/python-react-agent-simple-ui/client/src/App.tsx index a2a4198..7c8645e 100644 --- a/mcp/python-react-agent-simple-ui/client/src/App.tsx +++ b/mcp/python-react-agent-simple-ui/client/src/App.tsx @@ -35,12 +35,29 @@ interface Answer { answer_index?: number | null; } +/** + * One step of the Analytics Agent's own work, as the Spotter MCP server reports it: + * a progress note, a stretch of its reasoning, or a query it tried before settling + * on the answer. Not Claude's reasoning. + */ +interface WorkStep { + kind: "step" | "thought" | "query"; + text: string; + query?: string | null; + /** A query's chart, on live steps only - stored steps are resolved on expand. */ + iframe_url?: string | null; + frame_params?: FrameParams | null; +} + type Role = "user" | "assistant"; interface Message { role: Role; content: string; answers?: Answer[]; + work?: WorkStep[]; + /** Stored turns only: ThoughtSpot's id for each Show work query, in order. */ + workAnswerIds?: string[]; isError?: boolean; } @@ -52,7 +69,13 @@ interface ConversationSummary { interface ConversationDetail { id: string; analytical_session_id?: string | null; - turns?: { role: Role; content: string; answers?: Answer[] }[]; + turns?: { + role: Role; + content: string; + answers?: Answer[]; + work?: WorkStep[]; + work_answer_ids?: string[]; + }[]; } /** Server-sent events on the /api/chat stream. */ @@ -60,6 +83,7 @@ type StreamEvent = | { type: "delta"; text: string } | { type: "status"; message: string } | ({ type: "answer" } & Answer) + | ({ type: "work" } & WorkStep) | { type: "done"; response_id: string } | { type: "error"; message: string }; @@ -101,6 +125,7 @@ const embedDarkVariables: Record = { "--ts-var-spotterviz-text-primary": DARK_TEXT, "--ts-var-spotterviz-text-secondary": DARK_TEXT_MUTED, "--ts-var-spotterviz-border-color": DARK_BORDER, + "--ts-var-chip-color": DARK_SURFACE, }; // The conversational-answer content wrapper paints its own white background, and no @@ -196,6 +221,263 @@ const mcpFrameObserver = startAutoMCPFrameRenderer({ // duplicating the conversation-service calls once per reload since the page loaded. import.meta.hot?.dispose(() => mcpFrameObserver.disconnect()); +/** + * Add one streamed step to a turn's work. Reasoning arrives a few words per event, + * so a `thought` continues the previous one - the server's `append_work` applies the + * same rule, so a replayed turn matches what streamed live. + */ +const appendWork = (work: WorkStep[], step: WorkStep): WorkStep[] => { + const last = work[work.length - 1]; + if (step.kind === "thought" && last?.kind === "thought") { + return [...work.slice(0, -1), { ...last, text: last.text + step.text }]; + } + return [ + ...work, + { + kind: step.kind, + text: step.text, + query: step.query, + iframe_url: step.iframe_url, + frame_params: step.frame_params, + }, + ]; +}; + +/** One row of the Show work timeline. Consecutive queries fold into one row. */ +type WorkRow = + | { kind: "thought"; text: string } + | { kind: "step"; text: string } + | { kind: "queries"; queries: { step: WorkStep; index: number }[] }; + +const workRows = (work: WorkStep[]): WorkRow[] => { + const rows: WorkRow[] = []; + // A query's position among the turn's queries is how the server finds it again. + let queryIndex = 0; + for (const step of work) { + const text = step.text.trim(); + const last = rows[rows.length - 1]; + if (step.kind === "query") { + const query = { step, index: queryIndex++ }; + if (last?.kind === "queries") last.queries.push(query); + else rows.push({ kind: "queries", queries: [query] }); + } else if (text) { + // Reasoning often arrives with leading blank lines - drop them, and skip + // a thought that is only whitespace. + rows.push({ kind: step.kind, text }); + } + } + return rows; +}; + +const ICON_PATHS = { + search: "M7 12a5 5 0 1 0 0-10 5 5 0 0 0 0 10zm3.5-1.5L14 14", + memory: + "M4.5 12.5a3 3 0 0 1-.4-6 4 4 0 0 1 7.8 0 3 3 0 0 1-.4 6M8 8v6m-2-2 2 2 2-2", + chart: "M8 2v6h6A6 6 0 1 1 8 2zm2-.5A5 5 0 0 1 14.5 6H10z", + chevron: "M6 4l4 4-4 4", +}; + +function WorkIcon({ name }: { name: keyof typeof ICON_PATHS }) { + return ( + + ); +} + +/** Picks an icon for a progress step from its wording. */ +const stepIcon = (text: string): keyof typeof ICON_PATHS => { + if (/memor/i.test(text)) return "memory"; + if (/search|context|schema/i.test(text)) return "search"; + return "chart"; +}; + +/** + * The Analytics Agent's steps behind an answer, collapsed until asked for. Queries + * are shown as text, never as charts: only the settled answer is a real one. + */ +/** Where a stored turn's thinking answers are loaded from. */ +interface WorkSource { + conversationId: string | null; + answerIds: string[]; +} + +/** + * In-flight and finished thinking-answer loads, by conversation and answer id. + * Loading re-runs the answer on ThoughtSpot and takes tens of seconds, so a + * remount (React StrictMode, reopening the row) must not start it again. + */ +const workAnswerLoads = new Map>(); + +const loadWorkAnswer = (conversationId: string, answerId: string) => { + const key = `${conversationId}/${answerId}`; + let load = workAnswerLoads.get(key); + if (!load) { + load = fetch( + `${HISTORY_URL}/${conversationId}/work-answers/${encodeURIComponent(answerId)}`, + ).then(async (response) => { + if (!response.ok) throw new Error(String(response.status)); + return ((await response.json()) as { frame_params: FrameParams }).frame_params; + }); + // A failure is not kept, so expanding the row again retries. + load.catch(() => workAnswerLoads.delete(key)); + workAnswerLoads.set(key, load); + } + return load; +}; + +/** + * The chart for one query the Agent tried. A live step carries its URL. A stored + * one does not - it would have expired - so the server loads that thinking answer + * again by its position. Thinking answers are outside the SDK's stored-answer + * index, which counts only settled answers, so they cannot use that path. + */ +function WorkQueryChart({ + step, + index, + source, +}: { + step: WorkStep; + index: number; + source: WorkSource; +}) { + const liveSrc = answerSrc(step, null); + const answerId = source.answerIds[index]; + const conversationId = source.conversationId; + const [resolved, setResolved] = useState(null); + const [failed, setFailed] = useState(false); + // Bumped by Retry to run the load again; the failed load is already evicted. + const [attempt, setAttempt] = useState(0); + + useEffect(() => { + if (liveSrc || !conversationId || !answerId) return; + let cancelled = false; + setFailed(false); + loadWorkAnswer(conversationId, answerId) + .then((frame_params) => { + if (!cancelled) setResolved(answerHtml({ ...step, frame_params }, null)); + }) + .catch(() => !cancelled && setFailed(true)); + return () => { + cancelled = true; + }; + }, [liveSrc, conversationId, answerId, step, attempt]); + + const html = liveSrc ? answerHtml(step, null) : resolved; + // Nothing to load: a stored step ThoughtSpot has no answer for. + if (!liveSrc && (!conversationId || !answerId)) return null; + if (failed) { + return ( +

+ Could not load this chart.{" "} + +

+ ); + } + if (!html) return

Loading chart…

; + return ( +
+ +
+ ); +} + +/** + * One "Visualized data" row. Its charts mount the first time it is expanded, and + * stay mounted after - closing and reopening must not load them again. + */ +function WorkQueries({ + queries, + source, +}: { + queries: { step: WorkStep; index: number }[]; + source: WorkSource; +}) { + const [opened, setOpened] = useState(false); + const count = queries.length; + return ( +
e.currentTarget.open && setOpened(true)} + > + + {count === 1 ? "Visualized data" : `Created ${count} visualizations`} + + + {queries.map(({ step, index }) => ( +
+ {step.text &&
{step.text}
} + {step.query && {step.query}} + {opened && } +
+ ))} +
+ ); +} + +function ShowWork({ + work, + finished, + source, +}: { + work: WorkStep[]; + finished: boolean; + source: WorkSource; +}) { + const rows = workRows(work); + if (!rows.length) return null; + return ( +
+ + Show work + + +
    + {rows.map((row, i) => { + if (row.kind === "thought") { + return ( +
  • + +
    + + {row.text} + +
    +
  • + ); + } + if (row.kind === "step") { + return ( +
  • + + + +
    {row.text}
    +
  • + ); + } + return ( +
  • + + + + +
  • + ); + })} + {finished && ( +
  • + +
    Finished
    +
  • + )} +
+
+ ); +} + const escapeAttr = (value: unknown) => String(value ?? "").replace(/&/g, "&").replace(/"/g, """); @@ -344,6 +626,8 @@ function App() { role: turn.role, content: turn.content, answers: turn.answers || [], + work: turn.work || [], + workAnswerIds: turn.work_answer_ids || [], })), ); setResponseId(data.id); @@ -406,7 +690,7 @@ function App() { setMessages((prev) => [ ...prev, { role: "user", content: userMessage }, - { role: "assistant", content: "", answers: [] }, + { role: "assistant", content: "", answers: [], work: [] }, ]); try { @@ -478,6 +762,17 @@ function App() { }); break; } + case "work": + setMessages((prev) => { + const updated = [...prev]; + const last = updated[updated.length - 1]; + updated[updated.length - 1] = { + ...last, + work: appendWork(last.work || [], data), + }; + return updated; + }); + break; case "done": setResponseId(data.response_id); setStatus(""); @@ -492,6 +787,7 @@ function App() { role: "assistant", content: `Error: ${data.message}`, answers: last.answers || [], + work: last.work || [], isError: true, }; return updated; @@ -623,6 +919,17 @@ function App() {

{msg.content}

) : ( <> + {!!msg.work?.length && ( + + )} {(msg.answers || []).map((answer) => (
) : ( - !(msg.answers || []).length && ( + !(msg.answers || []).length && + !(msg.work || []).length && ( ) )} diff --git a/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server_and_chat_history.py b/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server_and_chat_history.py index a43efb3..6f21be2 100644 --- a/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server_and_chat_history.py +++ b/mcp/python-react-agent-simple-ui/server/claude_agent_with_spotter3_mcp_server_and_chat_history.py @@ -32,6 +32,7 @@ import asyncio import json +from urllib.parse import quote import os import sqlite3 import time @@ -678,6 +679,7 @@ async def _replace(self, failed: McpConnection) -> None: role TEXT NOT NULL, content TEXT NOT NULL DEFAULT '', answers TEXT NOT NULL DEFAULT '[]', + work TEXT NOT NULL DEFAULT '[]', created_at TEXT NOT NULL ); @@ -702,6 +704,11 @@ def db_connect() -> sqlite3.Connection: def db_init() -> None: with closing(db_connect()) as conn, conn: conn.executescript(SCHEMA) + # `CREATE TABLE IF NOT EXISTS` leaves a database from before `work` existed + # untouched, so add the column there. + columns = {row["name"] for row in conn.execute("PRAGMA table_info(turns)")} + if "work" not in columns: + conn.execute("ALTER TABLE turns ADD COLUMN work TEXT NOT NULL DEFAULT '[]'") print(f"[History] SQLite at {DB_PATH}") @@ -740,12 +747,14 @@ def db_start_conversation(conv_id: str, first_message: str) -> None: ) -def db_add_turn(conv_id: str, role: str, content: str, answers: list[dict]) -> None: +def db_add_turn( + conv_id: str, role: str, content: str, answers: list[dict], work: list[dict] | None = None +) -> None: with closing(db_connect()) as conn, conn: conn.execute( - "INSERT INTO turns (conversation_id, role, content, answers, created_at)" - " VALUES (?, ?, ?, ?, ?)", - (conv_id, role, content, json.dumps(answers), now_iso()), + "INSERT INTO turns (conversation_id, role, content, answers, work, created_at)" + " VALUES (?, ?, ?, ?, ?, ?)", + (conv_id, role, content, json.dumps(answers), json.dumps(work or []), now_iso()), ) conn.execute( "UPDATE conversations SET updated_at = ? WHERE id = ?", (now_iso(), conv_id) @@ -810,7 +819,7 @@ def db_get_conversation(conv_id: str) -> dict | None: if not row: return None turns = conn.execute( - "SELECT role, content, answers, created_at FROM turns" + "SELECT role, content, answers, work, created_at FROM turns" " WHERE conversation_id = ? ORDER BY id", (conv_id,), ).fetchall() @@ -826,6 +835,7 @@ def db_get_conversation(conv_id: str) -> dict | None: "role": turn["role"], "content": turn["content"], "answers": answers, + "work": json.loads(turn["work"] or "[]"), "created_at": turn["created_at"], } ) @@ -833,6 +843,14 @@ def db_get_conversation(conv_id: str) -> dict | None: return {**dict(row), "turns": replayed_turns} +def db_get_session_id(conv_id: str) -> str | None: + with closing(db_connect()) as conn: + row = conn.execute( + "SELECT analytical_session_id FROM conversations WHERE id = ?", (conv_id,) + ).fetchone() + return row["analytical_session_id"] if row else None + + def db_delete_conversation(conv_id: str) -> bool: with closing(db_connect()) as conn, conn: deleted = conn.execute( @@ -892,6 +910,8 @@ def __init__(self, queue: asyncio.Queue) -> None: self.queue = queue self.text_parts: list[str] = [] self.answers: list[dict] = [] + # The Analytics Agent's own steps, for the UI's "Show work" section. + self.work: list[dict] = [] self.persisted = False # guards against writing the assistant turn twice def send(self, event: dict) -> None: @@ -903,10 +923,11 @@ def send(self, event: dict) -> None: { "answer_id": event.get("answer_id"), "title": event.get("title"), - "query": event.get("query"), "iframe_url": event.get("iframe_url"), } ) + elif kind == "work": + append_work(self.work, event) self.queue.put_nowait(event) @property @@ -914,6 +935,24 @@ def text(self) -> str: return "".join(self.text_parts) +def append_work(work: list[dict], event: dict) -> None: + """Add one `work` event to a turn's work list, the way the client does. + + Reasoning prose streams a few words per event, so a `thought` continues the + previous one instead of starting a new step. The client applies the same rule, + which keeps the live and replayed views identical. + """ + kind = event.get("kind") + text = event.get("text") or "" + if kind == "thought" and work and work[-1].get("kind") == "thought": + work[-1]["text"] += text + return + step = {"kind": kind, "text": text} + if event.get("query"): + step["query"] = event["query"] + work.append(step) + + def storable_answers(answers: list[dict]) -> list[dict]: """Strip the parts of an answer that go stale, keeping what survives. @@ -925,9 +964,12 @@ def storable_answers(answers: list[dict]) -> list[dict]: We therefore persist only the durable parts. On replay the client asks the Visual Embed SDK to resolve a live URL from the conversation id plus the answer's position, so the chart is re-run against current data. + + `query` is dropped too: nothing reads it on replay. Stripping it here, rather than + only not recording it, also keeps it out of responses for rows written before. """ return [ - {k: v for k, v in answer.items() if k not in ("iframe_url", "answer_id")} + {k: v for k, v in answer.items() if k not in ("iframe_url", "answer_id", "query")} for answer in answers ] @@ -1034,11 +1076,54 @@ def merge_text_chunks(updates: list[dict]) -> list[dict]: return merged +def emit_work_event(update: dict, recorder: StreamRecorder) -> None: + """Record one update as a step of the Agent's work, for the "Show work" section. + + This is the Analytics Agent's own reasoning, as the MCP server reports it - not + Claude's. Steps are text only. A thinking answer in particular must never become + an iframe or join `recorder.answers`: stored answers are resolved by their + position among the non-thinking ones, so counting a thinking answer there would + shift every later chart onto the wrong answer. + """ + kind = update.get("type") + metadata = update_metadata(update) + if kind == "answer": + # The settled answer is the chart itself; only the attempts before it are work. + if not is_thinking_update(update): + return + # Sent even without a title: replay finds this answer by its position among + # the turn's thinking answers, so every one of them needs a step. + recorder.send( + { + "type": "work", + "kind": "query", + "text": (update.get("answer_title") or update.get("title") or "").strip(), + "query": (update.get("answer_query") or metadata.get("sage_query") or "").strip(), + # Live only - `append_work` drops these, since they expire with the + # answer. A replayed step is resolved through /work-answers instead. + "iframe_url": update.get("iframe_url"), + "frame_params": answer_frame_params(update), + } + ) + elif kind in ("step_notification", "notification"): + text = ( + update.get("text") or metadata.get("tool_title") or metadata.get("title") or "" + ).strip() + if text: + recorder.send({"type": "work", "kind": "step", "text": text}) + elif kind in ("text", *TEXT_CHUNK_TYPES) and is_thinking_update(update): + # Chunks are word-sized; keep their spacing, and let append_work join them. + text = update.get("text") or "" + if text: + recorder.send({"type": "work", "kind": "thought", "text": text}) + + def emit_update_events(updates: list[dict], recorder: StreamRecorder) -> None: """Stream Analytics Agent progress to the browser while we poll.""" for update in updates: if not isinstance(update, dict): continue + emit_work_event(update, recorder) kind = update.get("type") if kind == "answer": # The Agent emits an `answer` for each intermediate query it tries while @@ -1380,9 +1465,14 @@ async def persist_turn( return recorder.persisted = True - if recorder.text or recorder.answers: + if recorder.text or recorder.answers or recorder.work: await asyncio.to_thread( - db_add_turn, conv_id, "assistant", recorder.text, storable_answers(recorder.answers) + db_add_turn, + conv_id, + "assistant", + recorder.text, + storable_answers(recorder.answers), + recorder.work, ) if history is not None: await asyncio.to_thread( @@ -1450,42 +1540,90 @@ async def event_stream() -> AsyncGenerator[str, None]: # ── Chat history endpoints ────────────────────────────────────────────────────── -async def ts_answers_per_message(session_id: str) -> list[int] | None: - """How many real answers ThoughtSpot holds for each turn of a conversation. +def ts_headers(token: str) -> dict: + return {"Authorization": f"Bearer {token}", "Accept": "application/json"} - Returns one count per conversation message, oldest first, or None if the - conversation cannot be read. - """ + +async def ts_conversation_messages(session_id: str) -> list[dict] | None: + """A ThoughtSpot conversation's messages, oldest first, or None if unreadable.""" url = ( f"{TS_HOST.rstrip('/')}/api/rest/2.0/ai/agent/conversations/" f"{session_id}/messages" ) try: async with httpx2.AsyncClient(timeout=httpx2.Timeout(10.0, read=45.0)) as client: - response = await client.get( + response = await client.get(url, headers=ts_headers(await server_token())) + if response.status_code != 200: + print(f"[History] getConversation {session_id} -> {response.status_code}") + return None + messages = response.json().get("messages") + return messages if isinstance(messages, list) else [] + except (httpx2.HTTPError, ValueError, TypeError, AttributeError) as exc: + print(f"[History] getConversation {session_id} failed: {type(exc).__name__}: {exc}") + return None + + +def answer_items(message: dict, thinking: bool) -> list[dict]: + """A message's answer items that are (or are not) the Agent's thinking answers. + + Compared with `is` on purpose: an item missing `is_thinking` counts as neither. + Items without an `answer_id` are skipped too. Both are the rules the Visual Embed + SDK applies when it replays stored answers, so our counts match its indexes. + """ + items = message.get("response_items") if isinstance(message, dict) else None + return [ + item + for item in (items or []) + if isinstance(item, dict) + and item.get("type") == "answer" + and item.get("is_thinking") is thinking + and item.get("answer_id") + ] + + +async def ts_load_answer(session_id: str, answer_id: str) -> dict | None: + """Live embed ids for one answer of a conversation, or None if it cannot load. + + The answer object behind a stored URL expires after about 8 hours; loading it + again through the conversation service gives fresh ids for the same answer. This + is the call the Visual Embed SDK makes when it replays a stored answer. + """ + url = ( + f"{TS_HOST.rstrip('/')}/conversation/v2/{session_id}" + f"/message/{quote(answer_id, safe='')}/load/public" + ) + try: + # Loading re-runs the answer's query; 45s+ is common on a loaded cluster. + async with httpx2.AsyncClient(timeout=httpx2.Timeout(10.0, read=120.0)) as client: + response = await client.post( url, headers={ - "Authorization": f"Bearer {await server_token()}", - "Accept": "application/json", + **ts_headers(await server_token()), + "Content-Type": "application/json", + "x-requested-by": "ThoughtSpot", }, + json={"type": "TS_ANSWER"}, ) if response.status_code != 200: - print(f"[History] getConversation {session_id} -> {response.status_code}") + print(f"[History] load answer {answer_id} -> {response.status_code}") return None - return [ - sum( - 1 - for item in (message.get("response_items") or []) - if item.get("type") == "answer" and item.get("is_thinking") is False - ) - for message in (response.json().get("messages") or []) - ] - except (httpx2.HTTPError, ValueError, TypeError) as exc: - print(f"[History] getConversation {session_id} failed: {type(exc).__name__}: {exc}") + answer = (response.json() or {}).get("answer") or {} + except (httpx2.HTTPError, ValueError, TypeError, AttributeError) as exc: + print(f"[History] load answer {answer_id} failed: {type(exc).__name__}: {exc}") return None + ac_state = answer.get("ac_state") or {} + params = { + "session_id": answer.get("session_identifier"), + "gen_no": answer.get("generation_number"), + "ac_session_id": ac_state.get("transaction_identifier"), + "ac_gen_no": ac_state.get("generation_number"), + } + # All four are needed to build the embed route; a partial set renders an error. + return params if all(params.values()) else None -def reconcile_answers(conversation: dict, ts_counts: list[int] | None) -> dict: + +def reconcile_answers(conversation: dict, ts_messages: list[dict] | None) -> dict: """Align a stored conversation with the answers ThoughtSpot actually holds. Two things make the stored copy an unreliable source of truth: @@ -1501,7 +1639,7 @@ def reconcile_answers(conversation: dict, ts_counts: list[int] | None) -> dict: renders them, and the SDK resolves each one by its index. """ turns = conversation.get("turns") or [] - if ts_counts is None: + if ts_messages is None: # Fall back to the stored shape rather than dropping charts entirely. index = 0 for turn in turns: @@ -1515,7 +1653,13 @@ def reconcile_answers(conversation: dict, ts_counts: list[int] | None) -> dict: assistant_turns = [turn for turn in turns if turn.get("role") == "assistant"] index = 0 for position, turn in enumerate(assistant_turns): - expected = ts_counts[position] if position < len(ts_counts) else 0 + message = ts_messages[position] if position < len(ts_messages) else {} + expected = len(answer_items(message, thinking=False)) + # The thinking answers behind this turn's Show work queries, in the order the + # steps were recorded. The client loads one by id when its row is expanded. + turn["work_answer_ids"] = [ + item["answer_id"] for item in answer_items(message, thinking=True) + ] answers = turn.get("answers") or [] # Titles we recorded, padded out to the count ThoughtSpot reports. merged = answers[:expected] + [{} for _ in range(max(0, expected - len(answers)))] @@ -1538,8 +1682,26 @@ async def get_conversation(conv_id: str): raise HTTPException(status_code=404, detail="Conversation not found") session_id = conversation.get("analytical_session_id") - ts_counts = await ts_answers_per_message(session_id) if session_id else None - return reconcile_answers(conversation, ts_counts) + ts_messages = await ts_conversation_messages(session_id) if session_id else None + return reconcile_answers(conversation, ts_messages) + + +@app.get("/api/conversations/{conv_id}/work-answers/{answer_id}") +async def get_work_answer(conv_id: str, answer_id: str): + """Live embed ids for a thinking answer shown in a stored turn's Show work. + + `answer_id` comes from the turn's `work_answer_ids`. The client asks only when a + step is expanded, so a reopened chat does not load every intermediate chart up + front. ThoughtSpot checks the answer belongs to the conversation. + """ + session_id = await asyncio.to_thread(db_get_session_id, conv_id) + if not session_id: + raise HTTPException(status_code=404, detail="Conversation has no ThoughtSpot session") + + frame_params = await ts_load_answer(session_id, answer_id) + if not frame_params: + raise HTTPException(status_code=502, detail="Could not load that answer") + return {"frame_params": frame_params} @app.patch("/api/conversations/{conv_id}")