#![allow(missing_docs)]
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobState {
Queued,
Running,
Succeeded,
Failed,
Suspended,
Cancelled,
}
impl JobState {
pub fn is_terminal(self) -> bool {
matches!(
self,
Self::Succeeded | Self::Failed | Self::Suspended | Self::Cancelled
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobExecution {
Plan,
Run,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobFieldMap {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub id_field: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_field: Option<String>,
#[serde(
default,
deserialize_with = "crate::types::null_as_default",
skip_serializing_if = "Vec::is_empty"
)]
pub carry: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_type: Option<String>,
}
impl JobFieldMap {
pub(crate) fn is_empty(&self) -> bool {
*self == Self::default()
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobPreflight {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub estimated_credits: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub estimate_basis: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct JobChunk {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub seq: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub items: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub state: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub r#ref: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub units: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub credits_charged: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rate_book_version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<Value>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorJobValidation {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub identity: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sink: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorJobCapabilities {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub incremental_inference: bool,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub incremental_source_scan: bool,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub incremental_selection: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_scan: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_proof: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub checkpoint_profile: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub inference: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sink_targets: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ordering: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deletion_handling: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub publication: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorPlanOutputShape {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub result_kind: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub output_field: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub output_types: Vec<String>,
#[serde(default)]
pub dimensions: Option<u32>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ConnectorJobPlan {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub revision: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub expires_at: f64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub executable: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub executor_available: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub executor_availability: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub blocking_code: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub rows: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub mapped_bytes: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub input_bytes: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub eligible_count: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub eligible_count_quality: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub eligible_input_byte_count: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub matched_checkpoint_count: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub skipped_unchanged_count: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub deleted_preserved_count: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_dimensions: Option<u32>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub output: ConnectorPlanOutputShape,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cost_basis: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub max_reservation_credits: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub validation: ConnectorJobValidation,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub capabilities: ConnectorJobCapabilities,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorJobCheckpoint {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub profile_version: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub region: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub expected_generation: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub generation: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub published_revision: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorJobItemOutcomes {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub claimed: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub dispatched: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub inferred: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub staged: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub published: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub failed: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub reexecution_required: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub skipped_unchanged: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
pub struct ConnectorJobPublication {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub attempt_ordinal: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub revision: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub published: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub skipped_unchanged: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub deleted: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub failed: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub reexecuted: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub committed_at: f64,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConnectorJobOverlapOwner {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub job_id: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub attempt_ordinal: u64,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ConnectorJobAttempt {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub ordinal: u64,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub action: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub state: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub recovery_attempt_ordinal: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_code: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub replayed: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub billed_credits: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub overlap_owner: Option<ConnectorJobOverlapOwner>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub item_outcomes: ConnectorJobItemOutcomes,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub publication: Option<ConnectorJobPublication>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub created_at: f64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished_at: Option<f64>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
pub struct ConnectorJobRepairWindow {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expires_at: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempts_used: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempts_remaining: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempts_max: Option<u32>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ConnectorJobRecovery {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub required: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub state: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_code: Option<String>,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub reexecution_required: bool,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub repair: ConnectorJobRepairWindow,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct JobStatus {
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub id: String,
#[serde(
default,
deserialize_with = "crate::types::null_as_default",
rename = "object"
)]
pub object_kind: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub operation: String,
#[serde(default, deserialize_with = "crate::types::null_as_default")]
pub model: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub state: Option<JobState>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution: Option<JobExecution>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub phase: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_code: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub total_items: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub completed_items: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub chunks: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preflight: Option<JobPreflight>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub settled_credits: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plan_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plan_expires_at: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_expires_at: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plan: Option<ConnectorJobPlan>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub checkpoint: Option<ConnectorJobCheckpoint>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempt: Option<ConnectorJobAttempt>,
#[serde(
default,
deserialize_with = "crate::types::null_as_default",
skip_serializing_if = "Vec::is_empty"
)]
pub attempts: Vec<ConnectorJobAttempt>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub publication: Option<ConnectorJobPublication>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub recovery: Option<ConnectorJobRecovery>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub created_at: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished_at: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<Value>,
}
impl JobStatus {
pub fn chunks(&self) -> Vec<JobChunk> {
self.output
.as_ref()
.and_then(|output| output.get("chunks"))
.and_then(|chunks| serde_json::from_value(chunks.clone()).ok())
.unwrap_or_default()
}
pub fn is_settled(&self) -> bool {
self.state.is_some_and(JobState::is_terminal) || self.phase.as_deref() == Some("planned")
}
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct JobResultItem {
pub id: Option<String>,
pub success: Option<bool>,
pub units: Option<Value>,
pub dims: Option<u32>,
pub dense: Option<Vec<f32>>,
pub error: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct JobResults {
pub job_id: String,
pub state: Option<JobState>,
pub total_items: Option<u64>,
pub settled_credits: Option<u64>,
pub chunks: Vec<JobChunk>,
pub retrieved: usize,
pub dims: Option<u32>,
pub items: Vec<JobResultItem>,
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn terminal_states_are_the_ones_that_stop_a_poll() {
assert!(!JobState::Queued.is_terminal());
assert!(!JobState::Running.is_terminal());
for state in [
JobState::Succeeded,
JobState::Failed,
JobState::Suspended,
JobState::Cancelled,
] {
assert!(state.is_terminal(), "{state:?}");
}
}
#[test]
fn a_planned_connector_job_counts_as_settled() {
let planned = JobStatus {
state: Some(JobState::Running),
phase: Some("planned".to_string()),
..JobStatus::default()
};
assert!(planned.is_settled());
let running = JobStatus {
state: Some(JobState::Running),
..JobStatus::default()
};
assert!(!running.is_settled());
}
#[test]
fn chunks_are_read_out_of_the_output_envelope() {
let status = JobStatus {
output: Some(json!({"chunks": [
{"seq": 0, "items": 10, "state": "succeeded", "ref": "https://store/chunk-0"},
{"seq": 1, "items": 4, "state": "failed"}
]})),
..JobStatus::default()
};
let chunks = status.chunks();
assert_eq!(chunks.len(), 2);
assert_eq!(chunks[0].r#ref.as_deref(), Some("https://store/chunk-0"));
assert_eq!(chunks[1].r#ref, None);
assert!(JobStatus::default().chunks().is_empty());
}
#[test]
fn job_status_tolerates_a_minimal_payload() {
let status: JobStatus = serde_json::from_value(json!({
"id": "job_1", "object": "job", "operation": "encode", "model": "m", "state": "queued"
}))
.unwrap();
assert_eq!(status.state, Some(JobState::Queued));
assert!(status.plan.is_none());
assert!(status.attempts.is_empty());
}
}