Monorepo for Tangled tangled.org
3

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / xrpc / tests / coverage.rs
4.5 kB 126 lines
1use std::sync::Arc; 2 3use axum::body::{Body, to_bytes}; 4use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor, StateIndex}; 5use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; 6use bobbin_record_lru::{CacheCapacity, LruRecordStore}; 7use bobbin_resolver::RepoIdResolver; 8use bobbin_runtime::{RuntimeHasher, SystemClock}; 9use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; 10use bobbin_slingshot_client::SlingshotClient; 11use bobbin_xrpc::{AppState, router}; 12use http::{Request, StatusCode}; 13use serde_json::{Value, json}; 14use tower::ServiceExt; 15use url::Url; 16use wiremock::MockServer; 17 18struct Harness { 19 coverage: Arc<CoverageWatch>, 20 state: AppState, 21} 22 23impl Harness { 24 async fn new() -> Self { 25 let server = MockServer::start().await; 26 let coverage = Arc::new(CoverageWatch::new()); 27 let state = AppState::new( 28 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), 29 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), 30 Arc::new(EdgeStore::new(RuntimeHasher::default())), 31 Arc::new(StateIndex::new(RuntimeHasher::default())), 32 Arc::new(StateIndex::new(RuntimeHasher::default())), 33 coverage.clone(), 34 Arc::new( 35 KnotProxy::new( 36 KnotProxyConfig::default(), 37 KnotHttpConfig::default(), 38 Arc::new(SystemClock::new()), 39 RuntimeHasher::default(), 40 ) 41 .unwrap(), 42 ), 43 Arc::new( 44 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), 45 ) as Arc<dyn SearchReader>, 46 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 47 ); 48 Self { coverage, state } 49 } 50} 51 52fn coverage_request() -> Request<Body> { 53 Request::builder() 54 .uri("/xrpc/sh.tangled.bobbin.getCoverage") 55 .body(Body::empty()) 56 .unwrap() 57} 58 59async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { 60 let status = resp.status(); 61 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); 62 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); 63 (status, parsed) 64} 65 66#[tokio::test] 67async fn defaults_to_warming_at_zero() { 68 let h = Harness::new().await; 69 let app = router(h.state.clone()); 70 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 71 assert_eq!(status, StatusCode::OK); 72 assert_eq!(body["ready"], json!(false)); 73 assert_eq!(body["eventsProcessed"], json!(0)); 74 assert_eq!(body["lastCursor"], json!(0)); 75} 76 77#[tokio::test] 78async fn reflects_warming_state() { 79 let h = Harness::new().await; 80 h.coverage.update(|_| Coverage::Warming { 81 events_processed: 7, 82 last_cursor: HydrantCursor::new(21), 83 }); 84 let app = router(h.state.clone()); 85 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 86 assert_eq!(status, StatusCode::OK); 87 assert_eq!(body["ready"], json!(false)); 88 assert_eq!(body["eventsProcessed"], json!(7)); 89 assert_eq!(body["lastCursor"], json!(21)); 90} 91 92#[tokio::test] 93async fn reflects_ready_state() { 94 let h = Harness::new().await; 95 h.coverage.update(|_| Coverage::Ready { 96 events_processed: 99, 97 last_cursor: HydrantCursor::new(5000), 98 }); 99 let app = router(h.state.clone()); 100 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 101 assert_eq!(status, StatusCode::OK); 102 assert_eq!(body["ready"], json!(true)); 103 assert_eq!(body["eventsProcessed"], json!(99)); 104 assert_eq!(body["lastCursor"], json!(5000)); 105} 106 107#[tokio::test] 108async fn promotion_flips_ready_field() { 109 let h = Harness::new().await; 110 let app = router(h.state.clone()); 111 112 h.coverage.update(|_| Coverage::Warming { 113 events_processed: 1, 114 last_cursor: HydrantCursor::new(5), 115 }); 116 let (_, before) = json_response(app.clone().oneshot(coverage_request()).await.unwrap()).await; 117 assert_eq!(before["ready"], json!(false)); 118 119 h.coverage.update(|_| Coverage::Ready { 120 events_processed: 2, 121 last_cursor: HydrantCursor::new(9), 122 }); 123 let (_, after) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 124 assert_eq!(after["ready"], json!(true)); 125 assert_eq!(after["lastCursor"], json!(9)); 126}