diff --git a/server/test/feed_handler_test.gleam b/server/test/feed_handler_test.gleam index 32662df..7ec8250 100644 --- a/server/test/feed_handler_test.gleam +++ b/server/test/feed_handler_test.gleam @@ -30,8 +30,6 @@ const viewer_did = "did:plc:x" const session_cookie = "ar_oauth_sid" -const far_future = 9_999_999_999 - fn feed_path() -> String { support.xrpc("feed.getFeedSkeleton") } @@ -41,12 +39,12 @@ fn a_session_as(did: String) -> sessions.OauthSession { support.stub_session_with( access_token: "at", refresh_token: "rt", - expires_at: far_future, + expires_at: 9_999_999_999, ) sessions.OauthSession(..session, did:) } -/// Always an "acquired" row -- no test in this file varies the status. +// Always an "acquired" row: no test in this file varies the status. fn adoption( did did: String, rkey rkey: String, @@ -67,14 +65,31 @@ fn release_uri(rkey: String) -> String { "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/" <> rkey } -/// The nth day since the epoch, RFC3339 -- a cheap way to generate rows -/// spaced far enough apart (24h) that none of them ever collapse together. +// The nth day since the epoch, spacing rows far enough apart (24h) that none +// of them ever collapse together. fn nth_day(n: Int) -> String { timestamp.from_unix_seconds(0) |> timestamp.add(duration.hours(24 * n)) |> timestamp.to_rfc3339(calendar.utc_offset) } +// Each row is #(adopter did, day), newest first in the feed's own ordering. +fn store_with(rows: List(#(String, Int))) -> catalog_index.Store { + let assert Ok(store) = catalog_index.start() + rows + |> list.index_map(fn(row, i) { + let #(did, day) = row + let n = int.to_string(i + 1) + store.adoptions.upsert(adoption( + did:, + rkey: "e" <> n, + release_uri: release_uri("r" <> n), + created_at: nth_day(day), + )) + }) + store +} + fn follow_record_json(rkey: String, subject: String) -> json.Json { json.object([ #( @@ -94,21 +109,42 @@ fn follow_record_json(rkey: String, subject: String) -> json.Json { ]) } -fn list_records_body(records: List(json.Json)) -> String { - json.object([#("records", json.preprocessed_array(records))]) - |> json.to_string -} - -fn network_client(follows_body: String) -> xrpc.Client { +fn follows_client(subjects: List(String)) -> xrpc.Client { + let body = + subjects + |> list.index_map(fn(subject, i) { + follow_record_json("3f" <> int.to_string(i), subject) + }) + |> fn(records) { + json.object([#("records", json.preprocessed_array(records))]) + |> json.to_string + } xrpc.Client(send: fn(req) { case req.path { "/xrpc/com.atproto.repo.listRecords" -> - Ok(response.Response(200, [], bit_array.from_string(follows_body))) + Ok(response.Response(200, [], bit_array.from_string(body))) _ -> panic as { "unexpected path: " <> req.path } } }) } +// Any PDS call at all blows up: the assertion for a warm feed request. +fn panicking_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { panic as { "unexpected PDS call: " <> req.path } }) +} + +fn indexed_follow(rkey: String, subject: String) -> follow_index.Follow { + follow_index.Follow( + follow_uri: "at://" + <> viewer_did + <> "/dev.mokkenstorm.crate.graph.follow/" + <> rkey, + did: viewer_did, + subject:, + created_at: "2026-01-01T00:00:00Z", + ) +} + fn test_context( client: xrpc.Client, index: catalog_index.Store, @@ -127,20 +163,22 @@ fn test_context( #(ctx, cfg, ctx.follow_index) } -fn authed_get(path: String, cfg: config.Config) -> wisp.Request { - authed_get_as(path, cfg, viewer_did) -} - fn authed_get_as( - path: String, + query: String, cfg: config.Config, did: String, ) -> wisp.Request { let assert Ok(id) = session_store.create(cfg.sessions, a_session_as(did)) - simulate.request(http.Get, path) + simulate.request(http.Get, feed_path() <> query) |> simulate.cookie(session_cookie, id, wisp.Signed) } +fn get_feed(ctx: Context, cfg: config.Config, query: String) -> String { + let resp = feed.get_feed_skeleton(authed_get_as(query, cfg, viewer_did), ctx) + assert resp.status == 200 + simulate.read_body(resp) +} + fn item_actors(body: String) -> List(String) { let decoder = decode.at( @@ -161,298 +199,104 @@ fn item_count(body: String) -> Int { } pub fn followed_rows_are_included_and_viewer_and_non_followed_are_excluded_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:f1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(2), - )) - store.adoptions.upsert(adoption( - did: viewer_did, - rkey: "e3", - release_uri: release_uri("r3"), - created_at: nth_day(1), - )) - let client = - network_client( - list_records_body([follow_record_json("3aaa", "did:plc:f1")]), - ) - let #(ctx, cfg, follows) = test_context(client, store) - follows.edges.upsert(follow_index.Follow( - follow_uri: "at://" - <> viewer_did - <> "/dev.mokkenstorm.crate.graph.follow/stale", - did: viewer_did, - subject: "did:plc:stale", - created_at: "2026-01-01T00:00:00Z", - )) - let resp = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert resp.status == 200 - let body = simulate.read_body(resp) + let store = + store_with([#("did:plc:f1", 3), #("did:plc:n1", 2), #(viewer_did, 1)]) + let #(ctx, cfg, follows) = test_context(follows_client(["did:plc:f1"]), store) + follows.edges.upsert(indexed_follow("stale", "did:plc:stale")) + let body = get_feed(ctx, cfg, "") assert support.field_bool(body, ["fallback"]) == Ok(False) assert item_actors(body) == ["did:plc:f1"] assert follows.edges.following(viewer_did) == ["did:plc:f1"] } -pub fn empty_follow_graph_falls_back_to_network_wide_including_the_viewer_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: viewer_did, - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(2), - )) - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(1), - )) - let client = network_client(list_records_body([])) - let #(ctx, cfg, _follows) = test_context(client, store) - let resp = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert resp.status == 200 - let body = simulate.read_body(resp) - assert support.field_bool(body, ["fallback"]) == Ok(True) - let actors = item_actors(body) - assert list.contains(actors, viewer_did) - assert list.contains(actors, "did:plc:n1") -} - -pub fn nonempty_follows_with_no_followed_adoptions_falls_back_on_the_first_page_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(1), - )) - // Followed, but they have no adoptions of their own. - let client = - network_client( - list_records_body([follow_record_json("3bbb", "did:plc:f1")]), - ) - let #(ctx, cfg, _follows) = test_context(client, store) - let resp = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert resp.status == 200 - let body = simulate.read_body(resp) - assert support.field_bool(body, ["fallback"]) == Ok(True) - assert item_actors(body) == ["did:plc:n1"] -} - -pub fn network_fallback_cursor_continues_in_network_mode_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - store.adoptions.upsert(adoption( - did: "did:plc:n2", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(2), - )) - let client = - network_client( - list_records_body([follow_record_json("3bbb", "did:plc:f1")]), - ) - let #(ctx, cfg, _follows) = test_context(client, store) - let first = - feed.get_feed_skeleton(authed_get(feed_path() <> "?limit=1", cfg), ctx) - let first_body = simulate.read_body(first) - assert support.field_bool(first_body, ["fallback"]) == Ok(True) - assert item_actors(first_body) == ["did:plc:n1"] - let assert Ok(cursor) = support.field_string(first_body, ["cursor"]) - let second = - feed.get_feed_skeleton( - authed_get(feed_path() <> "?limit=1&cursor=" <> cursor, cfg), - ctx, - ) - let second_body = simulate.read_body(second) - assert support.field_bool(second_body, ["fallback"]) == Ok(True) - assert item_actors(second_body) == ["did:plc:n2"] - assert support.field_present(second_body, ["cursor"]) == False +// Nothing to show from the follow graph flips the first page to network-wide, +// whether the graph is empty or its members simply have no adoptions. +pub fn a_first_page_with_no_followed_rows_falls_back_to_network_wide_test() { + [#([], [viewer_did, "did:plc:n1"]), #(["did:plc:f1"], ["did:plc:n1"])] + |> list.each(fn(row) { + let #(follows, expected_actors) = row + let store = + store_with(list.index_map(expected_actors, fn(did, i) { #(did, 2 - i) })) + let #(ctx, cfg, _follows) = test_context(follows_client(follows), store) + let body = get_feed(ctx, cfg, "") + assert support.field_bool(body, ["fallback"]) == Ok(True) + assert item_actors(body) == expected_actors + }) } -pub fn network_fallback_cursor_is_scoped_to_the_viewer_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - store.adoptions.upsert(adoption( - did: "did:plc:n2", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(2), - )) - let client = - network_client( - list_records_body([follow_record_json("3bbb", "did:plc:f1")]), - ) - let #(ctx, cfg, _follows) = test_context(client, store) - let first = - feed.get_feed_skeleton(authed_get(feed_path() <> "?limit=1", cfg), ctx) - let first_body = simulate.read_body(first) - let assert Ok(cursor) = support.field_string(first_body, ["cursor"]) - let second = - feed.get_feed_skeleton( - authed_get_as( - feed_path() <> "?limit=1&cursor=" <> cursor, - cfg, - "did:plc:other", - ), - ctx, - ) - let second_body = simulate.read_body(second) - assert support.field_bool(second_body, ["fallback"]) == Ok(True) - assert item_actors(second_body) == ["did:plc:n1"] +// The namespace a cursor carries is what keeps followed and network pages from +// crossing: a cursor minted in one mode resumes in that same mode. +pub fn a_cursor_resumes_in_the_mode_that_minted_it_test() { + [ + #(["did:plc:f1"], ["did:plc:n1", "did:plc:n2"], True), + #(["did:plc:f1", "did:plc:f2"], ["did:plc:f1", "did:plc:f2"], False), + ] + |> list.each(fn(row) { + let #(follows, actors, fallback) = row + let store = store_with(list.index_map(actors, fn(did, i) { #(did, 3 - i) })) + let #(ctx, cfg, _follows) = test_context(follows_client(follows), store) + let first = get_feed(ctx, cfg, "?limit=1") + assert support.field_bool(first, ["fallback"]) == Ok(fallback) + assert item_actors(first) == list.take(actors, 1) + let assert Ok(cursor) = support.field_string(first, ["cursor"]) + let second = get_feed(ctx, cfg, "?limit=1&cursor=" <> cursor) + assert support.field_bool(second, ["fallback"]) == Ok(fallback) + assert item_actors(second) == list.drop(actors, 1) + assert support.field_present(second, ["cursor"]) == False + }) } -pub fn cursor_stays_in_the_followed_namespace_across_pages_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:f1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - store.adoptions.upsert(adoption( - did: "did:plc:f2", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(2), - )) - store.adoptions.upsert(adoption( - did: "did:plc:f3", - rkey: "e3", - release_uri: release_uri("r3"), - created_at: nth_day(1), - )) - let client = - network_client( - list_records_body([ - follow_record_json("3aaa", "did:plc:f1"), - follow_record_json("3bbb", "did:plc:f2"), - follow_record_json("3ccc", "did:plc:f3"), - ]), - ) - let #(ctx, cfg, _follows) = test_context(client, store) - let first = - feed.get_feed_skeleton(authed_get(feed_path() <> "?limit=1", cfg), ctx) - let first_body = simulate.read_body(first) - assert support.field_bool(first_body, ["fallback"]) == Ok(False) - assert item_actors(first_body) == ["did:plc:f1"] - let assert Ok(cursor) = support.field_string(first_body, ["cursor"]) +pub fn a_network_fallback_cursor_is_scoped_to_the_viewer_test() { + let store = store_with([#("did:plc:n1", 3), #("did:plc:n2", 2)]) + let #(ctx, cfg, _follows) = + test_context(follows_client(["did:plc:f1"]), store) + let first = get_feed(ctx, cfg, "?limit=1") + let assert Ok(cursor) = support.field_string(first, ["cursor"]) let second = feed.get_feed_skeleton( - authed_get(feed_path() <> "?limit=1&cursor=" <> cursor, cfg), + authed_get_as("?limit=1&cursor=" <> cursor, cfg, "did:plc:other"), ctx, ) - let second_body = simulate.read_body(second) - assert support.field_bool(second_body, ["fallback"]) == Ok(False) - assert item_actors(second_body) == ["did:plc:f2"] + |> simulate.read_body + assert support.field_bool(second, ["fallback"]) == Ok(True) + assert item_actors(second) == ["did:plc:n1"] } pub fn limit_is_capped_at_the_lexicon_default_test() { - let assert Ok(store) = catalog_index.start() - list.repeat(0, 105) - |> list.index_map(fn(_, i) { i + 1 }) - |> list.each(fn(n) { - store.adoptions.upsert(adoption( - did: "did:plc:f1", - rkey: "e" <> int.to_string(n), - release_uri: release_uri("r" <> int.to_string(n)), - created_at: nth_day(n), - )) - }) - let client = - network_client( - list_records_body([follow_record_json("3aaa", "did:plc:f1")]), + let store = + store_with( + list.index_map(list.repeat("did:plc:f1", 105), fn(did, i) { + #(did, i + 1) + }), ) - let #(ctx, cfg, _follows) = test_context(client, store) - let resp = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - let body = simulate.read_body(resp) + let #(ctx, cfg, _follows) = + test_context(follows_client(["did:plc:f1"]), store) + let body = get_feed(ctx, cfg, "") assert item_count(body) == 100 assert support.field_present(body, ["cursor"]) } -/// Any PDS call at all blows up: the assertion for a warm feed request. -fn panicking_client() -> xrpc.Client { - xrpc.Client(send: fn(req) { panic as { "unexpected PDS call: " <> req.path } }) -} - -fn indexed_follow(rkey: String, subject: String) -> follow_index.Follow { - follow_index.Follow( - follow_uri: "at://" - <> viewer_did - <> "/dev.mokkenstorm.crate.graph.follow/" - <> rkey, - did: viewer_did, - subject:, - created_at: "2026-01-01T00:00:00Z", - ) -} - pub fn an_indexed_follow_graph_is_served_without_touching_the_pds_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:f1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - store.adoptions.upsert(adoption( - did: "did:plc:n1", - rkey: "e2", - release_uri: release_uri("r2"), - created_at: nth_day(2), - )) + let store = store_with([#("did:plc:f1", 3), #("did:plc:n1", 2)]) let #(ctx, cfg, follows) = test_context(panicking_client(), store) follows.edges.upsert(indexed_follow("3aaa", "did:plc:f1")) follows.seen.mark_seen(viewer_did) - let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert res.status == 200 - assert item_actors(simulate.read_body(res)) == ["did:plc:f1"] + assert item_actors(get_feed(ctx, cfg, "")) == ["did:plc:f1"] } pub fn a_never_seen_viewer_is_seeded_from_their_pds_once_test() { - let assert Ok(store) = catalog_index.start() - store.adoptions.upsert(adoption( - did: "did:plc:f1", - rkey: "e1", - release_uri: release_uri("r1"), - created_at: nth_day(3), - )) - let client = - network_client( - list_records_body([follow_record_json("3aaa", "did:plc:f1")]), - ) - let #(ctx, cfg, follows) = test_context(client, store) + let store = store_with([#("did:plc:f1", 3)]) + let #(ctx, cfg, follows) = test_context(follows_client(["did:plc:f1"]), store) assert follows.seen.has_seen(viewer_did) == False - let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert res.status == 200 - assert item_actors(simulate.read_body(res)) == ["did:plc:f1"] + assert item_actors(get_feed(ctx, cfg, "")) == ["did:plc:f1"] // Seeded and marked, so the next request never asks the PDS again. assert follows.edges.following(viewer_did) == ["did:plc:f1"] assert follows.seen.has_seen(viewer_did) == True } pub fn a_failed_seed_leaves_the_viewer_unmarked_for_a_retry_test() { - let assert Ok(store) = catalog_index.start() - let #(ctx, cfg, follows) = test_context(support.unreachable_client(), store) - let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) - assert res.status == 200 + let #(ctx, cfg, follows) = + test_context(support.unreachable_client(), store_with([])) + let _ = get_feed(ctx, cfg, "") assert follows.seen.has_seen(viewer_did) == False } diff --git a/server/test/graph_handler_test.gleam b/server/test/graph_handler_test.gleam index 84086b4..e76b7fd 100644 --- a/server/test/graph_handler_test.gleam +++ b/server/test/graph_handler_test.gleam @@ -27,52 +27,34 @@ import wisp/simulate const target_did = "did:plc:target" -const session_cookie = "ar_oauth_sid" +const viewer_did = "did:plc:x" -const far_future = 9_999_999_999 +const session_cookie = "ar_oauth_sid" fn a_session() -> sessions.OauthSession { support.stub_session_with( access_token: "at", refresh_token: "rt", - expires_at: far_future, + expires_at: 9_999_999_999, ) } +fn follow_uri(rkey: String) -> String { + "at://" <> viewer_did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey +} + fn stored_follow(rkey: String, subject: String) -> StoredItem(GraphFollow) { StoredItem( - uri: "at://did:plc:x/dev.mokkenstorm.crate.graph.follow/" <> rkey, + uri: follow_uri(rkey), cid: "bafy" <> rkey, rkey:, value: GraphFollow(created_at: "2026-01-01T00:00:00Z", subject:), ) } -pub fn find_locates_the_stored_follow_or_returns_none_test() { - let stored = [ - stored_follow("3aaa", "did:plc:other"), - stored_follow("3bbb", target_did), - ] - let assert Some(found) = graph_follows.find(stored, target_did) - assert found.rkey == "3bbb" - assert graph_follows.find( - [stored_follow("3aaa", "did:plc:other")], - target_did, - ) - == None -} - -fn list_records_body(records: List(json.Json)) -> String { - json.object([#("records", json.preprocessed_array(records))]) - |> json.to_string -} - fn follow_record_json(rkey: String, subject: String) -> json.Json { json.object([ - #( - "uri", - json.string("at://did:plc:x/dev.mokkenstorm.crate.graph.follow/" <> rkey), - ), + #("uri", json.string(follow_uri(rkey))), #("cid", json.string("bafy" <> rkey)), #( "value", @@ -84,12 +66,14 @@ fn follow_record_json(rkey: String, subject: String) -> json.Json { ]) } +fn list_records_body(records: List(json.Json)) -> String { + json.object([#("records", json.preprocessed_array(records))]) + |> json.to_string +} + fn create_record_response(rkey: String) -> String { json.object([ - #( - "uri", - json.string("at://did:plc:x/dev.mokkenstorm.crate.graph.follow/" <> rkey), - ), + #("uri", json.string(follow_uri(rkey))), #("cid", json.string("bafy" <> rkey)), ]) |> json.to_string @@ -101,12 +85,6 @@ fn ok_response( Ok(response.Response(200, [], bit_array.from_string(body))) } -fn forbidden_response( - body: String, -) -> Result(response.Response(BitArray), xrpc.TransportError) { - Ok(response.Response(403, [], bit_array.from_string(body))) -} - fn network_client(list_body: String, create_body: String) -> xrpc.Client { xrpc.Client(send: fn(req) { case req.path { @@ -123,115 +101,60 @@ fn missing_scope_client() -> xrpc.Client { case req.path { "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_records_body([])) "/xrpc/com.atproto.repo.putRecord" -> - forbidden_response("{\"error\":\"InsufficientScope\"}") + Ok(response.Response( + 403, + [], + bit_array.from_string("{\"error\":\"InsufficientScope\"}"), + )) _ -> panic as { "unexpected path: " <> req.path } } }) } -fn test_context(client: xrpc.Client) -> #(Context, config.Config) { - let cfg = support.stub_config(client) - let ctx = - support.stub_context_with( - cfg, - support.unreachable_catalog_deps(), - fn(_req) { Error("unused") }, - [], - ) - #(ctx, cfg) -} - -fn authed_post( - path: String, - body: json.Json, - cfg: config.Config, -) -> wisp.Request { - let assert Ok(id) = session_store.create(cfg.sessions, a_session()) - simulate.request(http.Post, path) - |> simulate.json_body(body) - |> simulate.cookie(session_cookie, id, wisp.Signed) -} - -pub fn follow_creates_a_record_and_returns_its_uri_test() { - let #(ctx, cfg) = - test_context(network_client( - list_records_body([]), - create_record_response("3ccc"), - )) - let req = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - let resp = graph.follow(req, ctx) - assert resp.status == 200 - let body = simulate.read_body(resp) - assert support.field_string(body, ["uri"]) - == Ok("at://did:plc:x/dev.mokkenstorm.crate.graph.follow/3ccc") -} - -pub fn follow_requires_authentication_test() { - let #(ctx, _cfg) = test_context(network_client(list_records_body([]), "")) - let req = - simulate.request(http.Post, support.xrpc("graph.followUser")) - |> simulate.json_body(json.object([#("subject", json.string(target_did))])) - assert graph.follow(req, ctx).status == 401 -} - -pub fn follow_scope_failure_is_exposed_as_a_reauthorization_signal_test() { - let #(ctx, cfg) = test_context(missing_scope_client()) - let req = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - let resp = graph.follow(req, ctx) - assert resp.status == 403 - let body = simulate.read_body(resp) - assert support.field_string(body, ["error"]) == Ok("Forbidden") - assert support.field_string(body, ["message"]) - == Ok( - "your session needs permission to follow collectors; sign in again to approve it", - ) -} - -pub fn unfollow_deletes_the_matching_follow_and_returns_ok_test() { - let #(ctx, cfg) = - test_context(network_client( - list_records_body([follow_record_json("3ddd", target_did)]), - "", - )) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - let resp = graph.unfollow(req, ctx) - assert resp.status == 200 +// Panics on listRecords, so any lookup that reaches it failed to use the index. +fn no_list_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") + "/xrpc/com.atproto.repo.putRecord" -> + ok_response(create_record_response("3ccc")) + _ -> panic as { "unexpected PDS call: " <> req.path } + } + }) } -pub fn unfollow_returns_404_when_no_follow_exists_test() { - let #(ctx, cfg) = test_context(network_client(list_records_body([]), "")) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - let resp = graph.unfollow(req, ctx) - assert resp.status == 404 +fn no_create_client(list_body: String) -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_body) + "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") + "/xrpc/com.atproto.repo.createRecord" -> + panic as "idempotent follow attempted a create" + _ -> panic as { "unexpected PDS call: " <> req.path } + } + }) } -const viewer_did = "did:plc:x" - -fn follow_uri(rkey: String) -> String { - "at://" <> viewer_did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey +// Fails deleting the second (later-sorted) rkey, succeeds on the first. +fn partial_delete_failure_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.deleteRecord" -> + case bit_array.to_string(req.body) { + Ok(body) -> + case string.contains(body, "3bbb") { + True -> + Error(core_xrpc.ConnectionFailed("PDS unreachable for 3bbb")) + False -> ok_response("{}") + } + Error(_) -> panic as "expected a decodable delete body" + } + _ -> panic as { "unexpected PDS call: " <> req.path } + } + }) } -fn test_context_with( +fn test_context( client: xrpc.Client, follows: follow_index.Store, ) -> #(Context, config.Config) { @@ -249,176 +172,183 @@ fn test_context_with( #(ctx, cfg) } -/// Fails any listRecords: the assertion that a lookup came from the index. -fn no_list_client() -> xrpc.Client { - xrpc.Client(send: fn(req) { - case req.path { - "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") - "/xrpc/com.atproto.repo.putRecord" -> - ok_response(create_record_response("3ccc")) - _ -> panic as { "unexpected PDS call: " <> req.path } - } - }) +fn authed_post(path: String, cfg: config.Config) -> wisp.Request { + let assert Ok(id) = session_store.create(cfg.sessions, a_session()) + simulate.request(http.Post, support.xrpc(path)) + |> simulate.json_body(json.object([#("subject", json.string(target_did))])) + |> simulate.cookie(session_cookie, id, wisp.Signed) } -fn no_create_client(list_body: String) -> xrpc.Client { - xrpc.Client(send: fn(req) { - case req.path { - "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_body) - "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") - "/xrpc/com.atproto.repo.createRecord" -> - panic as "idempotent follow attempted a create" - _ -> panic as { "unexpected PDS call: " <> req.path } - } +fn call( + handler: fn(wisp.Request, Context) -> wisp.Response, + path: String, + client: xrpc.Client, + follows: follow_index.Store, +) -> wisp.Response { + let #(ctx, cfg) = test_context(client, follows) + handler(authed_post(path, cfg), ctx) +} + +fn seeded_index(edges: List(#(String, String))) -> follow_index.Store { + let follows = support.fresh_follow_index() + list.each(edges, fn(edge) { + let #(rkey, created_at) = edge + follows.edges.upsert(follow_index.Follow( + follow_uri: follow_uri(rkey), + did: viewer_did, + subject: target_did, + created_at:, + )) }) + follows.seen.mark_seen(viewer_did) + follows +} + +pub fn find_locates_the_stored_follow_or_returns_none_test() { + let stored = [ + stored_follow("3aaa", "did:plc:other"), + stored_follow("3bbb", target_did), + ] + let assert Some(found) = graph_follows.find(stored, target_did) + assert found.rkey == "3bbb" + assert graph_follows.find( + [stored_follow("3aaa", "did:plc:other")], + target_did, + ) + == None } -pub fn follow_is_idempotent_and_chooses_a_stable_legacy_record_test() { - let #(ctx, cfg) = +pub fn follow_requires_authentication_test() { + let #(ctx, _cfg) = test_context( + network_client(list_records_body([]), ""), + support.fresh_follow_index(), + ) + let req = + simulate.request(http.Post, support.xrpc("graph.followUser")) + |> simulate.json_body(json.object([#("subject", json.string(target_did))])) + assert graph.follow(req, ctx).status == 401 +} + +// A fresh follow writes a record; an existing one reuses the lowest-sorting +// legacy rkey instead of creating a second. +pub fn follow_returns_a_stable_record_uri_test() { + [ + #( + network_client(list_records_body([]), create_record_response("3ccc")), + "3ccc", + ), + #( no_create_client( list_records_body([ follow_record_json("3bbb", target_did), follow_record_json("3aaa", target_did), ]), ), + "3aaa", + ), + ] + |> list.each(fn(row) { + let #(client, rkey) = row + let resp = + call( + graph.follow, + "graph.followUser", + client, + support.fresh_follow_index(), + ) + assert resp.status == 200 + assert support.field_string(simulate.read_body(resp), ["uri"]) + == Ok(follow_uri(rkey)) + }) +} + +pub fn follow_scope_failure_is_exposed_as_a_reauthorization_signal_test() { + let resp = + call( + graph.follow, + "graph.followUser", + missing_scope_client(), + support.fresh_follow_index(), ) - let req = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, + assert resp.status == 403 + let body = simulate.read_body(resp) + assert support.field_string(body, ["error"]) == Ok("Forbidden") + assert support.field_string(body, ["message"]) + == Ok( + "your session needs permission to follow collectors; sign in again to approve it", ) - let body = simulate.read_body(graph.follow(req, ctx)) - assert support.field_string(body, ["uri"]) - == Ok("at://did:plc:x/dev.mokkenstorm.crate.graph.follow/3aaa") } pub fn follow_mirrors_the_written_record_into_the_index_test() { - let assert Ok(follows) = follow_index.start() - follows.seen.mark_seen(viewer_did) - let #(ctx, cfg) = test_context_with(no_list_client(), follows) - let req = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - assert graph.follow(req, ctx).status == 200 + let follows = seeded_index([]) + assert call(graph.follow, "graph.followUser", no_list_client(), follows).status + == 200 assert follows.edges.following(viewer_did) == [target_did] let assert Some(edge) = follows.edges.find(viewer_did, target_did) assert edge.follow_uri == follow_uri("3ccc") } -pub fn unfollow_finds_the_rkey_in_the_index_without_a_pds_list_test() { - let assert Ok(follows) = follow_index.start() - follows.edges.upsert(follow_index.Follow( - follow_uri: follow_uri("3ddd"), - did: viewer_did, - subject: target_did, - created_at: "2026-01-01T00:00:00Z", - )) - follows.seen.mark_seen(viewer_did) - let #(ctx, cfg) = test_context_with(no_list_client(), follows) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - assert graph.unfollow(req, ctx).status == 200 - assert follows.edges.following(viewer_did) == [] -} - -pub fn unfollow_removes_all_legacy_duplicate_index_edges_test() { - let assert Ok(follows) = follow_index.start() - follows.edges.upsert(follow_index.Follow( - follow_uri: follow_uri("3aaa"), - did: viewer_did, - subject: target_did, - created_at: "2026-01-01T00:00:00Z", - )) - follows.edges.upsert(follow_index.Follow( - follow_uri: follow_uri("3bbb"), - did: viewer_did, - subject: target_did, - created_at: "2026-01-02T00:00:00Z", - )) - follows.seen.mark_seen(viewer_did) - let #(ctx, cfg) = test_context_with(no_list_client(), follows) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - assert graph.unfollow(req, ctx).status == 200 - assert follows.edges.following(viewer_did) == [] -} - -pub fn unfollow_404s_from_the_index_when_the_edge_is_absent_test() { - let assert Ok(follows) = follow_index.start() - follows.seen.mark_seen(viewer_did) - let #(ctx, cfg) = test_context_with(no_list_client(), follows) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - assert graph.unfollow(req, ctx).status == 404 +// Unfollow resolves the rkey from the index rather than listing the PDS, and +// takes every duplicate legacy edge for the subject with it. +pub fn unfollow_clears_index_edges_without_a_pds_list_test() { + [ + #([#("3ddd", "2026-01-01T00:00:00Z")], 200), + #( + [#("3aaa", "2026-01-01T00:00:00Z"), #("3bbb", "2026-01-02T00:00:00Z")], + 200, + ), + #([], 404), + ] + |> list.each(fn(row) { + let #(edges, status) = row + let follows = seeded_index(edges) + let resp = + call(graph.unfollow, "graph.unfollowUser", no_list_client(), follows) + assert resp.status == status + assert follows.edges.following(viewer_did) == [] + }) } -/// Fails deleting the second (later-sorted) rkey, succeeds on the first. -fn partial_delete_failure_client() -> xrpc.Client { - xrpc.Client(send: fn(req) { - case req.path { - "/xrpc/com.atproto.repo.deleteRecord" -> - case bit_array.to_string(req.body) { - Ok(body) -> - case string.contains(body, "3bbb") { - True -> - Error(core_xrpc.ConnectionFailed("PDS unreachable for 3bbb")) - False -> ok_response("{}") - } - Error(_) -> panic as "expected a decodable delete body" - } - _ -> panic as { "unexpected PDS call: " <> req.path } - } +pub fn unfollow_without_the_index_falls_back_to_the_pds_test() { + [ + #(list_records_body([follow_record_json("3ddd", target_did)]), 200), + #(list_records_body([]), 404), + ] + |> list.each(fn(row) { + let #(list_body, status) = row + let resp = + call( + graph.unfollow, + "graph.unfollowUser", + network_client(list_body, ""), + support.fresh_follow_index(), + ) + assert resp.status == status }) } pub fn unfollow_index_keeps_only_edges_the_pds_delete_failed_on_test() { - let assert Ok(follows) = follow_index.start() - follows.edges.upsert(follow_index.Follow( - follow_uri: follow_uri("3aaa"), - did: viewer_did, - subject: target_did, - created_at: "2026-01-01T00:00:00Z", - )) - follows.edges.upsert(follow_index.Follow( - follow_uri: follow_uri("3bbb"), - did: viewer_did, - subject: target_did, - created_at: "2026-01-02T00:00:00Z", - )) - follows.seen.mark_seen(viewer_did) - let #(ctx, cfg) = test_context_with(partial_delete_failure_client(), follows) - let req = - authed_post( - support.xrpc("graph.unfollowUser"), - json.object([#("subject", json.string(target_did))]), - cfg, + let follows = + seeded_index([ + #("3aaa", "2026-01-01T00:00:00Z"), + #("3bbb", "2026-01-02T00:00:00Z"), + ]) + let resp = + call( + graph.unfollow, + "graph.unfollowUser", + partial_delete_failure_client(), + follows, ) - assert graph.unfollow(req, ctx).status == 502 + assert resp.status == 502 let assert Some(edge) = follows.edges.find(viewer_did, target_did) assert edge.follow_uri == follow_uri("3bbb") } pub fn deterministic_rkey_is_stable_and_atproto_safe_test() { let first = graph_follows.rkey_for_subject(target_did) - let second = graph_follows.rkey_for_subject(target_did) - assert first == second + assert first == graph_follows.rkey_for_subject(target_did) assert string.starts_with(first, "f_") assert string.contains(first, ":") == False } @@ -430,7 +360,7 @@ pub fn concurrent_follows_converge_on_the_same_put_rkey_test() { case req.path { "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_records_body([])) - "/xrpc/com.atproto.repo.putRecord" -> { + "/xrpc/com.atproto.repo.putRecord" -> case bit_array.to_string(req.body) { Ok(body) -> case string.contains(body, expected_rkey) { @@ -443,29 +373,19 @@ pub fn concurrent_follows_converge_on_the_same_put_rkey_test() { _ -> Error(core_xrpc.ConnectionFailed("non-deterministic follow rkey")) } - } _ -> panic as { "unexpected PDS call: " <> req.path } } }) - let #(ctx, cfg) = test_context(client) - let first = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) - let second = - authed_post( - support.xrpc("graph.followUser"), - json.object([#("subject", json.string(target_did))]), - cfg, - ) + let #(ctx, cfg) = test_context(client, support.fresh_follow_index()) let replies = process.new_subject() - process.spawn_unlinked(fn() { - process.send(replies, graph.follow(first, ctx).status) - }) - process.spawn_unlinked(fn() { - process.send(replies, graph.follow(second, ctx).status) + let requests = [ + authed_post("graph.followUser", cfg), + authed_post("graph.followUser", cfg), + ] + list.each(requests, fn(req) { + process.spawn_unlinked(fn() { + process.send(replies, graph.follow(req, ctx).status) + }) }) let assert Ok(a) = process.receive(replies, 1000) let assert Ok(b) = process.receive(replies, 1000)