diff --git a/server/src/crate_server/db.gleam b/server/src/crate_server/db.gleam index 06615bb..05ab191 100644 --- a/server/src/crate_server/db.gleam +++ b/server/src/crate_server/db.gleam @@ -1,7 +1,8 @@ //// The query shapes every Postgres store backend was spelling out by hand. -//// All of them degrade instead of failing the caller: these stores back -//// derived, rebuildable projections, so a failed read is a stale page and a -//// failed write is a row the next reconcile pass restores. +//// Most degrade instead of failing the caller: these stores back derived, +//// rebuildable projections, so a failed read is a stale page and a failed +//// write is a row the next reconcile pass restores. `try_all` is the +//// exception, for reads where empty and unreadable are not the same answer. import gleam/dynamic/decode.{type Decoder} import gleam/option.{type Option, None, Some} @@ -35,6 +36,18 @@ pub fn all(query: pog.Query(a), conn: pog.Connection) -> List(a) { } } +/// `all` for a read whose caller must tell "no rows" apart from "no answer", +/// a paged read being the usual case. +pub fn try_all( + query: pog.Query(a), + conn: pog.Connection, +) -> Result(List(a), Nil) { + case pog.execute(query, conn) { + Ok(pog.Returned(rows:, ..)) -> Ok(rows) + Error(_) -> Error(Nil) + } +} + pub fn first(query: pog.Query(a), conn: pog.Connection) -> Option(a) { case pog.execute(query, conn) { Ok(pog.Returned(rows: [row, ..], ..)) -> Some(row) diff --git a/server/src/crate_server/follow_index.gleam b/server/src/crate_server/follow_index.gleam index ea71727..636d4c0 100644 --- a/server/src/crate_server/follow_index.gleam +++ b/server/src/crate_server/follow_index.gleam @@ -34,11 +34,16 @@ pub type EdgeOps { /// The dids `did` follows. Empty is ambiguous alone; pair with /// `SeenOps.has_seen`. following: fn(String) -> List(String), - following_page: fn(String, Option(String), Int) -> List(String), + /// One page of `following`. Errors rather than degrading to empty: a + /// paged reader reads an empty page as the end of the list, so a failed + /// read served as `[]` is indistinguishable from an authoritative answer. + following_page: fn(String, Option(String), Int) -> + Result(List(String), String), /// The dids that follow `did`. This reverse lookup is appview-derived and /// may be incomplete until every known actor has been backfilled. followers: fn(String) -> List(String), - followers_page: fn(String, Option(String), Int) -> List(String), + followers_page: fn(String, Option(String), Int) -> + Result(List(String), String), following_among: fn(String, List(String)) -> List(String), followers_among: fn(String, List(String)) -> List(String), /// The unfollow path's lookup, served without a live `listRecords`. @@ -68,6 +73,8 @@ pub type Store { Store(edges: EdgeOps, seen: SeenOps) } +const unavailable = "follow index unavailable" + type State { State(follows: Dict(String, Follow), seen: Set(String)) } @@ -104,11 +111,13 @@ pub fn start() -> Result(Store, actor.StartError) { delete_for_did: fn(did) { process.send(subject, DeleteForDid(did)) }, following: fn(did) { parallel.ask(subject, Following(did, _), or: []) }, following_page: fn(did, after, limit) { - parallel.ask(subject, FollowingPage(did, after, limit, _), or: []) + parallel.try_ask(subject, FollowingPage(did, after, limit, _)) + |> result.replace_error(unavailable) }, followers: fn(did) { parallel.ask(subject, Followers(did, _), or: []) }, followers_page: fn(did, after, limit) { - parallel.ask(subject, FollowersPage(did, after, limit, _), or: []) + parallel.try_ask(subject, FollowersPage(did, after, limit, _)) + |> result.replace_error(unavailable) }, following_among: fn(did, candidates) { parallel.ask(subject, FollowingAmong(did, candidates, _), or: []) @@ -126,7 +135,7 @@ pub fn start() -> Result(Store, actor.StartError) { parallel.ask( subject, ReplaceForDid(did, follows, _), - or: Error("follow index unavailable"), + or: Error(unavailable), ) }, count: fn() { parallel.ask(subject, Count, or: 0) }, diff --git a/server/src/crate_server/follow_index_postgres.gleam b/server/src/crate_server/follow_index_postgres.gleam index f742b70..16454b5 100644 --- a/server/src/crate_server/follow_index_postgres.gleam +++ b/server/src/crate_server/follow_index_postgres.gleam @@ -167,7 +167,7 @@ fn following_page( did: String, after: Option(String), limit: Int, -) -> List(String) { +) -> Result(List(String), String) { connection_page(conn, "did", "subject", did, after, limit) } @@ -259,7 +259,7 @@ fn followers_page( subject: String, after: Option(String), limit: Int, -) -> List(String) { +) -> Result(List(String), String) { connection_page(conn, "subject", "did", subject, after, limit) } @@ -270,7 +270,7 @@ fn connection_page( scope: String, after: Option(String), limit: Int, -) -> List(String) { +) -> Result(List(String), String) { let row = { use value <- decode.field("value", decode.string) decode.success(value) @@ -284,7 +284,8 @@ fn connection_page( |> pog.parameter(pog.nullable(pog.text, after)) |> pog.parameter(pog.int(limit)) |> pog.returning(row) - |> db.all(conn) + |> db.try_all(conn) + |> result.replace_error("follow_edges page read failed") } fn following_among( diff --git a/server/src/crate_server/handlers/connections.gleam b/server/src/crate_server/handlers/connections.gleam index cc2af31..1741bfd 100644 --- a/server/src/crate_server/handlers/connections.gleam +++ b/server/src/crate_server/handlers/connections.gleam @@ -121,7 +121,7 @@ fn page_dids( direction: Direction, after: Option(String), limit: Int, -) -> List(String) { +) -> Result(List(String), String) { case direction { Following -> ctx.follow_index.edges.following_page(viewer_did, after, limit) Followers -> ctx.follow_index.edges.followers_page(viewer_did, after, limit) @@ -160,7 +160,15 @@ fn respond_page( cursor |> option.then(pagination.decode_cursor(namespace, _)) |> option.map(fn(point: pagination.ResumePoint) { point.last_id }) - let rows = page_dids(ctx, session_did, direction, after, limit + 1) + // An unreadable index must not reach the client as an empty final page: it + // would read as an authoritative end of the list. + use rows <- unless_unreadable(page_dids( + ctx, + session_did, + direction, + after, + limit + 1, + )) let page_dids = list.take(rows, limit) let related = case direction, following_complete { Following, _ -> @@ -211,6 +219,19 @@ fn respond_page( |> wisp.json_response(200) } +fn unless_unreadable( + read: Result(a, String), + continue: fn(a) -> Response, +) -> Response { + case read { + Ok(rows) -> continue(rows) + Error(message) -> { + wisp.log_warning("connections: " <> message) + error_json(503, "could not read your connections") + } + } +} + fn cursor_namespace(direction: Direction, viewer_did: String) -> String { let prefix = case direction { Following -> "connections-following-" diff --git a/server/src/crate_server/parallel.gleam b/server/src/crate_server/parallel.gleam index 675b886..e8d813c 100644 --- a/server/src/crate_server/parallel.gleam +++ b/server/src/crate_server/parallel.gleam @@ -39,7 +39,16 @@ pub fn ask( make_request: fn(Subject(reply)) -> msg, or default: reply, ) -> reply { - try_call(subject, ask_timeout, make_request) |> result.unwrap(default) + try_ask(subject, make_request) |> result.unwrap(default) +} + +/// `ask` for a read whose caller must tell "nothing there" apart from "could +/// not look". +pub fn try_ask( + subject: Subject(msg), + make_request: fn(Subject(reply)) -> msg, +) -> Result(reply, Nil) { + try_call(subject, ask_timeout, make_request) } type Event(b) { diff --git a/server/test/connections_handler_test.gleam b/server/test/connections_handler_test.gleam index 72113f8..d7779c4 100644 --- a/server/test/connections_handler_test.gleam +++ b/server/test/connections_handler_test.gleam @@ -264,3 +264,33 @@ pub fn a_connection_cursor_paginates_and_is_viewer_scoped_test() { assert support.field_nested(body, ["items"], ["did"]) == Ok(expected) }) } + +pub fn an_unreadable_connection_page_is_a_503_not_an_empty_final_page_test() { + let base = support.fresh_follow_index() + seed_followers(base, "did:plc:me", ["did:plc:a"]) + base.seen.mark_seen("did:plc:me") + let wedged = + follow_index.Store( + ..base, + edges: follow_index.EdgeOps( + ..base.edges, + following_page: fn(_did, _after, _limit) { Error("read failed") }, + followers_page: fn(_did, _after, _limit) { Error("read failed") }, + ), + ) + + ["following", "followers"] + |> list.each(fn(direction) { + let resp = + respond( + wedged, + support.unreachable_client(), + auth_session(), + "graph.listConnections?direction=" <> direction, + ) + assert resp.status == 503 + let body = simulate.read_body(resp) + assert support.field_present(body, ["items"]) == False + assert support.field_present(body, ["cursor"]) == False + }) +} diff --git a/server/test/follow_index_postgres_test.gleam b/server/test/follow_index_postgres_test.gleam index 4eca9db..6952379 100644 --- a/server/test/follow_index_postgres_test.gleam +++ b/server/test/follow_index_postgres_test.gleam @@ -44,10 +44,10 @@ pub fn follow_round_trip_test() { == ["did:plc:subj-a", "did:plc:subj-b"] store.edges.upsert(follow("did:plc:other", did, "fr3")) assert store.edges.followers(did) == ["did:plc:other"] - assert store.edges.following_page(did, None, 1) == ["did:plc:subj-a"] + assert store.edges.following_page(did, None, 1) == Ok(["did:plc:subj-a"]) assert store.edges.following_page(did, Some("did:plc:subj-a"), 1) - == ["did:plc:subj-b"] - assert store.edges.followers_page(did, None, 1) == ["did:plc:other"] + == Ok(["did:plc:subj-b"]) + assert store.edges.followers_page(did, None, 1) == Ok(["did:plc:other"]) assert store.edges.following_among(did, ["did:plc:subj-b"]) == ["did:plc:subj-b"] assert store.edges.followers_among(did, ["did:plc:other"]) @@ -82,7 +82,7 @@ pub fn connection_pages_resume_past_a_deleted_anchor_test() { store.edges.upsert(follow(follower, did, "fa4")) assert store.edges.following_page(did, Some("did:plc:anchor-a"), 2) - == ["did:plc:anchor-b", "did:plc:anchor-c"] + == Ok(["did:plc:anchor-b", "did:plc:anchor-c"]) // The anchor is unfollowed between pages: the next page must resume past it // rather than restart from the top or collapse to nothing. @@ -90,10 +90,10 @@ pub fn connection_pages_resume_past_a_deleted_anchor_test() { "at://" <> did <> "/dev.mokkenstorm.crate.graph.follow/fa2", ) assert store.edges.following_page(did, Some("did:plc:anchor-b"), 2) - == ["did:plc:anchor-c"] - assert store.edges.following_page(did, Some("did:plc:zzz"), 2) == [] + == Ok(["did:plc:anchor-c"]) + assert store.edges.following_page(did, Some("did:plc:zzz"), 2) == Ok([]) assert store.edges.followers_page(did, Some("did:plc:absent"), 2) - == [follower] + == Ok([follower]) store.edges.delete_for_did(did) store.edges.delete_for_did(follower) diff --git a/server/test/follow_index_test.gleam b/server/test/follow_index_test.gleam index 14f6c41..bd72181 100644 --- a/server/test/follow_index_test.gleam +++ b/server/test/follow_index_test.gleam @@ -90,15 +90,15 @@ pub fn connection_pages_are_distinct_bounded_and_resume_after_a_subject_test() { store.edges.upsert(follow("did:z", "did:a", "f5")) store.edges.upsert(follow("did:y", "did:a", "f6")) - assert store.edges.following_page("did:a", None, 2) == ["did:b", "did:c"] - assert store.edges.following_page("did:a", Some("did:c"), 2) == ["did:d"] - assert store.edges.followers_page("did:a", None, 1) == ["did:y"] + assert store.edges.following_page("did:a", None, 2) == Ok(["did:b", "did:c"]) + assert store.edges.following_page("did:a", Some("did:c"), 2) == Ok(["did:d"]) + assert store.edges.followers_page("did:a", None, 1) == Ok(["did:y"]) // A cursor whose anchor was unfollowed between pages resumes at the next // surviving subject rather than restarting, so page 2 never re-serves page 1. store.edges.delete("at://did:a/dev.mokkenstorm.crate.graph.follow/f3") - assert store.edges.following_page("did:a", Some("did:c"), 2) == ["did:d"] - assert store.edges.following_page("did:a", Some("did:zzz"), 2) == [] + assert store.edges.following_page("did:a", Some("did:c"), 2) == Ok(["did:d"]) + assert store.edges.following_page("did:a", Some("did:zzz"), 2) == Ok([]) } pub fn relationship_intersections_only_return_matching_candidates_test() {