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}
88
89/// DTO for a single event in responses
90#[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/// Request parameters for listing entities
124#[derive(Debug, Default, Deserialize)]
125pub struct ListEntitiesRequest {
126    /// Filter entities by event type prefix
127    pub event_type_prefix: Option<String>,
128    /// Filter by payload fields (JSON string)
129    pub payload_filter: Option<String>,
130    /// Limit number of entities returned
131    pub limit: Option<usize>,
132    /// Offset for pagination
133    pub offset: Option<usize>,
134    /// Sort direction by last-event time: `desc` (newest activity first, the
135    /// default) or `asc` (oldest activity first). Entities with the same
136    /// last-event time are tie-broken by `entity_id` ascending so offset
137    /// pagination is stable.
138    pub order: Option<String>,
139}
140
141/// A single entity summary in the list response
142#[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/// Response from listing entities
151#[derive(Debug, Serialize)]
152pub struct ListEntitiesResponse {
153    pub entities: Vec<EntitySummary>,
154    pub total: usize,
155    pub has_more: bool,
156}
157
158/// Request parameters for detecting duplicate entities
159#[derive(Debug, Deserialize)]
160pub struct DetectDuplicatesRequest {
161    /// Required: event type prefix to scope the search (no full-store scans)
162    pub event_type_prefix: String,
163    /// Comma-separated list of payload field names to group by (e.g., "name,user_id")
164    pub group_by: String,
165    /// Limit number of duplicate groups returned
166    pub limit: Option<usize>,
167    /// Offset for pagination
168    pub offset: Option<usize>,
169}
170
171/// A group of entities that share the same payload field values
172#[derive(Debug, Serialize)]
173pub struct DuplicateGroup {
174    /// The shared field values that define this group
175    pub key: serde_json::Value,
176    /// Entity IDs in this group
177    pub entity_ids: Vec<String>,
178    /// Number of entities in this group
179    pub count: usize,
180}
181
182/// Response from duplicate entity detection
183#[derive(Debug, Serialize)]
184pub struct DetectDuplicatesResponse {
185    /// Groups where count > 1
186    pub duplicates: Vec<DuplicateGroup>,
187    /// Total number of duplicate groups found
188    pub total: usize,
189    /// Whether more results are available
190    pub has_more: bool,
191}
192
193/// DTO for batch ingesting multiple events
194#[derive(Debug, Deserialize)]
195pub struct IngestEventsBatchRequest {
196    pub events: Vec<IngestEventRequest>,
197}
198
199/// DTO for batch ingestion response
200#[derive(Debug, Serialize)]
201pub struct IngestEventsBatchResponse {
202    /// Total number of events submitted
203    pub total: usize,
204    /// Number of events successfully ingested
205    pub ingested: usize,
206    /// Individual results for each event
207    pub events: Vec<IngestEventResponse>,
208}