Monorepo for Tangled
tangled.org
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}