use serde::Deserialize;
use serde_json::Value;
use std::time::Duration;
use tungstenite::Message;
use tungstenite::client::IntoClientRequest;
use tungstenite::stream::MaybeTlsStream;
use crate::ContentError;
use crate::config::INDEXER_WS_URL;
use crate::encode::{bytes_to_hex, hex_to_bytes};
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum QueryKey {
ItemId([u8; 32]),
AccountId([u8; 32]),
IpfsHash([u8; 32]),
ItemRevision { item_id: [u8; 32], revision_id: u32 },
Raw { name: String, value: Value },
}
impl QueryKey {
fn to_custom_key(&self) -> Value {
match self {
QueryKey::ItemId(b) => custom_key("item_id", "bytes32", &Value::from(bytes_to_hex(b))),
QueryKey::AccountId(b) => {
custom_key("account_id", "bytes32", &Value::from(bytes_to_hex(b)))
}
QueryKey::IpfsHash(b) => {
custom_key("ipfs_hash", "bytes32", &Value::from(bytes_to_hex(b)))
}
QueryKey::ItemRevision {
item_id,
revision_id,
} => custom_key(
"item_id_revision_id",
"composite",
&Value::Array(vec![
custom_scalar("bytes32", &Value::from(bytes_to_hex(item_id))),
custom_scalar("u32", &Value::from(*revision_id)),
]),
),
QueryKey::Raw { name, value } => {
let kind = value["kind"].clone();
let inner = value["value"].clone();
custom_key(name, kind.as_str().unwrap_or("bytes32"), &inner)
}
}
}
}
fn wire_key(custom: &Value) -> Value {
serde_json::json!({ "type": "Custom", "value": custom })
}
fn custom_key(name: &str, kind: &str, value: &Value) -> Value {
serde_json::json!({ "name": name, "kind": kind, "value": value })
}
fn custom_scalar(kind: &str, value: &Value) -> Value {
serde_json::json!({ "kind": kind, "value": value })
}
#[derive(Clone, Debug, Deserialize, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct DecodedEvent {
pub block_number: u32,
pub event_index: u32,
pub timestamp: u64,
pub event: StoredEvent,
}
#[derive(Clone, Debug, Deserialize, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct StoredEvent {
pub pallet_name: String,
pub event_name: String,
pub pallet_index: u8,
pub variant_index: u8,
pub event_index: u8,
pub fields: Value,
}
impl DecodedEvent {
#[must_use]
pub fn pallet_name(&self) -> &str {
&self.event.pallet_name
}
#[must_use]
pub fn event_name(&self) -> &str {
&self.event.event_name
}
#[must_use]
pub fn field(&self, name: &str) -> Option<&Value> {
self.event.fields.get(name)
}
pub fn field_str(&self, name: &str) -> Option<&str> {
self.field(name).and_then(Value::as_str)
}
#[must_use]
pub fn field_u64(&self, name: &str) -> Option<u64> {
self.field(name)
.and_then(|v| v.as_u64().or_else(|| v.as_str()?.parse().ok()))
}
}
#[derive(Clone, Debug, Deserialize, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct GetEventsResult {
pub events: Vec<DecodedEvent>,
}
#[derive(Clone, Debug, Deserialize, serde::Serialize)]
pub struct IndexStatusResult {
pub spans: Vec<Span>,
}
#[derive(Clone, Debug, Deserialize, serde::Serialize)]
pub struct Span {
pub start: u32,
pub end: u32,
}
struct Connection {
ws: tungstenite::WebSocket<tungstenite::stream::MaybeTlsStream<std::net::TcpStream>>,
next_id: u64,
}
impl Connection {
fn open() -> Result<Self, ContentError> {
let request = INDEXER_WS_URL
.into_client_request()
.map_err(|e| ContentError::Indexer(format!("invalid indexer url: {e}")))?;
let (ws, _resp) = tungstenite::connect(request)
.map_err(|e| ContentError::Indexer(format!("failed to connect to indexer: {e}")))?;
if let MaybeTlsStream::Plain(tcp) = ws.get_ref() {
let _ = tcp.set_read_timeout(Some(Duration::from_secs(15)));
let _ = tcp.set_write_timeout(Some(Duration::from_secs(10)));
}
Ok(Self { ws, next_id: 1 })
}
fn request(&mut self, method: &str, params: &Value) -> Result<Value, ContentError> {
let id = self.next_id;
self.next_id += 1;
let req =
serde_json::json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params });
self.ws
.send(Message::Text(req.to_string().into()))
.map_err(|e| ContentError::Indexer(format!("write failed: {e}")))?;
loop {
let msg = self
.ws
.read()
.map_err(|e| ContentError::Indexer(format!("read failed: {e}")))?;
match msg {
Message::Text(text) => {
let v: Value = serde_json::from_str(&text)
.map_err(|e| ContentError::Indexer(format!("bad json: {e}")))?;
if v.get("id").and_then(Value::as_u64) == Some(id) {
return if v.get("error").is_some() {
Err(ContentError::Indexer(format!(
"indexer error for {method}: {v}"
)))
} else {
Ok(v.get("result").cloned().unwrap_or_default())
};
}
}
Message::Ping(data) => {
let _ = self.ws.send(Message::Pong(data));
}
Message::Close(_) => {
return Err(ContentError::Indexer("indexer closed connection".into()));
}
_ => {}
}
}
}
fn close(&mut self) {
let _ = self.ws.close(None);
}
}
pub fn get_events(
key: &QueryKey,
limit: u16,
before: Option<(u32, u32)>,
) -> Result<Vec<DecodedEvent>, ContentError> {
let mut conn = Connection::open()?;
let key_json = wire_key(&key.to_custom_key());
let params = serde_json::json!({
"key": key_json,
"limit": limit,
"before": before.map(|(b, e)| serde_json::json!({ "blockNumber": b, "eventIndex": e })),
});
let result = conn.request("acuity_getEvents", ¶ms)?;
conn.close();
let parsed: GetEventsResult = serde_json::from_value(result)
.map_err(|e| ContentError::Indexer(format!("failed to decode get_events result: {e}")))?;
Ok(parsed.events)
}
pub fn index_status() -> Result<IndexStatusResult, ContentError> {
let mut conn = Connection::open()?;
let result = conn.request("acuity_indexStatus", &serde_json::json!({}))?;
conn.close();
serde_json::from_value(result)
.map_err(|e| ContentError::Indexer(format!("failed to decode index status: {e}")))
}
pub fn item_id_key(item_id_hex: &str) -> Result<QueryKey, ContentError> {
Ok(QueryKey::ItemId(hex_to_bytes(item_id_hex)?))
}
pub fn item_revision_key(item_id_hex: &str, revision_id: u32) -> Result<QueryKey, ContentError> {
Ok(QueryKey::ItemRevision {
item_id: hex_to_bytes(item_id_hex)?,
revision_id,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn wire_key_wraps_custom_key() {
let item_id = [0x11u8; 32];
let key = wire_key(&QueryKey::ItemId(item_id).to_custom_key());
assert_eq!(key["type"], "Custom");
assert_eq!(key["value"]["name"], "item_id");
assert_eq!(key["value"]["kind"], "bytes32");
assert_eq!(key["value"]["value"], bytes_to_hex(&item_id));
}
#[test]
fn query_key_wire_shapes() {
let item_id = [0x11u8; 32];
let key = QueryKey::ItemId(item_id).to_custom_key();
assert_eq!(key["name"], "item_id");
assert_eq!(key["kind"], "bytes32");
assert_eq!(key["value"], bytes_to_hex(&item_id));
let rev = QueryKey::ItemRevision {
item_id,
revision_id: 7,
}
.to_custom_key();
assert_eq!(rev["name"], "item_id_revision_id");
assert_eq!(rev["kind"], "composite");
assert_eq!(rev["value"][0]["kind"], "bytes32");
assert_eq!(rev["value"][1]["kind"], "u32");
assert_eq!(rev["value"][1]["value"], 7);
}
#[test]
fn decoded_event_field_helpers() {
let ev = DecodedEvent {
block_number: 1,
event_index: 2,
timestamp: 123,
event: StoredEvent {
pallet_name: "Content".into(),
event_name: "PublishRevision".into(),
pallet_index: 7,
variant_index: 3,
event_index: 2,
fields: serde_json::json!({
"ipfs_hash": "0xaa",
"revision_id": "7",
"n": 42,
}),
},
};
assert_eq!(ev.pallet_name(), "Content");
assert_eq!(ev.event_name(), "PublishRevision");
assert_eq!(ev.field_str("ipfs_hash"), Some("0xaa"));
assert_eq!(ev.field_u64("revision_id"), Some(7));
assert_eq!(ev.field_u64("n"), Some(42));
assert_eq!(ev.field_u64("missing"), None);
}
}