diff --git a/server/src/crate_server/context.gleam b/server/src/crate_server/context.gleam index 37583f2..8c9a8e0 100644 --- a/server/src/crate_server/context.gleam +++ b/server/src/crate_server/context.gleam @@ -94,6 +94,7 @@ fn error_name(status: Int) -> String { case status { 400 -> "InvalidRequest" 401 -> "AuthenticationRequired" + 403 -> "Forbidden" 404 -> "NotFound" 409 -> "Conflict" 413 -> "PayloadTooLarge" diff --git a/server/src/crate_server/follow_index.gleam b/server/src/crate_server/follow_index.gleam index 583cde5..11d3a6b 100644 --- a/server/src/crate_server/follow_index.gleam +++ b/server/src/crate_server/follow_index.gleam @@ -13,6 +13,7 @@ import gleam/option.{type Option} import gleam/otp/actor import gleam/result import gleam/set.{type Set} +import gleam/string /// One graph.follow record, keyed by its own at-uri. The unfollow path's rkey /// is derived from that uri rather than stored beside it, so they cannot @@ -32,8 +33,18 @@ pub type EdgeOps { /// The dids `did` follows. Empty is ambiguous alone; pair with /// `SeenOps.has_seen`. following: fn(String) -> 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), /// 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 + /// duplicates, so cleanup callers need all matching records. + find_all: fn(String, String) -> List(Follow), + /// 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, /// Total indexed edges: a cheap completeness figure for `/debug/index`. count: fn() -> Int, ) @@ -61,7 +72,10 @@ type Msg { Delete(String) DeleteForDid(String) Following(String, Subject(List(String))) + Followers(String, Subject(List(String))) Find(String, String, Subject(Option(Follow))) + FindAll(String, String, Subject(List(Follow))) + ReplaceForDid(String, List(Follow)) Count(Subject(Int)) HasSeen(String, Subject(Bool)) MarkSeen(String) @@ -80,9 +94,16 @@ 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: []) }, + followers: fn(did) { parallel.ask(subject, Followers(did, _), or: []) }, find: fn(did, subject_did) { parallel.ask(subject, Find(did, subject_did, _), or: option.None) }, + find_all: fn(did, subject_did) { + parallel.ask(subject, FindAll(did, subject_did, _), or: []) + }, + replace_for_did: fn(did, follows) { + process.send(subject, ReplaceForDid(did, follows)) + }, count: fn() { parallel.ask(subject, Count, or: 0) }, ), seen: SeenOps( @@ -115,7 +136,18 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { reply, dict.values(state.follows) |> list.filter(fn(f) { f.did == did }) - |> list.map(fn(f) { f.subject }), + |> list.map(fn(f) { f.subject }) + |> list.sort(string.compare), + ) + actor.continue(state) + } + 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), ) actor.continue(state) } @@ -123,11 +155,32 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { process.send( reply, dict.values(state.follows) - |> list.find(fn(f) { f.did == did && f.subject == subject_did }) + |> list.filter(fn(f) { f.did == did && f.subject == subject_did }) + |> list.sort(fn(a, b) { string.compare(a.follow_uri, b.follow_uri) }) + |> list.first |> option.from_result, ) actor.continue(state) } + FindAll(did, subject_did, reply) -> { + process.send( + reply, + dict.values(state.follows) + |> list.filter(fn(f) { f.did == did && f.subject == subject_did }) + |> list.sort(fn(a, b) { string.compare(a.follow_uri, b.follow_uri) }), + ) + actor.continue(state) + } + ReplaceForDid(did, follows) -> { + let retained = + dict.filter(state.follows, fn(_uri, follow) { follow.did != did }) + let indexed = + follows + |> list.fold(retained, fn(edges, follow) { + dict.insert(edges, follow.follow_uri, follow) + }) + actor.continue(State(..state, follows: indexed)) + } Count(reply) -> { process.send(reply, dict.size(state.follows)) actor.continue(state) diff --git a/server/src/crate_server/follow_index_postgres.gleam b/server/src/crate_server/follow_index_postgres.gleam index 68a157d..8a34840 100644 --- a/server/src/crate_server/follow_index_postgres.gleam +++ b/server/src/crate_server/follow_index_postgres.gleam @@ -10,9 +10,11 @@ import crate_server/follow_index.{ } import crate_server/provenance import gleam/dynamic/decode +import gleam/list import gleam/option.{type Option} import gleam/result import pog +import wisp const columns = "follow_uri, did, subject, created_at" @@ -25,6 +27,9 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { delete_for_did: fn(did) { delete_for_did(conn, did) }, following: fn(did) { following(conn, did) }, 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) }, count: fn() { count(conn) }, ), seen: SeenOps(has_seen: fn(did) { has_seen(conn, did) }, mark_seen: fn(did) { @@ -62,6 +67,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_idx on follow_edges (subject)", + ) + |> pog.execute(conn) + |> result.replace_error("follow_edges subject index migration failed"), + ) pog.query( "create table if not exists follow_index_seen ( did text primary key, @@ -118,7 +130,7 @@ fn following(conn: pog.Connection, did: String) -> List(String) { use subject <- decode.field("subject", decode.string) decode.success(subject) } - pog.query("select subject from follow_edges where did = $1") + pog.query("select subject from follow_edges where did = $1 order by subject") |> pog.parameter(pog.text(did)) |> pog.returning(row) |> db.all(conn) @@ -126,13 +138,75 @@ fn following(conn: pog.Connection, did: String) -> List(String) { fn find(conn: pog.Connection, did: String, subject: String) -> Option(Follow) { pog.query("select " <> columns <> " - from follow_edges where did = $1 and subject = $2 limit 1") + from follow_edges where did = $1 and subject = $2 + order by follow_uri limit 1") |> pog.parameter(pog.text(did)) |> pog.parameter(pog.text(subject)) |> pog.returning(row_decoder()) |> db.first(conn) } +fn find_all( + conn: pog.Connection, + did: String, + subject: String, +) -> List(Follow) { + pog.query("select " <> columns <> " + from follow_edges where did = $1 and subject = $2 + order by follow_uri") + |> pog.parameter(pog.text(did)) + |> pog.parameter(pog.text(subject)) + |> pog.returning(row_decoder()) + |> db.all(conn) +} + +fn replace_for_did( + conn: pog.Connection, + did: String, + follows: List(Follow), +) -> Nil { + let result = + 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) + 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) + }) + |> result.map(fn(_) { Nil }) + }) + case result { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning("follow_edges reconciliation failed for " <> did) + } +} + +fn followers(conn: pog.Connection, subject: String) -> List(String) { + let row = { + use did <- decode.field("did", decode.string) + decode.success(did) + } + pog.query("select did from follow_edges where subject = $1 order by did") + |> pog.parameter(pog.text(subject)) + |> 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/graph_follows.gleam b/server/src/crate_server/graph_follows.gleam index b3d357a..dc29537 100644 --- a/server/src/crate_server/graph_follows.gleam +++ b/server/src/crate_server/graph_follows.gleam @@ -7,10 +7,13 @@ import crate/gen/graph/follow.{type GraphFollow, GraphFollow} import crate/storage.{type StoredItem} import crate_server/oauth/sessions.{type OauthSession} import crate_server/provenance +import gleam/bit_array +import gleam/crypto import gleam/dynamic/decode import gleam/list import gleam/option.{type Option} import gleam/result +import gleam/string pub fn load( client: Client, @@ -26,10 +29,10 @@ pub fn load( ) } -/// Create a follow record for `subject_did`; duplicates are the caller's -/// problem to guard against, per the lexicon's doc comment. Returns the record -/// written, so a caller mirroring it into the index stores the timestamp that -/// actually reached the PDS. +/// Create a follow record for `subject_did` at a deterministic rkey. Repeated +/// requests, including concurrent requests, therefore converge on one PDS +/// record through putRecord. Legacy random-rkey records remain readable and +/// are removed together by the unfollow handler. pub fn create( client: Client, session: OauthSession, @@ -37,25 +40,45 @@ pub fn create( ) -> Result(#(repo.CreatedRecord, GraphFollow), XrpcError) { let record = GraphFollow(created_at: provenance.now_rfc3339(), subject: subject_did) - repo.create_record( + repo.put_record( client, session.pds, session.access_token, session.did, follow.collection, + rkey_for_subject(subject_did), follow.encode_graph_follow(record), ) |> result.map(fn(created) { #(created, record) }) } +/// Rkeys are URL-safe and stable across clients, while the hash avoids +/// characters that AT Protocol rkey validation rejects in a DID. +pub fn rkey_for_subject(subject_did: String) -> String { + "f_" + <> bit_array.base64_url_encode( + crypto.hash(crypto.Sha256, bit_array.from_string(subject_did)), + False, + ) +} + /// The stored follow record for `subject_did`, if any. pub fn find( stored: List(StoredItem(GraphFollow)), subject_did: String, ) -> Option(StoredItem(GraphFollow)) { + find_all(stored, subject_did) |> list.first |> option.from_result +} + +/// All stored records for a subject, in a stable URI order. Stable ordering +/// makes retries and cleanup deterministic when a legacy repo has duplicates. +pub fn find_all( + stored: List(StoredItem(GraphFollow)), + subject_did: String, +) -> List(StoredItem(GraphFollow)) { stored - |> list.find(fn(s) { s.value.subject == subject_did }) - |> option.from_result + |> list.filter(fn(s) { s.value.subject == subject_did }) + |> list.sort(fn(a, b) { string.compare(a.uri, b.uri) }) } pub fn delete( diff --git a/server/src/crate_server/handlers/connections.gleam b/server/src/crate_server/handlers/connections.gleam new file mode 100644 index 0000000..92ab363 --- /dev/null +++ b/server/src/crate_server/handlers/connections.gleam @@ -0,0 +1,223 @@ +//// Read-only connections for the signed-in viewer. Following is authoritative +//// after one successful read from the viewer's PDS. Followers come from the +//// appview index and deliberately report incomplete until a global index +//// watermark exists. + +import crate/gen/graph/follow.{type GraphFollow} +import crate/storage.{type StoredItem} +import crate_server/context.{ + type Context, error_json, require_session, with_pds_client, +} +import crate_server/follow_index +import crate_server/graph_follows +import crate_server/identity_cache.{type Identity} +import crate_server/identity_resolver +import crate_server/oauth/sessions.{type OauthSession} +import crate_server/pagination +import gleam/dict.{type Dict} +import gleam/int +import gleam/json +import gleam/list +import gleam/option.{None, Some} +import gleam/result +import gleam/set +import gleam/string +import wisp.{type Request, type Response} + +type Direction { + Following + Followers +} + +type Connection { + Connection( + did: String, + viewer_follows: Bool, + follows_viewer: Bool, + mutual: Bool, + ) +} + +pub fn list_connections(req: Request, ctx: Context) -> Response { + use id, session <- require_session(req, ctx) + case direction(wisp.get_query(req)) { + Error(Nil) -> error_json(400, "expected direction=following or followers") + Ok(Following) -> list_following(req, ctx, id, session) + Ok(Followers) -> list_followers(req, ctx, session) + } +} + +fn list_following( + req: Request, + ctx: Context, + 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)) + } + } +} + +/// Followers are appview-indexed and can be read without reaching the +/// viewer's PDS. If that viewer has not had a complete following read yet, +/// relationship flags remain conservative and the response says so. +fn list_followers( + req: Request, + ctx: Context, + 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)) +} + +fn direction(query: List(#(String, String))) -> Result(Direction, Nil) { + case list.key_find(query, "direction") { + Ok("following") -> Ok(Following) + Ok("followers") -> Ok(Followers) + _ -> Error(Nil) + } +} + +fn index_follows( + ctx: Context, + did: String, + stored: List(StoredItem(GraphFollow)), +) -> Nil { + let follows = + stored + |> list.map(fn(row) { + follow_index.Follow( + follow_uri: row.uri, + did:, + subject: row.value.subject, + created_at: row.value.created_at, + ) + }) + ctx.follow_index.edges.replace_for_did(did, follows) + ctx.follow_index.seen.mark_seen(did) +} + +fn respond( + 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) + }) + } + respond_page(ctx, viewer_did, direction, items, following_complete, query) +} + +fn connection( + did: String, + viewer_follows: Bool, + follows_viewer: Bool, +) -> Connection { + Connection( + did:, + viewer_follows:, + follows_viewer:, + mutual: viewer_follows && follows_viewer, + ) +} + +fn respond_page( + ctx: Context, + session_did: String, + direction: Direction, + items: List(Connection), + following_complete: Bool, + query: List(#(String, String)), +) -> Response { + let cursor = list.key_find(query, "cursor") |> option.from_result + let limit = + 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)), + ) + let identities = + page + |> list.map(fn(item) { item.did }) + |> identity_resolver.resolve_many(ctx.atproto.identity, _) + let complete = case direction { + Following -> True + Followers -> False + } + let relationship_complete = case direction { + Following -> False + Followers -> following_complete + } + json.object([ + #("items", json.array(page, encode_connection(identities, _))), + #("complete", json.bool(complete)), + #("relationshipComplete", json.bool(relationship_complete)), + ..case next_cursor { + Some(cursor) -> [#("cursor", json.string(cursor))] + None -> [] + } + ]) + |> json.to_string + |> wisp.json_response(200) +} + +fn cursor_namespace(direction: Direction, viewer_did: String) -> String { + let prefix = case direction { + Following -> "connections-following-" + Followers -> "connections-followers-" + } + prefix <> string.replace(viewer_did, ":", "-") +} + +fn encode_connection( + identities: Dict(String, Identity), + item: Connection, +) -> json.Json { + let handle = case dict.get(identities, item.did) { + Ok(identity) -> identity.handle + Error(Nil) -> item.did + } + json.object([ + #("did", json.string(item.did)), + #("handle", json.string(handle)), + #("viewerFollows", json.bool(item.viewer_follows)), + #("followsViewer", json.bool(item.follows_viewer)), + #("mutual", json.bool(item.mutual)), + ]) +} diff --git a/server/src/crate_server/handlers/feed.gleam b/server/src/crate_server/handlers/feed.gleam index f7f0558..c5bba37 100644 --- a/server/src/crate_server/handlers/feed.gleam +++ b/server/src/crate_server/handlers/feed.gleam @@ -28,7 +28,12 @@ import wisp.{type Request, type Response} // (a bare did always does, e.g. "did:plc:...") would never round-trip. const following_namespace_prefix = "feed-following-" -const network_namespace = "feed-network" +const network_namespace_prefix = "feed-network-" + +type FeedMode { + Following + Network +} pub fn get_feed_skeleton(req: Request, ctx: Context) -> Response { use id, session <- require_session(req, ctx) @@ -88,15 +93,17 @@ fn seed_follows( case graph_follows.load(client, session) { Error(_) -> [] Ok(stored) -> { - stored - |> list.each(fn(s) { - ctx.follow_index.edges.upsert(follow_index.Follow( - follow_uri: s.uri, - did: session.did, - subject: s.value.subject, - created_at: s.value.created_at, - )) - }) + let follows = + stored + |> list.map(fn(s) { + follow_index.Follow( + follow_uri: s.uri, + did: session.did, + subject: s.value.subject, + created_at: s.value.created_at, + ) + }) + ctx.follow_index.edges.replace_for_did(session.did, follows) ctx.follow_index.seen.mark_seen(session.did) list.map(stored, fn(s) { s.value.subject }) } @@ -106,9 +113,9 @@ fn seed_follows( /// Followed-mode first; falls back to network-wide (unfiltered, viewer /// included -- simpler than re-excluding them from their own discovery feed) /// when the follow graph is empty, or when a first request (no cursor yet) -/// finds nothing from people the caller follows. A cursor already scoped to -/// one namespace never resumes the other: `pagination.page` restarts a -/// foreign-namespace cursor from the top of whichever list it's handed. +/// finds nothing from people the caller follows. A cursor selects the mode it +/// was minted for, so a network continuation never passes through followed +/// pagination. Invalid or foreign cursors restart a fresh followed request. fn resolve_page( adoptions: List(Adoption), followed: List(String), @@ -123,33 +130,81 @@ fn resolve_page( a.did != viewer_did && set.contains(followed_set, a.did) }) |> sort_desc - let #(followed_page, followed_cursor) = + let #(mode, mode_cursor) = mode_for_cursor(viewer_did, cursor) + case mode { + Network -> network_page(adoptions, viewer_did, mode_cursor, limit) + Following -> { + let #(followed_page, followed_cursor) = + pagination.page( + followed_rows, + fn(a: Adoption) { a.entry_uri }, + following_namespace(viewer_did), + mode_cursor, + limit, + ) + let needs_fallback = mode_cursor == None && list.is_empty(followed_page) + case needs_fallback { + True -> network_page(adoptions, viewer_did, None, limit) + False -> #( + feed_skeleton.collapse(followed_page), + followed_cursor, + False, + ) + } + } + } +} + +fn network_page( + adoptions: List(Adoption), + viewer_did: String, + cursor: Option(String), + limit: Option(Int), +) -> #(List(feed_gen.FeedItem), Option(String), Bool) { + let #(page, next_cursor) = pagination.page( - followed_rows, + sort_desc(adoptions), fn(a: Adoption) { a.entry_uri }, - following_namespace_prefix <> string.replace(viewer_did, ":", "-"), + network_namespace(viewer_did), cursor, limit, ) - let needs_fallback = - list.is_empty(followed) - || { cursor == None && list.is_empty(followed_page) } - case needs_fallback { - True -> { - let #(network_page, network_cursor) = - pagination.page( - sort_desc(adoptions), - fn(a: Adoption) { a.entry_uri }, - network_namespace, - cursor, - limit, - ) - #(feed_skeleton.collapse(network_page), network_cursor, True) + #(feed_skeleton.collapse(page), next_cursor, True) +} + +fn mode_for_cursor( + viewer_did: String, + cursor: Option(String), +) -> #(FeedMode, Option(String)) { + case cursor { + None -> #(Following, None) + Some(token) -> { + case pagination.decode_cursor(network_namespace(viewer_did), token) { + Some(_) -> #(Network, Some(token)) + None -> + case + pagination.decode_cursor(following_namespace(viewer_did), token) + { + Some(_) -> #(Following, Some(token)) + None -> #(Following, None) + } + } } - False -> #(feed_skeleton.collapse(followed_page), followed_cursor, False) } } +fn following_namespace(viewer_did: String) -> String { + following_namespace_prefix <> cursor_scope(viewer_did) +} + +fn network_namespace(viewer_did: String) -> String { + network_namespace_prefix <> cursor_scope(viewer_did) +} + +fn cursor_scope(viewer_did: String) -> String { + string.replace(viewer_did, ":", "-") +} + fn sort_desc(rows: List(Adoption)) -> List(Adoption) { list.sort(rows, fn(a, b) { string.compare(b.created_at, a.created_at) }) } diff --git a/server/src/crate_server/handlers/graph.gleam b/server/src/crate_server/handlers/graph.gleam index 9c9314d..1f6f5b8 100644 --- a/server/src/crate_server/handlers/graph.gleam +++ b/server/src/crate_server/handlers/graph.gleam @@ -3,8 +3,10 @@ //// the feed sees it without waiting for the firehose to echo it back. import atproto/uri -import atproto/xrpc.{type Client} +import atproto/xrpc.{type Client, type XrpcError} +import atproto_core/xrpc as core_xrpc import crate/gen/graph/follow.{type GraphFollow} +import crate/storage.{type StoredItem} import crate_server/context.{ type Context, error_json, require_session, with_pds_client, } @@ -13,7 +15,8 @@ import crate_server/graph_follows import crate_server/oauth/sessions.{type OauthSession} import gleam/dynamic/decode import gleam/json -import gleam/option.{type Option, None, Some} +import gleam/list +import gleam/option.{Some} import wisp.{type Request, type Response} fn subject_decoder() -> decode.Decoder(String) { @@ -37,17 +40,63 @@ fn do_follow( subject_did: String, ) -> Response { use client, session <- with_pds_client(ctx, id, session) + case stored_follows(ctx, client, session, subject_did) { + Error(error) -> pds_error(error, "could not load your follows from PDS") + Ok(uris) -> + case list.first(uris) { + Ok(uri) -> follow_response(uri) + Error(Nil) -> create_follow(ctx, client, session, subject_did) + } + } +} + +fn create_follow( + ctx: Context, + client: Client, + session: OauthSession, + subject_did: String, +) -> Response { case graph_follows.create(client, session, subject_did) { - Error(_) -> error_json(502, "could not write follow record to PDS") + Error(error) -> pds_error(error, "could not write follow record to PDS") Ok(#(created, record)) -> { index_follow(ctx, session, created.uri, record) - json.object([#("uri", json.string(created.uri))]) - |> json.to_string - |> wisp.json_response(200) + follow_response(created.uri) } } } +fn follow_response(uri: String) -> Response { + json.object([#("uri", json.string(uri))]) + |> json.to_string + |> wisp.json_response(200) +} + +/// Existing sessions may not have the graph.follow scope that current OAuth +/// requests ask for. Preserve the authorization signal so the browser can +/// offer an explicit reauthorization instead of hiding it as an upstream +/// outage. +fn pds_error(error: XrpcError, fallback: String) -> Response { + case error { + core_xrpc.BadStatus(status, Some(code), _, _) + if status == 403 + || code == "InsufficientScope" + || code == "insufficient_scope" + -> + error_json( + 403, + "your session needs permission to follow collectors; sign in again to approve it", + ) + core_xrpc.BadStatus(status, _, _, _) if status == 403 -> + error_json( + 403, + "your session needs permission to follow collectors; sign in again to approve it", + ) + core_xrpc.BadStatus(status, _, _, _) if status == 401 -> + error_json(401, "your session expired while writing the follow record") + _ -> error_json(502, fallback) + } +} + fn index_follow( ctx: Context, session: OauthSession, @@ -78,49 +127,77 @@ fn do_unfollow( subject_did: String, ) -> Response { use client, session <- with_pds_client(ctx, id, session) - case stored_follow(ctx, client, session, subject_did) { - Error(Nil) -> error_json(502, "could not load your follows from PDS") - Ok(None) -> error_json(404, "not following that subject") - Ok(Some(follow_uri)) -> delete_and_confirm(ctx, client, session, follow_uri) + case stored_follows(ctx, client, session, subject_did) { + Error(error) -> pds_error(error, "could not load your follows from PDS") + Ok([]) -> error_json(404, "not following that subject") + Ok(follow_uris) -> delete_and_confirm(ctx, client, session, follow_uris) } } // Index-first; a caller the index has never seen still falls back to their -// own repo, same exception the feed makes. -fn stored_follow( +// own repo, same exception the feed makes. The index lookup returns every +// matching URI so legacy duplicate records are removed as one relationship. +fn stored_follows( ctx: Context, client: Client, session: OauthSession, subject_did: String, -) -> Result(Option(String), Nil) { +) -> Result(List(String), XrpcError) { case ctx.follow_index.seen.has_seen(session.did) { True -> Ok( - ctx.follow_index.edges.find(session.did, subject_did) - |> option.map(fn(f) { f.follow_uri }), + ctx.follow_index.edges.find_all(session.did, subject_did) + |> list.map(fn(f) { f.follow_uri }), ) False -> case graph_follows.load(client, session) { - Error(_) -> Error(Nil) - Ok(stored) -> + Error(error) -> Error(error) + Ok(stored) -> { + index_stored_follows(ctx, session, stored) Ok( - graph_follows.find(stored, subject_did) - |> option.map(fn(s) { s.uri }), + graph_follows.find_all(stored, subject_did) + |> list.map(fn(s) { s.uri }), ) + } } } } +fn index_stored_follows( + ctx: Context, + session: OauthSession, + stored: List(StoredItem(GraphFollow)), +) -> Nil { + let follows = + stored + |> list.map(fn(row) { + follow_index.Follow( + follow_uri: row.uri, + did: session.did, + subject: row.value.subject, + created_at: row.value.created_at, + ) + }) + ctx.follow_index.edges.replace_for_did(session.did, follows) + ctx.follow_index.seen.mark_seen(session.did) +} + fn delete_and_confirm( ctx: Context, client: Client, session: OauthSession, - follow_uri: String, + follow_uris: List(String), ) -> Response { - case graph_follows.delete(client, session, uri.rkey(follow_uri)) { - Error(_) -> error_json(502, "could not delete follow record on PDS") - Ok(Nil) -> { - ctx.follow_index.edges.delete(follow_uri) + case + list.try_each(follow_uris, fn(follow_uri) { + graph_follows.delete(client, session, uri.rkey(follow_uri)) + }) + { + Error(error) -> pds_error(error, "could not delete follow record on PDS") + Ok(_) -> { + list.each(follow_uris, fn(follow_uri) { + ctx.follow_index.edges.delete(follow_uri) + }) wisp.json_response("{}", 200) } } diff --git a/server/src/crate_server/oauth/client_metadata.gleam b/server/src/crate_server/oauth/client_metadata.gleam index 2ed1c30..22ad0e3 100644 --- a/server/src/crate_server/oauth/client_metadata.gleam +++ b/server/src/crate_server/oauth/client_metadata.gleam @@ -6,6 +6,7 @@ import crate/gen/catalog/artist as catalog_artist import crate/gen/catalog/edit as catalog_edit import crate/gen/catalog/genre as catalog_genre import crate/gen/catalog/release as catalog_release +import crate/gen/graph/follow as graph_follow import crate/gen/shelf/entry as shelf_entry import gleam/json import gleam/string @@ -31,6 +32,8 @@ pub const scopes = [ <> catalog_genre.collection, "repo:" <> catalog_edit.collection, + "repo:" + <> graph_follow.collection, "blob:image/*", ] diff --git a/server/src/crate_server/router.gleam b/server/src/crate_server/router.gleam index c1277b8..5672d78 100644 --- a/server/src/crate_server/router.gleam +++ b/server/src/crate_server/router.gleam @@ -4,6 +4,7 @@ import crate_server/context.{type Context, error_json} import crate_server/handlers/amend import crate_server/handlers/browse +import crate_server/handlers/connections import crate_server/handlers/cover_proxy import crate_server/handlers/crate_overlap import crate_server/handlers/debug_index @@ -88,6 +89,7 @@ fn crate_xrpc( "shelf.getCrateOverlap", Get -> crate_overlap.get_crate_overlap(req, ctx) "graph.followUser", Post -> graph.follow(req, ctx) "graph.unfollowUser", Post -> graph.unfollow(req, ctx) + "graph.listConnections", Get -> connections.list_connections(req, ctx) "feed.getFeedSkeleton", Get -> feed.get_feed_skeleton(req, ctx) "discogs.searchReleases", Get -> discogs.search(req, ctx) "discogs.searchArtists", Get -> discogs.search_artists(req, ctx) diff --git a/server/src/crate_server/user_backfill.gleam b/server/src/crate_server/user_backfill.gleam index 93319ca..4cb4827 100644 --- a/server/src/crate_server/user_backfill.gleam +++ b/server/src/crate_server/user_backfill.gleam @@ -120,11 +120,17 @@ fn backfill_follows( created_at: follow.created_at, )) }) - rows |> list.each(index.edges.upsert) case complete { - True -> index.seen.mark_seen(user.did) + True -> { + index.edges.replace_for_did(user.did, rows) + index.seen.mark_seen(user.did) + } False -> Nil } + case complete { + True -> Nil + False -> rows |> list.each(index.edges.upsert) + } list.length(rows) } @@ -298,7 +304,24 @@ fn fetch_all_checked( decode_row: fn(RecordEntry) -> Result(a, Nil), ) -> #(List(a), Bool) { let #(records, complete) = fetch_pages(client, user, collection, None, [], 0) - #(list.filter_map(records, decode_row), complete) + let #(rows, decoded) = decode_rows_checked(records, decode_row, [], True) + #(rows, complete && decoded) +} + +fn decode_rows_checked( + records: List(RecordEntry), + decode_row: fn(RecordEntry) -> Result(a, Nil), + acc: List(a), + complete: Bool, +) -> #(List(a), Bool) { + case records { + [] -> #(list.reverse(acc), complete) + [record, ..rest] -> + case decode_row(record) { + Ok(row) -> decode_rows_checked(rest, decode_row, [row, ..acc], complete) + Error(_) -> decode_rows_checked(rest, decode_row, acc, False) + } + } } /// `fetch_all`'s undecoded form: the raw paged rows, so a caller needing two diff --git a/server/test/connections_handler_test.gleam b/server/test/connections_handler_test.gleam new file mode 100644 index 0000000..6a2d93f --- /dev/null +++ b/server/test/connections_handler_test.gleam @@ -0,0 +1,280 @@ +//// 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} +import crate_server/follow_index +import crate_server/handlers/connections +import crate_server/oauth/config +import crate_server/oauth/session_store +import crate_server/oauth/sessions +import gleam/bit_array +import gleam/dynamic/decode +import gleam/http +import gleam/http/response +import gleam/int +import gleam/json +import gleam/list +import gleam/result +import support +import wisp +import wisp/simulate + +const session_cookie = "ar_oauth_sid" + +fn auth_session() -> sessions.OauthSession { + let session = support.stub_session() + sessions.OauthSession(..session, expires_at: 9_999_999_999) +} + +fn authed_get( + path: String, + cfg: config.Config, + session: sessions.OauthSession, +) -> wisp.Request { + let assert Ok(id) = session_store.create(cfg.sessions, session) + simulate.request(http.Get, path) + |> simulate.cookie(session_cookie, id, wisp.Signed) +} + +fn context_with( + client: xrpc.Client, + follows: follow_index.Store, +) -> #(Context, config.Config) { + let cfg = support.stub_config(client) + let ctx = + Context( + ..support.stub_context_with( + cfg, + support.unreachable_catalog_deps(), + fn(_req) { Error("unused") }, + [], + ), + follow_index: follows, + ) + #(ctx, cfg) +} + +fn ok_response( + body: String, +) -> Result(response.Response(BitArray), xrpc.TransportError) { + Ok(response.Response(200, [], bit_array.from_string(body))) +} + +fn list_records_client(body: String) -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> ok_response(body) + _ -> Error(core_xrpc.ConnectionFailed("unused")) + } + }) +} + +fn item_bools(body: String, field: String) -> Result(List(Bool), Nil) { + json.parse( + body, + decode.at( + ["items"], + decode.list(decode.field(field, decode.bool, decode.success)), + ), + ) + |> result.replace_error(Nil) +} + +fn follower_edge( + uri: String, + did: String, + subject: String, +) -> follow_index.Follow { + follow_index.Follow( + follow_uri: "at://" <> did <> "/dev.mokkenstorm.crate.graph.follow/" <> uri, + did:, + subject:, + created_at: "2026-01-01T00:00:00Z", + ) +} + +fn seed_followers( + store: follow_index.Store, + viewer: String, + dids: List(String), +) -> Nil { + dids + |> list.each(fn(did) { + store.edges.upsert(follower_edge("r" <> did, did, viewer)) + }) +} + +fn list_body(subjects: List(String)) -> String { + subjects + |> list.index_map(fn(subject, index) { + json.object([ + #( + "uri", + json.string( + "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/r" + <> int.to_string(index), + ), + ), + #("cid", json.string("cid" <> int.to_string(index))), + #( + "value", + json.object([ + #("subject", json.string(subject)), + #("createdAt", json.string("2026-01-01T00:00:00Z")), + ]), + ), + ]) + }) + |> json.preprocessed_array + |> fn(records) { json.object([#("records", records)]) } + |> json.to_string +} + +pub fn unauthenticated_connections_are_rejected_test() { + let #(ctx, _cfg) = + context_with(support.unreachable_client(), support.fresh_follow_index()) + let req = + simulate.request( + http.Get, + support.xrpc("graph.listConnections?direction=followers"), + ) + assert connections.list_connections(req, ctx).status == 401 +} + +pub fn followers_do_not_require_the_viewers_pds_test() { + let follows = support.fresh_follow_index() + seed_followers(follows, "did:plc:me", ["did:plc:follower"]) + follows.seen.mark_seen("did:plc:me") + let #(ctx, cfg) = context_with(support.unreachable_client(), follows) + let req = + authed_get( + support.xrpc("graph.listConnections?direction=followers"), + cfg, + auth_session(), + ) + let resp = connections.list_connections(req, ctx) + assert resp.status == 200 + let body = simulate.read_body(resp) + assert support.field_nested(body, ["items"], ["did"]) + == Ok(["did:plc:follower"]) + assert support.field_bool(body, ["relationshipComplete"]) == Ok(True) +} + +pub fn cold_followers_report_incomplete_relationships_test() { + let follows = support.fresh_follow_index() + seed_followers(follows, "did:plc:me", ["did:plc:follower"]) + let #(ctx, cfg) = context_with(support.unreachable_client(), follows) + let req = + authed_get( + support.xrpc("graph.listConnections?direction=followers"), + cfg, + auth_session(), + ) + let body = simulate.read_body(connections.list_connections(req, ctx)) + assert support.field_bool(body, ["relationshipComplete"]) == Ok(False) + assert item_bools(body, "viewerFollows") == Ok([False]) +} + +pub fn following_pds_failure_is_upstream_failure_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") + let #(ctx, cfg) = context_with(support.unreachable_client(), follows) + let req = + authed_get( + support.xrpc("graph.listConnections?direction=following"), + cfg, + auth_session(), + ) + assert connections.list_connections(req, ctx).status == 502 + assert follows.edges.following("did:plc:me") == ["did:plc:stale"] +} + +pub fn identity_failure_falls_back_to_the_did_test() { + let follows = support.fresh_follow_index() + seed_followers(follows, "did:plc:me", ["did:plc:unresolvable"]) + follows.seen.mark_seen("did:plc:me") + let #(ctx, cfg) = context_with(support.unreachable_client(), follows) + let req = + authed_get( + support.xrpc("graph.listConnections?direction=followers"), + cfg, + auth_session(), + ) + let body = simulate.read_body(connections.list_connections(req, ctx)) + assert support.field_nested(body, ["items"], ["handle"]) + == Ok(["did:plc:unresolvable"]) +} + +pub fn following_uses_pds_data_and_marks_the_viewer_complete_test() { + let follows = support.fresh_follow_index() + follows.edges.upsert(follower_edge("stale", "did:plc:me", "did:plc:stale")) + let #(ctx, cfg) = + context_with( + list_records_client(list_body(["did:plc:first", "did:plc:second"])), + follows, + ) + let req = + authed_get( + support.xrpc("graph.listConnections?direction=following&limit=1"), + cfg, + auth_session(), + ) + let resp = connections.list_connections(req, ctx) + assert resp.status == 200 + let body = simulate.read_body(resp) + assert support.field_nested(body, ["items"], ["did"]) == Ok(["did:plc:first"]) + assert support.field_bool(body, ["complete"]) == Ok(True) + assert support.field_bool(body, ["relationshipComplete"]) == Ok(False) + assert support.field_present(body, ["cursor"]) + assert follows.edges.following("did:plc:me") + == [ + "did:plc:first", + "did:plc:second", + ] +} + +pub fn 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"]) + follows.seen.mark_seen("did:plc:me") + let #(ctx, cfg) = context_with(support.unreachable_client(), follows) + let first_req = + authed_get( + support.xrpc("graph.listConnections?direction=followers&limit=1"), + cfg, + auth_session(), + ) + let first_body = + simulate.read_body(connections.list_connections(first_req, ctx)) + let assert Ok(cursor) = support.field_string(first_body, ["cursor"]) + let next_req = + authed_get( + support.xrpc( + "graph.listConnections?direction=followers&limit=1&cursor=" <> cursor, + ), + cfg, + auth_session(), + ) + let next_body = + simulate.read_body(connections.list_connections(next_req, ctx)) + assert support.field_nested(next_body, ["items"], ["did"]) + == Ok(["did:plc:b"]) + + let other = sessions.OauthSession(..auth_session(), did: "did:plc:other") + seed_followers(follows, "did:plc:other", ["did:plc:c"]) + follows.seen.mark_seen("did:plc:other") + let other_req = + authed_get( + support.xrpc( + "graph.listConnections?direction=followers&limit=1&cursor=" <> cursor, + ), + cfg, + other, + ) + let other_body = + simulate.read_body(connections.list_connections(other_req, ctx)) + assert support.field_nested(other_body, ["items"], ["did"]) + == Ok(["did:plc:c"]) +} diff --git a/server/test/feed_handler_test.gleam b/server/test/feed_handler_test.gleam index 211cc94..32662df 100644 --- a/server/test/feed_handler_test.gleam +++ b/server/test/feed_handler_test.gleam @@ -36,12 +36,14 @@ fn feed_path() -> String { support.xrpc("feed.getFeedSkeleton") } -fn a_session() -> sessions.OauthSession { - support.stub_session_with( - access_token: "at", - refresh_token: "rt", - expires_at: far_future, - ) +fn a_session_as(did: String) -> sessions.OauthSession { + let session = + support.stub_session_with( + access_token: "at", + refresh_token: "rt", + expires_at: far_future, + ) + sessions.OauthSession(..session, did:) } /// Always an "acquired" row -- no test in this file varies the status. @@ -126,7 +128,15 @@ fn test_context( } fn authed_get(path: String, cfg: config.Config) -> wisp.Request { - let assert Ok(id) = session_store.create(cfg.sessions, a_session()) + authed_get_as(path, cfg, viewer_did) +} + +fn authed_get_as( + path: 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.cookie(session_cookie, id, wisp.Signed) } @@ -174,12 +184,21 @@ pub fn followed_rows_are_included_and_viewer_and_non_followed_are_excluded_test( network_client( list_records_body([follow_record_json("3aaa", "did:plc:f1")]), ) - let #(ctx, cfg, _follows) = test_context(client, store) + 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) 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() { @@ -228,6 +247,79 @@ pub fn nonempty_follows_with_no_followed_adoptions_falls_back_on_the_first_page_ 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 +} + +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"] +} + pub fn cursor_stays_in_the_followed_namespace_across_pages_test() { let assert Ok(store) = catalog_index.start() store.adoptions.upsert(adoption( diff --git a/server/test/follow_index_postgres_test.gleam b/server/test/follow_index_postgres_test.gleam index fe9d16f..2d65897 100644 --- a/server/test/follow_index_postgres_test.gleam +++ b/server/test/follow_index_postgres_test.gleam @@ -42,6 +42,8 @@ pub fn follow_round_trip_test() { store.edges.upsert(follow(did, "did:plc:subj-b", "fr2")) assert list.sort(store.edges.following(did), string.compare) == ["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"] let assert Some(found) = store.edges.find(did, "did:plc:subj-a") assert found.created_at == "2024-01-01T00:00:00Z" diff --git a/server/test/follow_index_test.gleam b/server/test/follow_index_test.gleam index d574e0c..191f67f 100644 --- a/server/test/follow_index_test.gleam +++ b/server/test/follow_index_test.gleam @@ -26,6 +26,23 @@ pub fn following_lists_only_that_dids_subjects_test() { assert store.edges.count() == 3 } +pub fn followers_returns_authors_in_deterministic_order_test() { + let assert Ok(store) = follow_index.start() + store.edges.upsert(follow("did:z", "did:a", "f1")) + store.edges.upsert(follow("did:b", "did:a", "f2")) + store.edges.upsert(follow("did:c", "did:other", "f3")) + assert store.edges.followers("did:a") == ["did:b", "did:z"] +} + +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"] + store.edges.delete("at://did:a/dev.mokkenstorm.crate.graph.follow/f1") + assert store.edges.following("did:a") == ["did:b"] +} + pub fn following_is_empty_until_a_did_is_seen_test() { let assert Ok(store) = follow_index.start() assert store.seen.has_seen("did:a") == False @@ -51,6 +68,16 @@ pub fn deleting_a_follow_drops_only_that_edge_test() { assert store.edges.following("did:a") == ["did:c"] } +pub fn replacing_a_dids_edges_reconciles_stale_rows_atomically_test() { + let assert Ok(store) = follow_index.start() + 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.following("did:a") == ["did:new"] + assert store.edges.following("did:z") == ["did:old"] +} + 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 291c0f8..6a132c7 100644 --- a/server/test/graph_handler_test.gleam +++ b/server/test/graph_handler_test.gleam @@ -2,6 +2,7 @@ //// plus the pure `graph_follows.find` lookup. import atproto/xrpc +import atproto_core/xrpc as core_xrpc import crate/gen/graph/follow.{type GraphFollow, GraphFollow} import crate/storage.{type StoredItem, StoredItem} import crate_server/context.{type Context, Context} @@ -12,10 +13,14 @@ import crate_server/oauth/config import crate_server/oauth/session_store import crate_server/oauth/sessions import gleam/bit_array +import gleam/erlang/process import gleam/http import gleam/http/response +import gleam/int import gleam/json +import gleam/list import gleam/option.{None, Some} +import gleam/string import support import wisp import wisp/simulate @@ -96,17 +101,34 @@ 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 { "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_body) - "/xrpc/com.atproto.repo.createRecord" -> ok_response(create_body) + "/xrpc/com.atproto.repo.putRecord" -> ok_response(create_body) "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") _ -> panic as { "unexpected path: " <> req.path } } }) } +fn missing_scope_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> ok_response(list_records_body([])) + "/xrpc/com.atproto.repo.putRecord" -> + forbidden_response("{\"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 = @@ -149,6 +171,32 @@ pub fn follow_creates_a_record_and_returns_its_uri_test() { == 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( @@ -206,15 +254,49 @@ 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.createRecord" -> + "/xrpc/com.atproto.repo.putRecord" -> ok_response(create_record_response("3ccc")) _ -> panic as { "unexpected PDS call: " <> req.path } } }) } +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 } + } + }) +} + +pub fn follow_is_idempotent_and_chooses_a_stable_legacy_record_test() { + let #(ctx, cfg) = + test_context( + no_create_client( + list_records_body([ + follow_record_json("3bbb", target_did), + follow_record_json("3aaa", target_did), + ]), + ), + ) + let req = + authed_post( + support.xrpc("graph.followUser"), + json.object([#("subject", json.string(target_did))]), + cfg, + ) + 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( @@ -248,6 +330,32 @@ pub fn unfollow_finds_the_rkey_in_the_index_without_a_pds_list_test() { 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) @@ -260,3 +368,60 @@ pub fn unfollow_404s_from_the_index_when_the_edge_is_absent_test() { ) assert graph.unfollow(req, ctx).status == 404 } + +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 string.starts_with(first, "f_") + assert string.contains(first, ":") == False +} + +pub fn concurrent_follows_converge_on_the_same_put_rkey_test() { + let expected_rkey = graph_follows.rkey_for_subject(target_did) + let client = + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> + ok_response(list_records_body([])) + "/xrpc/com.atproto.repo.putRecord" -> { + case bit_array.to_string(req.body) { + Ok(body) -> + case string.contains(body, expected_rkey) { + True -> ok_response(create_record_response(expected_rkey)) + False -> + Error(core_xrpc.ConnectionFailed( + "non-deterministic follow rkey", + )) + } + _ -> + 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 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 assert Ok(a) = process.receive(replies, 1000) + let assert Ok(b) = process.receive(replies, 1000) + assert list.sort([a, b], int.compare) == [200, 200] +} diff --git a/server/test/oauth_test.gleam b/server/test/oauth_test.gleam index 4d34659..fc0dcfd 100644 --- a/server/test/oauth_test.gleam +++ b/server/test/oauth_test.gleam @@ -1,9 +1,11 @@ import atproto/xrpc +import crate/gen/graph/follow as graph_follow import crate_server/catalog_index import crate_server/context.{type Context, Context} import crate_server/handlers/oauth as oauth_handler import crate_server/oauth/assertion import crate_server/oauth/authed +import crate_server/oauth/client_metadata import crate_server/oauth/config import crate_server/oauth/dpop import crate_server/oauth/flow @@ -20,6 +22,7 @@ import gleam/erlang/process import gleam/http import gleam/http/request import gleam/http/response +import gleam/json import gleam/list import gleam/option.{None, Some} import gleam/result @@ -132,6 +135,15 @@ pub fn config_https_is_confidential_client_test() { assert cfg.redirect_uri == "https://app.example/api/oauth/callback" } +pub fn follow_scope_is_declared_in_client_metadata_test() { + let follow_scope = "repo:" <> graph_follow.collection + let scope = client_metadata.scope() + let body = client_metadata.document("https://app.example") |> json.to_string + + assert string.contains(scope, follow_scope) + assert support.field_string(body, ["scope"]) == Ok(scope) +} + // -- SSRF guard on the login path ----------------------------------------- // // Both the resolved PDS and the PAR endpoint the PDS's own discovery diff --git a/server/test/user_backfill_test.gleam b/server/test/user_backfill_test.gleam index c12c6eb..042f7bb 100644 --- a/server/test/user_backfill_test.gleam +++ b/server/test/user_backfill_test.gleam @@ -11,6 +11,7 @@ import crate/gen/catalog/release as catalog_release import crate/gen/graph/follow as graph_follow import crate/gen/shelf/entry as shelf_entry import crate_server/catalog_index +import crate_server/follow_index import crate_server/known_users.{KnownUser} import crate_server/shelf_index import crate_server/user_backfill @@ -395,6 +396,14 @@ fn follow_record(rkey: String, subject: String) -> json.Json { ) } +fn undecodable_follow_record(rkey: String) -> json.Json { + record( + "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/" <> rkey, + "bafy" <> rkey, + json.object([#("createdAt", json.string("2026-02-01T00:00:00Z"))]), + ) +} + fn follows_client(pages: fn(String) -> String) -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req) { @@ -409,6 +418,12 @@ pub fn backfill_seeds_follows_and_marks_the_did_seen_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() let follows = support.fresh_follow_index() + follows.edges.upsert(follow_index.Follow( + follow_uri: "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/stale", + did: "did:plc:me", + subject: "did:plc:stale", + created_at: "2026-01-01T00:00:00Z", + )) let client = follows_client(fn(_q) { page_body([follow_record("f1", "did:plc:b")], None) @@ -424,6 +439,12 @@ pub fn backfill_marks_seen_even_with_an_empty_follow_graph_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() let follows = support.fresh_follow_index() + follows.edges.upsert(follow_index.Follow( + follow_uri: "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/stale", + did: "did:plc:me", + subject: "did:plc:stale", + created_at: "2026-01-01T00:00:00Z", + )) let summary = user_backfill.backfill_user( follows_client(fn(_q) { empty_page_body() }), @@ -433,13 +454,51 @@ pub fn backfill_marks_seen_even_with_an_empty_follow_graph_test() { follows, ) assert summary.follows == 0 + assert follows.edges.following("did:plc:me") == [] assert follows.seen.has_seen("did:plc:me") == True } +pub fn malformed_follow_records_do_not_prune_edges_or_mark_seen_test() { + let assert Ok(store) = catalog_index.start() + let assert Ok(shelf) = shelf_index.start() + let follows = support.fresh_follow_index() + follows.edges.upsert(follow_index.Follow( + follow_uri: "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/stale", + did: "did:plc:me", + subject: "did:plc:stale", + created_at: "2026-01-01T00:00:00Z", + )) + let client = + follows_client(fn(_q) { + page_body( + [ + follow_record("f1", "did:plc:b"), + undecodable_follow_record("bad"), + ], + None, + ) + }) + let summary = + user_backfill.backfill_user(client, a_user(), store, shelf, follows) + assert summary.follows == 1 + assert follows.edges.following("did:plc:me") + == [ + "did:plc:b", + "did:plc:stale", + ] + assert follows.seen.has_seen("did:plc:me") == False +} + pub fn a_truncated_follow_run_keeps_rows_but_does_not_mark_seen_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() let follows = support.fresh_follow_index() + follows.edges.upsert(follow_index.Follow( + follow_uri: "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/stale", + did: "did:plc:me", + subject: "did:plc:stale", + created_at: "2026-01-01T00:00:00Z", + )) // First page has a cursor; the follow-up request fails, so the run is // partial and must not claim the graph is fully indexed. let client = @@ -458,6 +517,10 @@ pub fn a_truncated_follow_run_keeps_rows_but_does_not_mark_seen_test() { let summary = user_backfill.backfill_user(client, a_user(), store, shelf, follows) assert summary.follows == 1 - assert follows.edges.following("did:plc:me") == ["did:plc:b"] + assert follows.edges.following("did:plc:me") + == [ + "did:plc:b", + "did:plc:stale", + ] assert follows.seen.has_seen("did:plc:me") == False } diff --git a/web/css/15-misc.css b/web/css/15-misc.css index fc8d86c..d530020 100644 --- a/web/css/15-misc.css +++ b/web/css/15-misc.css @@ -49,6 +49,16 @@ flex: 1; min-width: 0; } +.notice__action { + flex-shrink: 0; + width: auto; + border: 1px solid currentColor; + background: transparent; + color: inherit; + cursor: pointer; + font: 700 11px/1 var(--mono); + padding: 6px 8px; +} .notice__dismiss { flex-shrink: 0; width: auto; @@ -88,4 +98,3 @@ font-size: 13px; color: var(--ink-muted); } - diff --git a/web/css/21-connections.css b/web/css/21-connections.css new file mode 100644 index 0000000..75aa9d0 --- /dev/null +++ b/web/css/21-connections.css @@ -0,0 +1,81 @@ +/* --- connections ------------------------------------------------------ */ + +.connections { + gap: 0; +} + +.connections-tabs { + display: flex; + gap: 22px; + padding: 18px 20px 0; + border-bottom: 1px solid var(--line); +} + +.connections-tab { + padding: 0 0 12px; + color: var(--ink-muted); + font: 700 11px/1.2 var(--mono); + letter-spacing: 0.7px; + text-decoration: none; + border-bottom: 3px solid transparent; +} + +.connections-tab.is-active { + color: var(--ink); + border-bottom-color: var(--mustard); +} + +.connections-list { + padding: 0 20px 24px; +} + +.connections-note { + margin: 16px 0 4px; + color: var(--ink-muted); + font: 400 11px/1.45 var(--mono); +} + +.connections-rows { + margin-top: 8px; +} + +.connection-row { + display: flex; + align-items: center; + justify-content: space-between; + gap: 12px; + padding: 15px 0; + border-bottom: 1px solid var(--line); +} + +.connection-row__identity { + display: flex; + align-items: center; + gap: 10px; + min-width: 0; +} + +.connection-row__handle { + min-width: 0; + overflow: hidden; + color: var(--ink); + font: 700 13px/1.3 var(--mono); + text-overflow: ellipsis; + text-decoration: none; +} + +.connection-row__follow { + flex-shrink: 0; +} + +.connections-more { + display: block; + margin: 18px auto 0; +} + +.connections-error { + display: grid; + justify-items: center; + gap: 12px; + padding: 28px 20px; +} diff --git a/web/src/crate_web.gleam b/web/src/crate_web.gleam index dee89f5..c62be99 100644 --- a/web/src/crate_web.gleam +++ b/web/src/crate_web.gleam @@ -69,6 +69,12 @@ fn init(_flags) -> #(Model, Effect(Msg)) { overlap: None, pressing: model.PressingLoading, feed: model.FeedLoading, + feed_generation: 0, + connections: model.ConnectionsLoading, + connections_generation: 0, + connection_write_generation: 0, + connection_pending: None, + connections_loading_more: False, feed_entries: dict.new(), feed_loading_more: False, nav_depth: 0, @@ -81,7 +87,12 @@ fn init(_flags) -> #(Model, Effect(Msg)) { let route_effects = case active_route { model.Add -> [effects.discogs_status()] model.Browse -> [effects.load_browse()] - model.Feed -> [effects.load_feed(None)] + model.Feed -> [ + effects.load_feed(model.FeedRequest(model.Feed, 0, None)), + ] + model.ConnectionsFollowing | model.ConnectionsFollowers -> [ + effects.load_connections(model.ConnectionsRequest(active_route, 0, None)), + ] model.EditInbox -> [effects.load_edit_inbox()] model.PublicCrate(handle) -> [effects.load_shelf(Some(handle), "current")] model.PublicRecord(handle, entry_id) -> [ diff --git a/web/src/crate_web/effects.gleam b/web/src/crate_web/effects.gleam index f8c8e2f..e29e30f 100644 --- a/web/src/crate_web/effects.gleam +++ b/web/src/crate_web/effects.gleam @@ -12,20 +12,23 @@ import crate/gen/shelf/list_entries import crate_web/appview import crate_web/browser import crate_web/model.{ - type AmendDraft, type Display, type Form, type ReleaseInfo, type Theme, + type AmendDraft, type ConnectionWrite, type ConnectionsRequest, type Display, + type FeedRequest, type Form, type ReleaseInfo, type Theme, Connection, DiscogsSearchPage, HandleSuggestion, ReleaseInfo, } import crate_web/money import crate_web/msg.{ - type ApiError, type EntryDetailData, type FeedSkeletonData, type Msg, - type ShelfData, AppliedProposal, BarcodeDetected, CameraUnsupported, - CoverUploaded, CrateOverlapData, EntryDetailData, FeedSkeletonData, GotAction, - GotActorShelf, GotAdd, GotAmend, GotApplyProposal, GotArtists, GotAvatar, - GotBrowse, GotBrowseAdd, GotBrowseMore, GotBrowseSearch, GotCrateOverlap, - GotDiscogs, GotDiscogsDisconnect, GotDiscogsImport, GotDiscogsStatus, - GotEditInbox, GotFeedMore, GotFeedSkeleton, GotFollow, GotHandleSuggestions, - GotLogout, GotPressing, GotScanResult, GotScanSeen, GotShelf, GotShelfMore, - GotUnfollow, LinkCopied, ScanLookup, ShelfData, + type ApiError, type ConnectionsData, type EntryDetailData, + type FeedSkeletonData, type Msg, type ShelfData, AppliedProposal, + BarcodeDetected, CameraUnsupported, ConnectionsData, CoverUploaded, + CrateOverlapData, EntryDetailData, FeedSkeletonData, GotAction, GotActorShelf, + GotAdd, GotAmend, GotApplyProposal, GotArtists, GotAvatar, GotBrowse, + GotBrowseAdd, GotBrowseMore, GotBrowseSearch, GotConnectionFollow, + GotConnectionUnfollow, GotConnections, GotCrateOverlap, GotDiscogs, + GotDiscogsDisconnect, GotDiscogsImport, GotDiscogsStatus, GotEditInbox, + GotFeedMore, GotFeedSkeleton, GotFollow, GotHandleSuggestions, GotLogout, + GotPressing, GotScanResult, GotScanSeen, GotShelf, GotShelfMore, GotUnfollow, + LinkCopied, ScanLookup, ShelfData, } import crate_web/prefs import gleam/dict @@ -743,20 +746,19 @@ fn release_info_decoder() -> decode.Decoder(ReleaseInfo) { }) } -/// One page of the network activity feed. Dispatches `GotFeedSkeleton` for a -/// fresh load (`cursor` is `None`) and `GotFeedMore` for an infinite-scroll -/// page (`cursor` is `Some`); the query shape is identical either way. -pub fn load_feed(cursor: Option(String)) -> Effect(Msg) { - let to_msg = case cursor { - None -> GotFeedSkeleton - Some(_) -> GotFeedMore +/// One page of the network activity feed. The request identity is carried +/// through the async completion so stale route visits or cursors are ignored. +pub fn load_feed(request: FeedRequest) -> Effect(Msg) { + let to_msg = case request.cursor { + None -> fn(result) { GotFeedSkeleton(request, result) } + Some(_) -> fn(result) { GotFeedMore(request, result) } } let url = xrpc( "feed.getFeedSkeleton", list.flatten([ [#("limit", int.to_string(model.feed_page_limit))], - opt_param("cursor", cursor), + opt_param("cursor", request.cursor), ]), ) bff_get(url, feed_skeleton_decoder(), to_msg) @@ -776,6 +778,78 @@ fn feed_skeleton_decoder() -> decode.Decoder(FeedSkeletonData) { decode.success(FeedSkeletonData(items:, fallback:, cursor:)) } +pub fn load_connections(request: ConnectionsRequest) -> Effect(Msg) { + let direction = case request.route { + model.ConnectionsFollowers -> "followers" + _ -> "following" + } + bff_get( + xrpc( + "graph.listConnections", + list.flatten([ + [#("direction", direction), #("limit", "50")], + opt_param("cursor", request.cursor), + ]), + ), + connections_decoder(), + fn(result) { GotConnections(request, result) }, + ) +} + +fn connections_decoder() -> decode.Decoder(ConnectionsData) { + let item_decoder = { + use did <- decode.field("did", decode.string) + use handle <- decode.field("handle", decode.string) + use viewer_follows <- decode.field("viewerFollows", decode.bool) + use follows_viewer <- decode.field("followsViewer", decode.bool) + use mutual <- decode.field("mutual", decode.bool) + decode.success(Connection( + did:, + handle:, + viewer_follows:, + follows_viewer:, + mutual:, + )) + } + use items <- decode.field("items", decode.list(item_decoder)) + use cursor <- decode.optional_field( + "cursor", + None, + decode.optional(decode.string), + ) + use complete <- decode.optional_field("complete", False, decode.bool) + use relationship_complete <- decode.optional_field( + "relationshipComplete", + False, + decode.bool, + ) + decode.success(ConnectionsData( + items:, + cursor:, + complete:, + relationship_complete:, + )) +} + +pub fn follow_connection(request: ConnectionWrite) -> Effect(Msg) { + let decoder = decode.field("uri", decode.string, decode.success) + bff_post( + xrpc("graph.followUser", []), + json.object([#("subject", json.string(request.did))]), + decoder, + fn(result) { GotConnectionFollow(request, result) }, + ) +} + +pub fn unfollow_connection(request: ConnectionWrite) -> Effect(Msg) { + bff_post( + xrpc("graph.unfollowUser", []), + json.object([#("subject", json.string(request.did))]), + nil_decoder(), + fn(result) { GotConnectionUnfollow(request, result) }, + ) +} + pub fn load_crate_overlap(handle: String) -> Effect(Msg) { let decoder = { use common_count <- decode.field("commonCount", decode.int) diff --git a/web/src/crate_web/model.gleam b/web/src/crate_web/model.gleam index 7df02b1..ca5c899 100644 --- a/web/src/crate_web/model.gleam +++ b/web/src/crate_web/model.gleam @@ -38,6 +38,8 @@ pub type Route { RecordAmend(entry_id: String) Browse Feed + ConnectionsFollowing + ConnectionsFollowers EditInbox /// One proposal's detail screen, off `EditInbox`. EditProposalDetail(id: String) @@ -50,6 +52,41 @@ pub type Route { PressingDetail(did: String, rkey: String) } +pub type Connection { + Connection( + did: String, + handle: String, + viewer_follows: Bool, + follows_viewer: Bool, + mutual: Bool, + ) +} + +pub type ConnectionsState { + ConnectionsLoading + ConnectionsLoaded( + items: List(Connection), + cursor: Option(String), + complete: Bool, + relationship_complete: Bool, + ) + ConnectionsFailed +} + +/// Identity of one connections list request. The generation is bumped for +/// every reload or page fetch, so a response from an earlier request cannot +/// replace a newer copy of the same route. +pub type ConnectionsRequest { + ConnectionsRequest(route: Route, generation: Int, cursor: Option(String)) +} + +/// Identity of one optimistic follow write. Keeping the route and generation +/// with the DID prevents a completion from an earlier visit from changing the +/// current visit's row. +pub type ConnectionWrite { + ConnectionWrite(route: Route, generation: Int, did: String, desired: Bool) +} + /// How many crate items render before the LOAD MORE button, and the amount /// each LOAD MORE reveals. Purely a render cap: the server may have fetched /// more than this already (see `crate_page_limit`), in which case pressing @@ -252,6 +289,13 @@ pub type FeedState { FeedFailed } +/// Identity of one feed skeleton or pagination request. The generation is +/// bumped when the feed route changes and when a page fetch starts, so an +/// earlier response cannot replace or append to the current feed. +pub type FeedRequest { + FeedRequest(route: Route, generation: Int, cursor: Option(String)) +} + /// The day-separator label for a feed item: the newest hydrated entry's /// `updatedAt` day (ISO dates sort lexicographically), or `None` while none /// of its refs are hydrated yet. @@ -316,6 +360,10 @@ pub type NoticeLevel { pub type Notice { Notice(level: NoticeLevel, message: String) + /// A permission failure with an explicit action. This is intentionally + /// distinct from a generic warning so it cannot silently reauthorize an + /// existing OAuth session. + ReauthorizeNotice(message: String) } pub type Form { @@ -774,6 +822,16 @@ 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). feed_entries: dict.Dict(#(String, String), EntryDetail), // True only while a cursor-driven feed page fetch is in flight. diff --git a/web/src/crate_web/msg.gleam b/web/src/crate_web/msg.gleam index 185999c..39ddb67 100644 --- a/web/src/crate_web/msg.gleam +++ b/web/src/crate_web/msg.gleam @@ -2,10 +2,11 @@ import atproto_core/xrpc import crate/gen/feed/get_feed_skeleton.{type FeedItem} import crate/gen/shelf/entry.{type ShelfEntry} import crate_web/model.{ - type ArtistHit, type BrowseRelease, type DiscogsResult, type DiscogsSearchPage, - type Display, type EditProposal, type Entry, type HandleSuggestion, - type ImportRun, type NetworkMatch, type ReleaseInfo, type Route, type ScanMode, - type Suggestion, type Theme, + type ArtistHit, type BrowseRelease, type Connection, type ConnectionWrite, + type ConnectionsRequest, type DiscogsResult, type DiscogsSearchPage, + type Display, type EditProposal, type Entry, type FeedRequest, + type HandleSuggestion, type ImportRun, type NetworkMatch, type ReleaseInfo, + type Route, type ScanMode, type Suggestion, type Theme, } import crate_web/photo_scan.{type DecodeError} import gleam/dict @@ -59,6 +60,15 @@ pub type FeedSkeletonData { ) } +pub type ConnectionsData { + ConnectionsData( + items: List(Connection), + cursor: Option(String), + complete: Bool, + relationship_complete: Bool, + ) +} + /// One page from `catalog.listReleases`, plus the opaque cursor for the next /// page when the filtered result set has more rows. pub type BrowseData { @@ -99,6 +109,7 @@ pub type Msg { GotHandleSuggestions(Result(List(HandleSuggestion), ApiError)) UseHandleSuggestion(HandleSuggestion) StartLogin + Reauthorize ArmLogout DisarmLogout Logout @@ -196,14 +207,29 @@ pub type Msg { GotApplyProposal(uri: String, result: Result(AppliedProposal, ApiError)) IgnoreProposal(uri: String) GotActorShelf(Result(ShelfData, ApiError)) - GotFeedSkeleton(Result(FeedSkeletonData, ApiError)) + GotFeedSkeleton( + request: FeedRequest, + result: Result(FeedSkeletonData, ApiError), + ) GotFeedEntry( actor: String, entry_id: String, result: Result(EntryDetailData, ApiError), ) FeedShowMore - GotFeedMore(Result(FeedSkeletonData, ApiError)) + GotFeedMore(request: FeedRequest, result: Result(FeedSkeletonData, ApiError)) + GotConnections( + request: ConnectionsRequest, + result: Result(ConnectionsData, ApiError), + ) + ConnectionsShowMore + RetryConnections + ToggleConnectionFollow(did: String, currently_following: Bool) + GotConnectionFollow( + request: ConnectionWrite, + result: Result(String, ApiError), + ) + GotConnectionUnfollow(request: ConnectionWrite, result: Result(Nil, ApiError)) GotCrateOverlap(Result(CrateOverlapData, ApiError)) /// Optimistically flips the public-crate page's follow state and fires /// the matching write; the target is the currently-open public crate's did. diff --git a/web/src/crate_web/pages/connections.gleam b/web/src/crate_web/pages/connections.gleam new file mode 100644 index 0000000..f02e758 --- /dev/null +++ b/web/src/crate_web/pages/connections.gleam @@ -0,0 +1,172 @@ +//// Following and followers for the signed-in viewer. Followers are an +//// appview index and may be incomplete, so an empty indexed result is never +//// presented as proof that nobody follows the viewer. + +import crate_web/model.{ + type Connection, type Model, ConnectionsFailed, ConnectionsLoaded, + ConnectionsLoading, +} +import crate_web/msg.{ + type Msg, ConnectionsShowMore, RetryConnections, ToggleConnectionFollow, +} +import crate_web/route +import crate_web/ui/controls as ctl +import crate_web/ui/states +import gleam/list +import gleam/option.{type Option, None, Some} +import lustre/attribute as attr +import lustre/element.{type Element, text} +import lustre/element/html +import lustre/event + +pub fn view(model: Model, followers: Bool) -> Element(Msg) { + html.div([attr.class("body page-scroll connections")], [ + tabs(followers), + case model.connections { + ConnectionsLoading -> loading() + ConnectionsFailed -> failed() + ConnectionsLoaded(items, cursor, complete, relationship_complete) -> + loaded(model, followers, items, cursor, complete, relationship_complete) + }, + ]) +} + +fn tabs(followers: Bool) -> Element(Msg) { + html.nav([attr.class("connections-tabs"), attr.aria_label("Connections")], [ + tab("FOLLOWING", route.to_path(model.ConnectionsFollowing), !followers), + tab("FOLLOWERS", route.to_path(model.ConnectionsFollowers), followers), + ]) +} + +fn tab(label: String, href: String, active: Bool) -> Element(Msg) { + html.a( + [ + attr.class(case active { + True -> "connections-tab is-active" + False -> "connections-tab" + }), + attr.href(href), + attr.attribute("aria-current", case active { + True -> "page" + False -> "false" + }), + ], + [text(label)], + ) +} + +fn loading() -> Element(Msg) { + html.p([attr.class("loading-status")], [text("◌ LOADING CONNECTIONS…")]) +} + +fn failed() -> Element(Msg) { + html.div([attr.class("connections-error")], [ + states.error_sticker("Couldn't load your connections right now."), + ctl.button("RETRY", ctl.Ghost, [event.on_click(RetryConnections)]), + ]) +} + +fn loaded( + model: Model, + followers: Bool, + items: List(Connection), + cursor: Option(String), + complete: Bool, + relationship_complete: Bool, +) -> Element(Msg) { + html.div([attr.class("connections-list")], [ + completeness_note(followers, items, complete), + case items { + [] -> empty(followers, complete) + _ -> + html.div( + [attr.class("connections-rows")], + list.map(items, row(model, relationship_complete, _)), + ) + }, + case cursor { + Some(_) -> + ctl.button("LOAD MORE", ctl.Ghost, [ + attr.class("connections-more"), + attr.disabled(model.connections_loading_more), + event.on_click(ConnectionsShowMore), + ]) + None -> element.none() + }, + ]) +} + +fn completeness_note( + followers: Bool, + items: List(Connection), + complete: Bool, +) -> Element(Msg) { + case followers && !complete { + True -> + html.p([attr.class("connections-note"), attr.role("status")], [ + text(case list.is_empty(items) { + True -> + "No indexed followers yet. The appview may still be catching up." + False -> "Followers may be incomplete while the appview catches up." + }), + ]) + False -> element.none() + } +} + +fn empty(followers: Bool, complete: Bool) -> Element(Msg) { + html.p([attr.class("empty")], [ + text(case followers { + True if !complete -> "No indexed followers yet." + True -> "No one follows you yet." + False -> "You are not following anyone yet." + }), + ]) +} + +fn row( + model: Model, + relationship_complete: Bool, + item: Connection, +) -> Element(Msg) { + let pending = case model.connection_pending { + Some(request) -> request.did == item.did && request.route == model.route + None -> False + } + html.div([attr.class("connection-row")], [ + html.div([attr.class("connection-row__identity")], [ + html.a( + [ + attr.class("connection-row__handle"), + attr.href(route.to_path(model.PublicCrate(item.handle))), + ], + [ + text("@" <> item.handle), + ], + ), + case relationship_complete && item.mutual { + True -> ctl.chip("MUTUAL", ctl.Accent) + False -> element.none() + }, + ]), + ctl.button( + case item.viewer_follows { + True -> "FOLLOWING" + False -> "FOLLOW" + }, + case item.viewer_follows { + True -> ctl.Ghost + False -> ctl.Primary + }, + [ + attr.class("connection-row__follow"), + attr.aria_label(case item.viewer_follows { + True -> "Unfollow @" <> item.handle + False -> "Follow @" <> item.handle + }), + attr.disabled(pending), + event.on_click(ToggleConnectionFollow(item.did, item.viewer_follows)), + ], + ), + ]) +} diff --git a/web/src/crate_web/pages/feed.gleam b/web/src/crate_web/pages/feed.gleam index 7e0cd52..b47dd47 100644 --- a/web/src/crate_web/pages/feed.gleam +++ b/web/src/crate_web/pages/feed.gleam @@ -23,6 +23,7 @@ import crate_web/ui/covers as cov import crate_web/ui/infinite_scroll as scroll import crate_web/ui/states import gleam/dict +import gleam/int import gleam/list import gleam/option.{type Option, None, Some} import lustre/attribute as attr @@ -134,16 +135,22 @@ fn single_row( } } -/// `batch_row`/`import_row`'s shared scaffold: a skeleton until the first -/// entry hydrates, else a strip row labeled by `label_for` off that entry. +/// `batch_row`/`import_row`'s shared scaffold: a skeleton until an entry +/// hydrates, else a strip row that counts only known entries. An entry that +/// hydrates out of order must not make its actor look responsible for the +/// entire unresolved batch. fn hydrated_strip_row( model: Model, item: FeedItem, - label_for: fn(EntryDetail) -> String, + label_for: fn(List(#(EntryRef, EntryDetail)), Int) -> String, ) -> Element(Msg) { - case hydrated_entries(model, item) { + let hydrated = hydrated_entries(model, item) + case hydrated { [] -> skeleton_row() - [#(_, fe), ..] -> strip_row(model, fe.handle, label_for(fe), item) + [#(_, fe), ..] -> { + let unresolved = list.length(item.entries) - list.length(hydrated) + strip_row(model, fe.handle, label_for(hydrated, unresolved), item) + } } } @@ -152,13 +159,26 @@ fn batch_row( item: FeedItem, reason: ReasonActorBatch, ) -> Element(Msg) { - hydrated_strip_row(model, item, fn(fe) { - "@" - <> fe.handle - <> " " - <> reason.action - <> " " - <> plural.count_noun(list.length(item.entries), "record") + hydrated_strip_row(model, item, fn(hydrated, unresolved) { + let known = list.length(hydrated) + let actors = hydrated_handles(hydrated) + let label = case unresolved { + 0 -> + format_handles(actors) + <> " " + <> reason.action + <> " " + <> plural.count_noun(known, "record") + _ -> + format_handles(actors) + <> " " + <> reason.action + <> " " + <> plural.count_noun(known, "record") + <> ", plus " + <> plural.count_noun(unresolved, "unresolved record") + } + label }) } @@ -167,16 +187,24 @@ fn import_row( item: FeedItem, reason: ReasonImport, ) -> Element(Msg) { - hydrated_strip_row(model, item, fn(fe) { + hydrated_strip_row(model, item, fn(hydrated, unresolved) { let via = case reason.source { Some(src) -> " via " <> provenance.provider_label(src) None -> "" } - "@" - <> fe.handle - <> " imported " - <> plural.count_noun(list.length(item.entries), "record") - <> via + let label = case unresolved { + 0 -> + format_handles(hydrated_handles(hydrated)) + <> " imported " + <> plural.count_noun(list.length(hydrated), "record") + _ -> + format_handles(hydrated_handles(hydrated)) + <> " imported " + <> plural.count_noun(list.length(hydrated), "record") + <> ", plus " + <> plural.count_noun(unresolved, "unresolved record") + } + label <> via }) } @@ -219,22 +247,20 @@ fn converge_row( item: FeedItem, reason: ReasonSubjectConverge, ) -> Element(Msg) { - case first_ref_hydrated(model, item) { - None -> skeleton_row() - Some(#(ref, fe)) -> { + case hydrated_entries(model, item) { + [] -> skeleton_row() + [#(ref, fe), ..] -> { let snap = model.entry_snapshot(fe.entry) let href = case model.split_release_uri(reason.subject.uri) { Ok(#(did, rkey)) -> route.to_path(PressingDetail(did, rkey)) Error(_) -> route.to_path(PublicRecord(fe.handle, ref.entry_id)) } + let label = converge_actors(model, item) <> " picked up the same pressing" html.div([attr.class("feed-row")], [ html.div([attr.class("feed-row__head")], [ avatar_cluster(model, item), - html.span([attr.class("feed-row__label")], [ - text( - plural.count_noun(list.length(item.entries), "collector") - <> " picked up the same pressing", - ), + html.span([attr.class("feed-row__label"), attr.aria_label(label)], [ + text(label), ]), ]), cov.cover_row( @@ -256,13 +282,77 @@ fn avatar_cluster(model: Model, item: FeedItem) -> Element(Msg) { |> list.take(3) |> list.map(fn(pair) { let #(_, fe) = pair - html.span([attr.class("avatar avatar--sm")], [ - bar.avatar_content(None, fe.handle), - ]) + html.span( + [attr.class("avatar avatar--sm"), attr.aria_label("@" <> fe.handle)], + [ + bar.avatar_content(None, fe.handle), + ], + ) }) html.div([attr.class("feed-avatar-cluster")], avatars) } +fn converge_actors(model: Model, item: FeedItem) -> String { + let hydrated = hydrated_entries(model, item) + let handles = + hydrated + |> list.map(fn(pair) { + let #(_, fe) = pair + "@" <> fe.handle + }) + |> list.unique + let known = format_handles(handles) + let unresolved = unresolved_actor_count(item, hydrated) + case unresolved { + 0 -> known + count if known == "" -> plural.count_noun(count, "unresolved collector") + count -> + known <> ", plus " <> plural.count_noun(count, "unresolved collector") + } +} + +fn format_handles(handles: List(String)) -> String { + case handles { + [] -> "" + [only] -> only + [first, second, ..rest] -> + first + <> " and " + <> second + <> case list.length(rest) { + 0 -> "" + count -> " and " <> int.to_string(count) <> " more" + } + } +} + +fn hydrated_handles(hydrated: List(#(EntryRef, EntryDetail))) -> List(String) { + hydrated + |> list.map(fn(pair) { + let #(_, fe) = pair + "@" <> fe.handle + }) + |> list.unique +} + +fn unresolved_actor_count( + item: FeedItem, + hydrated: List(#(EntryRef, EntryDetail)), +) -> Int { + let known_actors = + hydrated + |> list.map(fn(pair) { + let #(ref, _) = pair + ref.actor + }) + |> list.unique + item.entries + |> list.map(fn(ref) { ref.actor }) + |> list.unique + |> list.filter(fn(actor) { !list.contains(known_actors, actor) }) + |> list.length +} + fn row_head(seed: String, label: String) -> Element(Msg) { html.div([attr.class("feed-row__head")], [ html.span([attr.class("avatar avatar--sm")], [ @@ -320,8 +410,9 @@ fn hydrated_entries( }) } -/// The designed empty state: unchanged from the pre-feed stub, shown whenever -/// the server returns no items at all (regardless of the fallback flag). +/// Suggested collectors need a read model containing known users plus follow +/// state. Connections owns that contract, so this remains browse-only until +/// that appview response is available rather than inventing local suggestions. fn empty_state() -> Element(Msg) { html.div([attr.class("body page-scroll")], [ html.div([attr.class("empty-state")], [ diff --git a/web/src/crate_web/pages/settings.gleam b/web/src/crate_web/pages/settings.gleam index 4104525..554ce2a 100644 --- a/web/src/crate_web/pages/settings.gleam +++ b/web/src/crate_web/pages/settings.gleam @@ -2,7 +2,9 @@ //// here from Add), a link into the edit inbox with its pending count, and log //// out. Reached by tapping the app-bar avatar. -import crate_web/model.{type Model, type Theme, EditInbox, LoggedIn} +import crate_web/model.{ + type Model, type Theme, ConnectionsFollowing, EditInbox, LoggedIn, +} import crate_web/msg.{ type Msg, ArmLogout, DiscogsConnect, DiscogsDisconnect, DiscogsImport, DiscogsImportWantlist, Logout, SetTheme, @@ -23,12 +25,30 @@ pub fn view(model: Model) -> Element(Msg) { html.div([attr.class("body page-scroll")], [ account_header(model), theme_row(model.theme), + connections_row(), discogs_account_view(model), inbox_row(model), logout_button(model.confirm_logout), ]) } +fn connections_row() -> Element(Msg) { + html.a( + [attr.class("settings-row"), attr.href(route.to_path(ConnectionsFollowing))], + [ + html.span([attr.class("settings-row__label")], [text("CONNECTIONS")]), + html.span([attr.class("settings-row__spacer")], []), + html.span( + [ + attr.class("settings-row__chevron"), + attr.attribute("aria-hidden", "true"), + ], + [text("›")], + ), + ], + ) +} + fn account_header(model: Model) -> Element(Msg) { let handle = current_handle(model) html.div([attr.class("settings-account")], [ diff --git a/web/src/crate_web/route.gleam b/web/src/crate_web/route.gleam index ee65241..530f243 100644 --- a/web/src/crate_web/route.gleam +++ b/web/src/crate_web/route.gleam @@ -1,9 +1,9 @@ //// URL <-> Route mapping for modem. import crate_web/model.{ - type Route, Add, Browse, Crate, EditInbox, EditProposalDetail, Feed, - PressingDetail, PublicCrate, PublicRecord, Record, RecordAmend, Scan, ScanDone, - ScanReview, Settings, + type Route, Add, Browse, ConnectionsFollowers, ConnectionsFollowing, Crate, + EditInbox, EditProposalDetail, Feed, PressingDetail, PublicCrate, PublicRecord, + Record, RecordAmend, Scan, ScanDone, ScanReview, Settings, } import gleam/uri.{type Uri} @@ -15,6 +15,8 @@ pub fn parse(target: Uri) -> Route { ["scan", "done"] -> ScanDone ["browse"] -> Browse ["feed"] -> Feed + ["connections", "followers"] -> ConnectionsFollowers + ["connections"] -> ConnectionsFollowing ["inbox"] -> EditInbox ["inbox", id] -> EditProposalDetail(id) ["settings"] -> Settings @@ -36,6 +38,8 @@ pub fn to_path(route: Route) -> String { ScanDone -> "/scan/done" Browse -> "/browse" Feed -> "/feed" + ConnectionsFollowing -> "/connections" + ConnectionsFollowers -> "/connections/followers" EditInbox -> "/inbox" EditProposalDetail(id) -> "/inbox/" <> id Settings -> "/settings" @@ -57,6 +61,7 @@ pub fn section(route: Route) -> Route { Crate Browse | PressingDetail(_, _) -> Browse Feed -> Feed + ConnectionsFollowing | ConnectionsFollowers -> Settings Settings | EditInbox | EditProposalDetail(_) -> Settings PublicCrate(_) | PublicRecord(_, _) -> route } diff --git a/web/src/crate_web/update.gleam b/web/src/crate_web/update.gleam index 2f63f04..95a36be 100644 --- a/web/src/crate_web/update.gleam +++ b/web/src/crate_web/update.gleam @@ -3,6 +3,7 @@ import crate_web/msg.{type Msg} import crate_web/update/add import crate_web/update/amend import crate_web/update/browse +import crate_web/update/connections import crate_web/update/feed import crate_web/update/inbox import crate_web/update/routing @@ -49,6 +50,7 @@ pub fn update(model: Model, msg: Msg) -> #(Model, Effect(Msg)) { msg.UseHandleSuggestion(suggestion) -> session.use_handle_suggestion(model, suggestion) msg.StartLogin -> session.start_login(model) + msg.Reauthorize -> session.reauthorize(model) msg.ArmLogout -> session.arm_logout(model) msg.DisarmLogout -> session.disarm_logout(model) msg.Logout -> session.logout_msg(model) @@ -136,10 +138,22 @@ pub fn update(model: Model, msg: Msg) -> #(Model, Effect(Msg)) { inbox.got_apply_proposal(model, uri, result) msg.IgnoreProposal(uri) -> inbox.ignore_proposal(model, uri) - msg.GotFeedSkeleton(result) -> feed.got_feed_skeleton(model, result) + msg.GotFeedSkeleton(request, result) -> + feed.got_feed_skeleton(model, request, result) msg.GotFeedEntry(actor, entry_id, result) -> feed.got_feed_entry(model, actor, entry_id, result) msg.FeedShowMore -> feed.feed_show_more(model) - msg.GotFeedMore(result) -> feed.got_feed_more(model, result) + msg.GotFeedMore(request, result) -> + feed.got_feed_more(model, request, result) + msg.GotConnections(request, result) -> + connections.got_connections(model, request, result) + msg.ConnectionsShowMore -> connections.show_more(model) + msg.RetryConnections -> connections.retry(model) + msg.ToggleConnectionFollow(did, currently_following) -> + connections.toggle_follow(model, did, currently_following) + msg.GotConnectionFollow(request, result) -> + connections.got_follow(model, request, result) + msg.GotConnectionUnfollow(request, result) -> + connections.got_unfollow(model, request, result) } } diff --git a/web/src/crate_web/update/browse.gleam b/web/src/crate_web/update/browse.gleam index 29c9303..5867fdb 100644 --- a/web/src/crate_web/update/browse.gleam +++ b/web/src/crate_web/update/browse.gleam @@ -9,7 +9,7 @@ import crate_web/model.{ OwnCrate, PressingFailed, PressingLoaded, ShelfLoaded, } import crate_web/msg.{type Msg} -import crate_web/update/common.{failed, write_error} +import crate_web/update/common.{failed, follow_write_error, write_error} import gleam/list import gleam/option.{type Option, None, Some} import gleam/string @@ -320,7 +320,7 @@ pub fn got_follow( Error(e) -> case model.overlap { Some(overlap) -> - write_error( + follow_write_error( Model( ..model, overlap: Some( @@ -330,7 +330,7 @@ pub fn got_follow( e, "Could not follow that user.", ) - None -> write_error(model, e, "Could not follow that user.") + None -> follow_write_error(model, e, "Could not follow that user.") } } } @@ -345,7 +345,7 @@ pub fn got_unfollow( Error(e) -> case model.overlap { Some(overlap) -> - write_error( + follow_write_error( Model( ..model, overlap: Some(CrateOverlap(..overlap, viewer_follows: True)), @@ -353,7 +353,7 @@ pub fn got_unfollow( e, "Could not unfollow that user.", ) - None -> write_error(model, e, "Could not unfollow that user.") + None -> follow_write_error(model, e, "Could not unfollow that user.") } } } diff --git a/web/src/crate_web/update/common.gleam b/web/src/crate_web/update/common.gleam index 0f72c5b..dd0fa64 100644 --- a/web/src/crate_web/update/common.gleam +++ b/web/src/crate_web/update/common.gleam @@ -20,6 +20,10 @@ pub fn failed(message: String) -> Option(Notice) { Some(Notice(Failure, message)) } +pub fn reauthorize(message: String) -> Option(Notice) { + Some(model.ReauthorizeNotice(message)) +} + /// The shared 401-vs-everything-else handling for a failed write: a session /// expiry signs the visitor out and clears their shelf, anything else just /// surfaces `fallback` as a failure notice. @@ -45,6 +49,28 @@ pub fn write_error( } } +/// Follow writes use a distinct 403 response when an existing OAuth session +/// lacks the graph.follow scope. Other writes keep the generic error path. +pub fn follow_write_error( + model: Model, + error: msg.ApiError, + fallback: String, +) -> #(Model, Effect(Msg)) { + case error { + xrpc.BadStatus(status: 403, ..) -> #( + Model( + ..model, + busy: False, + notice: reauthorize( + "This session cannot follow collectors yet. Sign in again to approve the follow permission.", + ), + ), + effect.none(), + ) + _ -> write_error(model, error, fallback) + } +} + /// The one hydrated-entry payload the record pages and the feed cache share, /// built from the `getEntry` response. pub fn detail_of(data: msg.EntryDetailData) -> EntryDetail { diff --git a/web/src/crate_web/update/connections.gleam b/web/src/crate_web/update/connections.gleam new file mode 100644 index 0000000..c7f65e7 --- /dev/null +++ b/web/src/crate_web/update/connections.gleam @@ -0,0 +1,247 @@ +import crate_web/effects.{ + follow_connection, load_connections, unfollow_connection, +} +import crate_web/model.{ + type ConnectionWrite, type ConnectionsRequest, type ConnectionsState, + type Model, Connection, ConnectionWrite, ConnectionsFailed, + ConnectionsFollowers, ConnectionsFollowing, ConnectionsLoaded, + ConnectionsLoading, ConnectionsRequest, Model, +} +import crate_web/msg.{type ConnectionsData, type Msg} +import crate_web/update/common.{failed, follow_write_error} +import gleam/list +import gleam/option.{type Option, None, Some} +import lustre/effect.{type Effect} + +pub fn got_connections( + model: Model, + request: ConnectionsRequest, + result: Result(ConnectionsData, msg.ApiError), +) -> #(Model, Effect(Msg)) { + case + model.route == request.route + && model.connections_generation == request.generation + { + False -> #(model, effect.none()) + True -> + case result { + Ok(data) -> { + let items = case model.connections_loading_more { + True -> + case model.connections { + ConnectionsLoaded(existing, _, _, _) -> + list.append( + existing, + model.append_new(existing, data.items, fn(item) { + Some(item.did) + }), + ) + _ -> data.items + } + False -> data.items + } + #( + Model( + ..model, + connections: ConnectionsLoaded( + items:, + cursor: data.cursor, + complete: data.complete, + relationship_complete: data.relationship_complete, + ), + connections_loading_more: False, + ), + effect.none(), + ) + } + Error(_) -> #( + Model( + ..model, + connections: case model.connections_loading_more { + True -> model.connections + False -> ConnectionsFailed + }, + connections_loading_more: False, + notice: failed("Couldn't load your connections."), + ), + effect.none(), + ) + } + } +} + +pub fn retry(model: Model) -> #(Model, Effect(Msg)) { + case connections_route(model.route) { + Some(route) -> { + let generation = model.connections_generation + 1 + let write_generation = model.connection_write_generation + 1 + let request = ConnectionsRequest(route, generation, None) + #( + Model( + ..model, + connections: ConnectionsLoading, + connections_generation: generation, + connection_write_generation: write_generation, + connection_pending: None, + connections_loading_more: False, + ), + load_connections(request), + ) + } + None -> #(model, effect.none()) + } +} + +pub fn show_more(model: Model) -> #(Model, Effect(Msg)) { + case connections_route(model.route), model.connections { + Some(route), ConnectionsLoaded(_, Some(cursor), _, _) + if !model.connections_loading_more + -> { + let generation = model.connections_generation + 1 + let request = ConnectionsRequest(route, generation, Some(cursor)) + #( + Model( + ..model, + connections_generation: generation, + connections_loading_more: True, + ), + load_connections(request), + ) + } + _, _ -> #(model, effect.none()) + } +} + +pub fn toggle_follow( + model: Model, + did: String, + currently_following: Bool, +) -> #(Model, Effect(Msg)) { + case + model.connection_pending, + connections_route(model.route), + update_item(model.connections, did, !currently_following) + { + None, Some(route), Some(updated) -> { + let generation = model.connection_write_generation + 1 + let request = + ConnectionWrite( + route:, + generation:, + did:, + desired: !currently_following, + ) + #( + Model( + ..model, + connections: updated, + connection_write_generation: generation, + connection_pending: Some(request), + notice: None, + ), + case currently_following { + True -> unfollow_connection(request) + False -> follow_connection(request) + }, + ) + } + _, _, _ -> #(model, effect.none()) + } +} + +pub fn got_follow( + model: Model, + request: ConnectionWrite, + result: Result(String, msg.ApiError), +) -> #(Model, Effect(Msg)) { + case + model.connection_pending == Some(request) && model.route == request.route, + result + { + False, _ -> #(model, effect.none()) + True, Ok(_) -> #(Model(..model, connection_pending: None), effect.none()) + True, Error(error) -> + revert( + model, + request.did, + request.desired, + error, + "Could not follow that user.", + ) + } +} + +pub fn got_unfollow( + model: Model, + request: ConnectionWrite, + result: Result(Nil, msg.ApiError), +) -> #(Model, Effect(Msg)) { + case + model.connection_pending == Some(request) && model.route == request.route, + result + { + False, _ -> #(model, effect.none()) + True, Ok(Nil) -> #(Model(..model, connection_pending: None), effect.none()) + True, Error(error) -> + revert( + model, + request.did, + request.desired, + error, + "Could not unfollow that user.", + ) + } +} + +fn revert( + model: Model, + did: String, + desired: Bool, + error: msg.ApiError, + message: String, +) -> #(Model, Effect(Msg)) { + follow_write_error( + Model( + ..model, + connections: update_item(model.connections, did, !desired) + |> option.unwrap(model.connections), + connection_pending: None, + ), + error, + message, + ) +} + +fn update_item( + state: ConnectionsState, + did: String, + follows: Bool, +) -> Option(ConnectionsState) { + case state { + ConnectionsLoaded(items, cursor, complete, relationship_complete) -> + Some(ConnectionsLoaded( + items: list.map(items, fn(item) { + case item.did == did { + True -> + Connection( + ..item, + viewer_follows: follows, + mutual: follows && item.follows_viewer, + ) + False -> item + } + }), + cursor:, + complete:, + relationship_complete:, + )) + _ -> None + } +} + +fn connections_route(route: model.Route) -> Option(model.Route) { + case route { + ConnectionsFollowing | ConnectionsFollowers -> Some(route) + _ -> None + } +} diff --git a/web/src/crate_web/update/feed.gleam b/web/src/crate_web/update/feed.gleam index e14375c..d9b2313 100644 --- a/web/src/crate_web/update/feed.gleam +++ b/web/src/crate_web/update/feed.gleam @@ -1,23 +1,31 @@ import crate/gen/feed/get_feed_skeleton.{type FeedItem} import crate_web/effects.{load_entry, load_feed} -import crate_web/model.{type Model, FeedFailed, FeedLoaded, Model} +import crate_web/model.{ + type FeedRequest, type Model, Feed, FeedFailed, FeedLoaded, FeedRequest, Model, +} import crate_web/msg.{type FeedSkeletonData, type Msg, GotFeedEntry} import crate_web/update/common.{detail_of, failed} import gleam/dict import gleam/list -import gleam/option.{Some} +import gleam/option.{None, Some} import lustre/effect.{type Effect} pub fn got_feed_skeleton( model: Model, + request: FeedRequest, result: Result(FeedSkeletonData, msg.ApiError), ) -> #(Model, Effect(Msg)) { - case result { - Ok(data) -> #( - Model(..model, feed: FeedLoaded(data.items, data.cursor, data.fallback)), + case current_skeleton_request(model, request), result { + False, _ -> #(model, effect.none()) + True, Ok(data) -> #( + Model( + ..model, + feed: FeedLoaded(data.items, data.cursor, data.fallback), + feed_loading_more: False, + ), hydrate_uncached(data.items, model.feed_entries), ) - Error(_) -> #( + True, Error(_) -> #( Model( ..model, feed: FeedFailed, @@ -53,21 +61,27 @@ pub fn got_feed_entry( } pub fn feed_show_more(model: Model) -> #(Model, Effect(Msg)) { - case model.feed { - FeedLoaded(_, Some(cursor), _) if !model.feed_loading_more -> #( - Model(..model, feed_loading_more: True), - load_feed(Some(cursor)), - ) - _ -> #(model, effect.none()) + case model.route, model.feed { + Feed, FeedLoaded(_, Some(cursor), _) if !model.feed_loading_more -> { + let generation = model.feed_generation + 1 + let request = FeedRequest(Feed, generation, Some(cursor)) + #( + Model(..model, feed_generation: generation, feed_loading_more: True), + load_feed(request), + ) + } + _, _ -> #(model, effect.none()) } } pub fn got_feed_more( model: Model, + request: FeedRequest, result: Result(FeedSkeletonData, msg.ApiError), ) -> #(Model, Effect(Msg)) { - case result { - Ok(data) -> + case current_more_request(model, request), result { + False, _ -> #(model, effect.none()) + True, Ok(data) -> case model.feed { FeedLoaded(items, _, _) -> { let fresh = model.append_new(items, data.items, model.feed_item_key) @@ -86,7 +100,7 @@ pub fn got_feed_more( } _ -> #(Model(..model, feed_loading_more: False), effect.none()) } - Error(_) -> #( + True, Error(_) -> #( Model( ..model, feed_loading_more: False, @@ -97,6 +111,26 @@ pub fn got_feed_more( } } +fn current_skeleton_request(model: Model, request: FeedRequest) -> Bool { + model.route == request.route + && request.route == Feed + && model.feed_generation == request.generation + && request.cursor == None +} + +fn current_more_request(model: Model, request: FeedRequest) -> Bool { + model.route == request.route + && request.route == Feed + && model.feed_generation == request.generation + && request.cursor != None + && model.feed_loading_more + && case model.feed, request.cursor { + FeedLoaded(_, cursor, _), Some(request_cursor) -> + cursor == Some(request_cursor) + _, _ -> False + } +} + /// Fan out hydration for every entry ref in a batch of feed items not /// already sitting in the cache, so a re-fetched or overlapping item never /// re-hydrates a ref that's already loaded. diff --git a/web/src/crate_web/update/routing.gleam b/web/src/crate_web/update/routing.gleam index 7c51f95..c077c45 100644 --- a/web/src/crate_web/update/routing.gleam +++ b/web/src/crate_web/update/routing.gleam @@ -2,8 +2,9 @@ import crate_web/effects.{ discogs_status, load_browse, load_entry, load_feed, load_shelf, } import crate_web/model.{ - type Model, Add, Browse, EditInbox, EditProposalDetail, Feed, FeedLoading, - Model, PressingDetail, PressingLoading, PublicCrate, PublicRecord, Record, + type Model, Add, Browse, ConnectionsFollowers, ConnectionsFollowing, + ConnectionsLoading, EditInbox, EditProposalDetail, Feed, FeedLoading, Model, + PressingDetail, PressingLoading, PublicCrate, PublicRecord, Record, RecordAmend, Scan, ScanDone, crate_of, set_crate_loading_if_absent, } import crate_web/msg.{type Msg, GotEntry} @@ -125,6 +126,9 @@ pub fn on_route_change( // only on entry (never at login) and clears any stale in-flight add. route -> { let scan_state = scan.scan_for_route(model.scan, model.route, route) + let entering_connections = is_connections_route(route) + let leaving_connections = is_connections_route(model.route) + let feed_route_changed = route == Feed || model.route == Feed let updated = Model( ..model, @@ -164,6 +168,27 @@ pub fn on_route_change( Feed -> FeedLoading _ -> model.feed }, + feed_generation: case feed_route_changed { + True -> model.feed_generation + 1 + False -> model.feed_generation + }, + feed_loading_more: False, + connections: case route { + ConnectionsFollowing | ConnectionsFollowers -> ConnectionsLoading + _ -> model.connections + }, + connections_generation: case entering_connections { + True -> model.connections_generation + 1 + False -> model.connections_generation + }, + connection_write_generation: case + entering_connections || leaving_connections + { + True -> model.connection_write_generation + 1 + False -> model.connection_write_generation + }, + connection_pending: None, + connections_loading_more: False, ) // A cached actor crate stays put (stale-while-revalidate): the fetch // below still fires and refreshes it once it lands, but a fresh @@ -186,7 +211,14 @@ pub fn on_route_change( effects.scan_seen(), ]) Browse -> load_browse() - Feed -> load_feed(None) + Feed -> + load_feed(model.FeedRequest(Feed, updated.feed_generation, None)) + ConnectionsFollowing | ConnectionsFollowers -> + effects.load_connections(model.ConnectionsRequest( + route, + updated.connections_generation, + None, + )) EditInbox | EditProposalDetail(_) -> effects.load_edit_inbox() PublicCrate(handle) -> effect.batch([ @@ -203,6 +235,13 @@ pub fn on_route_change( } } +fn is_connections_route(route: model.Route) -> Bool { + case route { + ConnectionsFollowing | ConnectionsFollowers -> True + _ -> False + } +} + /// Both armed confirms reset on every route change; drop their outside-click /// watchers too so a stale one never lingers on a button that just left the /// page. diff --git a/web/src/crate_web/update/session.gleam b/web/src/crate_web/update/session.gleam index a466a95..27415fd 100644 --- a/web/src/crate_web/update/session.gleam +++ b/web/src/crate_web/update/session.gleam @@ -57,6 +57,20 @@ pub fn start_login(model: Model) -> #(Model, Effect(Msg)) { } } +/// Start a fresh OAuth flow with the current handle. Existing sessions cannot +/// be upgraded in place because the authorization server must grant the new +/// repository scope explicitly. +pub fn reauthorize(model: Model) -> #(Model, Effect(Msg)) { + case model.auth { + model.LoggedIn(handle) -> + case string.trim(handle) { + "" -> #(model, effect.none()) + _ -> #(Model(..model, notice: None), oauth_login(handle)) + } + _ -> #(model, effect.none()) + } +} + pub fn arm_logout(model: Model) -> #(Model, Effect(Msg)) { #( Model(..model, confirm_logout: True), diff --git a/web/src/crate_web/view.gleam b/web/src/crate_web/view.gleam index 5890d2f..88c7061 100644 --- a/web/src/crate_web/view.gleam +++ b/web/src/crate_web/view.gleam @@ -2,15 +2,17 @@ //// modules. import crate_web/model.{ - type Model, type Notice, type NoticeLevel, Add, Browse, Crate, EditInbox, - EditProposalDetail, EntryDetailFailed, EntryDetailLoaded, Failure, Feed, Info, - LoggedIn, LoggedOut, Notice, PressingDetail, PublicCrate, PublicRecord, Record, + type Model, type Notice, type NoticeLevel, Add, Browse, ConnectionsFollowers, + ConnectionsFollowing, Crate, EditInbox, EditProposalDetail, EntryDetailFailed, + EntryDetailLoaded, Failure, Feed, Info, LoggedIn, LoggedOut, Notice, + PressingDetail, PublicCrate, PublicRecord, ReauthorizeNotice, Record, RecordAmend, Scan, ScanDone, ScanReview, Settings, Success, Warning, inbox_pending_count, } -import crate_web/msg.{type Msg, Back, ClearNotice} +import crate_web/msg.{type Msg, Back, ClearNotice, Reauthorize} import crate_web/pages/add import crate_web/pages/browse +import crate_web/pages/connections import crate_web/pages/crate import crate_web/pages/edit_inbox import crate_web/pages/edit_proposal @@ -118,6 +120,10 @@ fn page(model: Model) -> Element(Msg) { Crate -> crate.view(model) Browse -> browse.view(model) Feed -> feed.view(model) + ConnectionsFollowing -> + sub_page("CONNECTIONS", element.none(), connections.view(model, False)) + ConnectionsFollowers -> + sub_page("CONNECTIONS", element.none(), connections.view(model, True)) Settings -> settings.view(model) Add -> sub_page("ADD RECORD", element.none(), add.view(model)) Scan -> @@ -232,30 +238,54 @@ fn bottom_bar(model: Model) -> Element(Msg) { fn notice_view(notice: Option(Notice)) -> Element(Msg) { case notice { - Some(Notice(level, message)) -> + Some(ReauthorizeNotice(message)) -> html.div( [ - attr.class("notice notice--" <> level_class(level)), + attr.class("notice notice--warning"), attr.role("alert"), ], [ - html.span([attr.class("notice__icon")], [text(level_icon(level))]), + html.span([attr.class("notice__icon")], [text("!")]), html.span([attr.class("notice__msg")], [text(message)]), html.button( [ - attr.class("notice__dismiss"), + attr.class("notice__action"), attr.type_("button"), - attr.attribute("aria-label", "Dismiss"), - event.on_click(ClearNotice), + event.on_click(Reauthorize), ], - [text("✕")], + [text("REAUTHORIZE")], ), + dismiss_notice(), + ], + ) + Some(Notice(level, message)) -> + html.div( + [ + attr.class("notice notice--" <> level_class(level)), + attr.role("alert"), + ], + [ + html.span([attr.class("notice__icon")], [text(level_icon(level))]), + html.span([attr.class("notice__msg")], [text(message)]), + dismiss_notice(), ], ) None -> element.none() } } +fn dismiss_notice() -> Element(Msg) { + html.button( + [ + attr.class("notice__dismiss"), + attr.type_("button"), + attr.attribute("aria-label", "Dismiss"), + event.on_click(ClearNotice), + ], + [text("✕")], + ) +} + fn level_class(level: NoticeLevel) -> String { case level { Success -> "success" diff --git a/web/test/connections_test.gleam b/web/test/connections_test.gleam new file mode 100644 index 0000000..e0ab1ba --- /dev/null +++ b/web/test/connections_test.gleam @@ -0,0 +1,279 @@ +import atproto_core/xrpc +import crate_web/model.{ + type Connection, type Model, Connection, ConnectionsFailed, + ConnectionsFollowers, ConnectionsFollowing, ConnectionsLoaded, + ConnectionsLoading, ConnectionsRequest, Model, ReauthorizeNotice, +} +import crate_web/msg.{ + ConnectionsData, GotConnectionFollow, GotConnections, OnRouteChange, + RetryConnections, ToggleConnectionFollow, +} +import crate_web/pages/connections +import crate_web/route +import crate_web/update.{update} +import gleam/option.{None, Some} +import gleam/string +import gleam/uri +import lustre/element +import support.{empty_effect, logged_in, network_error} + +pub fn connections_routes_round_trip_test() { + assert route.to_path(ConnectionsFollowers) == "/connections/followers" + let assert Ok(target) = uri.parse("/connections/followers") + assert route.parse(target) == ConnectionsFollowers + let assert Ok(default) = uri.parse("/connections") + assert route.parse(default) == model.ConnectionsFollowing +} + +pub fn followers_empty_state_explains_index_incompleteness_test() { + let html = + connections.view( + model.Model( + ..logged_in(), + connections: ConnectionsLoaded([], None, False, True), + ), + True, + ) + |> element.to_string + assert string.contains(html, "No indexed followers yet") + assert string.contains(html, "appview may still be catching up") + assert !string.contains(html, "No one follows you yet") +} + +pub fn connection_row_exposes_mutual_and_accessible_follow_action_test() { + let html = + connections.view( + model.Model( + ..logged_in(), + connections: ConnectionsLoaded( + [ + Connection( + did: "did:plc:bob", + handle: "bob.test", + viewer_follows: False, + follows_viewer: True, + mutual: False, + ), + ], + None, + False, + True, + ), + ), + True, + ) + |> element.to_string + assert string.contains(html, "@bob.test") + assert string.contains(html, "Follow @bob.test") + assert string.contains(html, "FOLLOW") + assert string.contains(html, "aria-current=\"page\"") +} + +pub fn connection_loading_and_failure_states_are_distinct_test() { + let loading = + connections.view( + model.Model(..logged_in(), connections: ConnectionsLoading), + False, + ) + |> element.to_string + let failed = + connections.view( + model.Model(..logged_in(), connections: ConnectionsFailed), + False, + ) + |> element.to_string + assert string.contains(loading, "LOADING CONNECTIONS") + assert string.contains(failed, "RETRY") + assert !string.contains(loading, "RETRY") +} + +pub fn stale_same_route_list_response_is_ignored_after_retry_test() { + let old = ConnectionsRequest(ConnectionsFollowing, 4, None) + let current = ConnectionsRequest(ConnectionsFollowing, 5, None) + let seeded = + Model( + ..logged_in(), + route: ConnectionsFollowing, + connections_generation: old.generation, + connections: ConnectionsLoaded([connection("old")], None, True, True), + ) + let #(reloading, _) = update(seeded, RetryConnections) + assert reloading.connections_generation == current.generation + let #(after_stale, _) = + update( + reloading, + GotConnections( + old, + Ok(ConnectionsData([connection("stale")], None, True, True)), + ), + ) + assert after_stale.connections == ConnectionsLoading + let #(loaded, _) = + update( + after_stale, + GotConnections( + current, + Ok(ConnectionsData([connection("fresh")], None, True, True)), + ), + ) + assert loaded.connections + == ConnectionsLoaded([connection("fresh")], None, True, True) +} + +pub fn pagination_response_appends_only_for_current_generation_test() { + let old = ConnectionsRequest(ConnectionsFollowing, 8, Some("cursor-1")) + let current = ConnectionsRequest(ConnectionsFollowing, 9, Some("cursor-1")) + let seeded = + Model( + ..logged_in(), + route: ConnectionsFollowing, + connections_generation: current.generation, + connections: ConnectionsLoaded( + [connection("first")], + Some("cursor-1"), + False, + True, + ), + connections_loading_more: True, + ) + let #(after_stale, _) = + update( + seeded, + GotConnections( + old, + Ok(ConnectionsData([connection("stale")], None, True, True)), + ), + ) + assert after_stale.connections == seeded.connections + let #(loaded, _) = + update( + after_stale, + GotConnections( + current, + Ok(ConnectionsData([connection("second")], None, True, True)), + ), + ) + assert loaded.connections + == ConnectionsLoaded( + [connection("first"), connection("second")], + None, + True, + True, + ) + assert loaded.connections_loading_more == False +} + +pub fn optimistic_follow_success_clears_pending_write_test() { + let seeded = loaded_follow_state(False) + let #(optimistic, _) = + update(seeded, ToggleConnectionFollow("did:bob", False)) + let assert Some(request) = optimistic.connection_pending + assert request.did == "did:bob" + assert request.desired == True + assert connection_follows(optimistic) == True + let #(confirmed, _) = + update(optimistic, GotConnectionFollow(request, Ok("at://alice/follow/1"))) + assert confirmed.connection_pending == None + assert connection_follows(confirmed) == True +} + +pub fn optimistic_follow_error_rolls_back_test() { + let seeded = loaded_follow_state(False) + let #(optimistic, _) = + update(seeded, ToggleConnectionFollow("did:bob", False)) + let assert Some(request) = optimistic.connection_pending + let #(rolled_back, _) = + update(optimistic, GotConnectionFollow(request, Error(network_error()))) + assert rolled_back.connection_pending == None + assert connection_follows(rolled_back) == False + assert rolled_back.notice != None +} + +pub fn stale_follow_completion_after_same_route_reload_is_ignored_test() { + let seeded = loaded_follow_state(False) + let #(optimistic, _) = + update(seeded, ToggleConnectionFollow("did:bob", False)) + let assert Some(request) = optimistic.connection_pending + let #(reloaded, _) = update(optimistic, OnRouteChange(ConnectionsFollowing)) + assert reloaded.connection_pending == None + let #(after_stale, _) = + update(reloaded, GotConnectionFollow(request, Error(network_error()))) + assert after_stale.connections == ConnectionsLoading + assert after_stale.connection_pending == None +} + +fn connection(handle: String) -> Connection { + Connection( + did: "did:" <> handle, + handle: handle <> ".test", + viewer_follows: False, + follows_viewer: False, + mutual: False, + ) +} + +fn loaded_follow_state(follows: Bool) -> Model { + Model( + ..logged_in(), + route: ConnectionsFollowing, + connections: ConnectionsLoaded( + [ + Connection( + did: "did:bob", + handle: "bob.test", + viewer_follows: follows, + follows_viewer: False, + mutual: False, + ), + ], + None, + True, + True, + ), + ) +} + +fn connection_follows(model: Model) -> Bool { + case model.connections { + ConnectionsLoaded([first, ..], _, _, _) -> first.viewer_follows + _ -> False + } +} + +pub fn missing_follow_permission_offers_reauthorization_in_connections_test() { + let did = "did:plc:bob" + let model = + model.Model( + ..logged_in(), + route: ConnectionsFollowing, + connections: ConnectionsLoaded( + [ + Connection( + did:, + handle: "bob.test", + viewer_follows: False, + follows_viewer: False, + mutual: False, + ), + ], + None, + True, + True, + ), + ) + let #(pending, effect) = update(model, ToggleConnectionFollow(did, False)) + assert effect != empty_effect() + let assert Some(request) = pending.connection_pending + let #(updated, _) = + update( + pending, + GotConnectionFollow( + request, + Error(xrpc.BadStatus(403, Some("InsufficientScope"), None, "")), + ), + ) + assert updated.notice + == Some(ReauthorizeNotice( + "This session cannot follow collectors yet. Sign in again to approve the follow permission.", + )) +} diff --git a/web/test/feed_test.gleam b/web/test/feed_test.gleam index 1493c73..e46cc5b 100644 --- a/web/test/feed_test.gleam +++ b/web/test/feed_test.gleam @@ -7,8 +7,8 @@ import crate/gen/feed/get_feed_skeleton.{ } import crate/gen/shelf/list_entries import crate_web/model.{ - type EntryDetail, type Model, EntryDetail, Feed, FeedFailed, FeedLoaded, - FeedLoading, Model, + type EntryDetail, type FeedRequest, type Model, Browse, EntryDetail, Feed, + FeedFailed, FeedLoaded, FeedLoading, FeedRequest, Model, } import crate_web/msg.{ EntryDetailData, FeedShowMore, FeedSkeletonData, GotFeedEntry, GotFeedMore, @@ -93,6 +93,14 @@ fn render(model: Model) -> String { feed.view(model) |> element.to_string } +fn feed_model() -> Model { + Model(..logged_in(), route: Feed) +} + +fn feed_request(cursor: option.Option(String)) -> FeedRequest { + FeedRequest(Feed, 0, cursor) +} + fn occurrences(haystack: String, needle: String) -> Int { list.length(string.split(haystack, needle)) - 1 } @@ -124,7 +132,10 @@ pub fn feed_failed_state_shows_the_couldnt_load_copy_test() { pub fn got_feed_skeleton_error_marks_failed_and_notices_test() { let #(model, _) = - update(logged_in(), GotFeedSkeleton(Error(support.network_error()))) + update( + feed_model(), + GotFeedSkeleton(feed_request(None), Error(support.network_error())), + ) assert model.feed == FeedFailed assert case model.notice { Some(model.Notice(model.Failure, _)) -> True @@ -161,6 +172,21 @@ pub fn actor_batch_hydrated_renders_count_and_strip_test() { assert string.contains(html, "class=\"feed-strip\"") } +pub fn actor_batch_partial_hydration_does_not_credit_later_actor_with_full_batch_test() { + let refs = [ref("did:a", "e0"), ref("did:b", "e1")] + let html = + render( + with_feed([batch(refs, "added")], False, [ + #(#("did:b", "e1"), fe("bob.test", "2026-07-18T11:00:00Z")), + ]), + ) + assert string.contains( + html, + "@bob.test added 1 record, plus 1 unresolved record", + ) + assert !string.contains(html, "@bob.test added 2 records") +} + pub fn subject_converge_hydrated_links_to_the_pressing_test() { let refs = [ref("did:a", "e0"), ref("did:b", "e9")] let subject = "at://did:x/dev.mokkenstorm.crate.catalog.release/rk1" @@ -168,9 +194,32 @@ pub fn subject_converge_hydrated_links_to_the_pressing_test() { render( with_feed([converge(refs, subject)], False, [ #(#("did:a", "e0"), fe("alice.test", "2026-07-18T10:00:00Z")), + #(#("did:b", "e9"), fe("bob.test", "2026-07-18T11:00:00Z")), ]), ) - assert string.contains(html, "2 collectors picked up the same pressing") + assert string.contains( + html, + "@alice.test and @bob.test picked up the same pressing", + ) + assert string.contains(html, "class=\"feed-avatar-cluster\"") + assert string.contains(html, "aria-label=\"@alice.test\"") + assert string.contains(html, "aria-label=\"@bob.test\"") + assert string.contains(html, "href=\"/pressing/did:x/rk1\"") +} + +pub fn subject_converge_partial_hydration_names_known_and_unresolved_test() { + let refs = [ref("did:a", "e0"), ref("did:b", "e9"), ref("did:c", "e8")] + let subject = "at://did:x/dev.mokkenstorm.crate.catalog.release/rk1" + let html = + render( + with_feed([converge(refs, subject)], False, [ + #(#("did:b", "e9"), fe("bob.test", "2026-07-18T11:00:00Z")), + ]), + ) + let label = + "@bob.test, plus 2 unresolved collectors picked up the same pressing" + assert string.contains(html, label) + assert string.contains(html, "aria-label=\"" <> label <> "\"") assert string.contains(html, "class=\"feed-avatar-cluster\"") assert string.contains(html, "href=\"/pressing/did:x/rk1\"") } @@ -191,6 +240,21 @@ pub fn import_reason_via_line_reflects_the_source_test() { assert !string.contains(without_source, " via ") } +pub fn import_partial_hydration_uses_unresolved_copy_test() { + let refs = [ref("did:a", "e0"), ref("did:b", "e1")] + let html = + render( + with_feed([import_item(refs, Some("discogs"))], False, [ + #(#("did:b", "e1"), fe("bob.test", "2026-07-18T11:00:00Z")), + ]), + ) + assert string.contains( + html, + "@bob.test imported 1 record, plus 1 unresolved record via Discogs", + ) + assert !string.contains(html, "@bob.test imported 2 records via Discogs") +} + // --- reason rows: unhydrated (skeleton, no crash) ----------------------- pub fn every_reason_renders_a_skeleton_when_unhydrated_test() { @@ -256,7 +320,8 @@ pub fn on_route_change_feed_sets_loading_and_fires_a_load_test() { pub fn got_feed_skeleton_stores_items_and_hydrates_test() { let data = FeedSkeletonData([single("did:a", "e0", "added")], False, None) - let #(model, effect) = update(logged_in(), GotFeedSkeleton(Ok(data))) + let #(model, effect) = + update(feed_model(), GotFeedSkeleton(feed_request(None), Ok(data))) let assert FeedLoaded(items, _, _) = model.feed assert list.length(items) == 1 // Uncached refs mean a real hydration effect fires. @@ -265,7 +330,8 @@ pub fn got_feed_skeleton_stores_items_and_hydrates_test() { pub fn got_feed_skeleton_with_no_items_fires_no_hydration_test() { let data = FeedSkeletonData([], False, None) - let #(_, effect) = update(logged_in(), GotFeedSkeleton(Ok(data))) + let #(_, effect) = + update(feed_model(), GotFeedSkeleton(feed_request(None), Ok(data))) assert effect == empty_effect() } @@ -295,7 +361,7 @@ pub fn got_feed_entry_error_leaves_the_cache_untouched_test() { pub fn got_feed_more_appends_and_dedupes_by_first_ref_test() { let seeded = Model( - ..logged_in(), + ..feed_model(), feed: FeedLoaded( [single("did:a", "e0", "added"), single("did:a", "e1", "added")], Some("cursor1"), @@ -310,7 +376,8 @@ pub fn got_feed_more_appends_and_dedupes_by_first_ref_test() { False, None, ) - let #(updated, _) = update(seeded, GotFeedMore(Ok(next))) + let #(updated, _) = + update(seeded, GotFeedMore(feed_request(Some("cursor1")), Ok(next))) let assert FeedLoaded(items, _, _) = updated.feed assert item_keys(items) == [#("did:a", "e0"), #("did:a", "e1"), #("did:a", "e2")] @@ -333,13 +400,64 @@ pub fn feed_show_more_fetches_only_with_a_cursor_and_not_already_loading_test() |> list.each(fn(row) { let #(feed_state, loading, expect_fetch) = row let seeded = - Model(..logged_in(), feed: feed_state, feed_loading_more: loading) + Model(..feed_model(), feed: feed_state, feed_loading_more: loading) let #(_, effect) = update(seeded, FeedShowMore) // A fetch fires only when a cursor exists and none is already in flight. assert { effect != empty_effect() } == expect_fetch }) } +pub fn stale_feed_skeleton_does_not_replace_the_current_route_test() { + let current = + Model( + ..feed_model(), + feed_generation: 2, + feed: FeedLoaded([single("did:a", "e0", "current")], None, False), + ) + let stale = FeedSkeletonData([single("did:b", "e1", "stale")], False, None) + let #(updated, effect) = + update(current, GotFeedSkeleton(FeedRequest(Feed, 1, None), Ok(stale))) + assert updated == current + assert effect == empty_effect() +} + +pub fn stale_feed_more_does_not_append_or_clear_loading_state_test() { + let current = + Model( + ..feed_model(), + feed_generation: 2, + feed: FeedLoaded( + [single("did:a", "e0", "current")], + Some("cursor2"), + False, + ), + feed_loading_more: True, + ) + let stale = FeedSkeletonData([single("did:b", "e1", "stale")], False, None) + let #(updated, effect) = + update( + current, + GotFeedMore(FeedRequest(Feed, 1, Some("cursor1")), Ok(stale)), + ) + assert updated == current + assert effect == empty_effect() +} + +pub fn route_transition_clears_feed_loading_state_test() { + let current = + Model( + ..feed_model(), + feed_loading_more: True, + feed: FeedLoaded( + [single("did:a", "e0", "current")], + Some("cursor"), + False, + ), + ) + let #(updated, _) = update(current, OnRouteChange(Browse)) + assert updated.feed_loading_more == False +} + // --- pure helpers -------------------------------------------------------- pub fn feed_item_day_takes_the_newest_hydrated_day_test() { diff --git a/web/test/login_test.gleam b/web/test/login_test.gleam index 8029c25..cdf66e2 100644 --- a/web/test/login_test.gleam +++ b/web/test/login_test.gleam @@ -1,5 +1,7 @@ -import crate_web/model.{LoggedOut, OwnCrate, ShelfLoading, crate_of} -import crate_web/msg.{GotAvatar, GotLogout, Logout} +import crate_web/model.{ + LoggedOut, Model, OwnCrate, ReauthorizeNotice, ShelfLoading, crate_of, +} +import crate_web/msg.{GotAvatar, GotLogout, Logout, Reauthorize} import crate_web/pages/login import crate_web/update.{update} import gleam/option.{None, Some} @@ -73,3 +75,11 @@ pub fn got_avatar_stores_it_and_logout_clears_it_test() { let #(after_logout, _) = update(with_avatar, Logout) assert after_logout.avatar == None } + +pub fn reauthorize_uses_the_current_handle_and_clears_the_notice_test() { + let model = + Model(..logged_in(), notice: Some(ReauthorizeNotice("needs scope"))) + let #(updated, effect) = update(model, Reauthorize) + assert updated.notice == None + assert effect != support.empty_effect() +} diff --git a/web/test/public_crate_test.gleam b/web/test/public_crate_test.gleam index 41a895b..a5dc133 100644 --- a/web/test/public_crate_test.gleam +++ b/web/test/public_crate_test.gleam @@ -1,9 +1,11 @@ //// Public-route tests: `/u/:handle` and `/u/:handle/record/:entryId`. +import atproto_core/xrpc import crate_web/model.{ type Model, Actor, ActorCrate, CrateOverlap, EntryDetailFailed, EntryDetailLoaded, EntryDetailLoading, LoggedIn, LoggedOut, Model, PublicCrate, - PublicRecord, ShelfFailed, ShelfLoaded, ShelfLoading, crate_of, set_crate, + PublicRecord, ReauthorizeNotice, ShelfFailed, ShelfLoaded, ShelfLoading, + crate_of, set_crate, } import crate_web/msg.{ CopyRecordLink, CrateOverlapData, EntryDetailData, GotActorShelf, @@ -367,6 +369,20 @@ pub fn got_follow_error_reverts_the_optimistic_flip_test() { assert updated.notice != None } +pub fn missing_follow_permission_offers_reauthorization_test() { + let model = model_with_overlap(CrateOverlap(0, None, 0, True, None)) + let #(updated, _) = + update( + model, + GotFollow(Error(xrpc.BadStatus(403, Some("InsufficientScope"), None, ""))), + ) + assert updated.overlap == Some(CrateOverlap(0, None, 0, False, None)) + assert updated.notice + == Some(ReauthorizeNotice( + "This session cannot follow collectors yet. Sign in again to approve the follow permission.", + )) +} + pub fn got_unfollow_ok_leaves_the_optimistic_clear_in_place_test() { let model = model_with_overlap(CrateOverlap(0, None, 0, False, None)) let #(updated, _) = update(model, GotUnfollow(Ok(Nil))) diff --git a/web/test/support.gleam b/web/test/support.gleam index 65cb35a..6accb64 100644 --- a/web/test/support.gleam +++ b/web/test/support.gleam @@ -78,6 +78,12 @@ pub fn base() -> Model { overlap: None, pressing: model.PressingLoading, feed: model.FeedLoading, + feed_generation: 0, + connections: model.ConnectionsLoading, + connections_generation: 0, + connection_write_generation: 0, + connection_pending: None, + connections_loading_more: False, feed_entries: dict.new(), feed_loading_more: False, nav_depth: 0,