Skip to main content

allsource_core/application/dto/
event_dto.rs

1use crate::domain::entities::Event;
2use chrono::{DateTime, Utc};
3use serde::{Deserialize, Serialize};
4use uuid::Uuid;
5
6/// DTO for ingesting a new event
7#[derive(Debug, Deserialize)]
8pub struct IngestEventRequest {
9    pub event_type: String,
10    pub entity_id: String,
11    pub tenant_id: Option<String>, // Optional, defaults to "default"
12    pub payload: serde_json::Value,
13    pub metadata: Option<serde_json::Value>,
14    /// Optional optimistic concurrency control: if set, the write is rejected
15    /// with 409 Conflict unless the entity's current version matches this value.
16    pub expected_version: Option<u64>,
17}
18
19/// DTO for event ingestion response
20#[derive(Debug, Serialize)]
21pub struct IngestEventResponse {
22    pub event_id: Uuid,
23    pub timestamp: DateTime<Utc>,
24    /// The entity's version after this event was appended
25    #[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/// DTO for querying events
40#[derive(Debug, Default, Deserialize)]
41pub struct QueryEventsRequest {
42    /// Filter by entity ID
43    pub entity_id: Option<String>,
44
45    /// Filter by event type
46    pub event_type: Option<String>,
47
48    /// Tenant ID (required for multi-tenancy)
49    pub tenant_id: Option<String>,
50
51    /// Time-travel: get events as of this timestamp
52    pub as_of: Option<DateTime<Utc>>,
53
54    /// Get events since this timestamp
55    pub since: Option<DateTime<Utc>>,
56
57    /// Get events until this timestamp
58    pub until: Option<DateTime<Utc>>,
59
60    /// Limit number of results
61    pub limit: Option<usize>,
62
63    /// Filter by event type prefix (e.g., "index." matches "index.created", "index.updated")
64    pub event_type_prefix: Option<String>,
65
66    /// Exclude events whose type starts with ANY of these prefixes
67    /// (comma-separated, e.g. `"audit.,service.,_system."`). Applied before the
68    /// limit so excluded events don't consume the result window — lets callers
69    /// drop high-frequency/operational namespaces from a recent-activity feed.
70    pub exclude_event_type_prefix: Option<String>,
71
72    /// Filter by payload fields (JSON string, e.g., `{"user_id":"abc-123"}`).
73    /// Matches events where payload contains ALL specified key-value pairs.
74    pub payload_filter: Option<String>,
75}
76
77/// DTO for query response
78#[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    /// Current version of the entity (present only when query filters by entity_id)
85    #[serde(skip_serializing_if = "Option::is_none")]
86    pub entity_version: Option<u64>,
87    /// Present only after a bounded, strict read of this retained entity.
88    #[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/// DTO for a single event in responses
100#[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/// Request parameters for listing entities
134#[derive(Debug, Default, Deserialize)]
135pub struct ListEntitiesRequest {
136    /// Filter entities by event type prefix
137    pub event_type_prefix: Option<String>,
138    /// Filter by payload fields (JSON string)
139    pub payload_filter: Option<String>,
140    /// Limit number of entities returned
141    pub limit: Option<usize>,
142    /// Offset for pagination
143    pub offset: Option<usize>,
144    /// Sort direction by last-event time: `desc` (newest activity first, the
145    /// default) or `asc` (oldest activity first). Entities with the same
146    /// last-event time are tie-broken by `entity_id` ascending so offset
147    /// pagination is stable.
148    pub order: Option<String>,
149}
150
151/// A single entity summary in the list response
152#[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/// Response from listing entities
161#[derive(Debug, Serialize)]
162pub struct ListEntitiesResponse {
163    pub entities: Vec<EntitySummary>,
164    pub total: usize,
165    pub has_more: bool,
166}
167
168/// Request parameters for detecting duplicate entities
169#[derive(Debug, Deserialize)]
170pub struct DetectDuplicatesRequest {
171    /// Required: event type prefix to scope the search (no full-store scans)
172    pub event_type_prefix: String,
173    /// Comma-separated list of payload field names to group by (e.g., "name,user_id")
174    pub group_by: String,
175    /// Limit number of duplicate groups returned
176    pub limit: Option<usize>,
177    /// Offset for pagination
178    pub offset: Option<usize>,
179}
180
181/// A group of entities that share the same payload field values
182#[derive(Debug, Serialize)]
183pub struct DuplicateGroup {
184    /// The shared field values that define this group
185    pub key: serde_json::Value,
186    /// Entity IDs in this group
187    pub entity_ids: Vec<String>,
188    /// Number of entities in this group
189    pub count: usize,
190}
191
192/// Response from duplicate entity detection
193#[derive(Debug, Serialize)]
194pub struct DetectDuplicatesResponse {
195    /// Groups where count > 1
196    pub duplicates: Vec<DuplicateGroup>,
197    /// Total number of duplicate groups found
198    pub total: usize,
199    /// Whether more results are available
200    pub has_more: bool,
201}
202
203/// DTO for batch ingesting multiple events
204#[derive(Debug, Deserialize)]
205pub struct IngestEventsBatchRequest {
206    pub events: Vec<IngestEventRequest>,
207}
208
209/// DTO for batch ingestion response
210#[derive(Debug, Serialize)]
211pub struct IngestEventsBatchResponse {
212    /// Total number of events submitted
213    pub total: usize,
214    /// Number of events successfully ingested
215    pub ingested: usize,
216    /// Individual results for each event
217    pub events: Vec<IngestEventResponse>,
218}