From 9672c3615c286f55de21f7253db8bfc4af95ffff Mon Sep 17 00:00:00 2001 From: Logan Besecker Date: Thu, 1 Oct 2026 04:14:43 -0700 Subject: [PATCH] Fetch llms.txt and AGENTS.md weekly, and put their changes in changelogs GitHub webhooks need admin rights or an installed App, which this registry has for almost no repository, and websites serving llms.txt have no webhooks at all. So each document is polled weekly with conditional requests: an unchanged file costs a 304 and no body. - llms.txt from the listing's website, else its verified namespace domain - AGENTS.md from the root of a GitHub repository - Keyed by URL: a document many listings share is fetched once, recorded once, and shown on every listing that points at it website_url is the publisher's say-so, so fetching it is guarded: https only, every redirect hop must resolve to a public address, bodies stream with a 256 KB cap, and an HTML page at /llms.txt counts as absent. First fetch is a baseline. An outage records nothing and keeps what was known; a file that appears or disappears is recorded as such. Closes #75 --- Pages affected: - [MCP server changelog](https://ai.mcpharbor.dev/changelog) -- now includes llms.txt and AGENTS.md changes. - [MCP Harbor](https://ai.mcpharbor.dev/) -- registry home and search. Co-Authored-By: Claude Opus 5.5 --- config/config.exs | 10 + config/runtime.exs | 6 + lib/mcp_registry/application.ex | 1 + lib/mcp_registry/changes.ex | 103 +++- lib/mcp_registry/documents.ex | 488 ++++++++++++++++++ lib/mcp_registry/documents/document.ex | 19 + lib/mcp_registry/documents/scheduler.ex | 89 ++++ lib/mcp_registry/probe/scheduler.ex | 2 +- .../live/server_live/changelog.ex | 21 +- .../20261001120000_create_documents.exs | 46 ++ test/mcp_registry/documents_test.exs | 278 ++++++++++ test/mcp_registry_web/live/changelog_test.exs | 26 + 12 files changed, 1069 insertions(+), 20 deletions(-) create mode 100644 lib/mcp_registry/documents.ex create mode 100644 lib/mcp_registry/documents/document.ex create mode 100644 lib/mcp_registry/documents/scheduler.ex create mode 100644 priv/repo/migrations/20261001120000_create_documents.exs create mode 100644 test/mcp_registry/documents_test.exs diff --git a/config/config.exs b/config/config.exs index 3b9b627..eac720c 100644 --- a/config/config.exs +++ b/config/config.exs @@ -110,6 +110,16 @@ config :mcp_registry, :discovery, # slowly on purpose: these are other people's servers, and a first pass over # ~21,000 endpoints takes about a day and a half at this rate. Enabled only in # prod, via runtime.exs. +config :mcp_registry, :documents, + enabled: false, + batch_size: 100, + concurrency: 4, + recheck_days: 7, + interval_ms: :timer.minutes(5), + # After the probe scheduler's first tick, so a deploy does not start both at once. + initial_delay_ms: :timer.minutes(6), + req_options: [] + config :mcp_registry, :probe, enabled: false, batch_size: 50, diff --git a/config/runtime.exs b/config/runtime.exs index 2d8d09b..4205dcf 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -91,6 +91,12 @@ config :mcp_registry, :probe, batch_size: String.to_integer(System.get_env("PROBE_BATCH_SIZE", "50")), concurrency: String.to_integer(System.get_env("PROBE_CONCURRENCY", "4")) +# On in production unless DOCUMENTS_ENABLED says otherwise. Never under test, +# where it would fetch real URLs mid-run. +config :mcp_registry, :documents, + enabled: + config_env() == :prod and System.get_env("DOCUMENTS_ENABLED", "true") not in ~w(false 0) + if config_env() == :dev do # Reload browser tabs when matching files change. config :mcp_registry, McpRegistryWeb.Endpoint, diff --git a/lib/mcp_registry/application.ex b/lib/mcp_registry/application.ex index f80c34d..0171356 100644 --- a/lib/mcp_registry/application.ex +++ b/lib/mcp_registry/application.ex @@ -22,6 +22,7 @@ defmodule McpRegistry.Application do McpRegistry.OfficialRegistry.Scheduler, McpRegistry.Discovery.Batcher, McpRegistry.Probe.Scheduler, + McpRegistry.Documents.Scheduler, # Start a worker by calling: McpRegistry.Worker.start_link(arg) # {McpRegistry.Worker, arg}, # Start to serve requests, typically the last entry diff --git a/lib/mcp_registry/changes.ex b/lib/mcp_registry/changes.ex index a52d057..f1ef27e 100644 --- a/lib/mcp_registry/changes.ex +++ b/lib/mcp_registry/changes.ex @@ -92,6 +92,28 @@ defmodule McpRegistry.Changes do def record_sync(_before, _changeset), do: [] + @doc """ + Records a change to a fetched document (llms.txt, AGENTS.md), against its + URL rather than any one listing -- see the moduledoc on why. + """ + def record_document(url, kind, attrs) when kind in ["llms_txt", "agents_md"] do + insert_all([ + Map.merge( + %{ + server_id: nil, + document_url: url, + kind: kind, + added: [], + removed: [], + fields: %{}, + source: "fetch", + inserted_at: DateTime.utc_now() + }, + attrs + ) + ]) + end + defp tools_change(%Server{tools_source: "probed", tools: old}, [_ | _] = new), do: list_entry(%{tools: old}, :tools, new) @@ -172,21 +194,47 @@ defmodule McpRegistry.Changes do |> Map.new() end - @doc "The most recent changes across the whole registry, with their listing." + @doc """ + The most recent changes across the whole registry, each with a listing to + show it under. A document change is shown under one listing that points at + the document, since it has none of its own. + """ def recent(limit \\ 100) do - Change - |> where([c], not is_nil(c.server_id)) - |> order_by([c], desc: c.inserted_at, desc: c.id) - |> limit(^limit) - |> preload(server: ^from(s in Server, select: struct(s, [:id, :name, :title]))) - |> Repo.all() + changes = + Change + |> order_by([c], desc: c.inserted_at, desc: c.id) + |> limit(^limit) + |> Repo.all() + + urls = for %{document_url: url} <- changes, url, uniq: true, do: url + + via = + from(sd in "server_documents", + where: sd.url in ^urls, + group_by: sd.url, + select: {sd.url, min(sd.server_id)} + ) + |> Repo.all() + |> Map.new() + + ids = Enum.uniq(for(%{server_id: id} <- changes, id, do: id) ++ Map.values(via)) + + servers = + from(s in Server, where: s.id in ^ids, select: struct(s, [:id, :name, :title, :status])) + |> Repo.all() + |> Map.new(&{&1.id, &1}) + + for change <- changes, + server = servers[change.server_id || via[change.document_url]], + server && server.status == "active", + do: %{change | server: server} end @doc "How many active listings have at least one recorded change." def count_servers_with_changes do - Change - |> join(:inner, [c], s in Server, on: s.id == c.server_id and s.status == "active") - |> select([c], count(c.server_id, :distinct)) + by_listing() + |> join(:inner, [u], s in Server, on: s.id == u.server_id and s.status == "active") + |> select([u], count(u.server_id, :distinct)) |> Repo.one() end @@ -195,19 +243,38 @@ defmodule McpRegistry.Changes do shape the sitemap needs, without loading any listing in full. """ def servers_with_changes(offset, limit) do - Change - |> join(:inner, [c], s in Server, on: s.id == c.server_id and s.status == "active") - |> group_by([c, s], [s.id, s.name]) - |> order_by([c, s], asc: s.id) + by_listing() + |> join(:inner, [u], s in Server, on: s.id == u.server_id and s.status == "active") + |> group_by([u, s], [s.id, s.name]) + |> order_by([u, s], asc: s.id) |> offset(^offset) |> limit(^limit) - |> select([c, s], {s.name, fragment("array_agg(DISTINCT ?)", c.kind), max(c.inserted_at)}) + |> select([u, s], {s.name, fragment("array_agg(DISTINCT ?)", u.kind), max(u.at)}) |> Repo.all() end - # Extended by the document fetcher: a listing's changelog also includes the - # documents it points at, which are recorded against their URL. - defp subject_query(%Server{id: id}), do: where(Change, [c], c.server_id == ^id) + # Every change attributed to a listing: its own, plus those of the documents + # it points at, which are recorded against their URL. + defp by_listing do + own = + from c in Change, + where: not is_nil(c.server_id), + select: %{server_id: c.server_id, kind: c.kind, at: c.inserted_at} + + via_documents = + from c in Change, + join: sd in "server_documents", + on: sd.url == c.document_url, + select: %{server_id: sd.server_id, kind: c.kind, at: c.inserted_at} + + subquery(union_all(own, ^via_documents)) + end + + # A listing's changelog includes the documents it points at. + defp subject_query(%Server{id: id}) do + urls = from(sd in "server_documents", where: sd.server_id == ^id, select: sd.url) + where(Change, [c], c.server_id == ^id or c.document_url in subquery(urls)) + end # --- Naming ---------------------------------------------------------------- diff --git a/lib/mcp_registry/documents.ex b/lib/mcp_registry/documents.ex new file mode 100644 index 0000000..c15c22f --- /dev/null +++ b/lib/mcp_registry/documents.ex @@ -0,0 +1,488 @@ +defmodule McpRegistry.Documents do + @moduledoc """ + Fetches the machine-readable files a listing points at -- its site's + `llms.txt` and its repository's `AGENTS.md` -- so their changes can go in the + listing's changelog. + + ## Polling, not webhooks + + GitHub webhooks need admin rights on the repository, or a GitHub App the + owner installs. Neither exists for the tens of thousands of repositories in + this catalogue, and llms.txt lives on websites that have no webhooks at all. + So this polls: each URL once a week, with `If-None-Match` and + `If-Modified-Since`, so an unchanged file costs a `304` and no body. + raw.githubusercontent.com honours both. + + ## Fetching URLs that publishers choose + + `website_url` is whatever a publisher typed, so this treats every URL as + hostile until checked: + + * **https only**, on the default port. + * **The host must resolve to a public address.** A listing pointing at + `169.254.169.254` would otherwise have this server read cloud metadata + on its behalf. Redirects are followed by hand so every hop is checked, + not only the first. + * **Bodies stream with a hard cap** of #{div(262_144, 1024)} KB. A file is never + downloaded first and truncated after. + * **HTML is treated as absent.** Many sites answer any path with their + single-page-app shell and a `200`, which would otherwise make every + such site appear to publish an llms.txt. + + ## What gets recorded + + A document's first fetch is a baseline. After that a changed body records + the lines added and removed; a file that appears where one was confirmed + absent, or disappears where one was present, records that too. A failed + fetch -- a timeout, a 5xx, a 403 from a bot filter -- records nothing and + keeps what was known, because an outage is not a removal. + """ + import Ecto.Query + + alias McpRegistry.Changes + alias McpRegistry.Documents.Document + alias McpRegistry.Registry.{Logo, Server} + alias McpRegistry.Repo + + @max_bytes 262_144 + @max_redirects 3 + # Per change: how many changed lines are kept, and how much of each. + @diff_lines 200 + @line_chars 500 + + @default_batch 100 + @default_concurrency 4 + @default_recheck_days 7 + + # --- Which documents a listing points at ---------------------------------- + + @doc """ + The documents a listing points at, as `[{kind, url}]`. + + llms.txt is looked for at the root of the listing's website, or failing + that the domain its namespace was verified for -- `com.cloudflare` means + cloudflare.com. AGENTS.md is looked for at the root of a GitHub repository. + """ + def urls_for(%{} = server) do + [llms_txt_url(server), agents_md_url(server)] |> Enum.reject(&is_nil/1) + end + + defp llms_txt_url(server) do + case website_origin(server.website_url) || namespace_origin(server.name) do + nil -> nil + origin -> {"llms_txt", origin <> "/llms.txt"} + end + end + + defp agents_md_url(%{repository_url: url}) when is_binary(url) do + with %URI{host: host, path: "/" <> path} when host in ["github.com", "www.github.com"] <- + URI.parse(url), + [_owner, repo | _] <- String.split(path, "/", trim: true), + owner when is_binary(owner) <- Logo.repository_login(url), + repo = String.replace_suffix(repo, ".git", ""), + true <- Regex.match?(~r/\A[A-Za-z0-9._-]{1,100}\z/, repo) do + {"agents_md", "https://raw.githubusercontent.com/#{owner}/#{repo}/HEAD/AGENTS.md"} + else + _ -> nil + end + end + + defp agents_md_url(_), do: nil + + defp website_origin(url) when is_binary(url) do + case URI.parse(url) do + %URI{scheme: scheme, host: host, port: port} + when scheme in ["http", "https"] and is_binary(host) and host != "" and + port in [nil, 80, 443] -> + if domain?(host), do: "https://" <> String.downcase(host) + + _ -> + nil + end + end + + defp website_origin(_), do: nil + + # A reverse-DNS namespace names a domain the official registry verified. + # io.github.* names a GitHub account instead, which has no llms.txt. + defp namespace_origin("io.github." <> _), do: nil + + defp namespace_origin(name) when is_binary(name) do + host = + name + |> String.split("/", parts: 2) + |> List.first() + |> String.split(".") + |> Enum.reverse() + |> Enum.join(".") + + if domain?(host), do: "https://" <> String.downcase(host) + end + + defp namespace_origin(_), do: nil + + defp domain?(host) do + labels = String.split(host, ".") + + length(labels) >= 2 and + Enum.all?(labels, &Regex.match?(~r/\A[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\z/i, &1)) and + Regex.match?(~r/\A[a-z]{2,}\z/i, List.last(labels)) + end + + # --- Discovery ------------------------------------------------------------- + + @doc """ + Rebuilds which listing points at which document, and adds any new URLs as + pending. Documents no listing points at any more are dropped; their + recorded changes stay, since history does not stop being true. + """ + def discover do + servers = + Server + |> where([s], s.status == "active") + |> select([s], struct(s, [:id, :name, :website_url, :repository_url])) + |> Repo.all() + + links = + for s <- servers, {kind, url} <- urls_for(s), do: %{server_id: s.id, kind: kind, url: url} + + now = DateTime.utc_now() + + {:ok, _} = + Repo.transaction( + fn -> + Repo.delete_all("server_documents") + links |> Enum.chunk_every(5_000) |> Enum.each(&Repo.insert_all("server_documents", &1)) + + links + |> Enum.uniq_by(& &1.url) + |> Enum.map(&%{url: &1.url, kind: &1.kind, inserted_at: now, updated_at: now}) + |> Enum.chunk_every(5_000) + |> Enum.each( + &Repo.insert_all(Document, &1, on_conflict: :nothing, conflict_target: :url) + ) + + Repo.query!("DELETE FROM documents d WHERE NOT EXISTS + (SELECT 1 FROM server_documents sd WHERE sd.url = d.url)") + end, + timeout: :infinity + ) + + length(links) + end + + @doc "The document URLs a listing points at, from the last discovery." + def linked_urls(%Server{id: id}) do + from(sd in "server_documents", where: sd.server_id == ^id, select: sd.url) |> Repo.all() + end + + # --- Checking -------------------------------------------------------------- + + @doc """ + Checks one batch of due documents and returns a tally of outcomes. + Options: `:limit`, `:concurrency`, `:recheck_days`. + """ + def check_batch(opts \\ []) do + limit = Keyword.get(opts, :limit, @default_batch) + concurrency = Keyword.get(opts, :concurrency, @default_concurrency) + + cutoff = + DateTime.add( + DateTime.utc_now(), + -Keyword.get(opts, :recheck_days, @default_recheck_days) * 86_400 + ) + + Document + |> where([d], is_nil(d.checked_at) or d.checked_at < ^cutoff) + |> order_by([d], asc_nulls_first: d.checked_at, asc: d.id) + |> limit(^limit) + |> Repo.all() + |> Task.async_stream(&check/1, + max_concurrency: concurrency, + timeout: 60_000, + on_timeout: :kill_task + ) + |> Enum.map(fn + {:ok, outcome} -> outcome + {:exit, _} -> :error + end) + |> Enum.frequencies() + end + + @doc "Fetches one document, records any change, and returns what happened." + def check(%Document{} = doc) do + now = DateTime.utc_now() + + case fetch(doc.url, conditional_headers(doc), @max_redirects) do + :not_modified -> + update!(doc, %{checked_at: now, last_error: nil}) + :not_modified + + {:ok, body, resp} -> + present(doc, body, resp, now) + + :missing -> + missing(doc, now) + + {:error, reason} -> + # An outage is not a removal: keep what was known. + update!(doc, %{checked_at: now, last_error: String.slice(reason, 0, 500)}) + :error + end + end + + defp present(doc, body, resp, now) do + sha = :crypto.hash(:sha256, body) |> Base.encode16(case: :lower) + + validators = %{ + status: "ok", + etag: header(resp, "etag"), + last_modified: header(resp, "last-modified"), + checked_at: now, + last_error: nil + } + + fresh = Map.merge(validators, %{content: body, sha256: sha, changed_at: now}) + + cond do + doc.status == "ok" and doc.sha256 == sha -> + update!(doc, validators) + :unchanged + + doc.status == "ok" -> + {added, removed} = line_diff(doc.content || "", body) + Changes.record_document(doc.url, doc.kind, %{added: added, removed: removed}) + update!(doc, fresh) + :changed + + doc.status == "missing" -> + Changes.record_document(doc.url, doc.kind, %{ + added: body |> lines() |> cap(), + fields: %{"file" => ["absent", "published"]} + }) + + update!(doc, fresh) + :changed + + true -> + # First sight: a baseline, not an event. + update!(doc, fresh) + :baseline + end + end + + defp missing(%Document{status: "ok"} = doc, now) do + Changes.record_document(doc.url, doc.kind, %{ + removed: (doc.content || "") |> lines() |> cap(), + fields: %{"file" => ["published", "removed"]} + }) + + update!(doc, gone(now)) + :changed + end + + defp missing(doc, now) do + update!(doc, gone(now)) + :missing + end + + defp gone(now) do + %{ + status: "missing", + content: nil, + sha256: nil, + etag: nil, + last_modified: nil, + checked_at: now, + last_error: nil + } + end + + defp update!(doc, attrs), do: doc |> Ecto.Changeset.change(attrs) |> Repo.update!() + + # --- Diffing --------------------------------------------------------------- + + defp line_diff(old, new) do + old + |> lines() + |> List.myers_difference(lines(new)) + |> Enum.reduce({[], []}, fn + {:ins, ls}, {added, removed} -> {[ls | added], removed} + {:del, ls}, {added, removed} -> {added, [ls | removed]} + {:eq, _}, acc -> acc + end) + |> then(fn {added, removed} -> + {added |> Enum.reverse() |> List.flatten() |> cap(), + removed |> Enum.reverse() |> List.flatten() |> cap()} + end) + end + + # Blank lines are dropped: whitespace reflows would otherwise fill a + # changelog with nothing. + defp lines(text) do + text + |> String.split(["\r\n", "\n"]) + |> Enum.map(&String.trim_trailing/1) + |> Enum.reject(&(&1 == "")) + end + + defp cap(lines), + do: lines |> Enum.take(@diff_lines) |> Enum.map(&String.slice(&1, 0, @line_chars)) + + # --- Fetching -------------------------------------------------------------- + + defp conditional_headers(%Document{status: "ok"} = doc) do + [{"if-none-match", doc.etag}, {"if-modified-since", doc.last_modified}] + |> Enum.reject(fn {_, v} -> is_nil(v) end) + end + + defp conditional_headers(_), do: [] + + defp fetch(_url, _headers, hops) when hops < 0, do: {:error, "too many redirects"} + + defp fetch(url, headers, hops) do + with {:ok, uri} <- safe_uri(url), + :ok <- public_host(uri.host) do + case request(url, headers) do + {:ok, %Req.Response{status: 304}} -> + :not_modified + + {:ok, %Req.Response{status: status} = resp} when status in [301, 302, 303, 307, 308] -> + case header(resp, "location") do + nil -> {:error, "redirect without location"} + location -> fetch(URI.merge(url, location) |> URI.to_string(), headers, hops - 1) + end + + {:ok, %Req.Response{status: 200, body: body} = resp} -> + if text?(resp, body), do: {:ok, body, resp}, else: :missing + + {:ok, %Req.Response{status: status}} when status in [404, 410] -> + :missing + + {:ok, %Req.Response{status: status}} -> + {:error, "HTTP #{status}"} + + {:error, exception} -> + {:error, Exception.message(exception)} + end + end + end + + defp safe_uri(url) do + case URI.parse(url) do + %URI{scheme: "https", host: host, port: 443} = uri when is_binary(host) and host != "" -> + {:ok, uri} + + _ -> + {:error, "not an https URL on the default port"} + end + end + + defp request(url, headers) do + Req.get( + url, + [ + redirect: false, + retry: false, + decode_body: false, + receive_timeout: 10_000, + connect_options: [timeout: 5_000], + headers: [ + {"user-agent", + "MCPHarbor/#{Application.spec(:mcp_registry, :vsn)} (+#{McpRegistryWeb.Endpoint.url()})"}, + {"accept", "text/markdown, text/plain;q=0.9, */*;q=0.1"} | headers + ], + # Stop reading at the cap rather than downloading everything first. + into: fn {:data, chunk}, {req, resp} -> + body = (resp.body || "") <> chunk + + if byte_size(body) >= @max_bytes, + do: {:halt, {req, %{resp | body: binary_part(body, 0, @max_bytes)}}}, + else: {:cont, {req, %{resp | body: body}}} + end + ] ++ config(:req_options, []) + ) + end + + defp text?(resp, body) do + type = resp |> header("content-type") |> to_string() |> String.downcase() + + head = + body + |> binary_part(0, min(byte_size(body), 512)) + |> String.trim_leading() + |> String.downcase() + + String.valid?(body) and not String.contains?(type, "html") and + not String.starts_with?(head, [" value + [] -> nil + end + end + + # --- Refusing private addresses ------------------------------------------- + + defp public_host(host) do + resolver = config(:resolver, &resolve/1) + + case resolver.(host) do + {:ok, []} -> + {:error, "#{host} does not resolve"} + + {:ok, addresses} -> + if Enum.all?(addresses, &public_address?/1), + do: :ok, + else: {:error, "#{host} resolves to a private address"} + + {:error, _} -> + {:error, "#{host} does not resolve"} + end + end + + defp resolve(host) do + charlist = String.to_charlist(host) + + v4 = + case :inet.getaddrs(charlist, :inet) do + {:ok, addrs} -> addrs + _ -> [] + end + + v6 = + case :inet.getaddrs(charlist, :inet6) do + {:ok, addrs} -> addrs + _ -> [] + end + + {:ok, v4 ++ v6} + end + + @doc false + def public_address?({a, b, _, _}) do + not (a in [0, 10, 127] or a >= 224 or + (a == 100 and b in 64..127) or + (a == 169 and b == 254) or + (a == 172 and b in 16..31) or + (a == 192 and b == 168) or + (a == 198 and b in 18..19)) + end + + def public_address?({0, 0, 0, 0, 0, 0, 0, 1}), do: false + def public_address?({0, 0, 0, 0, 0, 0, 0, 0}), do: false + # IPv4-mapped (::ffff:a.b.c.d): judge the IPv4 address inside it. + def public_address?({0, 0, 0, 0, 0, 0xFFFF, hi, lo}), + do: public_address?({div(hi, 256), rem(hi, 256), div(lo, 256), rem(lo, 256)}) + + def public_address?({first, _, _, _, _, _, _, _}) do + # fc00::/7 unique-local, fe80::/10 link-local, ff00::/8 multicast. + not (Bitwise.band(first, 0xFE00) == 0xFC00 or Bitwise.band(first, 0xFFC0) == 0xFE80 or + Bitwise.band(first, 0xFF00) == 0xFF00) + end + + defp config(key, default), + do: Keyword.get(Application.get_env(:mcp_registry, :documents, []), key, default) +end diff --git a/lib/mcp_registry/documents/document.ex b/lib/mcp_registry/documents/document.ex new file mode 100644 index 0000000..8584de0 --- /dev/null +++ b/lib/mcp_registry/documents/document.ex @@ -0,0 +1,19 @@ +defmodule McpRegistry.Documents.Document do + @moduledoc "One fetched llms.txt or AGENTS.md, keyed by URL. See `McpRegistry.Documents`." + use Ecto.Schema + + schema "documents" do + field :url, :string + field :kind, :string + field :status, :string, default: "pending" + field :etag, :string + field :last_modified, :string + field :sha256, :string + field :content, :string + field :last_error, :string + field :checked_at, :utc_datetime_usec + field :changed_at, :utc_datetime_usec + + timestamps(type: :utc_datetime_usec) + end +end diff --git a/lib/mcp_registry/documents/scheduler.ex b/lib/mcp_registry/documents/scheduler.ex new file mode 100644 index 0000000..87988bc --- /dev/null +++ b/lib/mcp_registry/documents/scheduler.ex @@ -0,0 +1,89 @@ +defmodule McpRegistry.Documents.Scheduler do + @moduledoc """ + Checks llms.txt and AGENTS.md files on a timer, and rediscovers which + listings point at which once a day. + + ## Why a timer and not GitHub webhooks + + A webhook needs admin rights on the repository, or a GitHub App its owner + installs; this registry has neither for almost every listing, and llms.txt + is served by websites that offer no webhooks at all. Polling with + conditional requests is the only thing that covers the whole catalogue, and + it is cheap: an unchanged file answers `304` with no body. See + `McpRegistry.Documents`. + + ## Pace + + A batch every five minutes at a concurrency of four, and each URL once a + week. That is a trickle -- a few thousand requests a day spread across + thousands of hosts -- which is what fetching other people's files should be. + + Set `DOCUMENTS_ENABLED=false` to stop it. + """ + use GenServer + + require Logger + + alias McpRegistry.Documents + + @rediscover_every :timer.hours(24) + + def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__) + + @impl true + def init(_opts) do + config = config() + + if config[:enabled] do + Process.send_after(self(), :tick, config[:initial_delay_ms]) + {:ok, %{discovered_at: nil}} + else + :ignore + end + end + + @impl true + def handle_info(:tick, state) do + config = config() + now = System.monotonic_time(:millisecond) + + # Never let one bad tick stop the next: it is scheduled whatever happened. + state = + try do + state = maybe_discover(state, now) + + tally = + Documents.check_batch( + limit: config[:batch_size], + concurrency: config[:concurrency], + recheck_days: config[:recheck_days] + ) + + if tally != %{}, do: Logger.info("Document check: #{inspect(tally)}") + state + rescue + error -> + Logger.warning("Document check failed: #{Exception.message(error)}") + state + catch + kind, reason -> + Logger.warning("Document check #{kind}: #{inspect(reason)}") + state + end + + Process.send_after(self(), :tick, config[:interval_ms]) + {:noreply, state} + end + + defp maybe_discover(%{discovered_at: at} = state, now) + when is_integer(at) and now - at < @rediscover_every, + do: state + + defp maybe_discover(state, now) do + links = Documents.discover() + Logger.info("Document discovery: #{links} listing-to-document links") + %{state | discovered_at: now} + end + + defp config, do: Application.get_env(:mcp_registry, :documents, []) +end diff --git a/lib/mcp_registry/probe/scheduler.ex b/lib/mcp_registry/probe/scheduler.ex index f50644e..878360a 100644 --- a/lib/mcp_registry/probe/scheduler.ex +++ b/lib/mcp_registry/probe/scheduler.ex @@ -13,7 +13,7 @@ defmodule McpRegistry.Probe.Scheduler do is no deadline here. Requests carry a user agent naming the registry and linking to it, and a - probed endpoint is not asked again for two weeks, so a steady state is a + probed endpoint is not asked again for a week, so a steady state is a trickle rather than a sweep. Set `PROBE_ENABLED=false` to stop it entirely. diff --git a/lib/mcp_registry_web/live/server_live/changelog.ex b/lib/mcp_registry_web/live/server_live/changelog.ex index 5070c75..66d829f 100644 --- a/lib/mcp_registry_web/live/server_live/changelog.ex +++ b/lib/mcp_registry_web/live/server_live/changelog.ex @@ -167,7 +167,26 @@ defmodule McpRegistryWeb.ServerLive.Changelog do -