use alloy::{
primitives::{keccak256, Bytes, B256},
sol_types::SolValue,
};
use newton_core::{
newton_prover_task_manager::{
INewtonPolicy, INewtonPolicyClient,
INewtonProverTaskManager::{Task, TaskResponse},
},
TaskId,
};
use newton_submission_protocol::{encode_message, ExecutionId};
use serde::{Deserialize, Serialize};
use std::{fmt, str::FromStr};
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct SubmissionId(Uuid);
impl SubmissionId {
pub fn new() -> Self {
Self(Uuid::new_v4())
}
pub const fn as_uuid(self) -> Uuid {
self.0
}
}
impl Default for SubmissionId {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for SubmissionId {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(formatter)
}
}
impl FromStr for SubmissionId {
type Err = uuid::Error;
fn from_str(value: &str) -> Result<Self, Self::Err> {
Uuid::parse_str(value).map(Self)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TaskOperation {
CombinedCreateAndRespond,
RespondOnly,
}
impl TaskOperation {
pub const fn as_str(self) -> &'static str {
match self {
Self::CombinedCreateAndRespond => "combined_create_and_respond",
Self::RespondOnly => "respond_only",
}
}
}
impl fmt::Display for TaskOperation {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
impl FromStr for TaskOperation {
type Err = ParseTaskEnumError;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"combined_create_and_respond" => Ok(Self::CombinedCreateAndRespond),
"respond_only" => Ok(Self::RespondOnly),
_ => Err(ParseTaskEnumError(value.to_string())),
}
}
}
#[derive(Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct TaskSubmissionPayload {
pub chain_id: u64,
pub operation: TaskOperation,
pub task_id: TaskId,
pub task_response_digest: B256,
pub task: Task,
pub task_response: TaskResponse,
pub signature_data: Bytes,
pub attestation_data: Bytes,
pub requested_at: u64,
}
impl fmt::Debug for TaskSubmissionPayload {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("TaskSubmissionPayload")
.field("chain_id", &self.chain_id)
.field("operation", &self.operation)
.field("task_id", &self.task_id)
.field("task_response_digest", &self.task_response_digest)
.field("signature_data_len", &self.signature_data.len())
.field("attestation_data_len", &self.attestation_data.len())
.field("requested_at", &self.requested_at)
.finish_non_exhaustive()
}
}
impl TaskSubmissionPayload {
pub fn validate_identity(&self) -> Result<(), PayloadError> {
if self.task.taskId != self.task_id {
return Err(PayloadError::TaskIdMismatch);
}
if self.task_response.taskId != self.task_id {
return Err(PayloadError::ResponseTaskIdMismatch);
}
if contract_response_hash(&self.task_response) != self.task_response_digest {
return Err(PayloadError::ResponseDigestMismatch);
}
if self.signature_data.is_empty() {
return Err(PayloadError::EmptySignatureData);
}
Ok(())
}
}
#[derive(Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct TaskSubmissionRequest {
pub producer_id: String,
pub idempotency_key: B256,
pub payload: TaskSubmissionPayload,
}
impl fmt::Debug for TaskSubmissionRequest {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("TaskSubmissionRequest")
.field("producer_id", &self.producer_id)
.field("idempotency_key", &self.idempotency_key)
.field("payload", &self.payload)
.finish()
}
}
impl TaskSubmissionRequest {
pub fn validate(&self) -> Result<(), PayloadError> {
if self.producer_id.trim().is_empty() {
return Err(PayloadError::EmptyProducerId);
}
self.payload.validate_identity()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SubmissionState {
BatchPending,
ReadyForSubmission,
Assigned,
Prepared,
Broadcast,
Mined,
Succeeded,
Failed,
}
impl SubmissionState {
pub const fn is_terminal(self) -> bool {
matches!(self, Self::Succeeded | Self::Failed)
}
pub const fn as_str(self) -> &'static str {
match self {
Self::BatchPending => "batch_pending",
Self::ReadyForSubmission => "ready_for_submission",
Self::Assigned => "assigned",
Self::Prepared => "prepared",
Self::Broadcast => "broadcast",
Self::Mined => "mined",
Self::Succeeded => "succeeded",
Self::Failed => "failed",
}
}
}
impl FromStr for SubmissionState {
type Err = ParseTaskEnumError;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"batch_pending" => Ok(Self::BatchPending),
"ready_for_submission" => Ok(Self::ReadyForSubmission),
"assigned" => Ok(Self::Assigned),
"prepared" => Ok(Self::Prepared),
"broadcast" => Ok(Self::Broadcast),
"mined" => Ok(Self::Mined),
"succeeded" => Ok(Self::Succeeded),
"failed" => Ok(Self::Failed),
_ => Err(ParseTaskEnumError(value.to_string())),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SubmissionResource {
pub submission_id: SubmissionId,
pub state: SubmissionState,
pub chain_id: u64,
pub operation: TaskOperation,
#[serde(alias = "jobId")]
pub execution_id: Option<ExecutionId>,
pub transaction_hashes: Vec<B256>,
pub mined_block: Option<u64>,
pub confirmations: u64,
pub task_created_block: u64,
pub deadline_at_ms: i64,
pub terminal_error: Option<String>,
pub accepted_at_ms: i64,
pub updated_at_ms: i64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(transparent)]
pub struct EventCursor(pub i64);
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SubmissionEvent {
pub cursor: EventCursor,
pub submission_id: SubmissionId,
pub producer_id: String,
pub chain_id: u64,
pub operation: TaskOperation,
pub task_id: TaskId,
pub state: SubmissionState,
#[serde(alias = "jobId")]
pub execution_id: Option<ExecutionId>,
pub transaction_hash: Option<B256>,
pub terminal_error: Option<String>,
pub created_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SubmissionEventPage {
pub events: Vec<SubmissionEvent>,
pub next_cursor: EventCursor,
pub latest_cursor: EventCursor,
pub retention_floor: EventCursor,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct CursorPrunedResponse {
pub code: String,
pub requested: EventCursor,
pub retention_floor: EventCursor,
}
pub fn normalized_payload_bytes(payload: &TaskSubmissionPayload) -> Result<Vec<u8>, PayloadError> {
encode_message(payload).map_err(|error| PayloadError::Encoding(error.to_string()))
}
pub fn payload_hash(payload: &TaskSubmissionPayload) -> Result<B256, PayloadError> {
normalized_payload_bytes(payload).map(|bytes| keccak256(&bytes))
}
pub fn contract_task_hash(task: &Task) -> B256 {
let policies_hash = keccak256(task.policies.abi_encode());
let wasm_args_hash = keccak256(task.wasmArgs.abi_encode());
keccak256(
(
task.taskId,
task.taskCreatedBlock,
task.intent.clone(),
task.intentSignature.clone(),
task.policyClient,
task.policyId,
task.policyRevision,
policies_hash,
wasm_args_hash,
task.quorumNumbers.clone(),
task.quorumThresholdPercentage,
task.initializationTimestamp,
)
.abi_encode_params(),
)
}
pub fn contract_response_hash(response: &TaskResponse) -> B256 {
keccak256(TaskResponse::abi_encode(response))
}
pub fn derive_idempotency_key(
producer_id: &[u8],
chain_id: u64,
operation: TaskOperation,
task_id: TaskId,
response_digest: B256,
) -> B256 {
let mut bytes = Vec::with_capacity(128 + producer_id.len());
bytes.extend_from_slice(b"newton-task-submission-v1");
bytes.extend_from_slice(producer_id);
bytes.extend_from_slice(&chain_id.to_be_bytes());
bytes.extend_from_slice(operation.as_str().as_bytes());
bytes.extend_from_slice(task_id.as_slice());
bytes.extend_from_slice(response_digest.as_slice());
keccak256(bytes)
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum PayloadError {
#[error("producer id must not be empty")]
EmptyProducerId,
#[error("task id does not match task.taskId")]
TaskIdMismatch,
#[error("task id does not match taskResponse.taskId")]
ResponseTaskIdMismatch,
#[error("task response digest does not match taskResponse")]
ResponseDigestMismatch,
#[error("signature data must not be empty")]
EmptySignatureData,
#[error("failed to encode normalized payload: {0}")]
Encoding(String),
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error("unknown task enum value '{0}'")]
pub struct ParseTaskEnumError(pub String);
#[cfg(test)]
mod tests {
use super::*;
use alloy::primitives::{Address, U256};
use newton_core::newton_prover_task_manager::{INewtonPolicy, NewtonMessage};
use newton_submission_protocol::decode_message;
fn golden_payload() -> TaskSubmissionPayload {
let task_id = B256::repeat_byte(1);
let policy_id = B256::repeat_byte(17);
let intent = NewtonMessage::Intent {
from: Address::repeat_byte(2),
to: Address::repeat_byte(3),
value: U256::from(4),
data: Bytes::from(vec![5, 6]),
chainId: U256::from(31_337),
functionSignature: Bytes::from(vec![7, 8, 9, 10]),
};
let task = Task {
taskId: task_id,
policyClient: Address::repeat_byte(11),
policyId: policy_id,
policyRevision: 1,
taskCreatedBlock: 12,
quorumThresholdPercentage: 67,
intent: intent.clone(),
intentSignature: Bytes::from(vec![13, 14]),
policies: vec![INewtonPolicyClient::PolicySpec {
policy: Address::repeat_byte(18),
config: INewtonPolicy::PolicyConfig {
policyParams: Bytes::from(vec![23]),
expireAfter: 24,
},
}],
wasmArgs: vec![Bytes::from(vec![15])],
quorumNumbers: Bytes::from(vec![0, 1]),
initializationTimestamp: U256::from(16),
};
let task_response = TaskResponse {
taskId: task_id,
policyClient: task.policyClient,
policyId: policy_id,
intent,
intentSignature: task.intentSignature.clone(),
allowed: true,
policyTaskData: vec![NewtonMessage::PolicyTaskData {
policyId: policy_id,
policyAddress: Address::repeat_byte(18),
policy: Bytes::from(vec![20]),
policyData: vec![NewtonMessage::PolicyData {
wasmArgs: Bytes::from(vec![15]),
data: Bytes::from(vec![21]),
expireBlock: 20,
}],
}],
initializationTimestamp: U256::from(16),
};
TaskSubmissionPayload {
chain_id: 31_337,
operation: TaskOperation::CombinedCreateAndRespond,
task_id,
task_response_digest: contract_response_hash(&task_response),
task,
task_response,
signature_data: Bytes::from(vec![25, 26]),
attestation_data: Bytes::from(vec![27]),
requested_at: 28,
}
}
#[test]
fn idempotency_key_binds_operation() {
let task_id = B256::repeat_byte(1);
let digest = B256::repeat_byte(2);
let combined = derive_idempotency_key(b"gateway", 1, TaskOperation::CombinedCreateAndRespond, task_id, digest);
let respond = derive_idempotency_key(b"gateway", 1, TaskOperation::RespondOnly, task_id, digest);
assert_ne!(combined, respond);
}
#[test]
fn terminal_states_are_closed() {
assert!(SubmissionState::Succeeded.is_terminal());
assert!(SubmissionState::Failed.is_terminal());
assert!(!SubmissionState::BatchPending.is_terminal());
assert!(!SubmissionState::Broadcast.is_terminal());
}
#[test]
fn payload_validation_binds_the_consensus_digest_to_response_bytes() {
let mut payload = golden_payload();
payload.task_response_digest = B256::ZERO;
assert_eq!(payload.validate_identity(), Err(PayloadError::ResponseDigestMismatch));
}
#[test]
fn task_payload_golden_vector_is_stable() {
let payload = golden_payload();
let normalized = normalized_payload_bytes(&payload).expect("normalized");
let decoded: TaskSubmissionPayload = decode_message(&normalized).expect("decode");
assert_eq!(
payload_hash(&decoded).expect("payload hash"),
payload_hash(&payload).expect("payload hash")
);
assert_eq!(
payload_hash(&payload).expect("payload hash"),
"0x36ade7e1abaa1abea54a4b0874419203e21422b5cdd66de4f7932ea83c25aec0"
.parse::<B256>()
.expect("hash")
);
assert_eq!(
contract_task_hash(&payload.task),
"0xe86e8439fbdbe7de292ef9130f7c9174c16d6ddbe9f17ad79da5432ca88d61da"
.parse::<B256>()
.expect("hash")
);
assert_eq!(
contract_response_hash(&payload.task_response),
"0x02a3ad5f5bf09e1a3ec417fe75493e646a3ca1909721832b31918d4b9b3cc17c"
.parse::<B256>()
.expect("hash")
);
assert_eq!(
derive_idempotency_key(
b"gateway-v1",
payload.chain_id,
payload.operation,
payload.task_id,
payload.task_response_digest
),
"0x975aa87fc522d4f770b72c503f6f60aac0fe088c30b204e949f1bbb6d9cb1dd6"
.parse::<B256>()
.expect("hash")
);
}
}