use crate::{SubmissionId, SubmissionState, TaskExecutionIntent, TaskOperation, TaskSubmissionPayload};
use newton_submission_protocol::{ExecutionId, ExecutionProgress, ExecutionRequest};
use std::{error::Error, future::Future};
#[derive(Debug, Clone)]
pub struct PendingTaskRecord {
pub submission_id: SubmissionId,
pub operation: TaskOperation,
pub payload: TaskSubmissionPayload,
pub accepted_at_ms: i64,
pub deadline_at_ms: i64,
pub effect_retry_count: u32,
pub max_effect_retries: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskPlanState {
Planned,
Submitted,
Projected,
Acknowledged,
}
impl TaskPlanState {
pub const fn as_str(self) -> &'static str {
match self {
Self::Planned => "planned",
Self::Submitted => "submitted",
Self::Projected => "projected",
Self::Acknowledged => "acknowledged",
}
}
}
#[derive(Debug, Clone)]
pub struct TaskPlanMemberRecord {
pub submission_id: SubmissionId,
pub effect_retry_count: u32,
pub max_effect_retries: u32,
}
#[derive(Debug, Clone)]
pub struct TaskPlanRecord {
pub plan_id: ExecutionId,
pub chain_id: u64,
pub operation: TaskOperation,
pub intent: TaskExecutionIntent,
pub deadline_at_ms: Option<i64>,
pub state: TaskPlanState,
pub members: Vec<TaskPlanMemberRecord>,
}
impl TaskPlanRecord {
pub fn execution(&self) -> ExecutionRequest<TaskExecutionIntent> {
ExecutionRequest {
execution_id: self.plan_id,
chain_id: self.chain_id,
intent: self.intent.clone(),
deadline_at_ms: self.deadline_at_ms,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct TaskPlanningSnapshot {
pub pending: Vec<PendingTaskRecord>,
pub plans: Vec<TaskPlanRecord>,
}
#[derive(Debug, Clone)]
pub struct TaskProjection {
pub submission_id: SubmissionId,
pub state: SubmissionState,
pub operation: Option<TaskOperation>,
pub increment_effect_retry: bool,
pub terminal_error: Option<String>,
}
#[derive(Debug, Clone)]
pub enum TaskPlanningWrite {
InsertPlan(TaskPlanRecord),
MarkSubmitted(ExecutionId),
ProjectProgress {
plan_id: ExecutionId,
state: SubmissionState,
progress: ExecutionProgress,
},
ProjectOutcome {
plan_id: ExecutionId,
progress: ExecutionProgress,
projections: Vec<TaskProjection>,
},
MarkAcknowledged(ExecutionId),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskPlanningCommit {
Applied,
Stale,
}
pub trait TaskPlanningStore: std::fmt::Debug + Send + Sync {
type Error: Error + Send + Sync + 'static;
fn load(
&self,
chain_id: u64,
pending_limit: usize,
plan_limit: usize,
) -> impl Future<Output = Result<TaskPlanningSnapshot, Self::Error>> + Send;
fn commit(&self, write: &TaskPlanningWrite)
-> impl Future<Output = Result<TaskPlanningCommit, Self::Error>> + Send;
}