1use 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#[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#[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
150pub 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
208pub 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
261pub 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
361fn 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(¶ms)
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
547pub 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 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 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
641pub 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 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
759const 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 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
840pub 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
947pub 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#[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}