use axum::{
Json,
extract::{Path, Query, State},
http::StatusCode,
response::IntoResponse,
};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sqlx::PgPool;
#[derive(Clone)]
pub struct ChangelogState {
pub pool: PgPool,
}
#[derive(Debug, Deserialize)]
pub struct ChangelogQuery {
#[serde(default)]
pub after_cursor: i64,
#[serde(default = "default_changelog_limit")]
pub limit: i64,
pub object_type: Option<String>,
}
const fn default_changelog_limit() -> i64 {
100
}
const MAX_CHANGELOG_LIMIT: i64 = 1_000;
#[derive(Debug, Serialize)]
pub struct ChangelogEntryResponse {
pub cursor: i64,
pub id: String,
pub org_id: Option<String>,
pub user_id: Option<String>,
pub object_type: String,
pub object_id: String,
pub modification_type: String,
pub status: Option<String>,
pub object_data: serde_json::Value,
pub metadata: Option<serde_json::Value>,
pub created_at: Option<String>,
}
#[derive(Debug, Serialize)]
pub struct ChangelogListResponse {
pub entries: Vec<ChangelogEntryResponse>,
pub next_cursor: Option<i64>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct CheckpointResponse {
pub listener_id: String,
pub last_cursor: i64,
pub updated_at: Option<String>,
}
#[derive(Debug, Deserialize)]
pub struct SaveCheckpointRequest {
pub last_cursor: i64,
}
type ChangelogRow = (
i64, uuid::Uuid, Option<String>, Option<String>, String, String, String, Option<String>, serde_json::Value, Option<serde_json::Value>, Option<DateTime<Utc>>, );
pub async fn changelog_list_handler(
State(state): State<ChangelogState>,
Query(query): Query<ChangelogQuery>,
) -> impl IntoResponse {
let limit = query.limit.clamp(1, MAX_CHANGELOG_LIMIT);
let result = if let Some(ref object_type) = query.object_type {
sqlx::query_as::<_, ChangelogRow>(
r"
SELECT
pk_entity_change_log, id, fk_customer_org, fk_contact,
object_type, object_id, modification_type, change_status,
object_data, extra_metadata, created_at
FROM core.tb_entity_change_log
WHERE pk_entity_change_log > $1 AND object_type = $3
ORDER BY pk_entity_change_log ASC
LIMIT $2
",
)
.bind(query.after_cursor)
.bind(limit)
.bind(object_type)
.fetch_all(&state.pool)
.await
} else {
sqlx::query_as::<_, ChangelogRow>(
r"
SELECT
pk_entity_change_log, id, fk_customer_org, fk_contact,
object_type, object_id, modification_type, change_status,
object_data, extra_metadata, created_at
FROM core.tb_entity_change_log
WHERE pk_entity_change_log > $1
ORDER BY pk_entity_change_log ASC
LIMIT $2
",
)
.bind(query.after_cursor)
.bind(limit)
.fetch_all(&state.pool)
.await
};
match result {
Ok(rows) => {
let entries: Vec<ChangelogEntryResponse> = rows
.into_iter()
.map(
|(pk, id, org, contact, obj_type, obj_id, mod_type, status, data, meta, ts)| {
ChangelogEntryResponse {
cursor: pk,
id: id.to_string(),
org_id: org,
user_id: contact,
object_type: obj_type,
object_id: obj_id,
modification_type: mod_type,
status,
object_data: data,
metadata: meta,
created_at: ts.map(|t| t.to_rfc3339()),
}
},
)
.collect();
let next_cursor = entries.last().map(|e| e.cursor);
(
StatusCode::OK,
Json(ChangelogListResponse {
entries,
next_cursor,
}),
)
.into_response()
},
Err(e) => {
tracing::error!("Failed to query changelog: {e}");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": format!("Database error: {e}") })),
)
.into_response()
},
}
}
pub async fn checkpoint_get_handler(
State(state): State<ChangelogState>,
Path(listener_id): Path<String>,
) -> impl IntoResponse {
let result = sqlx::query_as::<_, (String, i64, Option<DateTime<Utc>>)>(
r"
SELECT listener_id, last_processed_id, updated_at
FROM observer_checkpoints
WHERE listener_id = $1
",
)
.bind(&listener_id)
.fetch_optional(&state.pool)
.await;
match result {
Ok(Some((lid, cursor, updated))) => (
StatusCode::OK,
Json(CheckpointResponse {
listener_id: lid,
last_cursor: cursor,
updated_at: updated.map(|t| t.to_rfc3339()),
}),
)
.into_response(),
Ok(None) => (
StatusCode::NOT_FOUND,
Json(serde_json::json!({ "error": "Checkpoint not found" })),
)
.into_response(),
Err(e) => {
tracing::error!("Failed to read checkpoint: {e}");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": format!("Database error: {e}") })),
)
.into_response()
},
}
}
pub async fn checkpoint_save_handler(
State(state): State<ChangelogState>,
Path(listener_id): Path<String>,
Json(body): Json<SaveCheckpointRequest>,
) -> impl IntoResponse {
let result = sqlx::query(
r"
INSERT INTO observer_checkpoints
(listener_id, last_processed_id, last_processed_at, batch_size, event_count, updated_at)
VALUES ($1, $2, NOW(), 0, 0, NOW())
ON CONFLICT (listener_id) DO UPDATE SET
last_processed_id = $2,
last_processed_at = NOW(),
updated_at = NOW()
",
)
.bind(&listener_id)
.bind(body.last_cursor)
.execute(&state.pool)
.await;
match result {
Ok(_) => (
StatusCode::OK,
Json(serde_json::json!({
"listener_id": listener_id,
"last_cursor": body.last_cursor,
"message": "Checkpoint saved"
})),
)
.into_response(),
Err(e) => {
tracing::error!("Failed to save checkpoint: {e}");
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({ "error": format!("Database error: {e}") })),
)
.into_response()
},
}
}
#[cfg(test)]
mod tests;