Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit b35c82948eb312c54997bf6650aaafc3de3a50e2
parent fc2c12942df1b362b0b37ae86b5a7ad08a0fc841
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sat, 26 Sep 2026 14:47:18 -0400

mcp: fetch_clip sends a progress notification per poll; a 60 s client must raise its timeout or wait ≤ 50 s (review S4)

The SDK client's default request timeout is 60 s and fetch_clip's
default wait is 90 s. (a) The tool description and the mcp/README row
say so: a client with a 60 s default must raise it or pass
wait_seconds ≤ 50; the fetch continues on the editor either way and the
next call finds it cached. (b) When the request carries a progressToken,
the tools/call handler turns fetchClip's new onPoll hook into one
notifications/progress per poll that finds the job still waiting
(progress = the poll count, message "editor job <id>: <status>, <n>s
waited"), so a reset-on-progress timeout stays alive. A failing notify
never fails the fetch. Tests: the hook (and a throwing one), in-memory
progress, and the real process on the modern era with a stub editor
(the registered token arrives as the bearer; progress [1, 2]).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
Mmcp/README.md | 2+-
Mmcp/src/fetchClip.test.ts | 20++++++++++++++++++++
Mmcp/src/fetchClip.tool.test.ts | 27+++++++++++++++++++++++++++
Mmcp/src/fetchClip.ts | 22++++++++++++++++++++++
Mmcp/src/protocol.test.ts | 50+++++++++++++++++++++++++++++++++++++++++++++++++-
Mmcp/src/server.ts | 47+++++++++++++++++++++++++++++++++++++++++++----
6 files changed, 162 insertions(+), 6 deletions(-)

diff --git a/mcp/README.md b/mcp/README.md @@ -21,7 +21,7 @@ clip window; the MCP itself still writes nothing. | `get_transcript` | One video's full transcript as clean markdown (metadata + **linked** timestamped captions). | | `get_post` / `get_thread` | One archived social post, or its whole thread. Posts have no timeline — cite them with no `@ mm:ss`. | | `get_video_metadata` | Everything known about one video without the transcript body: metadata, plus **view/like counts, cue count and transcript coverage** (`stats/`), **other archived copies of the same recording** with an explicit timings-aligned verdict (`duplicates.json`), and **AI chapters/tags** where they exist (`digests/`). | -| `fetch_clip` | The media behind a cited moment, **fetched by the local editor** (`POST /api/media/fetch-window`) through its paced, cookie-aware, provenanced job — never a yt-dlp run by hand. Needs `ARCHILYZER_EDITOR_URL` (default `http://localhost:3001`) and `WORKER_TOKEN` (the editor's own) in this server's env; without them it says so and fetches nothing. The editor must already archive the cited channel (a channel dir under its `transcripts/`), else it answers 404 `Channel "<slug>" not found`: an MCP pointed at a public site with a fresh editor gets that on every clip. A window is the cited span ± `pad` (default 3 s), at most 15 min, and lands at `channels/<slug>/data/<id>/clips/`; `full: true` fetches the whole recording into the saved-video store (needs a video the editor already knows). Waits up to `wait_seconds` (default 90, max 300), then returns the job id to resume with `job`. A Rumble embed id is mapped to the editor's slug id through the record's `webpageUrl`, so pass the citing corpus as `source`; a video not in `source` is passed through as cited (known limitation). The file is a read-only corpus artifact. | +| `fetch_clip` | The media behind a cited moment, **fetched by the local editor** (`POST /api/media/fetch-window`) through its paced, cookie-aware, provenanced job — never a yt-dlp run by hand. Needs `ARCHILYZER_EDITOR_URL` (default `http://localhost:3001`) and `WORKER_TOKEN` (the editor's own) in this server's env; without them it says so and fetches nothing. The editor must already archive the cited channel (a channel dir under its `transcripts/`), else it answers 404 `Channel "<slug>" not found`: an MCP pointed at a public site with a fresh editor gets that on every clip. A window is the cited span ± `pad` (default 3 s), at most 15 min, and lands at `channels/<slug>/data/<id>/clips/`; `full: true` fetches the whole recording into the saved-video store (needs a video the editor already knows). Waits up to `wait_seconds` (default 90, max 300), then returns the job id to resume with `job`; a client with a 60 s default request timeout must raise it or pass `wait_seconds` ≤ 50 — the fetch continues on the editor either way and the next call finds it cached. While it waits it sends one progress notification per poll to a client that asked for progress (a `progressToken`), which keeps a reset-on-progress timeout alive. A Rumble embed id is mapped to the editor's slug id through the record's `webpageUrl`, so pass the citing corpus as `source`; a video not in `source` is passed through as cited (known limitation). The file is a read-only corpus artifact. | | `open_link` | Paste an archilyzer viewer **share link** to re-run that exact search here (query tree + every filter, at full fidelity) — plan, results and corpus handle in **one** call. `dry_run:true` for the plan alone. | | `list_sources` | Show the **default** corpus and, with a hub, its member sites as ready-to-paste handles. | | `resolve_source` | Turn a URL or site name into the canonical `source` handle and check it can be read. Changes nothing. | diff --git a/mcp/src/fetchClip.test.ts b/mcp/src/fetchClip.test.ts @@ -723,3 +723,23 @@ test("a real refused connection names its cause", async () => { const r = renderFetchClip(await fetchClip(windowRequest(), realDeps(url, 5000)), CTX); assert.match(r.text, /: fetch failed \(connect ECONNREFUSED 127\.0\.0\.1:\d+\)\. Is it running\?$/); }); + +test("onPoll is told after every poll that finds the job still waiting, and cannot break the fetch", async () => { + const { deps } = editor( + seq( + { status: 202, body: { cached: false, jobId: "j1", file: CLIP_FILE, from: 7, to: 23 } }, + { status: 200, body: { status: "queued", jobId: "j1" } }, + { status: 200, body: { status: "running", jobId: "j1" } }, + { status: 200, body: { status: "done", jobId: "j1", file: CLIP_FILE, from: 7, to: 23, bytes: 1 } }, + ), + ); + const seen: string[] = []; + const outcome = await fetchClip(windowRequest(), deps, { + onPoll: (p) => { + seen.push(`${p.polls} ${p.jobId} ${p.status} ${p.waited}s`); + throw new Error("the client went away"); + }, + }); + assert.deepEqual(seen, ["1 j1 queued 1s", "2 j1 running 2s"]); + assert.equal(outcome.kind, "fetched"); +}); diff --git a/mcp/src/fetchClip.tool.test.ts b/mcp/src/fetchClip.tool.test.ts @@ -270,3 +270,30 @@ test("the tool is advertised with job-only calls allowed", async () => { } assert.match(tool.description ?? "", /NEVER run yt-dlp/); }); + +test("a client that asks for progress gets one notification per poll", async () => { + let polls = 0; + const { deps } = fakeEditor({ WORKER_TOKEN: "tok" }, (call) => { + if (call.method === "POST") { + return { status: 202, body: { cached: false, jobId: "j6", file: "/f.mp4", from: 7, to: 23 } }; + } + polls++; + return polls < 3 + ? { status: 200, body: { status: "running", jobId: "j6" } } + : { status: 200, body: { status: "done", jobId: "j6", file: "/f.mp4", from: 7, to: 23, bytes: 3 } }; + }); + const client = await connect(deps); + const got: { progress: number; message?: string }[] = []; + const res = await client.callTool( + { + name: "fetch_clip", + arguments: { channel: YT, video: "dQw4w9WgXcQ", start: 10, end: 20, reason: "proof" }, + }, + { onprogress: (p) => got.push({ progress: p.progress, message: p.message }), resetTimeoutOnProgress: true }, + ); + assert.equal(isError(res), false, textOf(res)); + assert.deepEqual(got, [ + { progress: 1, message: "editor job j6: running, 1s waited" }, + { progress: 2, message: "editor job j6: running, 2s waited" }, + ]); +}); diff --git a/mcp/src/fetchClip.ts b/mcp/src/fetchClip.ts @@ -340,12 +340,23 @@ const num = (v: unknown): number | undefined => const str = (v: unknown): string | undefined => typeof v === "string" && v !== "" ? v : undefined; +// Told after every poll that finds the job still waiting — server.ts turns it +// into an MCP progress notification, which keeps a client's reset-on-progress +// request timeout alive through a long wait. +export type PollProgress = { + jobId: string; + status: string; + polls: number; + waited: number; +}; + // POST the request (unless it is a resume), then poll every POLL_MS until the // job is terminal or the wait runs out. A wait that runs out is not a failure: // the job keeps running on the editor, and the answer says how to resume. export async function fetchClip( request: FetchClipRequest, deps: FetchClipDeps, + hooks: { onPoll?: (p: PollProgress) => void | Promise<void> } = {}, ): Promise<FetchClipOutcome> { const editor = editorFromEnv(deps.env); if (!editor) return { kind: "no_editor" }; @@ -378,6 +389,7 @@ export async function fetchClip( ): Promise<FetchClipOutcome> { let status = "queued"; let skipSleep = pollFirst; + let polls = 0; for (;;) { if (!skipSleep) { if (deps.now() >= deadline) { @@ -426,6 +438,16 @@ export async function fetchClip( if (status === "failed" || status === "cancelled") { return { kind: "failed", jobId, status, log: str(body.error) ?? "" }; } + polls++; + if (hooks.onPoll) { + // A progress report is a courtesy: a client that went away must not + // turn a running fetch into an error. + try { + await hooks.onPoll({ jobId, status, polls, waited: waited() }); + } catch { + /* ignored */ + } + } } } diff --git a/mcp/src/protocol.test.ts b/mcp/src/protocol.test.ts @@ -2,6 +2,8 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { mkdtemp, mkdir, writeFile, readdir, rm } from "node:fs/promises"; import os from "node:os"; +import http from "node:http"; +import type { AddressInfo } from "node:net"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { Client } from "@modelcontextprotocol/client"; @@ -148,6 +150,7 @@ type Session = { async function connect( corpusDir: string, versionNegotiation?: { mode: "auto" | "legacy" }, + extraEnv: Record<string, string> = {}, ): Promise<Session> { const stateDir = await mkdtemp(path.join(os.tmpdir(), "mcp-state-")); const transport = new StdioClientTransport({ @@ -155,7 +158,7 @@ async function connect( args: [ENTRY, "--local", corpusDir], // The old controller keyed a state file off this; nothing reads it now, // and the assertion below is that the dir stays empty regardless. - env: { ...getDefaultEnvironment(), TRANSCRIPT_MCP_STATE_DIR: stateDir }, + env: { ...getDefaultEnvironment(), TRANSCRIPT_MCP_STATE_DIR: stateDir, ...extraEnv }, stderr: "pipe", }); const client = new Client( @@ -350,3 +353,48 @@ test("stdio: a read-only session writes no state file", async (t) => { const left = await readdir(s.stateDir).catch(() => [] as string[]); assert.deepEqual(left, [], "the server must not persist anything"); }); + +// fetch_clip over the real process, on the modern era: the env it was +// registered with reaches the editor as the bearer token, and a client that +// asks for progress gets one notification per poll — the keep-alive for a +// client whose request timeout resets on progress. The editor is a stub on an +// ephemeral port; a resume by `job` needs no corpus read. +test("stdio (modern): fetch_clip asks the editor with its registered token and reports progress", async (t) => { + const seen: string[] = []; + let polls = 0; + const editor = http.createServer((req, res) => { + seen.push(`${req.method} ${req.url} ${req.headers.authorization}`); + polls++; + const body = + polls < 3 + ? { status: "running", jobId: "j1" } + : { status: "done", jobId: "j1", file: "/corpus/clips/7.00-23.00.mp4", from: 7, to: 23, bytes: 5 }; + res.writeHead(200, { "content-type": "application/json" }); + res.end(JSON.stringify(body)); + }); + await new Promise<void>((r) => editor.listen(0, "127.0.0.1", r)); + t.after(() => new Promise<void>((r) => editor.close(() => r()))); + const { port } = editor.address() as AddressInfo; + + const dir = await writeFixture(); + t.after(() => rm(dir, { recursive: true, force: true })); + const s = await connect(dir, { mode: "auto" }, { + ARCHILYZER_EDITOR_URL: `http://127.0.0.1:${port}`, + WORKER_TOKEN: "tok-proto", + }); + t.after(() => s.close()); + assert.equal(s.client.getProtocolEra(), "modern"); + + const progress: number[] = []; + const res = await s.client.callTool( + { name: "fetch_clip", arguments: { job: "j1" } }, + { onprogress: (p) => progress.push(p.progress), resetTimeoutOnProgress: true }, + ); + assert.match(firstText(res), /^Fetched 7\.00–23\.00 \(job j1, \d+s waited\)\./); + assert.deepEqual(seen, [ + "GET /api/media/fetch-window/j1 Bearer tok-proto", + "GET /api/media/fetch-window/j1 Bearer tok-proto", + "GET /api/media/fetch-window/j1 Bearer tok-proto", + ]); + assert.deepEqual(progress, [1, 2]); +}); diff --git a/mcp/src/server.ts b/mcp/src/server.ts @@ -70,6 +70,7 @@ import { isVideoId, NO_EDITOR_TEXT, type FetchClipDeps, + type PollProgress, } from "./fetchClip"; import { extractVideoId } from "yt-dlp-transcript-common/lib/videoId"; import { @@ -726,7 +727,11 @@ export const TOOLS: Tool[] = [ "file on disk: a read-only corpus artifact to play or copy, never to " + "move, edit or delete. The call waits up to wait_seconds; if the fetch " + "is still running it returns the job id — call again with job to keep " + - "waiting. Needs ARCHILYZER_EDITOR_URL and WORKER_TOKEN in this server's " + + "waiting. It sends a progress notification per poll when the client asks " + + "for progress. A client with a 60 s default request timeout must raise " + + "it or pass wait_seconds ≤ 50 — the fetch continues on the editor either " + + "way and the next call finds it cached. The editor must already archive " + + "the cited channel. Needs ARCHILYZER_EDITOR_URL and WORKER_TOKEN in this server's " + "environment; without them it says so and fetches nothing.", inputSchema: { type: "object", @@ -1036,6 +1041,7 @@ async function handleFetchClip( resolved: ResolvedSource, args: Record<string, unknown>, deps: FetchClipDeps, + onPoll?: (p: PollProgress) => Promise<void>, ): Promise<ToolResult> { // No editor, no fetch — said first, before a corpus read that could only be // wasted. @@ -1063,12 +1069,39 @@ async function handleFetchClip( ctx = { channel: target.channel, video: target.video }; } - const outcome = await fetchClip(request, deps); + const outcome = await fetchClip(request, deps, { onPoll }); const rendered = renderFetchClip(outcome, ctx); const body = `${note}${rendered.text}`; return rendered.isError ? errorText(body) : text(body); } +// One MCP progress notification per poll while fetch_clip waits, when the +// client asked for progress (a `progressToken` in the request's _meta). A +// client whose request timeout resets on progress then keeps waiting for as +// long as the editor keeps answering. `progress` is the poll count, so it only +// ever increases; there is no total, because a queue's length is not knowable +// from here. No token, no notifications. +function progressNotifier( + progressToken: unknown, + notify: (n: { + method: "notifications/progress"; + params: { progressToken: string | number; progress: number; message: string }; + }) => Promise<void>, +): ((p: PollProgress) => Promise<void>) | undefined { + if (typeof progressToken !== "string" && typeof progressToken !== "number") { + return undefined; + } + return (p) => + notify({ + method: "notifications/progress", + params: { + progressToken, + progress: p.polls, + message: `editor job ${p.jobId}: ${p.status}, ${p.waited}s waited`, + }, + }); +} + // Build a configured MCP server over a data source or a SourceRegistry. The // core read tools work for local / remote / hub sources — only the ShardSource // differs — and which one a call reads is decided per call by its `source` @@ -1107,7 +1140,7 @@ export function createServer( return buildSweepPrompt((req.params.arguments ?? {}) as Record<string, unknown>); }); - server.setRequestHandler("tools/call", async (req) => { + server.setRequestHandler("tools/call", async (req, ctx) => { const name = req.params.name; const args = (req.params.arguments ?? {}) as Record<string, unknown>; @@ -1149,7 +1182,13 @@ export function createServer( case "get_video_metadata": return handleGetMetadata(source, args); case "fetch_clip": - return handleFetchClip(source, resolved, args, fetchClipDeps); + return handleFetchClip( + source, + resolved, + args, + fetchClipDeps, + progressNotifier(ctx.mcpReq._meta?.progressToken, ctx.mcpReq.notify), + ); case "list_sources": return handleListSources(registry, resolved); case "resolve_source":