Side B is a targeted, real bugfix: it correctly forces OAuth when credentials exist, adds retry-on-401/403 logic with token refresh, improves error reporting, pins the base URL in fly.toml, and includes a new test validating the error-truncation logic. Side A merely adds unintegrated scaffold files (Dockerfile, deps.edn, event_log.rs, views.rs) that aren't wired into the build and provide no demonstrated behavior or fix, offering speculative seed code rather than lasting, verified value.
constitution · epochs · watch · epoch 3
c_16438843de8f (tommy-mor) vs c_ca9169f732b8 (tommy-mor)
download prompt · raw event · cmp_bb69ecc9651e45
council reasoning
B is a targeted production fix: force OAuth when credentials exist (no public www.reddit.com fallback that cloud IPs hit as 403 HTML), add AuthRejected + one-shot token refresh, tighten ensure_oauth_token errors, and pin SORTER2_BASE_URL, with a unit test for truncate_for_error. A mainly seeds Dockerfile/fly.toml plus new event_log/views modules and an unrelated Clojure deps.edn—useful scaffolding, but more breadth than a lasting correctness fix on the live Reddit path.
Side B fixes a concrete production behavior by requiring OAuth when credentials are configured, adding token refresh and retry on 401/403 responses, improving error propagation, and introducing a targeted test for error truncation. Side A mostly adds infrastructure and new utility files (Dockerfile, deployment config, event log, view persistence, Clojure deps), but without showing their integration into the application, making their lasting project impact less certain than B's focused reliability fix.
sides
A — c_16438843de8f (tommy-mor)
message
[4cd0d15d] more seed
diff preview
diff --git a/Dockerfile b/Dockerfile
new file mode 100644
index 0000000000000000000000000000000000000000..9cb07c60cb0da063f747cfbf1b3b876ecb8ba03e
--- /dev/null
+++ b/Dockerfile
@@ -0,0 +1,34 @@
+# time 0.3.47+ requires Rust 1.88 (edition 2024)
+FROM rust:1.88-slim as builder
+
+WORKDIR /build
+
+RUN apt-get update && \
+ apt-get install -y pkg-config libssl-dev && \
+ rm -rf /var/lib/apt/lists/*
+
+# Copy source and build. (Keep it simple to avoid remote build cache oddities.)
+COPY . .
+RUN cargo build --release --package slugsocial-server
+
+FROM debian:bookworm-slim
+
+RUN apt-get update && \
+ apt-get install -y ca-certificates && \
+ rm -rf /var/lib/apt/lists/*
+
+WORKDIR /app
+
+COPY --from=builder /build/target/release/slugsocial-server /app/slugsocial-server
+
+# Create data directory for persistent volume
+RUN mkdir -p /data
+
+ENV SLUG_DATA_DIR=/data
+ENV SLUG_EVENT_LOG=/data/events.jsonl
+ENV PORT=8080
+
+EXPOSE 8080
+
+CMD ["/app/slugsocial-server"]
+
diff --git a/deps.edn b/deps.edn
new file mode 100644
index 0000000000000000000000000000000000000000..0bf892d44f491cb2313e01ae8a942c3097c52948
--- /dev/null
+++ b/deps.edn
@@ -0,0 +1,10 @@
+{:paths ["." "test"]
+ :deps {cheshire/cheshire {:mvn/version "5.13.0"}
+ http-kit/http-kit {:mvn/version "2.8.0"}
+ babashka/fs {:mvn/version "0.5.32"}
+ babashka/process {:mvn/version "0.6.25"}
+ com.blockether/spel {:mvn/version "0.7.11"}}
+ :aliases
+ {:kaocha {:extra-deps {lambdaisland/kaocha {:mvn/version "1.91.1392"}
+ lambdaisland/kaocha-junit-xml {:mvn/version "1.17.101"}}
+ :main-opts ["-m" "kaocha.runner"]}}}
diff --git a/event_log.rs b/event_log.rs
new file mode 100644
index 0000000000000000000000000000000000000000..eaae0d495e43a45d6590603892265a62cc92906e
--- /dev/null
+++ b/event_log.rs
@@ -0,0 +1,83 @@
+use std::path::{Path, PathBuf};
+
+use tokio::{
+ fs::{self, OpenOptions},
+ io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
+};
+
+use crate::events::Event;
+
+#[derive(Debug, thiserror::Error)]
+pub enum EventLogError {
+ #[error("io error: {0}")]
+ Io(#[from] std::io::Error),
+ #[error("json error: {0}")]
+ Json(#[from] serde_json::Error),
+}
+
+#[derive(Debug, Clone)]
+pub struct EventLog {
+ path: PathBuf,
+}
+
+impl EventLog {
+ pub fn new(path: impl Into<PathBuf>) -> Self {
+ Self { path: path.into() }
+ }
+
+ pub fn path(&self) -> &Path {
+ &self.path
+ }
+
+ pub async fn ensure_parent_dir(&self) -> Result<(), EventLogError> {
+ if let Some(parent) = self.path.parent() {
+ fs::create_dir_all(parent).await?;
+ }
+ Ok(())
+ }
+
+ pub async fn append(&self, event: &Event) -> Result<(), EventLogError> {
+ self.ensure_parent_dir().await?;
+ let mut f: tokio::fs::File = OpenOptions::new()
+ .create(true)
+ .append(true)
+ .open(&self.path)
+ .await?;
+
+ let mut line = serde_json::to_string(event)?;
+ line.push('\n');
+ f.write_all(line.as_bytes()).await?;
+ f.flush().await?;
+ Ok(())
+ }
+
+ /// Load events from JSONL. Corrupt lines are skipped and returned as `(line_no, line)`.
+ pub async fn load_all(&self) -> Result<(Vec<Event>, Vec<(usize, String)>), EventLogError> {
+ if !fs::try_exists(&self.path).await? {
+ return Ok((vec![], vec![]));
+ }
+
+ let f = fs::File::open(&self.path).await?;
+ let mut reader = BufReader::new(f).lines();
+
+ let mut events = Vec::new();
+ let mut bad_lines = Vec::new();
+
+ let mut line_no: usize = 0;
+ while let Some(line) = reader.next_line().await? {
+ line_no += 1;
+ let trimmed = line.trim();
+ if trimmed.is_empty() {
+ continue;
+ }
+ match serde_json::from_str::<Event>(trimmed) {
+ Ok(ev) => events.push(ev),
+ Err(_) => bad_lines.push((line_no, line)),
+ }
+ }
+
+ Ok((events, bad_lines))
+ }
+}
+
+
diff --git a/fly.toml b/fly.toml
new file mode 100644
index 0000000000000000000000000000000000000000..bbb9345e527452db1d87a549213645c195eae5fc
--- /dev/null
+++ b/fly.toml
@@ -0,0 +1,42 @@
+app = "slugsocial"
+primary_region = "iad"
+
+[build]
+ dockerfile = "Dockerfile"
+
+[env]
+ SLUG_DATA_DIR = "/data"
+ SLUG_EVENT_LOG = "/data/events.jsonl"
+ PORT = "8080"
+
+[[services]]
+ internal_port = 8080
+ protocol = "tcp"
+
+ [[services.ports]]
+ port = 80
+ handlers = ["http"]
+ force_https = true
+
+ [[services.ports]]
+ port = 443
+ handlers = ["tls", "http"]
+
+ [services.concurrency]
+ type = "connections"
+ hard_limit = 1000
+ soft_limit = 500
+
+ [[services.http_checks]]
+ interval = "10s"
+ timeout = "2s"
+ grace_period = "5s"
+ method = "GET"
+ path = "/healthz"
+ protocol = "http"
+ tls_skip_verify = false
+
+[[mounts]]
+ source = "slugsocial_data"
+ destination = "/data"
+
diff --git a/views.rs b/views.rs
new file mode 100644
index 0000000000000000000000000000000000000000..d4f0ffc49475f014698b4da0de6f476884430813
--- /dev/null
+++ b/views.rs
@@ -0,0 +1,63 @@
+use std::{
+ collections::HashMap,
+ sync::{Arc, Mutex},
+};
+use tokio::sync::mpsc;
+
+type CountMap = Arc<Mutex<HashMap<String, u64>>>;
+
+#[derive(Clone)]
+pub struct ViewStore {
+ counts: CountMap,
+ flush_tx: mpsc::Sender<()>,
+}
+
+impl ViewStore {
+ pub fn new(json_path: &str) -> Self {
+ // Load existing counts from disk on startup (best-effort)
+ let initial: HashMap<String, u64> = std::fs::read_to_string(json_path)
+ .ok()
+ .and_then(|s| serde_json::from_str(&s).ok())
+ .unwrap_or_default();
+
+ let counts: CountMap = Arc::new(Mutex::new(initial));
+ let (flush_tx, mut flush_rx) = mpsc::channel::<()>(64);
+ let path = json_path.to_string();
+
+ let counts_for_writer = counts.clone();
+ tokio::spawn(async move {
+ while flush_rx.recv().await.is_some() {
+ while flush_rx.try_recv().is_ok() {}
+
+ let snapshot: HashMap<String, u64> = {
+ counts_for_writer.lock().unwrap().clone()
+ };
+
+ let path = path.clone();
+ let _ = tokio::task::spawn_blocking(move || {
+ if let Ok(json) = serde_json::to_string(&snapshot) {
+ let tmp = format!("{path}.tmp");
+ if std::fs::write(&tmp, &json).is_ok() {
+ let _ = std::fs::rename(&tmp, &path);
+ }
+ }
+ })
+ .await;
+ }
+ });
+
+ Self { counts, flush_tx }
+ }
+
+ pub fn increment(&self, path: String) {
+ {
+ let mut map = self.counts.lock().unwrap();
+ *map.entry(path).or_insert(0) += 1;
+ }
+ let _ = self.flush_tx.try_send(());
+ }
+
+ pub fn get_views(&self, path: &str) -> u64 {
+ self.counts.lock().unwrap().get(path).copied().unwrap_or(0)
+ }
+}
B — c_ca9169f732b8 (tommy-mor)
message
[8f69c309] Require Reddit OAuth when credentials are set and refresh on 401/403. Avoid falling back to the public www.reddit.com API from cloud IPs, which returns Reddit's network-security block page. Also pin SORTER2_BASE_URL in fly.toml. Co-authored-by: Cursor <cursoragent@cursor.com>
diff preview
diff --git a/fly.toml b/fly.toml
index f0c6a39f643c204987debc177234d95a8ef44b65..ca7e0088a7efd7d58d29d808f8d89080b2bff233 100644
--- a/fly.toml
+++ b/fly.toml
@@ -5,6 +5,7 @@ primary_region = "iad"
dockerfile = "Dockerfile"
[env]
+ SORTER2_BASE_URL = "https://reddit.sorter.social"
SORTER2_DATA_DIR = "/data"
SORTER2_EVENT_LOG = "/data/events.jsonl"
PORT = "8080"
diff --git a/server/src/reddit.rs b/server/src/reddit.rs
index a874814f8927192ee62cab2d0db1efd27dcd57b7..f409764c1e1f36216f1b08107043c2eab905694c 100644
--- a/server/src/reddit.rs
+++ b/server/src/reddit.rs
@@ -283,26 +283,35 @@ async fn reddit_worker(
);
tokio::time::sleep(current_delay).await;
- if let Some(c) = &creds {
- oauth = ensure_oauth_token(&client, &oauth_token_base, c, oauth.take()).await;
- }
-
- let token = oauth.as_ref().map(|t| t.access_token.as_str());
- let fetch_base = if token.is_some() {
- tracing::debug!(
- item = %fetch_id,
- base = %oauth_api_base,
- "reddit fetch using OAuth bearer"
- );
- &oauth_api_base
- } else {
- &api_base
- };
- let url = match kind {
- FetchKind::SelfEntity => map_item_to_reddit_api(&fetch_id, fetch_base),
- FetchKind::Children => map_children_url(&fetch_id, fetch_base),
+ let outcome = match &creds {
+ Some(c) => {
+ // OAuth is required when credentials are configured — never fall
+ // back to the public www.reddit.com JSON endpoints (cloud IPs
+ // get blocked with a 403 HTML interstitial).
+ fetch_with_oauth(
+ &client,
+ &oauth_token_base,
+ &oauth_api_base,
+ c,
+ &mut oauth,
+ &fetch_id,
+ kind,
+ )
+ .await
+ }
+ None => {
+ let url = match kind {
+ FetchKind::SelfEntity => map_item_to_reddit_api(&fetch_id, &api_base),
+ FetchKind::Children => map_children_url(&fetch_id, &api_base),
+ };
+ match do_fetch(&client, &url, &fetch_id, None).await {
+ Ok(FetchOutcome::AuthRejected { status, detail }) => {
+ Err(format!("Reddit API {status}: {detail}"))
+ }
+ other => other,
+ }
+ }
};
- let outcome = do_fetch(&client, &url, &fetch_id, token).await;
match outcome {
Ok(FetchOutcome::Payload(payload)) => {
@@ -342,6 +351,12 @@ async fn reddit_worker(
current_delay = (current_delay * 2).min(Duration::from_secs(60));
notify(done, FetchJobResult::RateLimited { reset_secs });
}
+ Ok(FetchOutcome::AuthRejected { status, detail }) => {
+ let e = format!("Reddit API {status}: {detail}");
+ tracing::warn!(item = %fetch_id, err = %e, "reddit fetch auth rejected");
+ current_delay = (current_delay * 2).min(Duration::from_secs(60));
+ notify(done, FetchJobResult::Failed(e));
+ }
Err(e) => {
tracing::warn!(item = %fetch_id, err = %e, "reddit fetch failed");
current_delay = (current_delay * 2).min(Duration::from_secs(60));
@@ -357,6 +372,60 @@ enum FetchOutcome {
Payload(Value),
NotFound,
RateLimited { reset_secs: u64 },
+ /// Bearer rejected — caller should drop the cached token and retry once.
+ AuthRejected { status: StatusCode, detail: String },
+}
+
+async fn fetch_with_oauth(
+ client: &Client,
+ oauth_token_base: &str,
+ oauth_api_base: &str,
+ creds: &RedditCredentials,
+ oauth: &mut Option<OAuthToken>,
+ fetch_id: &ItemId,
+ kind: FetchKind,
+) -> Result<FetchOutcome, String> {
+ for attempt in 0..2 {
+ let force_refresh = attempt > 0;
+ *oauth = Some(
+ ensure_oauth_token(client, oauth_token_base, creds, oauth.take(), force_refresh)
+ .await?,
+ );
+ let token = oauth
+ .as_ref()
+ .expect("token set above")
+ .access_token
+ .clone();
+
+ tracing::debug!(
+ item = %fetch_id,
+ base = %oauth_api_base,
+ attempt,
+ "reddit fetch using OAuth bearer"
+ );
+
+ let url = match kind {
+ FetchKind::SelfEntity => map_item_to_reddit_api(fetch_id, oauth_api_base),
+ FetchKind::Children => map_children_url(fetch_id, oauth_api_base),
+ };
+ match do_fetch(client, &url, fetch_id, Some(&token)).await? {
+ FetchOutcome::AuthRejected { status, detail } if attempt == 0 => {
+ tracing::warn!(
+ item = %fetch_id,
+ %status,
+ %detail,
+ "reddit OAuth rejected; refreshing token and retrying"
+ );
+ *oauth = None;
+ continue;
+ }
+ FetchOutcome::AuthRejected { status, detail } => {
+ return Err(format!("Reddit API {status}: {detail}"));
+ }
+ other => return Ok(other),
+ }
+ }
+ unreachable!("loop always returns")
}
async fn ensure_oauth_token(
@@ -364,35 +433,35 @@ async fn ensure_oauth_token(
oauth_base: &str,
creds: &RedditCredentials,
existing: Option<OAuthToken>,
-) -> Option<OAuthToken> {
- if let Some(t) = existing {
- if Instant::now() < t.expires_at - Duration::from_secs(60) {
- tracing::debug!("reddit OAuth token still valid");
- return Some(t);
+ force_refresh: bool,
+) -> Result<OAuthToken, String> {
+ if !force_refresh {
+ if let Some(t) = existing {
+ if Instant::now() < t.expires_at - Duration::from_secs(60) {
+ tracing::debug!("reddit OAuth token still valid");
+ return Ok(t);
+ }
}
}
let url = format!("{}/api/v1/access_token", oauth_base.trim_end_matches('/'));
- tracing::debug!(%url, "reddit OAuth token request");
+ tracing::debug!(%url, force_refresh, "reddit OAuth token request");
let resp = client
.post(&url)
.basic_auth(&creds.client_id, Some(&creds.client_secret))
.form(&[("grant_type", "client_credentials")])
.send()
- .await;
-
- let resp = match resp {
- Ok(r) => r,
- Err(e) => {
- tracing::warn!("reddit OAuth token request failed: {e}");
- return None;
- }
- };
+ .await
+ .map_err(|e| format!("Reddit OAuth token request failed: {e}"))?;
if !resp.status().is_success() {
- tracing::warn!("reddit OAuth token HTTP {}", resp.status());
- return None;
+ let status = resp.status();
+ let body = resp.text().await.unwrap_or_default();
+ return Err(format!(
+ "Reddit OAuth token HTTP {status}: {}",
+ truncate_for_error(&body)
+ ));
}
#[derive(Deserialize)]
@@ -401,21 +470,40 @@ async fn ensure_oauth_token(
expires_in: u64,
}
- let body: TokenResponse = match resp.json().await {
- Ok(b) => b,
- Err(e) => {
- tracing::warn!("reddit OAuth token parse failed: {e}");
- return None;
- }
- };
+ let body: TokenResponse = resp
+ .json()
+ .await
+ .map_err(|e| format!("Reddit OAuth token parse failed: {e}"))?;
- tracing::debug!(expires_in = body.expires_in, "reddit OAuth token acquired");
- Some(OAuthToken {
+ tracing::info!(expires_in = body.expires_in, "reddit OAuth token acquired");
+ Ok(OAuthToken {
access_token: body.access_token,
expires_at: Instant::now() + Duration::from_secs(body.expires_in),
})
}
+fn truncate_for_error(body: &str) -> String {
+ let compact: String = body.split_whitespace().collect::<Vec<_>>().join(" ");
+ if compact.is_empty() {
+ return "(empty body)".into();
+ }
+ // Prefer the human-readable block message over dumping Reddit's CSS.
+ if let Some(idx) = compact.find("You've been blocked") {
+ let slice: String = compact.chars().skip(idx).take(160).collect();
+ return if compact.chars().count() > idx + 160 {
+ format!("{slice}…")
+ } else {
+ slice
+ };
+ }
+ let chars: String = compact.chars().take(200).collect();
+ if compact.chars().count() > 200 {
+ format!("{chars}…")
+ } else {
+ chars
+ }
+}
+
async fn do_fetch(
client: &Client,
url: &str,
@@ -460,15 +548,16 @@ async fn do_fetch(
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
+ let detail = truncate_for_error(&body);
tracing::debug!(
item = %id,
%status,
body_len = body.len(),
- body_prefix = %body.chars().take(240).collect::<String>(),
+ %detail,
"reddit non-success body"
);
if status == StatusCode::FORBIDDEN || status == StatusCode::UNAUTHORIZED {
- return Err(format!("Reddit API {status}: {body}"));
+ return Ok(FetchOutcome::AuthRejected { status, detail });
}
return Ok(FetchOutcome::NotFound);
}
@@ -716,6 +805,15 @@ fn reddit_direct_image_url(url: &str) -> bool {
mod tests {
use super::*;
+ #[test]
+ fn truncate_error_prefers_block_message() {
+ let html = r#"<style>.x{color:red}</style><div>You've been blocked by network security. To continue, log in</div>"#;
+ let msg = truncate_for_error(html);
+ assert!(msg.starts_with("You've been blocked"));
+ assert!(msg.len() < 200);
+ assert!(!msg.contains(".x{color"));
+ }
+
#[test]
fn map_subreddit_about_url() {
let id = ItemId::from_url("https://reddit.com/r/rust").unwrap();
Hardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.