Skip to main content

sie_sdk/types/
jobs.rs

1//! Batch jobs, including the connector-driven ones.
2
3// Wire-mirror types: field names are the API contract itself, and the ones whose
4// meaning is not obvious carry their own doc comment.
5#![allow(missing_docs)]
6
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9
10/// Where a job is in its lifecycle.
11#[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    /// Whether the job has stopped for good.
24    pub fn is_terminal(self) -> bool {
25        matches!(
26            self,
27            Self::Succeeded | Self::Failed | Self::Suspended | Self::Cancelled
28        )
29    }
30}
31
32/// Which half of a connector job's two-phase execution to run.
33#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
34#[serde(rename_all = "snake_case")]
35pub enum JobExecution {
36    /// Compute a plan and price it, without touching the sink.
37    Plan,
38    /// Execute a plan that was already computed.
39    Run,
40}
41
42/// How a connector source's rows map onto job items.
43#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
44pub struct JobFieldMap {
45    /// Source column holding the row's stable identifier.
46    #[serde(default, skip_serializing_if = "Option::is_none")]
47    pub id_field: Option<String>,
48    /// Source column holding the text or document to process.
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    pub input_field: Option<String>,
51    /// Source columns copied through to the sink unchanged.
52    #[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    /// How to interpret `input_field`: `text` or `document`.
59    #[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/// The cost the server estimated before starting the job.
70#[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/// One completed slice of a job's output.
79#[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    /// Where the chunk's msgpack payload is stored. Time-limited, and never the payload
88    /// itself.
89    #[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    /// Chunk charges sum exactly to the job's `settled_credits`.
94    #[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/// What a connector plan validated about its endpoints.
103#[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/// What the connector pair can do, which decides how much of a re-run can be skipped.
114#[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/// The shape a connector plan will write.
143#[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/// A priced, executable description of the work a connector job would do.
156#[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    /// Why the plan cannot execute, when it cannot.
169    #[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/// Where the sink's incremental state stands.
204#[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/// Per-item tallies for one attempt.
221#[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/// What one publication committed to the sink.
242#[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/// The job currently holding a sink region another job wants.
263#[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/// One execution or repair pass over a plan.
272#[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    /// `execute` or `repair`.
277    #[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/// How long a failed job may still be repaired.
304#[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/// Whether the job needs a repair pass, and what is left of its window.
317#[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/// A job as the server reports it.
334///
335/// The same shape answers submit, get, cancel, execute and repair. Connector source and
336/// sink URIs are deliberately absent: they can carry credentials.
337#[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    /// Holds `output.chunks`, the per-chunk result locations.
398    #[serde(default, skip_serializing_if = "Option::is_none")]
399    pub output: Option<Value>,
400}
401
402impl JobStatus {
403    /// The chunks recorded under `output`.
404    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    /// Whether the job has stopped for good, or has produced a plan awaiting execution.
413    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/// One decoded result row from a job chunk.
419#[derive(Debug, Clone, Default, PartialEq)]
420pub struct JobResultItem {
421    pub id: Option<String>,
422    pub success: Option<bool>,
423    /// Units billed for this row, as reported by the worker.
424    pub units: Option<Value>,
425    pub dims: Option<u32>,
426    pub dense: Option<Vec<f32>>,
427    /// The failure the worker reported, when this row failed.
428    pub error: Option<String>,
429}
430
431/// A job's results, with every retrievable chunk decoded.
432#[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    /// How many chunks were actually fetched; the rest had expired or had not succeeded.
440    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}