This repository has no description
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}