demoprep / tests /e2e_mcp.py
mikeboone's picture
fix(test): make e2e_mcp resilient to transient network blips
095d53e
Raw
History Blame Contribute Delete
19.7 kB
#!/usr/bin/env python3
"""
e2e_mcp.py — end-to-end test that drives the LIVE MCP server through N real demo
builds (default 8), exactly as AgentSpot does: fire `build_demo_from_brief`, poll
`status()`, collect each result by run_id from `recent_builds`, and grade it.
This replaces the local-pipeline e2e for the MCP path. One run exercises the whole
stack end to end — brief -> blueprint -> DDL -> Snowflake load -> connectivity
guard -> ThoughtSpot model -> liveboard — plus the server's bounded concurrency
and the 13122 connectivity guard.
The brief set deliberately includes two shapes that used to fail schema
validation (error 13122):
* "Summit Mutual" — TWO subject areas (underwriting + workforce): the guard
should keep the primary star and set the other aside.
* "GridLink" — time-series heavy: prone to an orphan date dimension the
guard should drop.
Usage:
source demoprep/bin/activate
python tests/e2e_mcp.py # all briefs against the default endpoint
python tests/e2e_mcp.py --limit 2 # first 2 briefs (quick check)
python tests/e2e_mcp.py --poll 20 --timeout 6000
Env (falls back to .env):
MCP_ENDPOINT default https://thoughtspot-demoprep-mcp.hf.space/mcp
MCP_ACCESS_TOKEN bearer token (REQUIRED)
MCP_TS_URL optional ts_url passed to each build (else server default)
MCP_OWNER_EMAIL optional owner_email passed to each build (else server default)
Exit code: 0 if every build PASSED, else 1.
"""
from __future__ import annotations
import argparse
import json
import os
import sys
import time
from pathlib import Path
import anyio
from dotenv import load_dotenv
ROOT = Path(__file__).resolve().parent.parent
load_dotenv(ROOT / ".env")
from mcp.client.streamable_http import streamablehttp_client
from mcp.client.session import ClientSession
DEFAULT_ENDPOINT = "https://thoughtspot-demoprep-mcp.hf.space/mcp"
# --------------------------------------------------------------------------- #
# Brief corpus. Each brief is intentionally rich (multiple dimensions with >=4
# concrete values, clear measures, a persona and a story) — thin briefs fail
# blueprint validation.
# --------------------------------------------------------------------------- #
BRIEFS = [
{
"name": "Northwind Retail",
"company_name": "Northwind Retail",
"use_case": "Retail Sales",
"brief": (
"Northwind Retail is a national apparel and home-goods retailer evaluating ThoughtSpot to give "
"merchandising and store-operations leaders self-service analytics. Personas: VP of Merchandising, "
"Regional Store Director, Category Manager. They sell across product categories (Apparel, Footwear, "
"Home, Accessories, Electronics), channels (In-Store, Online, Mobile App, Marketplace), regions "
"(Northeast, Southeast, Midwest, West, Southwest) and customer segments (New, Returning, Loyalty, VIP). "
"Key questions: which categories and regions drive margin, how promotions lift units vs. erode margin, "
"and how online is cannibalizing store sales. Measures: net sales, units sold, gross margin %, discount %, "
"average order value, return rate, same-store-sales growth. Story: a strong holiday quarter hides a "
"margin problem in one region driven by over-discounting."
),
},
{
"name": "Meridian Health",
"company_name": "Meridian Health System",
"use_case": "Healthcare Operations",
"brief": (
"Meridian Health System runs a network of hospitals and clinics and wants ThoughtSpot for operations and "
"patient-flow analytics. Personas: Chief Operating Officer, Department Administrator, Nurse Manager. "
"Dimensions: facility (Main Campus, North Clinic, South Clinic, Children's, Rehab), department (Emergency, "
"Surgery, Cardiology, Oncology, Orthopedics, Radiology), payer (Medicare, Medicaid, Commercial, "
"Self-Pay), admission type (Emergency, Elective, Transfer, Observation). Measures: patient volume, average "
"length of stay, bed occupancy %, ED wait time, readmission rate, cost per case, revenue per case. Story: "
"ED wait times spike at specific facilities and correlate with downstream readmissions."
),
},
{
"name": "Atlas Freight",
"company_name": "Atlas Freight",
"use_case": "Logistics Operations",
"brief": (
"Atlas Freight is a third-party logistics provider evaluating ThoughtSpot for shipment and hub performance. "
"Personas: VP of Operations, Regional Hub Manager, Carrier Account Lead. Dimensions: hub (Atlanta, Dallas, "
"Chicago, Newark, Los Angeles, Seattle), service level (Ground, Express, Overnight, Freight, Same-Day), "
"carrier (in-house fleet plus four partner carriers), lane region (Domestic-East, Domestic-West, Cross-Border, "
"International). Measures: shipments, on-time delivery %, average transit days, cost per shipment, damage "
"rate, SLA compliance %, capacity utilization %. Story: one hub is dragging network on-time performance "
"because of a single underperforming partner carrier."
),
},
{
"name": "Vantage Capital",
"company_name": "Vantage Capital",
"use_case": "Financial Services Risk",
"brief": (
"Vantage Capital is a commercial lender evaluating ThoughtSpot for portfolio and credit-risk analytics. "
"Personas: Chief Risk Officer, Portfolio Manager, Credit Analyst. Dimensions: product (Term Loan, Line of "
"Credit, Equipment Finance, Commercial Mortgage, SBA), risk grade (AAA, AA, A, BBB, BB, B, CCC), industry "
"(Manufacturing, Retail, Healthcare, Technology, Real Estate, Energy), region (Northeast, Southeast, "
"Midwest, West). Measures: outstanding balance, exposure at default, delinquency rate, charge-off rate, "
"weighted-average risk grade, net interest margin, provision coverage. Story: delinquency is quietly "
"concentrating in one industry within a single risk grade."
),
},
{
"name": "Streamly",
"company_name": "Streamly",
"use_case": "Media Engagement",
"brief": (
"Streamly is a subscription streaming service evaluating ThoughtSpot for content and engagement analytics. "
"Personas: VP of Content, Growth Lead, Retention Analyst. Dimensions: content genre (Drama, Comedy, "
"Documentary, Kids, Sports, Reality), device (Smart TV, Mobile, Web, Tablet, Console), plan tier (Free, "
"Basic, Standard, Premium), subscriber region (North America, Europe, LATAM, APAC). Measures: monthly active "
"users, watch hours, completion rate, churn rate, trial-to-paid conversion, ARPU, content cost per watch "
"hour. Story: a hit documentary drives sign-ups that churn fast because onboarding fails to surface similar titles."
),
},
{
"name": "Cloudscape",
"company_name": "Cloudscape",
"use_case": "SaaS Product & Revenue",
"brief": (
"Cloudscape is a B2B SaaS collaboration platform evaluating ThoughtSpot for product-usage and revenue "
"analytics. Personas: VP of Product, Head of Customer Success, RevOps Lead. Dimensions: plan (Free, Team, "
"Business, Enterprise), feature area (Docs, Chat, Video, Automation, Admin), industry (Technology, Finance, "
"Healthcare, Education, Retail), account region (AMER, EMEA, APAC, LATAM). Measures: active accounts, seats, "
"feature adoption %, weekly active users, expansion revenue, gross retention %, net revenue retention, "
"support tickets per account. Story: Enterprise accounts with low automation adoption are the ones at churn risk."
),
},
{
"name": "Summit Mutual",
"company_name": "Summit Mutual Insurance",
"use_case": "Insurance Underwriting",
"brief": (
"Summit Mutual Insurance is a mid-market insurer evaluating ThoughtSpot. They want TWO stories in one demo. "
"(1) Underwriting performance: personas Chief Underwriting Officer and Underwriting Manager; dimensions "
"underwriter, policy type (Auto, Home, Commercial, Umbrella, Marine), workflow stage (Submitted, Quoted, "
"Bound, Declined, Referred), region (Northeast, Southeast, Midwest, West); measures submissions, quote ratio, "
"hit ratio, written premium, loss ratio, rate change %. (2) Workforce / return-to-office: personas Chief "
"People Officer and Office Lead; dimensions employee, office location, department (Claims, Underwriting, "
"Sales, IT, Finance), employment type (Full-Time, Part-Time, Contract, Seasonal); measures headcount, PTO "
"utilization, office attendance rate, unplanned absence rate, training hours. Story: underwriting throughput "
"dips in offices with the lowest attendance."
),
},
{
"name": "GridLink",
"company_name": "GridLink Exchange",
"use_case": "Data Network & Platform Utilization",
"brief": (
"GridLink Exchange operates a B2B data-exchange network and is evaluating ThoughtSpot for network health and "
"platform-utilization analytics, including embedded dashboards for their participants. Personas: VP of "
"Network Operations, Platform Product Manager, Participant Success Lead. Dimensions: network participant "
"(dozens of member firms), service line (Clearing, Settlement, Reporting, Reference Data, Messaging), "
"geography (North America, Europe, APAC, LATAM), transaction quality band (Excellent, Good, Fair, Poor). "
"Track monthly trends over the last two years. Measures: transaction volume, transaction value, settlement "
"success rate, data-quality score, adoption rate, dashboard views, self-serve query rate, active users. "
"Story: a few high-volume participants with declining data-quality scores drag network-wide settlement success."
),
},
]
# --------------------------------------------------------------------------- #
# MCP client — one short session per call (robust for long-running polling).
# --------------------------------------------------------------------------- #
async def _call_async(endpoint: str, token: str, tool: str, args: dict):
async with streamablehttp_client(endpoint, headers={"Authorization": f"Bearer {token}"}) as (r, w, _):
async with ClientSession(r, w) as s:
await s.initialize()
res = await s.call_tool(tool, args)
return getattr(res, "structuredContent", None) or json.loads(res.content[0].text)
def call(endpoint: str, token: str, tool: str, args: dict | None = None,
retries: int = 4, backoff: float = 3.0):
"""One tool call, tolerant of transient network blips. Over a long run the
hf.space endpoint occasionally drops a connection / times out TLS; retry a few
times before giving up so a single blip doesn't kill the whole harness."""
last = None
for attempt in range(retries):
try:
return anyio.run(_call_async, endpoint, token, tool, args or {})
except Exception as e: # ConnectTimeout, ReadError, ExceptionGroup, etc.
last = e
if attempt < retries - 1:
time.sleep(backoff)
raise last
def _now() -> str:
return time.strftime("%H:%M:%S")
def log(msg: str) -> None:
print(f"[{_now()}] {msg}", flush=True)
def grade(res: dict) -> dict:
"""PASS = build succeeded, produced a liveboard, and hit no 13122 schema error."""
status = (res or {}).get("status")
errors = (res or {}).get("errors") or []
warnings = (res or {}).get("warnings") or []
has_13122 = any("13122" in str(e) for e in errors)
guard_acted = [w for w in warnings if "Connectivity guard" in str(w)]
passed = status == "success" and bool((res or {}).get("liveboard_url")) and not has_13122
return {
"grade": "PASS" if passed else "FAIL",
"status": status,
"model_url": (res or {}).get("model_url", ""),
"liveboard_url": (res or {}).get("liveboard_url", ""),
"elapsed_seconds": (res or {}).get("elapsed_seconds"),
"has_13122": has_13122,
"guard_acted": guard_acted,
"errors": errors[:3],
}
def run(endpoint: str, token: str, briefs: list, ts_url: str, owner: str,
poll: int, timeout: int) -> int:
log(f"endpoint: {endpoint}")
st = call(endpoint, token, "status")
cap = st.get("capacity", "?")
if "recent_builds" not in st:
log("WARNING: server has no `recent_builds` field — results may be missed under "
"concurrency. Deploy the recent_builds change for reliable collection.")
log(f"server up; capacity={cap}; firing {len(briefs)} builds")
pending = list(range(len(briefs))) # indices not yet started
inflight: dict[str, int] = {} # run_id -> brief index
results: dict[int, dict] = {} # brief index -> graded result
run_ids: dict[int, str] = {}
deadline = time.time() + timeout
while (pending or inflight) and time.time() < deadline:
# Fire as many pending builds as the server will accept right now.
while pending:
idx = pending[0]
b = briefs[idx]
try:
resp = call(endpoint, token, "build_demo_from_brief", {
"brief": b["brief"], "company_name": b["company_name"],
"use_case": b.get("use_case", ""),
"ts_url": ts_url, "owner_email": owner,
})
except Exception as e:
log(f" fire failed for [{b['name']}] ({type(e).__name__}); retry next cycle")
break
sc = resp.get("status")
if sc == "started":
rid = resp["run_id"]
inflight[rid] = idx
run_ids[idx] = rid
pending.pop(0)
log(f" started [{b['name']}] run={rid} ({len(inflight)} in flight)")
elif sc == "busy":
break # all slots full — wait and poll
else:
results[idx] = {"grade": "FAIL", "status": sc or "failed",
"errors": resp.get("errors", []), "has_13122": False,
"guard_acted": [], "model_url": "", "liveboard_url": ""}
run_ids[idx] = resp.get("run_id", "")
pending.pop(0)
log(f" FAILED to start [{b['name']}]: {resp.get('errors')}")
# Poll for completions + live progress. A failed poll must NOT end the run —
# builds keep running server-side and results persist in recent_builds.
try:
st = call(endpoint, token, "status")
except Exception as e:
log(f" status poll failed ({type(e).__name__}); retry next cycle "
"(builds keep running server-side)")
time.sleep(poll)
continue
recent = {r["run_id"]: r for r in st.get("recent_builds", [])}
if st.get("last_build"): # fallback for servers without recent_builds
recent.setdefault(st["last_build"]["run_id"], st["last_build"])
for rid in list(inflight):
if rid in recent:
idx = inflight.pop(rid)
results[idx] = grade(recent[rid])
g = results[idx]
extra = f" guard={g['guard_acted']}" if g["guard_acted"] else ""
log(f" done [{briefs[idx]['name']}] -> {g['grade']} ({g['status']}, "
f"{g.get('elapsed_seconds')}s){extra}")
if inflight:
active = {a["run_id"]: a for a in st.get("active_builds", [])}
parts = []
for rid, idx in inflight.items():
a = active.get(rid, {})
parts.append(f"{briefs[idx]['name']}:{a.get('phase','?')}({a.get('elapsed_seconds','?')}s)")
log(f" in flight: {', '.join(parts)} | pending={len(pending)}")
if pending or inflight:
time.sleep(poll)
# Anything still outstanding at the deadline is a timeout.
for rid, idx in inflight.items():
results[idx] = {"grade": "FAIL", "status": "timeout", "errors": ["exceeded harness timeout"],
"has_13122": False, "guard_acted": [], "model_url": "", "liveboard_url": ""}
for idx in pending:
results[idx] = {"grade": "FAIL", "status": "never-started", "errors": ["harness timed out before start"],
"has_13122": False, "guard_acted": [], "model_url": "", "liveboard_url": ""}
return _report(briefs, results, run_ids, endpoint)
def _report(briefs: list, results: dict, run_ids: dict, endpoint: str) -> int:
passed = sum(1 for r in results.values() if r["grade"] == "PASS")
total = len(briefs)
print("\n" + "=" * 78)
print(f"e2e_mcp results — {passed}/{total} PASSED")
print("=" * 78)
print(f"{'BRIEF':22} {'GRADE':6} {'STATUS':10} {'13122':6} LIVEBOARD/ERROR")
print("-" * 78)
for idx in range(total):
b = briefs[idx]
r = results.get(idx, {"grade": "FAIL", "status": "?", "has_13122": False,
"liveboard_url": "", "errors": []})
detail = r.get("liveboard_url") or (r.get("errors") or [""])[0]
print(f"{b['name'][:22]:22} {r['grade']:6} {str(r['status'])[:10]:10} "
f"{'YES' if r.get('has_13122') else '-':6} {str(detail)[:30]}")
if r.get("guard_acted"):
for w in r["guard_acted"]:
print(f"{'':40} guard: {str(w)[:60]}")
print("=" * 78)
ts = time.strftime("%Y%m%d_%H%M%S")
out_dir = ROOT / "tests" / "quality_results"
out_dir.mkdir(parents=True, exist_ok=True)
out_path = out_dir / f"e2e_mcp_{ts}.json"
payload = {
"endpoint": endpoint, "timestamp": ts, "passed": passed, "total": total,
"results": {briefs[i]["name"]: {**results.get(i, {}), "run_id": run_ids.get(i, "")}
for i in range(total)},
}
out_path.write_text(json.dumps(payload, indent=2, default=str))
print(f"report: {out_path}")
return 0 if passed == total else 1
def main() -> int:
ap = argparse.ArgumentParser(description="Drive the live MCP server through N real demo builds.")
ap.add_argument("--limit", type=int, default=0, help="run only the first N briefs (0 = all)")
ap.add_argument("--poll", type=int, default=15, help="status poll interval seconds")
ap.add_argument("--timeout", type=int, default=5400, help="overall timeout seconds (default 90m)")
ap.add_argument("--endpoint", default=os.getenv("MCP_ENDPOINT", DEFAULT_ENDPOINT))
args = ap.parse_args()
token = (os.getenv("MCP_ACCESS_TOKEN") or "").strip()
if not token:
print("ERROR: MCP_ACCESS_TOKEN is not set (env or .env).", file=sys.stderr)
return 2
briefs = BRIEFS[: args.limit] if args.limit else BRIEFS
ts_url = os.getenv("MCP_TS_URL", "").strip()
owner = os.getenv("MCP_OWNER_EMAIL", "").strip()
return run(args.endpoint, token, briefs, ts_url, owner, args.poll, args.timeout)
if __name__ == "__main__":
sys.exit(main())