Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,12 @@ pub(crate) mod tests {
*self.active_reads.lock().unwrap() += 1;
Ok(self.still_active_after_cancel)
}
async fn agreement_ended_by_indexer(
&self,
_agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError> {
Ok(false)
}
}

fn manager_conf(collector: Address) -> IndexingAgreementConfig {
Expand Down
14 changes: 14 additions & 0 deletions bin/dipper-service/src/chain_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,13 @@ pub trait ChainClient {
agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError>;

/// Read whether the indexer ended the agreement on-chain, which `getAgreementDetails`
/// reports with its BY_PROVIDER flag. False for one still live or ended by dipper.
async fn agreement_ended_by_indexer(
&self,
agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError>;

/// Read the latest block's unix timestamp from the chain. Lets agreement
/// deadlines be stamped from live chain time when the chain-clock bypass is
/// on, instead of a cached listener timestamp that can lag a fast chain.
Expand Down Expand Up @@ -182,6 +189,13 @@ impl<T: ChainClient + Send + Sync + ?Sized> ChainClient for Arc<T> {
(**self).agreement_still_active(agreement_id).await
}

async fn agreement_ended_by_indexer(
&self,
agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError> {
(**self).agreement_ended_by_indexer(agreement_id).await
}

async fn latest_block_timestamp(&self) -> Result<u64, ChainClientError> {
(**self).latest_block_timestamp().await
}
Expand Down
55 changes: 24 additions & 31 deletions bin/dipper-service/src/chain_client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,12 +58,13 @@ const RECEIPT_POLL_INTERVAL: Duration = Duration::from_millis(500);
const VERSION_CURRENT: u64 = 0;

/// `AgreementDetails.state` flags from `IAgreementCollector.sol` (REGISTERED=1,
/// ACCEPTED=2, NOTICE_GIVEN=4). `getAgreementDetails` keeps ACCEPTED set on a
/// canceled agreement and ORs in NOTICE_GIVEN, so a cancel must clear it, not
/// just lack it. REGISTERED without ACCEPTED is an offer still waiting.
/// ACCEPTED=2, NOTICE_GIVEN=4, BY_PROVIDER=32). `getAgreementDetails` keeps ACCEPTED set on
/// a canceled agreement and ORs in NOTICE_GIVEN, so a cancel must clear it, not just lack
/// it. REGISTERED without ACCEPTED is an offer still waiting. BY_PROVIDER: the indexer cancelled.
const STATE_REGISTERED: u16 = 1;
const STATE_ACCEPTED: u16 = 2;
const STATE_NOTICE_GIVEN: u16 = 4;
const STATE_BY_PROVIDER: u16 = 32;

/// Live iff the terms are accepted and no cancellation notice exists, or an
/// offer is still stored for the indexer to accept. A cancel sets NOTICE_GIVEN
Expand Down Expand Up @@ -602,6 +603,19 @@ impl AlloyChainClient {
}
}

/// The agreement's `AgreementDetails.state` flags for its current terms.
async fn agreement_state(&self, agreement_id: &[u8; 16]) -> Result<u16, ChainClientError> {
let call = IRecurringCollector::getAgreementDetailsCall {
agreementId: FixedBytes::<16>::from_slice(agreement_id),
index: thegraph_core::alloy::primitives::U256::from(VERSION_CURRENT),
};
let collector = self.inner.recurring_collector_address;
Ok(self
.view(collector, call, "get_agreement_details")
.await?
.state)
}

/// Run a read-only contract call and decode its return value.
async fn view<C: SolCall>(
&self,
Expand Down Expand Up @@ -789,35 +803,14 @@ impl ChainClient for AlloyChainClient {
&self,
agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError> {
let calldata = IRecurringCollector::getAgreementDetailsCall {
agreementId: FixedBytes::<16>::from_slice(agreement_id),
index: thegraph_core::alloy::primitives::U256::from(VERSION_CURRENT),
}
.abi_encode();

let collector = self.inner.recurring_collector_address;
let output = self
.inner
.rpc_pool
.execute("get_agreement_details", |provider| {
let calldata = calldata.clone();
async move {
let tx = TransactionRequest::default()
.to(collector)
.input(calldata.into());
provider.call(tx).await
}
})
.await?;

let details = IRecurringCollector::getAgreementDetailsCall::abi_decode_returns(&output)
.map_err(|err| {
ChainClientError::RpcError(anyhow::anyhow!(
"undecodable getAgreementDetails from {collector}: {err}"
))
})?;
Ok(still_live(self.agreement_state(agreement_id).await?))
}

Ok(still_live(details.state))
async fn agreement_ended_by_indexer(
&self,
agreement_id: &[u8; 16],
) -> Result<bool, ChainClientError> {
Ok(self.agreement_state(agreement_id).await? & STATE_BY_PROVIDER != 0)
}

async fn reconcile_provider(
Expand Down
Loading
Loading