use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueEnqueueRequest {
pub stream: String,
pub payload: Vec<u8>,
#[serde(default)]
pub priority: u8,
#[serde(default)]
pub not_before_ms: u64,
#[serde(default)]
pub shard_key: Option<Vec<u8>>,
#[serde(default)]
pub dedup_key: Option<Vec<u8>>,
#[serde(default)]
pub max_attempts: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueEnqueueReply {
pub job_id: Option<u64>,
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueBatchEnqueueJob {
pub payload: Vec<u8>,
#[serde(default)]
pub priority: u8,
#[serde(default)]
pub not_before_ms: u64,
#[serde(default)]
pub shard_key: Option<Vec<u8>>,
#[serde(default)]
pub dedup_key: Option<Vec<u8>>,
#[serde(default)]
pub max_attempts: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueEnqueueBatchRequest {
pub stream: String,
pub jobs: Vec<QueueBatchEnqueueJob>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueEnqueueBatchReply {
pub job_ids: Vec<u64>,
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueLeaseRequest {
pub stream: String,
pub worker_node: u64,
pub worker_instance: u32,
pub max: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueLeasedJobWire {
pub lease_id: u64,
pub job_id: u64,
pub payload: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueLeaseReply {
pub jobs: Vec<QueueLeasedJobWire>,
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueAckRequest {
pub stream: String,
pub worker_node: u64,
pub worker_instance: u32,
pub lease_id: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueAckReply {
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueAckBatchRequest {
pub stream: String,
pub worker_node: u64,
pub worker_instance: u32,
pub lease_ids: Vec<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueAckBatchReply {
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueNackRequest {
pub stream: String,
pub worker_node: u64,
pub worker_instance: u32,
pub lease_id: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueNackReply {
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueMetricsRequest {
pub stream: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueMetricsReply {
pub pending: u64,
pub leased: u64,
pub dead_letter: u64,
pub oldest_pending_age_ms: u64,
pub error: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum QueueJobLifecycleWire {
Pending = 0,
Leased = 1,
Delayed = 2,
DeadLetter = 3,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueJobStatusRequest {
pub stream: String,
pub job_id: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueJobStatusReply {
pub found: bool,
pub job_id: u64,
pub lifecycle: Option<QueueJobLifecycleWire>,
pub payload_len: u64,
pub priority: u8,
pub leased_worker_node: Option<u64>,
pub leased_worker_instance: Option<u32>,
pub attempts: u32,
pub max_attempts: u32,
pub error: Option<String>,
}
const fn default_true() -> bool {
true
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RecurringScheduleWire {
pub name: String,
pub cron: String,
pub payload: Vec<u8>,
#[serde(default)]
pub priority: u8,
#[serde(default)]
pub max_attempts: u32,
#[serde(default = "default_true")]
pub enabled: bool,
#[serde(default)]
pub next_run_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum QueueReplicateOp {
Enqueue {
job_id: u64,
payload: Vec<u8>,
enqueued_at_ms: u64,
next_job_id: u64,
#[serde(default)]
priority: u8,
#[serde(default)]
not_before_ms: u64,
#[serde(default)]
dedup_key: Option<Vec<u8>>,
#[serde(default)]
attempts: u32,
#[serde(default)]
max_attempts: u32,
},
Lease {
lease_id: u64,
job_id: u64,
worker_node: u64,
worker_instance: u32,
expires_at_ms: u64,
next_lease_id: u64,
},
Ack {
lease_id: u64,
job_id: u64,
},
Nack {
lease_id: u64,
job_id: u64,
#[serde(default)]
attempts: u32,
#[serde(default)]
dead_letter: bool,
#[serde(default)]
not_before_ms: u64,
},
Reclaim {
lease_id: u64,
job_id: u64,
#[serde(default)]
attempts: u32,
#[serde(default)]
dead_letter: bool,
#[serde(default)]
not_before_ms: u64,
},
RequeueDeadLetter {
job_id: u64,
#[serde(default)]
attempts: u32,
},
UpsertSchedule {
schedule: RecurringScheduleWire,
},
UpdateScheduleNextRun {
name: String,
next_run_ms: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueReplicateRequest {
pub stream: String,
pub ops: Vec<QueueReplicateOp>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueReplicateReply {
pub error: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueRequeueDeadLetterRequest {
pub stream: String,
pub job_id: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueRequeueDeadLetterReply {
pub error: Option<String>,
}