#![allow(clippy::panic)]
use axum::body::{to_bytes, Body};
use axum::http::{Method, Request, StatusCode};
use axum::Router;
use loonfs::{CreateNamespaceOptions, FsWriter, PutFileOptions};
use loonfs::{FsAdmin, FsReader};
use loonfs_api::v0::{
EnableGrepIndexResponse, GrepGcResponse, GrepIndexLifecycle, GrepIndexStatusResponse,
};
use loonfs_api::{
ApiError, CapabilityDocument, ChangeSeq, GrepRequest, GrepResponse, NamespaceId,
FEATURE_DOWNLOADS_DIRECT_GET, FEATURE_QUERY_GREP, FEATURE_UPLOADS_DIRECT_MULTIPART,
FEATURE_UPLOADS_DIRECT_PUT, LIMIT_QUERY_GREP_DEFAULT, LIMIT_QUERY_GREP_MAX,
LIMIT_QUERY_GREP_SCAN_BUDGET_FILES, LIMIT_QUERY_GREP_TAIL_BUDGET_FILES, PROFILE_QUERY_V0,
};
use loonfs_grep::root::{load_grep_root, GrepLifecycle};
use loonfs_grep::{GramIndexBuildPolicy, GrepBuildOutcome, GrepWorker, GREP_INDEX_JOB};
use loonfs_objectstore::local_fs_store::LocalFsStore;
use loonfs_objectstore::SharedObjectStore;
use loonfs_server::{
app, GrepConfig, GrepMode, MaintenanceMode, RuntimeCacheConfigOverrides, ServerConfig,
StoreConfig,
};
use serde::de::DeserializeOwned;
use std::num::NonZeroUsize;
use std::path::Path;
use std::sync::Arc;
use tempfile::tempdir;
use tower::ServiceExt;
#[tokio::test]
async fn disabled_mode_returns_not_supported_and_omits_grep_capabilities() {
let temp_dir = tempdir().expect("store tempdir");
let (_store, _writer, namespace_id) = seed_namespace(temp_dir.path(), "disabled").await;
let (router, server) = app(test_config(temp_dir.path(), GrepMode::Disabled))
.await
.expect("build app");
let capabilities: CapabilityDocument =
response_json(send(&router, Method::GET, "/v0/capabilities", None).await).await;
assert!(!capabilities.features.contains_key(FEATURE_QUERY_GREP));
assert!(
!capabilities.profiles.iter().any(|p| p == PROFILE_QUERY_V0),
"a deployment that answers `not_supported` on every query route must not advertise \
the plane"
);
for limit in grep_limits() {
assert!(!capabilities.limits.contains_key(limit));
}
assert!(!maintains_grep_index(&server));
for path in grep_paths(&namespace_id) {
let response = send(&router, Method::POST, &path, None).await;
assert_eq!(response.status(), StatusCode::NOT_IMPLEMENTED);
let error: ApiError = response_json(response).await;
assert_eq!(error.code, "not_supported");
assert_eq!(error.feature.as_deref(), Some(FEATURE_QUERY_GREP));
}
let status = send(&router, Method::GET, &status_path(&namespace_id), None).await;
assert_eq!(status.status(), StatusCode::NOT_IMPLEMENTED);
server.shutdown().await.expect("settle the server writer");
}
#[tokio::test]
async fn serving_and_maintaining_enables_queries_nudges_and_disables_per_namespace() {
let temp_dir = tempdir().expect("store tempdir");
let (store, writer, namespace_id) = seed_namespace(temp_dir.path(), "both").await;
let (router, server) = app(test_config(temp_dir.path(), GrepMode::ServeAndMaintain))
.await
.expect("build app");
assert!(maintains_grep_index(&server));
assert_eq!(
index_status(&router, &namespace_id).await.state,
GrepIndexLifecycle::Disabled
);
let enabled: EnableGrepIndexResponse = response_json(
send(
&router,
Method::POST,
&format!("/v0/admin/namespaces/{namespace_id}/grep/index/enable"),
None,
)
.await,
)
.await;
assert!(!enabled.already_enabled);
assert!(
matches!(
enabled.state,
GrepIndexLifecycle::Backfilling {
target_seq: ChangeSeq(0),
cursor_inode_id: None,
..
}
),
"{:?}",
enabled.state
);
assert_eq!(
index_status(&router, &namespace_id).await.state,
enabled.state,
"the status route and the enable response describe the same root"
);
settle(&server).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(0));
let steady = index_status(&router, &namespace_id).await;
assert_eq!(
steady.state,
GrepIndexLifecycle::Steady {
built_through_seq: ChangeSeq(0),
next_event_index: 0,
}
);
assert!(!steady.reorganize_pending);
let again: EnableGrepIndexResponse = response_json(
send(
&router,
Method::POST,
&format!("/v0/admin/namespaces/{namespace_id}/grep/index/enable"),
None,
)
.await,
)
.await;
assert!(again.already_enabled);
assert_eq!(again.state, steady.state);
writer
.put_file_bytes(
&namespace_id,
"/note.txt",
b"automatic needle\n",
PutFileOptions::default(),
)
.await
.expect("write file");
settle(&server).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(0));
let capabilities: CapabilityDocument =
response_json(send(&router, Method::GET, "/v0/capabilities", None).await).await;
assert!(capabilities.supports(FEATURE_QUERY_GREP));
assert!(capabilities.profiles.iter().any(|p| p == PROFILE_QUERY_V0));
for limit in grep_limits() {
assert!(capabilities.limits.contains_key(limit));
}
assert_served_document_covers_the_spec_example(&capabilities);
let response = grep(&router, &namespace_id, "automatic needle").await;
assert_eq!(response.matches.len(), 1);
assert_eq!(response.matches[0].absolute_path, "/note.txt");
settle(&server).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(1));
let caught_up = grep(&router, &namespace_id, "automatic needle").await;
assert_eq!(caught_up.matches.len(), 1);
assert_eq!(caught_up.built_through_seq, caught_up.head_seq);
assert_eq!(disable_grep(&router, &namespace_id).await, StatusCode::OK);
let disabled = load_grep_root(&*store, &namespace_id)
.await
.expect("load disabled root")
.expect("disabled root");
assert!(matches!(
disabled.manifest_state().lifecycle(),
GrepLifecycle::Disabled
));
settle(&server).await;
assert!(
matches!(
load_grep_root(&*store, &namespace_id)
.await
.expect("reload disabled root")
.expect("disabled root")
.manifest_state()
.lifecycle(),
GrepLifecycle::Disabled
),
"no step may resurrect a root the operator disabled"
);
let gc: GrepGcResponse = response_json(
send(
&router,
Method::POST,
&format!("/v0/admin/namespaces/{namespace_id}/grep/index/gc"),
Some(b"{}".to_vec()),
)
.await,
)
.await;
assert_eq!(gc.namespace_id, namespace_id);
assert_eq!(
gc.next_cursor, None,
"an unbudgeted pass walks the whole grep keyspace"
);
assert_eq!(enable_grep(&router, &namespace_id).await, StatusCode::OK);
settle(&server).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(1));
let reenabled = grep(&router, &namespace_id, "automatic needle").await;
assert_eq!(reenabled.matches.len(), 1);
server.shutdown().await.expect("settle the server writer");
}
#[tokio::test]
async fn first_query_after_restart_resumes_stale_and_mid_backfill_namespaces() {
let temp_dir = tempdir().expect("store tempdir");
let store = Arc::new(LocalFsStore::new(temp_dir.path()).expect("store")) as SharedObjectStore;
let writer = FsWriter::builder_with_store(store.clone())
.writer_id("restart-seed")
.min_publish_interval_ms(0)
.build()
.await
.expect("writer");
let stale = NamespaceId::parse("restart-stale").expect("namespace id");
let backfill = NamespaceId::parse("restart-backfill").expect("namespace id");
for namespace_id in [&stale, &backfill] {
writer
.create_namespace(namespace_id, CreateNamespaceOptions::default())
.await
.expect("create namespace");
}
writer
.put_file_bytes(
&stale,
"/indexed.txt",
b"indexed before restart\n",
PutFileOptions::default(),
)
.await
.expect("write indexed file");
for index in 0..3 {
writer
.put_file_bytes(
&backfill,
&format!("/backfill-{index}.txt"),
format!("mid-backfill needle {index}\n").as_bytes(),
PutFileOptions::default(),
)
.await
.expect("write backfill file");
}
let worker = grep_worker(&store, "restart-worker").await;
worker.enable(&stale).await.expect("enable stale namespace");
drive_worker_to_current(&worker, &stale, GramIndexBuildPolicy::default()).await;
writer
.put_file_bytes(
&stale,
"/tail.txt",
b"stale steady needle\n",
PutFileOptions::default(),
)
.await
.expect("write unindexed tail");
worker
.enable(&backfill)
.await
.expect("enable backfill namespace");
worker
.build_step(
&backfill,
GramIndexBuildPolicy {
max_files_per_step: NonZeroUsize::MIN,
..GramIndexBuildPolicy::default()
},
)
.await
.expect("leave mid-backfill root");
let root = load_grep_root(&*store, &backfill)
.await
.expect("load root")
.expect("backfill root");
assert!(matches!(
root.manifest_state().lifecycle(),
GrepLifecycle::Backfilling { .. }
));
writer.shutdown().await.expect("shutdown writer");
drop(writer);
drop(worker);
drop(store);
let (router, server) = app(test_config(temp_dir.path(), GrepMode::ServeAndMaintain))
.await
.expect("reopen app");
let stale_response = grep(&router, &stale, "stale steady needle").await;
assert_eq!(stale_response.matches.len(), 1);
settle(&server).await;
let store = Arc::new(LocalFsStore::new(temp_dir.path()).expect("store")) as SharedObjectStore;
assert_eq!(watermark(&store, &stale).await, ChangeSeq(2));
let not_materialized = send(
&router,
Method::POST,
&format!("/v0/namespaces/{backfill}/query/grep"),
Some(
serde_json::to_vec(&grep_request("mid-backfill needle"))
.expect("serialize grep request"),
),
)
.await;
assert!(
matches!(
not_materialized.status(),
StatusCode::OK | StatusCode::NOT_IMPLEMENTED
),
"first touch either observes backfill or its concurrently completed root"
);
settle(&server).await;
assert_eq!(watermark(&store, &backfill).await, ChangeSeq(3));
let resumed = grep(&router, &backfill, "mid-backfill needle").await;
assert_eq!(resumed.matches.len(), 3);
server.shutdown().await.expect("settle the server writer");
}
#[tokio::test]
async fn serve_only_answers_searches_over_an_index_it_refuses_to_administer() {
let temp_dir = tempdir().expect("store tempdir");
let (store, writer, namespace_id) = seed_namespace(temp_dir.path(), "serve-only").await;
let (router, server) = app(test_config(temp_dir.path(), GrepMode::ServeOnly))
.await
.expect("build app");
assert!(!maintains_grep_index(&server));
for path in admin_grep_paths(&namespace_id) {
let response = send(&router, Method::POST, &path, None).await;
assert_eq!(response.status(), StatusCode::NOT_IMPLEMENTED);
let error: ApiError = response_json(response).await;
assert_eq!(error.code, "not_supported");
assert!(
error.message.contains("does not maintain"),
"{}",
error.message
);
}
let status = send(&router, Method::GET, &status_path(&namespace_id), None).await;
assert_eq!(status.status(), StatusCode::NOT_IMPLEMENTED);
let error: ApiError = response_json(status).await;
assert!(
error.message.contains("does not maintain"),
"{}",
error.message
);
let worker = grep_worker(&store, "external-grep-worker").await;
worker.enable(&namespace_id).await.expect("enable grep");
writer
.put_file_bytes(
&namespace_id,
"/note.txt",
b"external needle\n",
PutFileOptions::default(),
)
.await
.expect("write file");
settle(&server).await;
assert_eq!(
lifecycle_of(&store, &namespace_id).await.steady_watermark(),
None,
"a deployment that maintains nothing must leave the backfill where it was"
);
drive_worker_to_current(&worker, &namespace_id, GramIndexBuildPolicy::default()).await;
let response = grep(&router, &namespace_id, "external needle").await;
assert_eq!(response.matches.len(), 1);
assert_eq!(response.built_through_seq, ChangeSeq(1));
server.shutdown().await.expect("settle the server writer");
}
#[tokio::test]
async fn maintain_only_keeps_the_index_built_without_serving_searches() {
let temp_dir = tempdir().expect("store tempdir");
let (store, writer, namespace_id) = seed_namespace(temp_dir.path(), "maintain-only").await;
let (router, server) = app(test_config(temp_dir.path(), GrepMode::MaintainOnly))
.await
.expect("build app");
assert!(maintains_grep_index(&server));
let capabilities: CapabilityDocument =
response_json(send(&router, Method::GET, "/v0/capabilities", None).await).await;
assert!(
!capabilities.features.contains_key(FEATURE_QUERY_GREP),
"a deployment that answers no searches must not advertise that it does"
);
let refused = send(
&router,
Method::POST,
&format!("/v0/namespaces/{namespace_id}/query/grep"),
Some(serde_json::to_vec(&grep_request("needle")).expect("serialize grep request")),
)
.await;
assert_eq!(refused.status(), StatusCode::NOT_IMPLEMENTED);
let error: ApiError = response_json(refused).await;
assert!(
error.message.contains("does not serve grep queries"),
"{}",
error.message
);
writer
.put_file_bytes(
&namespace_id,
"/note.txt",
b"unserved needle\n",
PutFileOptions::default(),
)
.await
.expect("write file");
assert_eq!(enable_grep(&router, &namespace_id).await, StatusCode::OK);
settle(&server).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(1));
let root = load_grep_root(&*store, &namespace_id)
.await
.expect("load root")
.expect("maintained root");
assert!(matches!(
root.manifest_state().lifecycle(),
GrepLifecycle::Steady { .. }
));
assert!(
!root.manifest_state().segments().is_empty(),
"the index this deployment maintains holds real segments"
);
server.shutdown().await.expect("settle the server writer");
}
#[tokio::test]
async fn manual_maintenance_registers_no_index_job_and_still_administers_one() {
let temp_dir = tempdir().expect("store tempdir");
let (store, writer, namespace_id) = seed_namespace(temp_dir.path(), "manual-maintenance").await;
let (router, server) = app(ServerConfig {
maintenance: MaintenanceMode::Manual,
..test_config(temp_dir.path(), GrepMode::ServeAndMaintain)
})
.await
.expect("build app");
assert!(
!maintains_grep_index(&server),
"a manual deployment registers no automatic index job, whatever the grep mode maintains"
);
writer
.put_file_bytes(
&namespace_id,
"/note.txt",
b"unscheduled needle\n",
PutFileOptions::default(),
)
.await
.expect("write file");
assert_eq!(enable_grep(&router, &namespace_id).await, StatusCode::OK);
assert!(
matches!(
index_status(&router, &namespace_id).await.state,
GrepIndexLifecycle::Backfilling { .. }
),
"the enable published a backfill this deployment left for someone else"
);
settle(&server).await;
assert_eq!(
lifecycle_of(&store, &namespace_id).await.steady_watermark(),
None,
"nothing here schedules the backfill it published"
);
let worker = grep_worker(&store, "assigned-grep-host").await;
drive_worker_to_current(&worker, &namespace_id, GramIndexBuildPolicy::default()).await;
assert_eq!(watermark(&store, &namespace_id).await, ChangeSeq(1));
assert_eq!(
grep(&router, &namespace_id, "unscheduled needle")
.await
.matches
.len(),
1
);
server.shutdown().await.expect("settle the server writer");
}
async fn seed_namespace(root: &Path, name: &str) -> (SharedObjectStore, FsWriter, NamespaceId) {
let store = Arc::new(LocalFsStore::new(root).expect("store")) as SharedObjectStore;
let writer = FsWriter::builder_with_store(store.clone())
.writer_id(format!("grep-mode-seed-{name}"))
.min_publish_interval_ms(0)
.build()
.await
.expect("writer");
let namespace_id = NamespaceId::parse(name).expect("namespace id");
writer
.create_namespace(&namespace_id, CreateNamespaceOptions::default())
.await
.expect("create namespace");
(store, writer, namespace_id)
}
fn test_config(store_root: &Path, mode: GrepMode) -> ServerConfig {
ServerConfig {
bind: "127.0.0.1:0".to_owned(),
auth_token: Some("test-token".into()),
content_token_secret: "test-content-token-secret".into(),
writer_id: format!("grep-mode-{mode:?}"),
runtime_cache: RuntimeCacheConfigOverrides::default(),
grep: GrepConfig {
mode,
..GrepConfig::default()
},
maintenance: MaintenanceMode::Automatic,
min_publish_interval_ms: 0,
max_upload_bytes: 1024 * 1024,
max_download_bytes: 1024 * 1024,
max_concurrent_uploads: 2,
max_concurrent_downloads: 2,
max_concurrent_maintenance: 2,
allow_unauthenticated_remote: false,
allow_remote_without_tls: false,
tls: None,
store: StoreConfig::LocalFs {
root: store_root.display().to_string(),
key_prefix: None,
},
}
}
async fn settle(server: &FsWriter) {
server.flush_background().await.expect("settle maintenance");
}
fn maintains_grep_index(server: &FsWriter) -> bool {
server.maintenance_job(GREP_INDEX_JOB).is_some()
}
async fn watermark(store: &SharedObjectStore, namespace_id: &NamespaceId) -> ChangeSeq {
lifecycle_of(store, namespace_id)
.await
.steady_watermark()
.expect("a steady grep root has a watermark")
.0
}
async fn lifecycle_of(store: &SharedObjectStore, namespace_id: &NamespaceId) -> GrepLifecycle {
load_grep_root(&**store, namespace_id)
.await
.expect("load grep root")
.expect("an enabled namespace has a grep root")
.manifest_state()
.lifecycle()
.clone()
}
fn grep_paths(namespace_id: &NamespaceId) -> Vec<String> {
let mut paths = vec![format!("/v0/namespaces/{namespace_id}/query/grep")];
paths.extend(admin_grep_paths(namespace_id));
paths
}
fn admin_grep_paths(namespace_id: &NamespaceId) -> Vec<String> {
["enable", "disable", "gc"]
.into_iter()
.map(|action| format!("/v0/admin/namespaces/{namespace_id}/grep/index/{action}"))
.collect()
}
fn status_path(namespace_id: &NamespaceId) -> String {
format!("/v0/admin/namespaces/{namespace_id}/grep/index")
}
async fn index_status(router: &Router, namespace_id: &NamespaceId) -> GrepIndexStatusResponse {
let response = send(router, Method::GET, &status_path(namespace_id), None).await;
assert_eq!(response.status(), StatusCode::OK);
response_json(response).await
}
async fn enable_grep(router: &Router, namespace_id: &NamespaceId) -> StatusCode {
send(
router,
Method::POST,
&format!("/v0/admin/namespaces/{namespace_id}/grep/index/enable"),
None,
)
.await
.status()
}
async fn disable_grep(router: &Router, namespace_id: &NamespaceId) -> StatusCode {
send(
router,
Method::POST,
&format!("/v0/admin/namespaces/{namespace_id}/grep/index/disable"),
None,
)
.await
.status()
}
async fn grep(router: &Router, namespace_id: &NamespaceId, pattern: &str) -> GrepResponse {
let request = grep_request(pattern);
response_json(
send(
router,
Method::POST,
&format!("/v0/namespaces/{namespace_id}/query/grep"),
Some(serde_json::to_vec(&request).expect("serialize grep request")),
)
.await,
)
.await
}
fn grep_request(pattern: &str) -> GrepRequest {
GrepRequest {
pattern: pattern.to_owned(),
case_insensitive: false,
path_prefix: None,
cursor: None,
limit: None,
allow_stale: false,
allow_scan: false,
}
}
async fn drive_worker_to_current(
worker: &GrepWorker<SharedObjectStore>,
namespace_id: &NamespaceId,
policy: GramIndexBuildPolicy,
) {
for _ in 0..64 {
let build = worker
.build_step(namespace_id, policy)
.await
.expect("build step");
let fold = worker
.reorganize_step(namespace_id, policy)
.await
.expect("fold step");
if matches!(build.outcome, GrepBuildOutcome::UpToDate { .. })
&& matches!(
fold.outcome,
loonfs_grep::GrepReorganizeOutcome::NotNeeded { .. }
)
{
return;
}
}
panic!("grep worker did not catch up");
}
async fn send(
router: &Router,
method: Method,
uri: &str,
body: Option<Vec<u8>>,
) -> axum::response::Response {
let mut request = Request::builder()
.method(method)
.uri(uri)
.header("authorization", "Bearer test-token");
if body.is_some() {
request = request.header("content-type", "application/json");
}
router
.clone()
.oneshot(
request
.body(body.map_or_else(Body::empty, Body::from))
.expect("request"),
)
.await
.expect("route request")
}
async fn response_json<T: DeserializeOwned>(response: axum::response::Response) -> T {
let bytes = to_bytes(response.into_body(), usize::MAX)
.await
.expect("read response body");
serde_json::from_slice(&bytes).expect("decode response JSON")
}
const API_SPEC_PATH: &str = concat!(env!("CARGO_MANIFEST_DIR"), "/../../docs/specs/api.md");
async fn grep_worker(store: &SharedObjectStore, actor: &str) -> GrepWorker<SharedObjectStore> {
let reader = FsReader::builder_with_store(store.clone())
.build()
.await
.expect("build reader");
let admin = FsAdmin::builder_with_store(store.clone())
.actor_id(actor)
.build()
.await
.expect("build admin");
GrepWorker::new(store.clone(), reader, admin)
}
fn assert_served_document_covers_the_spec_example(served: &CapabilityDocument) {
let spec = std::fs::read_to_string(API_SPEC_PATH).expect("read docs/specs/api.md");
let example = spec
.split("### 2.1")
.nth(1)
.expect("api.md section 2.1")
.split("### 2.2")
.next()
.expect("section end")
.split("```json")
.nth(1)
.expect("capability example block")
.split("```")
.next()
.expect("fenced block end");
let mut expected: CapabilityDocument =
serde_json::from_str(example).expect("spec capability example parses");
let mut served_features = served.features.clone();
for feature in [
FEATURE_UPLOADS_DIRECT_PUT,
FEATURE_UPLOADS_DIRECT_MULTIPART,
FEATURE_DOWNLOADS_DIRECT_GET,
] {
expected.features.remove(feature);
served_features.remove(feature);
}
served.validate().expect("served document is well-formed");
assert_eq!(served.protocol_version, expected.protocol_version);
assert_eq!(
served.profiles, expected.profiles,
"the served profiles drifted from the api.md section 2.1 example"
);
assert_eq!(
served_features, expected.features,
"the served features drifted from the api.md section 2.1 example"
);
for (limit, value) in &expected.limits {
assert_eq!(
served.limits.get(limit),
Some(value),
"the served `{limit}` limit drifted from the api.md section 2.1 example"
);
}
}
fn grep_limits() -> [&'static str; 4] {
[
LIMIT_QUERY_GREP_DEFAULT,
LIMIT_QUERY_GREP_MAX,
LIMIT_QUERY_GREP_SCAN_BUDGET_FILES,
LIMIT_QUERY_GREP_TAIL_BUDGET_FILES,
]
}