1use chrono::{DateTime, Utc};
2use schemars::JsonSchema;
3use serde::{Deserialize, Serialize};
4use utoipa::{IntoParams, ToSchema};
5use uuid::Uuid;
6
7use crate::Pagination;
8use crate::Visibility;
9
10#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
11#[serde(rename_all = "snake_case")]
12pub enum SourceOriginStatusKind {
13 Unknown,
14 Connected,
15 Unreachable,
16 Misconfigured,
17}
18
19#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
20#[serde(rename_all = "snake_case")]
21pub enum SourceSyncStrategy {
22 Cursor,
23 FullScan,
24}
25
26impl SourceSyncStrategy {
27 pub fn as_str(self) -> &'static str {
28 match self {
29 Self::Cursor => "cursor",
30 Self::FullScan => "full_scan",
31 }
32 }
33}
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
36#[serde(rename_all = "snake_case")]
37pub enum SourceConnectorType {
38 PostgresSql,
39}
40
41impl SourceConnectorType {
42 pub fn as_str(self) -> &'static str {
43 match self {
44 Self::PostgresSql => "postgres_sql",
45 }
46 }
47}
48
49#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
50pub struct SourceStatus {
51 pub group_key: String,
52 pub group_path: String,
53 pub visibility: Visibility,
54 pub source_key: String,
55 pub display_name: String,
56 #[serde(default, skip_serializing_if = "Option::is_none")]
57 pub description: Option<String>,
58 #[serde(default)]
59 pub example_queries: Vec<String>,
60 pub connection: String,
61 pub has_database_url: bool,
62 pub origin_status: SourceOriginStatusKind,
63 #[serde(default, skip_serializing_if = "Option::is_none")]
64 pub origin_message: Option<String>,
65 pub sync_strategy: SourceSyncStrategy,
66 pub connector_type: SourceConnectorType,
67 pub base_query: String,
68 pub batch_size: i64,
69 pub last_cursor_updated_at: Option<DateTime<Utc>>,
70 pub last_cursor_external_id: Option<String>,
71 pub last_success_at: Option<DateTime<Utc>>,
72}
73
74#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
75pub struct ListSourcesResponse {
76 pub sources: Vec<SourceStatus>,
77}
78
79#[derive(Debug, Clone, Deserialize, IntoParams, ToSchema)]
80#[into_params(parameter_in = Query)]
81pub struct SourcePageQuery {
82 #[serde(default = "default_page")]
83 pub page: u32,
84 #[serde(default = "default_page_size")]
85 pub page_size: u32,
86 #[serde(default)]
87 pub query: Option<String>,
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
91pub struct SourcePageResponse {
92 pub items: Vec<SourceStatus>,
93 pub pagination: Pagination,
94}
95
96const fn default_page() -> u32 {
97 1
98}
99
100const fn default_page_size() -> u32 {
101 50
102}
103
104#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
105pub struct SyncOutcome {
106 pub records_seen: usize,
107 pub records_changed: usize,
108 pub chunks_upserted: usize,
109}
110
111#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
112pub struct SourceConfigInput {
113 #[serde(default, skip_serializing_if = "Option::is_none")]
114 pub source_id: Option<Uuid>,
115 pub source_key: String,
116 #[serde(default, skip_serializing_if = "Option::is_none")]
117 pub display_name: Option<String>,
118 #[serde(default, skip_serializing_if = "Option::is_none")]
119 pub description: Option<String>,
120 #[serde(default)]
121 pub example_queries: Vec<String>,
122 pub connection: String,
123 #[serde(default, skip_serializing_if = "Option::is_none")]
124 pub database_url: Option<String>,
125 pub sync_strategy: SourceSyncStrategy,
126 pub connector_type: SourceConnectorType,
127 pub base_query: String,
128 pub batch_size: i64,
129 #[serde(default, skip_serializing_if = "Option::is_none")]
130 pub visibility: Option<Visibility>,
131}
132
133#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
134pub struct SourceConnectionResponse {
135 pub name: String,
136 pub has_database_url: bool,
137 pub origin_status: SourceOriginStatusKind,
138 #[serde(default, skip_serializing_if = "Option::is_none")]
139 pub origin_message: Option<String>,
140}
141
142#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
143pub struct UpsertSourceConnectionRequest {
144 pub name: String,
145 #[serde(default, skip_serializing_if = "Option::is_none")]
146 pub database_url: Option<String>,
147}
148
149#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
150pub struct CreateSourceFolderRequest {
151 #[serde(default, skip_serializing_if = "Option::is_none")]
152 pub parent_folder_id: Option<Uuid>,
153 pub folder_name: String,
154 pub source_config: SourceConfigInput,
155}
156
157#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
158pub struct SourceFolderResponse {
159 pub folder_id: Uuid,
160 pub source_config_file_id: Uuid,
161 pub records_folder_id: Uuid,
162 pub path: String,
163}