Skip to content

Commit

Permalink
[feature] #3468: add query_id to cursor
Browse files Browse the repository at this point in the history
Signed-off-by: Marin Veršić <marin.versic101@gmail.com>
  • Loading branch information
mversic committed Jul 21, 2023
1 parent 3f30478 commit a519cef
Show file tree
Hide file tree
Showing 21 changed files with 307 additions and 195 deletions.
10 changes: 10 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ ursa = "0.3.7"
aead = "0.3.2"

rand = "0.8.5"
warp = { version = "0.3.3", default-features = false }
warp = { version = "0.3.5", default-features = false }
wasmtime = "0.39.1"

tracing = "0.1.37"
Expand Down
10 changes: 7 additions & 3 deletions cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -67,13 +67,17 @@ thiserror = { workspace = true }
displaydoc = { workspace = true }
tokio = { workspace = true, features = ["sync", "time", "rt", "io-util", "rt-multi-thread", "macros", "fs", "signal"] }
warp = { workspace = true, features = ["multipart", "websocket"] }
dashmap = "5.4.0"
serial_test = "0.8.0"
once_cell = { workspace = true }
owo-colors = { workspace = true, features = ["supports-colors"] }
supports-color = { workspace = true }
thread-local-panic-hook = { version = "0.1.0", optional = true }
tempfile = { workspace = true }
dashmap = { workspace = true }

thread-local-panic-hook = { version = "0.1.0", optional = true }
uuid = { version = "1.4.1", features = ["v4"] }

[dev-dependencies]
serial_test = "0.8.0"

[build-dependencies]
iroha_wasm_builder = { workspace = true }
Expand Down
35 changes: 17 additions & 18 deletions cli/src/torii/cursor.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,4 @@
use std::num::NonZeroUsize;

use iroha_data_model::query::ForwardCursor;
use std::num::{NonZeroU64, NonZeroUsize};

use crate::torii::{Error, Result};

Expand All @@ -13,7 +11,7 @@ impl<I: IntoIterator> Batch for I {
Batched {
iter: self.into_iter(),
batch_size,
cursor: ForwardCursor::default(),
cursor: Some(0),
}
}
}
Expand All @@ -24,11 +22,11 @@ impl<I: IntoIterator> Batch for I {
pub struct Batched<I: IntoIterator> {
iter: I::IntoIter,
batch_size: NonZeroUsize,
cursor: ForwardCursor,
cursor: Option<u64>,
}

impl<I: IntoIterator + FromIterator<I::Item>> Batched<I> {
pub(crate) fn next_batch(&mut self, cursor: ForwardCursor) -> Result<(I, ForwardCursor)> {
pub(crate) fn next_batch(&mut self, cursor: Option<u64>) -> Result<(I, Option<NonZeroU64>)> {
if cursor != self.cursor {
return Err(Error::UnknownCursor);
}
Expand All @@ -41,7 +39,7 @@ impl<I: IntoIterator + FromIterator<I::Item>> Batched<I> {
.take(self.batch_size.get())
.collect();

self.cursor.cursor = if let Some(cursor) = self.cursor.cursor {
self.cursor = if let Some(cursor) = self.cursor {
if batch_size >= self.batch_size.get() {
let batch_size = self
.batch_size
Expand All @@ -57,23 +55,24 @@ impl<I: IntoIterator + FromIterator<I::Item>> Batched<I> {
None
}
} else if batch_size >= self.batch_size.get() {
Some(self.batch_size.try_into().expect("usize should fit in u64"))
Some(
self.batch_size
.get()
.try_into()
.expect("usize should fit in u64"),
)
} else {
None
};

Ok((batch, self.cursor))
Ok((
batch,
self.cursor
.map(|cursor| NonZeroU64::new(cursor).expect("Cursor is never 0")),
))
}

pub fn is_depleted(&self) -> bool {
self.cursor.cursor.is_none()
}
}

impl<I: Iterator> Iterator for Batched<I> {
type Item = I::Item;

fn next(&mut self) -> Option<Self::Item> {
self.iter.next()
self.cursor.is_none()
}
}
14 changes: 8 additions & 6 deletions cli/src/torii/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ use iroha_core::{
EventsSender,
};
use iroha_data_model::Value;
use parity_scale_codec::Encode;
use tokio::{sync::Notify, time::sleep};
use utils::*;
use warp::{
Expand All @@ -41,18 +42,19 @@ type LiveQuery = Batched<Vec<Value>>;

#[derive(Default)]
struct LiveQueryStore {
queries: DashMap<Vec<u8>, (LiveQuery, Instant)>,
queries: DashMap<(String, Vec<u8>), (LiveQuery, Instant)>,
}

impl LiveQueryStore {
fn insert(&self, request: Vec<u8>, live_query: LiveQuery) {
self.queries.insert(request, (live_query, Instant::now()));
fn insert<T: Encode>(&self, query_id: String, request: T, live_query: LiveQuery) {
self.queries
.insert((query_id, request.encode()), (live_query, Instant::now()));
}

fn remove(&self, request: &Vec<u8>) -> Option<(Vec<u8>, LiveQuery)> {
fn remove<T: Encode>(&self, query_id: &str, request: &T) -> Option<LiveQuery> {
self.queries
.remove(request)
.map(|(query_id, (query, _))| (query_id, query))
.remove(&(query_id.to_string(), request.encode()))
.map(|(_, (output, _))| output)
}

// TODO: Add notifier channel to enable graceful shutdown
Expand Down
72 changes: 42 additions & 30 deletions cli/src/torii/routing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,10 @@
// FIXME: This can't be fixed, because one trait in `warp` is private.
#![allow(opaque_hidden_inferred_bound)]

use std::{cmp::Ordering, num::NonZeroUsize};
use std::{
cmp::Ordering,
num::{NonZeroU64, NonZeroUsize},
};

use cursor::Batch;
use eyre::WrapErr;
Expand All @@ -28,13 +31,13 @@ use iroha_data_model::{
},
VersionedCommittedBlock,
},
http::{BatchedResponse, VersionedBatchedResponse},
prelude::*,
query::{ForwardCursor, Pagination, Sorting},
};
#[cfg(feature = "telemetry")]
use iroha_telemetry::metrics::Status;
use pagination::Paginate;
use parity_scale_codec::Encode;
use tokio::task;

use super::*;
Expand Down Expand Up @@ -90,39 +93,48 @@ async fn handle_queries(
pagination: Pagination,

cursor: ForwardCursor,
) -> Result<Scale<VersionedQueryResponse>> {
let (query_id, mut live_query) = if cursor.cursor.is_some() {
let query_id = (&request, &sorting, &pagination).encode();
query_store.remove(&query_id).ok_or(Error::UnknownCursor)?
) -> Result<Scale<VersionedBatchedResponse<Value>>> {
// TODO: Remove wsv clone
let mut wsv = sumeragi.wsv_clone();
let valid_request = ValidQueryRequest::validate(request, &mut wsv)?;
let request_id = (&valid_request, &sorting, &pagination);

let (query_id, curr_cursor, mut live_query) = if let Some(query_id) = cursor.query_id {
let live_query = query_store
.remove(&query_id, &request_id)
.ok_or(Error::UnknownCursor)?;

(query_id, cursor.cursor.map(NonZeroU64::get), live_query)
} else {
let mut wsv = sumeragi.wsv_clone();

let valid_request = ValidQueryRequest::validate(request, &mut wsv)?;
let res = valid_request.execute(&wsv).map_err(ValidationFail::from)?;

match res {
LazyValue::Value(result) => {
LazyValue::Value(batch) => {
let cursor = ForwardCursor::default();
let result = QueryResponse { result, cursor };
let result = BatchedResponse { batch, cursor };
return Ok(Scale(result.into()));
}
LazyValue::Iter(iter) => {
let query_id = (&valid_request, &sorting, &pagination).encode();
let query = apply_sorting_and_pagination(iter, &sorting, pagination);
(query_id, query.batched(fetch_size))
let live_query = apply_sorting_and_pagination(iter, &sorting, pagination);
let query_id = uuid::Uuid::new_v4().to_string();

(query_id, Some(0), live_query.batched(fetch_size))
}
}
};

let (batch, next_cursor) = live_query.next_batch(cursor)?;
let (batch, next_cursor) = live_query.next_batch(curr_cursor)?;

if !live_query.is_depleted() {
query_store.insert(query_id, live_query);
query_store.insert(query_id.clone(), request_id, live_query);
}

let query_response = QueryResponse {
result: Value::Vec(batch),
cursor: next_cursor,
let query_response = BatchedResponse {
batch: Value::Vec(batch),
cursor: ForwardCursor {
query_id: Some(query_id),
cursor: next_cursor,
},
};

Ok(Scale(query_response.into()))
Expand Down Expand Up @@ -192,16 +204,16 @@ async fn handle_pending_transactions(
) -> Result<Scale<Vec<VersionedSignedTransaction>>> {
// TODO: Don't clone wsv here
let wsv = sumeragi.wsv_clone();
Ok(Scale(
queue
.all_transactions(&wsv)
.map(Into::into)
.paginate(pagination)
.collect::<Vec<_>>(),
// TODO:
// NOTE: batching is done after collecting the result of pagination
//.batched(fetch_size)
))

let query_response = queue
.all_transactions(&wsv)
.map(Into::into)
.paginate(pagination)
.collect::<Vec<_>>();
// TODO:
//.batched(fetch_size)

Ok(Scale(query_response))
}

#[iroha_futures::telemetry_future]
Expand Down Expand Up @@ -495,7 +507,7 @@ impl Torii {
endpoint3(
handle_pending_transactions,
warp::path(uri::PENDING_TRANSACTIONS)
.and(add_state!(self.queue, self.sumeragi))
.and(add_state!(self.queue, self.sumeragi,))
.and(paginate()),
)
.or(endpoint2(
Expand Down
2 changes: 1 addition & 1 deletion cli/src/torii/utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,4 +66,4 @@ impl<O: Reply, E: Reply> Reply for WarpResult<O, E> {
}
}

iroha_cli_derive::generate_endpoints!(2, 3, 4, 5, 7);
iroha_cli_derive::generate_endpoints!(2, 3, 4, 5, 6, 7);
Loading

0 comments on commit a519cef

Please sign in to comment.