use crate::{IMirror, decode_log, mirror::MirrorQuery};
use alloy::{
contract::Event,
primitives::{Address as AlloyAddress, B256},
providers::{Provider, RootProvider},
rpc::types::eth::{Filter, Log, Topic},
sol_types::{Error, SolEvent},
};
use anyhow::Result;
use ethexe_common::events::{
MirrorEvent, MirrorRequestEvent,
mirror::{
ExecutableBalanceTopUpRequestedEvent, MessageCallFailedEvent, MessageEvent,
MessageQueueingRequestedEvent, OwnedBalanceTopUpRequestedEvent, ReplyCallFailedEvent,
ReplyEvent, ReplyQueueingRequestedEvent, ReplyTransferFailedEvent, StateChangedEvent,
TransferLockedValueToInheritorFailedEvent, ValueClaimFailedEvent, ValueClaimedEvent,
ValueClaimingRequestedEvent,
},
};
use futures::{Stream, StreamExt};
use gear_core::message::ReplyCode;
use gprimitives::ActorId;
use signatures::*;
pub mod signatures {
use super::*;
crate::signatures_consts! {
IMirror;
OWNED_BALANCE_TOP_UP_REQUESTED: OwnedBalanceTopUpRequested,
EXECUTABLE_BALANCE_TOP_UP_REQUESTED: ExecutableBalanceTopUpRequested,
MESSAGE_QUEUEING_REQUESTED: MessageQueueingRequested,
MESSAGE: Message,
MESSAGE_CALL_FAILED: MessageCallFailed,
REPLY_QUEUEING_REQUESTED: ReplyQueueingRequested,
REPLY: Reply,
REPLY_CALL_FAILED: ReplyCallFailed,
STATE_CHANGED: StateChanged,
VALUE_CLAIMED: ValueClaimed,
VALUE_CLAIMING_REQUESTED: ValueClaimingRequested,
TRANSFER_LOCKED_VALUE_TO_INHERITOR_FAILED: TransferLockedValueToInheritorFailed,
REPLY_TRANSFER_FAILED: ReplyTransferFailed,
VALUE_CLAIM_FAILED: ValueClaimFailed,
}
pub const REQUESTS: &[B256] = &[
OWNED_BALANCE_TOP_UP_REQUESTED,
EXECUTABLE_BALANCE_TOP_UP_REQUESTED,
MESSAGE_QUEUEING_REQUESTED,
REPLY_QUEUEING_REQUESTED,
VALUE_CLAIMING_REQUESTED,
];
}
pub fn try_extract_event(log: &Log) -> Result<Option<MirrorEvent>> {
let Some(topic0) = log.topic0().filter(|&v| ALL.contains(v)) else {
return Ok(None);
};
let event = match *topic0 {
OWNED_BALANCE_TOP_UP_REQUESTED => MirrorEvent::OwnedBalanceTopUpRequested(
decode_log::<IMirror::OwnedBalanceTopUpRequested>(log)?.into(),
),
EXECUTABLE_BALANCE_TOP_UP_REQUESTED => MirrorEvent::ExecutableBalanceTopUpRequested(
decode_log::<IMirror::ExecutableBalanceTopUpRequested>(log)?.into(),
),
MESSAGE_QUEUEING_REQUESTED => MirrorEvent::MessageQueueingRequested(
decode_log::<IMirror::MessageQueueingRequested>(log)?.into(),
),
MESSAGE => MirrorEvent::Message(decode_log::<IMirror::Message>(log)?.into()),
MESSAGE_CALL_FAILED => {
MirrorEvent::MessageCallFailed(decode_log::<IMirror::MessageCallFailed>(log)?.into())
}
REPLY_QUEUEING_REQUESTED => MirrorEvent::ReplyQueueingRequested(
decode_log::<IMirror::ReplyQueueingRequested>(log)?.into(),
),
REPLY => MirrorEvent::Reply(decode_log::<IMirror::Reply>(log)?.into()),
REPLY_CALL_FAILED => {
MirrorEvent::ReplyCallFailed(decode_log::<IMirror::ReplyCallFailed>(log)?.into())
}
STATE_CHANGED => {
MirrorEvent::StateChanged(decode_log::<IMirror::StateChanged>(log)?.into())
}
VALUE_CLAIMED => {
MirrorEvent::ValueClaimed(decode_log::<IMirror::ValueClaimed>(log)?.into())
}
VALUE_CLAIMING_REQUESTED => MirrorEvent::ValueClaimingRequested(
decode_log::<IMirror::ValueClaimingRequested>(log)?.into(),
),
TRANSFER_LOCKED_VALUE_TO_INHERITOR_FAILED => {
MirrorEvent::TransferLockedValueToInheritorFailed(
decode_log::<IMirror::TransferLockedValueToInheritorFailed>(log)?.into(),
)
}
REPLY_TRANSFER_FAILED => MirrorEvent::ReplyTransferFailed(
decode_log::<IMirror::ReplyTransferFailed>(log)?.into(),
),
VALUE_CLAIM_FAILED => {
MirrorEvent::ValueClaimFailed(decode_log::<IMirror::ValueClaimFailed>(log)?.into())
}
_ => unreachable!("filtered above"),
};
Ok(Some(event))
}
pub fn try_extract_request_event(log: &Log) -> Result<Option<MirrorRequestEvent>> {
if log.topic0().filter(|&v| REQUESTS.contains(v)).is_none() {
return Ok(None);
}
let request_event = try_extract_event(log)?
.and_then(|v| v.to_request())
.expect("filtered above");
Ok(Some(request_event))
}
pub struct AllEventsBuilder<'a> {
query: &'a MirrorQuery,
}
impl<'a> AllEventsBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self { query }
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<MirrorEvent>> + Unpin + use<>> {
let filter = Filter::new()
.address(*self.query.0.address())
.event_signature(Topic::from_iter([
signatures::OWNED_BALANCE_TOP_UP_REQUESTED,
signatures::EXECUTABLE_BALANCE_TOP_UP_REQUESTED,
signatures::MESSAGE_QUEUEING_REQUESTED,
signatures::MESSAGE,
signatures::MESSAGE_CALL_FAILED,
signatures::REPLY_QUEUEING_REQUESTED,
signatures::REPLY,
signatures::REPLY_CALL_FAILED,
signatures::STATE_CHANGED,
signatures::VALUE_CLAIMED,
signatures::VALUE_CLAIMING_REQUESTED,
signatures::TRANSFER_LOCKED_VALUE_TO_INHERITOR_FAILED,
signatures::REPLY_TRANSFER_FAILED,
signatures::VALUE_CLAIM_FAILED,
]));
Ok(self
.query
.0
.provider()
.subscribe_logs(&filter)
.await?
.into_stream()
.map(|log| try_extract_event(&log).transpose().expect("infallible")))
}
}
pub struct StateChangedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::StateChanged>,
}
impl<'a> StateChangedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.StateChanged_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(StateChangedEvent, Log), Error>> + Unpin + use<>> {
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct MessageQueueingRequestedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::MessageQueueingRequested>,
source: Option<ActorId>,
}
impl<'a> MessageQueueingRequestedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.MessageQueueingRequested_filter(),
source: None,
}
}
pub fn source(mut self, source: ActorId) -> Self {
self.source = Some(source);
self
}
pub async fn subscribe(
self,
) -> Result<
impl Stream<Item = Result<(MessageQueueingRequestedEvent, Log), Error>> + Unpin + use<>,
> {
let mut event = self.event;
if let Some(source) = self.source {
let source: AlloyAddress = source.into();
event = event.topic1(source);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ReplyQueueingRequestedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ReplyQueueingRequested>,
source: Option<ActorId>,
}
impl<'a> ReplyQueueingRequestedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ReplyQueueingRequested_filter(),
source: None,
}
}
pub fn source(mut self, source: ActorId) -> Self {
self.source = Some(source);
self
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ReplyQueueingRequestedEvent, Log), Error>> + Unpin + use<>>
{
let mut event = self.event;
if let Some(source) = self.source {
let source: AlloyAddress = source.into();
event = event.topic1(source);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ValueClaimingRequestedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ValueClaimingRequested>,
source: Option<ActorId>,
}
impl<'a> ValueClaimingRequestedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ValueClaimingRequested_filter(),
source: None,
}
}
pub fn source(mut self, source: ActorId) -> Self {
self.source = Some(source);
self
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ValueClaimingRequestedEvent, Log), Error>> + Unpin + use<>>
{
let mut event = self.event;
if let Some(source) = self.source {
let source: AlloyAddress = source.into();
event = event.topic1(source);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct OwnedBalanceTopUpRequestedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::OwnedBalanceTopUpRequested>,
}
impl<'a> OwnedBalanceTopUpRequestedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.OwnedBalanceTopUpRequested_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<
impl Stream<Item = Result<(OwnedBalanceTopUpRequestedEvent, Log), Error>> + Unpin + use<>,
> {
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ExecutableBalanceTopUpRequestedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ExecutableBalanceTopUpRequested>,
}
impl<'a> ExecutableBalanceTopUpRequestedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ExecutableBalanceTopUpRequested_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<
impl Stream<Item = Result<(ExecutableBalanceTopUpRequestedEvent, Log), Error>> + Unpin + use<>,
> {
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct MessageEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::Message>,
destination: Option<ActorId>,
}
impl<'a> MessageEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.Message_filter(),
destination: None,
}
}
pub fn with_destination(mut self, destination: ActorId) -> Self {
self.destination = Some(destination);
self
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(MessageEvent, Log), Error>> + Unpin + use<>> {
let mut event = self.event;
if let Some(destination) = self.destination {
let destination: AlloyAddress = destination.into();
event = event.topic1(destination);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct MessageCallFailedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::MessageCallFailed>,
destination: Option<ActorId>,
}
impl<'a> MessageCallFailedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.MessageCallFailed_filter(),
destination: None,
}
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(MessageCallFailedEvent, Log), Error>> + Unpin + use<>>
{
let mut event = self.event;
if let Some(destination) = self.destination {
let destination: AlloyAddress = destination.into();
event = event.topic1(destination);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ReplyEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::Reply>,
reply_code: Option<ReplyCode>,
}
impl<'a> ReplyEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.Reply_filter(),
reply_code: None,
}
}
pub fn reply_code(mut self, reply_code: ReplyCode) -> Self {
self.reply_code = Some(reply_code);
self
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ReplyEvent, Log), Error>> + Unpin + use<>> {
let mut event = self.event;
if let Some(reply_code) = self.reply_code {
let mut bytes32 = [0u8; 32];
bytes32[..4].copy_from_slice(&reply_code.to_bytes());
event = event.topic1(bytes32);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ReplyCallFailedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ReplyCallFailed>,
reply_code: Option<ReplyCode>,
}
impl<'a> ReplyCallFailedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ReplyCallFailed_filter(),
reply_code: None,
}
}
pub fn reply_code(mut self, reply_code: ReplyCode) -> Self {
self.reply_code = Some(reply_code);
self
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ReplyCallFailedEvent, Log), Error>> + Unpin + use<>>
{
let mut event = self.event;
if let Some(reply_code) = self.reply_code {
let mut bytes32 = [0u8; 32];
bytes32[..4].copy_from_slice(&reply_code.to_bytes());
event = event.topic1(bytes32);
}
Ok(event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ValueClaimedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ValueClaimed>,
}
impl<'a> ValueClaimedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ValueClaimed_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ValueClaimedEvent, Log), Error>> + Unpin + use<>> {
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct TransferLockedValueToInheritorFailedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::TransferLockedValueToInheritorFailed>,
}
impl<'a> TransferLockedValueToInheritorFailedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.TransferLockedValueToInheritorFailed_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<
impl Stream<Item = Result<(TransferLockedValueToInheritorFailedEvent, Log), Error>>
+ Unpin
+ use<>,
> {
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ReplyTransferFailedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ReplyTransferFailed>,
}
impl<'a> ReplyTransferFailedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ReplyTransferFailed_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ReplyTransferFailedEvent, Log), Error>> + Unpin + use<>>
{
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}
pub struct ValueClaimFailedEventBuilder<'a> {
event: Event<&'a RootProvider, IMirror::ValueClaimFailed>,
}
impl<'a> ValueClaimFailedEventBuilder<'a> {
pub(crate) fn new(query: &'a MirrorQuery) -> Self {
Self {
event: query.0.ValueClaimFailed_filter(),
}
}
pub async fn subscribe(
self,
) -> Result<impl Stream<Item = Result<(ValueClaimFailedEvent, Log), Error>> + Unpin + use<>>
{
Ok(self
.event
.subscribe()
.await?
.into_stream()
.map(|result| result.map(|(event, log)| (event.into(), log))))
}
}