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}
88
89#[derive(Debug, Serialize, Deserialize, Clone)]
91pub struct EventDto {
92 pub id: Uuid,
93 pub event_type: String,
94 pub entity_id: String,
95 pub tenant_id: String,
96 pub payload: serde_json::Value,
97 pub timestamp: DateTime<Utc>,
98 pub metadata: Option<serde_json::Value>,
99 pub version: i64,
100}
101
102impl From<&Event> for EventDto {
103 fn from(event: &Event) -> Self {
104 Self {
105 id: event.id(),
106 event_type: event.event_type().to_string(),
107 entity_id: event.entity_id().to_string(),
108 tenant_id: event.tenant_id().to_string(),
109 payload: event.payload().clone(),
110 timestamp: event.timestamp(),
111 metadata: event.metadata().cloned(),
112 version: event.version(),
113 }
114 }
115}
116
117impl From<Event> for EventDto {
118 fn from(event: Event) -> Self {
119 EventDto::from(&event)
120 }
121}
122
123#[derive(Debug, Default, Deserialize)]
125pub struct ListEntitiesRequest {
126 pub event_type_prefix: Option<String>,
128 pub payload_filter: Option<String>,
130 pub limit: Option<usize>,
132 pub offset: Option<usize>,
134 pub order: Option<String>,
139}
140
141#[derive(Debug, Serialize)]
143pub struct EntitySummary {
144 pub entity_id: String,
145 pub event_count: usize,
146 pub last_event_type: String,
147 pub last_event_at: DateTime<Utc>,
148}
149
150#[derive(Debug, Serialize)]
152pub struct ListEntitiesResponse {
153 pub entities: Vec<EntitySummary>,
154 pub total: usize,
155 pub has_more: bool,
156}
157
158#[derive(Debug, Deserialize)]
160pub struct DetectDuplicatesRequest {
161 pub event_type_prefix: String,
163 pub group_by: String,
165 pub limit: Option<usize>,
167 pub offset: Option<usize>,
169}
170
171#[derive(Debug, Serialize)]
173pub struct DuplicateGroup {
174 pub key: serde_json::Value,
176 pub entity_ids: Vec<String>,
178 pub count: usize,
180}
181
182#[derive(Debug, Serialize)]
184pub struct DetectDuplicatesResponse {
185 pub duplicates: Vec<DuplicateGroup>,
187 pub total: usize,
189 pub has_more: bool,
191}
192
193#[derive(Debug, Deserialize)]
195pub struct IngestEventsBatchRequest {
196 pub events: Vec<IngestEventRequest>,
197}
198
199#[derive(Debug, Serialize)]
201pub struct IngestEventsBatchResponse {
202 pub total: usize,
204 pub ingested: usize,
206 pub events: Vec<IngestEventResponse>,
208}