1#![allow(missing_docs)]
6
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
12#[serde(rename_all = "snake_case")]
13pub enum JobState {
14 Queued,
15 Running,
16 Succeeded,
17 Failed,
18 Suspended,
19 Cancelled,
20}
21
22impl JobState {
23 pub fn is_terminal(self) -> bool {
25 matches!(
26 self,
27 Self::Succeeded | Self::Failed | Self::Suspended | Self::Cancelled
28 )
29 }
30}
31
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
34#[serde(rename_all = "snake_case")]
35pub enum JobExecution {
36 Plan,
38 Run,
40}
41
42#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
44pub struct JobFieldMap {
45 #[serde(default, skip_serializing_if = "Option::is_none")]
47 pub id_field: Option<String>,
48 #[serde(default, skip_serializing_if = "Option::is_none")]
50 pub input_field: Option<String>,
51 #[serde(
53 default,
54 deserialize_with = "crate::types::null_as_default",
55 skip_serializing_if = "Vec::is_empty"
56 )]
57 pub carry: Vec<String>,
58 #[serde(default, skip_serializing_if = "Option::is_none")]
60 pub input_type: Option<String>,
61}
62
63impl JobFieldMap {
64 pub(crate) fn is_empty(&self) -> bool {
65 *self == Self::default()
66 }
67}
68
69#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
71pub struct JobPreflight {
72 #[serde(default, skip_serializing_if = "Option::is_none")]
73 pub estimated_credits: Option<u64>,
74 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub estimate_basis: Option<String>,
76}
77
78#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
80pub struct JobChunk {
81 #[serde(default, deserialize_with = "crate::types::null_as_default")]
82 pub seq: u64,
83 #[serde(default, deserialize_with = "crate::types::null_as_default")]
84 pub items: u64,
85 #[serde(default, deserialize_with = "crate::types::null_as_default")]
86 pub state: String,
87 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub r#ref: Option<String>,
91 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub units: Option<u64>,
93 #[serde(default, skip_serializing_if = "Option::is_none")]
95 pub credits_charged: Option<u64>,
96 #[serde(default, skip_serializing_if = "Option::is_none")]
97 pub rate_book_version: Option<String>,
98 #[serde(default, skip_serializing_if = "Option::is_none")]
99 pub error: Option<Value>,
100}
101
102#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
104pub struct ConnectorJobValidation {
105 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub source: Option<String>,
107 #[serde(default, skip_serializing_if = "Option::is_none")]
108 pub identity: Option<String>,
109 #[serde(default, skip_serializing_if = "Option::is_none")]
110 pub sink: Option<String>,
111}
112
113#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
115pub struct ConnectorJobCapabilities {
116 #[serde(default, deserialize_with = "crate::types::null_as_default")]
117 pub incremental_inference: bool,
118 #[serde(default, deserialize_with = "crate::types::null_as_default")]
119 pub incremental_source_scan: bool,
120 #[serde(default, deserialize_with = "crate::types::null_as_default")]
121 pub incremental_selection: bool,
122 #[serde(default, skip_serializing_if = "Option::is_none")]
123 pub source_scan: Option<String>,
124 #[serde(default, skip_serializing_if = "Option::is_none")]
125 pub source_proof: Option<String>,
126 #[serde(default, skip_serializing_if = "Option::is_none")]
127 pub checkpoint_profile: Option<String>,
128 #[serde(default, skip_serializing_if = "Option::is_none")]
129 pub inference: Option<String>,
130 #[serde(default, skip_serializing_if = "Option::is_none")]
131 pub sink_targets: Option<Value>,
132 #[serde(default, skip_serializing_if = "Option::is_none")]
133 pub snapshot: Option<String>,
134 #[serde(default, skip_serializing_if = "Option::is_none")]
135 pub ordering: Option<String>,
136 #[serde(default, skip_serializing_if = "Option::is_none")]
137 pub deletion_handling: Option<String>,
138 #[serde(default, skip_serializing_if = "Option::is_none")]
139 pub publication: Option<String>,
140}
141
142#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
144pub struct ConnectorPlanOutputShape {
145 #[serde(default, deserialize_with = "crate::types::null_as_default")]
146 pub result_kind: String,
147 #[serde(default, deserialize_with = "crate::types::null_as_default")]
148 pub output_field: String,
149 #[serde(default, deserialize_with = "crate::types::null_as_default")]
150 pub output_types: Vec<String>,
151 #[serde(default)]
152 pub dimensions: Option<u32>,
153}
154
155#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
157pub struct ConnectorJobPlan {
158 #[serde(default, deserialize_with = "crate::types::null_as_default")]
159 pub revision: u64,
160 #[serde(default, deserialize_with = "crate::types::null_as_default")]
161 pub expires_at: f64,
162 #[serde(default, deserialize_with = "crate::types::null_as_default")]
163 pub executable: bool,
164 #[serde(default, skip_serializing_if = "Option::is_none")]
165 pub executor_available: Option<bool>,
166 #[serde(default, skip_serializing_if = "Option::is_none")]
167 pub executor_availability: Option<String>,
168 #[serde(default, skip_serializing_if = "Option::is_none")]
170 pub blocking_code: Option<String>,
171 #[serde(default, deserialize_with = "crate::types::null_as_default")]
172 pub rows: u64,
173 #[serde(default, deserialize_with = "crate::types::null_as_default")]
174 pub mapped_bytes: u64,
175 #[serde(default, deserialize_with = "crate::types::null_as_default")]
176 pub input_bytes: u64,
177 #[serde(default, deserialize_with = "crate::types::null_as_default")]
178 pub eligible_count: u64,
179 #[serde(default, skip_serializing_if = "Option::is_none")]
180 pub eligible_count_quality: Option<String>,
181 #[serde(default, deserialize_with = "crate::types::null_as_default")]
182 pub eligible_input_byte_count: u64,
183 #[serde(default, deserialize_with = "crate::types::null_as_default")]
184 pub matched_checkpoint_count: u64,
185 #[serde(default, deserialize_with = "crate::types::null_as_default")]
186 pub skipped_unchanged_count: u64,
187 #[serde(default, deserialize_with = "crate::types::null_as_default")]
188 pub deleted_preserved_count: u64,
189 #[serde(default, skip_serializing_if = "Option::is_none")]
190 pub output_dimensions: Option<u32>,
191 #[serde(default, deserialize_with = "crate::types::null_as_default")]
192 pub output: ConnectorPlanOutputShape,
193 #[serde(default, skip_serializing_if = "Option::is_none")]
194 pub cost_basis: Option<String>,
195 #[serde(default, deserialize_with = "crate::types::null_as_default")]
196 pub max_reservation_credits: u64,
197 #[serde(default, deserialize_with = "crate::types::null_as_default")]
198 pub validation: ConnectorJobValidation,
199 #[serde(default, deserialize_with = "crate::types::null_as_default")]
200 pub capabilities: ConnectorJobCapabilities,
201}
202
203#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
205pub struct ConnectorJobCheckpoint {
206 #[serde(default, skip_serializing_if = "Option::is_none")]
207 pub profile: Option<String>,
208 #[serde(default, deserialize_with = "crate::types::null_as_default")]
209 pub profile_version: u64,
210 #[serde(default, skip_serializing_if = "Option::is_none")]
211 pub region: Option<String>,
212 #[serde(default, deserialize_with = "crate::types::null_as_default")]
213 pub expected_generation: u64,
214 #[serde(default, deserialize_with = "crate::types::null_as_default")]
215 pub generation: u64,
216 #[serde(default, deserialize_with = "crate::types::null_as_default")]
217 pub published_revision: u64,
218}
219
220#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
222pub struct ConnectorJobItemOutcomes {
223 #[serde(default, deserialize_with = "crate::types::null_as_default")]
224 pub claimed: u64,
225 #[serde(default, deserialize_with = "crate::types::null_as_default")]
226 pub dispatched: u64,
227 #[serde(default, deserialize_with = "crate::types::null_as_default")]
228 pub inferred: u64,
229 #[serde(default, deserialize_with = "crate::types::null_as_default")]
230 pub staged: u64,
231 #[serde(default, deserialize_with = "crate::types::null_as_default")]
232 pub published: u64,
233 #[serde(default, deserialize_with = "crate::types::null_as_default")]
234 pub failed: u64,
235 #[serde(default, deserialize_with = "crate::types::null_as_default")]
236 pub reexecution_required: u64,
237 #[serde(default, deserialize_with = "crate::types::null_as_default")]
238 pub skipped_unchanged: u64,
239}
240
241#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
243pub struct ConnectorJobPublication {
244 #[serde(default, deserialize_with = "crate::types::null_as_default")]
245 pub attempt_ordinal: u64,
246 #[serde(default, deserialize_with = "crate::types::null_as_default")]
247 pub revision: u64,
248 #[serde(default, deserialize_with = "crate::types::null_as_default")]
249 pub published: u64,
250 #[serde(default, deserialize_with = "crate::types::null_as_default")]
251 pub skipped_unchanged: u64,
252 #[serde(default, deserialize_with = "crate::types::null_as_default")]
253 pub deleted: u64,
254 #[serde(default, deserialize_with = "crate::types::null_as_default")]
255 pub failed: u64,
256 #[serde(default, deserialize_with = "crate::types::null_as_default")]
257 pub reexecuted: u64,
258 #[serde(default, deserialize_with = "crate::types::null_as_default")]
259 pub committed_at: f64,
260}
261
262#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
264pub struct ConnectorJobOverlapOwner {
265 #[serde(default, deserialize_with = "crate::types::null_as_default")]
266 pub job_id: String,
267 #[serde(default, deserialize_with = "crate::types::null_as_default")]
268 pub attempt_ordinal: u64,
269}
270
271#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
273pub struct ConnectorJobAttempt {
274 #[serde(default, deserialize_with = "crate::types::null_as_default")]
275 pub ordinal: u64,
276 #[serde(default, deserialize_with = "crate::types::null_as_default")]
278 pub action: String,
279 #[serde(default, deserialize_with = "crate::types::null_as_default")]
280 pub state: String,
281 #[serde(default, skip_serializing_if = "Option::is_none")]
282 pub recovery_attempt_ordinal: Option<u64>,
283 #[serde(default, skip_serializing_if = "Option::is_none")]
284 pub outcome: Option<String>,
285 #[serde(default, skip_serializing_if = "Option::is_none")]
286 pub error_code: Option<String>,
287 #[serde(default, deserialize_with = "crate::types::null_as_default")]
288 pub replayed: bool,
289 #[serde(default, skip_serializing_if = "Option::is_none")]
290 pub billed_credits: Option<u64>,
291 #[serde(default, skip_serializing_if = "Option::is_none")]
292 pub overlap_owner: Option<ConnectorJobOverlapOwner>,
293 #[serde(default, deserialize_with = "crate::types::null_as_default")]
294 pub item_outcomes: ConnectorJobItemOutcomes,
295 #[serde(default, skip_serializing_if = "Option::is_none")]
296 pub publication: Option<ConnectorJobPublication>,
297 #[serde(default, deserialize_with = "crate::types::null_as_default")]
298 pub created_at: f64,
299 #[serde(default, skip_serializing_if = "Option::is_none")]
300 pub finished_at: Option<f64>,
301}
302
303#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
305pub struct ConnectorJobRepairWindow {
306 #[serde(default, skip_serializing_if = "Option::is_none")]
307 pub expires_at: Option<f64>,
308 #[serde(default, skip_serializing_if = "Option::is_none")]
309 pub attempts_used: Option<u32>,
310 #[serde(default, skip_serializing_if = "Option::is_none")]
311 pub attempts_remaining: Option<u32>,
312 #[serde(default, skip_serializing_if = "Option::is_none")]
313 pub attempts_max: Option<u32>,
314}
315
316#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
318pub struct ConnectorJobRecovery {
319 #[serde(default, deserialize_with = "crate::types::null_as_default")]
320 pub required: bool,
321 #[serde(default, skip_serializing_if = "Option::is_none")]
322 pub state: Option<String>,
323 #[serde(default, skip_serializing_if = "Option::is_none")]
324 pub outcome: Option<String>,
325 #[serde(default, skip_serializing_if = "Option::is_none")]
326 pub error_code: Option<String>,
327 #[serde(default, deserialize_with = "crate::types::null_as_default")]
328 pub reexecution_required: bool,
329 #[serde(default, deserialize_with = "crate::types::null_as_default")]
330 pub repair: ConnectorJobRepairWindow,
331}
332
333#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
338pub struct JobStatus {
339 #[serde(default, deserialize_with = "crate::types::null_as_default")]
340 pub id: String,
341 #[serde(
342 default,
343 deserialize_with = "crate::types::null_as_default",
344 rename = "object"
345 )]
346 pub object_kind: String,
347 #[serde(default, deserialize_with = "crate::types::null_as_default")]
348 pub operation: String,
349 #[serde(default, deserialize_with = "crate::types::null_as_default")]
350 pub model: String,
351 #[serde(default, skip_serializing_if = "Option::is_none")]
352 pub state: Option<JobState>,
353 #[serde(default, skip_serializing_if = "Option::is_none")]
354 pub execution: Option<JobExecution>,
355 #[serde(default, skip_serializing_if = "Option::is_none")]
356 pub phase: Option<String>,
357 #[serde(default, skip_serializing_if = "Option::is_none")]
358 pub outcome: Option<String>,
359 #[serde(default, skip_serializing_if = "Option::is_none")]
360 pub error_code: Option<String>,
361 #[serde(default, skip_serializing_if = "Option::is_none")]
362 pub total_items: Option<u64>,
363 #[serde(default, skip_serializing_if = "Option::is_none")]
364 pub completed_items: Option<u64>,
365 #[serde(default, skip_serializing_if = "Option::is_none")]
366 pub chunks: Option<u64>,
367 #[serde(default, skip_serializing_if = "Option::is_none")]
368 pub preflight: Option<JobPreflight>,
369 #[serde(default, skip_serializing_if = "Option::is_none")]
370 pub settled_credits: Option<u64>,
371 #[serde(default, skip_serializing_if = "Option::is_none")]
372 pub plan_revision: Option<u64>,
373 #[serde(default, skip_serializing_if = "Option::is_none")]
374 pub plan_expires_at: Option<f64>,
375 #[serde(default, skip_serializing_if = "Option::is_none")]
376 pub idempotency_expires_at: Option<f64>,
377 #[serde(default, skip_serializing_if = "Option::is_none")]
378 pub plan: Option<ConnectorJobPlan>,
379 #[serde(default, skip_serializing_if = "Option::is_none")]
380 pub checkpoint: Option<ConnectorJobCheckpoint>,
381 #[serde(default, skip_serializing_if = "Option::is_none")]
382 pub attempt: Option<ConnectorJobAttempt>,
383 #[serde(
384 default,
385 deserialize_with = "crate::types::null_as_default",
386 skip_serializing_if = "Vec::is_empty"
387 )]
388 pub attempts: Vec<ConnectorJobAttempt>,
389 #[serde(default, skip_serializing_if = "Option::is_none")]
390 pub publication: Option<ConnectorJobPublication>,
391 #[serde(default, skip_serializing_if = "Option::is_none")]
392 pub recovery: Option<ConnectorJobRecovery>,
393 #[serde(default, skip_serializing_if = "Option::is_none")]
394 pub created_at: Option<f64>,
395 #[serde(default, skip_serializing_if = "Option::is_none")]
396 pub finished_at: Option<f64>,
397 #[serde(default, skip_serializing_if = "Option::is_none")]
399 pub output: Option<Value>,
400}
401
402impl JobStatus {
403 pub fn chunks(&self) -> Vec<JobChunk> {
405 self.output
406 .as_ref()
407 .and_then(|output| output.get("chunks"))
408 .and_then(|chunks| serde_json::from_value(chunks.clone()).ok())
409 .unwrap_or_default()
410 }
411
412 pub fn is_settled(&self) -> bool {
414 self.state.is_some_and(JobState::is_terminal) || self.phase.as_deref() == Some("planned")
415 }
416}
417
418#[derive(Debug, Clone, Default, PartialEq)]
420pub struct JobResultItem {
421 pub id: Option<String>,
422 pub success: Option<bool>,
423 pub units: Option<Value>,
425 pub dims: Option<u32>,
426 pub dense: Option<Vec<f32>>,
427 pub error: Option<String>,
429}
430
431#[derive(Debug, Clone, Default, PartialEq)]
433pub struct JobResults {
434 pub job_id: String,
435 pub state: Option<JobState>,
436 pub total_items: Option<u64>,
437 pub settled_credits: Option<u64>,
438 pub chunks: Vec<JobChunk>,
439 pub retrieved: usize,
441 pub dims: Option<u32>,
442 pub items: Vec<JobResultItem>,
443}
444
445#[cfg(test)]
446mod tests {
447 use super::*;
448 use serde_json::json;
449
450 #[test]
451 fn terminal_states_are_the_ones_that_stop_a_poll() {
452 assert!(!JobState::Queued.is_terminal());
453 assert!(!JobState::Running.is_terminal());
454 for state in [
455 JobState::Succeeded,
456 JobState::Failed,
457 JobState::Suspended,
458 JobState::Cancelled,
459 ] {
460 assert!(state.is_terminal(), "{state:?}");
461 }
462 }
463
464 #[test]
465 fn a_planned_connector_job_counts_as_settled() {
466 let planned = JobStatus {
467 state: Some(JobState::Running),
468 phase: Some("planned".to_string()),
469 ..JobStatus::default()
470 };
471 assert!(planned.is_settled());
472
473 let running = JobStatus {
474 state: Some(JobState::Running),
475 ..JobStatus::default()
476 };
477 assert!(!running.is_settled());
478 }
479
480 #[test]
481 fn chunks_are_read_out_of_the_output_envelope() {
482 let status = JobStatus {
483 output: Some(json!({"chunks": [
484 {"seq": 0, "items": 10, "state": "succeeded", "ref": "https://store/chunk-0"},
485 {"seq": 1, "items": 4, "state": "failed"}
486 ]})),
487 ..JobStatus::default()
488 };
489 let chunks = status.chunks();
490 assert_eq!(chunks.len(), 2);
491 assert_eq!(chunks[0].r#ref.as_deref(), Some("https://store/chunk-0"));
492 assert_eq!(chunks[1].r#ref, None);
493 assert!(JobStatus::default().chunks().is_empty());
494 }
495
496 #[test]
497 fn job_status_tolerates_a_minimal_payload() {
498 let status: JobStatus = serde_json::from_value(json!({
499 "id": "job_1", "object": "job", "operation": "encode", "model": "m", "state": "queued"
500 }))
501 .unwrap();
502 assert_eq!(status.state, Some(JobState::Queued));
503 assert!(status.plan.is_none());
504 assert!(status.attempts.is_empty());
505 }
506}