diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index 9546a3c..9e0a3cc 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -1,10 +1,9 @@ -//// Prototype in-memory catalog index: an actor over a Dict of -//// `browse.BrowseRow`, keyed by uri. Seeded by a one-time backfill over -//// `known_users` at boot (`main`), kept live by `jetstream_consumer`. Lost on -//// restart by design: rebuildable from the same backfill, and the dataset is -//// small enough at alpha scale to hold entirely in memory (no Postgres). +//// Prototype in-memory catalog index: an actor over a Dict of `BrowseRow` +//// keyed by uri, seeded by a boot backfill and kept live by `jetstream_consumer`. +//// Lost on restart by design (rebuildable, and small enough to hold in memory). import at_record_server/browse.{type BrowseRow} +import at_record_server/parallel import gleam/dict.{type Dict} import gleam/erlang/process.{type Subject} import gleam/otp/actor @@ -39,7 +38,7 @@ pub fn start() -> Result(Store, actor.StartError) { Nil }, list: fn() { - case try_call(subject, 1000, List) { + case parallel.try_call(subject, 1000, List) { Ok(rows) -> rows Error(Nil) -> [] } @@ -47,27 +46,6 @@ pub fn start() -> Result(Store, actor.StartError) { ) } -// process.call/actor.call panic on timeout or a dead callee; this mirrors -// known_users_memory's try_call so a read stays best-effort. -fn try_call( - subject: Subject(msg), - timeout: Int, - make_request: fn(Subject(reply)) -> msg, -) -> Result(reply, Nil) { - use callee <- result.try(process.subject_owner(subject)) - let reply_subject = process.new_subject() - let monitor = process.monitor(callee) - process.send(subject, make_request(reply_subject)) - let reply = - process.new_selector() - |> process.select_map(reply_subject, Ok) - |> process.select_specific_monitor(monitor, fn(_down) { Error(Nil) }) - |> process.selector_receive(timeout) - |> result.unwrap(Error(Nil)) - process.demonitor_process(monitor) - reply -} - fn handle( state: Dict(String, BrowseRow), msg: Msg, diff --git a/server/src/at_record_server/known_users_memory.gleam b/server/src/at_record_server/known_users_memory.gleam index 63cc7f6..db573bb 100644 --- a/server/src/at_record_server/known_users_memory.gleam +++ b/server/src/at_record_server/known_users_memory.gleam @@ -1,8 +1,8 @@ //// In-memory known_users.Store backend: an actor over a Dict keyed by did. -//// Lost on restart (use known_users_postgres to persist), but the data is -//// rebuildable, so that is acceptable. +//// Lost on restart, rebuildable; known_users_postgres persists instead. import at_record_server/known_users.{type KnownUser, type Store, Store} +import at_record_server/parallel import gleam/dict.{type Dict} import gleam/erlang/process.{type Subject} import gleam/otp/actor @@ -24,7 +24,7 @@ pub fn start() -> Result(Store, actor.StartError) { Nil }, list: fn() { - case try_call(subject, 1000, List) { + case parallel.try_call(subject, 1000, List) { Ok(users) -> users Error(Nil) -> [] } @@ -32,26 +32,6 @@ pub fn start() -> Result(Store, actor.StartError) { ) } -// process.call/actor.call panic on timeout or a dead callee; this mirrors process.call's own monitor-based body but returns Error(Nil) instead, so list() can honor its best-effort contract. -fn try_call( - subject: Subject(msg), - timeout: Int, - make_request: fn(Subject(reply)) -> msg, -) -> Result(reply, Nil) { - use callee <- result.try(process.subject_owner(subject)) - let reply_subject = process.new_subject() - let monitor = process.monitor(callee) - process.send(subject, make_request(reply_subject)) - let reply = - process.new_selector() - |> process.select_map(reply_subject, Ok) - |> process.select_specific_monitor(monitor, fn(_down) { Error(Nil) }) - |> process.selector_receive(timeout) - |> result.unwrap(Error(Nil)) - process.demonitor_process(monitor) - reply -} - fn handle( state: Dict(String, KnownUser), msg: Msg, diff --git a/server/src/at_record_server/parallel.gleam b/server/src/at_record_server/parallel.gleam new file mode 100644 index 0000000..f07cb78 --- /dev/null +++ b/server/src/at_record_server/parallel.gleam @@ -0,0 +1,106 @@ +//// Best-effort process-concurrency helpers. `try_call` is a request/reply to +//// one actor that returns `Error(Nil)` rather than panicking on a timeout or a +//// dead callee. `map` is scatter-gather where a crashed or timed-out worker +//// yields `Error(Nil)` in its (order-preserving) slot. + +import gleam/dict.{type Dict} +import gleam/erlang/process.{type Selector, type Subject} +import gleam/list +import gleam/result + +pub fn try_call( + subject: Subject(msg), + timeout: Int, + make_request: fn(Subject(reply)) -> msg, +) -> Result(reply, Nil) { + use callee <- result.try(process.subject_owner(subject)) + let reply_subject = process.new_subject() + let monitor = process.monitor(callee) + process.send(subject, make_request(reply_subject)) + let reply = + process.new_selector() + |> process.select_map(reply_subject, Ok) + |> process.select_specific_monitor(monitor, fn(_down) { Error(Nil) }) + |> process.selector_receive(timeout) + |> result.unwrap(Error(Nil)) + process.demonitor_process(monitor) + reply +} + +type Event(b) { + Done(Int, b) + Crashed(Int) + Deadline +} + +pub fn map( + over items: List(a), + within timeout: Int, + with work: fn(a) -> b, +) -> List(Result(b, Nil)) { + let reply = process.new_subject() + process.spawn_unlinked(fn() { + process.send(reply, coordinate(items, timeout, work)) + }) + process.receive(reply, timeout + 200) + |> result.unwrap(list.map(items, fn(_) { Error(Nil) })) +} + +fn coordinate( + items: List(a), + timeout: Int, + work: fn(a) -> b, +) -> List(Result(b, Nil)) { + let inbox = process.new_subject() + let deadline = process.new_subject() + let monitors = + list.index_map(items, fn(item, i) { + let pid = + process.spawn_unlinked(fn() { process.send(inbox, #(i, work(item))) }) + #(i, process.monitor(pid)) + }) + process.send_after(deadline, timeout, Nil) + let done = + gather(selector(inbox, deadline, monitors), dict.new(), list.length(items)) + list.index_map(items, fn(_, i) { + result.unwrap(dict.get(done, i), Error(Nil)) + }) +} + +fn selector( + inbox: Subject(#(Int, b)), + deadline: Subject(Nil), + monitors: List(#(Int, process.Monitor)), +) -> Selector(Event(b)) { + let base = + process.new_selector() + |> process.select_map(inbox, fn(pair) { Done(pair.0, pair.1) }) + |> process.select_map(deadline, fn(_) { Deadline }) + list.fold(monitors, base, fn(sel, entry) { + process.select_specific_monitor(sel, entry.1, fn(_down) { Crashed(entry.0) }) + }) +} + +fn gather( + selector: Selector(Event(b)), + acc: Dict(Int, Result(b, Nil)), + total: Int, +) -> Dict(Int, Result(b, Nil)) { + case dict.size(acc) == total { + True -> acc + False -> + case process.selector_receive_forever(selector) { + Done(i, value) -> + gather(selector, dict.insert(acc, i, Ok(value)), total) + Crashed(i) -> gather(selector, settle(acc, i), total) + Deadline -> acc + } + } +} + +fn settle(acc: Dict(Int, Result(b, Nil)), i: Int) -> Dict(Int, Result(b, Nil)) { + case dict.has_key(acc, i) { + True -> acc + False -> dict.insert(acc, i, Error(Nil)) + } +} diff --git a/server/test/parallel_test.gleam b/server/test/parallel_test.gleam new file mode 100644 index 0000000..8c05cfc --- /dev/null +++ b/server/test/parallel_test.gleam @@ -0,0 +1,39 @@ +import at_record_server/parallel +import gleam/erlang/process + +pub fn maps_and_preserves_order_test() { + assert parallel.map([1, 2, 3], 1000, fn(x) { x * 2 }) == [Ok(2), Ok(4), Ok(6)] +} + +pub fn empty_list_test() { + assert parallel.map([], 1000, fn(x) { x }) == [] +} + +pub fn slow_worker_times_out_test() { + let out = + parallel.map([10, 500], 100, fn(ms) { + process.sleep(ms) + ms + }) + assert out == [Ok(10), Error(Nil)] +} + +pub fn runs_concurrently_test() { + let out = + parallel.map([1, 2, 3], 200, fn(n) { + process.sleep(80) + n + }) + assert out == [Ok(1), Ok(2), Ok(3)] +} + +pub fn crashed_worker_yields_error_test() { + let out = + parallel.map([1, 2, 3], 2000, fn(n) { + case n { + 2 -> panic as "boom" + _ -> n + } + }) + assert out == [Ok(1), Error(Nil), Ok(3)] +}