use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskOffer {
pub task_id: String,
pub payload_b64: String,
pub max_price_per_hour: f64,
pub estimated_duration_secs: f64,
pub timeout_secs: f64,
#[serde(default = "default_cpus")]
pub cpus: f64,
#[serde(default = "default_memory")]
pub memory_bytes: u64,
#[serde(default)]
pub gpus: u32,
#[serde(default)]
pub worker_type: Option<String>,
#[serde(default)]
pub tags: Vec<String>,
pub requester_user_id: String,
pub source_broker: String,
#[serde(default)]
pub peer_key: String,
#[serde(default)]
pub two_phase: bool,
#[serde(default)]
pub voucher_signed_json: String,
#[serde(default)]
pub voucher_sig: String,
#[serde(default)]
pub target_worker: String,
}
pub fn norm_target(s: &str) -> String {
s.strip_prefix("zc://").unwrap_or(s).to_string()
}
fn default_cpus() -> f64 {
1.0
}
fn default_memory() -> u64 {
1024 * 1024 * 1024
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskResult {
pub task_id: String,
pub payload_b64: String,
pub duration_ms: f64,
pub actual_cost: f64,
pub worker_name: String,
pub worker_uri: String,
pub price_per_hour: f64,
pub executor_owner: String,
#[serde(default)]
pub worker_pid: Option<String>,
#[serde(default)]
pub executor_node_uri: String,
#[serde(default)]
pub executor_node_id: Option<String>,
#[serde(default)]
pub executor_sig: Option<String>,
}
pub fn receipt_bytes(r: &TaskResult) -> Vec<u8> {
format!(
"zc-receipt-v1\n{}\n{}\n{}\n{}\n{}\n{}",
r.task_id, r.duration_ms, r.actual_cost, r.price_per_hour, r.executor_owner, r.worker_name
)
.into_bytes()
}
pub fn verify_receipt(r: &TaskResult, roster: &super::roster_cache::RosterCache) -> bool {
let (node_id, sig) = match (&r.executor_node_id, &r.executor_sig) {
(Some(n), Some(s)) => (n, s),
_ => return false,
};
if !super::node_identity::verify_sig(node_id, &receipt_bytes(r), sig) {
return false;
}
match super::node_identity::fingerprint_of_pubkey_b64(node_id) {
Some(fp) => roster.fingerprint_authorized(&fp),
None => false,
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskReject {
pub task_id: String,
pub reason: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskAccept {
pub task_id: String,
pub estimated_cost: f64,
pub worker_name: String,
pub executor_owner: String,
pub price_per_hour: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskCommit {
pub task_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskCancel {
pub task_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SubscribeMessage {
pub action: String,
pub peer_url: String,
#[serde(default)]
pub peer_key: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerIdentity {
pub owner_user_id: String,
pub node_name: String,
pub verified: bool,
pub workers: Vec<WorkerSummary>,
#[serde(default)]
pub quic_port: u16,
#[serde(default)]
pub supports_two_phase: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerSummary {
pub name: String,
pub price_per_hour: f64,
pub status: String,
pub cpus: f64,
pub memory_bytes: u64,
pub gpus: u32,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn new_types_round_trip() {
let a = TaskAccept {
task_id: "t".into(),
estimated_cost: 1.5,
worker_name: "w".into(),
executor_owner: "o".into(),
price_per_hour: 3.6,
};
let s = serde_json::to_string(&a).unwrap();
let back: TaskAccept = serde_json::from_str(&s).unwrap();
assert_eq!(back.task_id, "t");
assert_eq!(back.estimated_cost, 1.5);
let c: TaskCommit = serde_json::from_str(r#"{"task_id":"t"}"#).unwrap();
assert_eq!(c.task_id, "t");
let x: TaskCancel = serde_json::from_str(r#"{"task_id":"t"}"#).unwrap();
assert_eq!(x.task_id, "t");
}
#[test]
fn task_offer_two_phase_defaults_false_for_old_senders() {
let json = r#"{"task_id":"t","payload_b64":"","max_price_per_hour":1.0,
"estimated_duration_secs":1.0,"timeout_secs":10.0,
"requester_user_id":"u","source_broker":"n"}"#;
let offer: TaskOffer = serde_json::from_str(json).unwrap();
assert!(!offer.two_phase);
}
#[test]
fn target_worker_defaults_empty_and_strips_scheme() {
let json = r#"{"task_id":"t","payload_b64":"","max_price_per_hour":1.0,
"estimated_duration_secs":1.0,"timeout_secs":10.0,
"requester_user_id":"u","source_broker":"n"}"#;
let offer: TaskOffer = serde_json::from_str(json).unwrap();
assert!(offer.target_worker.is_empty());
assert_eq!(norm_target("zc://worker-abc"), "worker-abc");
assert_eq!(norm_target("zc://node-i9"), "node-i9");
assert_eq!(norm_target("worker-abc"), "worker-abc");
}
#[test]
fn peer_identity_supports_two_phase_defaults_false_for_old_peers() {
let json = r#"{"owner_user_id":"o","node_name":"n","verified":true,"workers":[]}"#;
let id: PeerIdentity = serde_json::from_str(json).unwrap();
assert!(!id.supports_two_phase);
}
fn sample_result() -> TaskResult {
TaskResult {
task_id: "t1".into(),
payload_b64: "irrelevant-to-settlement".into(),
duration_ms: 1234.5,
actual_cost: 0.789,
worker_name: "worker-x".into(),
worker_uri: "http://127.0.0.1:9000".into(),
price_per_hour: 3.6,
executor_owner: "owner-1".into(),
worker_pid: None,
executor_node_uri: String::new(),
executor_node_id: None,
executor_sig: None,
}
}
#[derive(Debug, Deserialize)]
#[allow(dead_code)]
struct OldPeerTaskResult {
task_id: String,
payload_b64: String,
duration_ms: f64,
actual_cost: f64,
worker_name: String,
worker_uri: String,
price_per_hour: f64,
executor_owner: String,
#[serde(default)]
worker_pid: Option<String>,
#[serde(default)]
worker_ip: Option<String>,
#[serde(default)]
executor_node_id: Option<String>,
#[serde(default)]
executor_sig: Option<String>,
}
#[test]
fn task_result_carries_identity_not_addresses() {
let mut r = sample_result();
r.worker_uri = "zc://worker-aaaaaaaaaaaaaaaa-3960".into();
r.executor_node_uri = "zc://node-aaaaaaaaaaaaaaaa".into();
let json = serde_json::to_string(&r).unwrap();
assert!(
!json.contains("worker_ip"),
"worker_ip must be off the wire: {json}"
);
assert!(
!json.contains("http://"),
"no scheme+host on the wire: {json}"
);
assert!(
!json.contains("10.13.13."),
"no mesh address on the wire: {json}"
);
assert!(json.contains("zc://node-aaaaaaaaaaaaaaaa"));
let old = r#"{"task_id":"t","payload_b64":"","duration_ms":1.0,"actual_cost":0.1,
"worker_name":"w","worker_uri":"http://10.13.13.7:3960","price_per_hour":3.6,
"executor_owner":"o","worker_ip":"10.13.13.7"}"#;
let back: TaskResult = serde_json::from_str(old).unwrap();
assert!(back.executor_node_uri.is_empty());
let old_view: OldPeerTaskResult = serde_json::from_str(&json)
.expect("an un-upgraded peer must still deserialize a new-shape result");
assert!(old_view.worker_ip.is_none());
assert_eq!(old_view.worker_uri, "zc://worker-aaaaaaaaaaaaaaaa-3960");
}
#[test]
fn receipt_bytes_stable_and_sensitive_to_settlement_fields() {
let a = sample_result();
let b = sample_result();
assert_eq!(receipt_bytes(&a), receipt_bytes(&b));
let mut c = sample_result();
c.actual_cost = 999.0;
assert_ne!(receipt_bytes(&a), receipt_bytes(&c));
let mut d = sample_result();
d.payload_b64 = "totally-different-payload".into();
assert_eq!(receipt_bytes(&a), receipt_bytes(&d));
}
#[test]
fn verify_receipt_true_for_rostered_signer() {
use super::super::node_identity::NodeKey;
use super::super::roster_cache::RosterCache;
let key = NodeKey::generate();
let mut r = sample_result();
r.executor_node_id = Some(key.public_b64());
r.executor_sig = Some(key.sign(&receipt_bytes(&r)));
let roster = RosterCache::from_entries(vec![(key.public_b64(), false)]);
assert!(verify_receipt(&r, &roster));
}
#[test]
fn verify_receipt_false_for_unrostered_signer() {
use super::super::node_identity::NodeKey;
use super::super::roster_cache::RosterCache;
let key = NodeKey::generate();
let other = NodeKey::generate();
let mut r = sample_result();
r.executor_node_id = Some(key.public_b64());
r.executor_sig = Some(key.sign(&receipt_bytes(&r)));
let roster = RosterCache::from_entries(vec![(other.public_b64(), false)]);
assert!(!verify_receipt(&r, &roster));
let empty = RosterCache::from_entries(vec![]);
assert!(!verify_receipt(&r, &empty));
}
#[test]
fn verify_receipt_false_when_tampered_after_signing() {
use super::super::node_identity::NodeKey;
use super::super::roster_cache::RosterCache;
let key = NodeKey::generate();
let mut r = sample_result();
r.executor_node_id = Some(key.public_b64());
r.executor_sig = Some(key.sign(&receipt_bytes(&r)));
let roster = RosterCache::from_entries(vec![(key.public_b64(), false)]);
assert!(verify_receipt(&r, &roster));
r.actual_cost += 1000.0; assert!(!verify_receipt(&r, &roster));
}
#[test]
fn verify_receipt_false_when_sig_or_node_id_missing() {
use super::super::node_identity::NodeKey;
use super::super::roster_cache::RosterCache;
let key = NodeKey::generate();
let roster = RosterCache::from_entries(vec![(key.public_b64(), false)]);
let mut r = sample_result();
assert!(!verify_receipt(&r, &roster));
r.executor_node_id = Some(key.public_b64());
assert!(!verify_receipt(&r, &roster));
r.executor_node_id = None;
r.executor_sig = Some(key.sign(&receipt_bytes(&r)));
assert!(!verify_receipt(&r, &roster)); }
}