Skip to main content

ma_core/kubo/
kubo.rs

1//! Kubo RPC client for IPFS operations.
2//!
3//! HTTP helpers for the Kubo `/api/v0/` endpoints: data add/cat, DAG
4//! put/get, IPNS name publish/resolve, key management, and pinning.
5
6use crate::{Did, Document};
7use anyhow::{anyhow, Result};
8use reqwest::multipart;
9use serde::de::DeserializeOwned;
10use serde::{Deserialize, Serialize};
11use tokio::time::sleep;
12use tracing::warn;
13use web_time::Duration;
14
15// ─── Response types ─────────────────────────────────────────────────────────
16
17#[derive(Debug, Deserialize)]
18struct AddResponse {
19    #[serde(rename = "Hash")]
20    hash: String,
21}
22
23#[derive(Debug, Deserialize)]
24struct DagPutCid {
25    #[serde(rename = "/")]
26    slash: String,
27}
28
29#[derive(Debug, Deserialize)]
30struct DagPutResponse {
31    #[serde(default, rename = "Cid")]
32    cid_upper: Option<DagPutCid>,
33    #[serde(default)]
34    cid: Option<DagPutCid>,
35}
36
37#[derive(Debug, Deserialize)]
38struct NamePublishResponse {
39    #[serde(default, rename = "Value")]
40    value_upper: String,
41    #[serde(default, rename = "value")]
42    value_lower: String,
43}
44
45#[derive(Debug, Deserialize)]
46struct NameResolveResponse {
47    #[serde(default, rename = "Path")]
48    path_upper: String,
49    #[serde(default, rename = "path")]
50    path_lower: String,
51}
52
53#[derive(Debug, Deserialize)]
54struct VersionResponse {
55    #[serde(default, rename = "Version")]
56    version_upper: String,
57    #[serde(default, rename = "version")]
58    version_lower: String,
59}
60
61#[derive(Debug, Deserialize)]
62struct KeyListEntry {
63    #[serde(default, rename = "Name")]
64    name: String,
65    #[serde(default, rename = "name")]
66    name_lower: String,
67    #[serde(default, rename = "Id")]
68    id: String,
69    #[serde(default, rename = "id")]
70    id_lower: String,
71}
72
73#[derive(Debug, Deserialize)]
74struct KeyListResponse {
75    #[serde(default, rename = "Keys")]
76    keys: Vec<KeyListEntry>,
77}
78
79#[derive(Debug, Deserialize)]
80struct KeyImportResponse {
81    #[serde(default, rename = "Name")]
82    name_upper: String,
83    #[serde(default, rename = "name")]
84    name_lower: String,
85    #[serde(default, rename = "Id")]
86    id_upper: String,
87    #[serde(default, rename = "id")]
88    id_lower: String,
89}
90
91#[derive(Clone, Debug)]
92pub struct KuboKey {
93    pub name: String,
94    pub id: String,
95}
96
97// ─── Publish options ────────────────────────────────────────────────────────
98
99#[derive(Clone, Debug)]
100pub struct IpnsPublishOptions {
101    pub timeout: Duration,
102    pub allow_offline: bool,
103    pub lifetime: String,
104    pub ttl: Option<String>,
105    pub resolve: bool,
106    pub quieter: bool,
107}
108
109impl Default for IpnsPublishOptions {
110    fn default() -> Self {
111        Self {
112            timeout: Duration::from_mins(2),
113            allow_offline: true,
114            lifetime: "8760h".to_string(),
115            ttl: None,
116            resolve: false,
117            quieter: true,
118        }
119    }
120}
121
122// ─── Readiness ──────────────────────────────────────────────────────────────
123
124pub async fn wait_for_api(kubo_url: &str, attempts: u32) -> Result<()> {
125    if attempts == 0 {
126        return Err(anyhow!("kubo readiness attempts must be >= 1"));
127    }
128
129    let base = kubo_url.trim_end_matches('/');
130    let url = format!("{base}/api/v0/version");
131    let client = reqwest::Client::builder()
132        .timeout(Duration::from_secs(6))
133        .build()?;
134
135    let mut fib_prev = Duration::from_millis(0);
136    let mut fib_curr = Duration::from_millis(200);
137    let mut last_err: Option<anyhow::Error> = None;
138
139    for attempt in 1..=attempts {
140        let result = async {
141            let response = client.post(&url).send().await?.error_for_status()?;
142            let body = response.text().await?;
143            let parsed: VersionResponse = serde_json::from_str(&body)
144                .map_err(|e| anyhow!("failed parsing version response: {} body={}", e, body))?;
145            let version = if !parsed.version_upper.is_empty() {
146                parsed.version_upper
147            } else {
148                parsed.version_lower
149            };
150            if version.trim().is_empty() {
151                return Err(anyhow!("missing version field in response: {}", body));
152            }
153            Ok::<(), anyhow::Error>(())
154        }
155        .await;
156
157        match result {
158            Ok(()) => return Ok(()),
159            Err(err) => {
160                warn!("kubo readiness {}/{}: {}", attempt, attempts, err);
161                last_err = Some(err);
162                if attempt < attempts {
163                    sleep(fib_curr).await;
164                    let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
165                    fib_prev = fib_curr;
166                    fib_curr = Duration::from_millis(std::cmp::min(next_ms, 5_000) as u64);
167                }
168            }
169        }
170    }
171    Err(anyhow!(
172        "kubo API not ready after {} attempts: {}",
173        attempts,
174        last_err
175            .map(|e| e.to_string())
176            .unwrap_or_else(|| "unknown error".to_string())
177    ))
178}
179
180// ─── Add / Cat ──────────────────────────────────────────────────────────────
181
182pub async fn ipfs_add(kubo_url: &str, data: Vec<u8>) -> Result<String> {
183    let base = kubo_url.trim_end_matches('/');
184    let url = format!("{base}/api/v0/add");
185
186    let part = multipart::Part::bytes(data).file_name("data");
187    let form = multipart::Form::new().part("file", part);
188
189    let client = reqwest::Client::builder()
190        .timeout(Duration::from_secs(10))
191        .build()?;
192
193    let body = client
194        .post(url)
195        .query(&[("pin", "true")])
196        .multipart(form)
197        .send()
198        .await?
199        .error_for_status()?
200        .text()
201        .await?;
202
203    let parsed: AddResponse = serde_json::from_str(&body)
204        .map_err(|e| anyhow!("failed parsing add response: {} body={}", e, body))?;
205    Ok(parsed.hash)
206}
207
208pub async fn cat_bytes(kubo_url: &str, cid: &str) -> Result<Vec<u8>> {
209    let base = kubo_url.trim_end_matches('/');
210    let url = format!("{base}/api/v0/cat");
211
212    let client = reqwest::Client::builder()
213        .timeout(Duration::from_secs(30))
214        .build()?;
215
216    let bytes = client
217        .post(url)
218        .query(&[("arg", cid)])
219        .send()
220        .await?
221        .error_for_status()?
222        .bytes()
223        .await?;
224
225    Ok(bytes.to_vec())
226}
227
228pub async fn cat_text(kubo_url: &str, cid: &str) -> Result<String> {
229    let bytes = cat_bytes(kubo_url, cid).await?;
230    String::from_utf8(bytes).map_err(|e| anyhow!("non-utf8 content from {}: {}", cid, e))
231}
232
233// ─── DAG ────────────────────────────────────────────────────────────────────
234
235pub async fn dag_put<T: Serialize>(kubo_url: &str, value: &T) -> Result<String> {
236    let base = kubo_url.trim_end_matches('/');
237    let url = format!("{base}/api/v0/dag/put");
238    let payload = serde_json::to_vec(value)?;
239
240    let part = multipart::Part::bytes(payload)
241        .file_name("node.json")
242        .mime_str("application/json")?;
243    let form = multipart::Form::new().part("file", part);
244
245    let client = reqwest::Client::builder()
246        .timeout(Duration::from_secs(10))
247        .build()?;
248
249    let body = client
250        .post(url)
251        .query(&[
252            ("store-codec", "dag-cbor"),
253            ("input-codec", "dag-json"),
254            ("pin", "true"),
255        ])
256        .multipart(form)
257        .send()
258        .await?
259        .error_for_status()?
260        .text()
261        .await?;
262
263    let parsed: DagPutResponse = serde_json::from_str(&body)
264        .map_err(|e| anyhow!("failed parsing dag/put response: {} body={}", e, body))?;
265    parsed
266        .cid_upper
267        .or(parsed.cid)
268        .map(|c| c.slash)
269        .ok_or_else(|| anyhow!("missing CID in dag/put response: {}", body))
270}
271
272pub async fn dag_get<T: DeserializeOwned>(kubo_url: &str, cid: &str) -> Result<T> {
273    let base = kubo_url.trim_end_matches('/');
274    let url = format!("{base}/api/v0/dag/get");
275
276    let client = reqwest::Client::builder()
277        .timeout(Duration::from_secs(10))
278        .build()?;
279
280    let body = client
281        .post(url)
282        .query(&[("arg", cid)])
283        .send()
284        .await?
285        .error_for_status()?
286        .text()
287        .await?;
288
289    serde_json::from_str::<T>(&body).map_err(|e| {
290        anyhow!(
291            "failed parsing dag/get response for {}: {} body={}",
292            cid,
293            e,
294            body
295        )
296    })
297}
298
299// ─── Name publish / resolve ─────────────────────────────────────────────────
300
301fn normalize_ipfs_arg(cid_or_path: &str) -> String {
302    let mut value = cid_or_path.trim().to_string();
303    while let Some(rest) = value.strip_prefix("/ipfs/") {
304        value = rest.to_string();
305    }
306    while let Some(rest) = value.strip_prefix('/') {
307        value = rest.to_string();
308    }
309    format!("/ipfs/{value}")
310}
311
312fn normalize_cid_arg(cid_or_path: &str) -> String {
313    let mut value = cid_or_path.trim().to_string();
314    while let Some(rest) = value.strip_prefix("/ipfs/") {
315        value = rest.to_string();
316    }
317    while let Some(rest) = value.strip_prefix('/') {
318        value = rest.to_string();
319    }
320    value
321}
322
323pub async fn name_publish(kubo_url: &str, key_name: &str, cid: &str) -> Result<String> {
324    let options = IpnsPublishOptions::default();
325    name_publish_with_options(kubo_url, key_name, cid, &options).await
326}
327
328pub async fn name_publish_with_options(
329    kubo_url: &str,
330    key_name: &str,
331    cid: &str,
332    options: &IpnsPublishOptions,
333) -> Result<String> {
334    let base = kubo_url.trim_end_matches('/');
335    let url = format!("{base}/api/v0/name/publish");
336    let arg = normalize_ipfs_arg(cid);
337
338    let client = reqwest::Client::builder()
339        .timeout(options.timeout)
340        .build()?;
341
342    let allow_offline = if options.allow_offline {
343        "true"
344    } else {
345        "false"
346    };
347    let resolve = if options.resolve { "true" } else { "false" };
348    let quieter = if options.quieter { "true" } else { "false" };
349
350    let mut params: Vec<(&str, &str)> = vec![
351        ("arg", arg.as_str()),
352        ("key", key_name),
353        ("allow-offline", allow_offline),
354        ("lifetime", options.lifetime.as_str()),
355        ("resolve", resolve),
356        ("quieter", quieter),
357    ];
358    if let Some(ref ttl) = options.ttl {
359        params.push(("ttl", ttl.as_str()));
360    }
361
362    let body = client
363        .post(url)
364        .query(&params)
365        .send()
366        .await?
367        .error_for_status()?
368        .text()
369        .await?;
370
371    let parsed: NamePublishResponse = serde_json::from_str(&body)
372        .map_err(|e| anyhow!("failed parsing name/publish response: {} body={}", e, body))?;
373    let value = if !parsed.value_upper.is_empty() {
374        parsed.value_upper
375    } else {
376        parsed.value_lower
377    };
378    if value.is_empty() {
379        return Err(anyhow!("missing value in name/publish response: {}", body));
380    }
381    Ok(value)
382}
383
384pub async fn name_publish_with_retry(
385    kubo_url: &str,
386    key_name: &str,
387    cid: &str,
388    options: &IpnsPublishOptions,
389    attempts: u32,
390    initial_backoff: Duration,
391) -> Result<String> {
392    if attempts == 0 {
393        return Err(anyhow!("name publish attempts must be >= 1"));
394    }
395
396    let mut fib_prev = Duration::from_millis(0);
397    let mut fib_curr = initial_backoff;
398    let mut last_err: Option<anyhow::Error> = None;
399
400    for attempt in 1..=attempts {
401        match name_publish_with_options(kubo_url, key_name, cid, options).await {
402            Ok(value) => return Ok(value),
403            Err(err) => {
404                if let Ok(value) = verify_name_target_after_error(kubo_url, key_name, cid).await {
405                    warn!(
406                        "name publish attempt {}/{} reported error for key '{}' but resolve confirms target; accepting: {}",
407                        attempt, attempts, key_name, value
408                    );
409                    return Ok(value);
410                }
411                warn!(
412                    "name publish attempt {}/{} failed for key '{}' cid '{}': {}",
413                    attempt, attempts, key_name, cid, err
414                );
415                last_err = Some(err);
416                if attempt < attempts {
417                    sleep(fib_curr).await;
418                    let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
419                    fib_prev = fib_curr;
420                    fib_curr = Duration::from_millis(std::cmp::min(next_ms, 30_000) as u64);
421                }
422            }
423        }
424    }
425
426    Err(anyhow!(
427        "name publish failed after {} attempt(s): {}",
428        attempts,
429        last_err
430            .map(|e| e.to_string())
431            .unwrap_or_else(|| "unknown error".to_string())
432    ))
433}
434
435async fn verify_name_target_after_error(
436    kubo_url: &str,
437    key_name: &str,
438    cid: &str,
439) -> Result<String> {
440    let expected = normalize_ipfs_arg(cid);
441    let resolved = name_resolve(kubo_url, &format!("/ipns/{key_name}"), true).await?;
442    if resolved.trim() == expected {
443        return Ok(resolved);
444    }
445    Err(anyhow!(
446        "post-error resolve mismatch for key '{}': expected '{}' got '{}'",
447        key_name,
448        expected,
449        resolved
450    ))
451}
452
453pub async fn name_resolve(kubo_url: &str, path: &str, recursive: bool) -> Result<String> {
454    let base = kubo_url.trim_end_matches('/');
455    let url = format!("{base}/api/v0/name/resolve");
456
457    let client = reqwest::Client::builder()
458        .timeout(Duration::from_secs(15))
459        .build()?;
460
461    let recursive_flag = if recursive { "true" } else { "false" };
462    let body = client
463        .post(url)
464        .query(&[("arg", path), ("recursive", recursive_flag)])
465        .send()
466        .await?
467        .error_for_status()?
468        .text()
469        .await?;
470
471    let parsed: NameResolveResponse = serde_json::from_str(&body)
472        .map_err(|e| anyhow!("failed parsing name/resolve response: {} body={}", e, body))?;
473    let resolved = if !parsed.path_upper.is_empty() {
474        parsed.path_upper
475    } else {
476        parsed.path_lower
477    };
478    if resolved.is_empty() {
479        return Err(anyhow!("missing path in name/resolve response: {}", body));
480    }
481    Ok(resolved)
482}
483
484// ─── DID document fetch ─────────────────────────────────────────────────────
485
486pub async fn fetch_did_document(kubo_url: &str, did: &Did) -> Result<Document> {
487    let ipns_path = format!("/ipns/{}", did.ipns);
488    let mut fib_prev = Duration::from_millis(0);
489    let mut fib_curr = Duration::from_millis(150);
490    let mut last_err: Option<anyhow::Error> = None;
491    let mut document: Option<Document> = None;
492
493    for attempt in 1..=4 {
494        // DID documents are stored as DAG-CBOR via dag/put.
495        let dag_err = match dag_get::<Document>(kubo_url, &ipns_path).await {
496            Ok(doc) => {
497                document = Some(doc);
498                break;
499            }
500            Err(e) => e,
501        };
502
503        // Fallback: resolve IPNS manually then dag_get the CID.
504        match name_resolve(kubo_url, &ipns_path, true).await {
505            Err(resolve_err) => {
506                last_err = Some(anyhow!(
507                    "dag_get and name/resolve both failed for {}: dag={} resolve={}",
508                    ipns_path,
509                    dag_err,
510                    resolve_err
511                ));
512                if !should_retry_name_resolve_error(&resolve_err) {
513                    break;
514                }
515            }
516            Ok(resolved_path) => match dag_get::<Document>(kubo_url, &resolved_path).await {
517                Ok(doc) => {
518                    document = Some(doc);
519                    break;
520                }
521                Err(err) => {
522                    last_err = Some(anyhow!(
523                        "dag_get failed for {}: direct={} resolved={}",
524                        ipns_path,
525                        dag_err,
526                        err
527                    ));
528                }
529            },
530        }
531
532        if attempt < 4 {
533            sleep(fib_curr).await;
534            let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
535            fib_prev = fib_curr;
536            fib_curr = Duration::from_millis(std::cmp::min(next_ms, 2_000) as u64);
537        }
538    }
539
540    let document = document.ok_or_else(|| {
541        anyhow!(
542            "failed to fetch DID document for {} via {} after retries: {}",
543            did.id(),
544            ipns_path,
545            last_err
546                .map(|e| e.to_string())
547                .unwrap_or_else(|| "unknown error".to_string())
548        )
549    })?;
550
551    document.validate()?;
552    document.verify()?;
553
554    let doc_did = Did::try_from(document.id.as_str())
555        .map_err(|e| anyhow!("DID document has invalid id '{}': {}", document.id, e))?;
556    if doc_did.ipns != did.ipns {
557        return Err(anyhow!(
558            "DID document IPNS mismatch: expected {} but document id is {}",
559            did.base_id(),
560            document.id
561        ));
562    }
563
564    Ok(document)
565}
566
567fn should_retry_name_resolve_error(err: &anyhow::Error) -> bool {
568    let text = err.to_string().to_ascii_lowercase();
569    if text.contains("http status client error") {
570        return false;
571    }
572    if text.contains("missing path in name/resolve response") {
573        return false;
574    }
575    true
576}
577
578// ─── Pin ────────────────────────────────────────────────────────────────────
579
580pub async fn pin_add_named(kubo_url: &str, cid: &str, name: &str) -> Result<()> {
581    let base = kubo_url.trim_end_matches('/');
582    let url = format!("{base}/api/v0/pin/add");
583    let arg = normalize_ipfs_arg(cid);
584
585    let client = reqwest::Client::builder()
586        .timeout(Duration::from_secs(10))
587        .build()?;
588
589    client
590        .post(url)
591        .query(&[("arg", arg.as_str()), ("recursive", "true"), ("name", name)])
592        .send()
593        .await?
594        .error_for_status()?;
595
596    Ok(())
597}
598
599pub async fn pin_rm(kubo_url: &str, cid: &str) -> Result<()> {
600    let base = kubo_url.trim_end_matches('/');
601    let url = format!("{base}/api/v0/pin/rm");
602    let arg = normalize_ipfs_arg(cid);
603
604    let client = reqwest::Client::builder()
605        .timeout(Duration::from_secs(10))
606        .build()?;
607
608    client
609        .post(url)
610        .query(&[("arg", arg.as_str()), ("recursive", "true")])
611        .send()
612        .await?
613        .error_for_status()?;
614
615    Ok(())
616}
617
618pub async fn remote_pin_add_named(
619    kubo_url: &str,
620    service: &str,
621    cid: &str,
622    name: &str,
623) -> Result<()> {
624    let base = kubo_url.trim_end_matches('/');
625    let url = format!("{base}/api/v0/pin/remote/add");
626    let arg = normalize_ipfs_arg(cid);
627
628    let client = reqwest::Client::builder()
629        .timeout(Duration::from_secs(30))
630        .build()?;
631
632    let resp = client
633        .post(url)
634        .query(&[
635            ("arg", arg.as_str()),
636            ("service", service),
637            ("name", name),
638            ("background", "false"),
639        ])
640        .send()
641        .await?;
642    if resp.status().is_success() {
643        return Ok(());
644    }
645    let status = resp.status();
646    let body = resp.text().await.unwrap_or_default();
647    if status.as_u16() == 409
648        || body.contains("DUPLICATE_OBJECT")
649        || body.contains("already pinned")
650        || body.contains("already exists")
651    {
652        return Ok(());
653    }
654    Err(anyhow!(
655        "pin/remote/add {cid} to {service} as {name} failed: {body}"
656    ))
657}
658
659pub async fn remote_pin_rm(kubo_url: &str, service: &str, cid: &str) -> Result<()> {
660    let base = kubo_url.trim_end_matches('/');
661    let url = format!("{base}/api/v0/pin/remote/rm");
662    let arg = normalize_cid_arg(cid);
663
664    let client = reqwest::Client::builder()
665        .timeout(Duration::from_secs(30))
666        .build()?;
667
668    let resp = client
669        .post(url)
670        .query(&[("service", service), ("cid", arg.as_str())])
671        .send()
672        .await?;
673    if resp.status().is_success() {
674        return Ok(());
675    }
676    let status = resp.status();
677    let body = resp.text().await.unwrap_or_default();
678    if status.as_u16() == 404 || body.contains("not found") || body.contains("not pinned") {
679        return Ok(());
680    }
681    Err(anyhow!("pin/remote/rm {cid} from {service} failed: {body}"))
682}
683
684// ─── Key management ─────────────────────────────────────────────────────────
685
686pub async fn generate_key(kubo_url: &str, key_name: &str) -> Result<()> {
687    let base = kubo_url.trim_end_matches('/');
688    let url = format!("{base}/api/v0/key/gen");
689
690    reqwest::Client::builder()
691        .timeout(Duration::from_secs(10))
692        .build()?
693        .post(url)
694        .query(&[("arg", key_name), ("type", "ed25519")])
695        .send()
696        .await?
697        .error_for_status()?;
698
699    Ok(())
700}
701
702pub async fn import_key(kubo_url: &str, key_name: &str, key_bytes: Vec<u8>) -> Result<KuboKey> {
703    let base = kubo_url.trim_end_matches('/');
704    let url = format!("{base}/api/v0/key/import");
705
706    let part = multipart::Part::bytes(key_bytes)
707        .file_name("ipns.key")
708        .mime_str("application/octet-stream")?;
709    let form = multipart::Form::new().part("file", part);
710
711    let response = reqwest::Client::builder()
712        .timeout(Duration::from_secs(10))
713        .build()?
714        .post(url)
715        .query(&[
716            ("arg", key_name),
717            ("ipns-base", "base36"),
718            ("allow-any-key-type", "true"),
719        ])
720        .multipart(form)
721        .send()
722        .await?
723        .error_for_status()?;
724
725    let body = response.text().await?;
726    let parsed: KeyImportResponse = serde_json::from_str(&body)
727        .map_err(|e| anyhow!("failed parsing key/import response: {} body={}", e, body))?;
728
729    let name = if !parsed.name_upper.trim().is_empty() {
730        parsed.name_upper.trim().to_string()
731    } else {
732        parsed.name_lower.trim().to_string()
733    };
734    let id = if !parsed.id_upper.trim().is_empty() {
735        parsed.id_upper.trim().to_string()
736    } else {
737        parsed.id_lower.trim().to_string()
738    };
739
740    if name.is_empty() || id.is_empty() {
741        return Err(anyhow!("missing name/id in key/import response: {}", body));
742    }
743
744    Ok(KuboKey { name, id })
745}
746
747pub async fn list_keys(kubo_url: &str) -> Result<Vec<KuboKey>> {
748    let base = kubo_url.trim_end_matches('/');
749    let url = format!("{base}/api/v0/key/list");
750
751    let body = reqwest::Client::builder()
752        .timeout(Duration::from_secs(10))
753        .build()?
754        .post(url)
755        .send()
756        .await?
757        .error_for_status()?
758        .text()
759        .await?;
760
761    let parsed: KeyListResponse = serde_json::from_str(&body)
762        .map_err(|e| anyhow!("failed parsing key/list response: {} body={}", e, body))?;
763    Ok(parsed
764        .keys
765        .into_iter()
766        .filter_map(|k| {
767            let name = if !k.name.trim().is_empty() {
768                k.name.trim().to_string()
769            } else {
770                k.name_lower.trim().to_string()
771            };
772            let id = if !k.id.trim().is_empty() {
773                k.id.trim().to_string()
774            } else {
775                k.id_lower.trim().to_string()
776            };
777            if name.is_empty() {
778                None
779            } else {
780                Some(KuboKey { name, id })
781            }
782        })
783        .collect())
784}
785
786pub async fn list_key_names(kubo_url: &str) -> Result<Vec<String>> {
787    let keys = list_keys(kubo_url).await?;
788    Ok(keys.into_iter().map(|k| k.name).collect())
789}
790
791/// Remove a named key from the Kubo keystore.
792pub async fn remove_key(kubo_url: &str, key_name: &str) -> Result<()> {
793    let base = kubo_url.trim_end_matches('/');
794    let url = format!("{base}/api/v0/key/rm");
795
796    reqwest::Client::builder()
797        .timeout(Duration::from_secs(10))
798        .build()?
799        .post(url)
800        .query(&[("arg", key_name)])
801        .send()
802        .await?
803        .error_for_status()?;
804
805    Ok(())
806}
807
808// ─── Tests ──────────────────────────────────────────────────────────────────
809
810#[cfg(test)]
811mod tests {
812    use super::*;
813
814    #[test]
815    fn normalize_ipfs_arg_from_raw_cid() {
816        assert_eq!(
817            normalize_ipfs_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
818            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
819        );
820    }
821
822    #[test]
823    fn normalize_ipfs_arg_from_prefixed_path() {
824        assert_eq!(
825            normalize_ipfs_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
826            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
827        );
828    }
829
830    #[test]
831    fn normalize_ipfs_arg_from_double_prefixed_path() {
832        assert_eq!(
833            normalize_ipfs_arg(
834                "/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
835            ),
836            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
837        );
838    }
839
840    #[test]
841    fn normalize_cid_arg_from_raw_cid() {
842        assert_eq!(
843            normalize_cid_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
844            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
845        );
846    }
847
848    #[test]
849    fn normalize_cid_arg_from_prefixed_path() {
850        assert_eq!(
851            normalize_cid_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
852            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
853        );
854    }
855
856    #[test]
857    fn normalize_cid_arg_from_double_prefixed_path() {
858        assert_eq!(
859            normalize_cid_arg(
860                "/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
861            ),
862            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
863        );
864    }
865
866    #[test]
867    fn does_not_retry_http_client_status_errors() {
868        let err = anyhow!(
869            "HTTP status client error (404 Not Found) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
870        );
871        assert!(!should_retry_name_resolve_error(&err));
872    }
873
874    #[test]
875    fn retries_http_server_status_errors() {
876        let err = anyhow!(
877            "HTTP status server error (500 Internal Server Error) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
878        );
879        assert!(should_retry_name_resolve_error(&err));
880    }
881
882    #[test]
883    fn retries_network_errors() {
884        let err =
885            anyhow!("error sending request for url (http://127.0.0.1:5001/api/v0/name/resolve)");
886        assert!(should_retry_name_resolve_error(&err));
887    }
888}