This repository has no description
0

Configure Feed

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

core / bobbin / crates / xrpc / tests / coverage.rs
4.6 kB 127 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 Arc::new(bobbin_xrpc::default_directory()), 48 ); 49 Self { coverage, state } 50 } 51} 52 53fn coverage_request() -> Request<Body> { 54 Request::builder() 55 .uri("/xrpc/sh.tangled.bobbin.getCoverage") 56 .body(Body::empty()) 57 .unwrap() 58} 59 60async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { 61 let status = resp.status(); 62 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); 63 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); 64 (status, parsed) 65} 66 67#[tokio::test] 68async fn defaults_to_warming_at_zero() { 69 let h = Harness::new().await; 70 let app = router(h.state.clone()); 71 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 72 assert_eq!(status, StatusCode::OK); 73 assert_eq!(body["ready"], json!(false)); 74 assert_eq!(body["eventsProcessed"], json!(0)); 75 assert_eq!(body["lastCursor"], json!(0)); 76} 77 78#[tokio::test] 79async fn reflects_warming_state() { 80 let h = Harness::new().await; 81 h.coverage.update(|_| Coverage::Warming { 82 events_processed: 7, 83 last_cursor: HydrantCursor::new(21), 84 }); 85 let app = router(h.state.clone()); 86 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 87 assert_eq!(status, StatusCode::OK); 88 assert_eq!(body["ready"], json!(false)); 89 assert_eq!(body["eventsProcessed"], json!(7)); 90 assert_eq!(body["lastCursor"], json!(21)); 91} 92 93#[tokio::test] 94async fn reflects_ready_state() { 95 let h = Harness::new().await; 96 h.coverage.update(|_| Coverage::Ready { 97 events_processed: 99, 98 last_cursor: HydrantCursor::new(5000), 99 }); 100 let app = router(h.state.clone()); 101 let (status, body) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 102 assert_eq!(status, StatusCode::OK); 103 assert_eq!(body["ready"], json!(true)); 104 assert_eq!(body["eventsProcessed"], json!(99)); 105 assert_eq!(body["lastCursor"], json!(5000)); 106} 107 108#[tokio::test] 109async fn promotion_flips_ready_field() { 110 let h = Harness::new().await; 111 let app = router(h.state.clone()); 112 113 h.coverage.update(|_| Coverage::Warming { 114 events_processed: 1, 115 last_cursor: HydrantCursor::new(5), 116 }); 117 let (_, before) = json_response(app.clone().oneshot(coverage_request()).await.unwrap()).await; 118 assert_eq!(before["ready"], json!(false)); 119 120 h.coverage.update(|_| Coverage::Ready { 121 events_processed: 2, 122 last_cursor: HydrantCursor::new(9), 123 }); 124 let (_, after) = json_response(app.oneshot(coverage_request()).await.unwrap()).await; 125 assert_eq!(after["ready"], json!(true)); 126 assert_eq!(after["lastCursor"], json!(9)); 127}