1use crate::web5::{Web5Error, Web5Result};
7use serde::{Deserialize, Serialize};
8use std::collections::HashMap;
9use std::sync::{Arc, Mutex};
10use std::time::{Duration, SystemTime, UNIX_EPOCH};
11#[derive(Clone, Debug)]
16pub struct DWNConfig {
17 pub endpoint: Option<String>,
19 pub use_local_storage: bool,
21 pub max_message_size: usize,
23}
24
25impl Default for DWNConfig {
26 fn default() -> Self {
27 Self {
28 endpoint: None,
29 use_local_storage: true,
30 max_message_size: 1024 * 1024, }
32 }
33}
34
35#[derive(Clone, Debug)]
39pub struct DWNMessage {
40 pub id: String,
42 pub from: String,
44 pub to: String,
46 pub protocol: String,
48 pub message_type: String,
50 pub data: Vec<u8>,
52 pub timestamp: u64,
54 pub attestations: Vec<Attestation>,
56}
57
58pub struct DWNClient {
62 config: DWNConfig,
64 local_storage: Arc<Mutex<HashMap<String, DWNMessage>>>,
66 identity: Option<String>,
68}
69
70impl DWNClient {
71 pub fn new(config: DWNConfig) -> Self {
73 Self {
74 config,
75 local_storage: Arc::new(Mutex::new(HashMap::new())),
76 identity: None,
77 }
78 }
79
80 pub fn set_identity(&mut self, did: &str) {
82 self.identity = Some(did.to_string());
83 }
84
85 pub fn send_message(
87 &self,
88 to: &str,
89 protocol: &str,
90 message_type: &str,
91 data: &[u8],
92 ) -> Web5Result<String> {
93 let from = self
95 .identity
96 .as_ref()
97 .ok_or_else(|| Web5Error::Identity("Identity not set for DWN client".to_string()))?;
98
99 if data.len() > self.config.max_message_size {
101 return Err(Web5Error::Communication(format!(
102 "Message size exceeds maximum allowed: {} > {}",
103 data.len(),
104 self.config.max_message_size
105 )));
106 }
107
108 let id = format!("msg_{}", generate_id());
110
111 let message = DWNMessage {
113 id: id.clone(),
114 from: from.clone(),
115 to: to.to_string(),
116 protocol: protocol.to_string(),
117 message_type: message_type.to_string(),
118 data: data.to_vec(),
119 timestamp: current_time(),
120 attestations: Vec::new(),
121 };
122
123 if self.config.use_local_storage {
125 let mut storage = self
126 .local_storage
127 .lock()
128 .map_err(|e| format!("Mutex lock error: {e}"))?;
129 let message_for_storage = message.clone();
130 storage.insert(id.clone(), message_for_storage);
131 }
132
133 if let Some(endpoint) = &self.config.endpoint {
135 println!("Would send message to DWN at {endpoint}: {message:?}");
138 }
139
140 Ok(id)
141 }
142
143 pub fn get_messages(&self, protocol: Option<&str>) -> Web5Result<Vec<DWNMessage>> {
145 let _from = self
147 .identity
148 .as_ref()
149 .ok_or_else(|| Web5Error::Identity("Identity not set for DWN client".to_string()))?;
150
151 let storage = self
152 .local_storage
153 .lock()
154 .map_err(|e| format!("Mutex lock error: {e}"))?;
155
156 let messages: Vec<DWNMessage> = storage
158 .values()
159 .filter(|msg| msg.to == *_from && protocol.map_or(true, |p| msg.protocol == p))
160 .cloned()
161 .collect();
162
163 Ok(messages)
164 }
165}
166
167fn generate_id() -> String {
170 let now = SystemTime::now()
171 .duration_since(UNIX_EPOCH)
172 .unwrap_or_default()
173 .as_secs();
174
175 format!("{now:x}")
176}
177
178fn current_time() -> u64 {
180 SystemTime::now()
181 .duration_since(UNIX_EPOCH)
182 .map(|d| d.as_secs())
183 .unwrap_or(0)
184}
185
186#[derive(Debug)]
190pub struct DWNManager {
191 records: Arc<Mutex<HashMap<String, DWNRecord>>>,
193}
194
195#[derive(Debug, Clone, Serialize, Deserialize)]
199pub struct DWNRecord {
200 pub id: String,
202 pub owner: String,
204 pub schema: String,
206 pub data: serde_json::Value,
208 pub metadata: HashMap<String, String>,
210 pub attestations: Vec<Attestation>,
212}
213
214#[derive(Debug, Clone, Serialize, Deserialize)]
218pub struct Attestation {
219 pub issuer: String,
221 pub timestamp: u64,
223 pub signature: String,
225}
226
227#[derive(Debug, Clone, Serialize, Deserialize)]
231pub enum DWNMessageType {
232 #[serde(rename = "create")]
234 Create,
235 #[serde(rename = "read")]
237 Read,
238 #[serde(rename = "update")]
240 Update,
241 #[serde(rename = "delete")]
243 Delete,
244 #[serde(rename = "query")]
246 Query,
247}
248
249#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct DWNMessageDescriptor {
254 pub id: String,
256 pub author: String,
258 pub recipient: Option<String>,
260 pub protocol: Option<String>,
262 pub schema: String,
264 pub data_format: String,
266 pub timestamp: u64,
268}
269
270#[derive(Debug, Clone, Serialize, Deserialize)]
274pub struct DWNQuery {
275 pub filter: DWNQueryFilter,
277 pub pagination: Option<DWNQueryPagination>,
279}
280
281#[derive(Debug, Clone, Serialize, Deserialize)]
285pub struct DateRange {
286 pub from: Option<u64>,
288 pub to: Option<u64>,
290}
291
292#[derive(Debug, Clone, Serialize, Deserialize)]
296pub struct DWNQueryFilter {
297 pub owner: Option<String>,
299 pub schema: Option<String>,
301 pub metadata: Option<HashMap<String, String>>,
303 pub date_range: Option<DateRange>,
305 pub data_filter: Option<serde_json::Value>,
307}
308
309#[derive(Debug, Clone, Serialize, Deserialize)]
313pub struct DWNQueryPagination {
314 pub offset: Option<usize>,
316 pub limit: Option<usize>,
318 pub cursor: Option<String>,
320}
321
322impl Default for DWNQueryPagination {
323 fn default() -> Self {
324 Self {
325 offset: None,
326 limit: Some(100), cursor: None,
328 }
329 }
330}
331
332#[derive(Debug, Clone, Serialize, Deserialize)]
334pub struct DWNQueryResult {
335 pub records: Vec<DWNRecord>,
337 pub pagination: DWNQueryPaginationResult,
339}
340
341#[derive(Debug, Clone, Serialize, Deserialize)]
343pub struct DWNQueryPaginationResult {
344 pub total: usize,
346 pub count: usize,
348 pub has_more: bool,
350 pub next_cursor: Option<String>,
352}
353
354#[derive(Debug, Clone, Serialize, Deserialize)]
356pub struct AdvancedDWNQueryFilter {
357 pub base: DWNQueryFilter,
359 pub search: Option<String>,
361 pub geo_bounds: Option<GeoBounds>,
363 pub tags: Option<Vec<String>>,
365 pub numeric_ranges: Option<HashMap<String, NumericRange>>,
367}
368
369#[derive(Debug, Clone, Serialize, Deserialize)]
371pub struct GeoBounds {
372 pub min_lat: f64,
373 pub max_lat: f64,
374 pub min_lng: f64,
375 pub max_lng: f64,
376}
377
378#[derive(Debug, Clone, Serialize, Deserialize)]
380pub struct NumericRange {
381 pub min: Option<f64>,
382 pub max: Option<f64>,
383}
384
385#[derive(Debug, Clone, Serialize, Deserialize)]
387pub enum SyncStatus {
388 Synced,
389 Pending,
390 Failed(String),
391 Conflicted,
392}
393
394#[derive(Debug, Clone, Serialize, Deserialize)]
396pub struct SyncedDWNRecord {
397 pub record: DWNRecord,
398 pub sync_status: SyncStatus,
399 pub last_sync: u64,
400 pub sync_attempts: u32,
401}
402
403#[derive(Debug, Clone, Serialize, Deserialize)]
405pub enum ConflictResolution {
406 LastWriteWins,
407 FirstWriteWins,
408 Manual,
409 Custom(String),
410}
411
412impl Default for DWNManager {
413 fn default() -> Self {
414 Self {
415 records: Arc::new(Mutex::new(HashMap::new())),
416 }
417 }
418}
419
420impl DWNManager {
421 pub fn new() -> Self {
423 Self {
424 records: Arc::new(Mutex::new(HashMap::new())),
425 }
426 }
427
428 pub fn store_record(&self, record: DWNRecord) -> Web5Result<String> {
430 let mut storage = self
431 .records
432 .lock()
433 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
434 let record_id = record.id.clone();
435 storage.insert(record_id.clone(), record);
436 Ok(record_id)
437 }
438
439 pub fn query_records(&self, owner: &str, schema: &str) -> Web5Result<Vec<DWNRecord>> {
441 let storage = self
442 .records
443 .lock()
444 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
445 let records: Vec<DWNRecord> = storage
446 .values()
447 .filter(|r| r.owner == owner && r.schema == schema)
448 .cloned()
449 .collect();
450 Ok(records)
451 }
452
453 pub fn create_record(
455 &self,
456 owner: &str,
457 schema: &str,
458 data: serde_json::Value,
459 ) -> Web5Result<String> {
460 let record = DWNRecord {
461 id: generate_id(),
462 owner: owner.to_string(),
463 schema: schema.to_string(),
464 data,
465 metadata: HashMap::new(),
466 attestations: Vec::new(),
467 };
468 self.store_record(record)
469 }
470
471 pub fn read_record(&self, id: &str) -> Web5Result<DWNRecord> {
473 let storage = self
474 .records
475 .lock()
476 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
477 storage
478 .get(id)
479 .cloned()
480 .ok_or_else(|| Web5Error::NotFound(id.to_string()))
481 }
482
483 pub fn update_record(&self, id: &str, data: serde_json::Value) -> Web5Result<()> {
485 let mut storage = self
486 .records
487 .lock()
488 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
489 if let Some(record) = storage.get_mut(id) {
492 record.data = data;
493 record
494 .metadata
495 .insert("updated".to_string(), current_time().to_string());
496 Ok(())
497 } else {
498 Err(Web5Error::NotFound("Record not found".to_string()))
499 }
500 }
501
502 pub fn delete_record(&self, id: &str) -> Web5Result<()> {
504 self.records.lock().unwrap().remove(id);
507 Ok(())
508 }
509
510 pub fn send_message(&self, message: DWNMessage) -> Web5Result<DWNMessage> {
512 match message.message_type.as_str() {
516 "Create" => {
517 let data = message.data.clone();
519 let record = DWNRecord {
521 id: message.id.clone(),
522 owner: message.from.clone(),
523 schema: message.protocol.clone(),
524 data: serde_json::from_slice(&data).unwrap_or_else(|_| serde_json::Value::Null),
525 metadata: HashMap::new(),
526 attestations: Vec::new(),
527 };
528 self.store_record(record)?;
529 Ok(message)
530 }
531 "Read" => {
532 let id = message.id.clone();
534 if let Ok(records) = self.records.lock() {
535 if let Some(record) = records.get(&id) {
536 let mut response = message.clone();
537 response.data = serde_json::to_vec(&record.data).unwrap_or_default();
538 return Ok(response);
539 }
540 }
541 Err(Web5Error::DWNError(format!("Record not found: {id}")))
542 }
543 "Update" => {
544 let id = message.id.clone();
546 let data = message.data.clone();
547 if let Ok(mut records) = self.records.lock() {
548 if let Some(record) = records.get_mut(&id) {
549 record.data = match serde_json::from_slice(&data) {
550 Ok(value) => value,
551 Err(_) => serde_json::Value::Null,
552 };
553 record.attestations = message.attestations.clone();
554 return Ok(message);
555 }
556 }
557 Err(Web5Error::DWNError(format!("Record not found: {id}")))
558 }
559 "Delete" => {
560 let id = message.id.clone();
562 self.delete_record(&id)?;
563 Ok(message)
564 }
565 "Query" => {
566 let data = message.data.clone();
568 let query: DWNQuery = match serde_json::from_slice(&data) {
570 Ok(value) => match serde_json::from_value(value) {
571 Ok(query) => query,
572 Err(e) => return Err(Web5Error::SerializationError(e.to_string())),
573 },
574 Err(e) => return Err(Web5Error::SerializationError(e.to_string())),
575 };
576
577 let owner = query.filter.owner.unwrap_or_default();
578 let schema = query.filter.schema.unwrap_or_default();
579
580 let records = self.query_records(&owner, &schema)?;
581
582 let mut response = message.clone();
583 response.data = match serde_json::to_vec(&records) {
584 Ok(bytes) => bytes,
585 Err(e) => return Err(Web5Error::SerializationError(e.to_string())),
586 };
587
588 Ok(response)
589 }
590 _ => {
591 Err(Web5Error::DWNError(format!(
593 "Unsupported message type: {}",
594 message.message_type
595 )))
596 }
597 }
598 }
599
600 pub fn create_index(&self, schema: &str, fields: &[&str]) -> Web5Result<()> {
606 println!("Creating index for schema '{schema}' on fields: {fields:?}");
609 Ok(())
610 }
611
612 fn filter_records_by_base_filter(&self, filter: DWNQueryFilter) -> Web5Result<Vec<DWNRecord>> {
615 self.query_with_filter(filter)
617 }
618
619 pub fn query_with_filter(&self, filter: DWNQueryFilter) -> Web5Result<Vec<DWNRecord>> {
620 let storage = self
621 .records
622 .lock()
623 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
624
625 let mut filtered_records: Vec<DWNRecord> = storage
626 .values()
627 .filter(|record| {
628 if let Some(ref owner) = filter.owner {
630 if &record.owner != owner && owner != "*" {
631 return false;
632 }
633 }
634
635 if let Some(ref schema) = filter.schema {
637 if &record.schema != schema {
638 return false;
639 }
640 }
641
642 if let Some(ref metadata_filter) = filter.metadata {
644 for (key, value) in metadata_filter {
645 if record.metadata.get(key) != Some(value) {
646 return false;
647 }
648 }
649 }
650
651 if let Some(ref date_range) = filter.date_range {
653 if let Some(timestamp_str) = record.metadata.get("created_at") {
654 if let Ok(timestamp) = timestamp_str.parse::<u64>() {
655 if let Some(from) = date_range.from {
656 if timestamp < from {
657 return false;
658 }
659 }
660 if let Some(to) = date_range.to {
661 if timestamp > to {
662 return false;
663 }
664 }
665 }
666 }
667 }
668
669 if let Some(ref data_filter) = filter.data_filter {
671 if !self.matches_data_filter(&record.data, data_filter) {
672 return false;
673 }
674 }
675
676 true
677 })
678 .cloned()
679 .collect();
680
681 filtered_records.sort_by(|a, b| {
683 let a_timestamp = a
684 .metadata
685 .get("created_at")
686 .and_then(|s| s.parse::<u64>().ok())
687 .unwrap_or(0);
688 let b_timestamp = b
689 .metadata
690 .get("created_at")
691 .and_then(|s| s.parse::<u64>().ok())
692 .unwrap_or(0);
693 b_timestamp.cmp(&a_timestamp)
694 });
695
696 Ok(filtered_records)
697 }
698
699 pub fn aggregate(&self, pipeline: &[AggregationStage]) -> Web5Result<serde_json::Value> {
701 let storage = self
702 .records
703 .lock()
704 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
705
706 let mut records: Vec<DWNRecord> = storage.values().cloned().collect();
707
708 for stage in pipeline {
709 match stage {
710 AggregationStage::Match(filter) => {
711 records.retain(|record| self.matches_aggregation_filter(record, filter));
712 }
713 AggregationStage::Group {
714 id: _id,
715 fields: _fields,
716 } => {
717 return Ok(serde_json::json!({ "count": records.len() }));
720 }
721 AggregationStage::Sort(sort_fields) => {
722 records.sort_by(|a, b| {
723 for sort_field in sort_fields {
724 let a_value = self.extract_field_value(a, &sort_field.field);
725 let b_value = self.extract_field_value(b, &sort_field.field);
726 let cmp = if sort_field.ascending {
727 a_value.cmp(&b_value)
728 } else {
729 b_value.cmp(&a_value)
730 };
731 if cmp != std::cmp::Ordering::Equal {
732 return cmp;
733 }
734 }
735 std::cmp::Ordering::Equal
736 });
737 }
738 AggregationStage::Limit(limit) => {
739 records.truncate(*limit);
740 }
741 AggregationStage::Skip(skip) => {
742 if *skip < records.len() {
743 records.drain(0..*skip);
744 } else {
745 records.clear();
746 }
747 }
748 }
749 }
750
751 serde_json::to_value(records).map_err(|e| Web5Error::SerializationError(e.to_string()))
752 }
753
754 pub async fn batch_store(&self, records: Vec<DWNRecord>) -> Web5Result<Vec<String>> {
756 const BATCH_SIZE: usize = 50; let mut results = Vec::new();
759
760 for chunk in records.chunks(BATCH_SIZE) {
761 let mut chunk_results = Vec::new();
762 for record in chunk {
763 let result = self.store_record(record.clone())?;
764 chunk_results.push(result);
765 }
766 results.extend(chunk_results);
767
768 tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
770 }
771
772 Ok(results)
773 }
774
775 pub fn get_statistics(&self) -> Web5Result<serde_json::Value> {
777 let storage = self
778 .records
779 .lock()
780 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
781
782 let total_records = storage.len();
783 let mut schema_counts: HashMap<String, usize> = HashMap::new();
784 let mut owner_counts: HashMap<String, usize> = HashMap::new();
785
786 for record in storage.values() {
787 *schema_counts.entry(record.schema.clone()).or_insert(0) += 1;
788 *owner_counts.entry(record.owner.clone()).or_insert(0) += 1;
789 }
790
791 Ok(serde_json::json!({
792 "total_records": total_records,
793 "schemas": schema_counts,
794 "owners": owner_counts,
795 "timestamp": current_timestamp()
796 }))
797 }
798
799 fn matches_data_filter(&self, data: &serde_json::Value, filter: &serde_json::Value) -> bool {
801 match (data, filter) {
803 (serde_json::Value::Object(data_obj), serde_json::Value::Object(filter_obj)) => {
804 for (key, expected_value) in filter_obj {
805 if let Some(actual_value) = data_obj.get(key) {
806 if actual_value != expected_value {
807 return false;
808 }
809 } else {
810 return false;
811 }
812 }
813 true
814 }
815 _ => data == filter,
816 }
817 }
818
819 fn matches_aggregation_filter(
820 &self,
821 record: &DWNRecord,
822 filter: &HashMap<String, serde_json::Value>,
823 ) -> bool {
824 for (key, expected_value) in filter {
825 match key.as_str() {
826 "owner" => {
827 if serde_json::Value::String(record.owner.clone()) != *expected_value {
828 return false;
829 }
830 }
831 "schema" => {
832 if serde_json::Value::String(record.schema.clone()) != *expected_value {
833 return false;
834 }
835 }
836 _ => {
837 if let Some(metadata_value) = record.metadata.get(key) {
839 if serde_json::Value::String(metadata_value.clone()) != *expected_value {
840 return false;
841 }
842 } else if let Some(data_value) = record.data.get(key) {
843 if data_value != expected_value {
844 return false;
845 }
846 } else {
847 return false;
848 }
849 }
850 }
851 }
852 true
853 }
854
855 fn extract_field_value(&self, record: &DWNRecord, field: &str) -> String {
856 if let Some(value) = record.metadata.get(field) {
858 value.clone()
859 } else if let Some(value) = record.data.get(field) {
860 value.to_string()
861 } else {
862 String::new()
863 }
864 }
865
866 pub fn query_with_pagination(
868 &self,
869 filter: AdvancedDWNQueryFilter,
870 pagination: Option<DWNQueryPagination>,
871 ) -> Web5Result<DWNQueryResult> {
872 let pagination = pagination.unwrap_or_default();
873
874 let base_filter = filter.base.clone();
876
877 let filtered_records = self.filter_records_by_base_filter(base_filter)?;
879
880 let all_records = self.apply_advanced_filters(filtered_records, &filter)?;
882
883 let total = all_records.len();
884
885 let offset = pagination.offset.unwrap_or(0);
887 let limit = pagination.limit.unwrap_or(100);
888
889 let start = offset.min(total);
890 let end = (offset + limit).min(total);
891
892 let records = all_records
893 .into_iter()
894 .skip(start)
895 .take(end - start)
896 .collect::<Vec<_>>();
897 let count = records.len();
898 let has_more = end < total;
899
900 let next_cursor = if has_more {
902 Some(format!("cursor_{end}"))
903 } else {
904 None
905 };
906
907 Ok(DWNQueryResult {
908 records,
909 pagination: DWNQueryPaginationResult {
910 total,
911 count,
912 has_more,
913 next_cursor,
914 },
915 })
916 }
917
918 fn apply_advanced_filters(
920 &self,
921 mut records: Vec<DWNRecord>,
922 filter: &AdvancedDWNQueryFilter,
923 ) -> Web5Result<Vec<DWNRecord>> {
924 if let Some(ref search_query) = filter.search {
926 records.retain(|record| self.matches_search_query(record, search_query));
927 }
928
929 if let Some(ref tags) = filter.tags {
931 records.retain(|record| self.matches_tags(record, tags));
932 }
933
934 if let Some(ref numeric_ranges) = filter.numeric_ranges {
936 records.retain(|record| self.matches_numeric_ranges(record, numeric_ranges));
937 }
938
939 if let Some(ref geo_bounds) = filter.geo_bounds {
941 records.retain(|record| self.matches_geo_bounds(record, geo_bounds));
942 }
943
944 Ok(records)
945 }
946
947 fn matches_search_query(&self, record: &DWNRecord, search_query: &str) -> bool {
949 let search_lower = search_query.to_lowercase();
950
951 if let Ok(data_string) = serde_json::to_string(&record.data) {
953 if data_string.to_lowercase().contains(&search_lower) {
954 return true;
955 }
956 }
957
958 for (key, value) in &record.metadata {
960 if key.to_lowercase().contains(&search_lower)
961 || value.to_lowercase().contains(&search_lower)
962 {
963 return true;
964 }
965 }
966
967 if record.schema.to_lowercase().contains(&search_lower) {
969 return true;
970 }
971
972 false
973 }
974
975 fn matches_tags(&self, record: &DWNRecord, required_tags: &[String]) -> bool {
977 if let Some(record_tags) = record.metadata.get("tags") {
978 let record_tag_list: Result<Vec<String>, _> = serde_json::from_str(record_tags);
979 if let Ok(record_tag_list) = record_tag_list {
980 return required_tags
981 .iter()
982 .all(|tag| record_tag_list.contains(tag));
983 }
984 }
985 false
986 }
987
988 fn matches_numeric_ranges(
990 &self,
991 record: &DWNRecord,
992 ranges: &HashMap<String, NumericRange>,
993 ) -> bool {
994 for (field, range) in ranges {
995 if let Some(value_str) = record.metadata.get(field) {
997 if let Ok(value) = value_str.parse::<f64>() {
998 if !self.value_in_range(value, range) {
999 return false;
1000 }
1001 continue;
1002 }
1003 }
1004
1005 if let Some(value) = record.data.get(field) {
1007 if let Some(value_num) = value.as_f64() {
1008 if !self.value_in_range(value_num, range) {
1009 return false;
1010 }
1011 continue;
1012 }
1013 }
1014 return false;
1016 }
1017 true
1018 }
1019
1020 fn value_in_range(&self, value: f64, range: &NumericRange) -> bool {
1022 if let Some(min) = range.min {
1023 if value < min {
1024 return false;
1025 }
1026 }
1027 if let Some(max) = range.max {
1028 if value > max {
1029 return false;
1030 }
1031 }
1032 true
1033 }
1034
1035 fn matches_geo_bounds(&self, record: &DWNRecord, bounds: &GeoBounds) -> bool {
1037 let lat = self
1039 .extract_numeric_field(record, "latitude")
1040 .or_else(|| self.extract_numeric_field(record, "lat"));
1041 let lng = self
1042 .extract_numeric_field(record, "longitude")
1043 .or_else(|| self.extract_numeric_field(record, "lng"));
1044
1045 if let (Some(lat), Some(lng)) = (lat, lng) {
1046 lat >= bounds.min_lat
1047 && lat <= bounds.max_lat
1048 && lng >= bounds.min_lng
1049 && lng <= bounds.max_lng
1050 } else {
1051 false
1052 }
1053 }
1054
1055 fn extract_numeric_field(&self, record: &DWNRecord, field: &str) -> Option<f64> {
1057 if let Some(value_str) = record.metadata.get(field) {
1059 if let Ok(value) = value_str.parse::<f64>() {
1060 return Some(value);
1061 }
1062 }
1063
1064 record.data.get(field).and_then(|v| v.as_f64())
1066 }
1067
1068 pub async fn batch_delete(&self, record_ids: Vec<String>) -> Web5Result<Vec<String>> {
1070 let mut deleted_ids = Vec::new();
1071 let mut errors = Vec::new();
1072
1073 for record_id in record_ids {
1074 match self.delete_record(&record_id) {
1075 Ok(_) => deleted_ids.push(record_id),
1076 Err(e) => errors.push(format!("Failed to delete {record_id}: {e}")),
1077 }
1078 }
1079
1080 if !errors.is_empty() {
1081 return Err(Web5Error::DWNError(format!(
1082 "Batch delete errors: {}",
1083 errors.join(", ")
1084 )));
1085 }
1086
1087 Ok(deleted_ids)
1088 }
1089
1090 pub async fn sync_records(&self, remote_endpoint: &str) -> Web5Result<Vec<SyncedDWNRecord>> {
1092 let storage = self
1095 .records
1096 .lock()
1097 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
1098 let synced_records: Vec<SyncedDWNRecord> = storage
1099 .values()
1100 .map(|record| SyncedDWNRecord {
1101 record: record.clone(),
1102 sync_status: SyncStatus::Synced,
1103 last_sync: current_timestamp(),
1104 sync_attempts: 1,
1105 })
1106 .collect();
1107
1108 println!(
1109 "Would sync {} records with remote endpoint: {}",
1110 synced_records.len(),
1111 remote_endpoint
1112 );
1113
1114 Ok(synced_records)
1115 }
1116
1117 pub fn resolve_conflicts(
1119 &self,
1120 conflicts: Vec<(DWNRecord, DWNRecord)>,
1121 strategy: ConflictResolution,
1122 ) -> Web5Result<Vec<DWNRecord>> {
1123 let mut resolved = Vec::new();
1124
1125 for (local, remote) in conflicts {
1126 let winner = match strategy {
1127 ConflictResolution::LastWriteWins => {
1128 let local_timestamp = local
1129 .metadata
1130 .get("updated")
1131 .and_then(|s| s.parse::<u64>().ok())
1132 .unwrap_or(0);
1133 let remote_timestamp = remote
1134 .metadata
1135 .get("updated")
1136 .and_then(|s| s.parse::<u64>().ok())
1137 .unwrap_or(0);
1138
1139 if remote_timestamp > local_timestamp {
1140 remote
1141 } else {
1142 local
1143 }
1144 }
1145 ConflictResolution::FirstWriteWins => {
1146 let local_timestamp = local
1147 .metadata
1148 .get("created")
1149 .and_then(|s| s.parse::<u64>().ok())
1150 .unwrap_or(u64::MAX);
1151 let remote_timestamp = remote
1152 .metadata
1153 .get("created")
1154 .and_then(|s| s.parse::<u64>().ok())
1155 .unwrap_or(u64::MAX);
1156
1157 if local_timestamp <= remote_timestamp {
1158 local
1159 } else {
1160 remote
1161 }
1162 }
1163 ConflictResolution::Manual => {
1164 local
1167 }
1168 ConflictResolution::Custom(ref _strategy) => {
1169 local
1171 }
1172 };
1173
1174 resolved.push(winner);
1175 }
1176
1177 Ok(resolved)
1178 }
1179
1180 pub fn export_records(
1182 &self,
1183 format: &str,
1184 filter: Option<DWNQueryFilter>,
1185 ) -> Web5Result<String> {
1186 let records = if let Some(filter) = filter {
1187 self.query_with_filter(filter)?
1188 } else {
1189 let storage = self
1190 .records
1191 .lock()
1192 .map_err(|e| Web5Error::Storage(format!("Failed to acquire lock: {e}")))?;
1193 storage.values().cloned().collect()
1194 };
1195
1196 match format.to_lowercase().as_str() {
1197 "json" => serde_json::to_string_pretty(&records)
1198 .map_err(|e| Web5Error::SerializationError(e.to_string())),
1199 "csv" => {
1200 let mut csv_output = String::from("id,owner,schema,created_at\n");
1201 for record in records {
1202 let empty_string = "".to_string();
1203 let created_at = record.metadata.get("created_at").unwrap_or(&empty_string);
1204 csv_output.push_str(&format!(
1205 "{},{},{},{}\n",
1206 record.id, record.owner, record.schema, created_at
1207 ));
1208 }
1209 Ok(csv_output)
1210 }
1211 "xml" => {
1212 let mut xml_output =
1213 String::from("<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n<records>\n");
1214 for record in records {
1215 xml_output.push_str(&format!(
1216 " <record id=\"{}\" owner=\"{}\" schema=\"{}\">\n",
1217 record.id, record.owner, record.schema
1218 ));
1219 xml_output.push_str(" <data><![CDATA[");
1220 xml_output.push_str(&serde_json::to_string(&record.data).unwrap_or_default());
1221 xml_output.push_str("]]></data>\n");
1222 xml_output.push_str(" </record>\n");
1223 }
1224 xml_output.push_str("</records>");
1225 Ok(xml_output)
1226 }
1227 _ => Err(Web5Error::DWNError(format!(
1228 "Unsupported export format: {format}"
1229 ))),
1230 }
1231 }
1232
1233 pub fn import_records(&self, data: &str, format: &str) -> Web5Result<Vec<String>> {
1235 let records = match format.to_lowercase().as_str() {
1236 "json" => serde_json::from_str::<Vec<DWNRecord>>(data)
1237 .map_err(|e| Web5Error::SerializationError(e.to_string()))?,
1238 _ => {
1239 return Err(Web5Error::DWNError(format!(
1240 "Unsupported import format: {format}"
1241 )))
1242 }
1243 };
1244
1245 let mut imported_ids = Vec::new();
1246 for record in records {
1247 let id = self.store_record(record)?;
1248 imported_ids.push(id);
1249 }
1250
1251 Ok(imported_ids)
1252 }
1253}
1254
1255fn current_timestamp() -> u64 {
1257 SystemTime::now()
1258 .duration_since(UNIX_EPOCH)
1259 .unwrap_or_default()
1260 .as_secs()
1261}
1262
1263#[cfg(test)]
1264mod tests {
1265 use super::*;
1267
1268 #[test]
1269 fn test_store_record() -> Result<(), Box<dyn std::error::Error>> {
1270 let dwn_manager = DWNManager::new();
1271
1272 let record = DWNRecord {
1273 id: "record1".to_string(),
1274 owner: "did:ion:123".to_string(),
1275 schema: "https://schema.org/Person".to_string(),
1276 data: serde_json::json!({
1277 "name": "Alice",
1278 "email": "alice@example.com"
1279 }),
1280 metadata: HashMap::new(),
1281 attestations: Vec::new(),
1282 };
1283
1284 let id = dwn_manager.store_record(record.clone())?;
1285 assert_eq!(id, "record1");
1286
1287 let records = dwn_manager.query_records("did:ion:123", "https://schema.org/Person")?;
1288 assert_eq!(records.len(), 1);
1289 assert_eq!(records[0].id, "record1");
1290 assert_eq!(records[0].owner, "did:ion:123");
1291
1292 Ok(())
1293 }
1294
1295 #[test]
1296 fn test_create_and_read_record() -> Result<(), Box<dyn std::error::Error>> {
1297 let dwn_manager = DWNManager::new();
1298
1299 let data = serde_json::json!({
1300 "name": "Bob",
1301 "email": "bob@example.com"
1302 });
1303
1304 let id =
1305 dwn_manager.create_record("did:ion:456", "https://schema.org/Person", data.clone())?;
1306
1307 let record = dwn_manager.read_record(&id)?;
1308 assert_eq!(record.owner, "did:ion:456");
1309 assert_eq!(record.schema, "https://schema.org/Person");
1310 assert_eq!(record.data, data);
1311 Ok(())
1312 }
1313
1314 #[test]
1315 fn test_update_record() -> Result<(), Box<dyn std::error::Error>> {
1316 let dwn_manager = DWNManager::new();
1317
1318 let data = serde_json::json!({
1319 "name": "Charlie",
1320 "email": "charlie@example.com"
1321 });
1322
1323 let id =
1324 dwn_manager.create_record("did:ion:789", "https://schema.org/Person", data.clone())?;
1325
1326 let new_data = serde_json::json!({
1327 "name": "Charlie",
1328 "email": "charlie.updated@example.com"
1329 });
1330
1331 dwn_manager.update_record(&id, new_data.clone())?;
1332
1333 let record = dwn_manager.read_record(&id)?;
1334 assert_eq!(record.data, new_data);
1335
1336 Ok(())
1337 }
1338
1339 #[test]
1340 fn test_delete_record() -> Result<(), Box<dyn std::error::Error>> {
1341 let dwn_manager = DWNManager::new();
1342
1343 let data = serde_json::json!({
1344 "name": "Dave",
1345 "email": "dave@example.com"
1346 });
1347
1348 let id =
1349 dwn_manager.create_record("did:ion:abc", "https://schema.org/Person", data.clone())?;
1350
1351 dwn_manager.delete_record(&id)?;
1352
1353 let result = dwn_manager.read_record(&id);
1354 assert!(result.is_err());
1355
1356 Ok(())
1357 }
1358}
1359
1360#[cfg(test)]
1361mod advanced_tests {
1362 use super::*;
1363
1364 #[test]
1365 fn test_advanced_query_with_pagination() -> Result<(), Box<dyn std::error::Error>> {
1366 let dwn_manager = DWNManager::new();
1367
1368 for i in 0..25 {
1370 let record = DWNRecord {
1371 id: format!("record_{:02}", i),
1372 owner: "did:ion:test".to_string(),
1373 schema: "test/schema".to_string(),
1374 data: serde_json::json!({
1375 "name": format!("Test Record {}", i),
1376 "value": i,
1377 }),
1378 metadata: {
1379 let mut meta = HashMap::new();
1380 meta.insert(
1381 "created_at".to_string(),
1382 (1640000000 + i as u64).to_string(),
1383 );
1384 meta.insert(
1385 "category".to_string(),
1386 if i % 2 == 0 {
1387 "even".to_string()
1388 } else {
1389 "odd".to_string()
1390 },
1391 );
1392 meta
1393 },
1394 attestations: Vec::new(),
1395 };
1396 dwn_manager.store_record(record)?;
1397 }
1398
1399 let filter = AdvancedDWNQueryFilter {
1401 base: DWNQueryFilter {
1402 owner: Some("did:ion:test".to_string()),
1403 schema: Some("test/schema".to_string()),
1404 metadata: None,
1405 date_range: None,
1406 data_filter: None,
1407 },
1408 search: None,
1409 geo_bounds: None,
1410 tags: None,
1411 numeric_ranges: None,
1412 };
1413
1414 let pagination = DWNQueryPagination {
1415 offset: Some(10),
1416 limit: Some(5),
1417 cursor: None,
1418 };
1419
1420 let result = dwn_manager.query_with_pagination(filter, Some(pagination))?;
1421 assert_eq!(result.records.len(), 5);
1422 assert_eq!(result.pagination.total, 25);
1423 assert_eq!(result.pagination.count, 5);
1424 assert!(result.pagination.has_more);
1425 assert!(result.pagination.next_cursor.is_some());
1426
1427 Ok(())
1428 }
1429
1430 #[test]
1431 fn test_export_import_records() -> Result<(), Box<dyn std::error::Error>> {
1432 let dwn_manager = DWNManager::new();
1433
1434 let record = DWNRecord {
1436 id: "export_test".to_string(),
1437 owner: "did:ion:test".to_string(),
1438 schema: "test/export".to_string(),
1439 data: serde_json::json!({"test": "data"}),
1440 metadata: HashMap::new(),
1441 attestations: Vec::new(),
1442 };
1443
1444 dwn_manager.store_record(record)?;
1445
1446 let exported = dwn_manager.export_records("json", None)?;
1448 assert!(exported.contains("export_test"));
1449
1450 let csv_exported = dwn_manager.export_records("csv", None)?;
1452 assert!(csv_exported.contains("export_test"));
1453
1454 Ok(())
1455 }
1456}
1457
1458#[derive(Debug, Clone, Serialize, Deserialize)]
1462pub enum AggregationStage {
1463 Match(HashMap<String, serde_json::Value>),
1465 Group {
1467 id: String,
1469 fields: HashMap<String, String>,
1471 },
1472 Sort(Vec<SortField>),
1474 Limit(usize),
1476 Skip(usize),
1478}
1479
1480#[derive(Debug, Clone, Serialize, Deserialize)]
1484pub struct SortField {
1485 pub field: String,
1487 pub ascending: bool,
1489}
1490
1491#[allow(dead_code)]
1493trait DurationExt {
1494 fn from_mins(mins: u64) -> Duration;
1495 fn from_hours(hours: u64) -> Duration;
1496}
1497
1498impl DurationExt for Duration {
1499 fn from_mins(mins: u64) -> Duration {
1500 Duration::from_secs(mins * 60)
1501 }
1502
1503 fn from_hours(hours: u64) -> Duration {
1504 Duration::from_secs(hours * 3600)
1505 }
1506}