ferrox_database_dynamo/
lib.rs1use async_trait::async_trait;
2use ferrox_database_core::Repository;
3use ferrox_errors::AppError;
4use aws_sdk_dynamodb::{Client, types::AttributeValue};
5use serde::{Serialize, de::DeserializeOwned};
6use std::marker::PhantomData;
7
8#[derive(Clone, Debug)]
10pub struct DynamoId {
11 pub pk: String,
12 pub sk: Option<String>,
13}
14
15impl From<&str> for DynamoId {
16 fn from(s: &str) -> Self {
17 DynamoId { pk: s.to_string(), sk: None }
18 }
19}
20
21pub struct DynamoDbRepository<T> {
23 client: Client,
24 table_name: String,
25 partition_key_name: String,
26 sort_key_name: Option<String>,
27 _marker: PhantomData<T>,
28}
29
30impl<T> DynamoDbRepository<T> {
31 pub fn new(client: Client, table_name: String, partition_key_name: String, sort_key_name: Option<String>) -> Self {
32 Self {
33 client,
34 table_name,
35 partition_key_name,
36 sort_key_name,
37 _marker: PhantomData,
38 }
39 }
40}
41
42#[async_trait]
43impl<T> Repository<T, DynamoId> for DynamoDbRepository<T>
44where
45 T: Serialize + DeserializeOwned + Send + Sync + Clone,
46{
47 async fn find_by_id(&self, id: DynamoId) -> Result<Option<T>, AppError> {
48 let mut req = self.client.get_item()
49 .table_name(&self.table_name)
50 .key(&self.partition_key_name, AttributeValue::S(id.pk));
51
52 if let (Some(sk_name), Some(sk_val)) = (&self.sort_key_name, &id.sk) {
53 req = req.key(sk_name, AttributeValue::S(sk_val.clone()));
54 }
55
56 let res = req.send().await
57 .map_err(|e| AppError::InternalError(format!("DynamoDB GetItem Error: {}", e)))?;
58
59 if let Some(item) = res.item {
60 let parsed: T = serde_dynamo::aws_sdk_dynamodb_1::from_item(item)
61 .map_err(|e| AppError::InternalError(format!("DynamoDB Deserialization Error: {}", e)))?;
62 Ok(Some(parsed))
63 } else {
64 Ok(None)
65 }
66 }
67
68 async fn find_all(&self) -> Result<Vec<T>, AppError> {
69 let res = self.client.scan()
71 .table_name(&self.table_name)
72 .send()
73 .await
74 .map_err(|e| AppError::InternalError(format!("DynamoDB Scan Error: {}", e)))?;
75
76 let mut entities = Vec::new();
77 if let Some(items) = res.items {
78 for item in items {
79 let parsed: T = serde_dynamo::aws_sdk_dynamodb_1::from_item(item)
80 .map_err(|e| AppError::InternalError(format!("DynamoDB Deserialization Error: {}", e)))?;
81 entities.push(parsed);
82 }
83 }
84 Ok(entities)
85 }
86
87 async fn insert(&self, entity: T) -> Result<T, AppError> {
88 let item = serde_dynamo::aws_sdk_dynamodb_1::to_item(entity.clone())
89 .map_err(|e| AppError::InternalError(format!("DynamoDB Serialization Error: {}", e)))?;
90
91 self.client.put_item()
92 .table_name(&self.table_name)
93 .set_item(Some(item))
94 .send()
95 .await
96 .map_err(|e| AppError::InternalError(format!("DynamoDB PutItem Error: {}", e)))?;
97
98 Ok(entity)
99 }
100
101 async fn update(&self, _id: DynamoId, entity: T) -> Result<T, AppError> {
102 self.insert(entity).await
105 }
106
107 async fn delete(&self, id: DynamoId) -> Result<(), AppError> {
108 let mut req = self.client.delete_item()
109 .table_name(&self.table_name)
110 .key(&self.partition_key_name, AttributeValue::S(id.pk));
111
112 if let (Some(sk_name), Some(sk_val)) = (&self.sort_key_name, &id.sk) {
113 req = req.key(sk_name, AttributeValue::S(sk_val.clone()));
114 }
115
116 req.send().await
117 .map_err(|e| AppError::InternalError(format!("DynamoDB DeleteItem Error: {}", e)))?;
118
119 Ok(())
120 }
121}