1use serde::Deserialize;
17use serde_json::Value;
18use std::time::Duration;
19use tungstenite::Message;
20use tungstenite::client::IntoClientRequest;
21use tungstenite::stream::MaybeTlsStream;
22
23use crate::ContentError;
24use crate::config::INDEXER_WS_URL;
25use crate::encode::{bytes_to_hex, hex_to_bytes};
26
27#[derive(Clone, Debug, PartialEq, Eq)]
29pub enum QueryKey {
30 ItemId([u8; 32]),
32 AccountId([u8; 32]),
34 IpfsHash([u8; 32]),
36 ItemRevision { item_id: [u8; 32], revision_id: u32 },
38 Raw { name: String, value: Value },
41}
42
43impl QueryKey {
44 fn to_custom_key(&self) -> Value {
46 match self {
47 QueryKey::ItemId(b) => custom_key("item_id", "bytes32", &Value::from(bytes_to_hex(b))),
48 QueryKey::AccountId(b) => {
49 custom_key("account_id", "bytes32", &Value::from(bytes_to_hex(b)))
50 }
51 QueryKey::IpfsHash(b) => {
52 custom_key("ipfs_hash", "bytes32", &Value::from(bytes_to_hex(b)))
53 }
54 QueryKey::ItemRevision {
55 item_id,
56 revision_id,
57 } => custom_key(
58 "item_id_revision_id",
59 "composite",
60 &Value::Array(vec![
61 custom_scalar("bytes32", &Value::from(bytes_to_hex(item_id))),
62 custom_scalar("u32", &Value::from(*revision_id)),
63 ]),
64 ),
65 QueryKey::Raw { name, value } => {
66 let kind = value["kind"].clone();
68 let inner = value["value"].clone();
69 custom_key(name, kind.as_str().unwrap_or("bytes32"), &inner)
70 }
71 }
72 }
73}
74
75fn wire_key(custom: &Value) -> Value {
80 serde_json::json!({ "type": "Custom", "value": custom })
81}
82
83fn custom_key(name: &str, kind: &str, value: &Value) -> Value {
85 serde_json::json!({ "name": name, "kind": kind, "value": value })
86}
87
88fn custom_scalar(kind: &str, value: &Value) -> Value {
90 serde_json::json!({ "kind": kind, "value": value })
91}
92
93#[derive(Clone, Debug, Deserialize, serde::Serialize)]
95#[serde(rename_all = "camelCase")]
96pub struct DecodedEvent {
97 pub block_number: u32,
98 pub event_index: u32,
99 pub timestamp: u64,
101 pub event: StoredEvent,
103}
104
105#[derive(Clone, Debug, Deserialize, serde::Serialize)]
109#[serde(rename_all = "camelCase")]
110pub struct StoredEvent {
111 pub pallet_name: String,
112 pub event_name: String,
113 pub pallet_index: u8,
114 pub variant_index: u8,
115 pub event_index: u8,
116 pub fields: Value,
118}
119
120impl DecodedEvent {
121 #[must_use]
123 pub fn pallet_name(&self) -> &str {
124 &self.event.pallet_name
125 }
126 #[must_use]
128 pub fn event_name(&self) -> &str {
129 &self.event.event_name
130 }
131 #[must_use]
133 pub fn field(&self, name: &str) -> Option<&Value> {
134 self.event.fields.get(name)
135 }
136 pub fn field_str(&self, name: &str) -> Option<&str> {
138 self.field(name).and_then(Value::as_str)
139 }
140 #[must_use]
143 pub fn field_u64(&self, name: &str) -> Option<u64> {
144 self.field(name)
145 .and_then(|v| v.as_u64().or_else(|| v.as_str()?.parse().ok()))
146 }
147}
148
149#[derive(Clone, Debug, Deserialize, serde::Serialize)]
151#[serde(rename_all = "camelCase")]
152pub struct GetEventsResult {
153 pub events: Vec<DecodedEvent>,
154}
155
156#[derive(Clone, Debug, Deserialize, serde::Serialize)]
158pub struct IndexStatusResult {
159 pub spans: Vec<Span>,
160}
161
162#[derive(Clone, Debug, Deserialize, serde::Serialize)]
164pub struct Span {
165 pub start: u32,
166 pub end: u32,
167}
168
169struct Connection {
171 ws: tungstenite::WebSocket<tungstenite::stream::MaybeTlsStream<std::net::TcpStream>>,
172 next_id: u64,
173}
174
175impl Connection {
176 fn open() -> Result<Self, ContentError> {
177 let request = INDEXER_WS_URL
178 .into_client_request()
179 .map_err(|e| ContentError::Indexer(format!("invalid indexer url: {e}")))?;
180 let (ws, _resp) = tungstenite::connect(request)
181 .map_err(|e| ContentError::Indexer(format!("failed to connect to indexer: {e}")))?;
182 if let MaybeTlsStream::Plain(tcp) = ws.get_ref() {
186 let _ = tcp.set_read_timeout(Some(Duration::from_secs(15)));
187 let _ = tcp.set_write_timeout(Some(Duration::from_secs(10)));
188 }
189 Ok(Self { ws, next_id: 1 })
190 }
191
192 fn request(&mut self, method: &str, params: &Value) -> Result<Value, ContentError> {
193 let id = self.next_id;
194 self.next_id += 1;
195 let req =
196 serde_json::json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params });
197 self.ws
198 .send(Message::Text(req.to_string().into()))
199 .map_err(|e| ContentError::Indexer(format!("write failed: {e}")))?;
200 loop {
201 let msg = self
202 .ws
203 .read()
204 .map_err(|e| ContentError::Indexer(format!("read failed: {e}")))?;
205 match msg {
206 Message::Text(text) => {
207 let v: Value = serde_json::from_str(&text)
208 .map_err(|e| ContentError::Indexer(format!("bad json: {e}")))?;
209 if v.get("id").and_then(Value::as_u64) == Some(id) {
211 return if v.get("error").is_some() {
212 Err(ContentError::Indexer(format!(
213 "indexer error for {method}: {v}"
214 )))
215 } else {
216 Ok(v.get("result").cloned().unwrap_or_default())
217 };
218 }
219 }
220 Message::Ping(data) => {
221 let _ = self.ws.send(Message::Pong(data));
222 }
223 Message::Close(_) => {
224 return Err(ContentError::Indexer("indexer closed connection".into()));
225 }
226 _ => {}
227 }
228 }
229 }
230
231 fn close(&mut self) {
232 let _ = self.ws.close(None);
233 }
234}
235
236pub fn get_events(
245 key: &QueryKey,
246 limit: u16,
247 before: Option<(u32, u32)>,
248) -> Result<Vec<DecodedEvent>, ContentError> {
249 let mut conn = Connection::open()?;
250 let key_json = wire_key(&key.to_custom_key());
251 let params = serde_json::json!({
252 "key": key_json,
253 "limit": limit,
254 "before": before.map(|(b, e)| serde_json::json!({ "blockNumber": b, "eventIndex": e })),
255 });
256 let result = conn.request("acuity_getEvents", ¶ms)?;
257 conn.close();
258 let parsed: GetEventsResult = serde_json::from_value(result)
259 .map_err(|e| ContentError::Indexer(format!("failed to decode get_events result: {e}")))?;
260 Ok(parsed.events)
261}
262
263pub fn index_status() -> Result<IndexStatusResult, ContentError> {
271 let mut conn = Connection::open()?;
272 let result = conn.request("acuity_indexStatus", &serde_json::json!({}))?;
273 conn.close();
274 serde_json::from_value(result)
275 .map_err(|e| ContentError::Indexer(format!("failed to decode index status: {e}")))
276}
277
278pub fn item_id_key(item_id_hex: &str) -> Result<QueryKey, ContentError> {
285 Ok(QueryKey::ItemId(hex_to_bytes(item_id_hex)?))
286}
287
288pub fn item_revision_key(item_id_hex: &str, revision_id: u32) -> Result<QueryKey, ContentError> {
295 Ok(QueryKey::ItemRevision {
296 item_id: hex_to_bytes(item_id_hex)?,
297 revision_id,
298 })
299}
300
301#[cfg(test)]
302mod tests {
303 use super::*;
304
305 #[test]
306 fn wire_key_wraps_custom_key() {
307 let item_id = [0x11u8; 32];
310 let key = wire_key(&QueryKey::ItemId(item_id).to_custom_key());
311 assert_eq!(key["type"], "Custom");
312 assert_eq!(key["value"]["name"], "item_id");
313 assert_eq!(key["value"]["kind"], "bytes32");
314 assert_eq!(key["value"]["value"], bytes_to_hex(&item_id));
315 }
316
317 #[test]
318 fn query_key_wire_shapes() {
319 let item_id = [0x11u8; 32];
321 let key = QueryKey::ItemId(item_id).to_custom_key();
322 assert_eq!(key["name"], "item_id");
323 assert_eq!(key["kind"], "bytes32");
324 assert_eq!(key["value"], bytes_to_hex(&item_id));
325
326 let rev = QueryKey::ItemRevision {
328 item_id,
329 revision_id: 7,
330 }
331 .to_custom_key();
332 assert_eq!(rev["name"], "item_id_revision_id");
333 assert_eq!(rev["kind"], "composite");
334 assert_eq!(rev["value"][0]["kind"], "bytes32");
335 assert_eq!(rev["value"][1]["kind"], "u32");
336 assert_eq!(rev["value"][1]["value"], 7);
337 }
338
339 #[test]
340 fn decoded_event_field_helpers() {
341 let ev = DecodedEvent {
342 block_number: 1,
343 event_index: 2,
344 timestamp: 123,
345 event: StoredEvent {
346 pallet_name: "Content".into(),
347 event_name: "PublishRevision".into(),
348 pallet_index: 7,
349 variant_index: 3,
350 event_index: 2,
351 fields: serde_json::json!({
352 "ipfs_hash": "0xaa",
353 "revision_id": "7",
354 "n": 42,
355 }),
356 },
357 };
358 assert_eq!(ev.pallet_name(), "Content");
359 assert_eq!(ev.event_name(), "PublishRevision");
360 assert_eq!(ev.field_str("ipfs_hash"), Some("0xaa"));
361 assert_eq!(ev.field_u64("revision_id"), Some(7));
363 assert_eq!(ev.field_u64("n"), Some(42));
364 assert_eq!(ev.field_u64("missing"), None);
365 }
366}