Side B implements a substantial, working feature (SSE-based entity fetch streaming with proper job result tracking, oneshot notification, richer error handling/logging, and updated JS client + integration test), representing real functional value. Side A is merely deleting a leftover unused file, which is a trivial cleanup with no functional impact.
constitution · epochs · watch · epoch 3
c_b6adf338f405 (tommy-mor) vs c_06fce70179bc (tommy-mor)
download prompt · raw event · cmp_1f9e74025b1e21
council reasoning
A only removes an unused timeline.rs module (dead-code cleanup). B delivers lasting architecture: SSE-based entity fetch with oneshot job results, a dedicated fetch module, richer Reddit worker outcomes/error paths, and matching UI/JS/test updates—real product and design value beyond deletion.
Side B introduces substantial new functionality by refactoring entity fetching into an SSE-based streaming flow, adding a dedicated fetch module, asynchronous completion notifications with oneshot channels, client-side SSE handling, and corresponding integration tests. Side A only deletes an apparently vestigial `timeline.rs` file without replacing or improving functionality, so its lasting project value is much smaller.
sides
A — c_b6adf338f405 (tommy-mor)
message
[6c7377b6] deleted vestigal file
diff preview
diff --git a/server/src/timeline.rs b/server/src/timeline.rs
deleted file mode 100644
index 265ca4ae9946ecd29a7d1d4ec791de8465ce658e..0000000000000000000000000000000000000000
--- a/server/src/timeline.rs
+++ /dev/null
@@ -1,145 +0,0 @@
-//! Room admin lines merged into forum thread views.
-
-use crate::{
- canonical_path::canonicalize_tag,
- reducer::{ReducerState, RoomTimelineEntry, RoomTimelineKind},
-};
-
-fn cap_label(c: crate::events::ThreadCapability) -> &'static str {
- use crate::events::ThreadCapability::*;
- match c {
- View => "view",
- Post => "post",
- Vote => "vote",
- AddItem => "add_item",
- Manage => "manage",
- }
-}
-
-fn caps_list(caps: &[crate::events::ThreadCapability]) -> String {
- let mut v: Vec<_> = caps.iter().map(|c| cap_label(*c)).collect();
- v.sort();
- v.join(", ")
-}
-
-/// Human-readable system line for the thread feed.
-pub fn format_room_timeline_entry(e: &RoomTimelineEntry) -> String {
- match &e.kind {
- RoomTimelineKind::RoomCreated { owner, slug } => {
- format!("@{owner} created room #{slug}")
- }
- RoomTimelineKind::GrantAdded {
- username,
- granted_by,
- capabilities,
- } => {
- format!(
- "@{granted_by} granted @{} {}",
- username,
- caps_list(capabilities)
- )
- }
- RoomTimelineKind::GrantRevoked {
- username,
- revoked_by,
- capabilities,
- } => {
- format!(
- "@{revoked_by} revoked @{} {}",
- username,
- caps_list(capabilities)
- )
- }
- }
-}
-
-#[derive(Clone, Debug)]
-pub enum MergedThreadRow {
- System { ts: i64, text: String },
- Post {
- index: usize,
- id: String,
- ts: i64,
- principal: String,
- raw: String,
- },
-}
-
-/// Merge room admin lines with thread ingests for one room + tag. Oldest first.
-/// `actor_prefix` filters posts only (system lines always included).
-pub fn merge_thread_rows(
- reduced: &ReducerState,
- room_wire: &str,
- thread_tag: &str,
- since: Option<i64>,
- before: Option<i64>,
- actor_prefix: &str,
-) -> Vec<MergedThreadRow> {
- let scope = crate::reducer::scope_from_room_wire(room_wire);
- let tag = canonicalize_tag(thread_tag);
- let key = (scope.clone(), tag.clone());
-
- let mut rows: Vec<MergedThreadRow> = Vec::new();
-
- if let Some(entries) = reduced.room_timeline.get(room_wire.trim()) {
- for e in entries {
- if since.map_or(true, |s| e.ts >= s) && before.map_or(true, |b| e.ts < b) {
- rows.push(MergedThreadRow::System {
- ts: e.ts,
- text: format_room_timeline_entry(e),
- });
- }
- }
- }
-
- let all_ids: Vec<String> = reduced
- .ingests_by_scope_thread
- .get(&key)
- .map(|q| q.iter().rev().cloned().collect())
- .unwrap_or_default();
-
- for (idx, id) in all_ids.into_iter().enumerate() {
- let Some(ing) = reduced.ingests_by_id.get(&id) else {
- continue;
- };
- if since.map_or(true, |s| ing.ts >= s) && before.map_or(true, |b| ing.ts < b) {
- if !actor_prefix.is_empty()
- && !ing
- .principal
- .to_lowercase()
- .starts_with(actor_prefix)
- {
- continue;
- }
- rows.push(MergedThreadRow::Post {
- index: idx,
- id: ing.id.clone(),
- ts: ing.ts,
- principal: ing.principal.clone(),
- raw: ing.raw.clone(),
- });
- }
- }
-
- rows.sort_by(|a, b| {
- let ta = match a {
- MergedThreadRow::System { ts, .. } | MergedThreadRow::Post { ts, .. } => *ts,
- };
- let tb = match b {
- MergedThreadRow::System { ts, .. } | MergedThreadRow::Post { ts, .. } => *ts,
- };
- ta.cmp(&tb)
- });
- rows
-}
-
-/// Shared-site thread view (`room_wire == "public"` on the wire). Same merge as private rooms; private-room timeline is unused here.
-pub fn merge_public_thread_rows(
- reduced: &ReducerState,
- thread_tag: &str,
- since: Option<i64>,
- before: Option<i64>,
- actor_prefix: &str,
-) -> Vec<MergedThreadRow> {
- merge_thread_rows(reduced, "public", thread_tag, since, before, actor_prefix)
-}
B — c_06fce70179bc (tommy-mor)
message
[6d04afc2] refactor
diff preview
diff --git a/Cargo.lock b/Cargo.lock
index 2cea973082716e761ef6f5dd5886acc08ff9aac0..8c43fb75c472b102e6e1d3b837dce3355be898f2 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -17,6 +17,28 @@ version = "1.0.102"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c"
+[[package]]
+name = "async-stream"
+version = "0.3.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "0b5a71a6f37880a80d1d7f19efd781e4b5de42c88f0722cc13bcb6cc2cfe8476"
+dependencies = [
+ "async-stream-impl",
+ "futures-core",
+ "pin-project-lite",
+]
+
+[[package]]
+name = "async-stream-impl"
+version = "0.3.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c7c24de15d275a1ecfd47a380fb4d5ec9bfe0933f309ed5e705b775596a3574d"
+dependencies = [
+ "proc-macro2",
+ "quote",
+ "syn",
+]
+
[[package]]
name = "async-trait"
version = "0.1.89"
@@ -1242,9 +1264,11 @@ dependencies = [
name = "sorter2-server"
version = "0.0.1"
dependencies = [
+ "async-stream",
"axum",
"axum-extra",
"dotenvy",
+ "futures-util",
"maud",
"reqwest",
"serde",
diff --git a/server/Cargo.toml b/server/Cargo.toml
index bd600138b613bd0f546bdec217a5334cdcb20aa5..c940acb687fb141d21760a3d6656172013cf6f41 100644
--- a/server/Cargo.toml
+++ b/server/Cargo.toml
@@ -18,6 +18,8 @@ tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
reqwest = { version = "0.12", features = ["json"] }
dotenvy = "0.15"
+async-stream = "0.3"
+futures-util = { version = "0.3", default-features = false, features = ["std"] }
[dev-dependencies]
reqwest = { version = "0.12", features = ["json"] }
diff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs
index b33a84e8bb5e817b26592868d88090e6d664d950..7af6527d03c483f33f3469ce6766c01a554c5fe3 100644
--- a/server/src/api/ui_html.rs
+++ b/server/src/api/ui_html.rs
@@ -6,7 +6,8 @@ use axum::{
use std::collections::HashMap;
use crate::{
- html::{entity_section, input_panel, js_string_literal, ranking_panel, JsBuilder},
+ fetch,
+ html::{input_panel, js_string_literal, ranking_panel, JsBuilder},
parser::parse_reddit_url,
path_types::ItemId,
reddit::ensure_partial_tree,
@@ -89,18 +90,8 @@ pub async fn post_ui_html(
},
HtmlUiAction::FetchEntity { item } => {
let id = parse_item_param(&item);
- if id.is_root() {
- return ui_js_warn("nothing to fetch for the root").into_response();
- }
- state.queue_entity_fetch(id.clone());
- let tree = state.tree.read().await;
- let empty = crate::reducer::NodeState::default();
- let node = tree.get(&id).unwrap_or(&empty);
- let panel = entity_section(&id, node, true);
- JsBuilder::new()
- .morph_selector("#entity-section", panel)
- .into_response()
- },
+ fetch::fetch_entity_stream(state, id).into_response()
+ }
}
}
diff --git a/server/src/fetch/html.rs b/server/src/fetch/html.rs
new file mode 100644
index 0000000000000000000000000000000000000000..63634508496e224c38b9ec0308b7a6086462f925
--- /dev/null
+++ b/server/src/fetch/html.rs
@@ -0,0 +1,67 @@
+//! Markup for entity import / “Fetch from Reddit” (`POST /ui`, SSE response).
+
+use maud::{html, Markup};
+
+use crate::{
+ form_template::template_json_compact,
+ path_types::ItemId,
+ reddit::is_fetchable,
+ reducer::NodeState,
+ ui_action::UI_RPC_FIELD,
+};
+
+fn entity_panel(node: &NodeState) -> Markup {
+ html! {
+ @if let Some(data) = &node.data {
+ div id="entity-panel" class="entity-card" {
+ h2 { (data.title) }
+ @if let Some(author) = &data.author {
+ p class="muted small" { "by " (author) }
+ }
+ @if let Some(body) = &data.body_html {
+ div class="entity-body" { (maud::PreEscaped(body)) }
+ }
+ }
+ }
+ }
+}
+
+/// Reddit/API import — `POST /ui` with `fetch_entity` returns an SSE stream.
+pub fn fetch_entity_panel(item: &ItemId, has_data: bool, fetching: bool) -> Markup {
+ if !is_fetchable(item) {
+ return html! {};
+ }
+ let label = if fetching {
+ "Fetching…"
+ } else if has_data {
+ "Fetch more"
+ } else {
+ "Fetch from Reddit"
+ };
+ let rpc = template_json_compact(&serde_json::json!({
+ "action": "fetch_entity",
+ "item": item.as_str(),
+ }))
+ .expect("fetch_entity rpc template");
+ html! {
+ form method="post" action="/ui" id="fetch-entity-form" class="fetch-entity-form" {
+ input type="hidden" name=(UI_RPC_FIELD) value=(rpc);
+ @if fetching {
+ button type="submit" class="btn-secondary" disabled { (label) }
+ } @else {
+ button type="submit" class="btn-secondary" { (label) }
+ }
+ }
+ }
+}
+
+/// Entity card + fetch control (target `#entity-section` for Idiomorph / SSE).
+pub fn entity_section(item: &ItemId, node: &NodeState, fetching: bool) -> Markup {
+ let has_data = node.data.is_some();
+ html! {
+ section id="entity-section" class="demo-panel" {
+ (entity_panel(node))
+ (fetch_entity_panel(item, has_data, fetching))
+ }
+ }
+}
diff --git a/server/src/fetch/mod.rs b/server/src/fetch/mod.rs
new file mode 100644
index 0000000000000000000000000000000000000000..2290f9d3a0f1cbf1806c6339f82a4515c11cc3d3
--- /dev/null
+++ b/server/src/fetch/mod.rs
@@ -0,0 +1,115 @@
+//! Entity import over `POST /ui` as SSE (Reddit worker in [`crate::reddit`]).
+
+pub mod html;
+
+use std::convert::Infallible;
+use std::time::Duration;
+
+use async_stream::stream;
+use axum::response::sse::{Event, KeepAlive, Sse};
+use futures_util::Stream;
+use serde::Serialize;
+use tokio::sync::oneshot;
+
+use crate::{
+ path_types::ItemId,
+ reddit::FetchJobResult,
+ reducer::NodeState,
+ state::AppState,
+};
+
+pub fn now_ms() -> i64 {
+ let t = std::time::SystemTime::now()
+ .duration_since(std::time::UNIX_EPOCH)
+ .unwrap_or_default();
+ t.as_millis() as i64
+}
+
+#[derive(Serialize)]
+struct SseMorphPayload {
+ selector: &'static str,
+ html: String,
+}
+
+fn morph_complete_event(html: maud::Markup) -> Event {
+ let payload = SseMorphPayload {
+ selector: "#entity-section",
+ html: html.into_string(),
+ };
+ let data = serde_json::to_string(&payload).unwrap_or_else(|_| "{}".into());
+ Event::default().event("complete").data(data)
+}
+
+/// Stream `fetching` → `complete` / `error` for [`crate::ui_action::HtmlUiAction::FetchEntity`].
+pub fn fetch_entity_stream(
+ state: AppState,
+ id: ItemId,
+) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
+ tracing::debug!(item = %id, "fetch entity stream opened");
+
+ let stream = stream! {
+ if id.is_root() {
+ yield Ok(Event::default().event("error").data("{\"message\":\"nothing to fetch for the root\"}"));
+ return;
+ }
+
+ if !crate::reddit::is_fetchable(&id) {
+ tracing::debug!(item = %id, "fetch stream: not fetchable");
+ yield Ok(Event::default().event("error").data("{\"message\":\"this page cannot be fetched from Reddit\"}"));
+ return;
+ }
+
+ let fetching_html = {
+ let tree = state.tree.read().await;
+ let empty = NodeState::default();
+ let node = tree.get(&id).unwrap_or(&empty);
+ html::entity_section(&id, node, true).into_string()
+ };
+ let fetching_payload = serde_json::json!({
+ "selector": "#entity-section",
+ "html": fetching_html,
+ });
+ yield Ok(Event::default().event("fetching").data(fetching_payload.to_string()));
+
+ let (tx, rx) = oneshot::channel();
+ state.reddit.request_fetch(id.clone(), true, Some(tx));
+ tracing::debug!(item = %id, "fetch stream: queued reddit job");
+
+ let result = match rx.await {
+ Ok(r) => r,
+ Err(_) => {
+ tracing::warn!(item = %id, "fetch stream: worker dropped oneshot");
+ FetchJobResult::Failed("reddit worker stopped".into())
+ }
+ };
+
+ tracing::debug!(item = %id, ?result, "fetch stream: job finished");
+
+ match result {
+ FetchJobResult::Imported | FetchJobResult::NotFound => {
+ let tree = state.tree.read().await;
+ let empty = NodeState::default();
+ let node = tree.get(&id).unwrap_or(&empty);
+ yield Ok(morph_complete_event(html::entity_section(&id, node, false)));
+ }
+ FetchJobResult::SkippedCached | FetchJobResult::SkippedDuplicate => {
+ let tree = state.tree.read().await;
+ let empty = NodeState::default();
+ let node = tree.get(&id).unwrap_or(&empty);
+ yield Ok(morph_complete_event(html::entity_section(&id, node, false)));
+ }
+ FetchJobResult::RateLimited { reset_secs } => {
+ yield Ok(Event::default().event("error").data(
+ serde_json::json!({"message": format!("Reddit rate limit — retry in {reset_secs}s")}).to_string(),
+ ));
+ }
+ FetchJobResult::Failed(msg) => {
+ yield Ok(Event::default().event("error").data(
+ serde_json::json!({"message": msg}).to_string(),
+ ));
+ }
+ }
+ };
+
+ Sse::new(stream).keep_alive(KeepAlive::new().interval(Duration::from_secs(15)))
+}
diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs
index db5b4c7f06b0be64603981166835cde268234f67..9314a7556306ddab969b896dbf4126b542a46722 100644
--- a/server/src/html/mod.rs
+++ b/server/src/html/mod.rs
@@ -7,10 +7,10 @@ use axum::{
use maud::{html, Markup, DOCTYPE};
use crate::{
+ fetch::html::entity_section,
form_template::template_json_compact,
path_types::ItemId,
ranking::{top_bottom, RankedItem},
- reddit::is_fetchable,
reducer::{GroupState, NodeState},
state::AppState,
ui_action::UI_RPC_FIELD,
@@ -149,62 +149,6 @@ pub fn breadcrumb_path(item: &ItemId) -> Markup {
}
}
-fn entity_panel(node: &NodeState) -> Markup {
- html! {
- @if let Some(data) = &node.data {
- div id="entity-panel" class="entity-card" {
- h2 { (data.title) }
- @if let Some(author) = &data.author {
- p class="muted small" { "by " (author) }
- }
- @if let Some(body) = &data.body_html {
- div class="entity-body" { (maud::PreEscaped(body)) }
- }
- }
- }
- }
-}
-
-/// Reddit/API import control — only shown on fetchable pages; never auto-fires.
-pub fn fetch_entity_panel(item: &ItemId, has_data: bool, fetching: bool) -> Markup {
- if !is_fetchable(item) {
- return html! {};
- }
- let label = if fetching {
- "Fetching…"
- } else if has_data {
- "Fetch more"
- } else {
- "Fetch from Reddit"
- };
- let rpc = template_json_compact(&serde_json::json!({
- "action": "fetch_entity",
- "item": item.as_str(),
- }))
- .expect("fetch_entity rpc template");
- html! {
- form method="post" action="/ui" id="fetch-entity-form" class="fetch-entity-form" {
- input type="hidden" name=(UI_RPC_FIELD) value=(rpc);
- @if fetching {
- button type="submit" class="btn-secondary" disabled { (label) }
- } @else {
- button type="submit" class="btn-secondary" { (label) }
- }
- }
- }
-}
-
-/// Entity
… preview truncated; 22,618 characters omittedHardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.