use std::any::Any;
use std::fmt;
use std::sync::Arc;
use arrow::datatypes::SchemaRef;
use arrow::record_batch::RecordBatch;
use async_trait::async_trait;
use datafusion::catalog::{Session, TableProvider};
use datafusion::datasource::TableType;
use datafusion::error::{DataFusionError, Result as DFResult};
use datafusion::logical_expr::Expr;
use datafusion::physical_plan::ExecutionPlan;
use futures::stream;
use futures::{StreamExt, TryStreamExt};
use serde_json::{Map, Value};
use super::client::{GraphClient, QueryBounds};
use super::error::GraphError;
#[cfg(test)]
use super::udtf::GraphSourceHealth;
use super::udtf::{GraphScanExec, GraphScanKind, GraphSourceHandle, degraded_reason, mark_healthy};
use super::value::{DeclaredColumn, build_batch, declared_schema};
#[derive(Debug, Clone)]
pub struct ViewContract {
pub name: String,
pub cypher: String,
pub columns: Vec<DeclaredColumn>,
}
pub(crate) struct GraphViewProvider {
handle: Arc<GraphSourceHandle>,
view_name: String,
cypher: String,
columns: Arc<Vec<DeclaredColumn>>,
schema: SchemaRef,
}
impl GraphViewProvider {
pub(crate) fn new(
handle: Arc<GraphSourceHandle>,
view_name: String,
cypher: String,
columns: Vec<DeclaredColumn>,
) -> Self {
let schema = declared_schema(&columns);
Self {
handle,
view_name,
cypher,
columns: Arc::new(columns),
schema,
}
}
}
impl fmt::Debug for GraphViewProvider {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("GraphViewProvider")
.field("view", &self.view_name)
.field("columns", &self.columns.len())
.finish_non_exhaustive()
}
}
#[async_trait]
impl TableProvider for GraphViewProvider {
fn as_any(&self) -> &dyn Any {
self
}
fn schema(&self) -> SchemaRef {
Arc::clone(&self.schema)
}
fn table_type(&self) -> TableType {
TableType::Base
}
async fn scan(
&self,
_state: &dyn Session,
projection: Option<&Vec<usize>>,
_filters: &[Expr],
limit: Option<usize>,
) -> DFResult<Arc<dyn ExecutionPlan>> {
Ok(Arc::new(GraphScanExec::new(
GraphScanKind::View {
handle: Arc::clone(&self.handle),
view_name: self.view_name.clone(),
cypher: self.cypher.clone(),
columns: Arc::clone(&self.columns),
limit,
},
self.schema(),
projection.cloned(),
)?))
}
}
pub(crate) fn view_batches(
handle: Arc<GraphSourceHandle>,
view_name: String,
cypher: String,
columns: Arc<Vec<DeclaredColumn>>,
limit: Option<usize>,
) -> futures::stream::BoxStream<'static, DFResult<RecordBatch>> {
stream::once(async move {
ensure_healthy(&handle, &view_name).await?;
Ok::<_, DataFusionError>(super::udtf::cypher_batches(
handle,
cypher,
Value::Object(Map::new()),
columns,
limit,
))
})
.try_flatten()
.boxed()
}
pub(crate) enum RecoveryOutcome {
AlreadyHealthy,
Recovered { broken_view: Option<GraphError> },
BackoffActive {
registration_error: String,
remaining: std::time::Duration,
},
StillUnavailable {
registration_error: String,
error: GraphError,
},
}
pub(crate) async fn try_recover(handle: &Arc<GraphSourceHandle>) -> RecoveryOutcome {
let _gate = handle.recovery_gate.lock().await;
let Some(registration_error) = degraded_reason(handle) else {
return RecoveryOutcome::AlreadyHealthy;
};
if let Some(remaining) = recovery_backoff_remaining(handle) {
return RecoveryOutcome::BackoffActive {
registration_error,
remaining,
};
}
match revalidate_all_views(handle).await {
Ok(()) => {
mark_healthy(handle);
RecoveryOutcome::Recovered { broken_view: None }
}
Err(e) if !recovery_keeps_degraded(&e) => {
mark_healthy(handle);
RecoveryOutcome::Recovered {
broken_view: Some(e),
}
}
Err(e) => {
arm_recovery_backoff(handle);
RecoveryOutcome::StillUnavailable {
registration_error,
error: e,
}
}
}
}
pub(crate) async fn ensure_healthy(
handle: &Arc<GraphSourceHandle>,
view_name: &str,
) -> DFResult<()> {
if degraded_reason(handle).is_none() {
return Ok(());
}
match try_recover(handle).await {
RecoveryOutcome::AlreadyHealthy => Ok(()),
RecoveryOutcome::Recovered { broken_view: None } => Ok(()),
RecoveryOutcome::Recovered {
broken_view: Some(e),
} => {
tracing::warn!(
error = %e,
"graph source recovered (the backend answered), but a view failed \
re-validation — its own scans will report this"
);
Ok(())
}
RecoveryOutcome::BackoffActive {
registration_error,
remaining,
} => Err(DataFusionError::Execution(format!(
"graph source of view '{view_name}' is registered DEGRADED (registration \
error: {registration_error}); the last recovery attempt found the \
backend still unavailable, and the next retry is {}s away",
remaining.as_secs()
))),
RecoveryOutcome::StillUnavailable {
registration_error,
error,
} => Err(DataFusionError::Execution(format!(
"graph source of view '{view_name}' is registered DEGRADED (registration \
error: {registration_error}) and its first-scan re-validation \
failed: {error}"
))),
}
}
pub(crate) fn recovery_keeps_degraded(e: &GraphError) -> bool {
super::is_availability_artifact(e) || matches!(e, GraphError::RecoveryDeadlineExceeded { .. })
}
fn recovery_backoff_interval(handle: &GraphSourceHandle) -> std::time::Duration {
handle.bounds.timeout.clamp(
std::time::Duration::from_secs(30),
std::time::Duration::from_secs(300),
)
}
pub(crate) fn recovery_backoff_remaining(
handle: &GraphSourceHandle,
) -> Option<std::time::Duration> {
let last = (*handle
.last_failed_recovery
.lock()
.unwrap_or_else(|p| p.into_inner()))?;
recovery_backoff_interval(handle).checked_sub(last.elapsed())
}
pub(crate) fn arm_recovery_backoff(handle: &GraphSourceHandle) {
*handle
.last_failed_recovery
.lock()
.unwrap_or_else(|p| p.into_inner()) = Some(tokio::time::Instant::now());
}
pub(crate) async fn validate_view(
client: &Arc<dyn GraphClient>,
bounds: QueryBounds,
view_name: &str,
cypher: &str,
columns: &[DeclaredColumn],
) -> Result<(), GraphError> {
let fail = |e: GraphError| GraphError::ViewValidationFailed {
view: view_name.to_string(),
source: Box::new(e),
};
let rows = client
.execute(
cypher,
&Value::Object(Map::new()),
columns.len(),
bounds,
Some(1),
)
.await
.map_err(&fail)?
.try_collect::<Vec<_>>()
.await
.map_err(&fail)?;
build_batch(columns, &rows, 0).map_err(&fail)?;
Ok(())
}
pub(crate) async fn revalidate_all_views(
handle: &Arc<GraphSourceHandle>,
) -> Result<(), GraphError> {
let n = handle.view_contracts.len();
if n == 0 {
return Ok(());
}
let limit = handle.validation_limit.max(1);
let waves = n.div_ceil(limit) as u32;
let deadline = handle
.bounds
.timeout
.saturating_mul(2)
.saturating_add(std::time::Duration::from_secs(5))
.saturating_mul(waves);
let contracts = handle.view_contracts.iter().cloned().collect();
tokio::time::timeout(
deadline,
validate_views_concurrently(&handle.client, handle.bounds, contracts, limit),
)
.await
.map_err(|_| GraphError::RecoveryDeadlineExceeded {
seconds: deadline.as_secs(),
})?
}
pub(crate) async fn validate_views_concurrently(
client: &Arc<dyn GraphClient>,
bounds: QueryBounds,
contracts: Vec<ViewContract>,
limit: usize,
) -> Result<(), GraphError> {
stream::iter(contracts.into_iter().map(Ok))
.try_for_each_concurrent(limit.max(1), |contract| {
let client = Arc::clone(client);
async move {
validate_view(
&client,
bounds,
&contract.name,
&contract.cypher,
&contract.columns,
)
.await
}
})
.await
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use arrow::record_batch::RecordBatch;
use datafusion::prelude::SessionContext;
use futures::stream::{self, BoxStream};
use futures::{StreamExt, TryStreamExt};
use super::super::client::{GraphClient, QueryBounds};
use super::super::value::GraphType;
#[derive(Debug)]
struct CountingMock {
rows: Vec<Vec<Value>>,
error: Option<String>,
calls: Arc<AtomicUsize>,
}
#[async_trait]
impl GraphClient for CountingMock {
async fn execute(
&self,
_cypher: &str,
_params: &Value,
_arity: usize,
_bounds: QueryBounds,
limit: Option<usize>,
) -> Result<BoxStream<'static, Result<Vec<Value>, GraphError>>, GraphError> {
self.calls.fetch_add(1, Ordering::Relaxed);
if let Some(message) = &self.error {
return Err(GraphError::Unavailable {
source_name: "kg".to_string(),
reason: message.clone(),
});
}
let mut rows = self.rows.clone();
if let Some(l) = limit {
rows.truncate(l);
}
Ok(stream::iter(rows.into_iter().map(Ok)).boxed())
}
async fn labels(
&self,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<Vec<(String, String)>, GraphError> {
Ok(vec![])
}
}
fn handle_with(
health: GraphSourceHealth,
rows: Vec<Vec<Value>>,
error: Option<String>,
) -> (Arc<GraphSourceHandle>, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
let handle = Arc::new(GraphSourceHandle::new(
Arc::new(CountingMock {
rows,
error,
calls: Arc::clone(&calls),
}),
QueryBounds {
timeout: Duration::from_secs(5),
max_rows: 100,
},
health,
Arc::new(vec![ViewContract {
name: "user_posts".to_string(),
cypher: "MATCH (u:User) RETURN u.name".to_string(),
columns: columns(),
}]),
4,
));
(handle, calls)
}
fn columns() -> Vec<DeclaredColumn> {
vec![DeclaredColumn {
name: "name".to_string(),
ty: GraphType::String,
nullable: true,
}]
}
fn provider(handle: Arc<GraphSourceHandle>) -> GraphViewProvider {
GraphViewProvider::new(
handle,
"user_posts".to_string(),
"MATCH (u:User) RETURN u.name".to_string(),
columns(),
)
}
fn is_healthy(handle: &GraphSourceHandle) -> bool {
handle
.health
.read()
.unwrap_or_else(|p| p.into_inner())
.is_healthy()
}
async fn scan_rows(provider: GraphViewProvider) -> DFResult<Vec<RecordBatch>> {
let ctx = SessionContext::new();
ctx.register_table("user_posts", Arc::new(provider))?;
ctx.sql("SELECT name FROM user_posts")
.await?
.collect()
.await
}
#[derive(Debug)]
struct PickyMock {
fail_marker: &'static str,
}
#[async_trait]
impl GraphClient for PickyMock {
async fn execute(
&self,
cypher: &str,
_params: &Value,
_arity: usize,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<BoxStream<'static, Result<Vec<Value>, GraphError>>, GraphError> {
if cypher.contains(self.fail_marker) {
return Err(GraphError::backend(
"kg",
"42804",
"return row and column definition list do not match",
));
}
Ok(stream::iter(vec![vec![serde_json::json!("ada")]].into_iter().map(Ok)).boxed())
}
async fn labels(
&self,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<Vec<(String, String)>, GraphError> {
Ok(vec![])
}
}
#[tokio::test]
async fn physical_planning_never_touches_the_backend() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![vec![serde_json::json!("ada")]],
None,
);
let ctx = SessionContext::new();
ctx.register_table("user_posts", Arc::new(provider(Arc::clone(&handle))))
.expect("register");
let df = ctx.sql("SELECT name FROM user_posts").await.expect("plans");
let plan = df
.create_physical_plan()
.await
.expect("physical planning performs no network I/O");
assert_eq!(
calls.load(Ordering::Relaxed),
0,
"no backend call before execution"
);
let stream = plan.execute(0, ctx.task_ctx()).expect("execute");
let batches: Vec<RecordBatch> = stream.try_collect().await.expect("backend answers");
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert_eq!(calls.load(Ordering::Relaxed), 2, "validate + scan");
assert!(is_healthy(&handle), "recovered on execution");
}
#[tokio::test]
async fn a_degraded_backend_failure_surfaces_at_execution_not_planning() {
let (handle, _calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![],
Some("still refused".to_string()),
);
let ctx = SessionContext::new();
ctx.register_table("user_posts", Arc::new(provider(Arc::clone(&handle))))
.expect("register");
let df = ctx.sql("SELECT name FROM user_posts").await.expect("plans");
let plan = df
.create_physical_plan()
.await
.expect("physical planning succeeds against a down backend");
let err = plan
.execute(0, ctx.task_ctx())
.expect("execute")
.try_collect::<Vec<RecordBatch>>()
.await
.expect_err("the failure arrives at execution");
assert!(err.to_string().contains("DEGRADED"), "{err}");
}
#[tokio::test]
async fn a_broken_sibling_no_longer_blocks_recovery_when_the_backend_answers() {
let handle = Arc::new(GraphSourceHandle::new(
Arc::new(PickyMock {
fail_marker: "Broken",
}),
QueryBounds {
timeout: Duration::from_secs(5),
max_rows: 100,
},
GraphSourceHealth::Degraded("connection refused".to_string()),
Arc::new(vec![
ViewContract {
name: "good_view".to_string(),
cypher: "MATCH (g:Good) RETURN g.name".to_string(),
columns: columns(),
},
ViewContract {
name: "broken_view".to_string(),
cypher: "MATCH (b:Broken) RETURN b.name".to_string(),
columns: columns(),
},
]),
4,
));
let provider = GraphViewProvider::new(
handle.clone(),
"good_view".to_string(),
"MATCH (g:Good) RETURN g.name".to_string(),
columns(),
);
let batches = scan_rows(provider)
.await
.expect("the good view's correct answer is not discarded");
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert!(
is_healthy(&handle),
"the backend answered — the source recovers"
);
let broken = GraphViewProvider::new(
Arc::clone(&handle),
"broken_view".to_string(),
"MATCH (b:Broken) RETURN b.name".to_string(),
columns(),
);
let err = scan_rows(broken)
.await
.expect_err("the broken view still fails, on its own scans");
assert!(
!err.to_string().contains("DEGRADED"),
"a typed contract error, not a degraded wrap: {err}"
);
}
#[tokio::test]
async fn a_degraded_view_validates_on_first_scan_and_flips_healthy() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![vec![serde_json::json!("ada")]],
None,
);
let batches = scan_rows(provider(Arc::clone(&handle)))
.await
.expect("the backend answers now");
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert!(is_healthy(&handle), "a successful retry flips Healthy");
assert_eq!(
calls.load(Ordering::Relaxed),
2,
"validation execute + scan execute"
);
let batches = scan_rows(provider(Arc::clone(&handle)))
.await
.expect("healthy scans");
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert_eq!(calls.load(Ordering::Relaxed), 3, "no second validation");
}
#[tokio::test]
async fn a_degraded_view_fails_loudly_when_the_backend_is_still_gone() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![],
Some("still refused".to_string()),
);
let err = scan_rows(provider(Arc::clone(&handle)))
.await
.expect_err("the retry fails loudly");
let msg = err.to_string();
assert!(msg.contains("user_posts"), "the view is named: {msg}");
assert!(msg.contains("DEGRADED"), "{msg}");
assert!(msg.contains("connection refused"), "{msg}");
assert!(!is_healthy(&handle), "a failed retry stays Degraded");
assert_eq!(calls.load(Ordering::Relaxed), 1, "validation only");
}
#[tokio::test(start_paused = true)]
async fn failed_recovery_backs_off_instead_of_repaying_the_validation() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![],
Some("still refused".to_string()),
);
let err = scan_rows(provider(Arc::clone(&handle)))
.await
.expect_err("first attempt fails");
assert!(err.to_string().contains("DEGRADED"), "{err}");
assert_eq!(calls.load(Ordering::Relaxed), 1, "one validation probe");
let err = scan_rows(provider(Arc::clone(&handle)))
.await
.expect_err("inside the window: fast, cached failure");
let msg = err.to_string();
assert!(msg.contains("next retry"), "names the backoff: {msg}");
assert!(msg.contains("connection refused"), "keeps the cause: {msg}");
assert_eq!(
calls.load(Ordering::Relaxed),
1,
"no re-validation inside the backoff window"
);
tokio::time::advance(std::time::Duration::from_secs(31)).await;
let _ = scan_rows(provider(Arc::clone(&handle)))
.await
.expect_err("still down");
assert_eq!(
calls.load(Ordering::Relaxed),
2,
"the expired window re-arms a real attempt"
);
}
#[tokio::test]
async fn a_healthy_view_scans_without_revalidating() {
let (handle, calls) = handle_with(
GraphSourceHealth::Healthy,
vec![vec![serde_json::json!("ada")]],
None,
);
let batches = scan_rows(provider(Arc::clone(&handle)))
.await
.expect("healthy scans");
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert_eq!(
calls.load(Ordering::Relaxed),
1,
"no validation round-trip for a healthy source"
);
}
#[tokio::test]
async fn validation_catches_type_and_not_null_violations() {
let (handle, _) = handle_with(
GraphSourceHealth::Healthy,
vec![vec![serde_json::json!(7)]],
None,
);
let err = validate_view(
&handle.client,
handle.bounds,
"user_posts",
"MATCH (n) RETURN n",
&columns(),
)
.await
.unwrap_err();
let msg = err.to_string();
assert!(msg.contains("user_posts"), "{msg}");
assert!(msg.contains("declared 'string'"), "{msg}");
let (handle, _) = handle_with(GraphSourceHealth::Healthy, vec![vec![Value::Null]], None);
let strict = vec![DeclaredColumn {
name: "name".to_string(),
ty: GraphType::String,
nullable: false,
}];
let err = validate_view(
&handle.client,
handle.bounds,
"user_posts",
"MATCH (n) RETURN n",
&strict,
)
.await
.unwrap_err();
assert!(err.to_string().contains("nullable: false"), "{err}");
let (handle, _) = handle_with(GraphSourceHealth::Healthy, vec![], None);
validate_view(
&handle.client,
handle.bounds,
"user_posts",
"MATCH (n) RETURN n",
&columns(),
)
.await
.expect("empty is valid");
}
#[tokio::test]
async fn concurrent_first_scans_share_one_recovery_flight() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("connection refused".to_string()),
vec![vec![serde_json::json!("ada")]],
None,
);
let scans = (0..4).map(|_| {
let handle = Arc::clone(&handle);
async move { scan_rows(provider(Arc::clone(&handle))).await }
});
let results = futures::future::join_all(scans).await;
for result in results {
result.expect("every concurrent scan succeeds");
}
assert!(is_healthy(&handle), "the winner flipped the source");
assert_eq!(
calls.load(Ordering::Relaxed),
5,
"one shared validation flight + four scans — never 4 + 4"
);
}
#[derive(Debug)]
struct GaugeMock {
current: Arc<AtomicUsize>,
peak: Arc<AtomicUsize>,
calls: Arc<AtomicUsize>,
}
#[async_trait]
impl GraphClient for GaugeMock {
async fn execute(
&self,
_cypher: &str,
_params: &Value,
_arity: usize,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<BoxStream<'static, Result<Vec<Value>, GraphError>>, GraphError> {
let now = self.current.fetch_add(1, Ordering::SeqCst) + 1;
self.peak.fetch_max(now, Ordering::SeqCst);
self.calls.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(25)).await;
self.current.fetch_sub(1, Ordering::SeqCst);
Ok(stream::iter(vec![Ok(vec![serde_json::json!("x")])]).boxed())
}
async fn labels(
&self,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<Vec<(String, String)>, GraphError> {
Ok(vec![])
}
}
#[tokio::test]
async fn registration_validation_concurrency_is_bounded_by_the_limit() {
let current = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let calls = Arc::new(AtomicUsize::new(0));
let handle = Arc::new(GraphSourceHandle::new(
Arc::new(GaugeMock {
current: Arc::clone(¤t),
peak: Arc::clone(&peak),
calls: Arc::clone(&calls),
}),
QueryBounds {
timeout: Duration::from_secs(5),
max_rows: 100,
},
GraphSourceHealth::Healthy,
Arc::new(vec![]),
4,
));
let contracts: Vec<ViewContract> = (0..8)
.map(|i| ViewContract {
name: format!("v{i}"),
cypher: "MATCH (n) RETURN n.x".to_string(),
columns: columns(),
})
.collect();
validate_views_concurrently(&handle.client, handle.bounds, contracts, 2)
.await
.expect("all views validate");
assert_eq!(calls.load(Ordering::SeqCst), 8, "every view was proven");
assert!(
peak.load(Ordering::SeqCst) <= 2,
"at most `limit` probes in flight, got {}",
peak.load(Ordering::SeqCst)
);
}
#[derive(Debug)]
struct WedgedClient;
#[async_trait]
impl GraphClient for WedgedClient {
async fn execute(
&self,
_cypher: &str,
_params: &Value,
_arity: usize,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<BoxStream<'static, Result<Vec<Value>, GraphError>>, GraphError> {
std::future::pending().await
}
async fn labels(
&self,
_bounds: QueryBounds,
_limit: Option<usize>,
) -> Result<Vec<(String, String)>, GraphError> {
std::future::pending().await
}
}
#[tokio::test]
async fn diagnostics_carry_the_view_identity_never_the_cypher() {
let (handle, _) = handle_with(GraphSourceHealth::Healthy, vec![], None);
let provider = provider(handle);
assert_eq!(provider.table_type(), TableType::Base);
let dbg = format!("{provider:?}");
assert!(dbg.contains("user_posts"), "{dbg}");
assert!(
!dbg.contains("MATCH"),
"the Cypher text never appears in diagnostics: {dbg}"
);
}
#[tokio::test]
async fn a_gate_loser_wakes_to_a_healthy_source_and_skips_revalidation() {
let (handle, calls) = handle_with(
GraphSourceHealth::Degraded("down at startup".to_string()),
vec![],
None,
);
let gate = handle.recovery_gate.lock().await;
let loser = tokio::spawn({
let handle = Arc::clone(&handle);
async move { ensure_healthy(&handle, "user_posts").await }
});
tokio::time::sleep(Duration::from_millis(20)).await;
*handle.health.write().unwrap_or_else(|p| p.into_inner()) = GraphSourceHealth::Healthy;
drop(gate);
loser
.await
.expect("no panic")
.expect("the loser proceeds without error");
assert_eq!(
calls.load(Ordering::Relaxed),
0,
"the loser re-proved nothing"
);
}
#[tokio::test(start_paused = true)]
async fn a_wedged_revalidation_hits_the_backstop_deadline() {
let handle = Arc::new(GraphSourceHandle::new(
Arc::new(WedgedClient),
QueryBounds {
timeout: Duration::from_secs(5),
max_rows: 100,
},
GraphSourceHealth::Degraded("down at startup".to_string()),
Arc::new(vec![ViewContract {
name: "user_posts".to_string(),
cypher: "MATCH (u:User) RETURN u.name".to_string(),
columns: columns(),
}]),
4,
));
let err = revalidate_all_views(&handle)
.await
.expect_err("the wedge cannot pass validation");
assert!(
matches!(err, GraphError::RecoveryDeadlineExceeded { seconds: 15 }),
"{err}"
);
}
}