Side B fixes a concrete correctness/security bug (feed catch-up using ingest position instead of timestamp, preventing missed posts on clock rollback and leaked private-room content after permission changes) and backs it with thorough multi-user integration tests. Side A is a solid internal-tooling refactor (per-commit vs per-author ranking with rollup) but addresses a fairness/design limitation rather than a user-facing correctness bug, making B's impact more concrete and durable.
constitution · epochs · watch · epoch 3
c_fbeec5c4ad18 (tommy-mor) vs c_0a9a8eab32ba (tommy-mor)
download prompt · raw event · cmp_cfca245d62de28
council reasoning
B fixes real feed correctness and privacy failures by anchoring catch-up to durable ingest index (not colliding/backwards timestamps), centralizing logic in rpc_feed, and adding multi-user private-room/revoke regression tests plus a room field on posts. A is a solid redesign—pairwise-rank each commit and roll scores up, with matching UI/evidence/tests—but it is mostly an internal ranking granularity change rather than closing missed-update and leak bugs.
Side B fixes substantive feed correctness by introducing durable ingest-order anchoring instead of timestamp-only cutoffs, preserving permission checks, adding room metadata to feed responses, and covering multi-user/private-room edge cases with comprehensive integration tests. Side A mainly refactors ranking from contributor-level to commit-level with UI/evidence updates and tests, which is a meaningful feature but less fundamental than correcting feed consistency and visibility bugs.
sides
A — c_fbeec5c4ad18 (tommy-mor)
message
[c7ef287e] Rank every eligible commit with the LLM council. Stop short-circuiting on a single contributor; pairwise-sort commits, roll scores up for emission payouts, and surface commit rankings on epoch pages. Co-authored-by: Cursor <cursoragent@cursor.com>
diff preview
diff --git a/constitution.py b/constitution.py
index 26ba130e885e0e69fb7874ca5c3f07f42100a150..71fc46b4c9a4d7a9ea7bf319860b0db1acc4a673 100644
--- a/constitution.py
+++ b/constitution.py
@@ -503,20 +503,22 @@ def _epochs_in_ledger() -> list[int]:
def build_pairwise_prompt(side_a: dict, side_b: dict) -> str:
- return f"""You are ranking contributions to an open source project.
-Compare these two sides (each may be one or more commits). Decide which side contributed more.
+ return f"""You are ranking individual git commits to an open source project.
+Compare these two commits. Decide which commit contributed more.
Return ONLY a JSON object: {{"winner": "A" or "B", "ratio": "N:M", "explanation": "..."}}
-Side A — commit messages:
+Side A — contributor: {side_a.get('contributor', '?')}
+Side A — commit message:
{side_a['message']}
-Side A — unified diffs (full patches):
+Side A — unified diff (full patch):
{side_a['diff']}
-Side B — commit messages:
+Side B — contributor: {side_b.get('contributor', '?')}
+Side B — commit message:
{side_b['message']}
-Side B — unified diffs (full patches):
+Side B — unified diff (full patch):
{side_b['diff']}"""
@@ -1518,16 +1520,28 @@ async def broadcast_js(js: str):
await queue.put(js)
-def _author_side_for_llm(author: str, author_commits: dict) -> dict:
- cs = author_commits[author]
+def _commit_side_for_llm(row: dict) -> dict:
+ oid = row["oid"]
+ short = oid.split(":", 1)[1][:8] if ":" in oid else oid[:8]
return {
- "message": "\n".join(f"[{c['sha']}] {c['message']}" for c in cs),
- "diff": "\n\n".join(f"=== {c['sha']} ===\n{c['diff']}" for c in cs),
- "commit_ids": [c["commit_id"] for c in cs],
- "contributor": author,
+ "message": f"[{short}] {row['message']}",
+ "diff": row["patch"] or "",
+ "commit_id": commit_id_for_oid(oid),
+ "contributor": row["contributor"],
+ "oid": oid,
}
+def _rollup_contributor_scores(
+ ordered: list[dict], commit_scores: list[Decimal]
+) -> dict[str, Decimal]:
+ totals: dict[str, Decimal] = {}
+ for row, score in zip(ordered, commit_scores):
+ contributor = row["contributor"]
+ totals[contributor] = totals.get(contributor, Decimal("0")) + score
+ return totals
+
+
def _find_judgment(comparison_id: str, model_id: str) -> dict | None:
for e in evidence_by_kind("llm.judgment"):
p = e.payload
@@ -1546,101 +1560,106 @@ def _find_ranking_models(ranking_run_id: str) -> list[str] | None:
async def rank_commits(commits: list[dict], *, epoch: int = -1):
+ """Pairwise-rank every eligible commit; roll scores up to contributors."""
if not commits:
return {}, [], {"ranking_run_id": "", "ranking_event_id": ""}
- commit_ids = sorted(commit_id_for_oid(row["oid"]) for row in commits)
+ ordered = sorted(commits, key=lambda r: r["oid"])
+ commit_ids = [commit_id_for_oid(row["oid"]) for row in ordered]
ranking_run_id = _content_id("rank", {
"epoch": epoch,
- "commit_ids": commit_ids,
+ "commit_ids": sorted(commit_ids),
})
- contributors = sorted(set(c["contributor"] for c in commits))
+ contributors = sorted({c["contributor"] for c in ordered})
- if len(contributors) == 1:
+ # Nothing to compare: a single commit (not a single contributor).
+ if len(ordered) == 1:
await append_evidence(epoch, "ranking.started", {
"ranking_run_id": ranking_run_id,
"commit_ids": commit_ids,
"contributors": contributors,
"models": [],
- "summary": f"ranking epoch {epoch}: single contributor",
+ "summary": f"ranking epoch {epoch}: single commit",
})
- ranking = {contributors[0]: Decimal("1")}
+ commit_ranking = {commit_ids[0]: "1"}
+ contributor_ranking = {ordered[0]["contributor"]: Decimal("1")}
completed = await append_evidence(epoch, "ranking.completed", {
"ranking_run_id": ranking_run_id,
"models": [],
- "ranking": {contributors[0]: "1"},
+ "commit_ranking": commit_ranking,
+ "contributor_ranking": {ordered[0]["contributor"]: "1"},
+ "ranking": {ordered[0]["contributor"]: "1"},
"judgment_ids": [],
- "summary": f"Only {contributors[0]} is eligible; rank is 1.0",
+ "summary": f"Only one eligible commit; {ordered[0]['contributor']} rank 1.0",
})
await broadcast_audit(
"ranking",
- f"Only {contributors[0]} is eligible; rank is 1.0",
+ f"Only one eligible commit; {ordered[0]['contributor']} rank 1.0",
progress=90,
phase="finalizing",
evidence_event_id=completed.event_id,
evidence_url=_evidence_url("event", completed.event_id),
links={"epoch": _evidence_url("epoch", str(epoch))},
)
- return ranking, [], {
+ return contributor_ranking, [], {
"ranking_run_id": ranking_run_id,
"ranking_event_id": completed.event_id,
}
if not (OPENROUTER_API_KEY or "").strip():
raise RuntimeError(
- "OPENROUTER_API_KEY is required when multiple contributors need ranking"
+ "OPENROUTER_API_KEY is required when multiple commits need ranking"
)
models = _find_ranking_models(ranking_run_id)
if models is None:
models = await fetch_top_models(n=3)
if not models:
- raise RuntimeError("no council models available for contributor ranking")
+ raise RuntimeError("no council models available for commit ranking")
await append_evidence(epoch, "ranking.started", {
"ranking_run_id": ranking_run_id,
"commit_ids": commit_ids,
"contributors": contributors,
"models": models,
- "summary": f"Council selected: {', '.join(models)}",
+ "summary": (
+ f"Council selected: {', '.join(models)} — "
+ f"{len(ordered)} commits"
+ ),
})
await broadcast_audit(
"council",
- f"Council selected: {', '.join(models)}",
+ f"Council selected: {', '.join(models)} — ranking {len(ordered)} commits",
progress=35,
phase="ranking",
)
await broadcast_js(exec_event(Three[Selector("#emission-log")][PREPEND][
- ["div.log-council", f"Council: {', '.join(models)} — {len(commits)} commits"]
+ ["div.log-council",
+ f"Council: {', '.join(models)} — {len(ordered)} commits"]
]))
- authors = contributors
- author_commits = {a: [] for a in authors}
- for row in sorted(commits, key=lambda r: r["oid"]):
- author_commits[row["contributor"]].append({
- "message": row["message"],
- "sha": row["oid"].split(":", 1)[1][:8],
- "diff": row["patch"],
- "commit_id": commit_id_for_oid(row["oid"]),
- })
-
+ sides = [_commit_side_for_llm(row) for row in ordered]
judgment_ids: list[str] = []
async def compare_fn(i, j):
- a1, a2 = authors[i], authors[j]
- side_a = _author_side_for_llm(a1, author_commits)
- side_b = _author_side_for_llm(a2, author_commits)
+ side_a, side_b = sides[i], sides[j]
+ label_a = f"{side_a['commit_id'][:16]} ({side_a['contributor']})"
+ label_b = f"{side_b['commit_id'][:16]} ({side_b['contributor']})"
prompt = build_pairwise_prompt(side_a, side_b)
comparison_material = {
"ranking_run_id": ranking_run_id,
"side_a": {
- "contributor": a1,
- "commit_ids": side_a["commit_ids"],
+ "contributor": side_a["contributor"],
+ "commit_id": side_a["commit_id"],
+ "commit_ids": [side_a["commit_id"]],
+ "oid": side_a["oid"],
"message": _bytes_blob(side_a["message"]),
"diff": _bytes_blob(side_a["diff"]),
},
"side_b": {
- "contributor": a2,
- "commit_ids": side_b["commit_ids"],
+ "contributor": side_b["contributor"],
+ "commit_id": side_b["commit_id"],
+ "commit_ids": [side_b["commit_id"]],
+ "oid": side_b["oid"],
"message": _bytes_blob(side_b["message"]),
"diff": _bytes_blob(side_b["diff"]),
},
@@ -1650,22 +1669,24 @@ async def rank_commits(commits: list[dict], *, epoch: int = -1):
comparison_material = {
**comparison_material,
"comparison_id": comparison_id,
- "summary": f"Comparing {a1} with {a2}",
+ "summary": f"Comparing {label_a} with {label_b}",
}
cmp_ev = await append_evidence(epoch, "comparison.input", comparison_material)
await broadcast_audit(
"comparison",
- f"Comparing {a1} with {a2}",
+ f"Comparing commits {label_a} vs {label_b}",
phase="ranking",
evidence_event_id=cmp_ev.event_id,
evidence_url=_evidence_url("comparison", comparison_id),
links={
"comparison": _evidence_url("comparison", comparison_id),
+ "commit_a": _evidence_url("commit", side_a["commit_id"]),
+ "commit_b": _evidence_url("commit", side_b["commit_id"]),
"epoch": _evidence_url("epoch", str(epoch)),
},
)
await broadcast_js(exec_event(Three[Selector("#emission-status")][MORPH][
- ["div#emission-status", f"Comparing {a1} vs {a2}…"]
+ ["div#emission-status", f"Comparing {label_a} vs {label_b}…"]
]))
results = []
for model in models:
@@ -1697,9 +1718,15 @@ async def rank_commits(commits: list[dict], *, epoch: int = -1):
jud_id = (existing or _find_judgment(comparison_id, model) or {}).get(
"judgment_id"
)
+ win_label = (
+ f"{sides[w]['commit_id'][:16]} ({sides[w]['contributor']})"
+ )
+ lose_label = (
+ f"{sides[l]['commit_id'][:16]} ({sides[l]['contributor']})"
+ )
await broadcast_audit(
"vote",
- f"{model}: {authors[w]} over {authors[l]} ({result['ratio']})",
+ f"{model}: {win_label} over {lose_label} ({result['ratio']})",
phase="ranking",
evidence_url=(
_evidence_url("judgment", jud_id) if jud_id else None
@@ -1714,8 +1741,8 @@ async def rank_commits(commits: list[dict], *, epoch: int = -1):
await broadcast_js(exec_event(Three[Selector("#emission-log")][PREPEND][
["div.log-vote",
["span.model", model], " — ",
- ["span.winner", authors[w]], f" beat ",
- ["span.loser", authors[l]], f" ({result['ratio']}) ",
+ ["span.winner", win_label], f" beat ",
+ ["span.loser", lose_label], f" ({result['ratio']}) ",
["span.explanation", result["explanation"]],
]
]))
@@ -1745,24 +1772,46 @@ async def rank_commits(commits: list[dict], *, epoch: int = -1):
["div#emission-status", label]
]))
- pairs = await pairwise_rank(len(authors), compare_fn, progress_fn)
+ pairs = await pairwise_rank(len(ordered), compare_fn, progress_fn)
if not pairs:
- ranking = {authors[0]: Decimal("1")} if authors else {}
+ commit_score_list = [Decimal("1")]
else:
scores = rank_centrality(pairs)
- ranking = {authors[i]: Decimal(str(scores[i])) for i in range(len(authors))}
… preview truncated; 12,453 characters omittedB — c_0a9a8eab32ba (tommy-mor)
message
[c94456ff] Make feed catch-up stable and permission-aware Anchor implicit feeds to durable ingest order and cover multi-user private-room visibility so concurrent posts are not missed or leaked. Co-authored-by: Cursor <cursoragent@cursor.com>
diff preview
diff --git a/cli/src/main.rs b/cli/src/main.rs
index abb5a55b49f60fe28fbfd4ec02715cb94ea0b4ec..c4f1494df3aedbd8b883aea6249579aa8336abfe 100644
--- a/cli/src/main.rs
+++ b/cli/src/main.rs
@@ -878,6 +878,30 @@ mod tests {
"graph: 4 items, 3/6 pairs (50.0% density), 1 component, connected"
);
}
+
+ #[test]
+ fn feed_without_since_uses_logged_in_delegate_from_env() {
+ let key = "SLUG_DELEGATE";
+ let previous = std::env::var_os(key);
+ let expected = "00000000-0000-0000-0000-0000000000ee:test:local/model";
+ std::env::set_var(key, expected);
+
+ let cli = Cli::try_parse_from(["slugsocial", "feed"]).expect("parse feed");
+
+ match previous {
+ Some(value) => std::env::set_var(key, value),
+ None => std::env::remove_var(key),
+ }
+ match cli.cmd {
+ Some(Command::Feed {
+ delegate, since, ..
+ }) => {
+ assert_eq!(delegate.as_deref(), Some(expected));
+ assert!(since.is_none());
+ }
+ _ => panic!("expected feed command"),
+ }
+ }
}
async fn run_scoped(base: &str, room: &str, sub: ScopedCmd) -> Result<()> {
@@ -1626,7 +1650,15 @@ async fn run() -> Result<()> {
} else {
for p in &resp.posts {
let ago = slug_types::timeago::timeago(now_ms, p.ts);
- println!("<post id=\"{}\" ts=\"{}\">", p.id, ago);
+ let thread_attr = p
+ .thread
+ .as_deref()
+ .map(|thread| format!(" thread=\"{thread}\""))
+ .unwrap_or_default();
+ println!(
+ "<post id=\"{}\" ts=\"{}\" room=\"{}\"{}>",
+ p.id, ago, p.room, thread_attr
+ );
print!("{}", p.body);
if !p.body.ends_with('\n') { println!(); }
println!("</post>");
diff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs
index afc4f95bef160c1e38ecff2c096d6440cd94e2b3..46d748f918d9b225acc4ecedfe5a1793089407b5 100644
--- a/server/src/api/rpc.rs
+++ b/server/src/api/rpc.rs
@@ -72,6 +72,80 @@ fn can_view_scope(reduced: &ReducerState, scope: &ScopeId, principal: Option<&st
}
}
+/// Build a feed in durable ingest order.
+///
+/// An implicit feed boundary is an ingest position, not only its millisecond timestamp. Two users
+/// can post in the same millisecond, and wall-clock timestamps can move backwards during replay.
+/// Explicit `since` remains a timestamp query for API compatibility, but scans the whole ordered
+/// ledger rather than assuming timestamps are monotonic.
+fn rpc_feed(
+ reduced: &ReducerState,
+ viewer: &str,
+ delegate: Option<String>,
+ requested_since: Option<i64>,
+ implicit_anchor: Option<(usize, i64)>,
+ limit: usize,
+) -> FeedResponse {
+ let since = requested_since.or_else(|| implicit_anchor.map(|(_, ts)| ts));
+ let implicit_anchor_index = requested_since
+ .is_none()
+ .then(|| implicit_anchor.map(|(index, _)| index))
+ .flatten();
+
+ let matching: Vec<&str> = reduced
+ .ingests_ordered
+ .iter()
+ .enumerate()
+ .rev()
+ .filter(|(index, id)| {
+ reduced.ingests_by_id.get(id.as_str()).is_some_and(|ing| {
+ match requested_since {
+ Some(cutoff) => ing.ts > cutoff,
+ None => implicit_anchor_index.is_none_or(|anchor| *index > anchor),
+ }
+ })
+ })
+ .map(|(_, id)| id.as_str())
+ .filter(|id| {
+ reduced.ingests_by_id.get(*id).is_some_and(|ing| {
+ let scope = scope_from_room_wire(&ing.room_id);
+ can_view_scope(reduced, &scope, Some(viewer))
+ })
+ })
+ .filter(|id| !reduced.redacted_posts.contains(*id))
+ .collect();
+
+ let total = matching.len();
+ let posts = matching
+ .into_iter()
+ .take(limit)
+ .filter_map(|id| reduced.ingests_by_id.get(id))
+ .map(|ing| {
+ let scope = scope_from_room_wire(&ing.room_id);
+ let thread_post_index = reduced.try_thread_post_index_chronological(
+ &scope,
+ &ing.thread_tag,
+ &ing.id,
+ );
+ FeedPost {
+ ts: ing.ts,
+ id: ing.id.clone(),
+ room: ing.room_id.clone(),
+ thread: Some(ing.thread_tag.clone()),
+ thread_post_index,
+ body: ing.raw.clone(),
+ }
+ })
+ .collect();
+
+ FeedResponse {
+ delegate,
+ since,
+ posts,
+ total,
+ }
+}
+
fn principal_from_optional_bearer(headers: &HeaderMap, reduced: &ReducerState) -> Result<Option<String>, RpcErr> {
if headers.contains_key(axum::http::header::AUTHORIZATION) {
verify_bearer_principal(headers, reduced)
@@ -1506,60 +1580,27 @@ pub async fn handle_rpc_batch(
Some("this delegate is not bound to your signed-in account".into()),
)
} else {
- let since_default = reduced
+ let implicit_anchor = reduced
.ingests_ordered
.iter()
+ .enumerate()
.rev()
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .find(|ing| {
- if ing.delegate.as_deref() != Some(delegate_stored.as_str()) {
- return false;
- }
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
- })
- .map(|ing| ing.ts);
- let since = since.or(since_default);
- let cutoff = since.unwrap_or(0);
- let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
- let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()
- .map(|id| id.as_str())
- .take_while(|id| reduced.ingests_by_id.get(*id).is_some_and(|ing| ing.ts > cutoff))
- .filter(|id| {
- reduced.ingests_by_id.get(*id).is_some_and(|ing| {
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
+ .find_map(|(index, id)| {
+ reduced.ingests_by_id.get(id).and_then(|ing| {
+ (ing.delegate.as_deref()
+ == Some(delegate_stored.as_str()))
+ .then_some((index, ing.ts))
})
- })
- .filter(|id| !reduced.redacted_posts.contains(*id))
- .collect();
- let total = matching.len();
- let posts: Vec<FeedPost> = matching.into_iter()
- .take(limit)
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .map(|ing| {
- let scope = scope_from_room_wire(&ing.room_id);
- let thread_post_index = reduced
- .try_thread_post_index_chronological(
- &scope,
- &ing.thread_tag,
- &ing.id,
- );
- FeedPost {
- ts: ing.ts,
- id: ing.id.clone(),
- thread: Some(ing.thread_tag.clone()),
- thread_post_index,
- body: ing.raw.clone(),
- }
- })
- .collect();
- line_ok(RpcResult::Feed(FeedResponse {
- delegate: Some(delegate_stored),
+ });
+ let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
+ line_ok(RpcResult::Feed(rpc_feed(
+ &reduced,
+ &viewer,
+ Some(delegate_stored),
since,
- posts,
- total,
- }))
+ implicit_anchor,
+ limit,
+ )))
};
drop(reduced);
line
@@ -1567,60 +1608,25 @@ pub async fn handle_rpc_batch(
None => {
// Session catch-up: last time *you* posted anything (delegate or not), so revisiting
// an old chat with only a token still gets a sane cutoff.
- let since_default = reduced
+ let implicit_anchor = reduced
.ingests_ordered
.iter()
+ .enumerate()
.rev()
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .find(|ing| {
- if ing.principal != viewer {
- return false;
- }
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
- })
- .map(|ing| ing.ts);
- let since = since.or(since_default);
- let cutoff = since.unwrap_or(0);
- let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
- let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()
- .map(|id| id.as_str())
- .take_while(|id| reduced.inges
… preview truncated; 14,758 characters omittedHardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.