-
Notifications
You must be signed in to change notification settings - Fork 2
feat(mcp): funnel telemetry, health, rate-limit headers, playground, idempotency #19
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,22 @@ | ||
| import { createServerFn } from "@tanstack/react-start"; | ||
| import { z } from "zod"; | ||
| import { requireSupabaseAuth } from "@/integrations/supabase/auth-middleware"; | ||
|
|
||
| /** Aggregate counts per funnel event over the last N days. Admin only. */ | ||
| export const getMcpFunnelSummary = createServerFn({ method: "GET" }) | ||
| .middleware([requireSupabaseAuth]) | ||
| .inputValidator((d: unknown) => | ||
| z.object({ days: z.number().int().min(1).max(90).optional() }).parse(d ?? {}), | ||
| ) | ||
| .handler(async ({ context, data }) => { | ||
| const { supabase: _sb } = context as any; | ||
| const supabase = _sb as any; | ||
| const { data: rows, error } = await supabase.rpc("mcp_funnel_summary", { | ||
| _days: data.days ?? 7, | ||
| }); | ||
| if (error) throw new Response(error.message, { status: 403 }); | ||
| return { | ||
| days: data.days ?? 7, | ||
| events: (rows ?? []) as Array<{ event: string; count: number; distinct_users: number }>, | ||
| }; | ||
| }); | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,58 @@ | ||
| /** | ||
| * Idempotency helpers for MCP write tools (upload_packages, request_primitive). | ||
| * | ||
| * Agents retry on network blips. Without idempotency, a retried | ||
| * upload_packages creates duplicate drafts and burns the write quota. | ||
| * Callers may pass `idempotency_key` (any opaque string they generate | ||
| * once per logical operation); the server hashes it, scopes it to the | ||
| * user + tool, caches the response for 24h, and returns the cached | ||
| * result on subsequent retries with the same key. | ||
| * | ||
| * Implementation lives in two RPCs (mcp_idempotency_get / | ||
| * mcp_idempotency_put) so the contention is in the database and the | ||
| * tool execute() handlers stay small. | ||
| */ | ||
| import { createHash } from "crypto"; | ||
| import { supabaseAdmin as _supabaseAdmin } from "@/integrations/supabase/client.server"; | ||
| const supabaseAdmin = _supabaseAdmin as any; | ||
|
|
||
| function hashKey(key: string): string { | ||
| return createHash("sha256").update(key).digest("hex"); | ||
| } | ||
|
|
||
| export async function getIdempotent( | ||
| userId: string, | ||
| tool: string, | ||
| key: string | undefined | null, | ||
| ): Promise<unknown | null> { | ||
| if (!key) return null; | ||
| try { | ||
| const { data } = await supabaseAdmin.rpc("mcp_idempotency_get", { | ||
| _user_id: userId, | ||
| _key_hash: hashKey(key), | ||
| _tool: tool, | ||
| } as never); | ||
| return data ?? null; | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
|
|
||
| export async function putIdempotent( | ||
| userId: string, | ||
| tool: string, | ||
| key: string | undefined | null, | ||
| response: unknown, | ||
| ): Promise<void> { | ||
| if (!key) return; | ||
| try { | ||
| await supabaseAdmin.rpc("mcp_idempotency_put", { | ||
| _user_id: userId, | ||
| _key_hash: hashKey(key), | ||
| _tool: tool, | ||
| _response: response as never, | ||
| } as never); | ||
| } catch { | ||
| /* never fail the user's call on idempotency persistence */ | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -6,6 +6,7 @@ const supabaseAdmin = _supabaseAdmin as any; | |
| import { hashToken } from "@/lib/account/tokens.server"; | ||
| import { processBulkUpload } from "@/lib/uploads/uploads.server"; | ||
| import { getGatewayModel } from "@/lib/ai-gateway"; | ||
| import { getIdempotent, putIdempotent } from "@/lib/mcp/idempotency"; | ||
|
|
||
| const json = (v: unknown) => JSON.stringify(v, null, 2); | ||
|
|
||
|
|
@@ -185,24 +186,39 @@ export const searchRegistryTool = defineTool({ | |
| export const requestPrimitiveTool = defineTool({ | ||
| name: "request_primitive", | ||
| description: | ||
| "[PUBLISH] Submit a request for a primitive that does not yet exist. SuperAgentSkill researches and auto-creates it via the proprietary forge pipeline. Requires OAuth.", | ||
| "[PUBLISH] Submit a request for a primitive that does not yet exist. SuperAgentSkill researches and auto-creates it via the proprietary forge pipeline. Requires OAuth. Pass an `idempotency_key` (any opaque string you generate once per request) and retries return the original `request_id` instead of creating duplicates.", | ||
| parameters: z.object({ | ||
| type: z.enum(["skill", "playbook", "soul", "guardrail"]), | ||
| brief: z.string().min(20).max(2000).describe("What the primitive should do, with industry/context"), | ||
| industry: z.string().max(80).optional(), | ||
| idempotency_key: z | ||
| .string() | ||
| .min(8) | ||
| .max(200) | ||
| .optional() | ||
| .describe("Opaque string generated once per logical request. Repeats within 24h return the original response."), | ||
| }), | ||
| execute: async ({ type, brief, industry }) => { | ||
| execute: async ({ type, brief, industry, idempotency_key }, ctx) => { | ||
| const userId = (ctx?.auth?.claims as { user_id?: string } | undefined)?.user_id ?? null; | ||
| if (userId && idempotency_key) { | ||
| const cached = await getIdempotent(userId, "request_primitive", idempotency_key); | ||
| if (cached) return json({ ...(cached as object), replayed: true }); | ||
| } | ||
| const { data, error } = await supabaseAdmin | ||
| .from("package_requests") | ||
| .insert({ kind: type, brief, industry: industry ?? null, status: "queued" }) | ||
| .select("id,status") | ||
| .single(); | ||
| if (error) return json({ error: error.message }); | ||
| return json({ | ||
| const response = { | ||
| request_id: data.id, | ||
| status: data.status, | ||
| note: "Queued for the Super Agent Skill forge pipeline.", | ||
| }); | ||
| }; | ||
| if (userId && idempotency_key) { | ||
| await putIdempotent(userId, "request_primitive", idempotency_key, response); | ||
| } | ||
| return json(response); | ||
| }, | ||
| }); | ||
|
|
||
|
|
@@ -261,7 +277,7 @@ export const getTrustTool = defineTool({ | |
| export const uploadPackagesTool = defineTool({ | ||
| name: "upload_packages", | ||
| description: | ||
| "[PRIVATE UPLOAD] Push local primitive(s) into the author's PRIVATE workspace. Files are normalised by the SkillForge author pipeline and stored as private drafts owned by the token holder — NOT visible in the public marketplace, search, or trust leaderboard. To list a draft for sale on the marketplace, the author must explicitly publish it from the website UI (/account/packages). This MCP tool intentionally has no `publish` parameter so agents cannot expose a user's skill publicly without their consent. Authenticates via the OAuth bearer of the active MCP session — no extra personal token needed.", | ||
| "[PRIVATE UPLOAD] Push local primitive(s) into the author's PRIVATE workspace. Files are normalised by the SkillForge author pipeline and stored as private drafts owned by the token holder — NOT visible in the public marketplace, search, or trust leaderboard. To list a draft for sale on the marketplace, the author must submit it for admin review from the website UI (/account/packages). This MCP tool intentionally has no `publish` parameter so agents cannot expose a user's skill publicly without their consent. Authenticates via the OAuth bearer of the active MCP session — no extra personal token needed. Pass an `idempotency_key` (any opaque string you generate once per upload) so retries on network failures don't create duplicates.", | ||
| parameters: z.object({ | ||
| files: z | ||
| .array( | ||
|
|
@@ -278,26 +294,40 @@ export const uploadPackagesTool = defineTool({ | |
| .min(8) | ||
| .optional() | ||
| .describe("Deprecated. Ignored when the request already carries an OAuth bearer; only used as a fallback for legacy personal MCP tokens."), | ||
| idempotency_key: z | ||
| .string() | ||
| .min(8) | ||
| .max(200) | ||
| .optional() | ||
| .describe("Opaque string generated once per logical upload. Repeats within 24h return the original response and DO NOT re-process the files."), | ||
| }), | ||
| execute: async ({ auth_token, files }, ctx) => { | ||
| execute: async ({ auth_token, files, idempotency_key }, ctx) => { | ||
| const sessionUserId = (ctx?.auth?.claims as { user_id?: string } | undefined)?.user_id ?? null; | ||
| const userId = sessionUserId ?? (auth_token ? await resolveUserFromToken(auth_token) : null); | ||
| if (!userId) | ||
| return json({ | ||
| error: "unauthorized", | ||
| hint: "Connect via OAuth (the host opens https://superagentskill.com/oauth/authorize automatically) — no personal token needed.", | ||
| }); | ||
| if (idempotency_key) { | ||
| const cached = await getIdempotent(userId, "upload_packages", idempotency_key); | ||
| if (cached) return json({ ...(cached as object), replayed: true }); | ||
| } | ||
| try { | ||
| // Always private. Marketplace listing requires an explicit user action in the UI. | ||
| const results = await processBulkUpload(supabaseAdmin as any, userId, files); | ||
|
Comment on lines
+312
to
318
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
This flow checks cache first, performs the write, and only then stores the idempotent response. If two retries with the same Useful? React with 👍 / 👎. |
||
| const ok = results.filter((r) => r.ok).length; | ||
| return json({ | ||
| const response = { | ||
| uploaded: ok, | ||
| failed: results.length - ok, | ||
| visibility: "private_draft", | ||
| next_step: "Open /account/packages on superagentskill.com to list a draft on the marketplace.", | ||
| next_step: "Open /account/packages on superagentskill.com to submit a draft for admin review.", | ||
| results, | ||
| }); | ||
| }; | ||
| if (idempotency_key) { | ||
| await putIdempotent(userId, "upload_packages", idempotency_key, response); | ||
| } | ||
| return json(response); | ||
| } catch (e: any) { | ||
| return json({ error: e?.message ?? "upload_failed" }); | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,55 @@ | ||
| /** | ||
| * Funnel telemetry for the MCP connect flow. The browser calls | ||
| * `recordFunnelEvent` from the consent page, the success page and the | ||
| * /connect landing; the MCP route fires `mcp_first_call` server-side | ||
| * the first time a given user/anon hash invokes a tool. | ||
| * | ||
| * Telemetry is best-effort — the RPC swallows any error so a failed | ||
| * insert never breaks the user's flow. | ||
| */ | ||
| import { createServerFn } from "@tanstack/react-start"; | ||
| import { z } from "zod"; | ||
| import { supabaseAdmin as _supabaseAdmin } from "@/integrations/supabase/client.server"; | ||
| const supabaseAdmin = _supabaseAdmin as any; | ||
|
|
||
| const EVENTS = [ | ||
| "connect_viewed", | ||
| "oauth_authorize_viewed", | ||
| "oauth_authorize_approved", | ||
| "oauth_authorize_denied", | ||
| "oauth_success_shown", | ||
| "oauth_loopback_attempted", | ||
| "oauth_manual_code_copied", | ||
| "oauth_scheme_triggered", | ||
| "mcp_first_call", | ||
| "mcp_first_write", | ||
| "install_button_clicked", | ||
| "pat_minted", | ||
| ] as const; | ||
|
|
||
| const Input = z.object({ | ||
| event: z.enum(EVENTS), | ||
| client_id: z.string().max(200).optional(), | ||
| client_name: z.string().max(200).optional(), | ||
| anon_hash: z.string().max(64).optional(), | ||
| props: z.record(z.string(), z.unknown()).optional(), | ||
| }); | ||
|
|
||
| export const recordFunnelEvent = createServerFn({ method: "POST" }) | ||
| .inputValidator((d: unknown) => Input.parse(d)) | ||
| .handler(async ({ data }) => { | ||
| try { | ||
| await supabaseAdmin.rpc("record_mcp_funnel_event", { | ||
| _event: data.event, | ||
| _client_id: data.client_id ?? null, | ||
| _client_name: data.client_name ?? null, | ||
| _anon_hash: data.anon_hash ?? null, | ||
| _props: data.props ?? {}, | ||
| } as never); | ||
| } catch { | ||
| // never fail user flow on telemetry | ||
| } | ||
| return { ok: true }; | ||
| }); | ||
|
|
||
| export type FunnelEvent = (typeof EVENTS)[number]; |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This server function is labeled “Admin only” but only requires authentication, not admin role. As written, any signed-in user can call it;
mcp_funnel_summaryjust filters rows byauth.uid()and returns an empty result instead of a 403, so/admin/funnelis not actually access-controlled as intended. Use the existingrequireAdminmiddleware (or an explicit role check) here.Useful? React with 👍 / 👎.