use alloy::primitives::{FixedBytes, Uint, U256};
use ark_bn254::{G1Affine, G2Affine};
use ark_ec::AffineRepr;
use async_trait::async_trait;
use eigensdk::{
crypto_bls::{BlsG1Point, BlsG2Point, Signature},
crypto_bn254::utils::verify_message,
services_avsregistry::AvsRegistryService,
types::{
avs::{SignatureVerificationError, TaskResponseDigest},
avs_state::OperatorAvsState,
operator::{QuorumThresholdPercentage, QuorumThresholdPercentages},
},
};
use newton_core::{newton_prover_task_manager::IBLSSignatureCheckerTypes, TaskId};
use serde::{Deserialize, Serialize};
use std::{
collections::{HashMap, HashSet},
sync::Arc,
time::Instant,
};
use tokio::{
sync::{
mpsc::{self, UnboundedReceiver, UnboundedSender},
oneshot,
},
time::Duration,
};
use tracing::{debug, error, field, info, instrument, trace, warn};
#[derive(Debug, Clone)]
pub struct AggregatedOperators {
signers_apk_g2: BlsG2Point,
signers_agg_sig_g1: Signature,
signers_total_stake_per_quorum: HashMap<u8, U256>,
pub signers_operator_ids_set: HashMap<FixedBytes<32>, bool>,
}
#[derive(Debug, Clone)]
pub struct TaskMetadata {
pub task_id: TaskId,
quorum_numbers: Vec<u8>,
quorum_threshold_percentages: QuorumThresholdPercentages,
time_to_expiry: Duration,
window_duration: Duration,
task_created_block: u64,
}
impl TaskMetadata {
pub fn new(
task_id: TaskId,
task_created_block: u64,
quorum_numbers: Vec<u8>,
quorum_threshold_percentages: QuorumThresholdPercentages,
time_to_expiry: Duration,
) -> Self {
Self {
task_id,
task_created_block,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
window_duration: Duration::ZERO,
}
}
pub fn with_window_duration(mut self, window_duration: Duration) -> Self {
self.window_duration = window_duration;
self
}
}
#[derive(Debug, Clone)]
pub struct TaskSignature {
task_id: TaskId,
task_response_digest: TaskResponseDigest,
bls_signature: Signature,
operator_id: FixedBytes<32>,
}
impl TaskSignature {
pub fn new(
task_id: TaskId,
task_response_digest: TaskResponseDigest,
bls_signature: Signature,
operator_id: FixedBytes<32>,
) -> Self {
Self {
task_id,
task_response_digest,
bls_signature,
operator_id,
}
}
}
#[allow(unused)]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BlsAggregationServiceResponse {
pub task_id: TaskId,
pub task_created_block: u64,
pub task_response_digest: TaskResponseDigest,
pub signers_count: usize,
pub non_signers_pub_keys_g1: Vec<BlsG1Point>,
pub non_signers_operators_ids: Vec<FixedBytes<32>>,
pub quorum_apks_g1: Vec<BlsG1Point>,
pub signers_apk_g2: BlsG2Point,
pub signers_agg_sig_g1: Signature,
pub non_signer_quorum_bitmap_indices: Vec<u32>,
pub quorum_apk_indices: Vec<u32>,
pub total_stake_indices: Vec<u32>,
pub non_signer_stake_indices: Vec<Vec<u32>>,
}
use thiserror::Error;
use crate::avs::AvsRegistryServiceCaller;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskExpiryReason {
QuorumMet,
QuorumNotMet,
WindowOpenNoResponse,
WindowFinishedNoResponse,
}
impl std::fmt::Display for TaskExpiryReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
TaskExpiryReason::QuorumMet => write!(f, "task already finished (quorum reached and window closed)"),
TaskExpiryReason::QuorumNotMet => write!(f, "task expired without reaching quorum threshold"),
TaskExpiryReason::WindowOpenNoResponse => write!(
f,
"task expired while window was open but no aggregated response available"
),
TaskExpiryReason::WindowFinishedNoResponse => {
write!(f, "window finished but no aggregated response available")
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReceiverErrorReason {
OneshotReceiverClosed,
AggregateChannelClosed,
}
impl std::fmt::Display for ReceiverErrorReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ReceiverErrorReason::OneshotReceiverClosed => write!(f, "oneshot receiver closed"),
ReceiverErrorReason::AggregateChannelClosed => write!(f, "aggregate response channel closed"),
}
}
}
#[derive(Error, Debug, Clone, PartialEq, Eq)]
pub enum BlsAggregationServiceError {
#[error("task {task_id} expired: {reason}")]
TaskExpired {
task_id: TaskId,
reason: TaskExpiryReason,
},
#[error("task {task_id} not found: {reason}")]
TaskNotFound {
task_id: TaskId,
reason: String,
},
#[error("signature verification failed for task {task_id}, operator {operator_id}: {verification_error}")]
SignatureVerificationError {
task_id: TaskId,
operator_id: FixedBytes<32>,
verification_error: SignatureVerificationError,
},
#[error("signatures channel closed for task {task_id}: {reason}")]
SignaturesChannelClosed {
task_id: TaskId,
reason: String,
},
#[error("AVS registry error for task {task_id}{operator_context}: {reason}")]
RegistryError {
task_id: TaskId,
operator_context: String,
reason: String,
},
#[error("duplicate task id {task_id}: {reason}")]
DuplicateTaskId {
task_id: TaskId,
reason: String,
},
#[error("error sending {operation} message to service for task {task_id}{operator_context}: {reason}")]
SenderError {
operation: String,
task_id: TaskId,
operator_context: String,
reason: String,
},
#[error("error receiving {operation} response from service for task {task_id}{operator_context}: {reason}")]
ReceiverError {
operation: String,
task_id: TaskId,
operator_context: String,
reason: ReceiverErrorReason,
},
}
#[derive(Debug)]
pub enum AggregationMessage {
InitializeTask(
TaskMetadata,
oneshot::Sender<
Result<
UnboundedReceiver<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
BlsAggregationServiceError,
>,
>,
),
ProcessSignature(TaskSignature, oneshot::Sender<Result<(), BlsAggregationServiceError>>),
CancelAggregationLoop(TaskId),
}
#[async_trait]
pub trait BlsServiceHandle: Send + Sync {
async fn initialize_task(
&self,
metadata: TaskMetadata,
) -> Result<
UnboundedReceiver<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
BlsAggregationServiceError,
>;
async fn process_signature(&self, task_signature: TaskSignature) -> Result<(), BlsAggregationServiceError>;
fn cancel_aggregation_loop(&self, task_id: TaskId);
}
#[derive(Debug, Clone)]
pub struct ServiceHandle {
msg_sender: UnboundedSender<AggregationMessage>,
}
#[async_trait]
impl BlsServiceHandle for ServiceHandle {
#[instrument(skip(self), fields(task_id = %metadata.task_id, quorum_count = metadata.quorum_numbers.len()))]
async fn initialize_task(
&self,
metadata: TaskMetadata,
) -> Result<
UnboundedReceiver<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
BlsAggregationServiceError,
> {
let start_time = std::time::Instant::now();
info!(
"[ServiceHandle] initializing task {} with {} quorums, expires in {} ms",
metadata.task_id,
metadata.quorum_numbers.len(),
metadata.time_to_expiry.as_millis()
);
debug!(
"[ServiceHandle] task {} details - quorums: {}, thresholds: {}, window: {} ms",
metadata.task_id,
String::from_utf8(metadata.quorum_numbers.clone()).unwrap_or_default(),
String::from_utf8(metadata.quorum_threshold_percentages.clone()).unwrap_or_default(),
metadata.window_duration.as_millis()
);
let (tx, rx) = oneshot::channel();
if let Err(e) = self
.msg_sender
.send(AggregationMessage::InitializeTask(metadata.clone(), tx))
{
let reason = format!("channel closed: {:?}", e);
error!(
"[ServiceHandle] failed to send InitializeTask message for task {}: {}",
metadata.task_id, reason
);
return Err(BlsAggregationServiceError::SenderError {
operation: "InitializeTask".to_string(),
task_id: metadata.task_id,
operator_context: String::new(),
reason,
});
}
let response_receiver = rx.await.map_err(|e| {
error!(
"[ServiceHandle] failed to receive InitializeTask response for task {}: oneshot receiver closed: {:?}",
metadata.task_id, e
);
BlsAggregationServiceError::ReceiverError {
operation: "InitializeTask".to_string(),
task_id: metadata.task_id,
operator_context: String::new(),
reason: ReceiverErrorReason::OneshotReceiverClosed,
}
})?;
let duration = start_time.elapsed();
match &response_receiver {
Ok(_) => info!(
"[ServiceHandle] task {} initialization completed in {} ms",
metadata.task_id,
duration.as_millis()
),
Err(e) => error!(
"[ServiceHandle] task {} initialization failed in {} ms: {}",
metadata.task_id,
duration.as_millis(),
e.to_string()
),
}
response_receiver
}
#[instrument(skip(self, task_signature), fields(
task_id = %task_signature.task_id,
operator_id = %hex!(task_signature.operator_id.as_slice())
))]
async fn process_signature(&self, task_signature: TaskSignature) -> Result<(), BlsAggregationServiceError> {
let start_time = std::time::Instant::now();
let operator_id_hex = hex!(task_signature.operator_id.as_slice());
debug!(
"[ServiceHandle] processing signature from operator {} for task {}",
operator_id_hex, task_signature.task_id
);
debug!(
"[ServiceHandle] signature details - task_response_digest: {}, signature: {:?}",
hex!(task_signature.task_response_digest),
task_signature.bls_signature
);
let (tx, rx) = oneshot::channel();
if let Err(e) = self
.msg_sender
.send(AggregationMessage::ProcessSignature(task_signature.clone(), tx))
{
let reason = format!("channel closed: {:?}", e);
error!(
"[ServiceHandle] failed to send ProcessSignature message for task {} from operator {}: {}",
task_signature.task_id, operator_id_hex, reason
);
return Err(BlsAggregationServiceError::SenderError {
operation: "ProcessSignature".to_string(),
task_id: task_signature.task_id,
operator_context: format!(" from operator {}", operator_id_hex),
reason,
});
}
let result = rx.await.map_err(|e| {
error!(
"[ServiceHandle] failed to receive ProcessSignature response for task {} from operator {}: oneshot receiver closed: {:?}",
task_signature.task_id, operator_id_hex, e
);
BlsAggregationServiceError::ReceiverError {
operation: "ProcessSignature".to_string(),
task_id: task_signature.task_id,
operator_context: format!(" from operator {}", operator_id_hex),
reason: ReceiverErrorReason::OneshotReceiverClosed,
}
})?;
let duration = start_time.elapsed();
match &result {
Ok(_) => info!(
"[ServiceHandle] signature from operator {} for task {} processed in {} ms",
operator_id_hex,
task_signature.task_id,
duration.as_millis()
),
Err(e) => warn!(
"[ServiceHandle] signature from operator {} for task {} failed in {} ms: {}",
operator_id_hex,
task_signature.task_id,
duration.as_millis(),
e.to_string()
),
}
result
}
fn cancel_aggregation_loop(&self, task_id: TaskId) {
debug!(%task_id, "[ServiceHandle] sending CancelAggregationLoop");
if let Err(e) = self.msg_sender.send(AggregationMessage::CancelAggregationLoop(task_id)) {
warn!(
%task_id,
"[ServiceHandle] failed to send CancelAggregationLoop: channel closed: {:?}", e
);
}
}
}
#[derive(Debug)]
pub struct AggregateReceiver {
aggregate_receiver: UnboundedReceiver<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
}
impl AggregateReceiver {
#[instrument(skip(self))]
pub async fn receive_aggregated_response(
&mut self,
) -> Result<BlsAggregationServiceResponse, BlsAggregationServiceError> {
debug!("[AggregateReceiver] waiting for aggregated response...");
let start_time = std::time::Instant::now();
match self.aggregate_receiver.recv().await {
Some(Ok(response)) => {
let duration = start_time.elapsed();
info!(
"[AggregateReceiver] received successful aggregated response for task {} in {} ms",
response.task_id,
duration.as_millis()
);
debug!(
"[AggregateReceiver] response details - non-signer keys: {}",
response.non_signers_pub_keys_g1.len()
);
Ok(response)
}
Some(Err(e)) => {
let duration = start_time.elapsed();
error!(
"[AggregateReceiver] received error response in {} ms: {}",
duration.as_millis(),
e.to_string()
);
Err(e)
}
None => {
let duration = start_time.elapsed();
warn!(
"[AggregateReceiver] aggregate response channel closed after waiting {} ms",
duration.as_millis()
);
Err(BlsAggregationServiceError::ReceiverError {
operation: "ReceiveAggregatedResponse".to_string(),
task_id: TaskId::default(), operator_context: String::new(),
reason: ReceiverErrorReason::AggregateChannelClosed,
})
}
}
}
}
type TaskChannelEntry = (
UnboundedSender<SignedTaskResponseDigest>,
UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
Instant,
);
#[derive(Debug)]
pub struct BlsAggregatorService<A: AvsRegistryService>
where
A: Clone,
{
avs_registry_service: A,
}
#[derive(Debug)]
struct SignedTaskResponseDigest {
task_response_digest: TaskResponseDigest,
bls_signature: Signature,
operator_id: FixedBytes<32>,
result_channel: oneshot::Sender<Result<(), BlsAggregationServiceError>>,
}
impl<A: AvsRegistryService + Send + Sync + Clone + 'static> BlsAggregatorService<A> {
pub fn new(avs_registry_service: A) -> Self {
info!("[BlsAggregatorService] creating new BLS aggregator service");
Self { avs_registry_service }
}
#[instrument(skip(self))]
pub fn start(self) -> (ServiceHandle, AggregateReceiver) {
info!("[BlsAggregatorService] starting BLS aggregator service");
let (msg_tx, msg_rx) = mpsc::unbounded_channel();
let (agg_tx, agg_rx) = mpsc::unbounded_channel();
debug!("[BlsAggregatorService] created communication channels");
tokio::spawn(async move {
info!("[BlsAggregatorService] spawning main service loop");
self.run(msg_rx, agg_tx).await;
warn!("[BlsAggregatorService] main service loop ended");
});
let service_handler = ServiceHandle { msg_sender: msg_tx };
let aggregate_receiver = AggregateReceiver {
aggregate_receiver: agg_rx,
};
info!("[BlsAggregatorService] service started successfully");
(service_handler, aggregate_receiver)
}
const MAX_ACTIVE_TASKS: usize = 10000;
const MAX_AGGREGATED_OPERATORS_PER_TASK: usize = 100;
#[instrument(skip(self, msg_receiver, aggregate_sender))]
async fn run(
self,
mut msg_receiver: UnboundedReceiver<AggregationMessage>,
aggregate_sender: UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
) {
info!("[BlsAggregatorService] main service loop started");
let mut task_channels: HashMap<TaskId, TaskChannelEntry> = HashMap::new();
let mut finished_tasks: HashSet<TaskId> = HashSet::new();
let mut message_count = 0u64;
let mut active_tasks = 0u32;
let (cleanup_tx, mut cleanup_rx) = mpsc::unbounded_channel::<TaskId>();
loop {
tokio::select! {
message = msg_receiver.recv() => {
let message = match message {
Some(msg) => msg,
None => {
warn!(
"[BlsAggregatorService] message channel closed, shutting down main loop (processed {} messages, {} active tasks)",
message_count,
active_tasks
);
task_channels.clear();
break;
}
};
message_count += 1;
trace!(
"[BlsAggregatorService] processing message #{} (active tasks: {})",
message_count,
active_tasks
);
match message {
AggregationMessage::InitializeTask(metadata, result_sender) => {
let task_id = metadata.task_id;
info!(
"[BlsAggregatorService] received InitializeTask for task {} (message #{})",
task_id, message_count
);
if task_channels.contains_key(&task_id) {
let _ = result_sender.send(Err(BlsAggregationServiceError::DuplicateTaskId {
task_id,
reason: format!("task already exists in task_channels (message #{})", message_count),
}));
continue;
}
let (signature_tx, signature_rx) = mpsc::unbounded_channel::<SignedTaskResponseDigest>();
let (response_tx, response_rx) =
mpsc::unbounded_channel::<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>();
task_channels.insert(task_id, (signature_tx, response_tx.clone(), Instant::now()));
active_tasks += 1;
if result_sender.send(Ok(response_rx)).is_err() {
warn!(
"[BlsAggregatorService] failed to send response receiver for task {} (caller dropped)",
task_id
);
task_channels.remove(&task_id);
active_tasks = active_tasks.saturating_sub(1);
continue;
}
let avs_registry_service = self.avs_registry_service.clone();
let task_response_sender = response_tx;
let cleanup_tx_clone = cleanup_tx.clone();
let join_handle = tokio::spawn(async move {
let result = Self::single_task_aggregator(
avs_registry_service,
metadata,
task_response_sender,
signature_rx,
)
.await;
match result {
Ok(()) => {
info!(
task_id = %task_id,
"[BlsAggregatorService] task aggregator finished successfully"
);
}
Err(e) => {
error!(
task_id = %task_id,
error = %e,
"[BlsAggregatorService] task aggregator finished with error"
);
}
}
let _ = cleanup_tx_clone.send(task_id);
});
let task_id_for_monitor = task_id;
tokio::spawn(async move {
if let Err(e) = join_handle.await {
error!(
task_id = %task_id_for_monitor,
error = ?e,
"[BlsAggregatorService] task aggregator panicked or was cancelled"
);
}
});
}
AggregationMessage::ProcessSignature(task_signature, result_sender) => {
let task_id = task_signature.task_id;
if let Some((sig_sender, _response_tx, _)) = task_channels.get_mut(&task_signature.task_id) {
let signed_digest = SignedTaskResponseDigest {
task_response_digest: task_signature.task_response_digest,
bls_signature: task_signature.bls_signature,
operator_id: task_signature.operator_id,
result_channel: result_sender, };
debug!(
"[BlsAggregatorService] sending signed task response digest to task aggregator for task {}",
task_signature.task_id
);
if let Err(send_error) = sig_sender.send(signed_digest) {
warn!(
"[BlsAggregatorService] task {} aggregator channel closed (task finished), removing from task_channels",
task_signature.task_id
);
task_channels.remove(&task_signature.task_id);
active_tasks = active_tasks.saturating_sub(1);
finished_tasks.insert(task_signature.task_id);
error!(
"[BlsAggregatorService] failed to send signed task response digest to task aggregator for task {}: task already finished",
task_signature.task_id
);
let _ = send_error
.0
.result_channel
.send(Err(BlsAggregationServiceError::TaskExpired {
task_id: task_signature.task_id,
reason: TaskExpiryReason::QuorumMet,
}));
}
} else {
if finished_tasks.contains(&task_signature.task_id) {
warn!(
"[BlsAggregatorService] task {} not found in task_channels - task already finished",
task_signature.task_id
);
let _ = result_sender.send(Err(BlsAggregationServiceError::TaskExpired {
task_id: task_signature.task_id,
reason: TaskExpiryReason::QuorumMet,
}));
} else {
warn!(
"[BlsAggregatorService] task {} not found in task_channels for signature processing (task may not be initialized yet or was never created)",
task_signature.task_id
);
let _ = result_sender.send(Err(BlsAggregationServiceError::TaskNotFound {
task_id: task_signature.task_id,
reason:
"task not found in task_channels (task may not be initialized yet or was never created)"
.to_string(),
}));
}
}
}
AggregationMessage::CancelAggregationLoop(task_id) => {
if task_channels.remove(&task_id).is_some() {
active_tasks = active_tasks.saturating_sub(1);
finished_tasks.insert(task_id);
debug!(
%task_id,
"[BlsAggregatorService] cancelled aggregation loop (two-phase completed)"
);
}
}
}
}
completed_task_id = cleanup_rx.recv() => {
if let Some(task_id) = completed_task_id {
if task_channels.remove(&task_id).is_some() {
active_tasks = active_tasks.saturating_sub(1);
finished_tasks.insert(task_id);
trace!(
"[BlsAggregatorService] proactively cleaned up completed task {} (remaining: {})",
task_id,
active_tasks
);
}
}
}
}
}
}
#[instrument(skip(avs_registry_service, aggregated_response_sender, signatures_rx), fields(
task_id = %metadata.task_id,
quorum_count = metadata.quorum_numbers.len(),
))]
async fn single_task_aggregator(
avs_registry_service: A,
metadata: TaskMetadata,
aggregated_response_sender: UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
signatures_rx: UnboundedReceiver<SignedTaskResponseDigest>,
) -> Result<(), BlsAggregationServiceError> {
let start_time = std::time::Instant::now();
info!(
task_id = %metadata.task_id,
"[TaskAggregator] starting single task aggregator - quorums: {}, expires in: {} ms",
String::from_utf8(metadata.quorum_numbers.clone()).unwrap_or_default(),
metadata.time_to_expiry.as_millis()
);
debug!(task_id = %metadata.task_id, "[TaskAggregator] building quorum threshold map");
let quorum_threshold_percentage_map: HashMap<u8, u8> = metadata
.quorum_numbers
.iter()
.enumerate()
.map(|(i, quorum_number)| (*quorum_number, metadata.quorum_threshold_percentages[i]))
.collect();
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] quorum thresholds: {}",
quorum_threshold_percentage_map
.iter()
.map(|(k, v)| format!("{}:{}", k, v))
.collect::<Vec<String>>()
.join(", ")
);
debug!(task_id = %metadata.task_id, "[TaskAggregator] fetching operator AVS state...");
let operator_fetch_start = std::time::Instant::now();
let operator_state_avs = avs_registry_service
.get_operators_avs_state_at_block(metadata.task_created_block, &metadata.quorum_numbers)
.await
.map_err(|e| {
let duration = operator_fetch_start.elapsed();
error!(
task_id = %metadata.task_id,
block = metadata.task_created_block,
quorum_count = metadata.quorum_numbers.len(),
duration_ms = duration.as_millis(),
error = %e,
"Failed to get operator AVS state from registry service"
);
BlsAggregationServiceError::RegistryError {
task_id: metadata.task_id,
operator_context: String::new(),
reason: format!(
"failed to get operator AVS state at block {}: {}",
metadata.task_created_block, e
),
}
})?;
let operator_fetch_duration = operator_fetch_start.elapsed();
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] fetched {} operators in {} ms",
operator_state_avs.len(),
operator_fetch_duration.as_millis()
);
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] operator IDs: {:?}",
operator_state_avs
.keys()
.map(|k| hex!(k.as_slice()))
.collect::<Vec<_>>()
);
debug!(task_id = %metadata.task_id, "[TaskAggregator] fetching quorum AVS state...");
let quorum_fetch_start = std::time::Instant::now();
let quorums_avs_state = avs_registry_service
.get_quorums_avs_state_at_block(&metadata.quorum_numbers, metadata.task_created_block)
.await
.map_err(|e| {
let duration = quorum_fetch_start.elapsed();
error!(
task_id = %metadata.task_id,
block = metadata.task_created_block,
quorum_count = metadata.quorum_numbers.len(),
duration_ms = duration.as_millis(),
error = %e,
"Failed to get quorum AVS state from registry service"
);
BlsAggregationServiceError::RegistryError {
task_id: metadata.task_id,
operator_context: String::new(),
reason: format!(
"failed to get quorum AVS state at block {}: {}",
metadata.task_created_block, e
),
}
})?;
let quorum_fetch_duration = quorum_fetch_start.elapsed();
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] fetched quorum state in {} ms",
quorum_fetch_duration.as_millis()
);
for (quorum_num, state) in &quorums_avs_state {
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] quorum {} - total stake: {}, block: {}",
quorum_num, state.total_stake, state.block_num
);
}
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] computing total stakes per quorum",
);
let total_stake_per_quorum: HashMap<_, _> =
quorums_avs_state.iter().map(|(k, v)| (*k, v.total_stake)).collect();
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] total stakes: {}",
total_stake_per_quorum
.iter()
.map(|(quorum_num, stake)| format!("(quorum num: {}, stake: {})", quorum_num, stake))
.collect::<Vec<String>>()
.join(", ")
);
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] extracting quorum aggregate public keys",
);
let quorum_apks_g1: Vec<BlsG1Point> = metadata
.quorum_numbers
.iter()
.filter_map(|quorum_num| quorums_avs_state.get(quorum_num))
.map(|avs_state| avs_state.agg_pub_key_g1.clone())
.collect();
debug!(
task_id = %metadata.task_id,
"[TaskAggregator] extracted {} quorum aggregate keys",
quorum_apks_g1.len()
);
for (idx, apk) in quorum_apks_g1.iter().enumerate() {
if let (Some(x), Some(y)) = (apk.g1().x(), apk.g1().y()) {
debug!(
"[DEBUG] BLS_QUORUM_APK_EXTRACT: Task {} quorum_idx={} G1_X={} G1_Y={}",
metadata.task_id, idx, x, y
);
} else {
debug!(
"[DEBUG] BLS_QUORUM_APK_EXTRACT: Task {} quorum_idx={} APK is at infinity",
metadata.task_id, idx
);
}
}
let setup_duration = start_time.elapsed();
info!(
task_id = %metadata.task_id,
"[TaskAggregator] setup completed in {} ms, starting signature aggregation loop",
setup_duration.as_millis()
);
Self::loop_task_aggregator(
avs_registry_service,
metadata.task_id,
metadata.task_created_block,
metadata.time_to_expiry,
aggregated_response_sender,
signatures_rx,
operator_state_avs,
total_stake_per_quorum,
quorum_threshold_percentage_map,
quorum_apks_g1,
metadata.quorum_numbers,
metadata.window_duration,
)
.await
}
#[allow(clippy::too_many_arguments)]
#[instrument(skip(avs_registry_service, aggregated_response_sender, signatures_rx, operator_state_avs, total_stake_per_quorum, quorum_threshold_percentage_map, quorum_apks_g1), fields(
task_id = %task_id,
quorum_count = quorum_nums.len(),
operator_count = operator_state_avs.len(),
window_duration = ?window_duration
))]
async fn loop_task_aggregator(
avs_registry_service: A,
task_id: TaskId,
task_created_block: u64,
time_to_expiry: Duration,
aggregated_response_sender: UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
mut signatures_rx: UnboundedReceiver<SignedTaskResponseDigest>,
operator_state_avs: HashMap<FixedBytes<32>, OperatorAvsState>,
total_stake_per_quorum: HashMap<u8, Uint<256, 4>>,
quorum_threshold_percentage_map: HashMap<u8, u8>,
quorum_apks_g1: Vec<BlsG1Point>,
quorum_nums: Vec<u8>,
window_duration: Duration,
) -> Result<(), BlsAggregationServiceError> {
let start_time = std::time::Instant::now();
debug!(
task_id = %task_id,
"[TaskAggregator] starting signature aggregation loop - {} operators, {} quorums, window: {} ms",
operator_state_avs.len(),
quorum_nums.len(),
window_duration.as_millis()
);
debug!(
task_id = %task_id,
"[TaskAggregator] total stakes per quorum: {}",
total_stake_per_quorum
.iter()
.map(|(quorum_num, stake)| format!("quorum num: {}, stake: {}", quorum_num, stake))
.collect::<Vec<String>>()
.join(", ")
);
debug!(
task_id = %task_id,
"[TaskAggregator] threshold percentages: {}",
quorum_threshold_percentage_map
.iter()
.map(|(quorum_num, percentage)| format!("quorum num: {}, percentage: {}", quorum_num, percentage))
.collect::<Vec<String>>()
.join(", ")
);
let mut aggregated_operators: HashMap<FixedBytes<32>, AggregatedOperators> = HashMap::new();
let mut open_window = false;
let mut current_aggregated_response: Option<BlsAggregationServiceResponse> = None;
let (window_tx, mut window_rx) = tokio::sync::mpsc::unbounded_channel::<bool>();
let task_expired_timer = tokio::time::sleep(time_to_expiry);
tokio::pin!(task_expired_timer);
let mut signature_count = 0u32;
let mut duplicate_count = 0u32;
let mut invalid_count = 0u32;
debug!(
task_id = %task_id,
"[TaskAggregator] entering main aggregation loop, expires in {} ms",
time_to_expiry.as_millis()
);
loop {
tokio::select! {
_ = &mut task_expired_timer => {
let loop_duration = start_time.elapsed();
warn!(
task_id = %task_id,
"[TaskAggregator] task expired after {} ms - processed {} signatures ({} duplicates, {} invalid), quorum reached: {}",
loop_duration.as_millis(), signature_count, duplicate_count, invalid_count, current_aggregated_response.is_some()
);
Self::handle_task_expired(
&aggregated_response_sender,
task_id,
open_window,
¤t_aggregated_response,
)?;
return Ok(());
},
_ = window_rx.recv() => {
info!(
task_id = %task_id,
"[TaskAggregator] aggregation window closed - sending response",
);
Self::handle_window_finished(
&aggregated_response_sender,
task_id,
¤t_aggregated_response,
)?;
return Ok(());
},
signed_task_digest = signatures_rx.recv() => {
match signed_task_digest {
Some(digest) => {
signature_count += 1;
let operator_id_hex = hex!(digest.operator_id.as_slice());
debug!(
task_id = %task_id,
"[TaskAggregator] received signature #{} from operator {}",
signature_count, operator_id_hex
);
match Self::handle_new_signature(
&avs_registry_service,
&mut aggregated_operators,
&mut open_window,
&mut current_aggregated_response,
&window_tx,
task_id,
task_created_block,
&operator_state_avs,
&total_stake_per_quorum,
&quorum_threshold_percentage_map,
&quorum_apks_g1,
&quorum_nums,
window_duration,
Some(digest),
).await {
Ok(_) => {
trace!(task_id = %task_id, "[TaskAggregator] successfully processed signature from {}", operator_id_hex);
},
Err(BlsAggregationServiceError::SignatureVerificationError { verification_error: SignatureVerificationError::DuplicateSignature, .. }) => {
duplicate_count += 1;
debug!(task_id = %task_id, "[TaskAggregator] duplicate signature from {} (total duplicates: {})", operator_id_hex, duplicate_count);
},
Err(BlsAggregationServiceError::SignatureVerificationError { verification_error, .. }) => {
invalid_count += 1;
warn!(task_id = %task_id, "[TaskAggregator] invalid signature from {} (total invalid: {}, error: {:?})", operator_id_hex, invalid_count, verification_error);
},
Err(e) => {
error!(task_id = %task_id, "[TaskAggregator] error processing signature from {}: {}", operator_id_hex, e.to_string());
return Err(e);
}
}
},
None => {
warn!(task_id = %task_id, "[TaskAggregator] signature channel closed");
return Ok(());
}
}
}
}
}
}
#[allow(clippy::too_many_arguments)]
#[instrument(skip_all)]
async fn handle_new_signature(
avs_registry_service: &A,
aggregated_operators: &mut HashMap<FixedBytes<32>, AggregatedOperators>,
open_window: &mut bool,
current_aggregated_response: &mut Option<BlsAggregationServiceResponse>,
window_tx: &UnboundedSender<bool>,
task_id: TaskId,
task_created_block: u64,
operator_state_avs: &HashMap<FixedBytes<32>, OperatorAvsState>,
total_stake_per_quorum: &HashMap<u8, Uint<256, 4>>,
quorum_threshold_percentage_map: &HashMap<u8, u8>,
quorum_apks_g1: &[BlsG1Point],
quorum_nums: &[u8],
window_duration: Duration,
signed_task_digest: Option<SignedTaskResponseDigest>,
) -> Result<(), BlsAggregationServiceError> {
let start_time = std::time::Instant::now();
let signed_digest = signed_task_digest.ok_or_else(|| BlsAggregationServiceError::SignaturesChannelClosed {
task_id,
reason: "signature channel receiver dropped (task aggregator may have finished or expired)".to_string(),
})?;
if signed_digest.operator_id == FixedBytes::ZERO {
error!(
task_id = %task_id,
"Invalid operator_id: zero operator ID"
);
return Err(BlsAggregationServiceError::RegistryError {
task_id,
operator_context: String::new(),
reason: "invalid operator_id: zero operator ID".to_string(),
});
}
if signed_digest.task_response_digest == FixedBytes::ZERO {
error!(
task_id = %task_id,
operator_id = %hex!(signed_digest.operator_id.as_slice()),
"Invalid task_response_digest: zero digest"
);
return Err(BlsAggregationServiceError::RegistryError {
task_id,
operator_context: format!(" from operator {}", hex!(signed_digest.operator_id.as_slice())),
reason: "invalid task_response_digest: zero digest".to_string(),
});
}
if quorum_nums.is_empty() {
error!(
task_id = %task_id,
"Invalid quorum_nums: empty quorum numbers"
);
return Err(BlsAggregationServiceError::RegistryError {
task_id,
operator_context: String::new(),
reason: "invalid quorum_nums: empty quorum numbers".to_string(),
});
}
let operator_id_hex = hex!(signed_digest.operator_id.as_slice());
debug!(
task_id = %task_id,
"[TaskAggregator] processing signature from operator {} for digest {}",
operator_id_hex,
hex!(signed_digest.task_response_digest)
);
if Self::is_duplicate_signature(aggregated_operators, &signed_digest) {
debug!(
task_id = %task_id,
"[TaskAggregator] duplicate signature detected from operator {}",
operator_id_hex
);
if signed_digest
.result_channel
.send(Err(BlsAggregationServiceError::SignatureVerificationError {
task_id,
operator_id: signed_digest.operator_id,
verification_error: SignatureVerificationError::DuplicateSignature,
}))
.is_err()
{
warn!(
task_id = %task_id,
"[TaskAggregator] failed to send duplicate signature error to result channel for operator {} (receiver dropped - likely request timeout or cancellation)",
operator_id_hex
);
}
return Ok(());
}
debug!(
task_id = %task_id,
"[TaskAggregator] verifying signature from operator {}",
operator_id_hex
);
let verification_start = std::time::Instant::now();
let verification_result = verify_signature(task_id, &signed_digest, operator_state_avs)
.await
.map_err(|e| BlsAggregationServiceError::SignatureVerificationError {
task_id,
operator_id: signed_digest.operator_id,
verification_error: e,
});
let verification_duration = verification_start.elapsed();
let verification_has_error = verification_result.is_err();
match &verification_result {
Ok(_) => {
debug!(
task_id = %task_id,
"[TaskAggregator] signature verification passed for operator {} in {} ms",
operator_id_hex,
verification_duration.as_millis()
);
}
Err(e) => {
warn!(
task_id = %task_id,
"[TaskAggregator] signature verification failed for operator {} in {} ms: {}",
operator_id_hex,
verification_duration.as_millis(),
e.to_string()
);
}
}
if signed_digest.result_channel.send(verification_result).is_err() {
warn!(
task_id = %task_id,
"[TaskAggregator] failed to send verification result to result channel for operator {} (receiver dropped - likely request timeout or cancellation)",
operator_id_hex
);
return Ok(());
}
if verification_has_error {
return Ok(());
}
debug!(
task_id = %task_id,
"[TaskAggregator] looking up operator state for {}",
operator_id_hex
);
let operator_state = operator_state_avs.get(&signed_digest.operator_id).ok_or_else(|| {
let duration = start_time.elapsed();
error!(
task_id = %task_id,
operator_id = %operator_id_hex,
digest = %hex!(signed_digest.task_response_digest),
duration_ms = duration.as_millis(),
total_operators = operator_state_avs.len(),
"Operator state not found in operator_state_avs map"
);
BlsAggregationServiceError::RegistryError {
task_id,
operator_context: format!(" from operator {}", operator_id_hex),
reason: format!(
"operator state not found in operator_state_avs map (total operators: {}, duration: {} ms)",
operator_state_avs.len(),
duration.as_millis()
),
}
})?;
debug!(
task_id = %task_id,
"[TaskAggregator] operator {} stakes: {:?}",
operator_id_hex, operator_state.stake_per_quorum
);
debug!(
task_id = %task_id,
"[TaskAggregator] updating aggregated operators with signature from {}",
operator_id_hex
);
let update_start = std::time::Instant::now();
let updated_aggregated = update_aggregated_operators(
task_id,
aggregated_operators,
operator_state,
signed_digest.task_response_digest,
signed_digest.bls_signature,
signed_digest.operator_id,
)
.map_err(|e| {
error!(
task_id = %task_id,
operator_id = %operator_id_hex,
?e,
"Failed to update aggregated operators: missing operator public keys"
);
e
})?;
let is_new_digest = !aggregated_operators.contains_key(&signed_digest.task_response_digest);
if is_new_digest && aggregated_operators.len() >= Self::MAX_AGGREGATED_OPERATORS_PER_TASK {
warn!(
task_id = %task_id,
operator_id = %operator_id_hex,
digest = %hex!(signed_digest.task_response_digest.as_slice()),
current_size = aggregated_operators.len(),
max_size = Self::MAX_AGGREGATED_OPERATORS_PER_TASK,
"Ignoring new digest - aggregated_operators limit reached (this should never happen)"
);
return Ok(());
}
if is_new_digest {
aggregated_operators.insert(signed_digest.task_response_digest, updated_aggregated);
}
let update_duration = update_start.elapsed();
let aggregated = aggregated_operators
.get(&signed_digest.task_response_digest)
.expect("aggregated operators should exist after update");
debug!(
task_id = %task_id,
"[TaskAggregator] aggregated operators updated in {} ms - total signers: {}",
update_duration.as_millis(),
aggregated.signers_operator_ids_set.len()
);
let threshold_check_start = std::time::Instant::now();
let threshold_met = Self::check_if_stake_thresholds_met(
&aggregated.signers_total_stake_per_quorum,
total_stake_per_quorum,
quorum_threshold_percentage_map,
);
let threshold_check_duration = threshold_check_start.elapsed();
if !threshold_met {
debug!(
task_id = %task_id,
"[TaskAggregator] stake thresholds not yet met (checked in {} ms) - current stakes: {:?}",
threshold_check_duration.as_millis(),
aggregated.signers_total_stake_per_quorum
);
return Ok(());
}
info!(
task_id = %task_id,
"[TaskAggregator] stake thresholds met! ({} signers, checked in {} ms)",
aggregated.signers_operator_ids_set.len(),
threshold_check_duration.as_millis()
);
if !*open_window {
*open_window = true;
info!(
task_id = %task_id,
"[TaskAggregator] opening aggregation window for {} ms",
window_duration.as_millis()
);
Self::start_window(window_tx, window_duration, task_id);
} else {
debug!(
task_id = %task_id,
"[TaskAggregator] aggregation window already open, updating response",
);
}
debug!(task_id = %task_id, "[TaskAggregator] building aggregated response...");
let response_build_start = std::time::Instant::now();
*current_aggregated_response = Some(
Self::build_aggregated_response(
task_id,
task_created_block,
signed_digest.task_response_digest,
operator_state_avs,
aggregated.clone(),
avs_registry_service,
quorum_apks_g1,
quorum_nums,
)
.await?,
);
let response_build_duration = response_build_start.elapsed();
let total_duration = start_time.elapsed();
debug!(
task_id = %task_id,
"[TaskAggregator] signature processing completed in {} ms (response built in {} ms) for operator {}",
total_duration.as_millis(),
response_build_duration.as_millis(),
operator_id_hex
);
Ok(())
}
#[instrument(skip_all, fields(task_id = %task_id, open_window = open_window))]
fn handle_task_expired(
aggregated_response_sender: &UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
task_id: TaskId,
open_window: bool,
current_aggregated_response: &Option<BlsAggregationServiceResponse>,
) -> Result<(), BlsAggregationServiceError> {
if open_window {
info!(
task_id = %task_id,
"[TaskAggregator] task expired while aggregation window was open - sending current response",
);
if let Some(response) = current_aggregated_response {
debug!(
task_id = %task_id,
"[TaskAggregator] sending response with {} non-signer keys",
response.non_signers_pub_keys_g1.len()
);
aggregated_response_sender.send(Ok(response.clone())).map_err(|e| {
let reason = format!("failed to send expired task response: {:?}", e);
error!(task_id = %task_id, "[TaskAggregator] {}", reason);
BlsAggregationServiceError::SenderError {
operation: "SendExpiredTaskResponse".to_string(),
task_id,
operator_context: String::new(),
reason,
}
})?;
} else {
error!(task_id = %task_id, "[TaskAggregator] window was open but no response available");
if aggregated_response_sender
.send(Err(BlsAggregationServiceError::TaskExpired {
task_id,
reason: TaskExpiryReason::WindowOpenNoResponse,
}))
.is_err()
{
warn!(
task_id = %task_id,
"[TaskAggregator] handle_task_expired:window_open_no_response - failed to send task expired error",
);
return Ok(());
}
}
} else {
warn!(
task_id = %task_id,
"[TaskAggregator] task expired without reaching quorum threshold",
);
if aggregated_response_sender
.send(Err(BlsAggregationServiceError::TaskExpired {
task_id,
reason: TaskExpiryReason::QuorumNotMet,
}))
.is_err()
{
warn!(
task_id = %task_id,
"[TaskAggregator] handle_task_expired:quorum_not_reached - failed to send task expired error",
);
return Ok(());
}
}
Ok(())
}
#[instrument(skip_all, fields(task_id = %task_id))]
fn handle_window_finished(
aggregated_response_sender: &UnboundedSender<Result<BlsAggregationServiceResponse, BlsAggregationServiceError>>,
task_id: TaskId,
current_aggregated_response: &Option<BlsAggregationServiceResponse>,
) -> Result<(), BlsAggregationServiceError> {
info!(
task_id = %task_id,
"[TaskAggregator] aggregation window finished - sending final response",
);
if let Some(response) = current_aggregated_response {
debug!(
task_id = %task_id,
"[TaskAggregator] sending final response with {} non-signer keys",
response.non_signers_pub_keys_g1.len()
);
aggregated_response_sender.send(Ok(response.clone())).map_err(|e| {
let reason = format!("failed to send window finished response: {}", e);
error!(task_id = %task_id, "[TaskAggregator] {}", reason);
BlsAggregationServiceError::SenderError {
operation: "SendWindowFinishedResponse".to_string(),
task_id,
operator_context: String::new(),
reason,
}
})?;
} else {
error!(task_id = %task_id, "[TaskAggregator] window finished but no response available");
if aggregated_response_sender
.send(Err(BlsAggregationServiceError::TaskExpired {
task_id,
reason: TaskExpiryReason::WindowFinishedNoResponse,
}))
.is_err()
{
warn!(
task_id = %task_id,
"[TaskAggregator] handle_window_finished:no_response - failed to send task expired error",
);
return Ok(());
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
#[instrument(skip_all)]
pub async fn build_aggregated_response(
task_id: TaskId,
task_created_block: u64,
task_response_digest: FixedBytes<32>,
operator_state_avs: &HashMap<FixedBytes<32>, OperatorAvsState>,
digest_aggregated_operators: AggregatedOperators,
avs_registry_service: &A,
quorum_apks_g1: &[BlsG1Point],
quorum_nums: &[u8],
) -> Result<BlsAggregationServiceResponse, BlsAggregationServiceError> {
let start_time = std::time::Instant::now();
debug!(
task_id = %task_id,
"[TaskAggregator] building aggregated response for digest {}",
task_response_digest
);
debug!(
task_id = %task_id,
"[TaskAggregator] signers count: {}, quorums: {:?}",
digest_aggregated_operators.signers_operator_ids_set.len(),
quorum_nums
);
debug!(
task_id = %task_id,
"[TaskAggregator] computing non-signers from {} total operators",
operator_state_avs.len()
);
let mut non_signers_operators_ids: Vec<FixedBytes<32>> = operator_state_avs
.keys()
.filter(|operator_id| {
!digest_aggregated_operators
.signers_operator_ids_set
.contains_key(*operator_id)
})
.cloned()
.collect();
non_signers_operators_ids.sort();
debug!(
task_id = %task_id,
"[TaskAggregator] identified {} non-signers out of {} operators",
non_signers_operators_ids.len(),
operator_state_avs.len()
);
debug!(
task_id = %task_id,
"[TaskAggregator] non-signer IDs: {}",
non_signers_operators_ids
.iter()
.map(|id| hex!(id.as_slice()))
.collect::<Vec<_>>()
.join(", ")
);
debug!(
task_id = %task_id,
"[TaskAggregator] extracting public keys for {} non-signers",
non_signers_operators_ids.len()
);
let non_signers_pub_keys_g1: Vec<BlsG1Point> = non_signers_operators_ids
.iter()
.filter_map(|operator_id| {
let state = operator_state_avs.get(operator_id);
if state.is_none() {
warn!(
task_id = %task_id,
"[TaskAggregator] operator state not found for non-signer {}",
hex!(operator_id.as_slice())
);
}
state
})
.filter_map(|operator_avs_state| {
if operator_avs_state.operator_info.pub_keys.is_none() {
warn!(task_id = %task_id, "[TaskAggregator] public keys not found for non-signer");
}
operator_avs_state.operator_info.pub_keys.clone()
})
.map(|pub_keys| pub_keys.g1_pub_key)
.collect();
debug!(
task_id = %task_id,
"[TaskAggregator] extracted {} public keys for non-signers",
non_signers_pub_keys_g1.len()
);
debug!(
task_id = %task_id,
"[TaskAggregator] fetching signature check indices for block {}",
task_created_block
);
let indices_start = std::time::Instant::now();
let indices = avs_registry_service
.get_check_signatures_indices(
task_created_block,
quorum_nums.into(),
non_signers_operators_ids.clone(),
)
.await
.map_err(|err| {
let duration = indices_start.elapsed();
error!(
task_id = %task_id,
block = task_created_block,
quorum_count = quorum_nums.len(),
non_signer_count = non_signers_operators_ids.len(),
duration_ms = duration.as_millis(),
error = ?err,
"Failed to get check signatures indices from registry service"
);
BlsAggregationServiceError::RegistryError {
task_id,
operator_context: String::new(),
reason: format!(
"failed to get check signatures indices at block {}: {:?}",
task_created_block, err
),
}
})?;
let indices_duration = indices_start.elapsed();
debug!(
task_id = %task_id,
"[TaskAggregator] fetched signature check indices in {} ms",
indices_duration.as_millis()
);
let total_duration = start_time.elapsed();
info!(
task_id = %task_id,
"[TaskAggregator] aggregated response built in {} ms - signers: {}, non-signers: {}",
total_duration.as_millis(),
digest_aggregated_operators.signers_operator_ids_set.len(),
non_signers_operators_ids.len()
);
debug!(
task_id = %task_id,
"[TaskAggregator] aggregated response stakes: {:?}",
digest_aggregated_operators.signers_total_stake_per_quorum
);
let signers_apk_g2_is_infinity = digest_aggregated_operators.signers_apk_g2.g2().infinity;
let signers_agg_sig_is_infinity = digest_aggregated_operators.signers_agg_sig_g1.g1_point().g1().infinity;
debug!(
task_id = %task_id,
signers_apk_g2_is_infinity = signers_apk_g2_is_infinity,
signers_agg_sig_is_infinity = signers_agg_sig_is_infinity,
non_signers_pub_keys_count = non_signers_pub_keys_g1.len(),
quorum_apks_count = quorum_apks_g1.len(),
"[DEBUG] BLS aggregation result - checking for point-at-infinity issues"
);
if signers_apk_g2_is_infinity {
error!(
task_id = %task_id,
signers_count = digest_aggregated_operators.signers_operator_ids_set.len(),
"[DEBUG] CRITICAL: Signers APK G2 is point at infinity - BN254 pairing will fail. \
This likely means no valid BLS signatures were aggregated."
);
}
if signers_agg_sig_is_infinity {
error!(
task_id = %task_id,
signers_count = digest_aggregated_operators.signers_operator_ids_set.len(),
"[DEBUG] CRITICAL: Aggregated signature is point at infinity - BN254 verification will fail. \
This likely means no valid signatures were collected."
);
}
let signer_ids: Vec<String> = digest_aggregated_operators
.signers_operator_ids_set
.keys()
.map(|id| hex!(id.as_slice()))
.collect();
debug!(
task_id = %task_id,
signer_ids = ?signer_ids,
"[DEBUG] Signer operator IDs in aggregated response"
);
if tracing::enabled!(tracing::Level::DEBUG) {
use ark_ec::CurveGroup;
let mut signers_g1_sum = ark_bn254::G1Affine::identity();
for signer_id in digest_aggregated_operators.signers_operator_ids_set.keys() {
if let Some(state) = operator_state_avs.get(signer_id) {
if let Some(ref pub_keys) = state.operator_info.pub_keys {
signers_g1_sum =
(signers_g1_sum.into_group() + pub_keys.g1_pub_key.g1().into_group()).into_affine();
}
}
}
let mut quorum_apk_sum = ark_bn254::G1Affine::identity();
for apk in quorum_apks_g1.iter() {
quorum_apk_sum = (quorum_apk_sum.into_group() + apk.g1().into_group()).into_affine();
}
let mut non_signers_sum = ark_bn254::G1Affine::identity();
for ns_pk in non_signers_pub_keys_g1.iter() {
non_signers_sum = (non_signers_sum.into_group() + ns_pk.g1().into_group()).into_affine();
}
let expected_signers = (quorum_apk_sum.into_group() - non_signers_sum.into_group()).into_affine();
if let (Some(x1), Some(y1)) = (signers_g1_sum.x(), signers_g1_sum.y()) {
if let (Some(x2), Some(y2)) = (expected_signers.x(), expected_signers.y()) {
let matches = x1 == x2 && y1 == y2;
debug!(
"[DEBUG] BLS_APK_VERIFY: Task {} SIGNERS_G1_SUM X={} Y={} | EXPECTED (quorum-non_signers) X={} Y={} | MATCH={}",
task_id, x1, y1, x2, y2, matches
);
if !matches {
error!(
"[DEBUG] BLS_APK_MISMATCH: signers_G1_sum != quorum_APK - non_signers! This will cause BN254 failure."
);
}
}
}
let signers_apk_g2 = digest_aggregated_operators.signers_apk_g2.g2();
if let (Some(x), Some(y)) = (signers_apk_g2.x(), signers_apk_g2.y()) {
debug!(
"[DEBUG] BLS_SIGNERS_APK_G2: Task {} G2_X_c0={} G2_X_c1={} G2_Y_c0={} G2_Y_c1={}",
task_id, x.c0, x.c1, y.c0, y.c1
);
} else {
error!(
"[DEBUG] BLS_SIGNERS_APK_G2: Task {} G2 point is at infinity - BN254 pairing will fail",
task_id
);
}
}
Ok(BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: digest_aggregated_operators.signers_operator_ids_set.len(),
non_signers_pub_keys_g1,
non_signers_operators_ids,
quorum_apks_g1: quorum_apks_g1.into(),
signers_apk_g2: digest_aggregated_operators.signers_apk_g2,
signers_agg_sig_g1: digest_aggregated_operators.signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: indices.clone().nonSignerQuorumBitmapIndices,
quorum_apk_indices: indices.quorumApkIndices,
total_stake_indices: indices.totalStakeIndices,
non_signer_stake_indices: indices.nonSignerStakeIndices,
})
}
fn check_if_stake_thresholds_met(
signed_stake_per_quorum: &HashMap<u8, U256>,
total_stake_per_quorum: &HashMap<u8, U256>,
quorum_threshold_percentages_map: &HashMap<u8, QuorumThresholdPercentage>,
) -> bool {
for (quorum_num, quorum_threshold_percentage) in quorum_threshold_percentages_map {
let (Some(signed_stake_by_quorum), Some(total_stake_by_quorum)) = (
signed_stake_per_quorum.get(quorum_num),
total_stake_per_quorum.get(quorum_num),
) else {
return false;
};
let signed_stake = signed_stake_by_quorum * U256::from(100);
let threshold_stake = *total_stake_by_quorum * U256::from(*quorum_threshold_percentage);
if signed_stake < threshold_stake {
return false;
}
}
true
}
#[instrument(skip_all, fields(
operator_id = %hex!(signed_digest.operator_id.as_slice()),
digest = %hex!(signed_digest.task_response_digest.as_slice())
))]
fn is_duplicate_signature(
aggregated_operators: &HashMap<FixedBytes<32>, AggregatedOperators>,
signed_digest: &SignedTaskResponseDigest,
) -> bool {
let operator_id_hex = hex!(signed_digest.operator_id.as_slice());
let digest_hex = hex!(signed_digest.task_response_digest.as_slice());
let is_duplicate = aggregated_operators
.get(&signed_digest.task_response_digest)
.map(|ops| {
let has_signed = ops.signers_operator_ids_set.contains_key(&signed_digest.operator_id);
trace!(
"[BlsAggregatorService] operator {} duplicate check for digest {}: {}",
operator_id_hex,
digest_hex,
has_signed
);
has_signed
})
.unwrap_or_else(|| {
trace!(
"[BlsAggregatorService] no existing signatures for digest {} from operator {}",
digest_hex,
operator_id_hex
);
false
});
if is_duplicate {
debug!(
"[BlsAggregatorService] duplicate signature detected from operator {} for digest {}",
operator_id_hex, digest_hex
);
}
is_duplicate
}
#[instrument(skip_all, fields(task_id = %task_id, window_duration = ?window_duration))]
fn start_window(window_tx: &UnboundedSender<bool>, window_duration: Duration, task_id: TaskId) {
let sender = window_tx.clone();
if window_duration.is_zero() {
debug!(
task_id = %task_id,
"[TaskAggregator] window duration is zero - closing immediately",
);
if sender.send(true).is_err() {
error!(
task_id = %task_id,
"[TaskAggregator] failed to send immediate window close signal",
);
}
return;
}
info!(
task_id = %task_id,
"[TaskAggregator] starting aggregation window for {} ms",
window_duration.as_millis()
);
tokio::spawn(async move {
tokio::time::sleep(window_duration).await;
info!(task_id = %task_id, "[TaskAggregator] aggregation window expired");
if sender.send(true).is_err() {
error!(task_id = %task_id, "[TaskAggregator] failed to send window expiry signal");
}
});
}
}
#[instrument(skip_all)]
async fn verify_signature(
task_id: TaskId,
signed_task_response_digest: &SignedTaskResponseDigest,
operator_avs_state: &HashMap<FixedBytes<32>, OperatorAvsState>,
) -> Result<(), SignatureVerificationError> {
debug!(
"operator_avs_state: {}",
operator_avs_state
.iter()
.map(|(k, v)| format!(
"{}: [{}]",
hex!(k),
v.stake_per_quorum
.iter()
.map(|(quorum_num, stake)| format!("{}:{}", quorum_num, stake))
.collect::<Vec<String>>()
.join(", ")
))
.collect::<Vec<String>>()
.join(", ")
);
info!(
"signed_task_response_digest: {}",
hex!(signed_task_response_digest.task_response_digest.as_slice())
);
let Some(operator_state) = operator_avs_state.get(&signed_task_response_digest.operator_id) else {
error!("Operator Not Found for task index: {task_id}");
return Err(SignatureVerificationError::OperatorNotFound);
};
let Some(pub_keys) = &operator_state.operator_info.pub_keys else {
error!("Operator Public Key Not Found for task index: {task_id}");
return Err(SignatureVerificationError::OperatorPublicKeyNotFound);
};
let message = signed_task_response_digest
.task_response_digest
.as_slice()
.try_into()
.map_err(|_| SignatureVerificationError::IncorrectSignature)?;
verify_message(
pub_keys.g2_pub_key.g2(),
message,
signed_task_response_digest.bls_signature.g1_point().g1(),
)
.then_some(())
.ok_or(SignatureVerificationError::IncorrectSignature)
.inspect(|_| {
debug!("Signature verification successful for task index: {task_id}");
})
.inspect_err(|_| {
error!("Signature verification failed for task index: {task_id}");
})
}
pub fn update_aggregated_operators(
task_id: TaskId,
aggregated_operators: &mut HashMap<FixedBytes<32>, AggregatedOperators>,
operator_state: &OperatorAvsState,
task_response_digest: FixedBytes<32>,
bls_signature: Signature,
operator_id: FixedBytes<32>,
) -> Result<AggregatedOperators, BlsAggregationServiceError> {
debug!("Update aggregated operators");
let bls_signature_g1_point = bls_signature.g1_point().g1();
if let Some(existing) = aggregated_operators.get_mut(&task_response_digest) {
aggregate_new_operator(
task_id,
existing,
operator_state.clone(),
operator_id,
bls_signature_g1_point,
)?;
Ok(existing.clone())
} else {
let operator_pub_keys = operator_state.operator_info.pub_keys.clone().ok_or_else(|| {
error!(
task_id = %task_id,
operator_id = %hex!(operator_id.as_slice()),
"Operator public keys not found in operator state"
);
BlsAggregationServiceError::RegistryError {
task_id,
operator_context: format!(" from operator {}", hex!(operator_id.as_slice())),
reason: "operator public keys not found in operator state".to_string(),
}
})?;
let operator_g2_pubkey = operator_pub_keys.g2_pub_key.g2();
let operator_g1_pubkey = operator_pub_keys.g1_pub_key.g1();
if let (Some(x), Some(y)) = (operator_g1_pubkey.x(), operator_g1_pubkey.y()) {
debug!(
"[DEBUG] BLS_SIGNER_G1: First operator {} G1_X={} G1_Y={}",
hex!(operator_id.as_slice()),
x,
y
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_G1: First operator {} G1 point is at infinity",
hex!(operator_id.as_slice())
);
}
if let (Some(x), Some(y)) = (operator_g2_pubkey.x(), operator_g2_pubkey.y()) {
debug!(
"[DEBUG] BLS_SIGNER_G2: First operator {} G2_X_c0={} G2_X_c1={} G2_Y_c0={} G2_Y_c1={}",
hex!(operator_id.as_slice()),
x.c0,
x.c1,
y.c0,
y.c1
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_G2: First operator {} G2 point is at infinity",
hex!(operator_id.as_slice())
);
}
if let (Some(x), Some(y)) = (bls_signature_g1_point.x(), bls_signature_g1_point.y()) {
debug!(
"[DEBUG] BLS_SIGNER_SIG: First operator {} SIG_X={} SIG_Y={}",
hex!(operator_id.as_slice()),
x,
y
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_SIG: First operator {} signature is point at infinity",
hex!(operator_id.as_slice())
);
}
let quorum_count = operator_state.stake_per_quorum.len();
let mut signers_apk_g2 = BlsG2Point::new(G2Affine::zero());
let mut signers_agg_sig_g1 = Signature::new(G1Affine::zero());
for _ in 0..quorum_count {
signers_apk_g2 = BlsG2Point::new((signers_apk_g2.g2() + operator_g2_pubkey).into());
signers_agg_sig_g1 = Signature::new((signers_agg_sig_g1.g1_point().g1() + bls_signature_g1_point).into());
}
debug!(
"[DEBUG] BLS_FIRST_SIGNER_AGG: Operator {} added G2/sig {} times (quorum_count={})",
hex!(operator_id.as_slice()),
quorum_count,
quorum_count
);
Ok(AggregatedOperators {
signers_apk_g2,
signers_agg_sig_g1,
signers_operator_ids_set: HashMap::from([(operator_state.operator_id, true)]),
signers_total_stake_per_quorum: operator_state.stake_per_quorum.clone(),
})
}
}
#[instrument(skip_all)]
pub fn aggregate_new_operator(
task_id: TaskId,
aggregated_operators: &mut AggregatedOperators,
operator_state: OperatorAvsState,
operator_id: FixedBytes<32>,
signature_g1_point: G1Affine,
) -> Result<AggregatedOperators, BlsAggregationServiceError> {
let operator_pub_keys = operator_state.operator_info.pub_keys.clone().ok_or_else(|| {
error!(
task_id = %task_id,
operator_id = %hex!(operator_id.as_slice()),
"Operator public keys not found in operator state for aggregate_new_operator"
);
BlsAggregationServiceError::RegistryError {
task_id,
operator_context: format!(" from operator {}", hex!(operator_id.as_slice())),
reason: "operator public keys not found in operator state for aggregate_new_operator".to_string(),
}
})?;
let operator_g2_pubkey = operator_pub_keys.g2_pub_key.g2();
let operator_g1_pubkey = operator_pub_keys.g1_pub_key.g1();
aggregated_operators.signers_operator_ids_set.insert(operator_id, true);
debug!("operator {operator_id} inserted in signers_operator_ids_set");
if let (Some(x), Some(y)) = (operator_g1_pubkey.x(), operator_g1_pubkey.y()) {
debug!(
"[DEBUG] BLS_SIGNER_G1: Operator {} G1_X={} G1_Y={}",
hex!(operator_id.as_slice()),
x,
y
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_G1: Operator {} G1 point is at infinity",
hex!(operator_id.as_slice())
);
}
if let (Some(x), Some(y)) = (operator_g2_pubkey.x(), operator_g2_pubkey.y()) {
debug!(
"[DEBUG] BLS_SIGNER_G2: Operator {} G2_X_c0={} G2_X_c1={} G2_Y_c0={} G2_Y_c1={}",
hex!(operator_id.as_slice()),
x.c0,
x.c1,
y.c0,
y.c1
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_G2: Operator {} G2 point is at infinity",
hex!(operator_id.as_slice())
);
}
if let (Some(x), Some(y)) = (signature_g1_point.x(), signature_g1_point.y()) {
debug!(
"[DEBUG] BLS_SIGNER_SIG: Operator {} SIG_X={} SIG_Y={}",
hex!(operator_id.as_slice()),
x,
y
);
} else {
debug!(
"[DEBUG] BLS_SIGNER_SIG: Operator {} signature is point at infinity",
hex!(operator_id.as_slice())
);
}
let quorum_count = operator_state.stake_per_quorum.len();
for _ in 0..quorum_count {
aggregated_operators.signers_agg_sig_g1 =
Signature::new((aggregated_operators.signers_agg_sig_g1.g1_point().g1() + signature_g1_point).into());
aggregated_operators.signers_apk_g2 =
BlsG2Point::new((aggregated_operators.signers_apk_g2.g2() + operator_g2_pubkey).into());
}
debug!(
"[DEBUG] BLS_SIGNER_AGG: Operator {} added G2/sig {} times (quorum_count={})",
hex!(operator_id.as_slice()),
quorum_count,
quorum_count
);
for (quorum_num, stake) in operator_state.stake_per_quorum.iter() {
aggregated_operators
.signers_total_stake_per_quorum
.entry(*quorum_num)
.and_modify(|v| *v += stake)
.or_insert(*stake);
}
Ok(aggregated_operators.clone())
}
#[allow(clippy::result_large_err)]
pub fn convert_bls_response_to_contract_format(
response: &BlsAggregationServiceResponse,
) -> Result<IBLSSignatureCheckerTypes::NonSignerStakesAndSignature, crate::error::ChainIoError> {
use ark_ec::AffineRepr;
use eigensdk::crypto_bls::{convert_to_g1_point, convert_to_g2_point};
use newton_core::r#newton_prover_task_manager::BN254::{G1Point, G2Point};
let task_id_hex = newton_core::hex!(response.task_id);
debug!(
"[DEBUG] BLS_CONVERT_START: Task {} SIGNERS={} NON_SIGNERS={} QUORUM_APKS={} CREATED_BLOCK={} MSG_HASH={}",
task_id_hex,
response.signers_count,
response.non_signers_pub_keys_g1.len(),
response.quorum_apks_g1.len(),
response.task_created_block,
response.task_response_digest
);
let mut non_signer_pub_keys = Vec::<G1Point>::new();
for (idx, pub_key) in response.non_signers_pub_keys_g1.iter().enumerate() {
let has_x = pub_key.g1().x().is_some();
if has_x {
let g1 = convert_to_g1_point(pub_key.g1()).map_err(|e| {
error!(
task_id = %task_id_hex,
non_signer_idx = idx,
error = %e,
"[DEBUG] Failed to convert non-signer G1 point"
);
crate::error::ChainIoError::BlsResponseConversionError {
reason: format!("Failed to convert G1 point: {}", e),
}
})?;
debug!(
"[DEBUG] BLS_G1_NON_SIGNER: Task {} idx={} X={} Y={}",
task_id_hex, idx, g1.X, g1.Y
);
non_signer_pub_keys.push(G1Point { X: g1.X, Y: g1.Y });
} else {
debug!(
task_id = %task_id_hex,
non_signer_idx = idx,
"[DEBUG] Non-signer G1 public key has no X coordinate (skipping)"
);
}
}
let mut quorum_apks = Vec::<G1Point>::new();
for (idx, pub_key) in response.quorum_apks_g1.iter().enumerate() {
let is_infinity = pub_key.g1().infinity;
if is_infinity {
error!(
task_id = %task_id_hex,
quorum_idx = idx,
"[DEBUG] Quorum APK is point at infinity - this will cause BN254 failure"
);
return Err(crate::error::ChainIoError::BlsResponseConversionError {
reason: format!("Invalid BLS key found for task {}", task_id_hex),
});
}
let g1 = convert_to_g1_point(pub_key.g1()).map_err(|e| {
error!(
task_id = %task_id_hex,
quorum_idx = idx,
error = %e,
"[DEBUG] Failed to convert quorum APK G1 point"
);
crate::error::ChainIoError::BlsResponseConversionError {
reason: format!("Failed to convert G1 point: {}", e),
}
})?;
debug!(
"[DEBUG] BLS_G1_QUORUM_APK: Task {} idx={} X={} Y={}",
task_id_hex, idx, g1.X, g1.Y
);
quorum_apks.push(G1Point { X: g1.X, Y: g1.Y });
}
let signers_apk_g2_is_infinity = response.signers_apk_g2.g2().infinity;
debug!(
"[DEBUG] BLS_G2_CHECK: Task {} APK_G2_INFINITY={}",
task_id_hex, signers_apk_g2_is_infinity
);
if signers_apk_g2_is_infinity {
error!(
task_id = %task_id_hex,
"[DEBUG] Signers APK G2 is point at infinity - this will cause BN254 pairing failure"
);
}
let apk_g2 = convert_to_g2_point(response.signers_apk_g2.g2()).map_err(|e| {
error!(
task_id = %task_id_hex,
error = %e,
"[DEBUG] Failed to convert signers APK G2 point"
);
crate::error::ChainIoError::BlsResponseConversionError {
reason: format!("Failed to convert G2 point: {}", e),
}
})?;
debug!(
"[DEBUG] BLS_G2: Task {} APK_G2 X=[{}, {}] Y=[{}, {}]",
task_id_hex, apk_g2.X[0], apk_g2.X[1], apk_g2.Y[0], apk_g2.Y[1]
);
let agg_sig_is_infinity = response.signers_agg_sig_g1.g1_point().g1().infinity;
debug!(
"[DEBUG] BLS_G1_SIG_CHECK: Task {} SIG_INFINITY={}",
task_id_hex, agg_sig_is_infinity
);
if agg_sig_is_infinity {
error!(
task_id = %task_id_hex,
"[DEBUG] Aggregated signature is point at infinity - this will cause BN254 pairing failure"
);
}
let sigma = convert_to_g1_point(response.signers_agg_sig_g1.g1_point().g1()).map_err(|e| {
error!(
task_id = %task_id_hex,
error = %e,
"[DEBUG] Failed to convert aggregated signature G1 point"
);
crate::error::ChainIoError::BlsResponseConversionError {
reason: format!("Failed to convert G1 point: {}", e),
}
})?;
debug!("[DEBUG] BLS_G1_SIGMA: Task {} X={} Y={}", task_id_hex, sigma.X, sigma.Y);
debug!(
"[DEBUG] BLS_INDICES: Task {} NON_SIGNER_BITMAP_IDX={:?} QUORUM_APK_IDX={:?} TOTAL_STAKE_IDX={:?} NON_SIGNER_STAKE_IDX_COUNT={}",
task_id_hex,
response.non_signer_quorum_bitmap_indices,
response.quorum_apk_indices,
response.total_stake_indices,
response.non_signer_stake_indices.len()
);
let non_signer_stakes_and_signature =
newton_core::r#newton_prover_task_manager::IBLSSignatureCheckerTypes::NonSignerStakesAndSignature {
nonSignerPubkeys: non_signer_pub_keys.clone(),
nonSignerQuorumBitmapIndices: response.non_signer_quorum_bitmap_indices.clone(),
quorumApks: quorum_apks.clone(),
apkG2: G2Point {
X: apk_g2.X,
Y: apk_g2.Y,
},
sigma: G1Point { X: sigma.X, Y: sigma.Y },
quorumApkIndices: response.quorum_apk_indices.clone(),
totalStakeIndices: response.total_stake_indices.clone(),
nonSignerStakeIndices: response.non_signer_stake_indices.clone(),
};
info!(
task_id = %task_id_hex,
non_signer_pubkeys_count = non_signer_pub_keys.len(),
quorum_apks_count = quorum_apks.len(),
signers_apk_g2_is_infinity = signers_apk_g2_is_infinity,
agg_sig_is_infinity = agg_sig_is_infinity,
"[DEBUG] BLS_CONVERT: Complete for task {} - NON_SIGNER_PUBKEYS={}, QUORUM_APKS={}, APK_G2_INFINITY={}, SIG_INFINITY={}",
task_id_hex,
non_signer_pub_keys.len(),
quorum_apks.len(),
signers_apk_g2_is_infinity,
agg_sig_is_infinity
);
Ok(non_signer_stakes_and_signature)
}
#[cfg(test)]
mod tests {
use super::{
BlsAggregationServiceError, BlsAggregationServiceResponse, BlsAggregatorService, BlsServiceHandle,
TaskMetadata, TaskSignature,
};
use alloy::primitives::{keccak256, FixedBytes, B256, U256};
use eigensdk::crypto_bls::{BlsG1Point, BlsG2Point, BlsKeyPair, Signature};
use eigensdk::{
services_avsregistry::fake_avs_registry_service::FakeAvsRegistryService,
types::{
avs::{
SignatureVerificationError,
SignatureVerificationError::{DuplicateSignature, IncorrectSignature},
},
operator::{QuorumNum, QuorumThresholdPercentages},
test::TestOperator,
},
};
use newton_core::TaskId;
use sha2::{Digest, Sha256};
use std::{collections::HashMap, time::Duration, vec};
use tokio::time::{sleep, Instant};
const PRIVATE_KEY_1: &str = "13710126902690889134622698668747132666439281256983827313388062967626731803599";
const PRIVATE_KEY_2: &str = "14610126902690889134622698668747132666439281256983827313388062967626731803500";
const PRIVATE_KEY_3: &str = "15610126902690889134622698668747132666439281256983827313388062967626731803501";
fn hash(task_response: u64) -> B256 {
let mut hasher = Sha256::new();
hasher.update(task_response.to_be_bytes());
B256::from_slice(hasher.finalize().as_ref())
}
fn aggregate_g1_public_keys(operators: &[TestOperator]) -> BlsG1Point {
operators
.iter()
.map(|op| op.bls_keypair.public_key().g1())
.reduce(|a, b| (a + b).into())
.map(BlsG1Point::new)
.unwrap()
}
fn aggregate_g2_public_keys(operators: &[TestOperator]) -> BlsG2Point {
operators
.iter()
.map(|op| op.bls_keypair.public_key_g2().g2())
.reduce(|a, b| (a + b).into())
.map(BlsG2Point::new)
.unwrap()
}
fn aggregate_g1_signatures(signatures: &[Signature]) -> Signature {
let agg = signatures
.iter()
.map(|s| s.g1_point().g1())
.reduce(|a, b| (a + b).into())
.unwrap();
Signature::new(agg)
}
#[tokio::test]
async fn test_1_quorum_1_operator_1_correct_signature() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(1).into();
let task_created_block = 1;
let quorum_numbers = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123;
let task_response_digest = hash(task_response);
let bls_signature = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_signature,
test_operator_1.operator_id,
))
.await
.unwrap();
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 1,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![test_operator_1.bls_keypair.public_key()],
signers_apk_g2: test_operator_1.bls_keypair.public_key_g2(),
signers_agg_sig_g1: test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref()),
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_1_quorum_2_operator_2_duplicated_signatures() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(1).into();
let task_created_block = block_number;
let quorum_numbers = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_signature_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_signature_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let second_signature_processing_result = handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_signature_1.clone(),
test_operator_1.operator_id,
))
.await;
assert!(matches!(
second_signature_processing_result,
Err(BlsAggregationServiceError::SignatureVerificationError {
task_id: t,
operator_id: op_id,
verification_error: SignatureVerificationError::DuplicateSignature,
}) if t == task_id && op_id == test_operator_1.operator_id
));
let bls_signature_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_signature_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let quorum_apks_g1 = aggregate_g1_public_keys(&test_operators);
let signers_apk_g2 = aggregate_g2_public_keys(&test_operators);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_signature_1, bls_signature_2]);
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![quorum_apks_g1],
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_1_quorum_3_operator_3_correct_signatures() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operator_3 = TestOperator {
operator_id: U256::from(3).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(300)), (1u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_3.into()).unwrap(),
};
let test_operators = vec![
test_operator_1.clone(),
test_operator_2.clone(),
test_operator_3.clone(),
];
let block_number = 1;
let task_id: TaskId = U256::from(1).into();
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_op_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let bls_sig_op_3 = test_operator_3.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_3.clone(),
test_operator_3.operator_id,
))
.await
.unwrap();
let quorum_apks_g1 = aggregate_g1_public_keys(&test_operators);
let signers_apk_g2 = aggregate_g2_public_keys(&test_operators);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1, bls_sig_op_2, bls_sig_op_3]);
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block: block_number,
task_response_digest,
signers_count: 3,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![quorum_apks_g1],
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_2_quorum_2_operator_2_correct_signatures() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100u8, 100u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_op_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let quorum_apks_g1 = aggregate_g1_public_keys(&test_operators);
let signers_apk_g2 = aggregate_g2_public_keys(&[test_operators.clone(), test_operators].concat());
let signers_agg_sig_g1 =
aggregate_g1_signatures(&[bls_sig_op_1.clone(), bls_sig_op_1, bls_sig_op_2.clone(), bls_sig_op_2]);
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![quorum_apks_g1.clone(), quorum_apks_g1],
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
}
#[tokio::test]
async fn test_2_concurrent_tasks_2_quorum_2_operator_2_correct_signatures() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100u8, 100u8];
let time_to_expiry = Duration::from_secs(1);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let task_1_id: TaskId = U256::from(1).into();
let task_1_response = 123; let task_1_response_digest = hash(task_1_response);
let metadata1 = TaskMetadata::new(
task_1_id,
block_number,
quorum_numbers.clone(),
quorum_threshold_percentages.clone(),
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_1_receiver = handle.initialize_task(metadata1).await.unwrap();
let task_2_id: TaskId = U256::from(2).into();
let task_2_response = 234; let task_2_response_digest = hash(task_2_response);
let metadata2 = TaskMetadata::new(
task_2_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let mut task_2_receiver = handle.initialize_task(metadata2).await.unwrap();
let bls_sig_task_1_op_1 = test_operator_1
.bls_keypair
.sign_message(task_1_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_1_id,
task_1_response_digest,
bls_sig_task_1_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_task_1_op_2 = test_operator_2
.bls_keypair
.sign_message(task_1_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_1_id,
task_1_response_digest,
bls_sig_task_1_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let bls_sig_task_2_op_1 = test_operator_1
.bls_keypair
.sign_message(task_2_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_2_id,
task_2_response_digest,
bls_sig_task_2_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_task_2_op_2 = test_operator_2
.bls_keypair
.sign_message(task_2_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_2_id,
task_2_response_digest,
bls_sig_task_2_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let quorum_apks_g1 = aggregate_g1_public_keys(&test_operators);
let signers_apk_g2 = aggregate_g2_public_keys(&[test_operators.clone(), test_operators].concat());
let signers_agg_sig_g1_task_1 = aggregate_g1_signatures(&[
bls_sig_task_1_op_1.clone(),
bls_sig_task_1_op_1,
bls_sig_task_1_op_2.clone(),
bls_sig_task_1_op_2,
]);
let expected_response_task_1 = BlsAggregationServiceResponse {
task_id: task_1_id,
task_created_block: block_number,
task_response_digest: task_1_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![quorum_apks_g1.clone(), quorum_apks_g1.clone()],
signers_apk_g2: signers_apk_g2.clone(),
signers_agg_sig_g1: signers_agg_sig_g1_task_1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let signers_agg_sig_g1_task_2 = aggregate_g1_signatures(&[
bls_sig_task_2_op_1.clone(),
bls_sig_task_2_op_1,
bls_sig_task_2_op_2.clone(),
bls_sig_task_2_op_2,
]);
let expected_response_task_2 = BlsAggregationServiceResponse {
task_id: task_2_id,
task_created_block: block_number,
task_response_digest: task_2_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![quorum_apks_g1.clone(), quorum_apks_g1.clone()],
signers_apk_g2,
signers_agg_sig_g1: signers_agg_sig_g1_task_2,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let first_response = task_1_receiver
.recv()
.await
.expect("task 1 receiver channel should not be closed")
.expect("should receive response from task 1");
let second_response = task_2_receiver
.recv()
.await
.expect("task 2 receiver channel should not be closed")
.expect("should receive response from task 2");
let (task_1_response, task_2_response) = if first_response.task_id == task_1_id {
(first_response, second_response)
} else {
(second_response, first_response)
};
assert_eq!(expected_response_task_1, task_1_response);
assert_eq!(expected_response_task_2, task_2_response);
}
#[tokio::test]
async fn test_1_quorum_1_operator_0_signatures_task_expired() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100];
let time_to_expiry = Duration::from_secs(1);
let _task_response = 123;
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_1_quorum_2_operator_1_signatures_50_threshold() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![50u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let quorum_apks_g1 = aggregate_g1_public_keys(&test_operators);
let signers_apk_g2: BlsG2Point = test_operator_1.bls_keypair.public_key_g2();
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 1,
non_signers_pub_keys_g1: vec![test_operator_2.bls_keypair.public_key()],
non_signers_operators_ids: vec![test_operator_2.operator_id],
quorum_apks_g1: vec![quorum_apks_g1],
signers_apk_g2,
signers_agg_sig_g1: bls_sig_op_1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_1_quorum_2_operator_1_signatures_60_threshold() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![60u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await
.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_2_quorums_2_operators_which_just_take_1_quorum_2_correct_signatures() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100u8, 100u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_op_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let signers_apk_g2 = aggregate_g2_public_keys(&test_operators);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1, bls_sig_op_2]);
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![
test_operator_1.bls_keypair.public_key(),
test_operator_2.bls_keypair.public_key(),
],
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_2_quorums_3_operators_which_just_stake_1_quorum_50_threshold() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operator_3 = TestOperator {
operator_id: U256::from(3).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_3.into()).unwrap(),
};
let test_operators = vec![
test_operator_1.clone(),
test_operator_2.clone(),
test_operator_3.clone(),
];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![50u8, 50u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators.clone());
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_op_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let signers_apk_g2 = aggregate_g2_public_keys(&[test_operator_1.clone(), test_operator_2.clone()]);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1, bls_sig_op_2]);
let quorum_apks_g1 = vec![
aggregate_g1_public_keys(&[test_operator_1, test_operator_3.clone()]),
aggregate_g1_public_keys(&[test_operator_2, test_operator_3.clone()]),
];
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![test_operator_3.bls_keypair.public_key()],
non_signers_operators_ids: vec![test_operator_3.operator_id],
quorum_apks_g1,
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
}
#[tokio::test]
async fn test_2_quorums_3_operators_which_just_stake_1_quorum_60_threshold() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operator_3 = TestOperator {
operator_id: U256::from(3).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_3.into()).unwrap(),
};
let test_operators = vec![
test_operator_1.clone(),
test_operator_2.clone(),
test_operator_3.clone(),
];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![60u8, 60u8];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await
.unwrap();
let bls_sig_op_2 = test_operator_2.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_2,
test_operator_2.operator_id,
))
.await
.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_2_quorums_1_operator_which_just_take_1_quorum_1_signature_task_expired() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100, 100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await
.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_2_quorums_2_operators_where_1_operator_just_take_1_quorum_1_signature_task_expired() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers: Vec<QuorumNum> = vec![0, 1];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100, 100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123; let task_response_digest = hash(task_response);
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await
.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn send_signature_of_task_not_initialized() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_response = 123; let task_response_digest = hash(task_response);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let bls_sig_op_1 = test_operator_1.bls_keypair.sign_message(task_response_digest.as_ref());
let (handle, _) = bls_agg_service.start();
let result = handle
.process_signature(TaskSignature::new(
task_id,
task_response_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await;
assert!(matches!(
result,
Err(BlsAggregationServiceError::TaskNotFound { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_1_quorum_2_operator_2_signatures_on_2_different_msgs() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100u8];
let time_to_expiry = Duration::from_secs(1);
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let task_response_1 = 123; let task_response_1_digest = hash(task_response_1);
let bls_sig_op_1 = test_operator_1
.bls_keypair
.sign_message(task_response_1_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_1_digest,
bls_sig_op_1,
test_operator_1.operator_id,
))
.await
.unwrap();
let task_response_2 = 456; let task_response_2_digest = hash(task_response_2);
let bls_sig_op_2 = test_operator_1
.bls_keypair
.sign_message(task_response_2_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_2_digest,
bls_sig_op_2,
test_operator_1.operator_id,
))
.await
.unwrap();
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_1_quorum_1_operator_1_invalid_signature() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let quorum_numbers = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123;
let wrong_task_response_digest = hash(task_response + 1);
let bls_signature = test_operator_1.bls_keypair.sign_message(hash(task_response).as_ref());
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let result = handle
.process_signature(TaskSignature::new(
task_id,
wrong_task_response_digest,
bls_signature.clone(),
test_operator_1.operator_id,
))
.await;
assert!(matches!(
result,
Err(BlsAggregationServiceError::SignatureVerificationError {
task_id: t,
operator_id: op_id,
verification_error: SignatureVerificationError::IncorrectSignature,
}) if t == task_id && op_id == test_operator_1.operator_id
));
let response = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed");
assert!(matches!(
response,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
}
#[tokio::test]
async fn test_signatures_are_processed_during_window_after_quorum() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operator_3 = TestOperator {
operator_id: U256::from(3).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_3.into()).unwrap(),
};
let test_operators = vec![
test_operator_1.clone(),
test_operator_2.clone(),
test_operator_3.clone(),
];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let task_response = 123;
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![50_u8];
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let time_to_expiry = Duration::from_secs(5);
let window_duration = Duration::from_secs(1);
let start = Instant::now();
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
)
.with_window_duration(window_duration);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let task_response_1_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1
.bls_keypair
.sign_message(task_response_1_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_1_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let task_response_2_digest = hash(task_response);
let bls_sig_op_2 = test_operator_2
.bls_keypair
.sign_message(task_response_2_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_2_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
sleep(Duration::from_millis(500)).await;
let task_response_3_digest = hash(task_response);
let bls_sig_op_3 = test_operator_3
.bls_keypair
.sign_message(task_response_3_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_3_digest,
bls_sig_op_3.clone(),
test_operator_3.operator_id,
))
.await
.unwrap();
let signers_apk_g2 = aggregate_g2_public_keys(&[
test_operator_1.clone(),
test_operator_2.clone(),
test_operator_3.clone(),
]);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1, bls_sig_op_2, bls_sig_op_3]);
let quorum_apks_g1 = vec![aggregate_g1_public_keys(&[
test_operator_1,
test_operator_2,
test_operator_3,
])];
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest: task_response_3_digest,
signers_count: 3,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1,
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
let elapsed = start.elapsed();
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
assert!(elapsed < time_to_expiry);
assert!(elapsed >= window_duration);
}
#[tokio::test]
async fn test_if_quorum_has_been_reached_and_the_task_expires_during_window_the_response_is_sent() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let task_response = 123;
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![40_u8];
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let time_to_expiry = Duration::from_secs(2);
let window_duration = Duration::from_secs(10);
let start = Instant::now();
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
)
.with_window_duration(window_duration);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let task_response_1_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1
.bls_keypair
.sign_message(task_response_1_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_1_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let task_response_2_digest = hash(task_response);
let bls_sig_op_2 = test_operator_2
.bls_keypair
.sign_message(task_response_2_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_2_digest,
bls_sig_op_2.clone(),
test_operator_2.operator_id,
))
.await
.unwrap();
let signers_apk_g2 = aggregate_g2_public_keys(&[test_operator_1.clone(), test_operator_2.clone()]);
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1, bls_sig_op_2]);
let quorum_apks_g1 = vec![aggregate_g1_public_keys(&[test_operator_1, test_operator_2])];
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest: task_response_2_digest,
signers_count: 2,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1,
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
let elapsed = start.elapsed();
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
assert!(elapsed >= time_to_expiry);
assert!(elapsed < window_duration);
}
#[tokio::test]
async fn test_if_window_duration_is_zero_no_signatures_are_aggregated_after_reaching_quorum() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let task_response = 123;
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![40_u8];
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let time_to_expiry = Duration::from_secs(2);
let window_duration = Duration::ZERO;
let start = Instant::now();
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
)
.with_window_duration(window_duration);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let task_response_1_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1
.bls_keypair
.sign_message(task_response_1_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_1_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
sleep(Duration::from_millis(1)).await;
let task_response_2_digest = hash(task_response);
let bls_sig_op_2 = test_operator_2
.bls_keypair
.sign_message(task_response_2_digest.as_ref());
let process_signature_result = handle
.process_signature(TaskSignature::new(
task_id,
task_response_2_digest,
bls_sig_op_2,
test_operator_2.operator_id,
))
.await;
assert!(matches!(
process_signature_result,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
let signers_apk_g2 = aggregate_g2_public_keys(std::slice::from_ref(&test_operator_1));
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1]);
let quorum_apks_g1 = vec![aggregate_g1_public_keys(&[test_operator_1, test_operator_2.clone()])];
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest: task_response_1_digest,
signers_count: 1,
non_signers_pub_keys_g1: vec![test_operator_2.bls_keypair.public_key()],
non_signers_operators_ids: vec![test_operator_2.operator_id],
quorum_apks_g1,
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
let elapsed = start.elapsed();
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
assert!(elapsed < time_to_expiry);
}
#[tokio::test]
async fn test_no_signatures_are_aggregated_after_window() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let test_operator_2 = TestOperator {
operator_id: U256::from(2).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_2.into()).unwrap(),
};
let test_operators = vec![test_operator_1.clone(), test_operator_2.clone()];
let block_number = 1;
let task_id: TaskId = U256::from(0).into();
let task_created_block = block_number;
let task_response = 123;
let quorum_numbers: Vec<QuorumNum> = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![40_u8];
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, test_operators);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let time_to_expiry = Duration::from_secs(5);
let window_duration = Duration::from_secs(1);
let start = Instant::now();
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
)
.with_window_duration(window_duration);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
let task_response_1_digest = hash(task_response);
let bls_sig_op_1 = test_operator_1
.bls_keypair
.sign_message(task_response_1_digest.as_ref());
handle
.process_signature(TaskSignature::new(
task_id,
task_response_1_digest,
bls_sig_op_1.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
sleep(Duration::from_secs(2)).await;
let task_response_2_digest = hash(task_response);
let bls_sig_op_2 = test_operator_2
.bls_keypair
.sign_message(task_response_2_digest.as_ref());
let process_signature_result = handle
.process_signature(TaskSignature::new(
task_id,
task_response_2_digest,
bls_sig_op_2,
test_operator_2.operator_id,
))
.await;
assert!(matches!(
process_signature_result,
Err(BlsAggregationServiceError::TaskExpired { task_id: t, .. }) if t == task_id
));
let signers_apk_g2 = aggregate_g2_public_keys(std::slice::from_ref(&test_operator_1));
let signers_agg_sig_g1 = aggregate_g1_signatures(&[bls_sig_op_1]);
let quorum_apks_g1 = vec![aggregate_g1_public_keys(&[test_operator_1, test_operator_2.clone()])];
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest: task_response_1_digest,
signers_count: 1,
non_signers_pub_keys_g1: vec![test_operator_2.bls_keypair.public_key()],
non_signers_operators_ids: vec![test_operator_2.operator_id],
quorum_apks_g1,
signers_apk_g2,
signers_agg_sig_g1,
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
let elapsed = start.elapsed();
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
assert!(elapsed < time_to_expiry);
}
#[tokio::test]
async fn test_1_quorum_1_operator_1_correct_signature_multichain() {
let test_operator_1 = TestOperator {
operator_id: U256::from(1).into(),
stake_per_quorum: HashMap::from([(0u8, U256::from(100)), (1u8, U256::from(200))]),
bls_keypair: BlsKeyPair::new(PRIVATE_KEY_1.into()).unwrap(),
};
let block_number = 1;
let task_id: TaskId = U256::from(1).into();
let task_created_block = 1;
let quorum_numbers = vec![0];
let quorum_threshold_percentages: QuorumThresholdPercentages = vec![100];
let time_to_expiry = Duration::from_secs(1);
let task_response = 123;
let reference_timestamp: u32 = 1000;
let task_response_digest = hash(task_response);
let bn254_certificate_typehash = keccak256("BN254Certificate(uint32 referenceTimestamp,bytes32 messageHash)");
let mut hasher = Sha256::new();
hasher.update(bn254_certificate_typehash);
hasher.update([0u8; 28]);
hasher.update(reference_timestamp.to_be_bytes());
hasher.update(task_response_digest);
let digest = FixedBytes::from_slice(hasher.finalize().as_ref());
let bls_signature = test_operator_1.bls_keypair.sign_message(digest.as_ref());
let fake_avs_registry_service = FakeAvsRegistryService::new(block_number, vec![test_operator_1.clone()]);
let bls_agg_service = BlsAggregatorService::new(fake_avs_registry_service);
let metadata = TaskMetadata::new(
task_id,
block_number,
quorum_numbers,
quorum_threshold_percentages,
time_to_expiry,
);
let (handle, _aggregator_response) = bls_agg_service.start();
let mut task_receiver = handle.initialize_task(metadata).await.unwrap();
handle
.process_signature(TaskSignature::new(
task_id,
digest,
bls_signature.clone(),
test_operator_1.operator_id,
))
.await
.unwrap();
let expected_agg_service_response = BlsAggregationServiceResponse {
task_id,
task_created_block,
task_response_digest: digest,
signers_count: 1,
non_signers_pub_keys_g1: vec![],
non_signers_operators_ids: vec![],
quorum_apks_g1: vec![test_operator_1.bls_keypair.public_key()],
signers_apk_g2: test_operator_1.bls_keypair.public_key_g2(),
signers_agg_sig_g1: test_operator_1.bls_keypair.sign_message(digest.as_ref()),
non_signer_quorum_bitmap_indices: vec![],
quorum_apk_indices: vec![],
total_stake_indices: vec![],
non_signer_stake_indices: vec![],
};
let actual = task_receiver
.recv()
.await
.expect("task receiver channel should not be closed")
.expect("should receive successful response");
assert_eq!(expected_agg_service_response, actual);
assert_eq!(task_id, actual.task_id);
assert_eq!(actual.signers_apk_g2, test_operator_1.bls_keypair.public_key_g2());
assert_eq!(actual.signers_agg_sig_g1.g1_point(), bls_signature.g1_point());
assert!(actual.non_signers_pub_keys_g1.is_empty());
}
}