@@ -471,24 +471,11 @@ async fn main() -> anyhow::Result<()> {
471471 } ) ;
472472
473473 // Build the API router
474- let api_state = state. clone ( ) ;
475- let mut app = Router :: new ( )
476- . route ( "/health" , get ( health_check) )
477- . route ( "/status" , get ( status) )
478- . route ( "/stats" , get ( stats) )
479- // Mooring protocol endpoints
480- . route ( "/mooring/init" , post ( mooring_init) )
481- . route ( "/mooring/verify" , post ( mooring_verify) )
482- . route ( "/mooring/commit" , post ( mooring_commit) ) ;
483-
484- // Add metrics endpoint if enabled
474+ let app = build_api_router ( state. clone ( ) , args. metrics_enabled ) ;
485475 if args. metrics_enabled {
486- app = app. route ( "/metrics" , get ( prometheus_metrics) ) ;
487476 info ! ( "Prometheus metrics enabled at /metrics" ) ;
488477 }
489478
490- let app = app. with_state ( api_state) ;
491-
492479 // Bind API to localhost only (Nebula mesh provides external access)
493480 let api_addr = SocketAddr :: from ( ( [ 0 , 0 , 0 , 0 ] , args. api_port ) ) ;
494481 info ! ( "API listening on {}" , api_addr) ;
@@ -499,6 +486,29 @@ async fn main() -> anyhow::Result<()> {
499486 Ok ( ( ) )
500487}
501488
489+ // =============================================================================
490+ // ROUTER CONSTRUCTION
491+ // =============================================================================
492+
493+ /// Build the API router for the yacht-agent.
494+ ///
495+ /// Extracted for testability — the integration tests call this directly.
496+ fn build_api_router ( state : SharedState , metrics_enabled : bool ) -> Router {
497+ let mut app = Router :: new ( )
498+ . route ( "/health" , get ( health_check) )
499+ . route ( "/status" , get ( status) )
500+ . route ( "/stats" , get ( stats) )
501+ . route ( "/mooring/init" , post ( mooring_init) )
502+ . route ( "/mooring/verify" , post ( mooring_verify) )
503+ . route ( "/mooring/commit" , post ( mooring_commit) ) ;
504+
505+ if metrics_enabled {
506+ app = app. route ( "/metrics" , get ( prometheus_metrics) ) ;
507+ }
508+
509+ app. with_state ( state)
510+ }
511+
502512// =============================================================================
503513// DATABASE PROXY
504514// =============================================================================
@@ -1299,3 +1309,175 @@ async fn setup_nftables_firewall() -> Option<NftablesManager> {
12991309 }
13001310 }
13011311}
1312+
1313+ // =============================================================================
1314+ // INTEGRATION TESTS
1315+ // =============================================================================
1316+
1317+ #[ cfg( test) ]
1318+ mod tests {
1319+ use super :: * ;
1320+ use std:: collections:: HashMap ;
1321+ use wharf_core:: mooring:: { LayerManifest , MooringLayer } ;
1322+ use wharf_core:: mooring_client:: MooringClient ;
1323+
1324+ /// End-to-end test of the full mooring protocol flow:
1325+ /// init → verify → commit over real HTTP.
1326+ #[ tokio:: test]
1327+ async fn test_mooring_e2e_flow ( ) {
1328+ // Build the yacht-agent API on a random port
1329+ let state = Arc :: new ( RwLock :: new ( AgentState :: new ( None ) ) ) ;
1330+ let app = build_api_router ( state. clone ( ) , true ) ;
1331+
1332+ let listener = tokio:: net:: TcpListener :: bind ( "127.0.0.1:0" )
1333+ . await
1334+ . expect ( "Failed to bind test listener" ) ;
1335+ let addr = listener. local_addr ( ) . unwrap ( ) ;
1336+
1337+ tokio:: spawn ( async move {
1338+ axum:: serve ( listener, app) . await . unwrap ( ) ;
1339+ } ) ;
1340+
1341+ // Give the server a moment to start
1342+ tokio:: time:: sleep ( std:: time:: Duration :: from_millis ( 50 ) ) . await ;
1343+
1344+ // Create a MooringClient with a fresh keypair
1345+ let keypair = generate_hybrid_keypair ( ) . expect ( "keypair gen" ) ;
1346+ let client = MooringClient :: new ( & format ! ( "http://{}" , addr) , keypair) ;
1347+
1348+ // Step 1: Init mooring session
1349+ let layers = vec ! [ MooringLayer :: Config , MooringLayer :: Files ] ;
1350+ let init_resp = client
1351+ . init_session ( layers, false , false )
1352+ . await
1353+ . expect ( "init_session failed" ) ;
1354+
1355+ assert ! ( !init_resp. session_id. is_empty( ) , "session_id should not be empty" ) ;
1356+ assert_eq ! ( init_resp. accepted_layers. len( ) , 2 ) ;
1357+ assert ! ( init_resp. expires_at > 0 ) ;
1358+
1359+ // Step 2: Verify a layer
1360+ let manifest = LayerManifest {
1361+ files : HashMap :: from ( [
1362+ ( "index.html" . to_string ( ) , "abc123" . to_string ( ) ) ,
1363+ ( "style.css" . to_string ( ) , "def456" . to_string ( ) ) ,
1364+ ] ) ,
1365+ total_size : 2048 ,
1366+ file_count : 2 ,
1367+ root_hash : "root000" . to_string ( ) ,
1368+ } ;
1369+
1370+ let verify_resp = client
1371+ . verify_layer ( & init_resp. session_id , MooringLayer :: Config , manifest)
1372+ . await
1373+ . expect ( "verify_layer failed" ) ;
1374+
1375+ // Without site_root configured, verification passes (dev mode)
1376+ assert ! ( verify_resp. verified) ;
1377+ assert_eq ! ( verify_resp. matched_files, 2 ) ;
1378+
1379+ // Step 3: Commit
1380+ let commit_resp = client
1381+ . commit ( & init_resp. session_id , init_resp. accepted_layers )
1382+ . await
1383+ . expect ( "commit failed" ) ;
1384+
1385+ assert ! ( commit_resp. success) ;
1386+ assert ! ( commit_resp. snapshot_id. is_some( ) ) ;
1387+
1388+ // Verify state was updated
1389+ let s = state. read ( ) . await ;
1390+ assert ! ( s. moored) ;
1391+ assert ! ( s. last_mooring_time. is_some( ) ) ;
1392+ assert_eq ! ( s. mooring_session_count, 1 ) ;
1393+ }
1394+
1395+ /// Test that the health endpoint responds
1396+ #[ tokio:: test]
1397+ async fn test_health_endpoint ( ) {
1398+ let state = Arc :: new ( RwLock :: new ( AgentState :: new ( None ) ) ) ;
1399+ let app = build_api_router ( state, false ) ;
1400+
1401+ let listener = tokio:: net:: TcpListener :: bind ( "127.0.0.1:0" )
1402+ . await
1403+ . unwrap ( ) ;
1404+ let addr = listener. local_addr ( ) . unwrap ( ) ;
1405+
1406+ tokio:: spawn ( async move {
1407+ axum:: serve ( listener, app) . await . unwrap ( ) ;
1408+ } ) ;
1409+
1410+ tokio:: time:: sleep ( std:: time:: Duration :: from_millis ( 50 ) ) . await ;
1411+
1412+ let resp = reqwest:: get ( format ! ( "http://{}/health" , addr) )
1413+ . await
1414+ . expect ( "health request failed" ) ;
1415+
1416+ assert ! ( resp. status( ) . is_success( ) ) ;
1417+ assert_eq ! ( resp. text( ) . await . unwrap( ) , "OK" ) ;
1418+ }
1419+
1420+ /// Test that metrics endpoint returns real counters
1421+ #[ tokio:: test]
1422+ async fn test_metrics_endpoint ( ) {
1423+ let state = Arc :: new ( RwLock :: new ( AgentState :: new ( None ) ) ) ;
1424+ {
1425+ let mut s = state. write ( ) . await ;
1426+ s. queries_allowed = 42 ;
1427+ s. queries_blocked = 3 ;
1428+ }
1429+ let app = build_api_router ( state, true ) ;
1430+
1431+ let listener = tokio:: net:: TcpListener :: bind ( "127.0.0.1:0" )
1432+ . await
1433+ . unwrap ( ) ;
1434+ let addr = listener. local_addr ( ) . unwrap ( ) ;
1435+
1436+ tokio:: spawn ( async move {
1437+ axum:: serve ( listener, app) . await . unwrap ( ) ;
1438+ } ) ;
1439+
1440+ tokio:: time:: sleep ( std:: time:: Duration :: from_millis ( 50 ) ) . await ;
1441+
1442+ let resp = reqwest:: get ( format ! ( "http://{}/metrics" , addr) )
1443+ . await
1444+ . expect ( "metrics request failed" ) ;
1445+
1446+ let body = resp. text ( ) . await . unwrap ( ) ;
1447+ assert ! ( body. contains( "yacht_queries_total{status=\" allowed\" } 42" ) ) ;
1448+ assert ! ( body. contains( "yacht_queries_total{status=\" blocked\" } 3" ) ) ;
1449+ }
1450+
1451+ /// Test that stats endpoint returns real counters
1452+ #[ tokio:: test]
1453+ async fn test_stats_endpoint ( ) {
1454+ let state = Arc :: new ( RwLock :: new ( AgentState :: new ( None ) ) ) ;
1455+ {
1456+ let mut s = state. write ( ) . await ;
1457+ s. queries_allowed = 10 ;
1458+ s. queries_blocked = 5 ;
1459+ s. queries_audited = 2 ;
1460+ }
1461+ let app = build_api_router ( state, false ) ;
1462+
1463+ let listener = tokio:: net:: TcpListener :: bind ( "127.0.0.1:0" )
1464+ . await
1465+ . unwrap ( ) ;
1466+ let addr = listener. local_addr ( ) . unwrap ( ) ;
1467+
1468+ tokio:: spawn ( async move {
1469+ axum:: serve ( listener, app) . await . unwrap ( ) ;
1470+ } ) ;
1471+
1472+ tokio:: time:: sleep ( std:: time:: Duration :: from_millis ( 50 ) ) . await ;
1473+
1474+ let resp = reqwest:: get ( format ! ( "http://{}/stats" , addr) )
1475+ . await
1476+ . expect ( "stats request failed" ) ;
1477+
1478+ let body: serde_json:: Value = resp. json ( ) . await . unwrap ( ) ;
1479+ assert_eq ! ( body[ "queries" ] [ "allowed" ] , 10 ) ;
1480+ assert_eq ! ( body[ "queries" ] [ "blocked" ] , 5 ) ;
1481+ assert_eq ! ( body[ "queries" ] [ "audited" ] , 2 ) ;
1482+ }
1483+ }
0 commit comments