Spaces:
Running
Running
File size: 30,840 Bytes
5844a76 6ff6e4c 713e95e 5ce76f0 6ff6e4c 713e95e 5844a76 5ce76f0 5844a76 dcc1160 5844a76 713e95e 5844a76 713e95e 5844a76 4b20a09 5844a76 6ff6e4c 5844a76 5ce76f0 e27ea87 5ce76f0 e27ea87 5ce76f0 e27ea87 6ff6e4c 5ce76f0 6ff6e4c e27ea87 5ce76f0 e27ea87 5ce76f0 e27ea87 5ce76f0 5844a76 1fa623b 5844a76 296fba8 5844a76 1fa623b 296fba8 1fa623b dcc1160 1fa623b 296fba8 1fa623b 296fba8 daaa123 296fba8 daaa123 296fba8 daaa123 620405d cbd0efa 620405d daaa123 620405d daaa123 1fa623b 620405d 1fa623b 5844a76 1fa623b 5844a76 6ff6e4c 713e95e 1fa623b 5844a76 1fa623b 296fba8 1fa623b dcc1160 1fa623b 296fba8 1fa623b 6ff6e4c 296fba8 1fa623b dcc1160 1fa623b 6ff6e4c 713e95e 1fa623b 296fba8 1fa623b 5844a76 cbd0efa 5844a76 c9b5cac 5844a76 1fa623b 5844a76 6ff6e4c 713e95e cbd0efa 5844a76 296fba8 1fa623b cbd0efa 1fa623b dcc1160 296fba8 1fa623b 5844a76 296fba8 713e95e 5844a76 6ff6e4c 3d735c7 5844a76 296fba8 5844a76 cbd0efa 296fba8 5844a76 296fba8 daaa123 620405d 5844a76 296fba8 cbd0efa 620405d cbd0efa 5844a76 daaa123 5844a76 6ff6e4c 713e95e 5844a76 1fa623b 5844a76 1fa623b 5844a76 daaa123 5844a76 6ff6e4c 713e95e 5844a76 1fa623b 5844a76 1fa623b 713e95e 1fa623b c9b5cac 5844a76 daaa123 5844a76 6ff6e4c 713e95e 6ff6e4c 296fba8 5844a76 296fba8 5844a76 1fa623b daaa123 296fba8 1fa623b 296fba8 1fa623b 6ff6e4c 1fa623b 713e95e 1fa623b 296fba8 1fa623b 5844a76 1fa623b 5844a76 6ff6e4c 713e95e 5844a76 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 | """
mcp_server.py — Demo Wire MCP server (v1)
Exposes DemoPrep's demo-build pipeline as an MCP tool so an external agent
(AgentSpot / Spotter Assistant) can build a full ThoughtSpot demo from a prospect
*brief* without a human clicking through the Gradio UI.
This is a THIN ADAPTER. It does not reimplement any business logic — it drives the
exact same `ChatDemoInterface` controller the Gradio app uses, headlessly, skipping
the research phase by injecting the brief as the research output. All build logic
stays in one place (see docs/SINGLE_PIPELINE.md and CLAUDE.md).
v1 scope (simplest thing that works):
* ONE blocking tool — build_demo_from_brief(...) runs the whole pipeline
(blueprint -> dataset -> DDL -> Snowflake -> TS model -> liveboard) and returns
the resulting IDs when done.
* Everything is fixed/server-side; per-call config (env, model, sharing) is a
later version.
* Builds run for minutes. If the caller's connection times out, the build still
finishes server-side and DemoPrep's Slack notification delivers the liveboard
URL. (Async / job-backed delivery is the planned next step.)
Server-side config (env vars):
MCP_ACCESS_TOKEN (required in http mode) shared bearer secret gating the endpoint
TS_ENV_URL_DEFAULT ThoughtSpot URL used when a build doesn't pass ts_url
TS_USER_DEFAULT owner email used when a build doesn't pass owner_email
TS_ENV_<n>_LABEL / TS_ENV_<n>_URL / TS_ENV_<n>_KEY_VAR
the environment table — SAME numbered triplets the app's
dropdown uses. The trusted-auth key for a build is found
by matching its ts_url host against TS_ENV_<n>_URL.
MCP_TRANSPORT "stdio" (default, local dev) | "http" (HF Space)
MCP_HTTP_PORT http port (default 7860)
Also read from env: SUPABASE_URL / SUPABASE_ANON_KEY (Snowflake creds) + ANTHROPIC_API_KEY.
Both the TS target and the owner are resolved per build: the request may pass
ts_url / owner_email, else they fall back to TS_ENV_URL_DEFAULT / TS_USER_DEFAULT.
Local smoke test (needs the project venv with requirements installed):
TS_ENV_URL_DEFAULT=https://sebe.thoughtspotstaging.cloud python mcp_server.py
(the TS_ENV_* table comes from .env; the URL default just picks which entry)
Then point an MCP client (e.g. `mcp dev mcp_server.py` / MCP Inspector) at stdio.
"""
from __future__ import annotations
import os
import sys
import time
import types
import uuid
import threading
import collections
# --- env bootstrap: order matters --------------------------------------------
# .env first so MCP_* / TS_ENV_* / Supabase creds are available.
from dotenv import load_dotenv
load_dotenv()
# Owner identity is resolved PER BUILD (mirrors the TS target): the request may
# pass owner_email; otherwise it falls back to TS_USER_DEFAULT. Nothing is
# hardcoded, and boot does not require it — a build fails loud if neither is set.
os.environ.setdefault("DEMOPREP_NO_AUTH", "true")
_DEFAULT_OWNER = (os.getenv("TS_USER_DEFAULT") or "").strip()
if _DEFAULT_OWNER:
# Best-effort acting user for any headless path with no per-build owner.
os.environ.setdefault("DEMOPREP_DEV_USER_EMAIL", _DEFAULT_OWNER)
def _resolve_owner(owner_email: str) -> str:
"""Owner of the objects a build creates: the per-request owner_email, else
TS_USER_DEFAULT. Fails loud if neither is set (no silent blank owner)."""
owner = (owner_email or os.getenv("TS_USER_DEFAULT") or "").strip()
if not owner:
raise RuntimeError("No owner: pass owner_email in the request or set TS_USER_DEFAULT.")
return owner
# Snowflake creds + SNOWFLAKE_DATABASE from Supabase admin settings -> os.environ.
# Must run once at startup; the deploy path reads these via get_admin_setting/env.
from supabase_client import inject_admin_settings_to_env, log_mcp_payload
inject_admin_settings_to_env()
# Heavy imports (pull gradio etc.) — safe headless. DEMOPREP_NO_AUTH is set above,
# which must happen BEFORE importing chat_interface.
from chat_interface import ChatDemoInterface
from demo_builder_class import DemoBuilder
from demo_personas import parse_use_case, get_use_case_config
from llm_config import DEFAULT_LLM_MODEL
def _ts_env_table() -> list[tuple[str, str, str]]:
"""Read the numbered TS_ENV_<n>_LABEL/URL/KEY_VAR triplets — the SAME
environment table the app's dropdown uses (single source of truth for
which ThoughtSpot instances exist and their trusted-auth keys).
KEY_VAR normally holds the trusted-auth key itself. An `ENV:<NAME>` value
dereferences another env var instead — for referencing a secret that
already exists under its own name (e.g. TECHPARTNERS_TA_KEY on the Space)
without copying its value into the table."""
table = []
i = 1
while True:
label = (os.getenv(f"TS_ENV_{i}_LABEL") or "").strip()
if not label:
break
url = (os.getenv(f"TS_ENV_{i}_URL") or "").strip().rstrip("/")
key = (os.getenv(f"TS_ENV_{i}_KEY_VAR") or "").strip()
if key.startswith("ENV:"):
key = (os.getenv(key[4:].strip()) or "").strip()
table.append((label, url, key))
i += 1
return table
def _host_of(url: str) -> str:
"""Hostname of a URL; tolerates scheme-less input and strips any port."""
if "://" in url:
url = url.split("://", 1)[1]
return url.split("/", 1)[0].split(":", 1)[0].lower()
def _resolve_ts_target(ts_url: str) -> tuple[str, str]:
"""Resolve the (url, trusted_auth_key) for a build.
The URL comes from the build request; if absent it falls back to
TS_ENV_URL_DEFAULT. The trusted-auth key comes from the TS_ENV_<n>_* table
entry whose URL host matches. Fails loud on a missing URL, an unknown
instance, or a blank key (no silent blanks). Resolved per build (not at
boot) so the server stays up regardless of which environment a given
build targets.
"""
url = (ts_url or os.getenv("TS_ENV_URL_DEFAULT") or "").strip().rstrip("/")
if not url:
raise RuntimeError(
"No ThoughtSpot URL: pass ts_url in the request or set TS_ENV_URL_DEFAULT."
)
if "://" not in url:
url = f"https://{url}"
table = _ts_env_table()
if not table:
raise RuntimeError(
"No TS_ENV_<n>_LABEL/URL/KEY_VAR entries configured — the MCP server "
"reads the same numbered environment table as the app."
)
host = _host_of(url)
for label, env_url, key in table:
if _host_of(env_url) == host:
if not key:
raise RuntimeError(f"TS_ENV entry {label!r} matches {url} but its KEY_VAR is blank.")
return url, key
known = ", ".join(f"{label} ({env_url})" for label, env_url, _ in table)
raise RuntimeError(
f"No TS_ENV entry matches {url}. Known environments: {known}. "
f"Add a TS_ENV_<n>_LABEL/URL/KEY_VAR triplet for this instance."
)
from mcp.server.fastmcp import FastMCP
from mcp.server.transport_security import TransportSecuritySettings
# Bounded concurrency: the build pipeline is already concurrency-safe (the app/QA
# have run concurrent builds for months), so this is only a resource guardrail — a
# semaphore caps how many run at once so a caller can't OOM the Space by firing
# dozens; a call past the cap gets "busy". Owner/env/creds are per-instance, so
# concurrent builds share no mutable process state.
_MAX_CONCURRENT_BUILDS = max(1, int(os.getenv("MCP_MAX_CONCURRENT_BUILDS", "3")))
_build_sem = threading.BoundedSemaphore(_MAX_CONCURRENT_BUILDS)
# Lightweight hit log so we can watch incoming calls (who + what args). Writes to
# stderr (the server log) and, if MCP_HIT_LOG is set, appends to that file too.
_HIT_LOG = (os.getenv("MCP_HIT_LOG") or "").strip()
def _log_hit(msg: str) -> None:
line = f"[{time.strftime('%H:%M:%S')}] {msg}"
print(line, file=sys.stderr, flush=True)
if _HIT_LOG:
try:
with open(_HIT_LOG, "a") as fh:
fh.write(line + "\n")
except Exception:
pass
# --- build status tracking (for the status tool / polling long builds) ---
_SERVER_STARTED = time.time()
_state_lock = threading.Lock()
_build_count = 0
_active_builds = {} # run_id -> live build dict; supports multiple concurrent builds
_last_build = None # summary of the most recent completed build
# Compact result of each completed build, keyed by run_id (newest last), capped so
# a client that fired several concurrent builds can collect EACH result by run_id
# instead of racing the single _last_build slot. In-memory; resets on restart.
_recent_builds = collections.OrderedDict()
_RECENT_BUILDS_CAP = 50
def _set_progress(run_id: str, phase: str | None = None, detail: str | None = None) -> None:
"""Update a running build's coarse phase and/or latest raw progress line."""
with _state_lock:
b = _active_builds.get(run_id)
if b is not None:
if phase:
b["phase"] = phase
if detail is not None:
b["detail"] = detail
def _progress_text(item) -> str:
"""Best-effort human-readable line from a generator yield (str | dict | tuple)."""
if isinstance(item, dict):
for k in ("message", "text", "content", "status", "detail"):
v = item.get(k)
if isinstance(v, str) and v.strip():
item = v
break
else:
item = str(item)
elif isinstance(item, (tuple, list)):
item = next((x for x in item if isinstance(x, str) and x.strip()), str(item))
return " ".join(str(item).split())[:200] # collapse whitespace + cap length
def _phase_from_text(text: str) -> str | None:
"""Map a raw progress line to a granular phase label (None = keep current phase)."""
t = text.lower()
# [async] import lines say "…polling…" as normal operation — check them
# before the retry bucket so a healthy import isn't mislabeled "retrying".
if "[async]" in t:
return "importing tables"
if any(k in t for k in ("504", "retry", "retrying", "throttl", "timed out", "timeout")):
return "waiting on ThoughtSpot (retrying)"
if "connection" in t:
return "creating connection"
if "tml" in t or "enhance" in t or "post-process" in t or "applying style" in t:
return "liveboard TML phase"
if "liveboard" in t or "pinboard" in t:
return "creating liveboard"
if "model" in t:
return "model created" if any(k in t for k in ("created", "ready", "guid")) else "creating model"
if "verif" in t and "schema" in t:
return "verifying schema"
if "table" in t and ("creat" in t or "import" in t):
return "creating tables"
if any(k in t for k in ("import", "loading", "populat", "copy into", "batch", "chunk", " rows")):
return "loading data"
if "schema" in t and "creat" in t:
return "creating schema"
return None
# The app yields a periodic UI "spinner" banner while deploy_all runs in its
# worker thread. It is a heartbeat, not progress — and it happens to contain the
# word "connection" (…Creating connection & tables…), so if it reaches
# _set_progress it both freezes `phase` on "creating connection" and clobbers the
# real per-line detail arriving via on_progress. Drop it on the MCP side.
_BANNER_MARKERS = (
"deployment in progress",
"thoughtspot deploying",
"starting thoughtspot deployment",
)
def _is_heartbeat_banner(text: str) -> bool:
t = text.lower()
return any(m in t for m in _BANNER_MARKERS)
def _apply_progress(run_id: str, item) -> None:
"""Fold one generator yield / deploy line into the build's live status.
Ignores the UI heartbeat banner so the real per-line progress (delivered via
on_progress from the deploy thread) drives `phase`/`detail` instead of being
overwritten every 2s.
"""
text = _progress_text(item)
if not text or _is_heartbeat_banner(text):
return
_set_progress(run_id, phase=_phase_from_text(text), detail=text)
# DNS-rebinding protection validates the Host header and only trusts localhost
# by default, which returns 421 when reached through a tunnel or an hf.space host.
# We gate access on the bearer token instead (and this is server-to-server, not
# browser-driven), so disable the host check to allow any public host.
mcp = FastMCP(
"demoprep",
transport_security=TransportSecuritySettings(enable_dns_rebinding_protection=False),
)
@mcp.tool()
def ping() -> dict:
"""Health / connectivity check — returns server identity + config, no side effects.
Use this to confirm a client (e.g. AgentSpot) can reach and invoke the server
without triggering a full build.
"""
_log_hit("CALL ping")
return {
"ok": True,
"server": "demoprep-mcp v1",
"default_ts_url": os.getenv("TS_ENV_URL_DEFAULT", ""),
"default_owner": os.getenv("TS_USER_DEFAULT", ""),
"tools": ["ping", "status", "build_demo_from_brief"],
}
@mcp.tool()
def status() -> dict:
"""Report server + build status: whether the server is up, whether a build is
running (with its phase and elapsed time), and the most recent completed build.
Use this to poll a long build: call build_demo_from_brief, then poll status()
to watch progress and retrieve the result even if the original call's
connection dropped.
"""
_log_hit("CALL status")
now = time.time()
with _state_lock:
active = [dict(b) for b in _active_builds.values()]
last = dict(_last_build) if _last_build else None
recent = [dict(r) for r in _recent_builds.values()]
count = _build_count
for b in active:
b["elapsed_seconds"] = round(now - b.pop("started_at", now), 1)
active.sort(key=lambda b: b.get("elapsed_seconds", 0), reverse=True)
return {
"server": "up",
"default_ts_url": os.getenv("TS_ENV_URL_DEFAULT", ""),
"busy": len(active) > 0,
"active_count": len(active),
"capacity": _MAX_CONCURRENT_BUILDS,
"active_builds": active,
# Back-compat for single-build pollers: present only when exactly one runs.
"current_build": active[0] if len(active) == 1 else None,
"last_build": last,
# Compact result of each recently completed build (newest last), keyed data
# so a client that fired several builds can collect EACH result by run_id.
"recent_builds": recent,
"builds_started": count,
"uptime_seconds": round(now - _SERVER_STARTED, 1),
}
def _run_build(run_id: str, brief: str, company_name: str, use_case: str, company_url: str,
ts_target_url: str, ts_auth_key: str, owner: str) -> dict:
"""Drive the controller headlessly (runs in a BACKGROUND THREAD; the caller has
already registered the build in _active_builds and holds a _build_sem slot (released when this
returns). Mirrors tests/newvision_sample_runner.py but injects `brief` in place of
the research phase. Returns a structured dict and never raises — failures come back
as status 'failed' / 'partial'."""
started = time.time()
ddl_text = "" # surfaced via status/result once generated; stays "" if the build fails earlier
# Full-brief capture + durable persist — ONCE per accepted build (this only runs for
# builds that actually got a slot, so busy-retries of a queued build never duplicate
# rows). Logged with the run_id so the durable record correlates to the build.
_log_hit(f"PAYLOAD_BRIEF_BEGIN run_id={run_id} company={company_name!r} use_case={use_case!r} "
f"url={company_url!r} ts_url={ts_target_url!r} owner={owner!r} brief_len={len(brief or '')}")
print(brief or "", file=sys.stderr, flush=True)
_log_hit("PAYLOAD_BRIEF_END")
_persisted = log_mcp_payload({
"run_id": run_id, "company_name": company_name, "use_case": use_case,
"company_url": company_url, "ts_url": ts_target_url, "owner_email": owner, "brief": brief,
})
_log_hit(f"PAYLOAD_PERSISTED run_id={run_id} supabase={_persisted}")
def result(status: str, dc: dict | None = None, schema: str | None = None,
error: str | None = None) -> dict:
dc = dc or {}
out = {
"run_id": run_id,
"status": status, # success | partial | failed
"schema": dc.get("schema") or schema or "",
"model_guid": dc.get("model_guid", ""),
"liveboard_guid": dc.get("liveboard_guid", ""),
"model_url": dc.get("model_url", ""),
"liveboard_url": dc.get("liveboard_url", ""),
"ts_environment": ts_target_url,
"owner_email": owner,
"ddl": ddl_text,
"warnings": dc.get("warnings", []),
"errors": ([error] if error else dc.get("errors", [])),
"elapsed_seconds": round(time.time() - started, 1),
}
global _last_build
with _state_lock:
_last_build = {
"run_id": out["run_id"], "status": out["status"], "schema": out["schema"],
"model_url": out["model_url"], "liveboard_url": out["liveboard_url"],
"errors": out["errors"], "elapsed_seconds": out["elapsed_seconds"],
"ddl": out["ddl"],
"finished_at": time.strftime("%H:%M:%S"),
}
# Per-run result (compact, no ddl) so concurrent callers can collect EACH
# build by run_id without racing the single _last_build slot.
_recent_builds[run_id] = {
"run_id": out["run_id"], "status": out["status"], "schema": out["schema"],
"company_name": company_name,
"model_url": out["model_url"], "liveboard_url": out["liveboard_url"],
"warnings": out["warnings"], "errors": out["errors"],
"elapsed_seconds": out["elapsed_seconds"], "finished_at": time.strftime("%H:%M:%S"),
}
while len(_recent_builds) > _RECENT_BUILDS_CAP:
_recent_builds.popitem(last=False) # drop oldest
_active_builds.pop(run_id, None)
return out
controller = None
try:
# (i) controller — acting user = the per-build owner (per-instance; no
# process-global write, so concurrent builds never race on the owner).
controller = ChatDemoInterface(user_email=owner)
# (ii) settings: model + fixed TS env (exact key names per the wiring trace)
controller.settings["model"] = controller.settings.get("model") or DEFAULT_LLM_MODEL
controller.settings["thoughtspot_url"] = ts_target_url
controller.settings["thoughtspot_trusted_auth_key"] = ts_auth_key
# Always share the created model + liveboard with the operator, so every demo is
# visible no matter who ran it. The OWNER is unchanged — whoever was passed as
# owner_email (the build authenticates AS that user, so they own the objects);
# this just adds the operator as a viewer. Override target via MCP_SHARE_WITH.
controller.settings["share_with"] = os.getenv("MCP_SHARE_WITH", "mike.boone@thoughtspot.com")
# vertical / function / use_case_config exactly as the runner does
controller.vertical, controller.function = parse_use_case(use_case or "")
controller.use_case_config = get_use_case_config(
controller.vertical or "Generic", controller.function or "Generic"
)
# (iii) demo_builder with the BRIEF injected in place of research.
db = DemoBuilder(use_case=use_case, company_url=company_url)
db.company_analysis_results = brief # component field
db.combined_research_results = brief # <-- the field build_demo actually reads
db.company_summary = brief # <-- liveboard Spotter story reads this (Gotcha 2)
# Force the exact display name — extract_company_name() otherwise parses the
# domain from company_url. Minimal shim standing in for a scraped website.
db.website_data = types.SimpleNamespace(
title=company_name, url=company_url, text="", css_links=[], logo_candidates=[]
)
controller.demo_builder = db
controller.generic_use_case_context = ""
# (iv) DDL — returns a (response, ddl) tuple; NOT a generator.
_set_progress(run_id, phase="building dataset + DDL", detail="")
resp, ddl_text = controller.run_ddl_creation()
if not ddl_text or "CREATE TABLE" not in ddl_text.upper():
return result("failed", error=f"DDL generation failed: {str(resp)[:500]}")
# Surface the DDL immediately — it exists ~5 min in, well before the ~15-min
# TS deploy — so a caller polling status() gets the schema as soon as it's ready.
with _state_lock:
b = _active_builds.get(run_id)
if b is not None:
b["ddl"] = ddl_text
# (v) Snowflake load, then ThoughtSpot. Both are generators — draining them
# IS what runs the work. Decoupled from validation_mode: drain the Snowflake
# generator, read the schema it set, then run the TS deploy ourselves.
_set_progress(run_id, phase="loading Snowflake", detail="")
for _item in controller.run_deployment_streaming():
_apply_progress(run_id, _item)
schema = getattr(controller, "_deployed_schema_name", None)
if not schema:
return result(
"failed",
schema=getattr(controller, "_last_schema_name", None),
error="Snowflake load did not complete (no deployed schema).",
)
_set_progress(run_id, phase="deploying ThoughtSpot model + liveboard", detail="")
# Real-time status: on_progress fires from the deploy thread for EVERY
# progress line (incl. the [async] import + model/liveboard steps), so
# status.current_build.detail tracks live instead of freezing between yields.
def _dp(m):
_apply_progress(run_id, m)
for _item in controller._run_thoughtspot_deployment(schema, company_name, use_case, on_progress=_dp):
_dp(_item)
# (vi) structured result from the completion record.
dc = getattr(controller, "deployment_completion", None)
if not dc:
# deploy_all raised before the completion record was written — surface
# partial success: the Snowflake schema exists even if TS didn't finish.
return result(
"partial", schema=schema,
error="ThoughtSpot deploy did not complete; Snowflake schema exists.",
)
return result("success" if dc.get("success") else "partial", dc=dc, schema=schema)
except Exception as e: # never leak a raw exception to the MCP caller
schema = getattr(controller, "_deployed_schema_name", None) if controller else None
return result(
"partial" if schema else "failed",
schema=schema,
error=f"{type(e).__name__}: {e}",
)
@mcp.tool()
def build_demo_from_brief(
brief: str,
company_name: str,
use_case: str = "",
company_url: str = "",
ts_url: str = "",
owner_email: str = "",
) -> dict:
"""Start a full ThoughtSpot demo build from a prospect brief.
Skips DemoPrep's own research: the `brief` IS the research context. The build
(dataset -> DDL -> Snowflake -> model -> liveboard) runs in a BACKGROUND THREAD and
takes ~15-20 minutes, so this returns IMMEDIATELY with a run_id — it does NOT block.
Poll status() for progress (phase + elapsed) and the final result (schema, model,
liveboard URL); the build survives even if this call's connection drops.
Args:
brief: Prospect narrative — pain points, what they're evaluating, industry
context, goals. Becomes the research context the demo is built from.
company_name: Display name for the demo (e.g. "Acme Corporation").
use_case: (optional) The analytics story / label (e.g. "Retail Sales").
If omitted, the demo is authored purely from the brief and
labeled "Custom Analytics".
company_url: Optional company URL (used for context/branding; not scraped).
ts_url: Optional ThoughtSpot instance URL to deploy into. If omitted,
falls back to TS_ENV_URL_DEFAULT; the trusted-auth key is chosen
by instance (sebe vs SE primary).
owner_email: Optional ThoughtSpot user to own the created objects. If
omitted, falls back to TS_USER_DEFAULT.
Returns:
dict: {status: "started"|"busy"|"failed", run_id, message}. Poll status()
for progress and the eventual result.
"""
_log_hit(
f"CALL build_demo_from_brief company={company_name!r} use_case={use_case!r} "
f"url={company_url!r} ts_url={ts_url!r} owner={owner_email!r} brief_len={len(brief or '')} "
f"brief_preview={(brief or '')[:200]!r}"
)
# NOTE: the full-brief capture + durable persist happen ONCE per ACCEPTED build,
# inside _run_build — NOT here. Doing it here (before the slot/busy check) meant a
# client re-polling a queued build sprayed a duplicate durable row on every retry.
# The lightweight CALL line above still records every attempt in the (ephemeral) log.
if not brief or not brief.strip():
return {"status": "failed", "errors": ["brief is required"]}
if not (company_name or "").strip():
return {"status": "failed", "errors": ["company_name is required"]}
# Resolve the ThoughtSpot target up front (URL from the request or the default,
# key chosen by instance) so a bad URL / missing key fails fast — before we
# take the single build slot.
try:
ts_target_url, ts_auth_key = _resolve_ts_target(ts_url)
owner = _resolve_owner(owner_email)
except RuntimeError as e:
return {"status": "failed", "errors": [str(e)]}
# Bounded concurrency: grab a build slot (non-blocking). Only when ALL slots are
# in use does a call get 'busy' — otherwise it runs alongside the others.
if not _build_sem.acquire(blocking=False):
return {
"status": "busy",
"errors": [f"All {_MAX_CONCURRENT_BUILDS} build slots are in use. Poll status(); retry when one frees."],
}
run_id = uuid.uuid4().hex[:12]
b, c, u, url = brief.strip(), company_name.strip(), (use_case or "Custom Analytics").strip(), (company_url or "").strip()
global _build_count
with _state_lock:
_build_count += 1
_active_builds[run_id] = {
"run_id": run_id, "company_name": c, "use_case": u,
"ts_url": ts_target_url,
"started_at": time.time(), "phase": "starting",
}
# Run the build OFF the request thread so the server stays responsive to status()
# polls throughout — a synchronous build monopolizes the interpreter and starves them.
def _worker() -> None:
try:
_run_build(run_id, b, c, u, url, ts_target_url, ts_auth_key, owner)
finally:
_build_sem.release()
threading.Thread(target=_worker, daemon=True, name=f"build-{run_id}").start()
return {
"status": "started",
"run_id": run_id,
"message": "Build started in the background (~15-20 min). Poll status() for progress and the final result (schema, model, liveboard URL).",
}
def _run_http(token: str) -> None:
"""Serve over streamable HTTP behind a shared-bearer gate.
The auth check is a PURE-ASGI wrapper, NOT Starlette's BaseHTTPMiddleware —
the latter buffers responses and breaks the streamable-HTTP SSE stream. This
checks the bearer on every HTTP request and otherwise passes the raw ASGI
through untouched.
"""
import uvicorn
inner = mcp.streamable_http_app() # verified accessor on mcp 1.28.1
class _BearerGate:
def __init__(self, app):
self.app = app
async def __call__(self, scope, receive, send):
if scope.get("type") == "http":
path = scope.get("path", "")
if path in ("/", "/health"):
# Unauthenticated liveness check so the HF Space reports healthy.
# The MCP protocol itself lives at /mcp behind the bearer gate.
await send({
"type": "http.response.start",
"status": 200,
"headers": [(b"content-type", b"application/json")],
})
await send({"type": "http.response.body",
"body": b'{"status":"ok","service":"demoprep-mcp"}'})
return
headers = dict(scope.get("headers") or [])
ip = (headers.get(b"x-forwarded-for", b"") or headers.get(b"cf-connecting-ip", b"")).decode().split(",")[0].strip()
_log_hit(f"HIT {scope.get('method', '?')} {path} "
f"auth={'Y' if b'authorization' in headers else 'N'} ip={ip or '?'}")
if headers.get(b"authorization", b"").decode() != f"Bearer {token}":
await send({
"type": "http.response.start",
"status": 401,
"headers": [(b"content-type", b"application/json")],
})
await send({"type": "http.response.body", "body": b'{"error":"unauthorized"}'})
return
await self.app(scope, receive, send)
uvicorn.run(_BearerGate(inner), host="0.0.0.0", port=int(os.getenv("MCP_HTTP_PORT", "7860")))
def main() -> None:
transport = (os.getenv("MCP_TRANSPORT") or "stdio").strip().lower()
print(
f"[mcp_server] ready — default_ts_url={os.getenv('TS_ENV_URL_DEFAULT', '(unset)')} "
f"default_owner={os.getenv('TS_USER_DEFAULT', '(unset)')} transport={transport}",
file=sys.stderr, flush=True,
)
if transport in ("stdio", ""):
mcp.run(transport="stdio")
elif transport in ("http", "streamable-http", "streamable_http"):
token = (os.getenv("MCP_ACCESS_TOKEN") or "").strip()
if not token: # public endpoint must be gated
raise RuntimeError("MCP_ACCESS_TOKEN is required in http mode (public endpoint).")
_run_http(token)
else:
raise RuntimeError(f"Unknown MCP_TRANSPORT: {transport!r} (use 'stdio' or 'http').")
if __name__ == "__main__":
main()
|