pub mod policy_proto {
tonic::include_proto!("exchange.cex.policy");
}
use policy_proto::policy_internal_service_client::PolicyInternalServiceClient;
#[derive(Debug, Clone)]
pub struct GovernanceProposal {
pub id: String,
pub policy_id: String,
pub action: Option<String>,
pub status: String,
pub payload: Option<serde_json::Value>,
}
#[derive(Clone)]
pub struct PolicyGovernanceClient {
target: Option<String>,
}
impl PolicyGovernanceClient {
pub fn new(url: Option<String>) -> Self {
let target = url
.map(|url| url.trim().trim_end_matches('/').to_string())
.filter(|url| !url.is_empty())
.map(|url| {
if url.starts_with("http://") || url.starts_with("https://") {
url
} else {
format!("http://{url}")
}
});
Self { target }
}
pub fn configured(&self) -> bool {
self.target.is_some()
}
async fn client(
&self,
) -> anyhow::Result<PolicyInternalServiceClient<tonic::transport::Channel>> {
let target = self
.target
.clone()
.ok_or_else(|| anyhow::anyhow!("cex-policy gRPC target not configured"))?;
let channel = tonic::transport::Endpoint::from_shared(target)?
.timeout(std::time::Duration::from_secs(10))
.connect_timeout(std::time::Duration::from_secs(3))
.connect_lazy();
Ok(PolicyInternalServiceClient::new(channel))
}
pub async fn submit_proposal(
&self,
policy_id: &str,
proposer: &str,
payload: &serde_json::Value,
) -> anyhow::Result<String> {
let response = self
.client()
.await?
.create_proposal(policy_proto::CreateProposalRequest {
policy_id: policy_id.to_string(),
proposer: proposer.to_string(),
payload_json: serde_json::to_string(payload)?,
action: String::new(),
min_threshold: None,
context_json: String::new(),
})
.await?
.into_inner();
if response.id.is_empty() {
anyhow::bail!("cex-policy proposal submit returned no id");
}
Ok(response.id)
}
pub async fn submit_proposal_for_action(
&self,
action: &str,
proposer: &str,
payload: &serde_json::Value,
) -> anyhow::Result<Option<String>> {
self.submit_for_action(action, proposer, payload, None, None)
.await
}
pub async fn submit_proposal_for_action_with_floor(
&self,
action: &str,
proposer: &str,
payload: &serde_json::Value,
min_threshold: Option<u8>,
) -> anyhow::Result<Option<String>> {
self.submit_for_action(action, proposer, payload, min_threshold, None)
.await
}
pub async fn submit_proposal_for_action_with_context(
&self,
action: &str,
proposer: &str,
payload: &serde_json::Value,
context: Option<&serde_json::Value>,
) -> anyhow::Result<Option<String>> {
let context_json = context.map(|c| c.to_string());
self.submit_for_action(action, proposer, payload, None, context_json)
.await
}
async fn submit_for_action(
&self,
action: &str,
proposer: &str,
payload: &serde_json::Value,
min_threshold: Option<u8>,
context_json: Option<String>,
) -> anyhow::Result<Option<String>> {
let result = self
.client()
.await?
.create_proposal(policy_proto::CreateProposalRequest {
policy_id: String::new(),
proposer: proposer.to_string(),
payload_json: serde_json::to_string(payload)?,
action: action.to_string(),
min_threshold: min_threshold.map(u32::from),
context_json: context_json.unwrap_or_default(),
})
.await;
match result {
Ok(response) => {
let id = response.into_inner().id;
if id.is_empty() {
anyhow::bail!("cex-policy proposal submit returned no id");
}
Ok(Some(id))
}
Err(status) if status.code() == tonic::Code::NotFound => Ok(None),
Err(status) => Err(status.into()),
}
}
pub async fn fetch_proposal(
&self,
proposal_id: &str,
) -> anyhow::Result<Option<GovernanceProposal>> {
let result = self
.client()
.await?
.get_proposal(policy_proto::GetProposalRequest { id: proposal_id.to_string() })
.await;
match result {
Ok(response) => {
let proposal = response.into_inner();
let payload = serde_json::from_str(&proposal.payload_json).ok();
Ok(Some(GovernanceProposal {
id: proposal.id,
policy_id: proposal.policy_id,
action: proposal.action.filter(|action| !action.is_empty()),
status: proposal.status,
payload,
}))
}
Err(status) => {
#[cfg(feature = "observability")]
tracing::warn!(proposal_id, error = %status, "policy internal proposal fetch failed");
let _ = status;
Ok(None)
}
}
}
pub async fn fetch_bootstrap_state(&self) -> anyhow::Result<Option<String>> {
let result = self
.client()
.await?
.get_bootstrap_state(policy_proto::GetBootstrapStateRequest {})
.await;
match result {
Ok(response) => Ok(Some(response.into_inner().state)),
Err(status) => {
#[cfg(feature = "observability")]
tracing::warn!(error = %status, "policy internal bootstrap-state fetch failed");
let _ = status;
Ok(None)
}
}
}
}