mod common;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use axum::{
body::Body,
http::{Request, StatusCode},
};
use common::create_test_app_with_observer;
use serde_json::{json, Value};
use tempfile::TempDir;
use tower::ServiceExt;
use velesdb_core::collection::CollectionType;
use velesdb_core::DatabaseObserver;
const COLLECTION: &str = "observed";
const DIM: usize = 4;
const QUERY: [f32; DIM] = [1.0, 0.5, 0.25, 0.1];
#[derive(Default)]
struct CountingObserver {
created: AtomicUsize,
deleted: AtomicUsize,
upsert: AtomicUsize,
query: AtomicUsize,
}
impl DatabaseObserver for CountingObserver {
fn on_collection_created(&self, _name: &str, _kind: &CollectionType) {
self.created.fetch_add(1, Ordering::SeqCst);
}
fn on_collection_deleted(&self, _name: &str) {
self.deleted.fetch_add(1, Ordering::SeqCst);
}
fn on_upsert(&self, _collection: &str, _point_count: usize) {
self.upsert.fetch_add(1, Ordering::SeqCst);
}
fn on_query(&self, _collection: &str, _duration_us: u64) {
self.query.fetch_add(1, Ordering::SeqCst);
}
}
async fn post(app: &axum::Router, uri: &str, body: Value) -> StatusCode {
app.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.header("Content-Type", "application/json")
.body(Body::from(body.to_string()))
.expect("test: build POST request"),
)
.await
.expect("test: POST request")
.status()
}
async fn get(app: &axum::Router, uri: &str) -> StatusCode {
app.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(uri)
.body(Body::empty())
.expect("test: build GET request"),
)
.await
.expect("test: GET request")
.status()
}
async fn delete(app: &axum::Router, uri: &str) -> StatusCode {
app.clone()
.oneshot(
Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.expect("test: build DELETE request"),
)
.await
.expect("test: DELETE request")
.status()
}
#[tokio::test]
async fn observer_receives_full_lifecycle() {
let dir = TempDir::new().expect("test: dir");
let observer = Arc::new(CountingObserver::default());
let app = create_test_app_with_observer(&dir, observer.clone());
let status = post(
&app,
"/collections",
json!({"name": COLLECTION, "dimension": DIM, "metric": "cosine"}),
)
.await;
assert_eq!(status, StatusCode::CREATED, "create");
assert_eq!(
observer.created.load(Ordering::SeqCst),
1,
"on_collection_created"
);
let status = post(
&app,
&format!("/collections/{COLLECTION}/points"),
json!({"points": [{"id": 1, "vector": QUERY, "payload": {"k": "v"}}]}),
)
.await;
assert_eq!(status, StatusCode::OK, "upsert");
assert!(observer.upsert.load(Ordering::SeqCst) >= 1, "on_upsert");
let status = post(
&app,
&format!("/collections/{COLLECTION}/search"),
json!({"vector": QUERY, "top_k": 1}),
)
.await;
assert_eq!(status, StatusCode::OK, "search");
assert!(observer.query.load(Ordering::SeqCst) >= 1, "on_query");
let status = delete(&app, &format!("/collections/{COLLECTION}")).await;
assert_eq!(status, StatusCode::OK, "delete");
assert_eq!(
observer.deleted.load(Ordering::SeqCst),
1,
"on_collection_deleted"
);
}
#[tokio::test]
async fn query_endpoint_fires_on_query_exactly_once_per_dispatch_branch() {
let dir = TempDir::new().expect("test: dir");
let observer = Arc::new(CountingObserver::default());
let app = create_test_app_with_observer(&dir, observer.clone());
let status = post(
&app,
"/collections",
json!({"name": COLLECTION, "dimension": DIM, "metric": "cosine"}),
)
.await;
assert_eq!(status, StatusCode::CREATED, "create");
let status = post(
&app,
"/query",
json!({
"query": format!("SELECT * FROM {COLLECTION} WHERE vector NEAR $v LIMIT 1"),
"params": {"v": QUERY}
}),
)
.await;
assert_eq!(status, StatusCode::OK, "select via /query");
assert_eq!(
observer.query.load(Ordering::SeqCst),
1,
"on_query must fire exactly once for a SELECT dispatched via /query"
);
let status = post(&app, "/query", json!({"query": "SHOW COLLECTIONS"})).await;
assert_eq!(status, StatusCode::OK, "show collections via /query");
assert_eq!(
observer.query.load(Ordering::SeqCst),
2,
"on_query must fire exactly once more for an introspection query dispatched via /query"
);
}
struct DenyingReadObserver;
impl DatabaseObserver for DenyingReadObserver {
fn on_query_request(
&self,
_ctx: &velesdb_core::observer::QueryAccessContext,
) -> velesdb_core::Result<velesdb_core::observer::AccessDecision> {
Ok(velesdb_core::observer::AccessDecision::Deny(
velesdb_core::Error::Query("read denied by governance policy".to_string()),
))
}
}
#[tokio::test]
async fn read_gate_denies_rest_search_end_to_end() {
let dir = TempDir::new().expect("test: dir");
let app = create_test_app_with_observer(&dir, Arc::new(DenyingReadObserver));
let status = post(
&app,
"/collections",
json!({"name": COLLECTION, "dimension": DIM, "metric": "cosine"}),
)
.await;
assert_eq!(status, StatusCode::CREATED, "create must not be read-gated");
let status = post(
&app,
&format!("/collections/{COLLECTION}/points"),
json!({"points": [{"id": 1, "vector": QUERY, "payload": {"k": "v"}}]}),
)
.await;
assert_eq!(status, StatusCode::OK, "upsert must not be read-gated");
let status = post(
&app,
&format!("/collections/{COLLECTION}/search"),
json!({"vector": QUERY, "top_k": 1}),
)
.await;
assert!(
!status.is_success(),
"denied dense read must not return success, got {status}"
);
let status = post(
&app,
&format!("/collections/{COLLECTION}/search/text"),
json!({"query": "hello", "top_k": 1}),
)
.await;
assert!(
!status.is_success(),
"denied text read must not return success, got {status}"
);
let status = post(
&app,
&format!("/collections/{COLLECTION}/search/hybrid"),
json!({"vector": QUERY, "query": "hello", "top_k": 1, "vector_weight": 0.5}),
)
.await;
assert!(
!status.is_success(),
"denied hybrid read must not return success, got {status}"
);
}
async fn put(app: &axum::Router, uri: &str, body: Value) -> StatusCode {
app.clone()
.oneshot(
Request::builder()
.method("PUT")
.uri(uri)
.header("Content-Type", "application/json")
.body(Body::from(body.to_string()))
.expect("test: build PUT request"),
)
.await
.expect("test: PUT request")
.status()
}
fn assert_denied(status: StatusCode, what: &str) {
assert!(!status.is_success(), "{what} denied read, got {status}");
}
#[tokio::test]
async fn read_gate_denies_rest_graph_reads_end_to_end() {
const GRAPH: &str = "observed_graph";
let dir = TempDir::new().expect("test: dir");
let app = create_test_app_with_observer(&dir, Arc::new(DenyingReadObserver));
let status = post(
&app,
"/collections",
json!({"name": GRAPH, "collection_type": "graph"}),
)
.await;
assert_eq!(status, StatusCode::CREATED, "create graph collection");
for node_id in [1, 2] {
let uri = format!("/collections/{GRAPH}/graph/nodes/{node_id}/payload");
let status = put(&app, &uri, json!({"payload": {}})).await;
assert_eq!(
status,
StatusCode::NO_CONTENT,
"seed node {node_id} payload must not be read-gated"
);
}
let status = post(
&app,
&format!("/collections/{GRAPH}/graph/edges"),
json!({"id": 1, "source": 1, "target": 2, "label": "KNOWS"}),
)
.await;
assert_eq!(
status,
StatusCode::CREATED,
"add_edge must not be read-gated"
);
assert_denied(
get(
&app,
&format!("/collections/{GRAPH}/graph/edges?label=KNOWS"),
)
.await,
"get_edges",
);
assert_denied(
post(
&app,
&format!("/collections/{GRAPH}/graph/traverse"),
json!({"source": 1, "strategy": "bfs"}),
)
.await,
"traverse_graph",
);
assert_denied(
get(&app, &format!("/collections/{GRAPH}/graph/nodes/1/degree")).await,
"get_node_degree",
);
assert_denied(
get(&app, &format!("/collections/{GRAPH}/graph/edges/count")).await,
"get_edge_count",
);
assert_denied(
get(&app, &format!("/collections/{GRAPH}/graph/nodes")).await,
"list_nodes",
);
assert_denied(
get(&app, &format!("/collections/{GRAPH}/graph/nodes/1/edges")).await,
"get_node_edges",
);
assert_denied(
get(&app, &format!("/collections/{GRAPH}/graph/nodes/1/payload")).await,
"get_node_payload",
);
assert_denied(
post(
&app,
&format!("/collections/{GRAPH}/graph/traverse/parallel"),
json!({"sources": [1]}),
)
.await,
"traverse_parallel",
);
}