use std::collections::BTreeMap;
use crate::{Did, Document};
use anyhow::{anyhow, Result};
use reqwest::multipart;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use tokio::time::sleep;
use tracing::warn;
use web_time::Duration;
#[derive(Debug, Deserialize)]
struct AddResponse {
#[serde(rename = "Hash")]
hash: String,
}
#[derive(Debug, Deserialize)]
struct DagPutCid {
#[serde(rename = "/")]
slash: String,
}
#[derive(Debug, Deserialize)]
struct DagPutResponse {
#[serde(default, rename = "Cid")]
cid_upper: Option<DagPutCid>,
#[serde(default)]
cid: Option<DagPutCid>,
}
#[derive(Debug, Deserialize)]
struct NamePublishResponse {
#[serde(default, rename = "Value")]
value_upper: String,
#[serde(default, rename = "value")]
value_lower: String,
}
#[derive(Debug, Deserialize)]
struct NameResolveResponse {
#[serde(default, rename = "Path")]
path_upper: String,
#[serde(default, rename = "path")]
path_lower: String,
}
#[derive(Debug, Deserialize)]
struct VersionResponse {
#[serde(default, rename = "Version")]
version_upper: String,
#[serde(default, rename = "version")]
version_lower: String,
}
#[derive(Debug, Deserialize)]
struct KeyListEntry {
#[serde(default, rename = "Name")]
name: String,
#[serde(default, rename = "name")]
name_lower: String,
#[serde(default, rename = "Id")]
id: String,
#[serde(default, rename = "id")]
id_lower: String,
}
#[derive(Debug, Deserialize)]
struct KeyListResponse {
#[serde(default, rename = "Keys")]
keys: Vec<KeyListEntry>,
}
#[derive(Debug, Deserialize)]
struct KeyImportResponse {
#[serde(default, rename = "Name")]
name_upper: String,
#[serde(default, rename = "name")]
name_lower: String,
#[serde(default, rename = "Id")]
id_upper: String,
#[serde(default, rename = "id")]
id_lower: String,
}
#[derive(Clone, Debug)]
pub struct KuboKey {
pub name: String,
pub id: String,
}
#[derive(Debug, Deserialize)]
struct PinListResponse {
#[serde(default, rename = "Keys")]
keys: BTreeMap<String, PinListEntry>,
}
#[derive(Debug, Deserialize)]
struct PinListEntry {
#[serde(default, rename = "Type")]
pin_type: String,
#[serde(default, rename = "Name")]
name: String,
}
#[derive(Debug, Deserialize)]
struct RemotePinListEntry {
#[serde(default, rename = "Cid")]
cid_upper: String,
#[serde(default, rename = "cid")]
cid_lower: String,
#[serde(default, rename = "Name")]
name_upper: String,
#[serde(default, rename = "name")]
name_lower: String,
}
#[derive(Clone, Debug)]
pub struct IpnsPublishOptions {
pub timeout: Duration,
pub allow_offline: bool,
pub lifetime: String,
pub ttl: Option<String>,
pub resolve: bool,
pub quieter: bool,
}
impl Default for IpnsPublishOptions {
fn default() -> Self {
Self {
timeout: Duration::from_mins(2),
allow_offline: true,
lifetime: "8760h".to_string(),
ttl: None,
resolve: false,
quieter: true,
}
}
}
pub async fn wait_for_api(kubo_url: &str, attempts: u32) -> Result<()> {
if attempts == 0 {
return Err(anyhow!("kubo readiness attempts must be >= 1"));
}
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/version");
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(6))
.build()?;
let mut fib_prev = Duration::from_millis(0);
let mut fib_curr = Duration::from_millis(200);
let mut last_err: Option<anyhow::Error> = None;
for attempt in 1..=attempts {
let result = async {
let response = client.post(&url).send().await?.error_for_status()?;
let body = response.text().await?;
let parsed: VersionResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing version response: {} body={}", e, body))?;
let version = if !parsed.version_upper.is_empty() {
parsed.version_upper
} else {
parsed.version_lower
};
if version.trim().is_empty() {
return Err(anyhow!("missing version field in response: {}", body));
}
Ok::<(), anyhow::Error>(())
}
.await;
match result {
Ok(()) => return Ok(()),
Err(err) => {
warn!("kubo readiness {}/{}: {}", attempt, attempts, err);
last_err = Some(err);
if attempt < attempts {
sleep(fib_curr).await;
let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
fib_prev = fib_curr;
fib_curr = Duration::from_millis(std::cmp::min(next_ms, 5_000) as u64);
}
}
}
}
Err(anyhow!(
"kubo API not ready after {} attempts: {}",
attempts,
last_err
.map(|e| e.to_string())
.unwrap_or_else(|| "unknown error".to_string())
))
}
pub async fn ipfs_add(kubo_url: &str, data: Vec<u8>) -> Result<String> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/add");
let part = multipart::Part::bytes(data).file_name("data");
let form = multipart::Form::new().part("file", part);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
let body = client
.post(url)
.query(&[("pin", "true")])
.multipart(form)
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: AddResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing add response: {} body={}", e, body))?;
Ok(parsed.hash)
}
pub async fn cat_bytes(kubo_url: &str, cid: &str) -> Result<Vec<u8>> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/cat");
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()?;
let bytes = client
.post(url)
.query(&[("arg", cid)])
.send()
.await?
.error_for_status()?
.bytes()
.await?;
Ok(bytes.to_vec())
}
pub async fn cat_text(kubo_url: &str, cid: &str) -> Result<String> {
let bytes = cat_bytes(kubo_url, cid).await?;
String::from_utf8(bytes).map_err(|e| anyhow!("non-utf8 content from {}: {}", cid, e))
}
pub async fn dag_put<T: Serialize>(kubo_url: &str, value: &T) -> Result<String> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/dag/put");
let payload = serde_json::to_vec(value)?;
let part = multipart::Part::bytes(payload)
.file_name("node.json")
.mime_str("application/json")?;
let form = multipart::Form::new().part("file", part);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
let body = client
.post(url)
.query(&[
("store-codec", "dag-cbor"),
("input-codec", "dag-json"),
("pin", "true"),
])
.multipart(form)
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: DagPutResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing dag/put response: {} body={}", e, body))?;
parsed
.cid_upper
.or(parsed.cid)
.map(|c| c.slash)
.ok_or_else(|| anyhow!("missing CID in dag/put response: {}", body))
}
pub async fn dag_put_cbor(kubo_url: &str, data: Vec<u8>, pin: bool) -> Result<String> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/dag/put");
let part = multipart::Part::bytes(data)
.file_name("document.cbor")
.mime_str("application/octet-stream")?;
let form = multipart::Form::new().part("file", part);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
let body = client
.post(url)
.query(&[
("store-codec", "dag-cbor"),
("input-codec", "dag-cbor"),
("pin", if pin { "true" } else { "false" }),
])
.multipart(form)
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: DagPutResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing dag/put response: {} body={}", e, body))?;
parsed
.cid_upper
.or(parsed.cid)
.map(|c| c.slash)
.ok_or_else(|| anyhow!("missing CID in dag/put response: {}", body))
}
pub async fn dag_get<T: DeserializeOwned>(kubo_url: &str, cid: &str) -> Result<T> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/dag/get");
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
let body = client
.post(url)
.query(&[("arg", cid)])
.send()
.await?
.error_for_status()?
.text()
.await?;
serde_json::from_str::<T>(&body).map_err(|e| {
anyhow!(
"failed parsing dag/get response for {}: {} body={}",
cid,
e,
body
)
})
}
fn normalize_ipfs_arg(cid_or_path: &str) -> String {
let mut value = cid_or_path.trim().to_string();
while let Some(rest) = value.strip_prefix("/ipfs/") {
value = rest.to_string();
}
while let Some(rest) = value.strip_prefix('/') {
value = rest.to_string();
}
format!("/ipfs/{value}")
}
fn normalize_cid_arg(cid_or_path: &str) -> String {
let mut value = cid_or_path.trim().to_string();
while let Some(rest) = value.strip_prefix("/ipfs/") {
value = rest.to_string();
}
while let Some(rest) = value.strip_prefix('/') {
value = rest.to_string();
}
value
}
pub async fn name_publish(kubo_url: &str, key_name: &str, cid: &str) -> Result<String> {
let options = IpnsPublishOptions::default();
name_publish_with_options(kubo_url, key_name, cid, &options).await
}
pub async fn name_publish_with_options(
kubo_url: &str,
key_name: &str,
cid: &str,
options: &IpnsPublishOptions,
) -> Result<String> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/name/publish");
let arg = normalize_ipfs_arg(cid);
let client = reqwest::Client::builder()
.timeout(options.timeout)
.build()?;
let allow_offline = if options.allow_offline {
"true"
} else {
"false"
};
let resolve = if options.resolve { "true" } else { "false" };
let quieter = if options.quieter { "true" } else { "false" };
let mut params: Vec<(&str, &str)> = vec![
("arg", arg.as_str()),
("key", key_name),
("allow-offline", allow_offline),
("lifetime", options.lifetime.as_str()),
("resolve", resolve),
("quieter", quieter),
];
if let Some(ref ttl) = options.ttl {
params.push(("ttl", ttl.as_str()));
}
let body = client
.post(url)
.query(¶ms)
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: NamePublishResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing name/publish response: {} body={}", e, body))?;
let value = if !parsed.value_upper.is_empty() {
parsed.value_upper
} else {
parsed.value_lower
};
if value.is_empty() {
return Err(anyhow!("missing value in name/publish response: {}", body));
}
Ok(value)
}
pub async fn name_publish_with_retry(
kubo_url: &str,
key_name: &str,
ipns_id: &str,
cid: &str,
options: &IpnsPublishOptions,
attempts: u32,
initial_backoff: Duration,
) -> Result<String> {
if attempts == 0 {
return Err(anyhow!("name publish attempts must be >= 1"));
}
let mut fib_prev = Duration::from_millis(0);
let mut fib_curr = initial_backoff;
let mut last_err: Option<anyhow::Error> = None;
for attempt in 1..=attempts {
match name_publish_with_options(kubo_url, key_name, cid, options).await {
Ok(value) => return Ok(value),
Err(err) => {
if let Ok(value) = verify_name_target_after_error(kubo_url, ipns_id, cid).await {
warn!(
"name publish attempt {}/{} reported error for key '{}' but resolve confirms target; accepting: {}",
attempt, attempts, key_name, value
);
return Ok(value);
}
warn!(
"name publish attempt {}/{} failed for key '{}' cid '{}': {}",
attempt, attempts, key_name, cid, err
);
last_err = Some(err);
if attempt < attempts {
sleep(fib_curr).await;
let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
fib_prev = fib_curr;
fib_curr = Duration::from_millis(std::cmp::min(next_ms, 30_000) as u64);
}
}
}
}
Err(anyhow!(
"name publish failed after {} attempt(s): {}",
attempts,
last_err
.map(|e| e.to_string())
.unwrap_or_else(|| "unknown error".to_string())
))
}
async fn verify_name_target_after_error(
kubo_url: &str,
ipns_id: &str,
cid: &str,
) -> Result<String> {
let expected = normalize_ipfs_arg(cid);
let resolved = name_resolve(kubo_url, &format!("/ipns/{ipns_id}"), true).await?;
if resolved.trim() == expected {
return Ok(resolved);
}
Err(anyhow!(
"post-error resolve mismatch for IPNS id '{}': expected '{}' got '{}'",
ipns_id,
expected,
resolved
))
}
pub async fn name_resolve(kubo_url: &str, path: &str, recursive: bool) -> Result<String> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/name/resolve");
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(15))
.build()?;
let recursive_flag = if recursive { "true" } else { "false" };
let body = client
.post(url)
.query(&[("arg", path), ("recursive", recursive_flag)])
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: NameResolveResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing name/resolve response: {} body={}", e, body))?;
let resolved = if !parsed.path_upper.is_empty() {
parsed.path_upper
} else {
parsed.path_lower
};
if resolved.is_empty() {
return Err(anyhow!("missing path in name/resolve response: {}", body));
}
Ok(resolved)
}
pub async fn fetch_did_document(kubo_url: &str, did: &Did) -> Result<Document> {
let ipns_path = format!("/ipns/{}", did.ipns);
let mut fib_prev = Duration::from_millis(0);
let mut fib_curr = Duration::from_millis(150);
let mut last_err: Option<anyhow::Error> = None;
let mut document: Option<Document> = None;
for attempt in 1..=4 {
let dag_err = match dag_get::<Document>(kubo_url, &ipns_path).await {
Ok(doc) => {
document = Some(doc);
break;
}
Err(e) => e,
};
match name_resolve(kubo_url, &ipns_path, true).await {
Err(resolve_err) => {
last_err = Some(anyhow!(
"dag_get and name/resolve both failed for {}: dag={} resolve={}",
ipns_path,
dag_err,
resolve_err
));
if !should_retry_name_resolve_error(&resolve_err) {
break;
}
}
Ok(resolved_path) => match dag_get::<Document>(kubo_url, &resolved_path).await {
Ok(doc) => {
document = Some(doc);
break;
}
Err(err) => {
last_err = Some(anyhow!(
"dag_get failed for {}: direct={} resolved={}",
ipns_path,
dag_err,
err
));
}
},
}
if attempt < 4 {
sleep(fib_curr).await;
let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
fib_prev = fib_curr;
fib_curr = Duration::from_millis(std::cmp::min(next_ms, 2_000) as u64);
}
}
let document = document.ok_or_else(|| {
anyhow!(
"failed to fetch DID document for {} via {} after retries: {}",
did.id(),
ipns_path,
last_err
.map(|e| e.to_string())
.unwrap_or_else(|| "unknown error".to_string())
)
})?;
document.validate()?;
document.verify()?;
let doc_did = Did::try_from(document.id.as_str())
.map_err(|e| anyhow!("DID document has invalid id '{}': {}", document.id, e))?;
if doc_did.ipns != did.ipns {
return Err(anyhow!(
"DID document IPNS mismatch: expected {} but document id is {}",
did.base_id(),
document.id
));
}
Ok(document)
}
fn should_retry_name_resolve_error(err: &anyhow::Error) -> bool {
let text = err.to_string().to_ascii_lowercase();
if text.contains("http status client error") {
return false;
}
if text.contains("missing path in name/resolve response") {
return false;
}
true
}
pub async fn pin_add_named(kubo_url: &str, cid: &str, name: &str) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/add");
let arg = normalize_ipfs_arg(cid);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
client
.post(url)
.query(&[("arg", arg.as_str()), ("recursive", "true"), ("name", name)])
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn pin_rm(kubo_url: &str, cid: &str) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/rm");
let arg = normalize_ipfs_arg(cid);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?;
client
.post(url)
.query(&[("arg", arg.as_str()), ("recursive", "true")])
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn list_named_recursive_pins(kubo_url: &str, name: &str) -> Result<Vec<String>> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/ls");
let body = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()?
.post(url)
.query(&[("type", "recursive"), ("name", name), ("names", "true")])
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: PinListResponse = serde_json::from_str(&body)
.map_err(|error| anyhow!("failed parsing pin/ls response: {error} body={body}"))?;
Ok(parsed
.keys
.into_iter()
.filter_map(|(cid, pin)| (pin.pin_type == "recursive" && pin.name == name).then_some(cid))
.collect())
}
pub async fn remote_pin_add_named(
kubo_url: &str,
service: &str,
cid: &str,
name: &str,
) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/remote/add");
let arg = normalize_ipfs_arg(cid);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()?;
let resp = client
.post(url)
.query(&[
("arg", arg.as_str()),
("service", service),
("name", name),
("background", "false"),
])
.send()
.await?;
if resp.status().is_success() {
return Ok(());
}
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
if status.as_u16() == 409
|| body.contains("DUPLICATE_OBJECT")
|| body.contains("already pinned")
|| body.contains("already exists")
{
return Ok(());
}
Err(anyhow!(
"pin/remote/add {cid} to {service} as {name} failed: {body}"
))
}
pub async fn remote_pin_rm(kubo_url: &str, service: &str, cid: &str) -> Result<()> {
remote_pin_rm_query(kubo_url, service, cid, &[("force", "true")]).await
}
pub async fn remote_pin_rm_named(
kubo_url: &str,
service: &str,
cid: &str,
name: &str,
) -> Result<()> {
remote_pin_rm_query(kubo_url, service, cid, &[("name", name), ("force", "true")]).await
}
const REMOTE_PIN_STATUSES: [(&str, &str); 4] = [
("status", "queued"),
("status", "pinning"),
("status", "pinned"),
("status", "failed"),
];
async fn remote_pin_rm_query(
kubo_url: &str,
service: &str,
cid: &str,
extra: &[(&str, &str)],
) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/remote/rm");
let arg = normalize_cid_arg(cid);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()?;
let mut query: Vec<(&str, &str)> = vec![("service", service), ("cid", arg.as_str())];
query.extend_from_slice(&REMOTE_PIN_STATUSES);
query.extend_from_slice(extra);
let resp = client.post(url).query(&query).send().await?;
if resp.status().is_success() {
return Ok(());
}
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
if status.as_u16() == 404 || body.contains("not found") || body.contains("not pinned") {
return Ok(());
}
Err(anyhow!("pin/remote/rm {cid} from {service} failed: {body}"))
}
pub async fn list_named_remote_pins(
kubo_url: &str,
service: &str,
name: &str,
) -> Result<Vec<String>> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/pin/remote/ls");
let mut query: Vec<(&str, &str)> = vec![("service", service), ("name", name)];
query.extend_from_slice(&REMOTE_PIN_STATUSES);
let body = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()?
.post(url)
.query(&query)
.send()
.await?
.error_for_status()?
.text()
.await?;
let mut cids = Vec::new();
for line in body.lines().map(str::trim).filter(|l| !l.is_empty()) {
let pin: RemotePinListEntry = serde_json::from_str(line).map_err(|error| {
anyhow!("failed parsing pin/remote/ls response line: {error} line={line}")
})?;
let pin_name = if pin.name_upper.is_empty() {
pin.name_lower
} else {
pin.name_upper
};
let cid = if pin.cid_upper.is_empty() {
pin.cid_lower
} else {
pin.cid_upper
};
if pin_name == name && !cid.is_empty() {
cids.push(cid);
}
}
Ok(cids)
}
pub async fn generate_key(kubo_url: &str, key_name: &str) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/key/gen");
reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?
.post(url)
.query(&[("arg", key_name), ("type", "ed25519")])
.send()
.await?
.error_for_status()?;
Ok(())
}
pub async fn import_key(kubo_url: &str, key_name: &str, key_bytes: Vec<u8>) -> Result<KuboKey> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/key/import");
let part = multipart::Part::bytes(key_bytes)
.file_name("ipns.key")
.mime_str("application/octet-stream")?;
let form = multipart::Form::new().part("file", part);
let response = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?
.post(url)
.query(&[
("arg", key_name),
("ipns-base", "base36"),
("allow-any-key-type", "true"),
])
.multipart(form)
.send()
.await?
.error_for_status()?;
let body = response.text().await?;
let parsed: KeyImportResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing key/import response: {} body={}", e, body))?;
let name = if !parsed.name_upper.trim().is_empty() {
parsed.name_upper.trim().to_string()
} else {
parsed.name_lower.trim().to_string()
};
let id = if !parsed.id_upper.trim().is_empty() {
parsed.id_upper.trim().to_string()
} else {
parsed.id_lower.trim().to_string()
};
if name.is_empty() || id.is_empty() {
return Err(anyhow!("missing name/id in key/import response: {}", body));
}
Ok(KuboKey { name, id })
}
pub async fn list_keys(kubo_url: &str) -> Result<Vec<KuboKey>> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/key/list");
let body = reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?
.post(url)
.send()
.await?
.error_for_status()?
.text()
.await?;
let parsed: KeyListResponse = serde_json::from_str(&body)
.map_err(|e| anyhow!("failed parsing key/list response: {} body={}", e, body))?;
Ok(parsed
.keys
.into_iter()
.filter_map(|k| {
let name = if !k.name.trim().is_empty() {
k.name.trim().to_string()
} else {
k.name_lower.trim().to_string()
};
let id = if !k.id.trim().is_empty() {
k.id.trim().to_string()
} else {
k.id_lower.trim().to_string()
};
if name.is_empty() {
None
} else {
Some(KuboKey { name, id })
}
})
.collect())
}
pub async fn list_key_names(kubo_url: &str) -> Result<Vec<String>> {
let keys = list_keys(kubo_url).await?;
Ok(keys.into_iter().map(|k| k.name).collect())
}
pub async fn remove_key(kubo_url: &str, key_name: &str) -> Result<()> {
let base = kubo_url.trim_end_matches('/');
let url = format!("{base}/api/v0/key/rm");
reqwest::Client::builder()
.timeout(Duration::from_secs(10))
.build()?
.post(url)
.query(&[("arg", key_name)])
.send()
.await?
.error_for_status()?;
Ok(())
}
#[cfg(test)]
mod tests {
use std::io::{Read, Write};
use std::net::TcpListener;
use std::thread;
use super::*;
fn read_http_request(stream: &mut std::net::TcpStream) -> Vec<u8> {
let mut request = Vec::new();
let mut buffer = [0_u8; 4096];
let mut expected_length = None;
loop {
let count = stream.read(&mut buffer).expect("read request");
assert_ne!(count, 0, "request ended before its body arrived");
request.extend_from_slice(&buffer[..count]);
if expected_length.is_none() {
if let Some(header_end) = request.windows(4).position(|part| part == b"\r\n\r\n") {
let headers = String::from_utf8_lossy(&request[..header_end]);
expected_length = headers.lines().find_map(|line| {
line.split_once(':').and_then(|(name, value)| {
name.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().expect("content length"))
})
});
assert!(expected_length.is_some(), "request must have a body");
}
}
if let Some(length) = expected_length {
let header_end = request
.windows(4)
.position(|part| part == b"\r\n\r\n")
.expect("headers parsed")
+ 4;
if request.len() >= header_end + length {
return request;
}
}
}
}
#[tokio::test]
async fn dag_put_cbor_preserves_bytes_without_anonymous_pin() {
let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind Kubo mock");
let address = listener.local_addr().expect("mock address");
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().expect("accept Kubo request");
let request = read_http_request(&mut stream);
let body = r#"{"Cid":{"/":"bafy-local-pin"}}"#;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
);
stream
.write_all(response.as_bytes())
.expect("write Kubo response");
request
});
let document = vec![0xd8, 0x2a, 0x43, 0x01, 0x02, 0x03];
let cid = dag_put_cbor(&format!("http://{address}"), document.clone(), false)
.await
.expect("Kubo dag put");
assert_eq!(cid, "bafy-local-pin");
let request = server.join().expect("Kubo mock thread");
let request_text = String::from_utf8_lossy(&request);
assert!(request_text.starts_with(
"POST /api/v0/dag/put?store-codec=dag-cbor&input-codec=dag-cbor&pin=false HTTP/1.1"
));
assert!(request_text
.to_ascii_lowercase()
.contains("content-type: multipart/form-data;"));
assert!(request.windows(document.len()).any(|part| part == document));
}
#[test]
fn normalize_ipfs_arg_from_raw_cid() {
assert_eq!(
normalize_ipfs_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
"/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn normalize_ipfs_arg_from_prefixed_path() {
assert_eq!(
normalize_ipfs_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
"/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn normalize_ipfs_arg_from_double_prefixed_path() {
assert_eq!(
normalize_ipfs_arg(
"/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
),
"/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn normalize_cid_arg_from_raw_cid() {
assert_eq!(
normalize_cid_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
"bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn normalize_cid_arg_from_prefixed_path() {
assert_eq!(
normalize_cid_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
"bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn normalize_cid_arg_from_double_prefixed_path() {
assert_eq!(
normalize_cid_arg(
"/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
),
"bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
);
}
#[test]
fn does_not_retry_http_client_status_errors() {
let err = anyhow!(
"HTTP status client error (404 Not Found) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
);
assert!(!should_retry_name_resolve_error(&err));
}
#[test]
fn retries_http_server_status_errors() {
let err = anyhow!(
"HTTP status server error (500 Internal Server Error) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
);
assert!(should_retry_name_resolve_error(&err));
}
#[test]
fn retries_network_errors() {
let err =
anyhow!("error sending request for url (http://127.0.0.1:5001/api/v0/name/resolve)");
assert!(should_retry_name_resolve_error(&err));
}
}