use anyhow::Result; use axum::{Form, response::IntoResponse}; use axum_template::RenderHtml; use minijinja::context as template_context; use serde::Deserialize; use std::time::Duration; use crate::{ contextual_error, http::{ context::{AdminRequestContext, admin_template_context}, errors::WebError, }, select_template, }; /// TAP stats response types mod tap_stats { use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] pub struct RepoCount { pub repo_count: i64, } #[derive(Debug, Deserialize, Serialize)] pub struct RecordCount { pub record_count: i64, } #[derive(Debug, Deserialize, Serialize)] pub struct OutboxBuffer { pub outbox_buffer: i64, } #[derive(Debug, Deserialize, Serialize)] pub struct ResyncBuffer { pub resync_buffer: i64, } #[derive(Debug, Deserialize, Serialize)] pub struct Cursors { pub firehose: i64, pub list_repos: Option, } #[derive(Debug, Deserialize, Serialize)] pub struct RepoInfo { pub did: String, pub error: Option, pub handle: Option, pub records: i64, pub retries: i64, pub rev: Option, pub state: String, } } #[derive(Debug, Deserialize)] pub(crate) struct TapSubmitForm { subject: String, } #[derive(Debug, Deserialize)] pub(crate) struct TapInfoForm { subject: String, } /// GET /admin/tap - TAP instance stats pub(crate) async fn handle_admin_tap( admin_ctx: AdminRequestContext, ) -> Result { let canonical_url = format!( "https://{}/admin/tap", admin_ctx.web_context.config.external_base ); let default_context = admin_template_context(&admin_ctx, &canonical_url, "tap"); let render_template = select_template!("admin_tap", false, false, admin_ctx.language); let tap_enabled = admin_ctx.web_context.config.enable_tap; let tap_hostname = &admin_ctx.web_context.config.tap_hostname; // If TAP is not enabled, show a simple message if !tap_enabled { return Ok(RenderHtml( &render_template, admin_ctx.web_context.engine.clone(), template_context! { ..default_context, ..template_context! { tap_enabled => false, tap_hostname => tap_hostname, }}, ) .into_response()); } let base_url = format!("http://{}", tap_hostname); let client = &admin_ctx.web_context.http_client; // Fetch all stats endpoints concurrently let (repo_count, record_count, outbox_buffer, resync_buffer, cursors) = tokio::join!( fetch_tap_stat::(client, &base_url, "/stats/repo-count"), fetch_tap_stat::(client, &base_url, "/stats/record-count"), fetch_tap_stat::(client, &base_url, "/stats/outbox-buffer"), fetch_tap_stat::(client, &base_url, "/stats/resync-buffer"), fetch_tap_stat::(client, &base_url, "/stats/cursors"), ); // Determine connection status let connected = repo_count.is_some(); Ok(RenderHtml( &render_template, admin_ctx.web_context.engine.clone(), template_context! { ..default_context, ..template_context! { tap_enabled => true, tap_hostname => tap_hostname, connected => connected, repo_count => repo_count.map(|r| r.repo_count), record_count => record_count.map(|r| r.record_count), outbox_buffer => outbox_buffer.map(|r| r.outbox_buffer), resync_buffer => resync_buffer.map(|r| r.resync_buffer), firehose_cursor => cursors.as_ref().map(|c| c.firehose), list_repos_cursor => cursors.as_ref().and_then(|c| c.list_repos.clone()), }}, ) .into_response()) } /// Helper to fetch a TAP stats endpoint async fn fetch_tap_stat( client: &reqwest::Client, base_url: &str, path: &str, ) -> Option { let url = format!("{}{}", base_url, path); match client .get(&url) .timeout(Duration::from_secs(30)) .send() .await { Ok(response) if response.status().is_success() => response.json::().await.ok(), Ok(response) => { tracing::debug!(url = %url, status = %response.status(), "TAP stats request failed"); None } Err(e) => { tracing::debug!(url = %url, error = %e, "TAP stats request error"); None } } } /// POST /admin/tap/submit - Submit a DID to TAP for crawling pub(crate) async fn handle_admin_tap_submit( admin_ctx: AdminRequestContext, Form(form): Form, ) -> Result { let canonical_url = format!( "https://{}/admin/tap", admin_ctx.web_context.config.external_base ); let default_context = admin_template_context(&admin_ctx, &canonical_url, "tap"); let render_template = select_template!("admin_tap", false, false, admin_ctx.language); let error_template = select_template!(false, false, admin_ctx.language); let tap_enabled = admin_ctx.web_context.config.enable_tap; let tap_hostname = &admin_ctx.web_context.config.tap_hostname; if !tap_enabled { return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, "TAP is not enabled" ); } let subject = form.subject.trim(); // Resolve the subject (handle or DID) to a DID let document = match admin_ctx .web_context .identity_resolver .resolve(subject) .await { Ok(doc) => doc, Err(err) => { return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, err ); } }; let did = &document.id; let url = format!("http://{}/repos/add", tap_hostname); // Build the request payload let payload = serde_json::json!({ "dids": [did] }); // Send the request to TAP let response = match admin_ctx .web_context .http_client .post(&url) .json(&payload) .timeout(Duration::from_secs(10)) .send() .await { Ok(resp) => resp, Err(err) => { tracing::error!(error = %err, "TAP request failed"); return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, format!("TAP request failed: {}", err) ); } }; if !response.status().is_success() { let status = response.status(); let body = response.text().await.unwrap_or_default(); tracing::warn!(subject = %subject, did = %did, status = %status, body = %body, "TAP submit failed"); return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, format!("TAP returned error: {} - {}", status, body) ); } tracing::info!(subject = %subject, did = %did, "Submitted DID to TAP for crawling"); // Re-fetch stats and render the page with a success message let base_url = format!("http://{}", tap_hostname); let client = &admin_ctx.web_context.http_client; let (repo_count, record_count, outbox_buffer, resync_buffer, cursors) = tokio::join!( fetch_tap_stat::(client, &base_url, "/stats/repo-count"), fetch_tap_stat::(client, &base_url, "/stats/record-count"), fetch_tap_stat::(client, &base_url, "/stats/outbox-buffer"), fetch_tap_stat::(client, &base_url, "/stats/resync-buffer"), fetch_tap_stat::(client, &base_url, "/stats/cursors"), ); let connected = repo_count.is_some(); Ok(RenderHtml( &render_template, admin_ctx.web_context.engine.clone(), template_context! { ..default_context, ..template_context! { tap_enabled => true, tap_hostname => tap_hostname, connected => connected, repo_count => repo_count.map(|r| r.repo_count), record_count => record_count.map(|r| r.record_count), outbox_buffer => outbox_buffer.map(|r| r.outbox_buffer), resync_buffer => resync_buffer.map(|r| r.resync_buffer), firehose_cursor => cursors.as_ref().map(|c| c.firehose), list_repos_cursor => cursors.as_ref().and_then(|c| c.list_repos.clone()), submit_success => true, submit_subject => subject, submit_did => did.clone(), }}, ) .into_response()) } /// POST /admin/tap/info - Get info about a DID from TAP pub(crate) async fn handle_admin_tap_info( admin_ctx: AdminRequestContext, Form(form): Form, ) -> Result { let canonical_url = format!( "https://{}/admin/tap", admin_ctx.web_context.config.external_base ); let default_context = admin_template_context(&admin_ctx, &canonical_url, "tap"); let render_template = select_template!("admin_tap", false, false, admin_ctx.language); let error_template = select_template!(false, false, admin_ctx.language); let tap_enabled = admin_ctx.web_context.config.enable_tap; let tap_hostname = &admin_ctx.web_context.config.tap_hostname; if !tap_enabled { return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, "TAP is not enabled" ); } let subject = form.subject.trim(); // Resolve the subject (handle or DID) to a DID let document = match admin_ctx .web_context .identity_resolver .resolve(subject) .await { Ok(doc) => doc, Err(err) => { return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, err ); } }; let did = &document.id; let url = format!("http://{}/info/{}", tap_hostname, did); // Send the request to TAP let response = match admin_ctx .web_context .http_client .get(&url) .timeout(Duration::from_secs(10)) .send() .await { Ok(resp) => resp, Err(err) => { tracing::error!(error = %err, "TAP info request failed"); return contextual_error!( admin_ctx.web_context, admin_ctx.language, error_template, default_context, format!("TAP request failed: {}", err) ); } }; let info: Option = if response.status().is_success() { response.json().await.ok() } else { None }; // Re-fetch the TAP stats for the page let base_url = format!("http://{}", tap_hostname); let client = &admin_ctx.web_context.http_client; let (repo_count, record_count, outbox_buffer, resync_buffer, cursors) = tokio::join!( fetch_tap_stat::(client, &base_url, "/stats/repo-count"), fetch_tap_stat::(client, &base_url, "/stats/record-count"), fetch_tap_stat::(client, &base_url, "/stats/outbox-buffer"), fetch_tap_stat::(client, &base_url, "/stats/resync-buffer"), fetch_tap_stat::(client, &base_url, "/stats/cursors"), ); let connected = repo_count.is_some(); let info_context = info.as_ref().map(|i| { template_context! { did => i.did.clone(), error => i.error.clone(), handle => i.handle.clone(), records => i.records, retries => i.retries, rev => i.rev.clone(), state => i.state.clone(), } }); Ok(RenderHtml( &render_template, admin_ctx.web_context.engine.clone(), template_context! { ..default_context, ..template_context! { tap_enabled => true, tap_hostname => tap_hostname, connected => connected, repo_count => repo_count.map(|r| r.repo_count), record_count => record_count.map(|r| r.record_count), outbox_buffer => outbox_buffer.map(|r| r.outbox_buffer), resync_buffer => resync_buffer.map(|r| r.resync_buffer), firehose_cursor => cursors.as_ref().map(|c| c.firehose), list_repos_cursor => cursors.as_ref().and_then(|c| c.list_repos.clone()), info_subject => subject, info_result => info_context, info_not_found => info.is_none(), }}, ) .into_response()) }