use axum::Extension;
use axum::{extract::Query, http::StatusCode, response::IntoResponse, Json};
use chrono::offset::TimeZone;
use chrono::{DateTime, Duration, Utc};
use mongodb::bson::doc;
use mongodb::bson::Bson;
use mongodb::options::{FindOneAndUpdateOptions, FindOneOptions};
use mongodb::Collection;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::config::IConfig;
use crate::data::Data;
use crate::database::DBHandle;
use crate::routes::response::{ErrorResponse, QueueResponse};
use crate::subject::Subject;
use crate::user::User;
use crate::utils::deserialise_array::deserialise_array;
#[derive(Debug, Serialize, Deserialize)]
pub struct InternalQueueItem {
pub queue_id: String, pub platform_id: String,
pub platform: String,
pub last_processed: DateTime<Utc>,
pub lock_holder: Option<String>, pub lock_acquired_at: Option<DateTime<Utc>>,
pub references: u64,
pub confirmed_id: bool,
}
impl InternalQueueItem {
fn new(platform_id: String, platform: String) -> Self {
Self {
queue_id: Uuid::new_v4().to_string(),
platform_id,
platform,
last_processed: Utc.with_ymd_and_hms(1970, 1, 1, 0, 0, 0).unwrap(),
lock_holder: None,
lock_acquired_at: None,
references: 1,
confirmed_id: false,
}
}
}
#[derive(Deserialize)]
pub struct QueueQuery {
#[serde(deserialize_with = "deserialise_array")]
platforms: Vec<String>,
}
pub async fn queue(
user: User,
mut db: DBHandle,
Extension(config): Extension<IConfig>,
queue_query: Option<Query<QueueQuery>>,
) -> impl IntoResponse {
if queue_query.is_none() {
return error!(
BAD_REQUEST,
"You must provide a list of supported platforms."
);
}
let platforms = &queue_query.unwrap().platforms;
if platforms.is_empty() {
return error!(
BAD_REQUEST,
"You must specify which platforms you are performing jobs for."
);
}
if platforms.iter().any(|p| !config.valid_platform(p)) {
error!(
BAD_REQUEST,
"One or more of your given platforms is not valid.
See /types for supported platforms."
)
} else {
let filter_builder = FindOneAndUpdateOptions::builder()
.sort(doc! {"last_processed": -1_i32});
let filter = filter_builder.build();
let q_coll: Collection<InternalQueueItem> = db.collection("queue");
let result = q_coll
.find_one_and_update_with_session(
doc! {"lock_holder": Bson::Null, "platform": {"$in": &platforms}},
doc! {"$set":
{
"lock_holder": user.uuid,
"lock_acquired_at": Utc::now().to_string()
}
},
filter,
&mut db.session,
)
.await
.unwrap();
if let Some(queue_item) = result {
let username_hint: String = get_username_hint(
&queue_item.platform_id,
&queue_item.platform,
&mut db,
)
.await;
db.session.commit_transaction().await.unwrap();
ok!(
OK,
QueueResponse::new(
queue_item.queue_id,
queue_item.platform,
queue_item.platform_id,
username_hint,
)
)
} else {
error!(OK, "There are no jobs available. Please try again later.")
}
}
}
pub async fn process(
queue_id: &str,
id: &str,
platform: &str,
added_by: &str,
username: Option<String>,
db: &mut DBHandle,
) -> bool {
let q_coll: Collection<InternalQueueItem> = db.collection("queue");
if let Some(username) = username {
let find_result = q_coll
.find_one_with_session(
doc! {
"queue_id" : queue_id,
"platform": platform,
"platform_id": &username,
"lock_holder": added_by,
"confirmed_id": false
},
None,
&mut db.session,
)
.await
.unwrap();
if find_result.is_some() {
remove_queue_item(&username, platform, db).await;
add_queue_item(id, platform, db, true).await;
let subj_coll: Collection<Subject> = db.collection("subjects");
subj_coll
.update_one_with_session(
doc! {&format!("profiles.{platform}"): &username},
doc! {"$set": {&format!("profiles.{platform}.$"): id}},
None,
&mut db.session,
)
.await
.unwrap();
return true;
}
}
let q_update_result = q_coll
.update_one_with_session(
doc! {"queue_id" : queue_id, "lock_holder": added_by},
doc! {"$set":
{
"lock_holder": Bson::Null,
"lock_acquired_at": Bson::Null,
"last_processed": Utc::now().to_string()
}
},
None,
&mut db.session,
)
.await
.unwrap();
q_update_result.modified_count == 1
}
pub async fn add_queue_item(
platform_id: &str,
platform: &str,
db: &mut DBHandle,
confirmed_id: bool,
) {
let q_coll: Collection<InternalQueueItem> = db.collection("queue");
let queue_item = q_coll
.find_one_with_session(
doc! {"platform_id": platform_id, "platform": platform},
None,
&mut db.session,
)
.await
.unwrap();
if queue_item.is_some() {
q_coll
.update_one_with_session(
doc! {"platform_id": platform_id, "platform": platform},
doc! {"$inc": {"references": 1_u32}},
None,
&mut db.session,
)
.await
.unwrap();
if confirmed_id {
q_coll
.update_one_with_session(
doc! {"platform_id": platform_id, "platform": platform},
doc! {"$set": {"confirmed_id": true}},
None,
&mut db.session,
)
.await
.unwrap();
}
} else {
let queue_item: InternalQueueItem = InternalQueueItem::new(
platform_id.to_string(),
platform.to_string(),
);
q_coll
.insert_one_with_session(queue_item, None, &mut db.session)
.await
.unwrap();
}
}
pub async fn remove_queue_item(
platform_id: &str,
platform: &str,
db: &mut DBHandle,
) {
let q_coll: Collection<InternalQueueItem> = db.collection("queue");
let result = q_coll
.delete_one_with_session(
doc! {
"platform_id": platform_id,
"platform": platform,
"references": 1
},
None,
&mut db.session,
)
.await
.unwrap();
if result.deleted_count == 0 {
q_coll
.update_one_with_session(
doc! {"platform_id": platform_id, "platform": platform},
doc! {"$inc": {"references": -1_i32}},
None,
&mut db.session,
)
.await
.unwrap();
}
}
pub async fn get_username_hint(
platform_id: &str,
platform: &str,
db: &mut DBHandle,
) -> String {
async fn from_meta(
platform_id: &str,
platform: &str,
db: &mut DBHandle,
) -> Option<String> {
let filter = FindOneOptions::builder()
.projection(doc! {"username": 1_u32})
.build();
#[derive(Debug, Serialize, Deserialize)]
struct Username {
username: String,
}
let data_coll: Collection<Data> = db.collection("data");
let username = data_coll
.clone_with_type::<Username>()
.find_one_with_session(
doc! {"id": &platform_id, "platform": &platform},
filter,
&mut db.session,
)
.await;
if let Ok(Some(username)) = username {
return Some(username.username);
}
None
}
from_meta(platform_id, platform, db)
.await
.unwrap_or_else(|| platform_id.to_string())
}
pub async fn clear_old_locks(db: &mut DBHandle, timeout: Duration) {
let q_coll: Collection<InternalQueueItem> = db.collection("queue");
let thirty_seconds_ago = Utc::now() - timeout;
q_coll
.update_many(
doc! {"lock_acquired_at": {"$lt": thirty_seconds_ago.to_string()}},
doc! {"$set":
{"lock_acquired_at": Bson::Null,
"lock_holder": Bson::Null}
},
None,
)
.await
.unwrap();
}