mirror of
https://forgejo.ellis.link/continuwuation/continuwuity.git
synced 2026-05-26 20:49:55 +00:00
feat(federation): Restructure transaction_ids service
Adds two new in-memory maps to the service in to prepare for better handlers
This commit is contained in:
@@ -1,11 +1,22 @@
|
|||||||
use std::sync::Arc;
|
use std::{collections::HashMap, sync::Arc};
|
||||||
|
|
||||||
use conduwuit::{Result, implement};
|
use conduwuit::{Result, SyncRwLock};
|
||||||
use database::{Handle, Map};
|
use database::{Handle, Map};
|
||||||
use ruma::{DeviceId, TransactionId, UserId};
|
use ruma::{
|
||||||
|
DeviceId, OwnedServerName, OwnedTransactionId, TransactionId, UserId,
|
||||||
|
api::federation::transactions::send_transaction_message,
|
||||||
|
};
|
||||||
|
use tokio::sync::watch::{Receiver, Sender};
|
||||||
|
|
||||||
|
pub type TxnKey = (OwnedServerName, OwnedTransactionId);
|
||||||
|
pub type TxnChanType = (TxnKey, send_transaction_message::v1::Response);
|
||||||
|
pub type ActiveTxnsMapType = HashMap<TxnKey, (Sender<TxnChanType>, Receiver<TxnChanType>)>;
|
||||||
|
|
||||||
pub struct Service {
|
pub struct Service {
|
||||||
db: Data,
|
db: Data,
|
||||||
|
pub servername_txnid_response_cache:
|
||||||
|
Arc<SyncRwLock<HashMap<TxnKey, send_transaction_message::v1::Response>>>,
|
||||||
|
pub servername_txnid_active: Arc<SyncRwLock<ActiveTxnsMapType>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
struct Data {
|
struct Data {
|
||||||
@@ -18,13 +29,15 @@ impl crate::Service for Service {
|
|||||||
db: Data {
|
db: Data {
|
||||||
userdevicetxnid_response: args.db["userdevicetxnid_response"].clone(),
|
userdevicetxnid_response: args.db["userdevicetxnid_response"].clone(),
|
||||||
},
|
},
|
||||||
|
servername_txnid_response_cache: Arc::new(SyncRwLock::new(HashMap::new())),
|
||||||
|
servername_txnid_active: Arc::new(SyncRwLock::new(HashMap::new())),
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
|
fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
|
||||||
}
|
}
|
||||||
|
|
||||||
#[implement(Service)]
|
impl Service {
|
||||||
pub fn add_txnid(
|
pub fn add_txnid(
|
||||||
&self,
|
&self,
|
||||||
user_id: &UserId,
|
user_id: &UserId,
|
||||||
@@ -41,8 +54,6 @@ pub fn add_txnid(
|
|||||||
self.db.userdevicetxnid_response.insert(&key, data);
|
self.db.userdevicetxnid_response.insert(&key, data);
|
||||||
}
|
}
|
||||||
|
|
||||||
// If there's no entry, this is a new transaction
|
|
||||||
#[implement(Service)]
|
|
||||||
pub async fn existing_txnid(
|
pub async fn existing_txnid(
|
||||||
&self,
|
&self,
|
||||||
user_id: &UserId,
|
user_id: &UserId,
|
||||||
@@ -52,3 +63,4 @@ pub async fn existing_txnid(
|
|||||||
let key = (user_id, device_id, txn_id);
|
let key = (user_id, device_id, txn_id);
|
||||||
self.db.userdevicetxnid_response.qry(&key).await
|
self.db.userdevicetxnid_response.qry(&key).await
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user