Skip to main content

context69_contracts/
sources.rs

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}