use async_trait::async_trait;
use chrono::Utc;
use pubky_app_specs::ParsedUri;
use serde::{Deserialize, Serialize};
use nexus_common::db::RedisOps;
use nexus_common::types::DynError;
use crate::events::errors::EventProcessorError;
pub const RETRY_MANAGER_PREFIX: &str = "RetryManager";
pub const RETRY_MANAGER_EVENTS_INDEX: [&str; 1] = ["events"];
pub const RETRY_MANAGER_STATE_INDEX: [&str; 1] = ["state"];
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RetryEvent {
pub retry_count: u32,
pub error_type: EventProcessorError,
}
#[async_trait]
impl RedisOps for RetryEvent {
async fn prefix() -> String {
String::from(RETRY_MANAGER_PREFIX)
}
}
impl RetryEvent {
pub fn new(error_type: EventProcessorError) -> Self {
Self {
retry_count: 0,
error_type,
}
}
pub fn generate_index_key(event_uri: &str) -> Option<String> {
let parsed_uri = match ParsedUri::try_from(event_uri) {
Ok(parsed_uri) => parsed_uri,
Err(_) => return None,
};
let user_id = parsed_uri.user_id;
let key = match parsed_uri.resource.id() {
Some(id) => format!("{}:{}:{}", user_id, parsed_uri.resource, id),
None => format!("{}:{}", user_id, parsed_uri.resource),
};
Some(key)
}
pub async fn put_to_index(&self, event_line: String) -> Result<(), DynError> {
Self::put_index_sorted_set(
&RETRY_MANAGER_EVENTS_INDEX,
&[(Utc::now().timestamp_millis() as f64, &event_line)],
Some(RETRY_MANAGER_PREFIX),
None,
)
.await?;
let index = &[RETRY_MANAGER_STATE_INDEX, [&event_line]].concat();
self.put_index_json(index, None, None).await?;
Ok(())
}
pub async fn check_uri(event_index: &str) -> Result<Option<isize>, DynError> {
Self::check_sorted_set_member(
Some(RETRY_MANAGER_PREFIX),
&RETRY_MANAGER_EVENTS_INDEX,
&[event_index],
)
.await
}
pub async fn get_from_index(event_index: &str) -> Result<Option<Self>, DynError> {
let index: &Vec<&str> = &[RETRY_MANAGER_STATE_INDEX, [event_index]].concat();
Self::try_from_index_json(index, None).await
}
}