use async_trait::async_trait;
use chrono::DateTime;
use serde::{Deserialize, Serialize};
use std::cmp::Ordering;
use crate::core::{AppError, ErrorKind};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdObject {
pub id: String,
pub collection: String,
pub data: serde_json::Value,
pub created_at: String,
pub updated_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdEvent {
pub id: String,
pub stream: String,
pub event_type: String,
pub data: serde_json::Value,
pub timestamp: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdJob {
pub id: String,
pub queue: String,
pub payload: serde_json::Value,
pub state: JobState,
pub attempts: u32,
pub max_retries: u32,
pub lease_expires_at: Option<String>,
pub created_at: String,
pub updated_at: String,
pub available_at_ms: Option<i64>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct PushJobOptions {
pub idempotency_key: Option<String>,
pub delay_ms: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum JobState {
Queued,
Leased,
Completed,
Retrying,
Dead,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdLink {
pub id: String,
pub source_id: String,
pub target_id: String,
pub relation: String,
pub created_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdFilter {
pub field: String,
pub operator: FilterOperator,
pub value: serde_json::Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum FilterOperator {
Eq,
Ne,
Gt,
Lt,
Gte,
Lte,
Contains,
}
pub(crate) fn matches_filter(
object: &ThingdObject,
filter: &ThingdFilter,
) -> Result<bool, AppError> {
let Some(value) = object.data.get(&filter.field) else {
return Ok(false);
};
match filter.operator {
FilterOperator::Eq => Ok(*value == filter.value),
FilterOperator::Ne => Ok(*value != filter.value),
FilterOperator::Contains => match (value.as_str(), filter.value.as_str()) {
(Some(value), Some(needle)) => Ok(value.contains(needle)),
_ => Err(invalid_filter(
&filter.field,
"Contains requires string values",
)),
},
FilterOperator::Gt | FilterOperator::Lt | FilterOperator::Gte | FilterOperator::Lte => {
let ordering = compare_filter_values(value, &filter.value, &filter.field)?;
Ok(match filter.operator {
FilterOperator::Gt => ordering == Ordering::Greater,
FilterOperator::Lt => ordering == Ordering::Less,
FilterOperator::Gte => ordering != Ordering::Less,
FilterOperator::Lte => ordering != Ordering::Greater,
_ => unreachable!(),
})
}
}
}
pub(crate) fn filter_objects(
objects: Vec<ThingdObject>,
filters: &[ThingdFilter],
) -> Result<Vec<ThingdObject>, AppError> {
objects
.into_iter()
.filter_map(|object| {
let result = filters.iter().try_fold(true, |matches, filter| {
if matches {
matches_filter(&object, filter)
} else {
Ok(false)
}
});
match result {
Ok(true) => Some(Ok(object)),
Ok(false) => None,
Err(error) => Some(Err(error)),
}
})
.collect()
}
fn compare_filter_values(
value: &serde_json::Value,
filter_value: &serde_json::Value,
field: &str,
) -> Result<Ordering, AppError> {
if let (Some(left), Some(right)) = (value.as_f64(), filter_value.as_f64()) {
return left
.partial_cmp(&right)
.ok_or_else(|| invalid_filter(field, "numeric values must be finite"));
}
if let (Some(left), Some(right)) = (value.as_str(), filter_value.as_str()) {
if let (Ok(left), Ok(right)) = (
DateTime::parse_from_rfc3339(left),
DateTime::parse_from_rfc3339(right),
) {
return Ok(left.cmp(&right));
}
return Ok(left.cmp(right));
}
Err(invalid_filter(
field,
"range comparisons require two numbers or two strings",
))
}
fn invalid_filter(field: &str, message: &str) -> AppError {
AppError::new(
ErrorKind::Validation,
format!("invalid filter for field '{field}': {message}"),
)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SearchOptions {
pub limit: usize,
pub offset: usize,
pub filters: Vec<ThingdFilter>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct QueryOptions {
pub filters: Vec<ThingdFilter>,
pub limit: Option<usize>,
pub offset: usize,
}
impl QueryOptions {
pub fn filtered(filters: Vec<ThingdFilter>) -> Self {
Self {
filters,
limit: None,
offset: 0,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SearchResults {
pub items: Vec<ThingdObject>,
pub total: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ThingdOperation {
Put {
collection: String,
id: String,
data: serde_json::Value,
},
Delete { collection: String, id: String },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ThingdOperationResult {
pub success: bool,
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ThingdCompatibilityReport {
pub api_version: String,
pub status: String,
}
#[async_trait]
pub trait ThingdBackend: Send + Sync {
async fn check_compatibility(
&self,
) -> Result<ThingdCompatibilityReport, crate::core::AppError> {
Ok(ThingdCompatibilityReport {
api_version: "local".to_string(),
status: "ok".to_string(),
})
}
async fn get_object(
&self,
collection: &str,
id: &str,
) -> Result<Option<ThingdObject>, crate::core::AppError>;
async fn put_object(
&self,
collection: &str,
id: &str,
data: serde_json::Value,
) -> Result<ThingdObject, crate::core::AppError>;
async fn delete_object(&self, collection: &str, id: &str) -> Result<(), crate::core::AppError>;
async fn query_objects(
&self,
collection: &str,
options: QueryOptions,
) -> Result<Vec<ThingdObject>, crate::core::AppError>;
async fn count_objects(&self, collection: &str) -> Result<usize, crate::core::AppError>;
async fn batch_write(
&self,
operations: Vec<ThingdOperation>,
) -> Result<Vec<ThingdOperationResult>, crate::core::AppError>;
async fn append_event(
&self,
stream: &str,
event_type: &str,
data: serde_json::Value,
) -> Result<ThingdEvent, crate::core::AppError>;
async fn read_events(
&self,
stream: &str,
from: Option<String>,
limit: usize,
) -> Result<Vec<ThingdEvent>, crate::core::AppError>;
async fn push_job(
&self,
queue: &str,
payload: serde_json::Value,
max_retries: u32,
) -> Result<ThingdJob, crate::core::AppError>;
async fn push_job_with_options(
&self,
queue: &str,
payload: serde_json::Value,
max_retries: u32,
options: PushJobOptions,
) -> Result<ThingdJob, crate::core::AppError> {
let _ = options;
self.push_job(queue, payload, max_retries).await
}
async fn claim_job(
&self,
queue: &str,
worker_id: &str,
lease_seconds: u32,
) -> Result<Option<ThingdJob>, crate::core::AppError>;
async fn complete_job(&self, queue: &str, job_id: &str) -> Result<(), crate::core::AppError>;
async fn nack_job(&self, queue: &str, job_id: &str) -> Result<(), crate::core::AppError>;
async fn dead_letter_job(&self, queue: &str, job_id: &str)
-> Result<(), crate::core::AppError>;
async fn search(
&self,
query: &str,
options: SearchOptions,
) -> Result<SearchResults, crate::core::AppError>;
async fn create_link(
&self,
source_id: &str,
target_id: &str,
relation: &str,
) -> Result<ThingdLink, crate::core::AppError>;
async fn get_links(
&self,
source_id: &str,
relation: Option<&str>,
) -> Result<Vec<ThingdLink>, crate::core::AppError>;
async fn delete_link(&self, link_id: &str) -> Result<(), crate::core::AppError>;
async fn reset(&self) -> Result<(), crate::core::AppError>;
async fn seed(&self) -> Result<(), crate::core::AppError>;
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn object(value: serde_json::Value) -> ThingdObject {
ThingdObject {
id: "one".into(),
collection: "items".into(),
data: value,
created_at: String::new(),
updated_at: String::new(),
}
}
#[test]
fn range_filters_compare_rfc3339_and_boundaries() {
let value = object(json!({"expiresAt": "2026-08-10T12:00:00Z", "price": 10}));
for (operator, expected) in [
(FilterOperator::Lt, true),
(FilterOperator::Lte, true),
(FilterOperator::Gt, false),
(FilterOperator::Gte, false),
] {
let filter = ThingdFilter {
field: "expiresAt".into(),
operator,
value: json!("2026-08-11T12:00:00+00:00"),
};
assert_eq!(matches_filter(&value, &filter).unwrap(), expected);
}
let equal = ThingdFilter {
field: "price".into(),
operator: FilterOperator::Lte,
value: json!(10),
};
assert!(matches_filter(&value, &equal).unwrap());
}
#[test]
fn invalid_range_filter_fails_loudly() {
let filter = ThingdFilter {
field: "price".into(),
operator: FilterOperator::Lt,
value: json!("not-a-number"),
};
let error = matches_filter(&object(json!({"price": 10})), &filter).unwrap_err();
assert_eq!(error.kind, ErrorKind::Validation);
}
}