use axum::Extension;
use axum::Json;
use axum::extract::Query;
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use sqlx::QueryBuilder;
use sqlx::{Pool, Sqlite};
use crate::xrpc_gateway::XrpcAuthClaims;
use super::XrpcGatewayState;
use super::projections::subject_status::{
ProjectedSubjectStatus, SubjectActionRow, project_subject_status,
};
const MAX_LIMIT: u32 = 100;
const DEFAULT_LIMIT: u32 = 50;
#[derive(Debug, Deserialize, Default)]
#[serde(rename_all = "camelCase")]
pub struct QueryStatusesParams {
#[serde(default)]
pub subject: Option<String>,
#[serde(default)]
pub comment: Option<String>,
#[serde(default)]
pub reported_after: Option<String>,
#[serde(default)]
pub reported_before: Option<String>,
#[serde(default)]
pub review_state: Option<String>,
#[serde(default)]
pub ignore_subjects: Option<String>,
#[serde(default)]
pub last_reviewed_by: Option<String>,
#[serde(default)]
pub sort_field: Option<String>,
#[serde(default)]
pub sort_direction: Option<String>,
#[serde(default)]
pub takendown: Option<bool>,
#[serde(default)]
pub appealed: Option<bool>,
#[serde(default)]
pub limit: Option<u32>,
#[serde(default)]
pub cursor: Option<String>,
#[serde(default)]
pub tags: Option<String>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct QueryStatusesResponse {
#[serde(skip_serializing_if = "Option::is_none")]
pub cursor: Option<String>,
pub subject_statuses: Vec<SubjectStatusView>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SubjectStatusView {
pub id: i64,
pub subject: Value,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub subject_blob_cids: Vec<String>,
pub updated_at: String,
pub created_at: String,
pub review_state: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub comment: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_reviewed_by: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_reviewed_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_reported_at: Option<String>,
pub takendown: bool,
pub appealed: bool,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub tags: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Cursor {
updated_at_ms: i64,
subject_did: String,
subject_uri: String,
}
impl Cursor {
fn encode(&self) -> String {
let json = serde_json::to_vec(self).expect("cursor serializes");
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(json)
}
fn decode(s: &str) -> Result<Self, String> {
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(s)
.map_err(|e| format!("malformed cursor (base64): {e}"))?;
serde_json::from_slice(&bytes).map_err(|e| format!("malformed cursor (json): {e}"))
}
}
pub(crate) async fn handler(
Extension(state): Extension<XrpcGatewayState>,
Extension(_claims): Extension<XrpcAuthClaims>,
Query(params): Query<QueryStatusesParams>,
) -> Response {
let parsed = match parse_params(params) {
Ok(p) => p,
Err(msg) => return invalid_request(msg),
};
if parsed.appealed == Some(true) {
return Json(QueryStatusesResponse {
cursor: None,
subject_statuses: Vec::new(),
})
.into_response();
}
let page = match fetch_page(&state.pool, &parsed).await {
Ok(p) => p,
Err(e) => {
tracing::error!(error = %e, "xrpc_gateway queryStatuses: page query failed");
return internal_server_error();
}
};
let mut views = Vec::with_capacity(page.rows.len());
for row in &page.rows {
let actions = match fetch_actions_for_subject(
&state.pool,
&row.subject_did,
row.subject_uri.as_deref(),
)
.await
{
Ok(a) => a,
Err(e) => {
tracing::error!(
error = %e,
subject = row.subject_did,
"xrpc_gateway queryStatuses: per-subject actions query failed"
);
return internal_server_error();
}
};
let tags = match load_active_label_vals(
&state.pool,
&row.subject_did,
row.subject_uri.as_deref(),
&state.service_did,
)
.await
{
Ok(t) => t,
Err(e) => {
tracing::error!(error = %e, "xrpc_gateway queryStatuses: tags query failed");
return internal_server_error();
}
};
let last_reported_at =
match load_last_reported_at(&state.pool, &row.subject_did, row.subject_uri.as_deref())
.await
{
Ok(t) => t,
Err(e) => {
tracing::error!(
error = %e,
"xrpc_gateway queryStatuses: last_reported_at query failed"
);
return internal_server_error();
}
};
if let Some(tag_filter) = &parsed.tags
&& !tag_filter.is_empty()
&& !tags.iter().any(|t| tag_filter.iter().any(|f| f == t))
{
continue;
}
let projected = project_subject_status(
&row.subject_did,
row.subject_uri.as_deref(),
&actions,
tags,
last_reported_at,
);
views.push(serialize_view(projected));
}
let cursor = if page.has_more {
page.rows.last().map(|r| {
Cursor {
updated_at_ms: r.latest_created_at,
subject_did: r.subject_did.clone(),
subject_uri: r.subject_uri.clone().unwrap_or_default(),
}
.encode()
})
} else {
None
};
Json(QueryStatusesResponse {
cursor,
subject_statuses: views,
})
.into_response()
}
#[derive(Debug)]
struct ParsedParams {
subject_did: Option<String>,
subject_uri: Option<String>,
limit: u32,
cursor: Option<Cursor>,
descending: bool,
takendown: Option<bool>,
appealed: Option<bool>,
tags: Option<Vec<String>>,
}
fn parse_params(p: QueryStatusesParams) -> Result<ParsedParams, String> {
if p.comment.is_some() {
return Err("filter 'comment' is not supported in v1.7".to_string());
}
if p.reported_after.is_some() {
return Err("filter 'reportedAfter' is not supported in v1.7".to_string());
}
if p.reported_before.is_some() {
return Err("filter 'reportedBefore' is not supported in v1.7".to_string());
}
if p.review_state.is_some() {
return Err("filter 'reviewState' is not supported in v1.7".to_string());
}
if p.ignore_subjects.is_some() {
return Err("filter 'ignoreSubjects' is not supported in v1.7".to_string());
}
if p.last_reviewed_by.is_some() {
return Err("filter 'lastReviewedBy' is not supported in v1.7".to_string());
}
if let Some(sf) = &p.sort_field
&& sf != "lastReviewedAt"
{
return Err(format!(
"sortField '{sf}' is not supported in v1.7; only 'lastReviewedAt' is accepted"
));
}
let descending = match p.sort_direction.as_deref() {
None | Some("desc") => true,
Some("asc") => false,
Some(other) => {
return Err(format!(
"sortDirection '{other}' is invalid; expected 'asc' or 'desc'"
));
}
};
let limit = match p.limit {
None => DEFAULT_LIMIT,
Some(0) => return Err("limit must be at least 1".to_string()),
Some(n) if n > MAX_LIMIT => {
return Err(format!("limit {n} exceeds maximum {MAX_LIMIT}"));
}
Some(n) => n,
};
let cursor = match p.cursor.as_deref() {
Some(s) => Some(Cursor::decode(s)?),
None => None,
};
let (subject_did, subject_uri) = match p.subject.as_deref() {
None => (None, None),
Some(s) if s.starts_with("at://") => {
let did = s
.strip_prefix("at://")
.and_then(|rest| rest.split('/').next())
.filter(|d| d.starts_with("did:"))
.ok_or_else(|| format!("subject AT-URI {s:?} missing DID authority"))?
.to_string();
(Some(did), Some(s.to_string()))
}
Some(s) if s.starts_with("did:") => (Some(s.to_string()), None),
Some(s) => return Err(format!("subject {s:?} is neither a DID nor an AT-URI")),
};
let tags = match p.tags.as_deref() {
None => None,
Some(s) => {
let v: Vec<String> = s
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.map(String::from)
.collect();
if v.is_empty() { None } else { Some(v) }
}
};
Ok(ParsedParams {
subject_did,
subject_uri,
limit,
cursor,
descending,
takendown: p.takendown,
appealed: p.appealed,
tags,
})
}
#[derive(Debug)]
struct PageRow {
subject_did: String,
subject_uri: Option<String>,
latest_created_at: i64,
}
#[derive(Debug)]
struct Page {
rows: Vec<PageRow>,
has_more: bool,
}
async fn fetch_page(pool: &Pool<Sqlite>, parsed: &ParsedParams) -> sqlx::Result<Page> {
let mut qb: QueryBuilder<Sqlite> = QueryBuilder::new(
r#"SELECT
subject_did,
subject_uri,
MAX(created_at) AS latest_created_at
FROM subject_actions
WHERE 1=1"#,
);
if let Some(did) = &parsed.subject_did {
qb.push(" AND subject_did = ").push_bind(did.clone());
}
match parsed.subject_uri.as_deref() {
Some(uri) => {
qb.push(" AND subject_uri = ").push_bind(uri.to_string());
}
None if parsed.subject_did.is_some() => {
qb.push(" AND subject_uri IS NULL");
}
None => {}
}
qb.push(
r#"
GROUP BY subject_did, subject_uri
HAVING 1=1"#,
);
if let Some(td) = parsed.takendown {
if td {
qb.push(
r#" AND EXISTS (
SELECT 1 FROM subject_actions sa2
WHERE sa2.subject_did = subject_actions.subject_did
AND (sa2.subject_uri IS subject_actions.subject_uri
OR (sa2.subject_uri IS NULL AND subject_actions.subject_uri IS NULL))
AND sa2.action_type = 'takedown'
AND sa2.revoked_at IS NULL
)"#,
);
} else {
qb.push(
r#" AND NOT EXISTS (
SELECT 1 FROM subject_actions sa2
WHERE sa2.subject_did = subject_actions.subject_did
AND (sa2.subject_uri IS subject_actions.subject_uri
OR (sa2.subject_uri IS NULL AND subject_actions.subject_uri IS NULL))
AND sa2.action_type = 'takedown'
AND sa2.revoked_at IS NULL
)"#,
);
}
}
if let Some(c) = &parsed.cursor {
let comparator = if parsed.descending { "<" } else { ">" };
qb.push(format!(
" AND (MAX(created_at), subject_did, COALESCE(subject_uri, '')) {comparator} ("
));
qb.push_bind(c.updated_at_ms);
qb.push(", ");
qb.push_bind(c.subject_did.clone());
qb.push(", ");
qb.push_bind(c.subject_uri.clone());
qb.push(")");
}
let order = if parsed.descending { "DESC" } else { "ASC" };
qb.push(format!(
" ORDER BY latest_created_at {order}, subject_did {order}, subject_uri {order}"
));
qb.push(" LIMIT ").push_bind((parsed.limit + 1) as i64);
let mut rows: Vec<PageRow> = qb
.build_query_as::<(String, Option<String>, i64)>()
.fetch_all(pool)
.await?
.into_iter()
.map(|(subject_did, subject_uri, latest_created_at)| PageRow {
subject_did,
subject_uri,
latest_created_at,
})
.collect();
let has_more = rows.len() > parsed.limit as usize;
if has_more {
rows.truncate(parsed.limit as usize);
}
Ok(Page { rows, has_more })
}
async fn fetch_actions_for_subject(
pool: &Pool<Sqlite>,
subject_did: &str,
subject_uri: Option<&str>,
) -> sqlx::Result<Vec<SubjectActionRow>> {
let rows = match subject_uri {
Some(uri) => sqlx::query!(
r#"SELECT
id AS "id!: i64",
action_type AS "action_type!: String",
actor_did AS "actor_did!: String",
notes,
revoked_at,
created_at AS "created_at!: i64"
FROM subject_actions
WHERE subject_did = ?1 AND subject_uri = ?2
ORDER BY id ASC"#,
subject_did,
uri,
)
.fetch_all(pool)
.await?
.into_iter()
.map(|r| SubjectActionRow {
id: r.id,
action_type: r.action_type,
actor_did: r.actor_did,
notes: r.notes,
revoked_at: r.revoked_at,
created_at: r.created_at,
})
.collect(),
None => sqlx::query!(
r#"SELECT
id AS "id!: i64",
action_type AS "action_type!: String",
actor_did AS "actor_did!: String",
notes,
revoked_at,
created_at AS "created_at!: i64"
FROM subject_actions
WHERE subject_did = ?1 AND subject_uri IS NULL
ORDER BY id ASC"#,
subject_did,
)
.fetch_all(pool)
.await?
.into_iter()
.map(|r| SubjectActionRow {
id: r.id,
action_type: r.action_type,
actor_did: r.actor_did,
notes: r.notes,
revoked_at: r.revoked_at,
created_at: r.created_at,
})
.collect(),
};
Ok(rows)
}
async fn load_active_label_vals(
pool: &Pool<Sqlite>,
subject_did: &str,
subject_uri: Option<&str>,
service_did: &str,
) -> sqlx::Result<Vec<String>> {
let label_uri = subject_uri.unwrap_or(subject_did);
let rows = sqlx::query!(
r#"SELECT DISTINCT l1.val AS "val!: String"
FROM labels l1
WHERE l1.src = ?1 AND l1.uri = ?2
AND l1.seq = (
SELECT MAX(l2.seq) FROM labels l2
WHERE l2.src = l1.src AND l2.uri = l1.uri AND l2.val = l1.val
)
AND l1.neg = 0
ORDER BY l1.val ASC"#,
service_did,
label_uri,
)
.fetch_all(pool)
.await?;
Ok(rows.into_iter().map(|r| r.val).collect())
}
async fn load_last_reported_at(
pool: &Pool<Sqlite>,
subject_did: &str,
subject_uri: Option<&str>,
) -> sqlx::Result<Option<String>> {
let row = match subject_uri {
Some(uri) => sqlx::query!(
r#"SELECT MAX(created_at) AS "last_reported_at: String"
FROM reports
WHERE subject_did = ?1 AND subject_uri = ?2"#,
subject_did,
uri,
)
.fetch_optional(pool)
.await?
.and_then(|r| r.last_reported_at),
None => sqlx::query!(
r#"SELECT MAX(created_at) AS "last_reported_at: String"
FROM reports
WHERE subject_did = ?1 AND subject_uri IS NULL"#,
subject_did,
)
.fetch_optional(pool)
.await?
.and_then(|r| r.last_reported_at),
};
Ok(row)
}
fn serialize_view(p: ProjectedSubjectStatus) -> SubjectStatusView {
SubjectStatusView {
id: p.id,
subject: p.subject,
subject_blob_cids: Vec::new(),
updated_at: epoch_ms_to_rfc3339(p.updated_at),
created_at: epoch_ms_to_rfc3339(p.created_at),
review_state: p.review_state,
comment: p.comment,
last_reviewed_by: p.last_reviewed_by,
last_reviewed_at: p.last_reviewed_at.map(epoch_ms_to_rfc3339),
last_reported_at: p.last_reported_at,
takendown: p.takendown,
appealed: p.appealed,
tags: p.tags,
}
}
fn epoch_ms_to_rfc3339(ms: i64) -> String {
crate::writer::rfc3339_from_epoch_ms(ms)
.unwrap_or_else(|_| String::from("1970-01-01T00:00:00.000Z"))
}
fn invalid_request(message: String) -> Response {
let body = ErrorEnvelope {
error: "InvalidRequest",
message,
};
(StatusCode::BAD_REQUEST, Json(body)).into_response()
}
fn internal_server_error() -> Response {
let body = ErrorEnvelope {
error: "InternalServerError",
message: "service temporarily unavailable".to_string(),
};
(StatusCode::INTERNAL_SERVER_ERROR, Json(body)).into_response()
}
#[derive(Serialize)]
struct ErrorEnvelope {
error: &'static str,
message: String,
}
#[cfg(test)]
mod tests {
use super::*;
fn empty_params() -> QueryStatusesParams {
QueryStatusesParams::default()
}
#[test]
fn unsupported_comment_filter_rejected() {
let p = QueryStatusesParams {
comment: Some("foo".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("comment"), "{err}");
}
#[test]
fn unsupported_reported_after_filter_rejected() {
let p = QueryStatusesParams {
reported_after: Some("2026-01-01".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("reportedAfter"), "{err}");
}
#[test]
fn unsupported_review_state_filter_rejected() {
let p = QueryStatusesParams {
review_state: Some("open".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("reviewState"), "{err}");
}
#[test]
fn invalid_sort_field_rejected() {
let p = QueryStatusesParams {
sort_field: Some("createdAt".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("sortField"), "{err}");
}
#[test]
fn invalid_sort_direction_rejected() {
let p = QueryStatusesParams {
sort_direction: Some("sideways".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("sortDirection"), "{err}");
}
#[test]
fn default_limit_is_50() {
let p = parse_params(empty_params()).unwrap();
assert_eq!(p.limit, 50);
}
#[test]
fn limit_capped_at_100() {
let p = QueryStatusesParams {
limit: Some(200),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("100"), "{err}");
}
#[test]
fn limit_zero_rejected() {
let p = QueryStatusesParams {
limit: Some(0),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("at least 1"), "{err}");
}
#[test]
fn account_level_subject_parses_to_did_only() {
let p = QueryStatusesParams {
subject: Some("did:plc:abc".into()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert_eq!(parsed.subject_did.as_deref(), Some("did:plc:abc"));
assert!(parsed.subject_uri.is_none());
}
#[test]
fn record_level_subject_parses_to_did_plus_uri() {
let p = QueryStatusesParams {
subject: Some("at://did:plc:author/c/r".into()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert_eq!(parsed.subject_did.as_deref(), Some("did:plc:author"));
assert_eq!(
parsed.subject_uri.as_deref(),
Some("at://did:plc:author/c/r")
);
}
#[test]
fn malformed_subject_rejected() {
let p = QueryStatusesParams {
subject: Some("https://example.com".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("DID"), "{err}");
}
#[test]
fn at_uri_without_did_authority_rejected() {
let p = QueryStatusesParams {
subject: Some("at://example.com/c/r".into()),
..empty_params()
};
let err = parse_params(p).unwrap_err();
assert!(err.contains("DID authority"), "{err}");
}
#[test]
fn tags_split_on_comma() {
let p = QueryStatusesParams {
tags: Some("spam,harassment,nsfw".into()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert_eq!(
parsed.tags.as_deref(),
Some(
&[
"spam".to_string(),
"harassment".to_string(),
"nsfw".to_string()
][..]
)
);
}
#[test]
fn empty_tags_string_yields_none() {
let p = QueryStatusesParams {
tags: Some("".into()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert!(parsed.tags.is_none());
}
#[test]
fn tags_with_whitespace_trimmed() {
let p = QueryStatusesParams {
tags: Some("spam, harassment ,nsfw".into()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert_eq!(
parsed.tags.as_deref(),
Some(
&[
"spam".to_string(),
"harassment".to_string(),
"nsfw".to_string()
][..]
)
);
}
#[test]
fn cursor_round_trip_is_stable() {
let c = Cursor {
updated_at_ms: 1_700_000_000_000,
subject_did: "did:plc:abc".into(),
subject_uri: "at://did:plc:abc/c/r".into(),
};
let encoded = c.encode();
let decoded = Cursor::decode(&encoded).unwrap();
assert_eq!(decoded.updated_at_ms, c.updated_at_ms);
assert_eq!(decoded.subject_did, c.subject_did);
assert_eq!(decoded.subject_uri, c.subject_uri);
}
#[test]
fn cursor_round_trip_with_empty_uri() {
let c = Cursor {
updated_at_ms: 1_700_000_000_000,
subject_did: "did:plc:abc".into(),
subject_uri: "".into(),
};
let encoded = c.encode();
let decoded = Cursor::decode(&encoded).unwrap();
assert_eq!(decoded.subject_uri, "");
}
#[test]
fn malformed_cursor_rejected() {
let err = Cursor::decode("not-base64-!!!").unwrap_err();
assert!(err.contains("base64") || err.contains("malformed"), "{err}");
}
#[test]
fn cursor_with_invalid_json_payload_rejected() {
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(b"not-json");
let err = Cursor::decode(&bytes).unwrap_err();
assert!(err.contains("json") || err.contains("malformed"), "{err}");
}
#[test]
fn cursor_round_trip_via_param_parsing() {
let c = Cursor {
updated_at_ms: 1_700_000_000_000,
subject_did: "did:plc:abc".into(),
subject_uri: "".into(),
};
let p = QueryStatusesParams {
cursor: Some(c.encode()),
..empty_params()
};
let parsed = parse_params(p).unwrap();
assert_eq!(
parsed.cursor.as_ref().map(|c| c.updated_at_ms),
Some(1_700_000_000_000)
);
}
}