allsource_core/application/dto/
event_dto.rs1use crate::domain::entities::Event;
2use chrono::{DateTime, Utc};
3use serde::{Deserialize, Serialize};
4use uuid::Uuid;
5
6#[derive(Debug, Deserialize)]
8pub struct IngestEventRequest {
9 pub event_type: String,
10 pub entity_id: String,
11 pub tenant_id: Option<String>, pub payload: serde_json::Value,
13 pub metadata: Option<serde_json::Value>,
14 pub expected_version: Option<u64>,
17}
18
19#[derive(Debug, Serialize)]
21pub struct IngestEventResponse {
22 pub event_id: Uuid,
23 pub timestamp: DateTime<Utc>,
24 #[serde(skip_serializing_if = "Option::is_none")]
26 pub version: Option<u64>,
27}
28
29impl IngestEventResponse {
30 pub fn from_event(event: &Event) -> Self {
31 Self {
32 event_id: event.id(),
33 timestamp: event.timestamp(),
34 version: None,
35 }
36 }
37}
38
39#[derive(Debug, Default, Deserialize)]
41pub struct QueryEventsRequest {
42 pub entity_id: Option<String>,
44
45 pub event_type: Option<String>,
47
48 pub tenant_id: Option<String>,
50
51 pub as_of: Option<DateTime<Utc>>,
53
54 pub since: Option<DateTime<Utc>>,
56
57 pub until: Option<DateTime<Utc>>,
59
60 pub limit: Option<usize>,
62
63 pub event_type_prefix: Option<String>,
65
66 pub exclude_event_type_prefix: Option<String>,
71
72 pub payload_filter: Option<String>,
75}
76
77#[derive(Debug, Serialize)]
79pub struct QueryEventsResponse {
80 pub events: Vec<EventDto>,
81 pub count: usize,
82 pub total_count: usize,
83 pub has_more: bool,
84 #[serde(skip_serializing_if = "Option::is_none")]
86 pub entity_version: Option<u64>,
87 #[serde(skip_serializing_if = "Option::is_none")]
89 pub archive_integrity: Option<RetainedEntityIntegrity>,
90}
91
92#[derive(Debug, Serialize)]
93pub struct RetainedEntityIntegrity {
94 pub protocol: &'static str,
95 pub tenant_id: String,
96 pub entity_id: String,
97}
98
99#[derive(Debug, Serialize, Deserialize, Clone)]
101pub struct EventDto {
102 pub id: Uuid,
103 pub event_type: String,
104 pub entity_id: String,
105 pub tenant_id: String,
106 pub payload: serde_json::Value,
107 pub timestamp: DateTime<Utc>,
108 pub metadata: Option<serde_json::Value>,
109 pub version: i64,
110}
111
112impl From<&Event> for EventDto {
113 fn from(event: &Event) -> Self {
114 Self {
115 id: event.id(),
116 event_type: event.event_type().to_string(),
117 entity_id: event.entity_id().to_string(),
118 tenant_id: event.tenant_id().to_string(),
119 payload: event.payload().clone(),
120 timestamp: event.timestamp(),
121 metadata: event.metadata().cloned(),
122 version: event.version(),
123 }
124 }
125}
126
127impl From<Event> for EventDto {
128 fn from(event: Event) -> Self {
129 EventDto::from(&event)
130 }
131}
132
133#[derive(Debug, Default, Deserialize)]
135pub struct ListEntitiesRequest {
136 pub event_type_prefix: Option<String>,
138 pub payload_filter: Option<String>,
140 pub limit: Option<usize>,
142 pub offset: Option<usize>,
144 pub order: Option<String>,
149}
150
151#[derive(Debug, Serialize)]
153pub struct EntitySummary {
154 pub entity_id: String,
155 pub event_count: usize,
156 pub last_event_type: String,
157 pub last_event_at: DateTime<Utc>,
158}
159
160#[derive(Debug, Serialize)]
162pub struct ListEntitiesResponse {
163 pub entities: Vec<EntitySummary>,
164 pub total: usize,
165 pub has_more: bool,
166}
167
168#[derive(Debug, Deserialize)]
170pub struct DetectDuplicatesRequest {
171 pub event_type_prefix: String,
173 pub group_by: String,
175 pub limit: Option<usize>,
177 pub offset: Option<usize>,
179}
180
181#[derive(Debug, Serialize)]
183pub struct DuplicateGroup {
184 pub key: serde_json::Value,
186 pub entity_ids: Vec<String>,
188 pub count: usize,
190}
191
192#[derive(Debug, Serialize)]
194pub struct DetectDuplicatesResponse {
195 pub duplicates: Vec<DuplicateGroup>,
197 pub total: usize,
199 pub has_more: bool,
201}
202
203#[derive(Debug, Deserialize)]
205pub struct IngestEventsBatchRequest {
206 pub events: Vec<IngestEventRequest>,
207}
208
209#[derive(Debug, Serialize)]
211pub struct IngestEventsBatchResponse {
212 pub total: usize,
214 pub ingested: usize,
216 pub events: Vec<IngestEventResponse>,
218}