#![recursion_limit = "256"]
#![allow(private_interfaces)]
#![allow(private_bounds)]
#![allow(mismatched_lifetime_syntaxes)]
#![deny(missing_docs)]
#![deny(rustdoc::broken_intra_doc_links)]
#![doc = include_str!("../README.md")]
mod client;
mod connection;
mod connection_info;
#[cfg(feature = "embedded")]
mod embedded;
mod error;
mod graph;
mod graph_schema;
#[cfg(any(feature = "tracing", feature = "metrics"))]
mod observability;
mod parser;
mod response;
mod retry;
mod value;
pub type FalkorResult<T> = Result<T, FalkorDBError>;
pub use client::{blocking::FalkorSyncClient, builder::FalkorClientBuilder, ConnectionStrategy};
pub use connection_info::FalkorConnectionInfo;
pub use error::FalkorDBError;
pub use graph::{
batch::{BatchBuilder, BatchItemResult, BatchQuery, BatchResult},
blocking::SyncGraph,
ops::{ConstraintOpBuilder, CopyGraphBuilder, IndexOpBuilder, WaitOperation, WaitOptions},
query_builder::{ProcedureQueryBuilder, QueryBuilder},
VectorSimilarity,
};
pub use graph_schema::{GraphSchema, SchemaType};
pub use response::{
constraint::{Constraint, ConstraintStatus, ConstraintType},
execution_plan::ExecutionPlan,
index::{FalkorIndex, IndexStatus, IndexType},
lazy_result_set::LazyResultSet,
row::Row,
slowlog_entry::SlowlogEntry,
QueryResult,
};
pub use retry::{Backoff, RetryPolicy, RetryScope};
pub use value::{
config::ConfigValue,
graph_entities::{Edge, EntityType, Node},
path::Path,
point::Point,
temporal::{Date, DateTime, Duration, Seconds, Time},
to_cypher_param, FalkorParams, FalkorValue, FromFalkorValue, IntoFalkorParam, IntoFalkorParams,
RawParam,
};
#[cfg(feature = "tokio")]
pub use response::row_stream::RowStream;
#[cfg(feature = "serde")]
pub use response::typed_result_set::TypedLazyResultSet;
#[cfg(all(feature = "serde", feature = "tokio"))]
pub use response::typed_row_stream::TypedRowStream;
#[cfg(feature = "serde")]
pub use value::{from_falkor_row, from_falkor_value, FalkorValueDeserializer};
#[cfg(feature = "tokio")]
pub use client::asynchronous::FalkorAsyncClient;
#[cfg(feature = "tokio")]
pub use graph::asynchronous::AsyncGraph;
#[cfg(feature = "tokio")]
pub use graph::ops::{AsyncConstraintOpBuilder, AsyncCopyGraphBuilder, AsyncIndexOpBuilder};
#[cfg(feature = "embedded")]
pub use embedded::{EmbeddedConfig, EmbeddedServer};
#[cfg(test)]
pub(crate) mod test_utils {
use super::*;
pub(crate) struct TestSyncGraphHandle {
pub(crate) inner: SyncGraph,
}
impl Drop for TestSyncGraphHandle {
fn drop(&mut self) {
self.inner.delete().ok();
}
}
#[cfg(feature = "tokio")]
pub(crate) struct TestAsyncGraphHandle {
pub(crate) inner: AsyncGraph,
}
#[cfg(feature = "tokio")]
impl Drop for TestAsyncGraphHandle {
fn drop(&mut self) {
tokio::task::block_in_place(|| {
let mut graph_handle =
AsyncGraph::new(self.inner.get_client().clone(), self.inner.graph_name());
tokio::runtime::Handle::current().block_on(async move {
graph_handle.delete().await.ok();
})
})
}
}
pub(crate) fn create_test_client() -> FalkorSyncClient {
FalkorClientBuilder::new()
.build()
.expect("Could not create client")
}
#[cfg(feature = "tokio")]
pub(crate) async fn create_async_test_client() -> FalkorAsyncClient {
FalkorClientBuilder::new_async()
.build()
.await
.expect("Could not create client")
}
pub(crate) const IMDB_FIXTURE_GRAPH: &str = "imdb";
const IMDB_FIXTURE_HINT: &str = "the shared `imdb` test graph is empty or missing. Run \
`resources/populate_graph.py` against the test server first (CI does this in the \
coverage job's \"Populate test graph\" step). This indicates a missing test fixture, \
NOT a code regression.";
fn imdb_actor_count(graph: &mut SyncGraph) -> i64 {
let mut result = graph
.ro_query("MATCH (a:actor) RETURN count(a)")
.execute()
.expect("failed to query the imdb actor count (connection or query regression)");
result
.data
.next()
.expect("imdb actor count query returned no rows")
.expect("imdb actor count row failed to parse")
.try_get_at::<i64>(0)
.expect("imdb actor count column was not an i64")
}
pub(crate) fn imdb_test_client() -> FalkorSyncClient {
let client = create_test_client();
let mut graph = client.select_graph(IMDB_FIXTURE_GRAPH);
assert!(imdb_actor_count(&mut graph) > 0, "{IMDB_FIXTURE_HINT}");
client
}
#[cfg(feature = "tokio")]
pub(crate) async fn imdb_async_test_client() -> FalkorAsyncClient {
use futures::StreamExt;
let client = create_async_test_client().await;
let mut graph = client.select_graph(IMDB_FIXTURE_GRAPH);
let mut result = graph
.ro_query("MATCH (a:actor) RETURN count(a)")
.execute()
.await
.expect("failed to query the imdb actor count (connection or query regression)");
let count = result
.data
.next()
.await
.expect("imdb actor count query returned no rows")
.expect("imdb actor count row failed to parse")
.try_get_at::<i64>(0)
.expect("imdb actor count column was not an i64");
assert!(count > 0, "{IMDB_FIXTURE_HINT}");
client
}
pub(crate) fn open_empty_test_graph(graph_name: &str) -> TestSyncGraphHandle {
let client = create_test_client();
TestSyncGraphHandle {
inner: client.select_graph(graph_name),
}
}
#[cfg(feature = "tokio")]
pub(crate) async fn open_empty_async_test_graph(graph_name: &str) -> TestAsyncGraphHandle {
let client = create_async_test_client().await;
TestAsyncGraphHandle {
inner: client.select_graph(graph_name),
}
}
const RETRY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
const RETRY_INTERVAL: std::time::Duration = std::time::Duration::from_millis(50);
pub(crate) const COPY_RETRY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
pub(crate) fn retry_until<T>(
op: impl FnMut() -> T,
done: impl Fn(&T) -> bool,
) -> T {
retry_until_with_timeout(RETRY_TIMEOUT, op, done)
}
pub(crate) fn retry_until_with_timeout<T>(
timeout: std::time::Duration,
mut op: impl FnMut() -> T,
done: impl Fn(&T) -> bool,
) -> T {
let deadline = std::time::Instant::now() + timeout;
loop {
let value = op();
if done(&value) || std::time::Instant::now() >= deadline {
return value;
}
std::thread::sleep(RETRY_INTERVAL);
}
}
#[cfg(feature = "tokio")]
pub(crate) async fn retry_until_async<T>(
graph: &mut AsyncGraph,
mut op: impl for<'a> FnMut(
&'a mut AsyncGraph,
)
-> std::pin::Pin<Box<dyn std::future::Future<Output = T> + 'a>>,
done: impl Fn(&T) -> bool,
) -> T {
let deadline = std::time::Instant::now() + RETRY_TIMEOUT;
loop {
let value = op(graph).await;
if done(&value) || std::time::Instant::now() >= deadline {
return value;
}
tokio::time::sleep(RETRY_INTERVAL).await;
}
}
#[cfg(feature = "tokio")]
pub(crate) async fn retry_until_async_fn<T, Fut>(
op: impl FnMut() -> Fut,
done: impl Fn(&T) -> bool,
) -> T
where
Fut: std::future::Future<Output = T>,
{
retry_until_async_fn_with_timeout(RETRY_TIMEOUT, op, done).await
}
#[cfg(feature = "tokio")]
pub(crate) async fn retry_until_async_fn_with_timeout<T, Fut>(
timeout: std::time::Duration,
mut op: impl FnMut() -> Fut,
done: impl Fn(&T) -> bool,
) -> T
where
Fut: std::future::Future<Output = T>,
{
let deadline = std::time::Instant::now() + timeout;
loop {
let value = op().await;
if done(&value) || std::time::Instant::now() >= deadline {
return value;
}
tokio::time::sleep(RETRY_INTERVAL).await;
}
}
pub(crate) fn default_on_connection_down<T: Default>(
context: &str,
result: FalkorResult<T>,
) -> T {
match result {
Ok(value) => value,
Err(FalkorDBError::ConnectionDown) => T::default(),
Err(err) => panic!("{context}: {err}"),
}
}
#[cfg(feature = "tokio")]
pub(crate) async fn read_copied_actors(
copy: impl std::future::Future<Output = FalkorResult<AsyncGraph>>
) -> FalkorResult<Vec<FalkorResult<Row>>> {
use futures::StreamExt as _;
let mut graph = copy.await?;
Ok(graph
.query("MATCH (a:actor) RETURN a")
.execute()
.await?
.data
.collect::<Vec<_>>()
.await)
}
}
#[cfg(test)]
mod retry_tests {
use super::test_utils::retry_until;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn retry_until_returns_immediately_when_already_done() {
let calls = AtomicUsize::new(0);
let value = retry_until(|| calls.fetch_add(1, Ordering::Relaxed) + 1, |v| *v == 1);
assert_eq!(value, 1);
assert_eq!(calls.load(Ordering::Relaxed), 1);
}
#[test]
fn retry_until_polls_until_condition_is_met() {
let calls = AtomicUsize::new(0);
let value = retry_until(|| calls.fetch_add(1, Ordering::Relaxed) + 1, |v| *v == 3);
assert_eq!(value, 3);
assert_eq!(calls.load(Ordering::Relaxed), 3);
}
#[cfg(feature = "tokio")]
#[tokio::test(flavor = "multi_thread")]
async fn retry_until_async_fn_polls_until_condition_is_met() {
use super::test_utils::retry_until_async_fn;
let calls = AtomicUsize::new(0);
let value = retry_until_async_fn(
|| async { calls.fetch_add(1, Ordering::Relaxed) + 1 },
|v| *v == 3,
)
.await;
assert_eq!(value, 3);
assert_eq!(calls.load(Ordering::Relaxed), 3);
}
#[test]
fn default_on_connection_down_passes_ok_through() {
use super::test_utils::default_on_connection_down;
let value = default_on_connection_down("ctx", Ok::<_, crate::FalkorDBError>(vec![1, 2, 3]));
assert_eq!(value, vec![1, 2, 3]);
}
#[test]
fn default_on_connection_down_returns_default_on_connection_down() {
use super::test_utils::default_on_connection_down;
let value: Vec<i32> =
default_on_connection_down("ctx", Err(crate::FalkorDBError::ConnectionDown));
assert!(value.is_empty());
}
#[test]
#[should_panic(expected = "could not copy")]
fn default_on_connection_down_panics_on_other_error() {
use super::test_utils::default_on_connection_down;
let _: Vec<i32> =
default_on_connection_down("could not copy", Err(crate::FalkorDBError::ParsingArray));
}
}