Skip to main content

ferrox_database_dynamo/
lib.rs

1use 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/// DynamoDB Primary Key definition, supporting both simple (PK only) and composite (PK + SK) keys.
9#[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
21/// A DynamoDB-backed repository implementation.
22pub 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        // Warning: Scans are expensive in DynamoDB. Used here to satisfy the generic Repository trait.
70        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        // In DynamoDB, PutItem overwrites completely. For a true update, UpdateItem is used, 
103        // but PutItem satisfies the standard generic repository update definition via replacement.
104        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}