Skip to main content

context69_contracts/
tasks.rs

1use chrono::{DateTime, Utc};
2use schemars::JsonSchema;
3use serde::{Deserialize, Serialize};
4use utoipa::{IntoParams, ToSchema};
5use uuid::Uuid;
6
7use crate::{
8    GroupResponse, ImportLibraryFileFromUrlRequest, LibraryFileUploadMetadata,
9    UpsertLibraryTextRequest,
10};
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
13#[serde(rename_all = "snake_case")]
14pub enum TaskKind {
15    SourceSync,
16    TextBatch,
17    FileBatch,
18    UrlBatch,
19    DeleteBatch,
20    Translation,
21    VectorRebuild,
22}
23
24impl TaskKind {
25    pub fn as_str(self) -> &'static str {
26        match self {
27            Self::SourceSync => "source_sync",
28            Self::TextBatch => "text_batch",
29            Self::FileBatch => "file_batch",
30            Self::UrlBatch => "url_batch",
31            Self::DeleteBatch => "delete_batch",
32            Self::Translation => "translation",
33            Self::VectorRebuild => "vector_rebuild",
34        }
35    }
36}
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
39#[serde(rename_all = "snake_case")]
40pub enum TaskStatus {
41    Queued,
42    Running,
43    Waiting,
44    Succeeded,
45    Failed,
46    Cancelled,
47}
48
49impl TaskStatus {
50    pub fn as_str(self) -> &'static str {
51        match self {
52            Self::Queued => "queued",
53            Self::Running => "running",
54            Self::Waiting => "waiting",
55            Self::Succeeded => "succeeded",
56            Self::Failed => "failed",
57            Self::Cancelled => "cancelled",
58        }
59    }
60}
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
63#[serde(rename_all = "snake_case")]
64pub enum TaskItemStatus {
65    Queued,
66    Running,
67    Waiting,
68    Succeeded,
69    Failed,
70    Cancelled,
71}
72
73impl TaskItemStatus {
74    pub fn as_str(self) -> &'static str {
75        match self {
76            Self::Queued => "queued",
77            Self::Running => "running",
78            Self::Waiting => "waiting",
79            Self::Succeeded => "succeeded",
80            Self::Failed => "failed",
81            Self::Cancelled => "cancelled",
82        }
83    }
84}
85
86#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
87pub struct TaskRef {
88    pub task_id: Uuid,
89    #[serde(default)]
90    pub item_ids: Vec<Uuid>,
91}
92
93#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
94pub struct TaskProgress {
95    pub total: i64,
96    pub queued: i64,
97    pub running: i64,
98    pub waiting: i64,
99    pub succeeded: i64,
100    pub failed: i64,
101    pub cancelled: i64,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
105#[serde(rename_all = "snake_case")]
106pub enum TaskOrigin {
107    Manual,
108    Rerun,
109}
110
111impl TaskOrigin {
112    pub fn as_str(self) -> &'static str {
113        match self {
114            Self::Manual => "manual",
115            Self::Rerun => "rerun",
116        }
117    }
118}
119
120#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
121pub struct TaskResponse {
122    pub task_id: Uuid,
123    pub kind: TaskKind,
124    pub status: TaskStatus,
125    pub origin: TaskOrigin,
126    pub group_path: Option<String>,
127    pub source_key: Option<String>,
128    pub stage: Option<String>,
129    pub waiting_reason: Option<String>,
130    pub dependency_key: Option<String>,
131    pub progress: TaskProgress,
132    pub failure_stage: Option<String>,
133    pub error_summary: Option<String>,
134    pub eta_seconds: Option<i64>,
135    pub created_at: DateTime<Utc>,
136    pub started_at: Option<DateTime<Utc>>,
137    pub finished_at: Option<DateTime<Utc>>,
138    pub updated_at: DateTime<Utc>,
139}
140
141#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
142pub struct ExternalJobInfo {
143    pub provider: String,
144    pub remote_task_id: String,
145    pub status: String,
146    pub remote_status: Option<String>,
147    pub submitted_at: DateTime<Utc>,
148    pub last_polled_at: Option<DateTime<Utc>>,
149    pub next_poll_at: Option<DateTime<Utc>>,
150    pub deadline_at: Option<DateTime<Utc>>,
151    pub error_message: Option<String>,
152}
153
154#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
155pub struct TaskItemResponse {
156    pub item_id: Uuid,
157    pub ordinal: i32,
158    pub status: TaskItemStatus,
159    pub resource_id: Option<String>,
160    pub file_id: Option<Uuid>,
161    pub stage: Option<String>,
162    pub waiting_reason: Option<String>,
163    pub dependency_key: Option<String>,
164    pub next_attempt_at: Option<DateTime<Utc>>,
165    pub failure_stage: Option<String>,
166    pub error_message: Option<String>,
167    pub attempt_count: i32,
168    pub retryable: bool,
169    pub created_at: DateTime<Utc>,
170    pub started_at: Option<DateTime<Utc>>,
171    pub finished_at: Option<DateTime<Utc>>,
172    #[serde(skip_serializing_if = "Option::is_none")]
173    pub external_job: Option<ExternalJobInfo>,
174}
175
176#[derive(Debug, Clone, Serialize, Deserialize, IntoParams, ToSchema)]
177#[into_params(parameter_in = Query)]
178pub struct TaskListQuery {
179    #[serde(default = "default_page")]
180    pub page: u32,
181    #[serde(default = "default_page_size")]
182    pub page_size: u32,
183    #[serde(default)]
184    pub query: Option<String>,
185    #[serde(default)]
186    pub kind: Option<TaskKind>,
187    #[serde(default)]
188    pub status: Option<TaskStatus>,
189    #[serde(default)]
190    pub stage: Option<String>,
191    #[serde(default)]
192    pub waiting_reason: Option<String>,
193    #[serde(default)]
194    pub dependency_key: Option<String>,
195}
196
197#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
198pub struct TaskPageResponse {
199    pub items: Vec<TaskResponse>,
200    pub pagination: crate::Pagination,
201}
202
203#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
204pub struct TaskItemsResponse {
205    pub items: Vec<TaskItemResponse>,
206    pub next_cursor: Option<String>,
207}
208
209#[derive(Debug, Clone, Serialize, Deserialize, IntoParams, ToSchema)]
210#[into_params(parameter_in = Query)]
211pub struct TaskItemsQuery {
212    #[serde(default = "default_item_limit")]
213    pub limit: u32,
214    #[serde(default)]
215    pub cursor: Option<String>,
216}
217
218#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
219pub struct ScopeSpec {
220    pub group_path: String,
221    pub name: String,
222    pub visibility: crate::Visibility,
223    #[serde(default)]
224    pub kind: Option<crate::GroupKind>,
225    #[serde(default)]
226    pub metadata_indexes: Vec<ScopeMetadataIndex>,
227}
228
229#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
230pub struct ScopeMetadataIndex {
231    pub source_key: String,
232    #[serde(flatten)]
233    pub definition: crate::CreateMetadataIndexRequest,
234}
235
236#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
237pub struct EnsureScopeResponse {
238    pub group: GroupResponse,
239    pub metadata_indexes: Vec<crate::MetadataIndexResponse>,
240}
241
242#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
243pub struct TextBatchRequest {
244    pub items: Vec<UpsertLibraryTextRequest>,
245}
246
247#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
248pub struct UrlBatchRequest {
249    pub items: Vec<ImportLibraryFileFromUrlRequest>,
250}
251
252#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
253pub struct DeleteBatchRequest {
254    pub items: Vec<crate::DocumentKey>,
255}
256
257#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
258pub struct FileBatchItem {
259    pub filename: String,
260    pub media_type: String,
261    pub content_base64: String,
262    #[serde(default, skip_serializing_if = "Option::is_none")]
263    pub declared_sha256: Option<String>,
264    #[serde(default)]
265    pub folder_id: Option<Uuid>,
266    #[serde(default)]
267    pub metadata: Option<LibraryFileUploadMetadata>,
268    #[serde(default)]
269    pub translation: Option<crate::TranslationDirective>,
270    #[serde(default)]
271    pub extraction: Option<crate::ExtractionDirective>,
272}
273
274#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
275pub struct FileBatchRequest {
276    pub items: Vec<FileBatchItem>,
277}
278
279#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
280#[serde(tag = "kind", rename_all = "snake_case")]
281pub enum TaskSubmitRequest {
282    /// Re-process existing library files that already have a file id.
283    RetryFileBatch {
284        #[serde(default)]
285        group_path: Option<String>,
286        items: Vec<FileRetryItem>,
287    },
288    /// Ingest file contents uploaded inline as base64.
289    FileBatch {
290        #[serde(default)]
291        group_path: Option<String>,
292        items: Vec<FileBatchItem>,
293    },
294    TextBatch {
295        #[serde(default)]
296        group_path: Option<String>,
297        items: Vec<UpsertLibraryTextRequest>,
298    },
299    UrlBatch {
300        #[serde(default)]
301        group_path: Option<String>,
302        items: Vec<ImportLibraryFileFromUrlRequest>,
303    },
304    DeleteBatch {
305        #[serde(default)]
306        group_path: Option<String>,
307        items: Vec<crate::DocumentKey>,
308    },
309    SourceSync {
310        #[serde(default)]
311        group_path: Option<String>,
312        source_key: String,
313    },
314    TranslationBatch {
315        #[serde(default)]
316        group_path: Option<String>,
317        items: Vec<TranslationSubmitItem>,
318    },
319    VectorRebuild,
320}
321
322#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
323pub struct FileRetryItem {
324    pub file_id: Uuid,
325}
326
327#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
328pub struct TranslationSubmitItem {
329    pub document_id: i64,
330    #[serde(default)]
331    pub target_locales: Vec<String>,
332}
333
334#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
335pub struct TaskRetryResponse {
336    pub task: TaskRef,
337    pub retried_items: i64,
338}
339
340#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
341pub struct RerunTaskResponse {
342    pub task: TaskRef,
343}
344
345#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema, JsonSchema)]
346#[serde(rename_all = "snake_case")]
347pub enum TaskPurgeMode {
348    Expired,
349    AllTerminal,
350}
351
352#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
353pub struct TaskMaintenanceSettings {
354    pub cleanup_enabled: bool,
355    pub retention_days: i64,
356    pub updated_at: DateTime<Utc>,
357}
358
359#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
360pub struct TaskMaintenanceStats {
361    pub total: i64,
362    pub queued: i64,
363    pub running: i64,
364    pub waiting: i64,
365    pub succeeded: i64,
366    pub failed: i64,
367    pub cancelled: i64,
368    pub active: i64,
369    pub expired_terminal: i64,
370}
371
372#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
373pub struct TaskMaintenanceOverview {
374    pub settings: TaskMaintenanceSettings,
375    pub stats: TaskMaintenanceStats,
376}
377
378#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
379pub struct UpdateTaskMaintenanceSettingsRequest {
380    pub cleanup_enabled: bool,
381    pub retention_days: i64,
382}
383
384#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
385pub struct CancelActiveTasksResponse {
386    pub cancelled_tasks: i64,
387}
388
389#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
390pub struct PurgeTasksRequest {
391    pub mode: TaskPurgeMode,
392}
393
394#[derive(Debug, Clone, Serialize, Deserialize, ToSchema, JsonSchema)]
395pub struct PurgeTasksResponse {
396    pub deleted_tasks: i64,
397}
398
399fn default_page() -> u32 {
400    1
401}
402
403fn default_page_size() -> u32 {
404    50
405}
406
407fn default_item_limit() -> u32 {
408    100
409}