//! LIVE READ-ONLY integration test harness. //! //! Drives the app's real HTTP stack against a real Proxmox cluster to verify //! every read endpoint and cluster connectivity. All tests are `#[ignore]`d //! and only run with `--ignored`; they read `LIVE_PVE_URL` / `LIVE_PVE_USER` / //! `LIVE_PVE_PASS` from the environment and skip (with an explanation) when //! `LIVE_PVE_PASS` is unset. //! //! HARD RULE: no state-changing HTTP request is ever issued against Proxmox. //! The only non-GET request is the authentication handshake //! `POST /api2/json/access/ticket`, which only issues a session ticket. //! //! The OS keyring is deliberately never touched: `set_session_ticket` injects //! the password-mode session into the in-memory connection state, so //! `login_with_password`/`connect` (which write to the keyring) are NOT used. //! //! Run with: //! TMPDIR=/home/user/.cache/clustri-build-tmp \ //! LIVE_PVE_URL=https://10.0.0.22:8006 LIVE_PVE_USER=ro_opencode@pve \ //! LIVE_PVE_PASS='...' \ //! cargo test --test live -- --ignored --nocapture use std::collections::BTreeSet; use std::time::Duration; use clustri::{ api_request, derive_node_url, AuthContext, AuthMode, ConnectionConfig, ConnectionManager, EndpointConfig, ServerType, }; use reqwest::Method; // --------------------------------------------------------------------------- // Environment + shared helpers // --------------------------------------------------------------------------- /// Reads the live environment. Returns `None` (with an explanation printed) /// when `LIVE_PVE_PASS` is unset or empty so the test can skip cleanly. fn live_env() -> Option<(String, String, String)> { let pass = std::env::var("LIVE_PVE_PASS").ok().filter(|p| !p.is_empty())?; let url = std::env::var("LIVE_PVE_URL") .unwrap_or_else(|_| "https://10.0.0.22:8006".to_string()); let user = std::env::var("LIVE_PVE_USER") .unwrap_or_else(|_| "ro_opencode@pve".to_string()); Some((url, user, pass)) } fn skip_reason() -> String { "LIVE_PVE_PASS env var is not set — skipping live test. Set LIVE_PVE_PASS (and optionally LIVE_PVE_URL / LIVE_PVE_USER) to run it.".to_string() } /// Prints a section label to stderr (visible with `--nocapture`). fn label(msg: &str) { eprintln!("\n========== {} ==========", msg); } /// Builds the same permissive HTTP client the app uses at the transport layer /// (self-signed Proxmox certs are the norm; trust is an application-level /// concern that does not apply to a throwaway read-only harness). fn live_client() -> Result { reqwest::Client::builder() .danger_accept_invalid_certs(true) .connect_timeout(Duration::from_secs(10)) .timeout(Duration::from_secs(45)) .build() .map_err(|e| format!("failed to build HTTP client: {}", e)) } /// Performs the password authentication handshake against the real server: /// `POST {url}/api2/json/access/ticket` with a form-encoded body. Returns /// `(ticket, csrf_token)`. This mirrors the app's own `login_with_password` /// transport (which is unusable here only because it also writes to the /// keyring). This is the ONLY state-adjacent request in the harness and it /// only issues a session ticket. async fn login(url: &str, user: &str, pass: &str) -> Result<(String, String), String> { let client = live_client()?; let body = url::form_urlencoded::Serializer::new(String::new()) .append_pair("username", user) .append_pair("password", pass) .finish(); let login_url = format!("{}/api2/json/access/ticket", url.trim_end_matches('/')); let response = client .post(&login_url) .header(reqwest::header::CONTENT_TYPE, "application/x-www-form-urlencoded") .body(body) .send() .await .map_err(|e| format!("login POST failed: {}", e))?; let status = response.status(); let body: serde_json::Value = response .json() .await .map_err(|e| format!("login response parse failed: {}", e))?; if !status.is_success() { return Err(format!("login failed (HTTP {}): {}", status, body)); } let ticket = body["data"]["ticket"] .as_str() .ok_or_else(|| format!("no ticket in login response: {}", body))? .to_string(); let csrf = body["data"]["CSRFPreventionToken"] .as_str() .unwrap_or("") .to_string(); Ok((ticket, csrf)) } /// Builds a `ConnectionManager` with one password-mode connection (`"conn"`, /// `accept_untrusted: true`) and injects the session ticket in memory via /// `set_session_ticket` (no keyring). Returns the manager, the temp dir /// (kept alive so the connections file stays valid), and the raw ticket. async fn manager_with_session( url: &str, user: &str, pass: &str, ) -> Result<(ConnectionManager, tempfile::TempDir, String), String> { let (ticket, csrf) = login(url, user, pass).await?; let dir = tempfile::tempdir().map_err(|e| format!("tempdir failed: {}", e))?; let path = dir.path().join("connections.json"); let mut manager = ConnectionManager::new(); let config = ConnectionConfig { id: "conn".to_string(), name: "conn".to_string(), primary: EndpointConfig { url: url.trim_end_matches('/').to_string(), node: None, token: None, }, fallbacks: vec![], cert_fingerprint: None, trusted: false, accept_untrusted: true, status: "disconnected".to_string(), cluster_name: None, is_cluster: false, auth_mode: "password".to_string(), username: Some(user.to_string()), nodes: vec![], cluster_id: None, server_type: "pve".to_string(), }; manager .add_connection(config, &path) .await .map_err(|e| format!("add_connection failed: {}", e))?; manager .set_session_ticket("conn", &ticket, &csrf) .await .map_err(|e| format!("set_session_ticket failed: {}", e))?; Ok((manager, dir, ticket)) } /// A password-mode `AuthContext` carrying the login ticket (cookie auth). fn ticket_auth(ticket: &str) -> AuthContext { AuthContext { mode: AuthMode::Password, token: None, ticket: Some(ticket.to_string()), csrf_token: None, server_type: ServerType::Pve, } } /// Fetches the raw (unparsed) VM/container config object for printing the /// verbatim disk / net value strings. async fn raw_config( url: &str, ticket: &str, node: &str, vmid: u32, vm_type: &str, ) -> Result { let client = live_client()?; let path = format!("/nodes/{}/{}/{}/config", node, vm_type, vmid); api_request(&client, url, Method::GET, &path, &ticket_auth(ticket), &[], None) .await .map_err(|e| e.to_string()) } /// Raw authenticated GET that keeps the HTTP status and body even on a /// non-success response, so parameter-verification failures can be inspected /// verbatim. Cookie auth mirrors `api_request` password mode. async fn raw_fetch( url: &str, ticket: &str, path: &str, query: &[(&str, &str)], ) -> Result<(reqwest::StatusCode, serde_json::Value), String> { let client = live_client()?; let mut target = format!("{}/api2/json{}", url.trim_end_matches('/'), path); if !query.is_empty() { let q = query .iter() .map(|(k, v)| format!("{}={}", k, v)) .collect::>() .join("&"); target.push('?'); target.push_str(&q); } let response = client .get(&target) .header(reqwest::header::COOKIE, format!("PVEAuthCookie={}", ticket)) .send() .await .map_err(|e| format!("GET {} failed: {}", target, e))?; let status = response.status(); let body: serde_json::Value = response .json() .await .unwrap_or_else(|_| serde_json::Value::Null); Ok((status, body)) } /// Prints a compact view of a raw JSON value: the `data` payload when the /// standard `{ "data": ... }` envelope is present (an array becomes a count + /// sample; an object prints in full), otherwise the whole envelope so error /// bodies (`errors` / `message`) stay visible. fn print_raw(label: &str, status: reqwest::StatusCode, body: &serde_json::Value) { let payload = match body.get("data") { Some(serde_json::Value::Null) | None => body.clone(), Some(inner) => inner.clone(), }; let text = serde_json::to_string(&payload).unwrap_or_default(); match &payload { serde_json::Value::Array(items) => { eprintln!(" {} -> HTTP {} array ({} entries)", label, status, items.len()); for item in items.iter().take(3) { eprintln!(" {}", serde_json::to_string(item).unwrap_or_default()); } if items.len() > 3 { eprintln!(" ... ({} more)", items.len() - 3); } } serde_json::Value::Object(_) => { eprintln!(" {} -> HTTP {}: {}", label, status, text); } other => eprintln!(" {} -> HTTP {}: {}", label, status, other), } } /// The entries of a raw response's `data` field when it is an array. fn data_array(body: &serde_json::Value) -> Vec { body.get("data") .and_then(|d| d.as_array()) .cloned() .unwrap_or_default() } fn is_disk_key(key: &str) -> bool { let bus_key = |bus: &str| { key.strip_prefix(bus) .map_or(false, |suffix| !suffix.is_empty() && suffix.bytes().all(|b| b.is_ascii_digit())) }; ["scsi", "virtio", "ide", "sata", "nvme"].iter().any(|bus| bus_key(bus)) || key == "rootfs" || bus_key("mp") // LXC mount points (mp0, mp1, ...) } fn is_net_key(key: &str) -> bool { key.strip_prefix("net") .map_or(false, |suffix| !suffix.is_empty() && suffix.bytes().all(|b| b.is_ascii_digit())) } // --------------------------------------------------------------------------- // Test 1: every read endpoint parses against real data // --------------------------------------------------------------------------- #[tokio::test] #[ignore] async fn live_all_read_endpoints_parse() { let Some((url, user, pass)) = live_env() else { eprintln!("{}", skip_reason()); return; }; eprintln!("Live target: {} (user {})", url, user); let (manager, _dir, ticket) = match manager_with_session(&url, &user, &pass).await { Ok(v) => v, Err(e) => panic!("setup failed: {}", e), }; let mut failures: Vec = Vec::new(); let mut first_online_node: Option = None; // ---- get_nodes -------------------------------------------------------- label("get_nodes"); match manager.get_nodes("conn").await { Ok(nodes) => { for n in &nodes { eprintln!( " node={:<10} status={:<8} cpu={:<6} maxcpu={:<4} uptime={}", n.node, n.status, n.cpu, n.maxcpu, n.uptime ); } first_online_node = nodes.iter().find(|n| n.status == "online").map(|n| n.node.clone()); eprintln!(" first_online_node = {:?}", first_online_node); } Err(e) => failures.push(format!("get_nodes: {}", e)), } // ---- get_cluster_status ------------------------------------------------ label("get_cluster_status"); let mut cluster_nodes_len = 0usize; let mut cluster_status_ok = false; match manager.get_cluster_status("conn").await { Ok(status) => { cluster_status_ok = true; let nodes = status.nodes.as_ref(); cluster_nodes_len = nodes.map(|n| n.len()).unwrap_or(0); eprintln!(" cluster name={:?} id={:?} node_count={}", status.name, status.id, cluster_nodes_len); for cn in nodes.into_iter().flatten() { eprintln!( " node name={:<10} nodeid={:<4} online={} local={:?} ip={:?}", cn.name, cn.nodeid, cn.online, cn.local, cn.ip ); } } Err(e) => { failures.push(format!("get_cluster_status: {}", e)); eprintln!(" ERROR: {} (the endpoint was denied / errored)", e); } } if cluster_status_ok { assert!( cluster_nodes_len >= 2, "this is a real cluster: expected >= 2 nodes in /cluster/status, got {}", cluster_nodes_len ); } else { eprintln!( " [note] /cluster/status was not readable for this user, so the '>= 2 cluster nodes' assertion was NOT evaluated" ); } // ---- get_vms ----------------------------------------------------------- label("get_vms"); let mut vms = Vec::new(); match manager.get_vms("conn").await { Ok(list) => { eprintln!(" total VMs: {}", list.len()); eprintln!(" {:<6} {:<5} {:<9} {:<8} {}", "vmid", "type", "status", "node", "name"); for v in &list { eprintln!( " {:<6} {:<5} {:<9} {:<8} {:?}", v.vmid, v.r#type, v.status, v.node, v.name ); } vms = list; } Err(e) => failures.push(format!("get_vms: {}", e)), } // ---- get_storage --------------------------------------------------------- label("get_storage"); let mut storages = Vec::new(); match manager.get_storage("conn").await { Ok(list) => { eprintln!(" total storages: {}", list.len()); for s in &list { eprintln!( " storage={:<20} type={:<6} content={:<30} enabled={} active={} shared={} node={}", s.storage, s.r#type, s.content, s.enabled, s.active, s.shared, s.node ); } storages = list; } Err(e) => failures.push(format!("get_storage: {}", e)), } if storages.is_empty() { eprintln!(" [note] no storages are visible to this user — per-storage content/detail checks below will be skipped"); } // ---- per-storage content + detail ----------------------------------------- label("per-storage content & detail"); for s in &storages { // Probe each storage through its OWN node (the `node` field of the // storage list entry). Node-local storages are only servable by their // owning node — the first-online-node path fails with "storage not // available on node". Shared storages respond from any node. let node = if s.node.is_empty() { first_online_node.as_deref().unwrap_or("") } else { s.node.as_str() }; if node.is_empty() { failures.push(format!( "get_storage_content({}): no node available for storage", s.storage )); failures.push(format!( "get_storage_detail({}): no node available for storage detail", s.storage )); continue; } // Content via the storage's own node. match manager .get_storage_content("conn", &s.storage, Some(node)) .await { Ok(items) => { eprintln!(" storage={:<20} content(Some({})) -> {} entries", s.storage, node, items.len()); if items.is_empty() { eprintln!(" (empty)"); } } Err(e) => { failures.push(format!("get_storage_content({}, Some({})): {}", s.storage, node, e)) } } // Detail via the storage's own node. match manager.get_storage_detail("conn", node, &s.storage).await { Ok(detail) => { eprintln!( " storage={:<20} detail({}) -> storage={} type={} content={} enabled={} active={} shared={} used={} total={} avail={}", s.storage, node, detail.storage, detail.r#type, detail.content, detail.enabled, detail.active, detail.shared, detail.used, detail.total, detail.avail ); } Err(e) => { failures.push(format!("get_storage_detail({}, {}): {}", s.storage, node, e)); // Raw capture: the exact field shape the struct mis-parses. let path = format!("/nodes/{}/storage/{}/status", node, s.storage); match raw_fetch(&url, &ticket, &path, &[]).await { Ok((status, raw)) => print_raw( &format!("RAW /storage/{}/status on {}", s.storage, node), status, &raw, ), Err(e2) => eprintln!(" RAW /storage/{}/status fetch failed: {}", s.storage, e2), } } } } // ---- raw /cluster/resources (type=storage) field shapes -------------------- label("raw /cluster/resources (storage entries)"); match raw_fetch(&url, &ticket, "/cluster/resources", &[("type", "storage")]).await { Ok((status, body)) => { let entries = data_array(&body); eprintln!(" HTTP {} -> {} storage entries", status, entries.len()); if let Some(first) = entries.first() { if let Some(obj) = first.as_object() { let mut keys: Vec<_> = obj.keys().cloned().collect(); keys.sort(); eprintln!(" field keys: {:?}", keys); } } for entry in entries.iter().take(3) { eprintln!(" {}", serde_json::to_string(entry).unwrap_or_default()); } } Err(e) => eprintln!(" ERROR: raw /cluster/resources fetch failed: {}", e), } // ---- get_tasks ------------------------------------------------------------- label("get_tasks"); match manager.get_tasks("conn").await { Ok(tasks) => { eprintln!(" task count: {}", tasks.len()); for t in tasks.iter().take(5) { eprintln!( " upid={} node={} type={} id={} user={} status={:?} exitstatus={:?}", t.upid, t.node, t.r#type, t.id, t.user, t.status, t.exitstatus ); } if tasks.len() > 5 { eprintln!(" ... ({} more)", tasks.len() - 5); } } Err(e) => { eprintln!(" ERROR: get_tasks failed: {}", e); failures.push(format!("get_tasks: {}", e)); } } // ---- raw task probes ---------------------------------------------------------- // The app now sends plain `/cluster/tasks` (the server rejects any // `limit=` query with HTTP 400); probe variants to confirm the rejection // and the node-scoped `/nodes/{node}/tasks` alternative. label("raw task probes"); let probe_node = first_online_node.as_deref().unwrap_or("conn"); let probes: Vec<(String, Vec<(&str, &str)>)> = vec![ ("/cluster/tasks".to_string(), vec![]), ("/cluster/tasks".to_string(), vec![("limit", "50")]), ("/cluster/tasks".to_string(), vec![("limit", "10")]), (format!("/nodes/{}/tasks", probe_node), vec![("limit", "50")]), ]; for (path, query) in &probes { match raw_fetch(&url, &ticket, path, query).await { Ok((status, body)) => print_raw(&format!("GET {}", path), status, &body), Err(e) => eprintln!(" GET {} failed: {}", path, e), } } // ---- get_backup_jobs (parsed + full raw JSON) -------------------------------- label("get_backup_jobs"); match manager.get_backup_jobs("conn").await { Ok(jobs) => { eprintln!(" parsed job count: {}", jobs.len()); for j in &jobs { eprintln!( " id={} store={} schedule={:?} all={} enabled={} node={:?} vmid={:?} compress={:?} mode={:?} quiet={:?}", j.id, j.store, j.schedule, j.all, j.enabled, j.node, j.vmid, j.compress, j.mode, j.quiet ); } } Err(e) => { eprintln!(" ERROR: get_backup_jobs failed: {}", e); failures.push(format!("get_backup_jobs: {}", e)); } } // Raw JSON of every job entry (captures the full field shape the parsed // struct drops). let client = match live_client() { Ok(c) => Some(c), Err(e) => { failures.push(format!("raw backup jobs client: {}", e)); None } }; if let Some(client) = client { let auth = ticket_auth(&ticket); match api_request(&client, &url, Method::GET, "/cluster/backup", &auth, &[], None).await { Ok(raw) => { let entries = raw.as_array().cloned().unwrap_or_default(); eprintln!(" RAW entries ({}):", entries.len()); for entry in &entries { eprintln!(" {}", serde_json::to_string(entry).unwrap_or_default()); } } Err(e) => { eprintln!(" ERROR: raw /cluster/backup fetch failed: {}", e); failures.push(format!("raw /cluster/backup: {}", e)); } } } // ---- get_backups(conn, None) aggregation -------------------------------------- label("get_backups(conn, None)"); match manager.get_backups("conn", None).await { Ok(backups) => { eprintln!(" aggregated backup count: {}", backups.len()); let mut storages = BTreeSet::new(); for b in &backups { storages.insert(b.storage.clone()); if let Some(prefix) = b.volid.split(':').next() { storages.insert(prefix.to_string()); } } eprintln!(" storages aggregated (from storage field / volid prefix): {:?}", storages); for b in backups.iter().take(8) { eprintln!( " volid={} backupid={} type={} id={} time={} size={} storage={}", b.volid, b.backupid, b.backup_type, b.backup_id, b.backup_time, b.size, b.storage ); } if backups.len() > 8 { eprintln!(" ... ({} more)", backups.len() - 8); } } Err(e) => failures.push(format!("get_backups(conn, None): {}", e)), } // ---- raw backup-storage content ------------------------------------------------ // The backup jobs target kashyyk / bb-pve-b2 / ots-pve-b2; dump the raw // node-scoped content so the actual backup volumeids are visible (the // parsed content(None) path returned empty). label("raw backup storage content"); let backup_storages = ["kashyyk", "bb-pve-b2", "ots-pve-b2"]; let content_node = first_online_node.as_deref().unwrap_or("conn"); for storage in backup_storages { let path = format!("/nodes/{}/storage/{}/content", content_node, storage); match raw_fetch(&url, &ticket, &path, &[("content", "backup")]).await { Ok((status, body)) => print_raw( &format!("GET {}/content?content=backup", path), status, &body, ), Err(e) => eprintln!(" GET {} failed: {}", path, e), } } // ---- per-guest disk / network / snapshot raw data -------------------------------- label("guest disks / nics / snapshots"); let first_qemu = vms.iter().find(|v| v.r#type == "qemu"); let first_lxc = vms.iter().find(|v| v.r#type == "lxc"); let guests: Vec<(&str, u32, &str)> = vec![ ("qemu", first_qemu.map(|v| v.vmid).unwrap_or(0), first_qemu.map(|v| v.node.as_str()).unwrap_or("")), ("lxc", first_lxc.map(|v| v.vmid).unwrap_or(0), first_lxc.map(|v| v.node.as_str()).unwrap_or("")), ]; for (vm_type, vmid, node) in guests { if vmid == 0 || node.is_empty() { eprintln!(" no {} guest found — skipping", vm_type); continue; } eprintln!(" --- {} guest vmid={} on node={} ---", vm_type, vmid, node); match manager.get_disks("conn", node, vmid, vm_type).await { Ok(disks) => { eprintln!(" parsed disks ({}):", disks.len()); for d in &disks { eprintln!( " {}: storage={} size={} format={}", d.device, d.storage, d.size, d.format ); } } Err(e) => failures.push(format!("get_disks({}, {}, {}): {}", node, vmid, vm_type, e)), } // Raw disk value strings from the VM config. match raw_config(&url, &ticket, node, vmid, vm_type).await { Ok(cfg) => { eprintln!(" RAW disk config value strings:"); if let Some(obj) = cfg.as_object() { let mut keys: Vec<_> = obj.keys().cloned().collect(); keys.sort(); for k in keys { if is_disk_key(&k) { if let Some(v) = obj.get(&k).and_then(|v| v.as_str()) { eprintln!(" {} = {}", k, v); } } } } } Err(e) => failures.push(format!("raw_config({}, {}, {}): {}", node, vmid, vm_type, e)), } match manager.get_network_interfaces("conn", node, vmid, vm_type).await { Ok(nics) => { eprintln!(" parsed nics ({}):", nics.len()); for nic in &nics { eprintln!( " {}: model={} mac={} bridge={:?} tag={:?} firewall={:?} link_down={:?}", nic.name, nic.model, nic.macaddr, nic.bridge, nic.tag, nic.firewall, nic.link_down ); } } Err(e) => { failures.push(format!("get_network_interfaces({}, {}, {}): {}", node, vmid, vm_type, e)) } } match raw_config(&url, &ticket, node, vmid, vm_type).await { Ok(cfg) => { eprintln!(" RAW net config value strings:"); if let Some(obj) = cfg.as_object() { let mut keys: Vec<_> = obj.keys().cloned().collect(); keys.sort(); for k in keys { if is_net_key(&k) { if let Some(v) = obj.get(&k).and_then(|v| v.as_str()) { eprintln!(" {} = {}", k, v); } } } } } Err(e) => failures.push(format!("raw_config net({}, {}, {}): {}", node, vmid, vm_type, e)), } match manager.get_snapshots("conn", node, vmid, vm_type).await { Ok(snaps) => { eprintln!(" snapshots ({}):", snaps.len()); for s in &snaps { eprintln!( " name={} vmstate={} snaptime={} parent={:?}", s.name, s.vmstate, s.snaptime, s.parent ); } } Err(e) => failures.push(format!("get_snapshots({}, {}, {}): {}", node, vmid, vm_type, e)), } } // ---- get_websocket_url --------------------------------------------------------- label("get_websocket_url"); let ws_node = first_online_node.as_deref().unwrap_or("conn"); match manager.get_websocket_url("conn", ws_node).await { Ok(ws_url) => eprintln!(" websocket base origin: {}", ws_url), Err(e) => failures.push(format!("get_websocket_url: {}", e)), } // ---- summary -------------------------------------------------------------------- label("SUMMARY"); if failures.is_empty() { eprintln!(" ALL READ ENDPOINTS PASSED"); } else { eprintln!(" {} endpoint failures:", failures.len()); for f in &failures { eprintln!(" - {}", f); } panic!("{} endpoint failure(s):\n - {}", failures.len(), failures.join("\n - ")); } } // --------------------------------------------------------------------------- // Test 2: node discovery + per-node reachability // --------------------------------------------------------------------------- #[tokio::test] #[ignore] async fn live_discover_nodes_and_reachability() { let Some((url, user, pass)) = live_env() else { eprintln!("{}", skip_reason()); return; }; eprintln!("Live target: {} (user {})", url, user); let (mut manager, _dir, ticket) = match manager_with_session(&url, &user, &pass).await { Ok(v) => v, Err(e) => panic!("setup failed: {}", e), }; label("discover_nodes"); let primary_url = url.trim_end_matches('/').to_string(); let discovered = match manager.discover_nodes("conn").await { Ok(d) => Some(d), Err(e) => { eprintln!(" discover_nodes FAILED: {}", e); eprintln!( " NOTE: this is a credentials/permission limitation — /cluster/status needs Sys.Audit on /,\n which this user lacks, so cluster node IPs are unavailable for URL derivation.\n Falling back to name-derived URLs via derive_node_url." ); None } }; // Candidate node URLs: real discovered URLs when discovery worked, // otherwise name-derived URLs (best effort). let mut candidates: Vec<(String, bool)> = Vec::new(); // (url, is_primary) match &discovered { Some(nodes) => { for n in nodes { eprintln!( " node={:<10} status={:<8} is_primary={:<5} local={} url={}", n.name, n.status, n.is_primary, n.local, n.url ); candidates.push((n.url.clone(), n.is_primary)); } } None => { let nodes = manager .get_nodes("conn") .await .expect("get_nodes must still work even without /cluster/status"); for n in &nodes { let cand_url = derive_node_url(&primary_url, None, &n.node); let is_primary = cand_url.eq_ignore_ascii_case(&primary_url); eprintln!( " node={:<10} url={}{}", n.node, cand_url, if is_primary { " [PRIMARY]" } else { " (name-derived; cluster IP unknown)" } ); candidates.push((cand_url, is_primary)); } if !candidates.iter().any(|(_, is_primary)| *is_primary) { eprintln!( " (no candidate matched the primary URL — treating the connection's primary URL itself as the primary node)" ); candidates.push((primary_url.clone(), true)); } } } label("per-node reachability (GET /api2/json/version)"); let client = match live_client() { Ok(c) => c, Err(e) => panic!("client build failed: {}", e), }; let auth = ticket_auth(&ticket); let mut reachable = 0usize; let mut unreachable: Vec = Vec::new(); // Direct reachability against each candidate node URL. for (cand_url, is_primary) in &candidates { match api_request(&client, cand_url, Method::GET, "/version", &auth, &[], None).await { Ok(value) => { reachable += 1; eprintln!( " {} REACHABLE -> version={} release={} (primary={})", cand_url, value["version"].as_str().unwrap_or("?"), value["release"].as_str().unwrap_or("?"), is_primary ); } Err(e) => { unreachable.push(format!("{}: {}", cand_url, e)); eprintln!(" {} UNREACHABLE -> {} (primary={})", cand_url, e, is_primary); } } } // Cluster-proxy reachability: every node queried through the primary // endpoint. This works even when the node IPs are unknown (pveproxy // forwards node-scoped requests to the target node) and proves each node // is up and servable by the cluster. let node_names: Vec = match &discovered { Some(nodes) => nodes.iter().map(|n| n.name.clone()).collect(), None => manager .get_nodes("conn") .await .expect("get_nodes must work") .into_iter() .map(|n| n.node) .collect(), }; let mut proxied_reachable = 0usize; for name in &node_names { let path = format!("/nodes/{}/version", name); match api_request(&client, &primary_url, Method::GET, &path, &auth, &[], None).await { Ok(value) => { proxied_reachable += 1; eprintln!( " primary-proxy /nodes/{}/version -> REACHABLE version={}", name, value["version"].as_str().unwrap_or("?") ); } Err(e) => eprintln!(" primary-proxy /nodes/{}/version -> FAILED: {}", name, e), } } // Assert the primary node/endpoint is reachable. assert!( !candidates.is_empty(), "no node candidates at all — cluster unreachable?" ); let primary_target = candidates .iter() .find(|(_, is_primary)| *is_primary) .map(|(u, _)| u.clone()) .unwrap_or_else(|| primary_url.clone()); let primary_result = api_request(&client, &primary_target, Method::GET, "/version", &auth, &[], None).await; assert!( primary_result.is_ok(), "primary node {} must be reachable, got: {:?}", primary_target, primary_result.err() ); eprintln!( " primary node {} REACHABLE (asserted) — version={}", primary_target, primary_result.unwrap()["version"].as_str().unwrap_or("?") ); eprintln!( " direct reachability: {} of {} candidate URLs reachable; cluster-proxy: {} of {} nodes reachable", reachable, candidates.len(), proxied_reachable, node_names.len() ); if !unreachable.is_empty() { eprintln!(" NOTE (not asserted): unreachable candidates:"); for u in &unreachable { eprintln!(" - {}", u); } } } // --------------------------------------------------------------------------- // Test 3: failover to a configured fallback endpoint // --------------------------------------------------------------------------- #[tokio::test] #[ignore] async fn live_failover_to_fallback() { let Some((url, user, pass)) = live_env() else { eprintln!("{}", skip_reason()); return; }; eprintln!("Live target: {} (user {})", url, user); let (ticket, csrf) = match login(&url, &user, &pass).await { Ok(v) => v, Err(e) => panic!("login failed: {}", e), }; let dir = match tempfile::tempdir() { Ok(d) => d, Err(e) => panic!("tempdir failed: {}", e), }; let path = dir.path().join("connections.json"); let mut manager = ConnectionManager::new(); let config = ConnectionConfig { id: "bad".to_string(), name: "bad".to_string(), // 127.0.0.1:1 is always connection-refused (nothing listens on port 1). primary: EndpointConfig { url: "https://127.0.0.1:1".to_string(), node: None, token: None, }, fallbacks: vec![EndpointConfig { url: url.trim_end_matches('/').to_string(), node: None, token: None, }], cert_fingerprint: None, trusted: false, accept_untrusted: true, status: "disconnected".to_string(), cluster_name: None, is_cluster: false, auth_mode: "password".to_string(), username: Some(user.to_string()), nodes: vec![], cluster_id: None, server_type: "pve".to_string(), }; manager .add_connection(config, &path) .await .expect("connection 'bad' should be added"); manager .set_session_ticket("bad", &ticket, &csrf) .await .expect("session ticket should be injected"); label("failover: get_nodes via fallback"); let nodes = match manager.get_nodes("bad").await { Ok(n) => n, Err(e) => panic!("get_nodes('bad') failed — fallback did not serve the request: {}", e), }; eprintln!(" get_nodes succeeded via the real fallback endpoint; node count = {}", nodes.len()); for n in &nodes { eprintln!(" node={} status={}", n.node, n.status); } let status = manager .runtime_status("bad") .expect("runtime_status should be readable"); eprintln!(" runtime_status('bad') = {:?}", status); assert_eq!(status, "failover", "expected runtime_status 'failover' after the primary was refused, got {:?}", status); } // --------------------------------------------------------------------------- // Test 4: cluster identity is consistent across every discovered node // --------------------------------------------------------------------------- #[tokio::test] #[ignore] async fn live_cluster_identity_consistent_across_nodes() { let Some((url, user, pass)) = live_env() else { eprintln!("{}", skip_reason()); return; }; eprintln!("Live target: {} (user {})", url, user); let (mut manager, dir, ticket) = match manager_with_session(&url, &user, &pass).await { Ok(v) => v, Err(e) => panic!("setup failed: {}", e), }; let path = dir.path().join("connections.json"); let discovered = match manager.discover_nodes("conn").await { Ok(d) => d, Err(e) => { eprintln!(" discover_nodes FAILED: {}", e); eprintln!( " Cannot verify cluster identity: /cluster/status is permission-denied for this user\n (needs Sys.Audit on /), so per-node URLs cannot be derived and every node would report the\n same 403. There is no identity data to compare." ); panic!( "live_cluster_identity_consistent_across_nodes: cannot derive per-node cluster identity because /cluster/status is denied: {}", e ); } }; eprintln!("discovered {} nodes", discovered.len()); label("per-node cluster identity"); let mut identities: Vec<(String, String, String)> = Vec::new(); // (name, id, name) let mut direct_failures: Vec = Vec::new(); for (i, node) in discovered.iter().enumerate() { let id = format!("node-{}", i); let config = ConnectionConfig { id: id.clone(), name: id.clone(), primary: EndpointConfig { url: node.url.clone(), node: Some(node.name.clone()), token: None, }, fallbacks: vec![], cert_fingerprint: None, trusted: false, accept_untrusted: true, status: "disconnected".to_string(), cluster_name: None, is_cluster: false, auth_mode: "password".to_string(), username: Some(user.to_string()), nodes: vec![], cluster_id: None, server_type: "pve".to_string(), }; manager .add_connection(config, &path) .await .expect("node connection should be added"); manager .set_session_ticket(&id, &ticket, "") .await .expect("node session ticket should be injected"); match manager.get_cluster_status(&id).await { Ok(status) => { let node_count = status.nodes.as_ref().map(|n| n.len()).unwrap_or(0); eprintln!( " node={:<10} url={} cluster_id={:?} cluster_name={:?} node_count={}", node.name, node.url, status.id, status.name, node_count ); identities.push((node.name.clone(), status.id.clone(), status.name.clone())); } Err(e) => { eprintln!(" node={:<10} url={} FAILED: {}", node.name, node.url, e); direct_failures.push(e.to_string()); } } } if identities.is_empty() { // When every per-node attempt failed at the transport layer, the // derived node URLs are not routable from this host (cluster-internal // IPs). That is an environment limitation, not an identity mismatch: // verify the cluster identity through the primary endpoint (every // node-scoped request is forwarded by pveproxy to the same cluster) // and skip the direct comparison. let all_transport = !direct_failures.is_empty() && direct_failures.iter().all(|f| f.contains("Cannot connect")); if all_transport { eprintln!( " [skip] all {} per-node direct connections failed with transport errors — derived\n\ node URLs are not routable from this host. Verifying cluster identity through\n\ the primary endpoint instead.", direct_failures.len() ); let proxy_nodes = manager .get_nodes("conn") .await .expect("get_nodes must still work through the primary"); match manager.get_cluster_status("conn").await { Ok(status) => { let node_count = status.nodes.as_ref().map(|n| n.len()).unwrap_or(0); eprintln!( " primary-proxy /cluster/status -> cluster_id={:?} cluster_name={:?} node_count={}", status.id, status.name, node_count ); assert!( !status.id.is_empty(), "cluster id from primary /cluster/status is empty — not a real cluster?" ); eprintln!( " primary-proxy per-node reachability: {} nodes served through the primary", proxy_nodes.len() ); eprintln!( " CONSISTENT (via primary proxy): cluster id {:?} serves all {} nodes", status.id, proxy_nodes.len() ); } Err(e) => panic!("cluster identity via primary proxy failed: {}", e), } return; } assert!( !identities.is_empty(), "no node reported a cluster status — cannot compare cluster identity (failures: {:?})", direct_failures ); } let (first_name, first_id, _) = &identities[0]; assert!( !first_id.is_empty(), "cluster id reported by node {} is empty — not a real cluster?", first_name ); for (name, cid, _) in identities.iter().skip(1) { assert_eq!( cid, first_id, "cluster id mismatch: node {} reports {:?} but node {} reports {:?} — same-cluster merging would break", name, cid, first_name, first_id ); } eprintln!( " CONSISTENT: all {} nodes report cluster id {:?}", identities.len(), first_id ); }