diff --git a/server/src/crate_server/follow_index.gleam b/server/src/crate_server/follow_index.gleam index 11d3a6b..686383a 100644 --- a/server/src/crate_server/follow_index.gleam +++ b/server/src/crate_server/follow_index.gleam @@ -10,6 +10,7 @@ import gleam/dict.{type Dict} import gleam/erlang/process.{type Subject} import gleam/list import gleam/option.{type Option} +import gleam/order import gleam/otp/actor import gleam/result import gleam/set.{type Set} @@ -33,9 +34,13 @@ pub type EdgeOps { /// The dids `did` follows. Empty is ambiguous alone; pair with /// `SeenOps.has_seen`. following: fn(String) -> List(String), + following_page: fn(String, Option(String), Int) -> List(String), /// The dids that follow `did`. This reverse lookup is appview-derived and /// may be incomplete until every known actor has been backfilled. followers: fn(String) -> List(String), + followers_page: fn(String, Option(String), Int) -> List(String), + following_among: fn(String, List(String)) -> List(String), + followers_among: fn(String, List(String)) -> List(String), /// The unfollow path's lookup, served without a live `listRecords`. find: fn(String, String) -> Option(Follow), /// Every edge in a logical relationship. Legacy repos may contain @@ -44,7 +49,7 @@ pub type EdgeOps { /// Replace every indexed edge authored by `did` with a complete PDS read. /// The operation is atomic for the in-memory backend and transactional for /// the Postgres backend, so stale rows cannot survive a successful read. - replace_for_did: fn(String, List(Follow)) -> Nil, + replace_for_did: fn(String, List(Follow)) -> Result(Nil, String), /// Total indexed edges: a cheap completeness figure for `/debug/index`. count: fn() -> Int, ) @@ -72,10 +77,14 @@ type Msg { Delete(String) DeleteForDid(String) Following(String, Subject(List(String))) + FollowingPage(String, Option(String), Int, Subject(List(String))) Followers(String, Subject(List(String))) + FollowersPage(String, Option(String), Int, Subject(List(String))) + FollowingAmong(String, List(String), Subject(List(String))) + FollowersAmong(String, List(String), Subject(List(String))) Find(String, String, Subject(Option(Follow))) FindAll(String, String, Subject(List(Follow))) - ReplaceForDid(String, List(Follow)) + ReplaceForDid(String, List(Follow), Subject(Result(Nil, String))) Count(Subject(Int)) HasSeen(String, Subject(Bool)) MarkSeen(String) @@ -94,7 +103,19 @@ pub fn start() -> Result(Store, actor.StartError) { delete: fn(follow_uri) { process.send(subject, Delete(follow_uri)) }, delete_for_did: fn(did) { process.send(subject, DeleteForDid(did)) }, following: fn(did) { parallel.ask(subject, Following(did, _), or: []) }, + following_page: fn(did, after, limit) { + parallel.ask(subject, FollowingPage(did, after, limit, _), or: []) + }, followers: fn(did) { parallel.ask(subject, Followers(did, _), or: []) }, + followers_page: fn(did, after, limit) { + parallel.ask(subject, FollowersPage(did, after, limit, _), or: []) + }, + following_among: fn(did, candidates) { + parallel.ask(subject, FollowingAmong(did, candidates, _), or: []) + }, + followers_among: fn(did, candidates) { + parallel.ask(subject, FollowersAmong(did, candidates, _), or: []) + }, find: fn(did, subject_did) { parallel.ask(subject, Find(did, subject_did, _), or: option.None) }, @@ -102,7 +123,11 @@ pub fn start() -> Result(Store, actor.StartError) { parallel.ask(subject, FindAll(did, subject_did, _), or: []) }, replace_for_did: fn(did, follows) { - process.send(subject, ReplaceForDid(did, follows)) + parallel.ask( + subject, + ReplaceForDid(did, follows, _), + or: Error("follow index unavailable"), + ) }, count: fn() { parallel.ask(subject, Count, or: 0) }, ), @@ -141,6 +166,17 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { ) actor.continue(state) } + FollowingPage(did, after, limit, reply) -> { + process.send( + reply, + state.follows + |> matching_values(fn(follow) { follow.did == did }, fn(follow) { + follow.subject + }) + |> page_after(after, limit), + ) + actor.continue(state) + } Followers(did, reply) -> { process.send( reply, @@ -151,6 +187,45 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { ) actor.continue(state) } + FollowersPage(did, after, limit, reply) -> { + process.send( + reply, + state.follows + |> matching_values(fn(follow) { follow.subject == did }, fn(follow) { + follow.did + }) + |> page_after(after, limit), + ) + actor.continue(state) + } + FollowingAmong(did, candidates, reply) -> { + let candidate_set = set.from_list(candidates) + process.send( + reply, + state.follows + |> matching_values( + fn(follow) { + follow.did == did && set.contains(candidate_set, follow.subject) + }, + fn(follow) { follow.subject }, + ), + ) + actor.continue(state) + } + FollowersAmong(did, candidates, reply) -> { + let candidate_set = set.from_list(candidates) + process.send( + reply, + state.follows + |> matching_values( + fn(follow) { + follow.subject == did && set.contains(candidate_set, follow.did) + }, + fn(follow) { follow.did }, + ), + ) + actor.continue(state) + } Find(did, subject_did, reply) -> { process.send( reply, @@ -171,7 +246,7 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { ) actor.continue(state) } - ReplaceForDid(did, follows) -> { + ReplaceForDid(did, follows, reply) -> { let retained = dict.filter(state.follows, fn(_uri, follow) { follow.did != did }) let indexed = @@ -179,6 +254,7 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { |> list.fold(retained, fn(edges, follow) { dict.insert(edges, follow.follow_uri, follow) }) + process.send(reply, Ok(Nil)) actor.continue(State(..state, follows: indexed)) } Count(reply) -> { @@ -193,3 +269,35 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { actor.continue(State(..state, seen: set.insert(state.seen, did))) } } + +fn matching_values( + follows: Dict(String, Follow), + matches: fn(Follow) -> Bool, + value: fn(Follow) -> String, +) -> List(String) { + follows + |> dict.values + |> list.filter(matches) + |> list.map(value) + |> list.unique + |> list.sort(string.compare) +} + +fn page_after( + values: List(String), + after: Option(String), + limit: Int, +) -> List(String) { + let remaining = case after { + option.None -> values + option.Some(marker) -> + case list.contains(values, marker) { + True -> + list.drop_while(values, fn(value) { + string.compare(value, marker) != order.Gt + }) + False -> values + } + } + list.take(remaining, limit) +} diff --git a/server/src/crate_server/follow_index_postgres.gleam b/server/src/crate_server/follow_index_postgres.gleam index 8a34840..1758b3c 100644 --- a/server/src/crate_server/follow_index_postgres.gleam +++ b/server/src/crate_server/follow_index_postgres.gleam @@ -10,7 +10,7 @@ import crate_server/follow_index.{ } import crate_server/provenance import gleam/dynamic/decode -import gleam/list +import gleam/json import gleam/option.{type Option} import gleam/result import pog @@ -26,10 +26,22 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { delete: fn(follow_uri) { delete(conn, follow_uri) }, delete_for_did: fn(did) { delete_for_did(conn, did) }, following: fn(did) { following(conn, did) }, + following_page: fn(did, after, limit) { + following_page(conn, did, after, limit) + }, find: fn(did, subject) { find(conn, did, subject) }, find_all: fn(did, subject) { find_all(conn, did, subject) }, replace_for_did: fn(did, follows) { replace_for_did(conn, did, follows) }, followers: fn(did) { followers(conn, did) }, + followers_page: fn(did, after, limit) { + followers_page(conn, did, after, limit) + }, + following_among: fn(did, candidates) { + following_among(conn, did, candidates) + }, + followers_among: fn(did, candidates) { + followers_among(conn, did, candidates) + }, count: fn() { count(conn) }, ), seen: SeenOps(has_seen: fn(did) { has_seen(conn, did) }, mark_seen: fn(did) { @@ -59,6 +71,13 @@ fn migrate(conn: pog.Connection) -> Result(Nil, String) { |> pog.execute(conn) |> result.replace_error("follow_edges migration failed"), ) + use _ <- result.try( + pog.query( + "create index if not exists follow_edges_did_subject_idx on follow_edges (did, subject)", + ) + |> pog.execute(conn) + |> result.replace_error("follow_edges did subject index migration failed"), + ) // Every read is scoped to one follower. use _ <- result.try( pog.query( @@ -67,6 +86,13 @@ fn migrate(conn: pog.Connection) -> Result(Nil, String) { |> pog.execute(conn) |> result.replace_error("follow_edges did index migration failed"), ) + use _ <- result.try( + pog.query( + "create index if not exists follow_edges_subject_did_idx on follow_edges (subject, did)", + ) + |> pog.execute(conn) + |> result.replace_error("follow_edges subject did index migration failed"), + ) use _ <- result.try( pog.query( "create index if not exists follow_edges_subject_idx on follow_edges (subject)", @@ -136,6 +162,15 @@ fn following(conn: pog.Connection, did: String) -> List(String) { |> db.all(conn) } +fn following_page( + conn: pog.Connection, + did: String, + after: Option(String), + limit: Int, +) -> List(String) { + connection_page(conn, "did", "subject", did, after, limit) +} + fn find(conn: pog.Connection, did: String, subject: String) -> Option(Follow) { pog.query("select " <> columns <> " from follow_edges where did = $1 and subject = $2 @@ -164,35 +199,47 @@ fn replace_for_did( conn: pog.Connection, did: String, follows: List(Follow), -) -> Nil { - let result = +) -> Result(Nil, String) { + let encoded = + follows + |> json.array(fn(follow) { + json.object([ + #("follow_uri", json.string(follow.follow_uri)), + #("subject", json.string(follow.subject)), + #("created_at", json.string(follow.created_at)), + ]) + }) + |> json.to_string + let replaced = pog.transaction(conn, fn(transaction) { use _ <- result.try( pog.query("delete from follow_edges where did = $1") |> pog.parameter(pog.text(did)) |> pog.execute(transaction), ) - list.try_each(follows, fn(follow: Follow) { - pog.query( - "insert into follow_edges (follow_uri, did, subject, created_at) - values ($1, $2, $3, $4) + pog.query( + "insert into follow_edges (follow_uri, did, subject, created_at) + select follow_uri, $2, subject, created_at + from json_to_recordset($1::json) as rows( + follow_uri text, subject text, created_at text + ) on conflict (follow_uri) do update set did = excluded.did, subject = excluded.subject, created_at = excluded.created_at", - ) - |> pog.parameter(pog.text(follow.follow_uri)) - |> pog.parameter(pog.text(follow.did)) - |> pog.parameter(pog.text(follow.subject)) - |> pog.parameter(pog.text(follow.created_at)) - |> pog.execute(transaction) - }) + ) + |> pog.parameter(pog.text(encoded)) + |> pog.parameter(pog.text(did)) + |> pog.execute(transaction) |> result.map(fn(_) { Nil }) }) - case result { - Ok(_) -> Nil - Error(_) -> - wisp.log_warning("follow_edges reconciliation failed for " <> did) + |> result.replace_error("follow_edges reconciliation failed for " <> did) + case replaced { + Ok(Nil) -> Ok(Nil) + Error(message) -> { + wisp.log_warning(message) + Error(message) + } } } @@ -207,6 +254,82 @@ fn followers(conn: pog.Connection, subject: String) -> List(String) { |> db.all(conn) } +fn followers_page( + conn: pog.Connection, + subject: String, + after: Option(String), + limit: Int, +) -> List(String) { + connection_page(conn, "subject", "did", subject, after, limit) +} + +fn connection_page( + conn: pog.Connection, + scope_column: String, + value_column: String, + scope: String, + after: Option(String), + limit: Int, +) -> List(String) { + let row = { + use value <- decode.field("value", decode.string) + decode.success(value) + } + pog.query("with marker as ( + select $2::text as value + where $2::text is null or exists ( + select 1 from follow_edges + where " <> scope_column <> " = $1 and " <> value_column <> " = $2 + ) + ) + select distinct " <> value_column <> " as value + from follow_edges cross join marker + where " <> scope_column <> " = $1 + and (marker.value is null or " <> value_column <> " > marker.value) + order by value limit $3") + |> pog.parameter(pog.text(scope)) + |> pog.parameter(pog.nullable(pog.text, after)) + |> pog.parameter(pog.int(limit)) + |> pog.returning(row) + |> db.all(conn) +} + +fn following_among( + conn: pog.Connection, + did: String, + candidates: List(String), +) -> List(String) { + relationship_intersection(conn, "did", "subject", did, candidates) +} + +fn followers_among( + conn: pog.Connection, + subject: String, + candidates: List(String), +) -> List(String) { + relationship_intersection(conn, "subject", "did", subject, candidates) +} + +fn relationship_intersection( + conn: pog.Connection, + scope_column: String, + value_column: String, + scope: String, + candidates: List(String), +) -> List(String) { + let row = { + use value <- decode.field("value", decode.string) + decode.success(value) + } + pog.query("select distinct " <> value_column <> " as value from follow_edges + where " <> scope_column <> " = $1 and " <> value_column <> " = any($2) + order by value") + |> pog.parameter(pog.text(scope)) + |> pog.parameter(pog.array(pog.text, candidates)) + |> pog.returning(row) + |> db.all(conn) +} + fn count(conn: pog.Connection) -> Int { pog.query("select count(*)::int as count from follow_edges") |> db.count(conn) diff --git a/server/src/crate_server/handlers/connections.gleam b/server/src/crate_server/handlers/connections.gleam index 92ab363..15ecbb1 100644 --- a/server/src/crate_server/handlers/connections.gleam +++ b/server/src/crate_server/handlers/connections.gleam @@ -18,7 +18,7 @@ import gleam/dict.{type Dict} import gleam/int import gleam/json import gleam/list -import gleam/option.{None, Some} +import gleam/option.{type Option, None, Some} import gleam/result import gleam/set import gleam/string @@ -53,13 +53,25 @@ fn list_following( id: String, session: OauthSession, ) -> Response { - use client, session <- with_pds_client(ctx, id, session) - case graph_follows.load(client, session) { - Error(_) -> error_json(502, "could not load your follows from PDS") - Ok(stored) -> { - let following = list.map(stored, fn(row) { row.value.subject }) - index_follows(ctx, session.did, stored) - respond(ctx, session.did, Following, following, True, wisp.get_query(req)) + case ctx.follow_index.seen.has_seen(session.did) { + True -> respond_page(ctx, session.did, Following, True, wisp.get_query(req)) + False -> { + use client, session <- with_pds_client(ctx, id, session) + case graph_follows.load(client, session) { + Error(_) -> error_json(502, "could not load your follows from PDS") + Ok(stored) -> + case index_follows(ctx, session.did, stored) { + Error(_) -> error_json(503, "could not index your follows") + Ok(Nil) -> + respond_page( + ctx, + session.did, + Following, + True, + wisp.get_query(req), + ) + } + } } } } @@ -73,11 +85,7 @@ fn list_followers( session: OauthSession, ) -> Response { let known = ctx.follow_index.seen.has_seen(session.did) - let following = case known { - True -> ctx.follow_index.edges.following(session.did) - False -> [] - } - respond(ctx, session.did, Followers, following, known, wisp.get_query(req)) + respond_page(ctx, session.did, Followers, known, wisp.get_query(req)) } fn direction(query: List(#(String, String))) -> Result(Direction, Nil) { @@ -92,7 +100,7 @@ fn index_follows( ctx: Context, did: String, stored: List(StoredItem(GraphFollow)), -) -> Nil { +) -> Result(Nil, String) { let follows = stored |> list.map(fn(row) { @@ -103,39 +111,21 @@ fn index_follows( created_at: row.value.created_at, ) }) - ctx.follow_index.edges.replace_for_did(did, follows) + use _ <- result.map(ctx.follow_index.edges.replace_for_did(did, follows)) ctx.follow_index.seen.mark_seen(did) } -fn respond( +fn page_dids( ctx: Context, viewer_did: String, direction: Direction, - following: List(String), - following_complete: Bool, - query: List(#(String, String)), -) -> Response { - let followers = - ctx.follow_index.edges.followers(viewer_did) - |> list.unique - |> list.sort(string.compare) - let following_set = set.from_list(following) - let follower_set = set.from_list(followers) - let items = case direction { - Following -> - following - |> list.unique - |> list.sort(string.compare) - |> list.map(fn(did) { - connection(did, True, set.contains(follower_set, did)) - }) - Followers -> - followers - |> list.map(fn(did) { - connection(did, set.contains(following_set, did), True) - }) + after: Option(String), + limit: Int, +) -> List(String) { + case direction { + Following -> ctx.follow_index.edges.following_page(viewer_did, after, limit) + Followers -> ctx.follow_index.edges.followers_page(viewer_did, after, limit) } - respond_page(ctx, viewer_did, direction, items, following_complete, query) } fn connection( @@ -155,7 +145,6 @@ fn respond_page( ctx: Context, session_did: String, direction: Direction, - items: List(Connection), following_complete: Bool, query: List(#(String, String)), ) -> Response { @@ -164,14 +153,35 @@ fn respond_page( list.key_find(query, "limit") |> result.try(int.parse) |> option.from_result - let #(page, next_cursor) = - pagination.page( - items, - fn(item) { item.did }, - cursor_namespace(direction, session_did), - cursor, - Some(option.unwrap(limit, pagination.default_limit)), - ) + |> pagination.clamp_limit + let namespace = cursor_namespace(direction, session_did) + let after = cursor |> option.then(pagination.decode_cursor(namespace, _)) + let rows = page_dids(ctx, session_did, direction, after, limit + 1) + let page_dids = list.take(rows, limit) + let related = case direction, following_complete { + Following, _ -> + ctx.follow_index.edges.followers_among(session_did, page_dids) + Followers, True -> + ctx.follow_index.edges.following_among(session_did, page_dids) + Followers, False -> [] + } + let related_set = set.from_list(related) + let page = + page_dids + |> list.map(fn(did) { + case direction { + Following -> connection(did, True, set.contains(related_set, did)) + Followers -> connection(did, set.contains(related_set, did), True) + } + }) + let next_cursor = case list.length(rows) > limit { + False -> None + True -> + page_dids + |> list.last + |> result.map(fn(did) { pagination.encode_cursor(namespace, did) }) + |> option.from_result + } let identities = page |> list.map(fn(item) { item.did }) diff --git a/server/src/crate_server/handlers/feed.gleam b/server/src/crate_server/handlers/feed.gleam index c5bba37..9d0acb6 100644 --- a/server/src/crate_server/handlers/feed.gleam +++ b/server/src/crate_server/handlers/feed.gleam @@ -103,8 +103,10 @@ fn seed_follows( created_at: s.value.created_at, ) }) - ctx.follow_index.edges.replace_for_did(session.did, follows) - ctx.follow_index.seen.mark_seen(session.did) + case ctx.follow_index.edges.replace_for_did(session.did, follows) { + Ok(Nil) -> ctx.follow_index.seen.mark_seen(session.did) + Error(_) -> Nil + } list.map(stored, fn(s) { s.value.subject }) } } diff --git a/server/src/crate_server/handlers/graph.gleam b/server/src/crate_server/handlers/graph.gleam index 23bd83f..77e2780 100644 --- a/server/src/crate_server/handlers/graph.gleam +++ b/server/src/crate_server/handlers/graph.gleam @@ -176,8 +176,10 @@ fn index_stored_follows( created_at: row.value.created_at, ) }) - ctx.follow_index.edges.replace_for_did(session.did, follows) - ctx.follow_index.seen.mark_seen(session.did) + case ctx.follow_index.edges.replace_for_did(session.did, follows) { + Ok(Nil) -> ctx.follow_index.seen.mark_seen(session.did) + Error(_) -> Nil + } } fn delete_and_confirm( diff --git a/server/src/crate_server/pagination.gleam b/server/src/crate_server/pagination.gleam index ae36f5f..70b3ce9 100644 --- a/server/src/crate_server/pagination.gleam +++ b/server/src/crate_server/pagination.gleam @@ -100,7 +100,7 @@ fn resume_after( } } -fn clamp_limit(limit: Option(Int)) -> Int { +pub fn clamp_limit(limit: Option(Int)) -> Int { case limit { Some(n) if n >= 1 && n <= default_limit -> n _ -> default_limit diff --git a/server/src/crate_server/user_backfill.gleam b/server/src/crate_server/user_backfill.gleam index 4cb4827..64f9740 100644 --- a/server/src/crate_server/user_backfill.gleam +++ b/server/src/crate_server/user_backfill.gleam @@ -121,10 +121,11 @@ fn backfill_follows( )) }) case complete { - True -> { - index.edges.replace_for_did(user.did, rows) - index.seen.mark_seen(user.did) - } + True -> + case index.edges.replace_for_did(user.did, rows) { + Ok(Nil) -> index.seen.mark_seen(user.did) + Error(_) -> Nil + } False -> Nil } case complete { diff --git a/server/test/connections_handler_test.gleam b/server/test/connections_handler_test.gleam index f641bfc..72113f8 100644 --- a/server/test/connections_handler_test.gleam +++ b/server/test/connections_handler_test.gleam @@ -1,5 +1,3 @@ -//// Handler-level coverage for the read-only connections contract. - import atproto/xrpc import atproto_core/xrpc as core_xrpc import crate_server/context.{type Context, Context} @@ -175,9 +173,7 @@ pub fn an_unresolvable_identity_falls_back_to_the_did_test() { == Ok(["did:plc:unresolvable"]) } -// Following reads through to the PDS, so its failure is upstream failure -// rather than a silent fall back to whatever the index still remembers. -pub fn following_reports_a_pds_failure_and_keeps_the_index_test() { +pub fn known_following_is_served_from_the_index_without_the_pds_test() { let follows = support.fresh_follow_index() follows.edges.upsert(follower_edge("stale", "did:plc:me", "did:plc:stale")) follows.seen.mark_seen("did:plc:me") @@ -188,11 +184,13 @@ pub fn following_reports_a_pds_failure_and_keeps_the_index_test() { auth_session(), "graph.listConnections?direction=following", ) - assert resp.status == 502 + assert resp.status == 200 + assert support.field_nested(simulate.read_body(resp), ["items"], ["did"]) + == Ok(["did:plc:stale"]) assert follows.edges.following("did:plc:me") == ["did:plc:stale"] } -pub fn following_replaces_the_index_with_pds_data_test() { +pub fn unseen_following_reconciles_the_index_from_the_pds_test() { let follows = support.fresh_follow_index() follows.edges.upsert(follower_edge("stale", "did:plc:me", "did:plc:stale")) let resp = @@ -212,6 +210,28 @@ pub fn following_replaces_the_index_with_pds_data_test() { == ["did:plc:first", "did:plc:second"] } +pub fn unseen_following_does_not_mark_a_failed_reconciliation_seen_test() { + let base = support.fresh_follow_index() + let failing = + follow_index.Store( + ..base, + edges: follow_index.EdgeOps( + ..base.edges, + replace_for_did: fn(_did, _follows) { Error("write failed") }, + ), + ) + let resp = + respond( + failing, + list_records_client(["did:plc:first"]), + auth_session(), + "graph.listConnections?direction=following", + ) + + assert resp.status == 503 + assert base.seen.has_seen("did:plc:me") == False +} + pub fn a_connection_cursor_paginates_and_is_viewer_scoped_test() { let follows = support.fresh_follow_index() seed_followers(follows, "did:plc:me", ["did:plc:a", "did:plc:b"]) diff --git a/server/test/follow_index_postgres_test.gleam b/server/test/follow_index_postgres_test.gleam index 2d65897..8c35810 100644 --- a/server/test/follow_index_postgres_test.gleam +++ b/server/test/follow_index_postgres_test.gleam @@ -44,13 +44,28 @@ pub fn follow_round_trip_test() { == ["did:plc:subj-a", "did:plc:subj-b"] store.edges.upsert(follow("did:plc:other", did, "fr3")) assert store.edges.followers(did) == ["did:plc:other"] + assert store.edges.following_page(did, None, 1) == ["did:plc:subj-a"] + assert store.edges.following_page(did, Some("did:plc:subj-a"), 1) + == ["did:plc:subj-b"] + assert store.edges.followers_page(did, None, 1) == ["did:plc:other"] + assert store.edges.following_among(did, ["did:plc:subj-b"]) + == ["did:plc:subj-b"] + assert store.edges.followers_among(did, ["did:plc:other"]) + == ["did:plc:other"] + assert store.edges.following_among(did, []) == [] + assert store.edges.followers_among(did, []) == [] let assert Some(found) = store.edges.find(did, "did:plc:subj-a") assert found.created_at == "2024-01-01T00:00:00Z" assert store.edges.find(did, "did:plc:nobody") == None - store.edges.delete(follow(did, "did:plc:subj-a", "fr1").follow_uri) - assert store.edges.following(did) == ["did:plc:subj-b"] + assert store.edges.replace_for_did(did, [ + follow(did, "did:plc:replacement", "fr4"), + ]) + == Ok(Nil) + assert store.edges.following(did) == ["did:plc:replacement"] + assert store.edges.replace_for_did(did, []) == Ok(Nil) + assert store.edges.following(did) == [] store.edges.delete_for_did(did) assert store.edges.following(did) == [] } diff --git a/server/test/follow_index_test.gleam b/server/test/follow_index_test.gleam index 191f67f..bcacba6 100644 --- a/server/test/follow_index_test.gleam +++ b/server/test/follow_index_test.gleam @@ -73,11 +73,41 @@ pub fn replacing_a_dids_edges_reconciles_stale_rows_atomically_test() { store.edges.upsert(follow("did:a", "did:old", "f1")) store.edges.upsert(follow("did:a", "did:keep", "f2")) store.edges.upsert(follow("did:z", "did:old", "f3")) - store.edges.replace_for_did("did:a", [follow("did:a", "did:new", "f4")]) + assert store.edges.replace_for_did("did:a", [ + follow("did:a", "did:new", "f4"), + ]) + == Ok(Nil) assert store.edges.following("did:a") == ["did:new"] assert store.edges.following("did:z") == ["did:old"] } +pub fn connection_pages_are_distinct_bounded_and_resume_after_a_subject_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")) + store.edges.upsert(follow("did:a", "did:c", "f3")) + store.edges.upsert(follow("did:a", "did:d", "f4")) + store.edges.upsert(follow("did:z", "did:a", "f5")) + store.edges.upsert(follow("did:y", "did:a", "f6")) + + assert store.edges.following_page("did:a", None, 2) == ["did:b", "did:c"] + assert store.edges.following_page("did:a", Some("did:c"), 2) == ["did:d"] + assert store.edges.following_page("did:a", Some("did:missing"), 2) + == ["did:b", "did:c"] + assert store.edges.followers_page("did:a", None, 1) == ["did:y"] +} + +pub fn relationship_intersections_only_return_matching_candidates_test() { + let assert Ok(store) = follow_index.start() + store.edges.upsert(follow("did:a", "did:b", "f1")) + store.edges.upsert(follow("did:a", "did:c", "f2")) + store.edges.upsert(follow("did:b", "did:a", "f3")) + + assert store.edges.following_among("did:a", ["did:b", "did:z"]) == ["did:b"] + assert store.edges.followers_among("did:a", ["did:b", "did:z"]) == ["did:b"] + assert store.edges.following_among("did:a", []) == [] +} + pub fn delete_for_did_drops_that_dids_follows_and_seen_mark_test() { let assert Ok(store) = follow_index.start() store.edges.upsert(follow("did:a", "did:b", "f1")) diff --git a/server/test/graph_handler_test.gleam b/server/test/graph_handler_test.gleam index e76b7fd..f3de9e6 100644 --- a/server/test/graph_handler_test.gleam +++ b/server/test/graph_handler_test.gleam @@ -135,7 +135,6 @@ fn no_create_client(list_body: String) -> xrpc.Client { }) } -// 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 { diff --git a/web/src/crate_web/model.gleam b/web/src/crate_web/model.gleam index ca5c899..0ee7848 100644 --- a/web/src/crate_web/model.gleam +++ b/web/src/crate_web/model.gleam @@ -822,14 +822,10 @@ pub type Model { pressing: PressingState, // The network feed's load state, for the /feed route. feed: FeedState, - // Monotonic identity for feed skeleton and pagination requests. feed_generation: Int, - // The signed-in viewer's following/followers list. connections: ConnectionsState, - // Monotonic identities for list requests and optimistic writes. connections_generation: Int, connection_write_generation: Int, - // The follow/unfollow write currently in flight, if any. connection_pending: Option(ConnectionWrite), connections_loading_more: Bool, // Hydrated feed entry refs, keyed #(actor_did, entry_id).