diff --git a/server/test/shelf_index_postgres_test.gleam b/server/test/shelf_index_postgres_test.gleam index 402676b..11a7f99 100644 --- a/server/test/shelf_index_postgres_test.gleam +++ b/server/test/shelf_index_postgres_test.gleam @@ -9,7 +9,9 @@ import crate_server/shelf_index_postgres import envoy import exception import gleam/erlang/process +import gleam/int import gleam/option.{None, Some} +import gleam/time/timestamp // A single-connection pool, killed as soon as this call returns (even if // `run` itself crashes): each test opens its own fresh pool, and a normal @@ -59,6 +61,12 @@ fn blank_entry( ) } +fn fresh_did(tag: String) -> String { + let #(_, nanoseconds) = + timestamp.to_unix_seconds_and_nanoseconds(timestamp.system_time()) + "did:plc:pgtest-" <> tag <> "-" <> int.to_string(nanoseconds) +} + pub fn event_round_trip_test() { use store <- with_store() let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg1" @@ -120,7 +128,7 @@ pub fn record_and_fold_round_trips_flattened_columns_test() { pub fn seen_round_trip_test() { use store <- with_store() - let did = "did:plc:pgtest-seen" + let did = fresh_did("seen") assert store.seen.has_seen(did) == False store.seen.mark_seen(did) assert store.seen.has_seen(did) == True @@ -128,7 +136,7 @@ pub fn seen_round_trip_test() { pub fn counts_reflect_stored_rows_and_go_back_down_on_delete_test() { use store <- with_store() - let did = "did:plc:pgtest-counts" + let did = fresh_did("counts") let entries_before = store.entries.count() let seen_before = store.seen.count() let entry_uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/pgc1" @@ -157,7 +165,7 @@ pub fn counts_reflect_stored_rows_and_go_back_down_on_delete_test() { pub fn delete_for_did_round_trip_test() { use store <- with_store() - let did = "did:plc:pgtest-purge" + let did = fresh_did("purge") let entry_uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/pg3" shelf_index.record_and_fold( store, -- 2.51.2 From a304cc34c423cee06db1a05d855b83a0d783d79e Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Tue, 11 Aug 2026 10:42:06 +0200 Subject: [PATCH 2/3] fix: dedupe the unpaged follow reads so they agree with the paged ones --- server/src/crate_server/follow_index.gleam | 16 ++++++++-------- server/test/follow_index_test.gleam | 5 ++++- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/server/src/crate_server/follow_index.gleam b/server/src/crate_server/follow_index.gleam index 636d4c0..14e1845 100644 --- a/server/src/crate_server/follow_index.gleam +++ b/server/src/crate_server/follow_index.gleam @@ -168,10 +168,10 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { Following(did, reply) -> { process.send( reply, - dict.values(state.follows) - |> list.filter(fn(f) { f.did == did }) - |> list.map(fn(f) { f.subject }) - |> list.sort(string.compare), + state.follows + |> matching_values(fn(follow) { follow.did == did }, fn(follow) { + follow.subject + }), ) actor.continue(state) } @@ -189,10 +189,10 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { Followers(did, reply) -> { process.send( reply, - dict.values(state.follows) - |> list.filter(fn(f) { f.subject == did }) - |> list.map(fn(f) { f.did }) - |> list.sort(string.compare), + state.follows + |> matching_values(fn(follow) { follow.subject == did }, fn(follow) { + follow.did + }), ) actor.continue(state) } diff --git a/server/test/follow_index_test.gleam b/server/test/follow_index_test.gleam index bd72181..997e3d6 100644 --- a/server/test/follow_index_test.gleam +++ b/server/test/follow_index_test.gleam @@ -38,9 +38,12 @@ pub fn duplicate_records_are_preserved_until_their_uri_is_deleted_test() { let assert Ok(store) = follow_index.start() store.edges.upsert(follow("did:a", "did:b", "f1")) store.edges.upsert(follow("did:a", "did:b", "f2")) - assert store.edges.following("did:a") == ["did:b", "did:b"] + // One subject however many records point at it, matching the paged read. + assert store.edges.following("did:a") == ["did:b"] store.edges.delete("at://did:a/dev.mokkenstorm.crate.graph.follow/f1") assert store.edges.following("did:a") == ["did:b"] + store.edges.delete("at://did:a/dev.mokkenstorm.crate.graph.follow/f2") + assert store.edges.following("did:a") == [] } pub fn following_is_empty_until_a_did_is_seen_test() { -- 2.51.2 From a865614c1c400d05474b818b055ae7c563eac9ab Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Tue, 11 Aug 2026 10:42:06 +0200 Subject: [PATCH 3/3] fix: cap the identity resolver fan-out per batch --- server/src/crate_server/identity_resolver.gleam | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/server/src/crate_server/identity_resolver.gleam b/server/src/crate_server/identity_resolver.gleam index 76e356c..41b6116 100644 --- a/server/src/crate_server/identity_resolver.gleam +++ b/server/src/crate_server/identity_resolver.gleam @@ -14,6 +14,11 @@ import gleam/set /// on the slowest concurrent miss. const fetch_timeout_ms = 5000 +/// Misses are fetched in chunks of this size. A cold cache on a full +/// connections page would otherwise open one outbound request per miss, with +/// nothing bounding the total across concurrent viewers. +const max_concurrent_fetches = 20 + /// `cache` backs both the positive and negative TTL lookups; `fetch` resolves /// one identifier (did or handle) to an `Identity`, `Error(Nil)` on any /// failure (not found, transport, decode -- all indistinguishable here). @@ -47,7 +52,11 @@ pub fn resolve_many( && !set.contains(resolution.suppressed, id) }) let fetched = - parallel.map(misses, fetch_timeout_ms, fn(id) { #(id, resolver.fetch(id)) }) + misses + |> list.sized_chunk(max_concurrent_fetches) + |> list.flat_map(fn(chunk) { + parallel.map(chunk, fetch_timeout_ms, fn(id) { #(id, resolver.fetch(id)) }) + }) let resolved = fetched |> list.filter_map(record_outcome(resolver, _)) dict.merge(resolution.hits, dict.from_list(resolved)) }