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 std::collections::BTreeMap;
7
8use crate::{Did, Document};
9use anyhow::{anyhow, Result};
10use reqwest::multipart;
11use serde::de::DeserializeOwned;
12use serde::{Deserialize, Serialize};
13use tokio::time::sleep;
14use tracing::warn;
15use web_time::Duration;
16
17// ─── Response types ─────────────────────────────────────────────────────────
18
19#[derive(Debug, Deserialize)]
20struct AddResponse {
21    #[serde(rename = "Hash")]
22    hash: String,
23}
24
25#[derive(Debug, Deserialize)]
26struct DagPutCid {
27    #[serde(rename = "/")]
28    slash: String,
29}
30
31#[derive(Debug, Deserialize)]
32struct DagPutResponse {
33    #[serde(default, rename = "Cid")]
34    cid_upper: Option<DagPutCid>,
35    #[serde(default)]
36    cid: Option<DagPutCid>,
37}
38
39#[derive(Debug, Deserialize)]
40struct NamePublishResponse {
41    #[serde(default, rename = "Value")]
42    value_upper: String,
43    #[serde(default, rename = "value")]
44    value_lower: String,
45}
46
47#[derive(Debug, Deserialize)]
48struct NameResolveResponse {
49    #[serde(default, rename = "Path")]
50    path_upper: String,
51    #[serde(default, rename = "path")]
52    path_lower: String,
53}
54
55#[derive(Debug, Deserialize)]
56struct VersionResponse {
57    #[serde(default, rename = "Version")]
58    version_upper: String,
59    #[serde(default, rename = "version")]
60    version_lower: String,
61}
62
63#[derive(Debug, Deserialize)]
64struct KeyListEntry {
65    #[serde(default, rename = "Name")]
66    name: String,
67    #[serde(default, rename = "name")]
68    name_lower: String,
69    #[serde(default, rename = "Id")]
70    id: String,
71    #[serde(default, rename = "id")]
72    id_lower: String,
73}
74
75#[derive(Debug, Deserialize)]
76struct KeyListResponse {
77    #[serde(default, rename = "Keys")]
78    keys: Vec<KeyListEntry>,
79}
80
81#[derive(Debug, Deserialize)]
82struct KeyImportResponse {
83    #[serde(default, rename = "Name")]
84    name_upper: String,
85    #[serde(default, rename = "name")]
86    name_lower: String,
87    #[serde(default, rename = "Id")]
88    id_upper: String,
89    #[serde(default, rename = "id")]
90    id_lower: String,
91}
92
93#[derive(Clone, Debug)]
94pub struct KuboKey {
95    pub name: String,
96    pub id: String,
97}
98
99#[derive(Debug, Deserialize)]
100struct PinListResponse {
101    #[serde(default, rename = "Keys")]
102    keys: BTreeMap<String, PinListEntry>,
103}
104
105#[derive(Debug, Deserialize)]
106struct PinListEntry {
107    #[serde(default, rename = "Type")]
108    pin_type: String,
109    #[serde(default, rename = "Name")]
110    name: String,
111}
112
113#[derive(Debug, Deserialize)]
114struct RemotePinListEntry {
115    #[serde(default, rename = "Cid")]
116    cid_upper: String,
117    #[serde(default, rename = "cid")]
118    cid_lower: String,
119    #[serde(default, rename = "Name")]
120    name_upper: String,
121    #[serde(default, rename = "name")]
122    name_lower: String,
123}
124
125// ─── Publish options ────────────────────────────────────────────────────────
126
127#[derive(Clone, Debug)]
128pub struct IpnsPublishOptions {
129    pub timeout: Duration,
130    pub allow_offline: bool,
131    pub lifetime: String,
132    pub ttl: Option<String>,
133    pub resolve: bool,
134    pub quieter: bool,
135}
136
137impl Default for IpnsPublishOptions {
138    fn default() -> Self {
139        Self {
140            timeout: Duration::from_mins(2),
141            allow_offline: true,
142            lifetime: "8760h".to_string(),
143            ttl: None,
144            resolve: false,
145            quieter: true,
146        }
147    }
148}
149
150// ─── Readiness ──────────────────────────────────────────────────────────────
151
152pub async fn wait_for_api(kubo_url: &str, attempts: u32) -> Result<()> {
153    if attempts == 0 {
154        return Err(anyhow!("kubo readiness attempts must be >= 1"));
155    }
156
157    let base = kubo_url.trim_end_matches('/');
158    let url = format!("{base}/api/v0/version");
159    let client = reqwest::Client::builder()
160        .timeout(Duration::from_secs(6))
161        .build()?;
162
163    let mut fib_prev = Duration::from_millis(0);
164    let mut fib_curr = Duration::from_millis(200);
165    let mut last_err: Option<anyhow::Error> = None;
166
167    for attempt in 1..=attempts {
168        let result = async {
169            let response = client.post(&url).send().await?.error_for_status()?;
170            let body = response.text().await?;
171            let parsed: VersionResponse = serde_json::from_str(&body)
172                .map_err(|e| anyhow!("failed parsing version response: {} body={}", e, body))?;
173            let version = if !parsed.version_upper.is_empty() {
174                parsed.version_upper
175            } else {
176                parsed.version_lower
177            };
178            if version.trim().is_empty() {
179                return Err(anyhow!("missing version field in response: {}", body));
180            }
181            Ok::<(), anyhow::Error>(())
182        }
183        .await;
184
185        match result {
186            Ok(()) => return Ok(()),
187            Err(err) => {
188                warn!("kubo readiness {}/{}: {}", attempt, attempts, err);
189                last_err = Some(err);
190                if attempt < attempts {
191                    sleep(fib_curr).await;
192                    let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
193                    fib_prev = fib_curr;
194                    fib_curr = Duration::from_millis(std::cmp::min(next_ms, 5_000) as u64);
195                }
196            }
197        }
198    }
199    Err(anyhow!(
200        "kubo API not ready after {} attempts: {}",
201        attempts,
202        last_err
203            .map(|e| e.to_string())
204            .unwrap_or_else(|| "unknown error".to_string())
205    ))
206}
207
208// ─── Add / Cat ──────────────────────────────────────────────────────────────
209
210pub async fn ipfs_add(kubo_url: &str, data: Vec<u8>) -> Result<String> {
211    let base = kubo_url.trim_end_matches('/');
212    let url = format!("{base}/api/v0/add");
213
214    let part = multipart::Part::bytes(data).file_name("data");
215    let form = multipart::Form::new().part("file", part);
216
217    let client = reqwest::Client::builder()
218        .timeout(Duration::from_secs(10))
219        .build()?;
220
221    let body = client
222        .post(url)
223        .query(&[("pin", "true")])
224        .multipart(form)
225        .send()
226        .await?
227        .error_for_status()?
228        .text()
229        .await?;
230
231    let parsed: AddResponse = serde_json::from_str(&body)
232        .map_err(|e| anyhow!("failed parsing add response: {} body={}", e, body))?;
233    Ok(parsed.hash)
234}
235
236pub async fn cat_bytes(kubo_url: &str, cid: &str) -> Result<Vec<u8>> {
237    let base = kubo_url.trim_end_matches('/');
238    let url = format!("{base}/api/v0/cat");
239
240    let client = reqwest::Client::builder()
241        .timeout(Duration::from_secs(30))
242        .build()?;
243
244    let bytes = client
245        .post(url)
246        .query(&[("arg", cid)])
247        .send()
248        .await?
249        .error_for_status()?
250        .bytes()
251        .await?;
252
253    Ok(bytes.to_vec())
254}
255
256pub async fn cat_text(kubo_url: &str, cid: &str) -> Result<String> {
257    let bytes = cat_bytes(kubo_url, cid).await?;
258    String::from_utf8(bytes).map_err(|e| anyhow!("non-utf8 content from {}: {}", cid, e))
259}
260
261// ─── DAG ────────────────────────────────────────────────────────────────────
262
263pub async fn dag_put<T: Serialize>(kubo_url: &str, value: &T) -> Result<String> {
264    let base = kubo_url.trim_end_matches('/');
265    let url = format!("{base}/api/v0/dag/put");
266    let payload = serde_json::to_vec(value)?;
267
268    let part = multipart::Part::bytes(payload)
269        .file_name("node.json")
270        .mime_str("application/json")?;
271    let form = multipart::Form::new().part("file", part);
272
273    let client = reqwest::Client::builder()
274        .timeout(Duration::from_secs(10))
275        .build()?;
276
277    let body = client
278        .post(url)
279        .query(&[
280            ("store-codec", "dag-cbor"),
281            ("input-codec", "dag-json"),
282            ("pin", "true"),
283        ])
284        .multipart(form)
285        .send()
286        .await?
287        .error_for_status()?
288        .text()
289        .await?;
290
291    let parsed: DagPutResponse = serde_json::from_str(&body)
292        .map_err(|e| anyhow!("failed parsing dag/put response: {} body={}", e, body))?;
293    parsed
294        .cid_upper
295        .or(parsed.cid)
296        .map(|c| c.slash)
297        .ok_or_else(|| anyhow!("missing CID in dag/put response: {}", body))
298}
299
300pub async fn dag_put_cbor(kubo_url: &str, data: Vec<u8>, pin: bool) -> Result<String> {
301    let base = kubo_url.trim_end_matches('/');
302    let url = format!("{base}/api/v0/dag/put");
303    let part = multipart::Part::bytes(data)
304        .file_name("document.cbor")
305        .mime_str("application/octet-stream")?;
306    let form = multipart::Form::new().part("file", part);
307
308    let client = reqwest::Client::builder()
309        .timeout(Duration::from_secs(10))
310        .build()?;
311    let body = client
312        .post(url)
313        .query(&[
314            ("store-codec", "dag-cbor"),
315            ("input-codec", "dag-cbor"),
316            ("pin", if pin { "true" } else { "false" }),
317        ])
318        .multipart(form)
319        .send()
320        .await?
321        .error_for_status()?
322        .text()
323        .await?;
324
325    let parsed: DagPutResponse = serde_json::from_str(&body)
326        .map_err(|e| anyhow!("failed parsing dag/put response: {} body={}", e, body))?;
327    parsed
328        .cid_upper
329        .or(parsed.cid)
330        .map(|c| c.slash)
331        .ok_or_else(|| anyhow!("missing CID in dag/put response: {}", body))
332}
333
334pub async fn dag_get<T: DeserializeOwned>(kubo_url: &str, cid: &str) -> Result<T> {
335    let base = kubo_url.trim_end_matches('/');
336    let url = format!("{base}/api/v0/dag/get");
337
338    let client = reqwest::Client::builder()
339        .timeout(Duration::from_secs(10))
340        .build()?;
341
342    let body = client
343        .post(url)
344        .query(&[("arg", cid)])
345        .send()
346        .await?
347        .error_for_status()?
348        .text()
349        .await?;
350
351    serde_json::from_str::<T>(&body).map_err(|e| {
352        anyhow!(
353            "failed parsing dag/get response for {}: {} body={}",
354            cid,
355            e,
356            body
357        )
358    })
359}
360
361// ─── Name publish / resolve ─────────────────────────────────────────────────
362
363fn normalize_ipfs_arg(cid_or_path: &str) -> String {
364    let mut value = cid_or_path.trim().to_string();
365    while let Some(rest) = value.strip_prefix("/ipfs/") {
366        value = rest.to_string();
367    }
368    while let Some(rest) = value.strip_prefix('/') {
369        value = rest.to_string();
370    }
371    format!("/ipfs/{value}")
372}
373
374fn normalize_cid_arg(cid_or_path: &str) -> String {
375    let mut value = cid_or_path.trim().to_string();
376    while let Some(rest) = value.strip_prefix("/ipfs/") {
377        value = rest.to_string();
378    }
379    while let Some(rest) = value.strip_prefix('/') {
380        value = rest.to_string();
381    }
382    value
383}
384
385pub async fn name_publish(kubo_url: &str, key_name: &str, cid: &str) -> Result<String> {
386    let options = IpnsPublishOptions::default();
387    name_publish_with_options(kubo_url, key_name, cid, &options).await
388}
389
390pub async fn name_publish_with_options(
391    kubo_url: &str,
392    key_name: &str,
393    cid: &str,
394    options: &IpnsPublishOptions,
395) -> Result<String> {
396    let base = kubo_url.trim_end_matches('/');
397    let url = format!("{base}/api/v0/name/publish");
398    let arg = normalize_ipfs_arg(cid);
399
400    let client = reqwest::Client::builder()
401        .timeout(options.timeout)
402        .build()?;
403
404    let allow_offline = if options.allow_offline {
405        "true"
406    } else {
407        "false"
408    };
409    let resolve = if options.resolve { "true" } else { "false" };
410    let quieter = if options.quieter { "true" } else { "false" };
411
412    let mut params: Vec<(&str, &str)> = vec![
413        ("arg", arg.as_str()),
414        ("key", key_name),
415        ("allow-offline", allow_offline),
416        ("lifetime", options.lifetime.as_str()),
417        ("resolve", resolve),
418        ("quieter", quieter),
419    ];
420    if let Some(ref ttl) = options.ttl {
421        params.push(("ttl", ttl.as_str()));
422    }
423
424    let body = client
425        .post(url)
426        .query(&params)
427        .send()
428        .await?
429        .error_for_status()?
430        .text()
431        .await?;
432
433    let parsed: NamePublishResponse = serde_json::from_str(&body)
434        .map_err(|e| anyhow!("failed parsing name/publish response: {} body={}", e, body))?;
435    let value = if !parsed.value_upper.is_empty() {
436        parsed.value_upper
437    } else {
438        parsed.value_lower
439    };
440    if value.is_empty() {
441        return Err(anyhow!("missing value in name/publish response: {}", body));
442    }
443    Ok(value)
444}
445
446pub async fn name_publish_with_retry(
447    kubo_url: &str,
448    key_name: &str,
449    ipns_id: &str,
450    cid: &str,
451    options: &IpnsPublishOptions,
452    attempts: u32,
453    initial_backoff: Duration,
454) -> Result<String> {
455    if attempts == 0 {
456        return Err(anyhow!("name publish attempts must be >= 1"));
457    }
458
459    let mut fib_prev = Duration::from_millis(0);
460    let mut fib_curr = initial_backoff;
461    let mut last_err: Option<anyhow::Error> = None;
462
463    for attempt in 1..=attempts {
464        match name_publish_with_options(kubo_url, key_name, cid, options).await {
465            Ok(value) => return Ok(value),
466            Err(err) => {
467                if let Ok(value) = verify_name_target_after_error(kubo_url, ipns_id, cid).await {
468                    warn!(
469                        "name publish attempt {}/{} reported error for key '{}' but resolve confirms target; accepting: {}",
470                        attempt, attempts, key_name, value
471                    );
472                    return Ok(value);
473                }
474                warn!(
475                    "name publish attempt {}/{} failed for key '{}' cid '{}': {}",
476                    attempt, attempts, key_name, cid, err
477                );
478                last_err = Some(err);
479                if attempt < attempts {
480                    sleep(fib_curr).await;
481                    let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
482                    fib_prev = fib_curr;
483                    fib_curr = Duration::from_millis(std::cmp::min(next_ms, 30_000) as u64);
484                }
485            }
486        }
487    }
488
489    Err(anyhow!(
490        "name publish failed after {} attempt(s): {}",
491        attempts,
492        last_err
493            .map(|e| e.to_string())
494            .unwrap_or_else(|| "unknown error".to_string())
495    ))
496}
497
498async fn verify_name_target_after_error(
499    kubo_url: &str,
500    ipns_id: &str,
501    cid: &str,
502) -> Result<String> {
503    let expected = normalize_ipfs_arg(cid);
504    let resolved = name_resolve(kubo_url, &format!("/ipns/{ipns_id}"), true).await?;
505    if resolved.trim() == expected {
506        return Ok(resolved);
507    }
508    Err(anyhow!(
509        "post-error resolve mismatch for IPNS id '{}': expected '{}' got '{}'",
510        ipns_id,
511        expected,
512        resolved
513    ))
514}
515
516pub async fn name_resolve(kubo_url: &str, path: &str, recursive: bool) -> Result<String> {
517    let base = kubo_url.trim_end_matches('/');
518    let url = format!("{base}/api/v0/name/resolve");
519
520    let client = reqwest::Client::builder()
521        .timeout(Duration::from_secs(15))
522        .build()?;
523
524    let recursive_flag = if recursive { "true" } else { "false" };
525    let body = client
526        .post(url)
527        .query(&[("arg", path), ("recursive", recursive_flag)])
528        .send()
529        .await?
530        .error_for_status()?
531        .text()
532        .await?;
533
534    let parsed: NameResolveResponse = serde_json::from_str(&body)
535        .map_err(|e| anyhow!("failed parsing name/resolve response: {} body={}", e, body))?;
536    let resolved = if !parsed.path_upper.is_empty() {
537        parsed.path_upper
538    } else {
539        parsed.path_lower
540    };
541    if resolved.is_empty() {
542        return Err(anyhow!("missing path in name/resolve response: {}", body));
543    }
544    Ok(resolved)
545}
546
547// ─── DID document fetch ─────────────────────────────────────────────────────
548
549pub async fn fetch_did_document(kubo_url: &str, did: &Did) -> Result<Document> {
550    let ipns_path = format!("/ipns/{}", did.ipns);
551    let mut fib_prev = Duration::from_millis(0);
552    let mut fib_curr = Duration::from_millis(150);
553    let mut last_err: Option<anyhow::Error> = None;
554    let mut document: Option<Document> = None;
555
556    for attempt in 1..=4 {
557        // DID documents are stored as DAG-CBOR via dag/put.
558        let dag_err = match dag_get::<Document>(kubo_url, &ipns_path).await {
559            Ok(doc) => {
560                document = Some(doc);
561                break;
562            }
563            Err(e) => e,
564        };
565
566        // Fallback: resolve IPNS manually then dag_get the CID.
567        match name_resolve(kubo_url, &ipns_path, true).await {
568            Err(resolve_err) => {
569                last_err = Some(anyhow!(
570                    "dag_get and name/resolve both failed for {}: dag={} resolve={}",
571                    ipns_path,
572                    dag_err,
573                    resolve_err
574                ));
575                if !should_retry_name_resolve_error(&resolve_err) {
576                    break;
577                }
578            }
579            Ok(resolved_path) => match dag_get::<Document>(kubo_url, &resolved_path).await {
580                Ok(doc) => {
581                    document = Some(doc);
582                    break;
583                }
584                Err(err) => {
585                    last_err = Some(anyhow!(
586                        "dag_get failed for {}: direct={} resolved={}",
587                        ipns_path,
588                        dag_err,
589                        err
590                    ));
591                }
592            },
593        }
594
595        if attempt < 4 {
596            sleep(fib_curr).await;
597            let next_ms = fib_prev.as_millis().saturating_add(fib_curr.as_millis());
598            fib_prev = fib_curr;
599            fib_curr = Duration::from_millis(std::cmp::min(next_ms, 2_000) as u64);
600        }
601    }
602
603    let document = document.ok_or_else(|| {
604        anyhow!(
605            "failed to fetch DID document for {} via {} after retries: {}",
606            did.id(),
607            ipns_path,
608            last_err
609                .map(|e| e.to_string())
610                .unwrap_or_else(|| "unknown error".to_string())
611        )
612    })?;
613
614    document.validate()?;
615    document.verify()?;
616
617    let doc_did = Did::try_from(document.id.as_str())
618        .map_err(|e| anyhow!("DID document has invalid id '{}': {}", document.id, e))?;
619    if doc_did.ipns != did.ipns {
620        return Err(anyhow!(
621            "DID document IPNS mismatch: expected {} but document id is {}",
622            did.base_id(),
623            document.id
624        ));
625    }
626
627    Ok(document)
628}
629
630fn should_retry_name_resolve_error(err: &anyhow::Error) -> bool {
631    let text = err.to_string().to_ascii_lowercase();
632    if text.contains("http status client error") {
633        return false;
634    }
635    if text.contains("missing path in name/resolve response") {
636        return false;
637    }
638    true
639}
640
641// ─── Pin ────────────────────────────────────────────────────────────────────
642
643pub async fn pin_add_named(kubo_url: &str, cid: &str, name: &str) -> Result<()> {
644    let base = kubo_url.trim_end_matches('/');
645    let url = format!("{base}/api/v0/pin/add");
646    let arg = normalize_ipfs_arg(cid);
647
648    let client = reqwest::Client::builder()
649        .timeout(Duration::from_secs(10))
650        .build()?;
651
652    client
653        .post(url)
654        .query(&[("arg", arg.as_str()), ("recursive", "true"), ("name", name)])
655        .send()
656        .await?
657        .error_for_status()?;
658
659    Ok(())
660}
661
662pub async fn pin_rm(kubo_url: &str, cid: &str) -> Result<()> {
663    let base = kubo_url.trim_end_matches('/');
664    let url = format!("{base}/api/v0/pin/rm");
665    let arg = normalize_ipfs_arg(cid);
666
667    let client = reqwest::Client::builder()
668        .timeout(Duration::from_secs(10))
669        .build()?;
670
671    client
672        .post(url)
673        .query(&[("arg", arg.as_str()), ("recursive", "true")])
674        .send()
675        .await?
676        .error_for_status()?;
677
678    Ok(())
679}
680
681pub async fn list_named_recursive_pins(kubo_url: &str, name: &str) -> Result<Vec<String>> {
682    let base = kubo_url.trim_end_matches('/');
683    let url = format!("{base}/api/v0/pin/ls");
684    let body = reqwest::Client::builder()
685        .timeout(Duration::from_secs(30))
686        .build()?
687        .post(url)
688        .query(&[("type", "recursive"), ("name", name), ("names", "true")])
689        .send()
690        .await?
691        .error_for_status()?
692        .text()
693        .await?;
694    let parsed: PinListResponse = serde_json::from_str(&body)
695        .map_err(|error| anyhow!("failed parsing pin/ls response: {error} body={body}"))?;
696
697    // Kubo's name filter is a partial match; keep exact matches only.
698    Ok(parsed
699        .keys
700        .into_iter()
701        .filter_map(|(cid, pin)| (pin.pin_type == "recursive" && pin.name == name).then_some(cid))
702        .collect())
703}
704
705pub async fn remote_pin_add_named(
706    kubo_url: &str,
707    service: &str,
708    cid: &str,
709    name: &str,
710) -> Result<()> {
711    let base = kubo_url.trim_end_matches('/');
712    let url = format!("{base}/api/v0/pin/remote/add");
713    let arg = normalize_ipfs_arg(cid);
714
715    let client = reqwest::Client::builder()
716        .timeout(Duration::from_secs(30))
717        .build()?;
718
719    let resp = client
720        .post(url)
721        .query(&[
722            ("arg", arg.as_str()),
723            ("service", service),
724            ("name", name),
725            ("background", "false"),
726        ])
727        .send()
728        .await?;
729    if resp.status().is_success() {
730        return Ok(());
731    }
732    let status = resp.status();
733    let body = resp.text().await.unwrap_or_default();
734    if status.as_u16() == 409
735        || body.contains("DUPLICATE_OBJECT")
736        || body.contains("already pinned")
737        || body.contains("already exists")
738    {
739        return Ok(());
740    }
741    Err(anyhow!(
742        "pin/remote/add {cid} to {service} as {name} failed: {body}"
743    ))
744}
745
746pub async fn remote_pin_rm(kubo_url: &str, service: &str, cid: &str) -> Result<()> {
747    remote_pin_rm_query(kubo_url, service, cid, &[("force", "true")]).await
748}
749
750pub async fn remote_pin_rm_named(
751    kubo_url: &str,
752    service: &str,
753    cid: &str,
754    name: &str,
755) -> Result<()> {
756    remote_pin_rm_query(kubo_url, service, cid, &[("name", name), ("force", "true")]).await
757}
758
759// Kubo defaults pin/remote/ls and pin/remote/rm to status=pinned; include the
760// transient states so stale queued/failed requests are visible and removable.
761const REMOTE_PIN_STATUSES: [(&str, &str); 4] = [
762    ("status", "queued"),
763    ("status", "pinning"),
764    ("status", "pinned"),
765    ("status", "failed"),
766];
767
768async fn remote_pin_rm_query(
769    kubo_url: &str,
770    service: &str,
771    cid: &str,
772    extra: &[(&str, &str)],
773) -> Result<()> {
774    let base = kubo_url.trim_end_matches('/');
775    let url = format!("{base}/api/v0/pin/remote/rm");
776    let arg = normalize_cid_arg(cid);
777
778    let client = reqwest::Client::builder()
779        .timeout(Duration::from_secs(30))
780        .build()?;
781
782    let mut query: Vec<(&str, &str)> = vec![("service", service), ("cid", arg.as_str())];
783    query.extend_from_slice(&REMOTE_PIN_STATUSES);
784    query.extend_from_slice(extra);
785    let resp = client.post(url).query(&query).send().await?;
786    if resp.status().is_success() {
787        return Ok(());
788    }
789    let status = resp.status();
790    let body = resp.text().await.unwrap_or_default();
791    if status.as_u16() == 404 || body.contains("not found") || body.contains("not pinned") {
792        return Ok(());
793    }
794    Err(anyhow!("pin/remote/rm {cid} from {service} failed: {body}"))
795}
796
797pub async fn list_named_remote_pins(
798    kubo_url: &str,
799    service: &str,
800    name: &str,
801) -> Result<Vec<String>> {
802    let base = kubo_url.trim_end_matches('/');
803    let url = format!("{base}/api/v0/pin/remote/ls");
804    let mut query: Vec<(&str, &str)> = vec![("service", service), ("name", name)];
805    query.extend_from_slice(&REMOTE_PIN_STATUSES);
806    let body = reqwest::Client::builder()
807        .timeout(Duration::from_secs(30))
808        .build()?
809        .post(url)
810        .query(&query)
811        .send()
812        .await?
813        .error_for_status()?
814        .text()
815        .await?;
816
817    // pin/remote/ls streams NDJSON: one JSON record per pin, per line.
818    let mut cids = Vec::new();
819    for line in body.lines().map(str::trim).filter(|l| !l.is_empty()) {
820        let pin: RemotePinListEntry = serde_json::from_str(line).map_err(|error| {
821            anyhow!("failed parsing pin/remote/ls response line: {error} line={line}")
822        })?;
823        let pin_name = if pin.name_upper.is_empty() {
824            pin.name_lower
825        } else {
826            pin.name_upper
827        };
828        let cid = if pin.cid_upper.is_empty() {
829            pin.cid_lower
830        } else {
831            pin.cid_upper
832        };
833        if pin_name == name && !cid.is_empty() {
834            cids.push(cid);
835        }
836    }
837    Ok(cids)
838}
839
840// ─── Key management ─────────────────────────────────────────────────────────
841
842pub async fn generate_key(kubo_url: &str, key_name: &str) -> Result<()> {
843    let base = kubo_url.trim_end_matches('/');
844    let url = format!("{base}/api/v0/key/gen");
845
846    reqwest::Client::builder()
847        .timeout(Duration::from_secs(10))
848        .build()?
849        .post(url)
850        .query(&[("arg", key_name), ("type", "ed25519")])
851        .send()
852        .await?
853        .error_for_status()?;
854
855    Ok(())
856}
857
858pub async fn import_key(kubo_url: &str, key_name: &str, key_bytes: Vec<u8>) -> Result<KuboKey> {
859    let base = kubo_url.trim_end_matches('/');
860    let url = format!("{base}/api/v0/key/import");
861
862    let part = multipart::Part::bytes(key_bytes)
863        .file_name("ipns.key")
864        .mime_str("application/octet-stream")?;
865    let form = multipart::Form::new().part("file", part);
866
867    let response = reqwest::Client::builder()
868        .timeout(Duration::from_secs(10))
869        .build()?
870        .post(url)
871        .query(&[
872            ("arg", key_name),
873            ("ipns-base", "base36"),
874            ("allow-any-key-type", "true"),
875        ])
876        .multipart(form)
877        .send()
878        .await?
879        .error_for_status()?;
880
881    let body = response.text().await?;
882    let parsed: KeyImportResponse = serde_json::from_str(&body)
883        .map_err(|e| anyhow!("failed parsing key/import response: {} body={}", e, body))?;
884
885    let name = if !parsed.name_upper.trim().is_empty() {
886        parsed.name_upper.trim().to_string()
887    } else {
888        parsed.name_lower.trim().to_string()
889    };
890    let id = if !parsed.id_upper.trim().is_empty() {
891        parsed.id_upper.trim().to_string()
892    } else {
893        parsed.id_lower.trim().to_string()
894    };
895
896    if name.is_empty() || id.is_empty() {
897        return Err(anyhow!("missing name/id in key/import response: {}", body));
898    }
899
900    Ok(KuboKey { name, id })
901}
902
903pub async fn list_keys(kubo_url: &str) -> Result<Vec<KuboKey>> {
904    let base = kubo_url.trim_end_matches('/');
905    let url = format!("{base}/api/v0/key/list");
906
907    let body = reqwest::Client::builder()
908        .timeout(Duration::from_secs(10))
909        .build()?
910        .post(url)
911        .send()
912        .await?
913        .error_for_status()?
914        .text()
915        .await?;
916
917    let parsed: KeyListResponse = serde_json::from_str(&body)
918        .map_err(|e| anyhow!("failed parsing key/list response: {} body={}", e, body))?;
919    Ok(parsed
920        .keys
921        .into_iter()
922        .filter_map(|k| {
923            let name = if !k.name.trim().is_empty() {
924                k.name.trim().to_string()
925            } else {
926                k.name_lower.trim().to_string()
927            };
928            let id = if !k.id.trim().is_empty() {
929                k.id.trim().to_string()
930            } else {
931                k.id_lower.trim().to_string()
932            };
933            if name.is_empty() {
934                None
935            } else {
936                Some(KuboKey { name, id })
937            }
938        })
939        .collect())
940}
941
942pub async fn list_key_names(kubo_url: &str) -> Result<Vec<String>> {
943    let keys = list_keys(kubo_url).await?;
944    Ok(keys.into_iter().map(|k| k.name).collect())
945}
946
947/// Remove a named key from the Kubo keystore.
948pub async fn remove_key(kubo_url: &str, key_name: &str) -> Result<()> {
949    let base = kubo_url.trim_end_matches('/');
950    let url = format!("{base}/api/v0/key/rm");
951
952    reqwest::Client::builder()
953        .timeout(Duration::from_secs(10))
954        .build()?
955        .post(url)
956        .query(&[("arg", key_name)])
957        .send()
958        .await?
959        .error_for_status()?;
960
961    Ok(())
962}
963
964// ─── Tests ──────────────────────────────────────────────────────────────────
965
966#[cfg(test)]
967mod tests {
968    use std::io::{Read, Write};
969    use std::net::TcpListener;
970    use std::thread;
971
972    use super::*;
973
974    fn read_http_request(stream: &mut std::net::TcpStream) -> Vec<u8> {
975        let mut request = Vec::new();
976        let mut buffer = [0_u8; 4096];
977        let mut expected_length = None;
978
979        loop {
980            let count = stream.read(&mut buffer).expect("read request");
981            assert_ne!(count, 0, "request ended before its body arrived");
982            request.extend_from_slice(&buffer[..count]);
983
984            if expected_length.is_none() {
985                if let Some(header_end) = request.windows(4).position(|part| part == b"\r\n\r\n") {
986                    let headers = String::from_utf8_lossy(&request[..header_end]);
987                    expected_length = headers.lines().find_map(|line| {
988                        line.split_once(':').and_then(|(name, value)| {
989                            name.eq_ignore_ascii_case("content-length")
990                                .then(|| value.trim().parse::<usize>().expect("content length"))
991                        })
992                    });
993                    assert!(expected_length.is_some(), "request must have a body");
994                }
995            }
996
997            if let Some(length) = expected_length {
998                let header_end = request
999                    .windows(4)
1000                    .position(|part| part == b"\r\n\r\n")
1001                    .expect("headers parsed")
1002                    + 4;
1003                if request.len() >= header_end + length {
1004                    return request;
1005                }
1006            }
1007        }
1008    }
1009
1010    #[tokio::test]
1011    async fn dag_put_cbor_preserves_bytes_without_anonymous_pin() {
1012        let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind Kubo mock");
1013        let address = listener.local_addr().expect("mock address");
1014        let server = thread::spawn(move || {
1015            let (mut stream, _) = listener.accept().expect("accept Kubo request");
1016            let request = read_http_request(&mut stream);
1017            let body = r#"{"Cid":{"/":"bafy-local-pin"}}"#;
1018            let response = format!(
1019                "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
1020                body.len()
1021            );
1022            stream
1023                .write_all(response.as_bytes())
1024                .expect("write Kubo response");
1025            request
1026        });
1027
1028        let document = vec![0xd8, 0x2a, 0x43, 0x01, 0x02, 0x03];
1029        let cid = dag_put_cbor(&format!("http://{address}"), document.clone(), false)
1030            .await
1031            .expect("Kubo dag put");
1032
1033        assert_eq!(cid, "bafy-local-pin");
1034        let request = server.join().expect("Kubo mock thread");
1035        let request_text = String::from_utf8_lossy(&request);
1036        assert!(request_text.starts_with(
1037            "POST /api/v0/dag/put?store-codec=dag-cbor&input-codec=dag-cbor&pin=false HTTP/1.1"
1038        ));
1039        assert!(request_text
1040            .to_ascii_lowercase()
1041            .contains("content-type: multipart/form-data;"));
1042        assert!(request.windows(document.len()).any(|part| part == document));
1043    }
1044
1045    #[test]
1046    fn normalize_ipfs_arg_from_raw_cid() {
1047        assert_eq!(
1048            normalize_ipfs_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
1049            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1050        );
1051    }
1052
1053    #[test]
1054    fn normalize_ipfs_arg_from_prefixed_path() {
1055        assert_eq!(
1056            normalize_ipfs_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
1057            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1058        );
1059    }
1060
1061    #[test]
1062    fn normalize_ipfs_arg_from_double_prefixed_path() {
1063        assert_eq!(
1064            normalize_ipfs_arg(
1065                "/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1066            ),
1067            "/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1068        );
1069    }
1070
1071    #[test]
1072    fn normalize_cid_arg_from_raw_cid() {
1073        assert_eq!(
1074            normalize_cid_arg("bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
1075            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1076        );
1077    }
1078
1079    #[test]
1080    fn normalize_cid_arg_from_prefixed_path() {
1081        assert_eq!(
1082            normalize_cid_arg("/ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"),
1083            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1084        );
1085    }
1086
1087    #[test]
1088    fn normalize_cid_arg_from_double_prefixed_path() {
1089        assert_eq!(
1090            normalize_cid_arg(
1091                "/ipfs//ipfs/bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1092            ),
1093            "bafkreieg3hp4tr3iv7rj24uohuqlcocojxcakzgbvqh6ul2cpbftg3wb7q"
1094        );
1095    }
1096
1097    #[test]
1098    fn does_not_retry_http_client_status_errors() {
1099        let err = anyhow!(
1100            "HTTP status client error (404 Not Found) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
1101        );
1102        assert!(!should_retry_name_resolve_error(&err));
1103    }
1104
1105    #[test]
1106    fn retries_http_server_status_errors() {
1107        let err = anyhow!(
1108            "HTTP status server error (500 Internal Server Error) for url (http://127.0.0.1:5001/api/v0/name/resolve)"
1109        );
1110        assert!(should_retry_name_resolve_error(&err));
1111    }
1112
1113    #[test]
1114    fn retries_network_errors() {
1115        let err =
1116            anyhow!("error sending request for url (http://127.0.0.1:5001/api/v0/name/resolve)");
1117        assert!(should_retry_name_resolve_error(&err));
1118    }
1119}