Skip to content

Commit 771bfd3

Browse files
Merge pull request #328 from originalworks/blobs-batch-sender-final-implementation
blobs_batch_sender - final implementation
2 parents b1ca6dd + 301718e commit 771bfd3

13 files changed

Lines changed: 389 additions & 72 deletions

File tree

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎aws/blobs_batch_sender/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,3 +14,4 @@ owen = { workspace = true, features = ["aws-integration"] }
1414
serde = { workspace = true }
1515
serde_json = { workspace = true }
1616
aws-sdk-s3 = "1.98.0"
17+
alloy = { version = "1.0.32", features = ["full"] }
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
use crate::contract::sEOA::SubmitNewBlobInput;
2+
use alloy::{
3+
consensus::BlobTransactionSidecar,
4+
network::EthereumWallet,
5+
primitives::{Address, Bytes, FixedBytes},
6+
providers::{
7+
Identity, ProviderBuilder,
8+
fillers::{
9+
BlobGasFiller, ChainIdFiller, FillProvider, GasFiller, JoinFill, NonceFiller,
10+
WalletFiller,
11+
},
12+
},
13+
sol,
14+
};
15+
use blobs_batch_sender::BlobsBatchSenderConfig;
16+
use lambda_runtime::Error;
17+
use owen::blobs_queue::BlobsQueueS3JsonFile;
18+
19+
sol!(
20+
#[allow(missing_docs)]
21+
#[sol(rpc)]
22+
sEOA,
23+
"../../submodules/account-abstraction/artifacts/contracts/sEOA.sol/sEOA.json"
24+
);
25+
26+
type HardlyTypedProvider = FillProvider<
27+
JoinFill<
28+
JoinFill<
29+
Identity,
30+
JoinFill<GasFiller, JoinFill<BlobGasFiller, JoinFill<NonceFiller, ChainIdFiller>>>,
31+
>,
32+
WalletFiller<EthereumWallet>,
33+
>,
34+
alloy::providers::RootProvider,
35+
>;
36+
37+
struct SendBatchTxInput {
38+
tx_params: Vec<SubmitNewBlobInput>,
39+
sidecar: BlobTransactionSidecar,
40+
}
41+
42+
pub struct SmartEoaManager {
43+
ddex_sequencer_address: Address,
44+
s_eoa: sEOA::sEOAInstance<HardlyTypedProvider>,
45+
}
46+
47+
impl SmartEoaManager {
48+
pub fn build(config: &BlobsBatchSenderConfig, wallet: EthereumWallet) -> Result<Self, Error> {
49+
let provider = ProviderBuilder::new()
50+
.wallet(wallet)
51+
.connect_http(config.rpc_url.parse()?);
52+
53+
let s_eoa = sEOA::new(config.s_eoa_address, provider);
54+
55+
Ok(Self {
56+
ddex_sequencer_address: config.ddex_sequencer_address,
57+
s_eoa,
58+
})
59+
}
60+
61+
fn parse_batch_input(
62+
blob_tx_data_vec: Vec<BlobsQueueS3JsonFile>,
63+
) -> Result<SendBatchTxInput, Error> {
64+
let tx_params: Vec<SubmitNewBlobInput> = blob_tx_data_vec
65+
.iter()
66+
.map(|blob_json_file| SubmitNewBlobInput {
67+
imageId: blob_json_file.image_id,
68+
commitment: Bytes::from(blob_json_file.tx_data.kzg_commitment.to_vec()),
69+
blobSha2: FixedBytes::<32>::from(blob_json_file.tx_data.blob_sha2),
70+
})
71+
.collect();
72+
73+
let mut sidecar = BlobTransactionSidecar {
74+
blobs: Vec::new(),
75+
commitments: Vec::new(),
76+
proofs: Vec::new(),
77+
};
78+
79+
for blob_json_file in blob_tx_data_vec {
80+
sidecar.blobs.push(
81+
blob_json_file
82+
.tx_data
83+
.blob_sidecar
84+
.blobs
85+
.first()
86+
.expect("Blob missing in input")
87+
.clone(),
88+
);
89+
sidecar.commitments.push(
90+
blob_json_file
91+
.tx_data
92+
.blob_sidecar
93+
.commitments
94+
.first()
95+
.expect("commitment missing in input")
96+
.clone(),
97+
);
98+
sidecar.proofs.push(
99+
blob_json_file
100+
.tx_data
101+
.blob_sidecar
102+
.proofs
103+
.first()
104+
.expect("proof missing in input")
105+
.clone(),
106+
);
107+
}
108+
Ok(SendBatchTxInput { tx_params, sidecar })
109+
}
110+
111+
pub async fn send_batch(
112+
&self,
113+
blob_tx_data_vec: Vec<BlobsQueueS3JsonFile>,
114+
) -> Result<(), Error> {
115+
let batch_input = Self::parse_batch_input(blob_tx_data_vec)?;
116+
let mut _tx_builder = self
117+
.s_eoa
118+
.submitNewBlobBatch(batch_input.tx_params, self.ddex_sequencer_address)
119+
.sidecar(batch_input.sidecar)
120+
.max_fee_per_blob_gas(1000000001);
121+
Ok(())
122+
}
123+
}

‎aws/blobs_batch_sender/src/event_handler.rs‎

Lines changed: 22 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,29 @@
11
use aws_lambda_events::event::sqs::SqsEvent;
2+
use blobs_batch_sender::BlobsBatchSenderConfig;
23
use lambda_runtime::{Error, LambdaEvent};
3-
use owen::blobs_queue::BlobsQueueMessageBody;
4+
use owen::{
5+
blobs_queue::BlobsQueueMessageBody,
6+
wallet::{OwenWallet, OwenWalletConfig},
7+
};
48

5-
use crate::s3;
9+
use crate::{contract::SmartEoaManager, s3::BlobsStorage};
610

711
pub(crate) async fn function_handler(event: LambdaEvent<SqsEvent>) -> Result<(), Error> {
12+
let config = BlobsBatchSenderConfig::build()?;
13+
let blobs_storage = BlobsStorage::build(&config).await?;
14+
let owen_wallet_config = OwenWalletConfig::from(&config)?;
15+
let owen_wallet = OwenWallet::build(&owen_wallet_config).await?;
16+
let smart_eoa_manager = SmartEoaManager::build(&config, owen_wallet.wallet)?;
17+
18+
let blobhashes = extract_blobhashes(event)?;
19+
let blob_tx_data_vec = blobs_storage.read(blobhashes).await?;
20+
21+
smart_eoa_manager.send_batch(blob_tx_data_vec).await?;
22+
23+
Ok(())
24+
}
25+
26+
fn extract_blobhashes(event: LambdaEvent<SqsEvent>) -> Result<Vec<String>, Error> {
827
let messages: Vec<BlobsQueueMessageBody> = event
928
.payload
1029
.records
@@ -24,8 +43,5 @@ pub(crate) async fn function_handler(event: LambdaEvent<SqsEvent>) -> Result<(),
2443
.iter()
2544
.map(|message| message.blobhash.clone())
2645
.collect();
27-
28-
s3::read_blobs(blobhashes).await?;
29-
30-
Ok(())
46+
Ok(blobhashes)
3147
}

‎aws/blobs_batch_sender/src/lib.rs‎

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,75 @@
1+
use alloy::primitives::Address;
2+
use lambda_runtime::Error;
3+
use owen::constants;
4+
use owen::wallet::HasOwenWalletFields;
5+
use std::env;
6+
use std::str::FromStr;
7+
8+
impl HasOwenWalletFields for BlobsBatchSenderConfig {
9+
fn use_kms(&self) -> bool {
10+
self.use_kms
11+
}
12+
fn rpc_url(&self) -> String {
13+
self.rpc_url.clone()
14+
}
15+
fn private_key(&self) -> Option<String> {
16+
self.private_key.clone()
17+
}
18+
fn signer_kms_id(&self) -> Option<String> {
19+
self.signer_kms_id.clone()
20+
}
21+
}
22+
23+
pub struct BlobsBatchSenderConfig {
24+
pub ddex_sequencer_address: Address,
25+
pub s_eoa_address: Address,
26+
pub rpc_url: String,
27+
pub blobs_temp_storage_bucket_name: String,
28+
pub use_kms: bool,
29+
pub private_key: Option<String>,
30+
pub signer_kms_id: Option<String>,
31+
}
32+
33+
impl BlobsBatchSenderConfig {
34+
fn get_env_var(key: &str) -> String {
35+
env::var(key).expect(format!("Missing env variable: {key}").as_str())
36+
}
37+
pub fn build() -> Result<BlobsBatchSenderConfig, Error> {
38+
let rpc_url = Self::get_env_var("RPC_URL");
39+
let ddex_sequencer_address = Address::from_str(
40+
std::env::var("DDEX_SEQUENCER_ADDRESS")
41+
.unwrap_or_else(|_| constants::DDEX_SEQUENCER_ADDRESS.to_string())
42+
.as_str(),
43+
)
44+
.expect("Could not parse ddex sequencer address");
45+
46+
let s_eoa_address = Address::from_str(Self::get_env_var("S_EOA_ADDRESS").as_str())?;
47+
48+
let blobs_temp_storage_bucket_name = Self::get_env_var("BLOBS_TEMP_STORAGE_BUCKET_NAME");
49+
50+
let mut signer_kms_id = None;
51+
let mut private_key = None;
52+
let use_kms = matches!(
53+
std::env::var("USE_KMS")
54+
.unwrap_or_else(|_| "false".to_string())
55+
.as_str(),
56+
"1" | "true"
57+
);
58+
59+
if use_kms {
60+
signer_kms_id = Some(Self::get_env_var("SIGNER_KMS_ID"));
61+
} else {
62+
private_key = Some(Self::get_env_var("PRIVATE_KEY"));
63+
}
64+
65+
Ok(BlobsBatchSenderConfig {
66+
ddex_sequencer_address,
67+
rpc_url,
68+
s_eoa_address,
69+
blobs_temp_storage_bucket_name,
70+
private_key,
71+
use_kms,
72+
signer_kms_id,
73+
})
74+
}
75+
}

‎aws/blobs_batch_sender/src/main.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
mod contract;
12
mod event_handler;
23
mod s3;
34
use event_handler::function_handler;

‎aws/blobs_batch_sender/src/s3.rs‎

Lines changed: 38 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -1,37 +1,48 @@
1-
use std::env;
2-
31
use aws_config::meta::region::RegionProviderChain;
42
use aws_sdk_s3::Client;
3+
use blobs_batch_sender::BlobsBatchSenderConfig;
54
use lambda_runtime::Error;
6-
use owen::blob::BlobTransactionData;
5+
use owen::blobs_queue::BlobsQueueS3JsonFile;
76
use tokio::io::AsyncReadExt;
87

9-
pub async fn read_blobs(blobhashes: Vec<String>) -> Result<(), Error> {
10-
let region_provider = RegionProviderChain::default_provider().or_else("us-east-1");
11-
let config = aws_config::from_env().region(region_provider).load().await;
12-
let client = Client::new(&config);
13-
14-
let bucket = env::var("BLOBS_TEMP_STORAGE_BUCKET_NAME")
15-
.expect(format!("Missing env variable: BLOBS_TEMP_STORAGE_BUCKET_NAME").as_str());
16-
17-
let mut blobs: Vec<BlobTransactionData> = Vec::new();
8+
pub struct BlobsStorage {
9+
client: Client,
10+
bucket_name: String,
11+
}
1812

19-
for blobhash in blobhashes {
20-
let key = format!("blobs/{}.json", blobhash);
21-
let resp = client.get_object().bucket(&bucket).key(key).send().await?;
13+
impl BlobsStorage {
14+
pub async fn build(config: &BlobsBatchSenderConfig) -> Result<Self, Error> {
15+
let region_provider = RegionProviderChain::default_provider().or_else("us-east-1");
16+
let aws_config = aws_config::from_env().region(region_provider).load().await;
17+
let client: Client = Client::new(&aws_config);
2218

23-
let mut body = resp.body.into_async_read();
24-
let mut contents = String::new();
25-
body.read_to_string(&mut contents).await?;
26-
let blob_transaction_data: BlobTransactionData = serde_json::from_str(&contents).unwrap();
27-
blobs.push(blob_transaction_data);
19+
Ok(Self {
20+
client,
21+
bucket_name: config.blobs_temp_storage_bucket_name.clone(),
22+
})
2823
}
2924

30-
let commitments: Vec<Vec<u8>> = blobs
31-
.iter()
32-
.map(|blob| blob.kzg_commitment.clone())
33-
.collect();
34-
println!("{commitments:?}");
35-
36-
Ok(())
25+
pub async fn read(&self, blobhashes: Vec<String>) -> Result<Vec<BlobsQueueS3JsonFile>, Error> {
26+
let mut blobs: Vec<BlobsQueueS3JsonFile> = Vec::new();
27+
28+
for blobhash in blobhashes {
29+
let key = format!("blobs/{}.json", blobhash);
30+
let resp = self
31+
.client
32+
.get_object()
33+
.bucket(self.bucket_name.clone())
34+
.key(key)
35+
.send()
36+
.await?;
37+
38+
let mut body = resp.body.into_async_read();
39+
let mut contents = String::new();
40+
body.read_to_string(&mut contents).await?;
41+
let blob_transaction_data: BlobsQueueS3JsonFile =
42+
serde_json::from_str(&contents).unwrap();
43+
blobs.push(blob_transaction_data);
44+
}
45+
46+
Ok(blobs)
47+
}
3748
}

‎owen/src/blob.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ use c_kzg::{ethereum_kzg_settings, Blob};
44
use log_macros::{format_error, log_info};
55
use serde::{Deserialize, Serialize};
66

7-
#[derive(Deserialize, Serialize)]
7+
#[derive(Deserialize, Serialize, Clone)]
88
pub struct BlobTransactionData {
99
pub kzg_commitment: Vec<u8>,
1010
pub blob_sidecar: BlobTransactionSidecar,

0 commit comments

Comments
 (0)