diff --git a/Cargo.lock b/Cargo.lock index 169f668d..7e937e07 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -30,6 +30,12 @@ dependencies = [ "gimli", ] +[[package]] +name = "adler" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f26201604c87b1e01bd3d98f8d5d9a8fcbb815e8cedb41ffccbeb4bf593a35fe" + [[package]] name = "adler2" version = "2.0.1" @@ -57,6 +63,15 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "250f629c0161ad8107cf89319e990051fae62832fd343083bea452d93e2205fd" +[[package]] +name = "aligned-vec" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc890384c8602f339876ded803c97ad529f3842aba97f6392b3dba0dd171769b" +dependencies = [ + "equator", +] + [[package]] name = "alloc-no-stdlib" version = "2.0.4" @@ -81,6 +96,15 @@ dependencies = [ "libc", ] +[[package]] +name = "ansi_colours" +version = "1.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "14eec43e0298190790f41679fe69ef7a829d2a2ddd78c8c00339e84710e435fe" +dependencies = [ + "rgb", +] + [[package]] name = "anstream" version = "0.6.20" @@ -137,6 +161,29 @@ version = "1.0.100" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61" +[[package]] +name = "arbitrary" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d036a3c4ab069c7b410a2ce876bd74808d2d0888a82667669f8e783a898bf1" + +[[package]] +name = "arg_enum_proc_macro" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ae92a5119aa49cdbcf6b9f893fe4e1d98b04ccbf82ee0584ad948a44a734dea" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + +[[package]] +name = "arrayvec" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50" + [[package]] name = "ascii" version = "1.1.0" @@ -179,6 +226,29 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" +[[package]] +name = "av1-grain" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4f3efb2ca85bc610acfa917b5aaa36f3fcbebed5b3182d7f877b02531c4b80c8" +dependencies = [ + "anyhow", + "arrayvec", + "log", + "nom", + "num-rational", + "v_frame", +] + +[[package]] +name = "avif-serialize" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47c8fbc0f831f4519fe8b810b6a7a91410ec83031b8233f730a0480029f6a23f" +dependencies = [ + "arrayvec", +] + [[package]] name = "axum" version = "0.8.6" @@ -280,7 +350,7 @@ dependencies = [ "addr2line", "cfg-if", "libc", - "miniz_oxide", + "miniz_oxide 0.8.9", "object", "rustc-demangle", "windows-link 0.2.0", @@ -325,12 +395,24 @@ version = "1.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "55248b47b0caf0546f7988906588779981c43bb1bc9d0c44087278f80cdb44ba" +[[package]] +name = "bit_field" +version = "0.10.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e4b40c7323adcfc0a41c4b88143ed58346ff65a288fc144329c5c45e05d70c6" + [[package]] name = "bitflags" version = "2.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2261d10cca569e4643e526d8dc2e62e433cc8aba21ab764233731f8d369bf394" +[[package]] +name = "bitstream-io" +version = "2.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6099cdc01846bc367c4e7dd630dc5966dccf36b652fae7a74e17b640411a91b2" + [[package]] name = "block-buffer" version = "0.10.4" @@ -429,18 +511,36 @@ dependencies = [ "safemem", ] +[[package]] +name = "built" +version = "0.7.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56ed6191a7e78c36abdb16ab65341eefd73d64d303fffccdbb00d51e4205967b" + [[package]] name = "bumpalo" version = "3.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43" +[[package]] +name = "bytemuck" +version = "1.24.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fbdf580320f38b612e485521afda1ee26d10cc9884efaaa750d383e13e3c5f4" + [[package]] name = "byteorder" version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" +[[package]] +name = "byteorder-lite" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f1fe948ff07f4bd06c30984e69f5b4899c516a3ef74f34df92a2df2ab535495" + [[package]] name = "bytes" version = "1.10.1" @@ -472,6 +572,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e1354349954c6fc9cb0deab020f27f783cf0b604e8bb754dc4658ecf0d29c35f" dependencies = [ "find-msvc-tools", + "jobserver", + "libc", "shlex", ] @@ -490,6 +592,16 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6d43a04d8753f35258c91f8ec639f792891f748a1edbd759cf1dcea3382ad83c" +[[package]] +name = "cfg-expr" +version = "0.15.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d067ad48b8650848b989a59a86c6c36a995d02d2bf778d45c3c5d57bc2718f02" +dependencies = [ + "smallvec", + "target-lexicon", +] + [[package]] name = "cfg-if" version = "1.0.3" @@ -603,6 +715,12 @@ version = "0.7.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b94f61472cee1439c0b966b47e3aca9ae07e45d070759512cd390ea2bebc6675" +[[package]] +name = "color_quant" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3d7b894f5411737b7867f4827955924d7c254fc9f4d91a6aad6b097804b1018b" + [[package]] name = "colorchoice" version = "1.0.4" @@ -636,6 +754,18 @@ version = "0.4.29" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e47641d3deaf41fb1538ac1f54735925e275eaf3bf4d55c81b137fba797e5cbb" +[[package]] +name = "console" +version = "0.15.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "054ccb5b10f9f2cbf51eb355ca1d05c2d279ce1804688d0db74b4733a5aeafd8" +dependencies = [ + "encode_unicode", + "libc", + "once_cell", + "windows-sys 0.59.0", +] + [[package]] name = "const-oid" version = "0.9.6" @@ -705,12 +835,53 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "crossbeam-deque" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9dd111b7b7f7d55b72c0a6ae361660ee5853c9af73f70c3c2ef6858b950e2e51" +dependencies = [ + "crossbeam-epoch", + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-epoch" +version = "0.9.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b82ac4a3c2ca9c3460964f020e1402edd5753411d7737aa39c3714ad1b5420e" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-utils" version = "0.8.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" +[[package]] +name = "crossterm" +version = "0.28.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "829d955a0bb380ef178a640b91779e3987da38c9aea133b20614cfed8cdea9c6" +dependencies = [ + "bitflags", + "crossterm_winapi", + "parking_lot", + "rustix 0.38.44", + "winapi", +] + +[[package]] +name = "crossterm_winapi" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "acdd7c62a3665c7f6830a51635d9ac9b23ed385797f70a83bb8bafe9c572ab2b" +dependencies = [ + "winapi", +] + [[package]] name = "crunchy" version = "0.2.4" @@ -1004,6 +1175,12 @@ dependencies = [ "serde", ] +[[package]] +name = "encode_unicode" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34aa73646ffb006b8f5147f3dc182bd4bcb190227ce861fc4a4844bf8e3cb2c0" + [[package]] name = "encoding_rs" version = "0.8.35" @@ -1025,6 +1202,26 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "equator" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4711b213838dfee0117e3be6ac926007d7f433d7bbe33595975d4190cb07e6fc" +dependencies = [ + "equator-macro", +] + +[[package]] +name = "equator-macro" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "44f23cf4b44bfce11a86ace86f8a73ffdec849c9fd00a386a53d278bd9e81fb3" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + [[package]] name = "equivalent" version = "1.0.2" @@ -1080,12 +1277,56 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "exr" +version = "1.73.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f83197f59927b46c04a183a619b7c29df34e63e63c7869320862268c0ef687e0" +dependencies = [ + "bit_field", + "half", + "lebe", + "miniz_oxide 0.8.9", + "rayon-core", + "smallvec", + "zune-inflate", +] + [[package]] name = "fastrand" version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" +[[package]] +name = "fax" +version = "0.2.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f05de7d48f37cd6730705cbca900770cab77a89f413d23e100ad7fad7795a0ab" +dependencies = [ + "fax_derive", +] + +[[package]] +name = "fax_derive" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a0aca10fb742cb43f9e7bb8467c91aa9bcb8e3ffbc6a6f7389bb93ffc920577d" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + +[[package]] +name = "fdeflate" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e6853b52649d4ac5c0bd02320cddc5ba956bdb407c4b75a2c6b75bf51500f8c" +dependencies = [ + "simd-adler32", +] + [[package]] name = "ff" version = "0.13.1" @@ -1127,7 +1368,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4a3d7db9596fecd151c5f638c0ee5d5bd487b6e0ea232e5dc96d5250f6f94b1d" dependencies = [ "crc32fast", - "miniz_oxide", + "miniz_oxide 0.8.9", ] [[package]] @@ -1343,6 +1584,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "gif" +version = "0.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ae047235e33e2829703574b54fdec96bfbad892062d97fed2f76022287de61b" +dependencies = [ + "color_quant", + "weezl", +] + [[package]] name = "gimli" version = "0.32.3" @@ -1764,6 +2015,46 @@ dependencies = [ "icu_properties", ] +[[package]] +name = "image" +version = "0.25.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "529feb3e6769d234375c4cf1ee2ce713682b8e76538cb13f9fc23e1400a591e7" +dependencies = [ + "bytemuck", + "byteorder-lite", + "color_quant", + "exr", + "gif", + "image-webp", + "moxcms", + "num-traits", + "png", + "qoi", + "ravif", + "rayon", + "rgb", + "tiff 0.10.3", + "zune-core", + "zune-jpeg", +] + +[[package]] +name = "image-webp" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "525e9ff3e1a4be2fbea1fdf0e98686a6d98b4d8f937e1bf7402245af1909e8c3" +dependencies = [ + "byteorder-lite", + "quick-error 2.0.1", +] + +[[package]] +name = "imgref" +version = "1.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7c5cedc30da3a610cac6b4ba17597bdf7152cf974e8aab3afb3d54455e371c8" + [[package]] name = "indexmap" version = "1.9.3" @@ -1793,6 +2084,17 @@ version = "2.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f4c7245a08504955605670dbf141fceab975f15ca21570696aebe9d2e71576bd" +[[package]] +name = "interpolate_name" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c34819042dc3d3971c46c2190835914dfbe0c3c13f61449b2997f4e9722dfa60" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + [[package]] name = "inventory" version = "0.3.21" @@ -1864,6 +2166,15 @@ version = "1.70.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7943c866cc5cd64cbc25b2e01621d07fa8eb2a1a23160ee81ce38704e97b8ecf" +[[package]] +name = "itertools" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba291022dbbd398a455acf126c1e341954079855bc60dfdda641363bd6922569" +dependencies = [ + "either", +] + [[package]] name = "itertools" version = "0.14.0" @@ -1886,8 +2197,10 @@ dependencies = [ "bon", "bytes", "clap", + "futures", "getrandom 0.2.16", "http", + "image", "jacquard-api 0.5.5", "jacquard-common 0.5.4", "jacquard-derive 0.5.4", @@ -1895,6 +2208,7 @@ dependencies = [ "jacquard-oauth", "jose-jwk", "miette", + "n0-future", "p256", "percent-encoding", "rand_core 0.6.4", @@ -1905,10 +2219,12 @@ dependencies = [ "serde_json", "smol_str", "thiserror 2.0.17", + "tiff 0.6.1", "tokio", "tracing", "trait-variant", "url", + "viuer", ] [[package]] @@ -2050,7 +2366,7 @@ version = "0.5.1" source = "git+https://tangled.org/@nonbinary.computer/jacquard#77915fd4920b282b4b1342749dcdad9dce30cadf" dependencies = [ "heck 0.5.0", - "itertools", + "itertools 0.14.0", "prettyplease", "proc-macro2", "quote", @@ -2107,6 +2423,7 @@ dependencies = [ "jacquard-api 0.5.5", "jacquard-common 0.5.4", "miette", + "n0-future", "percent-encoding", "reqwest", "serde", @@ -2163,6 +2480,7 @@ dependencies = [ "jose-jwa", "jose-jwk", "miette", + "n0-future", "p256", "rand 0.8.5", "rand_core 0.6.4", @@ -2204,6 +2522,16 @@ version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8eaf4bc02d17cbdd7ff4c7438cafcdf7fb9a4613313ad11b4f8fefe7d3fa0130" +[[package]] +name = "jobserver" +version = "0.1.34" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9afb3de4395d6b3e67a780b6de64b51c978ecf11cb9a462c66be7d4ca9039d33" +dependencies = [ + "getrandom 0.3.4", + "libc", +] + [[package]] name = "jose-b64" version = "0.1.2" @@ -2240,6 +2568,12 @@ dependencies = [ "zeroize", ] +[[package]] +name = "jpeg-decoder" +version = "0.1.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "229d53d58899083193af11e15917b5640cd40b29ff475a1fe4ef725deb02d0f2" + [[package]] name = "js-sys" version = "0.3.81" @@ -2272,7 +2606,7 @@ checksum = "81a29e7b50079ff44549f68c0becb1c73d7f6de2a4ea952da77966daf3d4761e" dependencies = [ "miette", "num", - "winnow", + "winnow 0.6.24", ] [[package]] @@ -2295,12 +2629,28 @@ dependencies = [ "spin 0.9.8", ] +[[package]] +name = "lebe" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a79a3332a6609480d7d0c9eab957bca6b455b91bb84e66d19f5ff66294b85b8" + [[package]] name = "libc" version = "0.2.176" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "58f929b4d672ea937a23a1ab494143d968337a5f47e56d0815df1e0890ddf174" +[[package]] +name = "libfuzzer-sys" +version = "0.4.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5037190e1f70cbeef565bd267599242926f724d3b8a9f510fd7e0b540cfa4404" +dependencies = [ + "arbitrary", + "cc", +] + [[package]] name = "libm" version = "0.2.15" @@ -2324,6 +2674,12 @@ version = "0.5.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0717cef1bc8b636c6e1c1bbdefc09e6322da8a9321966e8928ef80d20f7f770f" +[[package]] +name = "linux-raw-sys" +version = "0.4.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d26c52dbd32dccf2d10cac7725f8eae5296885fb5703b261f7d0a0739ec807ab" + [[package]] name = "linux-raw-sys" version = "0.11.0" @@ -2364,6 +2720,15 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "loop9" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fae87c125b03c1d2c0150c90365d7d6bcc53fb73a9acaef207d2d065860f062" +dependencies = [ + "imgref", +] + [[package]] name = "lru-cache" version = "0.1.2" @@ -2379,6 +2744,12 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "make-cmd" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8ca8afbe8af1785e09636acb5a41e08a765f5f0340568716c18a8700ba3c0d3" + [[package]] name = "malloc_buf" version = "0.0.6" @@ -2403,6 +2774,16 @@ version = "0.8.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3" +[[package]] +name = "maybe-rayon" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ea1f30cedd69f0a2954655f7188c6a834246d2bcf1e315e2ac40c4b24dc9519" +dependencies = [ + "cfg-if", + "rayon", +] + [[package]] name = "memchr" version = "2.7.6" @@ -2461,6 +2842,16 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "68354c5c6bd36d73ff3feceb05efa59b6acb7626617f4962be322a825e61f79a" +[[package]] +name = "miniz_oxide" +version = "0.4.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a92518e98c078586bc6c934028adcca4c92a53d6a958196de835170a01d84e4b" +dependencies = [ + "adler", + "autocfg", +] + [[package]] name = "miniz_oxide" version = "0.8.9" @@ -2468,6 +2859,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1fa76a2c86f704bdb222d66965fb3d63269ce38518b83cb0575fca855ebb6316" dependencies = [ "adler2", + "simd-adler32", ] [[package]] @@ -2481,6 +2873,16 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "moxcms" +version = "0.7.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c588e11a3082784af229e23e8e4ecf5bcc6fbe4f69101e0421ce8d79da7f0b40" +dependencies = [ + "num-traits", + "pxfm", +] + [[package]] name = "multibase" version = "0.9.1" @@ -2514,7 +2916,7 @@ dependencies = [ "log", "mime", "mime_guess", - "quick-error", + "quick-error 1.2.3", "rand 0.8.5", "safemem", "tempfile", @@ -2548,6 +2950,12 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "27b02d87554356db9e9a873add8782d4ea6e3e58ea071a9adb9a2e8ddb884a8b" +[[package]] +name = "new_debug_unreachable" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" + [[package]] name = "nom" version = "7.1.3" @@ -2558,6 +2966,12 @@ dependencies = [ "minimal-lexical", ] +[[package]] +name = "noop_proc_macro" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0676bb32a98c1a483ce53e500a81ad9c3d5b3f7c920c28c24e9cb0980d0b5bc8" + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -2623,6 +3037,17 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "51d515d32fb182ee37cda2ccdcb92950d6a3c2893aa280e540671c2cd0f3b1d9" +[[package]] +name = "num-derive" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed3955f1a9c7c0c15e092f9c887db08b1fc683305fdf6eb6684f22555355e202" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + [[package]] name = "num-integer" version = "0.1.46" @@ -2794,6 +3219,12 @@ dependencies = [ "windows-link 0.2.0", ] +[[package]] +name = "paste" +version = "1.0.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "57c0d7b74b563b49d38dae00a0c37d4d6de9b432382b2892f0574ddcae73fd0a" + [[package]] name = "pem-rfc7468" version = "0.7.0" @@ -2862,6 +3293,25 @@ dependencies = [ "spki", ] +[[package]] +name = "pkg-config" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" + +[[package]] +name = "png" +version = "0.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97baced388464909d42d89643fe4361939af9b7ce7a31ee32a168f832a70f2a0" +dependencies = [ + "bitflags", + "crc32fast", + "fdeflate", + "flate2", + "miniz_oxide 0.8.9", +] + [[package]] name = "potential_utf" version = "0.1.3" @@ -2993,12 +3443,55 @@ dependencies = [ "yansi", ] +[[package]] +name = "profiling" +version = "1.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3eb8486b569e12e2c32ad3e204dbaba5e4b5b216e9367044f25f1dba42341773" +dependencies = [ + "profiling-procmacros", +] + +[[package]] +name = "profiling-procmacros" +version = "1.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52717f9a02b6965224f95ca2a81e2e0c5c43baacd28ca057577988930b6c3d5b" +dependencies = [ + "quote", + "syn 2.0.106", +] + +[[package]] +name = "pxfm" +version = "0.1.25" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a3cbdf373972bf78df4d3b518d07003938e2c7d1fb5891e55f9cb6df57009d84" +dependencies = [ + "num-traits", +] + +[[package]] +name = "qoi" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f6d64c71eb498fe9eae14ce4ec935c555749aef511cca85b5568910d6e48001" +dependencies = [ + "bytemuck", +] + [[package]] name = "quick-error" version = "1.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a1d01941d82fa2ab50be1e79e6714289dd7cde78eba4c074bc5a4374f650dfe0" +[[package]] +name = "quick-error" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a993555f31e5a609f617c12db6250dedcac1b0a85076912c436e6fc9b2c8e6a3" + [[package]] name = "quinn" version = "0.11.9" @@ -3134,12 +3627,82 @@ version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d20581732dd76fa913c7dff1a2412b714afe3573e94d41c34719de73337cc8ab" +[[package]] +name = "rav1e" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cd87ce80a7665b1cce111f8a16c1f3929f6547ce91ade6addf4ec86a8dda5ce9" +dependencies = [ + "arbitrary", + "arg_enum_proc_macro", + "arrayvec", + "av1-grain", + "bitstream-io", + "built", + "cfg-if", + "interpolate_name", + "itertools 0.12.1", + "libc", + "libfuzzer-sys", + "log", + "maybe-rayon", + "new_debug_unreachable", + "noop_proc_macro", + "num-derive", + "num-traits", + "once_cell", + "paste", + "profiling", + "rand 0.8.5", + "rand_chacha 0.3.1", + "simd_helpers", + "system-deps", + "thiserror 1.0.69", + "v_frame", + "wasm-bindgen", +] + +[[package]] +name = "ravif" +version = "0.11.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5825c26fddd16ab9f515930d49028a630efec172e903483c94796cfe31893e6b" +dependencies = [ + "avif-serialize", + "imgref", + "loop9", + "quick-error 2.0.1", + "rav1e", + "rayon", + "rgb", +] + [[package]] name = "raw-window-handle" version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2ff9a1f06a88b01621b7ae906ef0211290d1c8a168a15542486a8f61c0833b9" +[[package]] +name = "rayon" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "368f01d005bf8fd9b1206fb6fa653e6c4a81ceb1466406b81792d87c5677a58f" +dependencies = [ + "either", + "rayon-core", +] + +[[package]] +name = "rayon-core" +version = "1.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22e18b0f0062d30d4230b2e85ff77fdfe4326feb054b9783a3460d8435c8ab91" +dependencies = [ + "crossbeam-deque", + "crossbeam-utils", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -3268,6 +3831,15 @@ dependencies = [ "subtle", ] +[[package]] +name = "rgb" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c6a884d2998352bb4daf0183589aec883f16a6da1f4dde84d8e2e9a5409a1ce" +dependencies = [ + "bytemuck", +] + [[package]] name = "ring" version = "0.17.14" @@ -3362,6 +3934,19 @@ dependencies = [ "semver", ] +[[package]] +name = "rustix" +version = "0.38.44" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fdb5bc1ae2baa591800df16c9ca78619bf65c0488b41b96ccec5d11220d8c154" +dependencies = [ + "bitflags", + "errno", + "libc", + "linux-raw-sys 0.4.15", + "windows-sys 0.59.0", +] + [[package]] name = "rustix" version = "1.1.2" @@ -3371,7 +3956,7 @@ dependencies = [ "bitflags", "errno", "libc", - "linux-raw-sys", + "linux-raw-sys 0.11.0", "windows-sys 0.60.2", ] @@ -3598,6 +4183,15 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "serde_spanned" +version = "0.6.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf41e0cfaf7226dca15e8197172c295a782857fcb97fad1808a166870dee75a3" +dependencies = [ + "serde", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -3704,6 +4298,39 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "simd-adler32" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d66dc143e6b11c1eddc06d5c423cfc97062865baf299914ab64caa38182078fe" + +[[package]] +name = "simd_helpers" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "95890f873bec569a0362c235787f3aca6e1e887302ba4840839bcc6459c42da6" +dependencies = [ + "quote", +] + +[[package]] +name = "sixel-rs" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfa95c014543113a192d906e5971d0c8d1e8b4cc1e61026539687a7016644ce5" +dependencies = [ + "sixel-sys", +] + +[[package]] +name = "sixel-sys" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fb46e0cd5569bf910390844174a5a99d52dd40681fff92228d221d9f8bf87dea" +dependencies = [ + "make-cmd", +] + [[package]] name = "slab" version = "0.4.11" @@ -3907,6 +4534,25 @@ dependencies = [ "libc", ] +[[package]] +name = "system-deps" +version = "6.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a3e535eb8dded36d55ec13eddacd30dec501792ff23a0b1682c38601b8cf2349" +dependencies = [ + "cfg-expr", + "heck 0.5.0", + "pkg-config", + "toml", + "version-compare", +] + +[[package]] +name = "target-lexicon" +version = "0.12.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c41af27dd6d1e27b1b16b489db798443478cef1f06a660c96db617ba5de3b1" + [[package]] name = "tempfile" version = "3.23.0" @@ -3916,17 +4562,26 @@ dependencies = [ "fastrand", "getrandom 0.3.4", "once_cell", - "rustix", + "rustix 1.1.2", "windows-sys 0.60.2", ] +[[package]] +name = "termcolor" +version = "1.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06794f8f6c5c898b3275aebefa6b8a1cb24cd2c6c79397ab15774837a0bc5755" +dependencies = [ + "winapi-util", +] + [[package]] name = "terminal_size" version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "60b8cb979cb11c32ce1603f8137b22262a9d131aaa5c37b5678025f22b8becd0" dependencies = [ - "rustix", + "rustix 1.1.2", "windows-sys 0.60.2", ] @@ -3998,6 +4653,31 @@ dependencies = [ "num_cpus", ] +[[package]] +name = "tiff" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a53f4706d65497df0c4349241deddf35f84cee19c87ed86ea8ca590f4464437" +dependencies = [ + "jpeg-decoder", + "miniz_oxide 0.4.4", + "weezl", +] + +[[package]] +name = "tiff" +version = "0.10.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af9605de7fee8d9551863fd692cce7637f548dbd9db9180fcc07ccc6d26c336f" +dependencies = [ + "fax", + "flate2", + "half", + "quick-error 2.0.1", + "weezl", + "zune-jpeg", +] + [[package]] name = "time" version = "0.3.44" @@ -4153,6 +4833,40 @@ dependencies = [ "tokio", ] +[[package]] +name = "toml" +version = "0.8.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc1beb996b9d83529a9e75c17a1686767d148d70663143c7854d8b4a09ced362" +dependencies = [ + "serde", + "serde_spanned", + "toml_datetime", + "toml_edit", +] + +[[package]] +name = "toml_datetime" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22cddaf88f4fbc13c51aebbf5f8eceb5c7c5a9da2ac40a13519eb5b0a0e8f11c" +dependencies = [ + "serde", +] + +[[package]] +name = "toml_edit" +version = "0.22.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a" +dependencies = [ + "indexmap 2.11.4", + "serde", + "serde_spanned", + "toml_datetime", + "winnow 0.7.13", +] + [[package]] name = "tower" version = "0.5.2" @@ -4437,18 +5151,52 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "v_frame" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "666b7727c8875d6ab5db9533418d7c764233ac9c0cff1d469aec8fa127597be2" +dependencies = [ + "aligned-vec", + "num-traits", + "wasm-bindgen", +] + [[package]] name = "valuable" version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" +[[package]] +name = "version-compare" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "852e951cb7832cb45cb1169900d19760cfa39b82bc0ea9c0e5a14ae88411c98b" + [[package]] name = "version_check" version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "viuer" +version = "0.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ae7c6870b98c838123f22cac9a594cbe2d74ea48d79271c08f8c9e680b40fac" +dependencies = [ + "ansi_colours", + "base64 0.22.1", + "console", + "crossterm", + "image", + "lazy_static", + "sixel-rs", + "tempfile", + "termcolor", +] + [[package]] name = "walkdir" version = "2.5.0" @@ -4614,12 +5362,34 @@ dependencies = [ "rustls-pki-types", ] +[[package]] +name = "weezl" +version = "0.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a751b3277700db47d3e574514de2eced5e54dc8a5436a3bf7a0b248b2cee16f3" + [[package]] name = "widestring" version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dd7cf3379ca1aac9eea11fba24fd7e315d621f8dfe35c8d7d2be8b793726e07d" +[[package]] +name = "winapi" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c839a674fcd7a98952e593242ea400abe93992746761e38641405d28b00f419" +dependencies = [ + "winapi-i686-pc-windows-gnu", + "winapi-x86_64-pc-windows-gnu", +] + +[[package]] +name = "winapi-i686-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" + [[package]] name = "winapi-util" version = "0.1.11" @@ -4629,6 +5399,12 @@ dependencies = [ "windows-sys 0.60.2", ] +[[package]] +name = "winapi-x86_64-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" + [[package]] name = "windows" version = "0.61.3" @@ -5085,6 +5861,15 @@ dependencies = [ "memchr", ] +[[package]] +name = "winnow" +version = "0.7.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "21a0236b59786fed61e2a80582dd500fe61f18b5dca67a4a067d0bc9039339cf" +dependencies = [ + "memchr", +] + [[package]] name = "winreg" version = "0.50.0" @@ -5219,3 +6004,27 @@ dependencies = [ "quote", "syn 2.0.106", ] + +[[package]] +name = "zune-core" +version = "0.4.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f423a2c17029964870cfaabb1f13dfab7d092a62a29a89264f4d36990ca414a" + +[[package]] +name = "zune-inflate" +version = "0.2.54" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73ab332fe2f6680068f3582b16a24f90ad7096d5d39b974d1c0aff0125116f02" +dependencies = [ + "simd-adler32", +] + +[[package]] +name = "zune-jpeg" +version = "0.4.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29ce2c8a9384ad323cf564b67da86e21d3cfdff87908bc1223ed5c99bc792713" +dependencies = [ + "zune-core", +] diff --git a/Cargo.toml b/Cargo.toml index fb99291c..f4e7ba30 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -61,6 +61,7 @@ reqwest = { version = "0.12", default-features = false } # Async and runtimes tokio = { version = "1", default-features = false } +n0-future = "0.1" # Observability tracing = "0.1" diff --git a/crates/jacquard-common/Cargo.toml b/crates/jacquard-common/Cargo.toml index e155cc11..8cf4354a 100644 --- a/crates/jacquard-common/Cargo.toml +++ b/crates/jacquard-common/Cargo.toml @@ -42,7 +42,7 @@ tracing = { workspace = true, optional = true } tokio = { workspace = true, default-features = false, features = ["sync"] } # Streaming support (optional) -n0-future = { version = "0.1", optional = true } +n0-future = { workspace = true, optional = true } futures = { version = "0.3", optional = true } tokio-tungstenite-wasm = { version = "0.4", optional = true } genawaiter = { version = "0.99.1", features = ["futures03"] } diff --git a/crates/jacquard-common/src/stream.rs b/crates/jacquard-common/src/stream.rs index 4fe38658..9ab14cb4 100644 --- a/crates/jacquard-common/src/stream.rs +++ b/crates/jacquard-common/src/stream.rs @@ -49,9 +49,10 @@ use std::fmt; pub type BoxError = Box; /// Error type for streaming operations -#[derive(Debug)] +#[derive(Debug, thiserror::Error, miette::Diagnostic)] pub struct StreamError { kind: StreamErrorKind, + #[source] source: Option, } @@ -156,14 +157,6 @@ impl fmt::Display for StreamError { } } -impl Error for StreamError { - fn source(&self) -> Option<&(dyn Error + 'static)> { - self.source - .as_ref() - .map(|e| e.as_ref() as &(dyn Error + 'static)) - } -} - use bytes::Bytes; use n0_future::stream::Boxed; @@ -204,6 +197,48 @@ impl ByteStream { pub fn into_inner(self) -> Boxed> { self.inner } + + /// Split this stream into two streams that both receive all chunks + /// + /// Chunks are cloned (cheaply via Bytes rc). Spawns a forwarder task. + /// Both returned streams will receive all chunks from the original stream. + /// The forwarder continues as long as at least one stream is alive. + /// If the underlying stream errors, both teed streams will end. + pub fn tee(self) -> (ByteStream, ByteStream) { + use futures::channel::mpsc; + use n0_future::StreamExt as _; + + let (tx1, rx1) = mpsc::unbounded(); + let (tx2, rx2) = mpsc::unbounded(); + + n0_future::task::spawn(async move { + let mut stream = self.inner; + while let Some(result) = stream.next().await { + match result { + Ok(chunk) => { + // Clone chunk (cheap - Bytes is rc'd) + let chunk2 = chunk.clone(); + + // Send to both channels, continue if at least one succeeds + let send1 = tx1.unbounded_send(Ok(chunk)); + let send2 = tx2.unbounded_send(Ok(chunk2)); + + // Only stop if both channels are closed + if send1.is_err() && send2.is_err() { + break; + } + } + Err(_e) => { + // Underlying stream errored, stop forwarding. + // Both channels will close, ending both streams. + break; + } + } + } + }); + + (ByteStream::new(rx1), ByteStream::new(rx2)) + } } impl fmt::Debug for ByteStream { diff --git a/crates/jacquard-common/src/types/blob.rs b/crates/jacquard-common/src/types/blob.rs index ec5a8e1e..13ea162f 100644 --- a/crates/jacquard-common/src/types/blob.rs +++ b/crates/jacquard-common/src/types/blob.rs @@ -1,4 +1,7 @@ -use crate::{CowStr, IntoStatic, types::cid::Cid}; +use crate::{ + CowStr, IntoStatic, + types::cid::{Cid, CidLink}, +}; #[allow(unused)] use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error}; use smol_str::ToSmolStr; @@ -24,7 +27,7 @@ use std::{ #[serde(rename_all = "camelCase")] pub struct Blob<'b> { /// CID (Content Identifier) reference to the blob data - pub r#ref: Cid<'b>, + pub r#ref: CidLink<'b>, /// MIME type of the blob (e.g., "image/png", "video/mp4") #[serde(borrow)] pub mime_type: MimeType<'b>, diff --git a/crates/jacquard-common/src/types/value/convert.rs b/crates/jacquard-common/src/types/value/convert.rs index 778e0cae..99af866a 100644 --- a/crates/jacquard-common/src/types/value/convert.rs +++ b/crates/jacquard-common/src/types/value/convert.rs @@ -1,4 +1,5 @@ use crate::IntoStatic; +use crate::types::cid::CidLink; use crate::types::{ DataModelType, cid::Cid, @@ -298,7 +299,7 @@ impl<'s> TryFrom> for Data<'s> { } }; return Ok(Data::Blob(crate::types::blob::Blob { - r#ref: cid.clone(), + r#ref: CidLink::str(cid).into_static(), mime_type: crate::types::blob::MimeType::from(mime.clone()), size: size_val, })); diff --git a/crates/jacquard-common/src/types/value/parsing.rs b/crates/jacquard-common/src/types/value/parsing.rs index c50074fe..61eb6f59 100644 --- a/crates/jacquard-common/src/types/value/parsing.rs +++ b/crates/jacquard-common/src/types/value/parsing.rs @@ -251,7 +251,7 @@ pub fn cbor_to_blob<'b>(blob: &'b BTreeMap) -> Option> { }); if let (Some(mime_type), Some(size)) = (mime_type, size) { return Some(Blob { - r#ref: Cid::ipld(*value), + r#ref: CidLink::ipld(*value), mime_type: MimeType::raw(mime_type), size: size as usize, }); @@ -259,7 +259,7 @@ pub fn cbor_to_blob<'b>(blob: &'b BTreeMap) -> Option> { } else if let Some(Ipld::String(value)) = blob.get("cid") { if let Some(mime_type) = mime_type { return Some(Blob { - r#ref: Cid::str(value), + r#ref: CidLink::str(value), mime_type: MimeType::raw(mime_type), size: 0, }); @@ -281,7 +281,7 @@ pub fn json_to_blob<'b>(blob: &'b serde_json::Map) -> let size = blob.get("size").and_then(|v| v.as_u64()); if let (Some(mime_type), Some(size)) = (mime_type, size) { return Some(Blob { - r#ref: Cid::str(value), + r#ref: CidLink::str(value), mime_type: MimeType::raw(mime_type), size: size as usize, }); @@ -290,7 +290,7 @@ pub fn json_to_blob<'b>(blob: &'b serde_json::Map) -> } else if let Some(value) = blob.get("cid").and_then(|v| v.as_str()) { if let Some(mime_type) = mime_type { return Some(Blob { - r#ref: Cid::str(value), + r#ref: CidLink::str(value), mime_type: MimeType::raw(mime_type), size: 0, }); diff --git a/crates/jacquard-common/src/types/value/serde_impl.rs b/crates/jacquard-common/src/types/value/serde_impl.rs index f0f9249c..b30fd80a 100644 --- a/crates/jacquard-common/src/types/value/serde_impl.rs +++ b/crates/jacquard-common/src/types/value/serde_impl.rs @@ -320,7 +320,7 @@ fn apply_type_inference<'s>(mut map: BTreeMap>) -> Result( if let (Some(ref_cid), Some(mime_cowstr), Some(size)) = (ref_cid, mime_type, size) { return Ok(RawData::Blob(Blob { - r#ref: ref_cid, + r#ref: CidLink::str(ref_cid.as_str()).into_static(), mime_type: MimeType::from(mime_cowstr), size, })); diff --git a/crates/jacquard-common/src/xrpc.rs b/crates/jacquard-common/src/xrpc.rs index 5085cd3f..6fb2dca6 100644 --- a/crates/jacquard-common/src/xrpc.rs +++ b/crates/jacquard-common/src/xrpc.rs @@ -15,7 +15,9 @@ pub mod streaming; use ipld_core::ipld::Ipld; #[cfg(feature = "streaming")] -pub use streaming::StreamingResponse; +pub use streaming::{ + StreamingResponse, XrpcProcedureSend, XrpcProcedureStream, XrpcResponseStream, XrpcStreamResp, +}; #[cfg(feature = "websocket")] pub mod subscription; @@ -44,10 +46,7 @@ use crate::{AuthorizationToken, error::AuthError}; use crate::{CowStr, error::XrpcResult}; use crate::{IntoStatic, error::DecodeError}; #[cfg(feature = "streaming")] -use crate::{ - StreamError, - xrpc::streaming::{XrpcProcedureSend, XrpcProcedureStream, XrpcResponseStream, XrpcStreamResp}, -}; +use crate::StreamError; use crate::{error::TransportError, types::value::RawData}; /// Error type for encoding XRPC requests @@ -272,7 +271,7 @@ pub type XrpcResponse = Response<::Response>; #[cfg_attr(not(target_arch = "wasm32"), trait_variant::make(Send))] pub trait XrpcClient: HttpClient { /// Get the base URI for the client. - fn base_uri(&self) -> Url; + fn base_uri(&self) -> impl Future; /// Get the call options for the client. fn opts(&self) -> impl Future> { @@ -316,6 +315,53 @@ pub trait XrpcClient: HttpClient { where R: XrpcRequest + Send + Sync, ::Response: Send + Sync; + +} + +/// Stateful XRPC streaming client trait +#[cfg(feature = "streaming")] +pub trait XrpcStreamingClient: XrpcClient + HttpClientExt { + /// Send an XRPC request and stream the response + #[cfg(not(target_arch = "wasm32"))] + fn download( + &self, + request: R, + ) -> impl Future> + Send + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync, + Self: Sync; + + /// Send an XRPC request and stream the response + #[cfg(target_arch = "wasm32")] + fn download( + &self, + request: R, + ) -> impl Future> + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync; + + /// Stream an XRPC procedure call and its response + #[cfg(not(target_arch = "wasm32"))] + fn stream( + &self, + stream: XrpcProcedureSend>, + ) -> impl Future::Response as XrpcStreamResp>::Frame<'static>>, StreamError>> + where + S: XrpcProcedureStream + 'static, + <::Response as XrpcStreamResp>::Frame<'static>: XrpcStreamResp, + Self: Sync; + + /// Stream an XRPC procedure call and its response + #[cfg(target_arch = "wasm32")] + fn stream( + &self, + stream: XrpcProcedureSend>, + ) -> impl Future::Response as XrpcStreamResp>::Frame<'static>>, StreamError>> + where + S: XrpcProcedureStream + 'static, + <::Response as XrpcStreamResp>::Frame<'static>: XrpcStreamResp; } /// Stateless XRPC call builder. @@ -947,7 +993,7 @@ impl<'a, C: HttpClient + HttpClientExt> XrpcCall<'a, C> { /// Stream an XRPC procedure call and its response /// /// Useful for streaming upload of large payloads, or for "pipe-through" operations - /// where you processing a large payload. + /// where you are processing a large payload. pub async fn stream( self, stream: XrpcProcedureSend>, diff --git a/crates/jacquard-common/src/xrpc/streaming.rs b/crates/jacquard-common/src/xrpc/streaming.rs index 1b44878e..6a44f98c 100644 --- a/crates/jacquard-common/src/xrpc/streaming.rs +++ b/crates/jacquard-common/src/xrpc/streaming.rs @@ -208,7 +208,7 @@ impl XrpcResponseStream { } } -/// XRPC streaming response +/// HTTP streaming response /// /// Similar to `Response` but holds a streaming body instead of a buffer. pub struct StreamingResponse { diff --git a/crates/jacquard-common/src/xrpc/subscription.rs b/crates/jacquard-common/src/xrpc/subscription.rs index 0085cd2a..6c6750f3 100644 --- a/crates/jacquard-common/src/xrpc/subscription.rs +++ b/crates/jacquard-common/src/xrpc/subscription.rs @@ -472,7 +472,7 @@ impl<'a, C: WebSocketClient> SubscriptionCall<'a, C> { #[cfg_attr(not(target_arch = "wasm32"), trait_variant::make(Send))] pub trait SubscriptionClient: WebSocketClient { /// Get the base URI for the client. - fn base_uri(&self) -> Url; + fn base_uri(&self) -> impl Future; /// Get the subscription options for the client. fn subscription_opts(&self) -> impl Future> { @@ -570,7 +570,7 @@ impl WebSocketClient for BasicSubscriptionClient { } impl SubscriptionClient for BasicSubscriptionClient { - fn base_uri(&self) -> Url { + async fn base_uri(&self) -> Url { self.base_uri.clone() } @@ -613,7 +613,7 @@ impl SubscriptionClient for BasicSubscriptionClient { Sub: XrpcSubscription + Send + Sync, Self: Sync, { - let base = self.base_uri(); + let base = self.base_uri().await; self.subscription(base) .with_options(opts) .subscribe(params) diff --git a/crates/jacquard-identity/Cargo.toml b/crates/jacquard-identity/Cargo.toml index e60fe223..f043a8ea 100644 --- a/crates/jacquard-identity/Cargo.toml +++ b/crates/jacquard-identity/Cargo.toml @@ -15,6 +15,7 @@ description = "ATProto identity resolution utilities for Jacquard" [features] dns = ["dep:hickory-resolver"] tracing = ["dep:tracing"] +streaming = ["jacquard-common/streaming", "dep:n0-future"] [dependencies] trait-variant.workspace = true @@ -33,7 +34,7 @@ http.workspace = true serde_html_form.workspace = true urlencoding.workspace = true tracing = { workspace = true, optional = true } - +n0-future = { workspace = true, optional = true } [target.'cfg(not(target_family = "wasm"))'.dependencies] hickory-resolver = { optional = true, version = "0.24", default-features = false, features = ["system-config", "tokio-runtime"]} diff --git a/crates/jacquard-identity/src/lib.rs b/crates/jacquard-identity/src/lib.rs index d7e2678b..c2f59621 100644 --- a/crates/jacquard-identity/src/lib.rs +++ b/crates/jacquard-identity/src/lib.rs @@ -77,6 +77,8 @@ use crate::resolver::{ use bytes::Bytes; use jacquard_api::com_atproto::identity::resolve_did; use jacquard_api::com_atproto::identity::resolve_handle::ResolveHandle; +#[cfg(feature = "streaming")] +use jacquard_common::ByteStream; use jacquard_common::error::TransportError; use jacquard_common::http_client::HttpClient; use jacquard_common::types::did::Did; @@ -89,7 +91,10 @@ use reqwest::StatusCode; use url::{ParseError, Url}; #[cfg(all(feature = "dns", not(target_family = "wasm")))] -use {hickory_resolver::{TokioAsyncResolver, config::ResolverConfig}, std::sync::Arc}; +use { + hickory_resolver::{TokioAsyncResolver, config::ResolverConfig}, + std::sync::Arc, +}; /// Default resolver implementation with configurable fallback order. #[derive(Clone)] @@ -501,6 +506,31 @@ impl HttpClient for JacquardResolver { type Error = reqwest::Error; } +#[cfg(feature = "streaming")] +impl jacquard_common::http_client::HttpClientExt for JacquardResolver { + /// Send HTTP request and return streaming response + fn send_http_streaming( + &self, + request: http::Request>, + ) -> impl Future, Self::Error>> { + self.http.send_http_streaming(request) + } + + /// Send HTTP request with streaming body and receive streaming response + fn send_http_bidirectional( + &self, + parts: http::request::Parts, + body: S, + ) -> impl Future, Self::Error>> + where + S: n0_future::Stream> + + Send + + 'static, + { + self.http.send_http_bidirectional(parts, body) + } +} + /// Warnings produced during identity checks that are not fatal #[derive(Debug, Clone, PartialEq, Eq)] pub enum IdentityWarning { diff --git a/crates/jacquard-oauth/Cargo.toml b/crates/jacquard-oauth/Cargo.toml index 33000d77..83225338 100644 --- a/crates/jacquard-oauth/Cargo.toml +++ b/crates/jacquard-oauth/Cargo.toml @@ -37,6 +37,7 @@ dashmap = "6.1.0" tokio = { workspace = true, default-features = false, features = ["sync"] } reqwest.workspace = true trait-variant.workspace = true +n0-future = { workspace = true, optional = true } webbrowser = { version = "0.8", optional = true } tracing = { workspace = true, optional = true } @@ -50,3 +51,4 @@ loopback = ["dep:rouille"] browser-open = ["dep:webbrowser"] tracing = ["dep:tracing"] websocket = ["jacquard-common/websocket"] +streaming = ["jacquard-common/streaming", "dep:n0-future"] diff --git a/crates/jacquard-oauth/src/client.rs b/crates/jacquard-oauth/src/client.rs index 179d6688..97500f4a 100644 --- a/crates/jacquard-oauth/src/client.rs +++ b/crates/jacquard-oauth/src/client.rs @@ -29,7 +29,7 @@ use jacquard_identity::{ resolver::{DidDocResponse, IdentityError, IdentityResolver, ResolverOptions}, }; use jose_jwk::JwkSet; -use std::sync::Arc; +use std::{future::Future, sync::Arc}; use tokio::sync::RwLock; use url::Url; @@ -458,15 +458,8 @@ where T: OAuthResolver + DpopExt + XrpcExt + Send + Sync + 'static, W: Send + Sync, { - fn base_uri(&self) -> Url { - // base_uri is a synchronous trait method; we must avoid async `.read().await`. - // Use `block_in_place` under Tokio runtime to perform a blocking RwLock read safely. - #[cfg(not(target_arch = "wasm32"))] - if tokio::runtime::Handle::try_current().is_ok() { - return tokio::task::block_in_place(|| self.data.blocking_read().host_url.clone()); - } - - self.data.blocking_read().host_url.clone() + async fn base_uri(&self) -> Url { + self.data.read().await.host_url.clone() } async fn opts(&self) -> CallOptions<'_> { @@ -491,7 +484,7 @@ where R: XrpcRequest + Send + Sync, ::Response: Send + Sync, { - let base_uri = self.base_uri(); + let base_uri = self.base_uri().await; opts.auth = Some(self.access_token().await); let guard = self.data.read().await; let mut dpop = guard.dpop_data.clone(); @@ -524,6 +517,195 @@ where } } +#[cfg(feature = "streaming")] +impl jacquard_common::http_client::HttpClientExt for OAuthSession +where + S: ClientAuthStore + Send + Sync + 'static, + T: OAuthResolver + + DpopExt + + XrpcExt + + jacquard_common::http_client::HttpClientExt + + Send + + Sync + + 'static, + W: Send + Sync, +{ + async fn send_http_streaming( + &self, + request: http::Request>, + ) -> core::result::Result, Self::Error> + { + self.client.send_http_streaming(request).await + } + + async fn send_http_bidirectional( + &self, + parts: http::request::Parts, + body: Str, + ) -> core::result::Result, Self::Error> + where + Str: n0_future::Stream< + Item = core::result::Result, + > + Send + + 'static, + { + self.client.send_http_bidirectional(parts, body).await + } +} + +#[cfg(feature = "streaming")] +impl jacquard_common::xrpc::XrpcStreamingClient for OAuthSession +where + S: ClientAuthStore + Send + Sync + 'static, + T: OAuthResolver + + DpopExt + + XrpcExt + + jacquard_common::http_client::HttpClientExt + + Send + + Sync + + 'static, + W: Send + Sync, +{ + async fn download( + &self, + request: R, + ) -> core::result::Result + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync, + { + use jacquard_common::StreamError; + + let base_uri = ::base_uri(self).await; + let mut opts = self.options.read().await.clone(); + opts.auth = Some(self.access_token().await); + let http_request = build_http_request(&base_uri, &request, &opts) + .map_err(|e| StreamError::protocol(e.to_string()))?; + let guard = self.data.read().await; + let mut dpop = guard.dpop_data.clone(); + let result = self + .client + .dpop_call(&mut dpop) + .send_streaming(http_request) + .await; + drop(guard); + + match result { + Ok(response) => Ok(response), + Err(_e) => { + // Check if it's an auth error and retry + opts.auth = Some( + self.refresh() + .await + .map_err(|e| StreamError::transport(e))?, + ); + let http_request = build_http_request(&base_uri, &request, &opts) + .map_err(|e| StreamError::protocol(e.to_string()))?; + let guard = self.data.read().await; + let mut dpop = guard.dpop_data.clone(); + self.client + .dpop_call(&mut dpop) + .send_streaming(http_request) + .await + .map_err(StreamError::transport) + } + } + } + + async fn stream( + &self, + stream: jacquard_common::xrpc::streaming::XrpcProcedureSend>, + ) -> core::result::Result< + jacquard_common::xrpc::streaming::XrpcResponseStream< + <::Response as jacquard_common::xrpc::streaming::XrpcStreamResp>::Frame<'static>, + >, + jacquard_common::StreamError, + > + where + Str: jacquard_common::xrpc::streaming::XrpcProcedureStream + 'static, + <::Response as jacquard_common::xrpc::streaming::XrpcStreamResp>::Frame<'static>: jacquard_common::xrpc::streaming::XrpcStreamResp, + { + use jacquard_common::StreamError; + use n0_future::{StreamExt, TryStreamExt}; + + let base_uri = self.base_uri().await; + let mut opts = self.options.read().await.clone(); + opts.auth = Some(self.access_token().await); + + let mut url = base_uri; + let mut path = url.path().trim_end_matches('/').to_owned(); + path.push_str("/xrpc/"); + path.push_str(::NSID); + url.set_path(&path); + + let mut builder = http::Request::post(url.to_string()); + + if let Some(token) = &opts.auth { + use jacquard_common::AuthorizationToken; + let hv = match token { + AuthorizationToken::Bearer(t) => { + http::HeaderValue::from_str(&format!("Bearer {}", t.as_ref())) + } + AuthorizationToken::Dpop(t) => { + http::HeaderValue::from_str(&format!("DPoP {}", t.as_ref())) + } + } + .map_err(|e| StreamError::protocol(format!("Invalid authorization token: {}", e)))?; + builder = builder.header(http::header::AUTHORIZATION, hv); + } + + if let Some(proxy) = &opts.atproto_proxy { + builder = builder.header("atproto-proxy", proxy.as_ref()); + } + if let Some(labelers) = &opts.atproto_accept_labelers { + if !labelers.is_empty() { + let joined = labelers + .iter() + .map(|s| s.as_ref()) + .collect::>() + .join(", "); + builder = builder.header("atproto-accept-labelers", joined); + } + } + for (name, value) in &opts.extra_headers { + builder = builder.header(name, value); + } + + let (parts, _) = builder + .body(()) + .map_err(|e| StreamError::protocol(e.to_string()))? + .into_parts(); + + let body_stream = + jacquard_common::stream::ByteStream::new(stream.0.map_ok(|f| f.buffer).boxed()); + + let guard = self.data.read().await; + let mut dpop = guard.dpop_data.clone(); + let result = self + .client + .dpop_call(&mut dpop) + .send_bidirectional(parts, body_stream) + .await; + drop(guard); + + match result { + Ok(response) => { + let (resp_parts, resp_body) = response.into_parts(); + Ok( + jacquard_common::xrpc::streaming::XrpcResponseStream::from_typed_parts( + resp_parts, resp_body, + ), + ) + } + Err(e) => { + // OAuth token refresh and retry is handled by dpop wrapper + // If we get here, it's a real error + Err(StreamError::transport(e)) + } + } + } +} + fn is_invalid_token_response(response: &XrpcResult>) -> bool { match response { Err(ClientError::Auth(AuthError::InvalidToken)) => true, @@ -592,7 +774,7 @@ where T: OAuthResolver + Send + Sync + 'static, W: WebSocketClient + Send + Sync, { - fn base_uri(&self) -> Url { + async fn base_uri(&self) -> Url { #[cfg(not(target_arch = "wasm32"))] if tokio::runtime::Handle::try_current().is_ok() { return tokio::task::block_in_place(|| self.data.blocking_read().host_url.clone()); @@ -608,10 +790,8 @@ where AuthorizationToken::Bearer(t) => format!("Bearer {}", t.as_ref()), AuthorizationToken::Dpop(t) => format!("DPoP {}", t.as_ref()), }; - opts.headers.push(( - CowStr::from("Authorization"), - CowStr::from(auth_value), - )); + opts.headers + .push((CowStr::from("Authorization"), CowStr::from(auth_value))); opts } diff --git a/crates/jacquard-oauth/src/dpop.rs b/crates/jacquard-oauth/src/dpop.rs index 04418cfa..e06448c7 100644 --- a/crates/jacquard-oauth/src/dpop.rs +++ b/crates/jacquard-oauth/src/dpop.rs @@ -109,6 +109,76 @@ impl<'r, C: HttpClient, N: DpopDataSource> DpopCall<'r, C, N> { ) .await } + + #[cfg(feature = "streaming")] + pub async fn send_streaming( + self, + request: Request>, + ) -> Result + where + C: jacquard_common::http_client::HttpClientExt, + { + wrap_request_with_dpop_streaming( + self.client, + self.data_source, + self.is_to_auth_server, + request, + ) + .await + } + + #[cfg(feature = "streaming")] + pub async fn send_bidirectional( + self, + parts: http::request::Parts, + body: jacquard_common::stream::ByteStream, + ) -> Result + where + C: jacquard_common::http_client::HttpClientExt, + { + wrap_request_with_dpop_bidirectional( + self.client, + self.data_source, + self.is_to_auth_server, + parts, + body, + ) + .await + } +} + +/// Extract authorization hash from request headers +fn extract_ath(headers: &http::HeaderMap) -> Option> { + headers + .get("Authorization") + .filter(|v| v.to_str().is_ok_and(|s| s.starts_with("DPoP "))) + .map(|auth| { + URL_SAFE_NO_PAD + .encode(sha2::Sha256::digest(&auth.as_bytes()[5..])) + .into() + }) +} + +/// Get nonce from data source based on target +fn get_nonce(data_source: &N, is_to_auth_server: bool) -> Option> { + if is_to_auth_server { + data_source.authserver_nonce() + } else { + data_source.host_nonce() + } +} + +/// Store nonce in data source based on target +fn store_nonce( + data_source: &mut N, + is_to_auth_server: bool, + nonce: CowStr<'static>, +) { + if is_to_auth_server { + data_source.set_authserver_nonce(nonce); + } else { + data_source.set_host_nonce(nonce); + } } pub async fn wrap_request_with_dpop( @@ -124,22 +194,9 @@ where let uri = request.uri().clone(); let method = request.method().to_cowstr().into_static(); let uri = uri.to_cowstr(); - // https://datatracker.ietf.org/doc/html/rfc9449#section-4.2 - let ath = request - .headers() - .get("Authorization") - .filter(|v| v.to_str().is_ok_and(|s| s.starts_with("DPoP "))) - .map(|auth| { - URL_SAFE_NO_PAD - .encode(sha2::Sha256::digest(&auth.as_bytes()[5..])) - .into() - }); + let ath = extract_ath(request.headers()); - let init_nonce = if is_to_auth_server { - data_source.authserver_nonce() - } else { - data_source.host_nonce() - }; + let init_nonce = get_nonce(data_source, is_to_auth_server); let init_proof = build_dpop_proof( data_source.key(), method.clone(), @@ -157,19 +214,12 @@ where .headers() .get("DPoP-Nonce") .and_then(|v| v.to_str().ok()) - .map(|c| c.to_cowstr()); + .map(|c| CowStr::from(c.to_string())); match &next_nonce { Some(s) if next_nonce != init_nonce => { - // Store the fresh nonce for future requests - if is_to_auth_server { - data_source.set_authserver_nonce(s.clone()); - } else { - data_source.set_host_nonce(s.clone()); - } + store_nonce(data_source, is_to_auth_server, s.clone()); } _ => { - // No nonce was returned or it is the same as the one we sent. No need to - // update the nonce store, or retry the request. return Ok(response); } } @@ -186,6 +236,159 @@ where Ok(response) } +#[cfg(feature = "streaming")] +pub async fn wrap_request_with_dpop_streaming( + client: &T, + data_source: &mut N, + is_to_auth_server: bool, + mut request: Request>, +) -> Result +where + T: jacquard_common::http_client::HttpClientExt, + N: DpopDataSource, +{ + use jacquard_common::xrpc::StreamingResponse; + + let uri = request.uri().clone(); + let method = request.method().to_cowstr().into_static(); + let uri = uri.to_cowstr(); + let ath = extract_ath(request.headers()); + + let init_nonce = get_nonce(data_source, is_to_auth_server); + let init_proof = build_dpop_proof( + data_source.key(), + method.clone(), + uri.clone(), + init_nonce.clone(), + ath.clone(), + )?; + request.headers_mut().insert("DPoP", init_proof.parse()?); + let http_response = client + .send_http_streaming(request.clone()) + .await + .map_err(|e| Error::Inner(e.into()))?; + + let (parts, body) = http_response.into_parts(); + let next_nonce = parts + .headers + .get("DPoP-Nonce") + .and_then(|v| v.to_str().ok()) + .map(|c| CowStr::from(c.to_string())); + match &next_nonce { + Some(s) if next_nonce != init_nonce => { + store_nonce(data_source, is_to_auth_server, s.clone()); + } + _ => { + return Ok(StreamingResponse::new(parts, body)); + } + } + + // For streaming responses, we can't easily check the body for use_dpop_nonce error + // We check status code + headers only + if !is_use_dpop_nonce_error_streaming(is_to_auth_server, parts.status, &parts.headers) { + return Ok(StreamingResponse::new(parts, body)); + } + + let next_proof = build_dpop_proof(data_source.key(), method, uri, next_nonce, ath)?; + request.headers_mut().insert("DPoP", next_proof.parse()?); + let http_response = client + .send_http_streaming(request) + .await + .map_err(|e| Error::Inner(e.into()))?; + let (parts, body) = http_response.into_parts(); + Ok(StreamingResponse::new(parts, body)) +} + +#[cfg(feature = "streaming")] +pub async fn wrap_request_with_dpop_bidirectional( + client: &T, + data_source: &mut N, + is_to_auth_server: bool, + mut parts: http::request::Parts, + body: jacquard_common::stream::ByteStream, +) -> Result +where + T: jacquard_common::http_client::HttpClientExt, + N: DpopDataSource, +{ + use jacquard_common::xrpc::StreamingResponse; + + let uri = parts.uri.clone(); + let method = parts.method.to_cowstr().into_static(); + let uri = uri.to_cowstr(); + let ath = extract_ath(&parts.headers); + + let init_nonce = get_nonce(data_source, is_to_auth_server); + let init_proof = build_dpop_proof( + data_source.key(), + method.clone(), + uri.clone(), + init_nonce.clone(), + ath.clone(), + )?; + parts.headers.insert("DPoP", init_proof.parse()?); + + // Clone the stream for potential retry + let (body1, body2) = body.tee(); + + let http_response = client + .send_http_bidirectional(parts.clone(), body1.into_inner()) + .await + .map_err(|e| Error::Inner(e.into()))?; + + let (resp_parts, resp_body) = http_response.into_parts(); + let next_nonce = resp_parts + .headers + .get("DPoP-Nonce") + .and_then(|v| v.to_str().ok()) + .map(|c| CowStr::from(c.to_string())); + match &next_nonce { + Some(s) if next_nonce != init_nonce => { + store_nonce(data_source, is_to_auth_server, s.clone()); + } + _ => { + return Ok(StreamingResponse::new(resp_parts, resp_body)); + } + } + + // For streaming responses, we can't easily check the body for use_dpop_nonce error + // We check status code + headers only + if !is_use_dpop_nonce_error_streaming(is_to_auth_server, resp_parts.status, &resp_parts.headers) + { + return Ok(StreamingResponse::new(resp_parts, resp_body)); + } + + let next_proof = build_dpop_proof(data_source.key(), method, uri, next_nonce, ath)?; + parts.headers.insert("DPoP", next_proof.parse()?); + let http_response = client + .send_http_bidirectional(parts, body2.into_inner()) + .await + .map_err(|e| Error::Inner(e.into()))?; + let (parts, body) = http_response.into_parts(); + Ok(StreamingResponse::new(parts, body)) +} + +#[cfg(feature = "streaming")] +fn is_use_dpop_nonce_error_streaming( + is_to_auth_server: bool, + status: http::StatusCode, + headers: &http::HeaderMap, +) -> bool { + if is_to_auth_server && status == 400 { + // Can't check body for streaming, so we rely on DPoP-Nonce header presence + return false; + } + if !is_to_auth_server && status == 401 { + if let Some(www_auth) = headers + .get("WWW-Authenticate") + .and_then(|v| v.to_str().ok()) + { + return www_auth.starts_with("DPoP") && www_auth.contains(r#"error="use_dpop_nonce""#); + } + } + false +} + #[inline] fn is_use_dpop_nonce_error(is_to_auth_server: bool, response: &Response>) -> bool { // https://datatracker.ietf.org/doc/html/rfc9449#name-authorization-server-provid diff --git a/crates/jacquard/Cargo.toml b/crates/jacquard/Cargo.toml index 3ae3f492..8ce7230d 100644 --- a/crates/jacquard/Cargo.toml +++ b/crates/jacquard/Cargo.toml @@ -12,7 +12,7 @@ exclude.workspace = true license.workspace = true [features] -default = ["api_full", "dns", "loopback", "derive"] +default = ["api_full", "dns", "loopback", "derive", "streaming"] derive = ["dep:jacquard-derive"] # Minimal API bindings api = ["jacquard-api/minimal"] @@ -38,7 +38,13 @@ tracing = [ "jacquard-identity/tracing", ] dns = ["jacquard-identity/dns"] -streaming = ["jacquard-common/streaming"] +streaming = [ + "jacquard-common/streaming", + "jacquard-oauth/streaming", + "jacquard-identity/streaming", + "dep:n0-future", + "dep:futures" +] websocket = ["jacquard-common/websocket"] [[example]] @@ -73,6 +79,11 @@ path = "../../examples/read_whitewind_post.rs" name = "read_tangled_repo" path = "../../examples/read_tangled_repo.rs" +[[example]] +name = "stream_get_blob" +path = "../../examples/stream_get_blob.rs" +required-features = ["api_bluesky", "streaming"] + [[example]] name = "resolve_did" path = "../../examples/resolve_did.rs" @@ -127,6 +138,8 @@ jose-jwk = { workspace = true, features = ["p256"] } p256 = { workspace = true, features = ["ecdsa"] } rand_core.workspace = true tracing = { workspace = true, optional = true } +n0-future = { workspace = true, optional = true } +futures = { version = "0.3", optional = true } [target.'cfg(not(target_arch = "wasm32"))'.dependencies] reqwest = { workspace = true, features = [ @@ -142,6 +155,9 @@ getrandom = { version = "0.2", features = ["js"] } [dev-dependencies] clap.workspace = true miette = { workspace = true, features = ["fancy"] } +viuer = { version = "0.9", features = ["print-file", "sixel"] } +tiff = { version = "0.6.0-alpha" } +image = { version = "0.25" } [package.metadata.docs.rs] features = ["api_all", "derive", "dns", "loopback"] diff --git a/crates/jacquard/src/client.rs b/crates/jacquard/src/client.rs index 091e5936..260eabcf 100644 --- a/crates/jacquard/src/client.rs +++ b/crates/jacquard/src/client.rs @@ -789,6 +789,8 @@ pub trait AgentSessionExt: AgentSession + IdentityResolver { })?, )); let response = self.send_with_opts(request, opts).await?; + let debug: serde_json::Value = serde_json::from_slice(response.buffer()).unwrap(); + println!("json: {}", serde_json::to_string_pretty(&debug).unwrap()); let output = response.into_output().map_err(|e| match e { XrpcError::Auth(auth) => AgentError::Auth(auth), XrpcError::Generic(g) => AgentError::Generic(g), @@ -912,9 +914,80 @@ impl HttpClient for Agent { } } +#[cfg(feature = "streaming")] +impl jacquard_common::http_client::HttpClientExt for Agent +where + A: AgentSession + jacquard_common::http_client::HttpClientExt, +{ + #[cfg(not(target_arch = "wasm32"))] + fn send_http_streaming( + &self, + request: http::Request>, + ) -> impl Future< + Output = core::result::Result< + http::Response, + Self::Error, + >, + > + Send { + self.inner.send_http_streaming(request) + } + + #[cfg(target_arch = "wasm32")] + fn send_http_streaming( + &self, + request: http::Request>, + ) -> impl Future< + Output = core::result::Result< + http::Response, + Self::Error, + >, + > { + self.inner.send_http_streaming(request) + } + + #[cfg(not(target_arch = "wasm32"))] + fn send_http_bidirectional( + &self, + parts: http::request::Parts, + body: Str, + ) -> impl Future< + Output = core::result::Result< + http::Response, + Self::Error, + >, + > + Send + where + Str: n0_future::Stream< + Item = core::result::Result, + > + Send + + 'static, + { + self.inner.send_http_bidirectional(parts, body) + } + + #[cfg(target_arch = "wasm32")] + fn send_http_bidirectional( + &self, + parts: http::request::Parts, + body: Str, + ) -> impl Future< + Output = core::result::Result< + http::Response, + Self::Error, + >, + > + where + Str: n0_future::Stream< + Item = core::result::Result, + > + 'static, + { + self.inner.send_http_bidirectional(parts, body) + } +} + impl XrpcClient for Agent { - fn base_uri(&self) -> url::Url { - self.inner.base_uri() + async fn base_uri(&self) -> url::Url { + self.inner.base_uri().await } fn opts(&self) -> impl Future> { self.inner.opts() @@ -943,6 +1016,82 @@ impl XrpcClient for Agent { } } +#[cfg(feature = "streaming")] +impl jacquard_common::xrpc::XrpcStreamingClient for Agent +where + A: AgentSession + jacquard_common::xrpc::XrpcStreamingClient, +{ + #[cfg(not(target_arch = "wasm32"))] + fn download( + &self, + request: R, + ) -> impl Future< + Output = core::result::Result< + jacquard_common::xrpc::StreamingResponse, + jacquard_common::StreamError, + >, + > + Send + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync, + Self: Sync, + { + self.inner.download(request) + } + + #[cfg(target_arch = "wasm32")] + fn download( + &self, + request: R, + ) -> impl Future< + Output = core::result::Result< + jacquard_common::xrpc::StreamingResponse, + jacquard_common::StreamError, + >, + > + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync, + { + self.inner.download(request) + } + + #[cfg(not(target_arch = "wasm32"))] + fn stream( + &self, + stream: jacquard_common::xrpc::XrpcProcedureSend>, + ) -> impl Future< + Output = core::result::Result< + jacquard_common::xrpc::XrpcResponseStream<<::Response as jacquard_common::xrpc::XrpcStreamResp>::Frame<'static>>, + jacquard_common::StreamError, + >, + > + where + S: jacquard_common::xrpc::XrpcProcedureStream + 'static, + <::Response as jacquard_common::xrpc::XrpcStreamResp>::Frame<'static>: jacquard_common::xrpc::XrpcStreamResp, + Self: Sync, + { + self.inner.stream::(stream) + } + + #[cfg(target_arch = "wasm32")] + fn stream( + &self, + stream: jacquard_common::xrpc::XrpcProcedureSend>, + ) -> impl Future< + Output = core::result::Result< + jacquard_common::xrpc::XrpcResponseStream<<::Response as jacquard_common::xrpc::XrpcStreamResp>::Frame<'static>>, + jacquard_common::StreamError, + >, + > + where + S: jacquard_common::xrpc::XrpcProcedureStream + 'static, + <::Response as jacquard_common::xrpc::XrpcStreamResp>::Frame<'static>: jacquard_common::xrpc::XrpcStreamResp, + { + self.inner.stream::(stream) + } +} + impl IdentityResolver for Agent { fn options(&self) -> &ResolverOptions { self.inner.options() diff --git a/crates/jacquard/src/client/credential_session.rs b/crates/jacquard/src/client/credential_session.rs index b651cd42..2a930b65 100644 --- a/crates/jacquard/src/client/credential_session.rs +++ b/crates/jacquard/src/client/credential_session.rs @@ -433,29 +433,10 @@ where T: HttpClient + XrpcExt + Send + Sync + 'static, W: Send + Sync, { - fn base_uri(&self) -> Url { - // base_uri is a synchronous trait method; avoid `.await` here. - // Under Tokio, use `block_in_place` to make a blocking RwLock read safe. - #[cfg(not(target_arch = "wasm32"))] - if tokio::runtime::Handle::try_current().is_ok() { - tokio::task::block_in_place(|| { - self.endpoint.blocking_read().clone().unwrap_or( - Url::parse("https://public.bsky.app") - .expect("public appview should be valid url"), - ) - }) - } else { - self.endpoint.blocking_read().clone().unwrap_or( - Url::parse("https://public.bsky.app").expect("public appview should be valid url"), - ) - } - - #[cfg(target_arch = "wasm32")] - { - self.endpoint.blocking_read().clone().unwrap_or( - Url::parse("https://public.bsky.app").expect("public appview should be valid url"), - ) - } + async fn base_uri(&self) -> Url { + self.endpoint.read().await.clone().unwrap_or( + Url::parse("https://public.bsky.app").expect("public appview should be valid url"), + ) } async fn send(&self, request: R) -> XrpcResult> @@ -476,7 +457,7 @@ where R: XrpcRequest + Send + Sync, ::Response: Send + Sync, { - let base_uri = self.base_uri(); + let base_uri = self.base_uri().await; let auth = self.access_token().await; opts.auth = auth; let resp = self @@ -512,6 +493,218 @@ fn is_expired(response: &XrpcResult>) -> bool { } } +#[cfg(feature = "streaming")] +impl jacquard_common::http_client::HttpClientExt for CredentialSession +where + S: SessionStore + Send + Sync + 'static, + T: HttpClient + XrpcExt + jacquard_common::http_client::HttpClientExt + Send + Sync + 'static, + W: Send + Sync, +{ + async fn send_http_streaming( + &self, + request: http::Request>, + ) -> core::result::Result, Self::Error> { + self.client.send_http_streaming(request).await + } + + async fn send_http_bidirectional( + &self, + parts: http::request::Parts, + body: Str, + ) -> core::result::Result, Self::Error> + where + Str: n0_future::Stream> + + Send + + 'static, + { + self.client.send_http_bidirectional(parts, body).await + } +} + +#[cfg(feature = "streaming")] +impl jacquard_common::xrpc::XrpcStreamingClient for CredentialSession +where + S: SessionStore + Send + Sync + 'static, + T: HttpClient + XrpcExt + jacquard_common::http_client::HttpClientExt + Send + Sync + 'static, + W: Send + Sync, +{ + async fn download( + &self, + request: R, + ) -> core::result::Result + where + R: XrpcRequest + Send + Sync, + ::Response: Send + Sync, + { + use jacquard_common::{StreamError, xrpc::build_http_request}; + + let base_uri = ::base_uri(self).await; + let mut opts = self.options.read().await.clone(); + opts.auth = self.access_token().await; + + let http_request = build_http_request(&base_uri, &request, &opts) + .map_err(|e| StreamError::protocol(e.to_string()))?; + + let response = self + .client + .send_http_streaming(http_request.clone()) + .await + .map_err(StreamError::transport)?; + + let (parts, body) = response.into_parts(); + let status = parts.status; + + // Check if expired based on status code + if status == http::StatusCode::UNAUTHORIZED || status == http::StatusCode::BAD_REQUEST { + // Try to refresh + let auth = self.refresh().await.map_err(StreamError::transport)?; + opts.auth = Some(auth); + + let http_request = build_http_request(&base_uri, &request, &opts) + .map_err(|e| StreamError::protocol(e.to_string()))?; + + let response = self + .client + .send_http_streaming(http_request) + .await + .map_err(StreamError::transport)?; + let (parts, body) = response.into_parts(); + Ok(jacquard_common::xrpc::StreamingResponse::new(parts, body)) + } else { + Ok(jacquard_common::xrpc::StreamingResponse::new(parts, body)) + } + } + + async fn stream( + &self, + stream: jacquard_common::xrpc::streaming::XrpcProcedureSend>, + ) -> core::result::Result< + jacquard_common::xrpc::streaming::XrpcResponseStream< + <::Response as jacquard_common::xrpc::streaming::XrpcStreamResp>::Frame<'static>, + >, + jacquard_common::StreamError, + > + where + Str: jacquard_common::xrpc::streaming::XrpcProcedureStream + 'static, + <::Response as jacquard_common::xrpc::streaming::XrpcStreamResp>::Frame<'static>: jacquard_common::xrpc::streaming::XrpcStreamResp, + { + use jacquard_common::StreamError; + use n0_future::{StreamExt, TryStreamExt}; + + let base_uri = self.base_uri().await; + let mut opts = self.options.read().await.clone(); + opts.auth = self.access_token().await; + + let mut url = base_uri; + let mut path = url.path().trim_end_matches('/').to_owned(); + path.push_str("/xrpc/"); + path.push_str(::NSID); + url.set_path(&path); + + let mut builder = http::Request::post(url.to_string()); + + if let Some(token) = &opts.auth { + use jacquard_common::AuthorizationToken; + let hv = match token { + AuthorizationToken::Bearer(t) => { + http::HeaderValue::from_str(&format!("Bearer {}", t.as_ref())) + } + AuthorizationToken::Dpop(t) => { + http::HeaderValue::from_str(&format!("DPoP {}", t.as_ref())) + } + } + .map_err(|e| StreamError::protocol(format!("Invalid authorization token: {}", e)))?; + builder = builder.header(http::header::AUTHORIZATION, hv); + } + + if let Some(proxy) = &opts.atproto_proxy { + builder = builder.header("atproto-proxy", proxy.as_ref()); + } + if let Some(labelers) = &opts.atproto_accept_labelers { + if !labelers.is_empty() { + let joined = labelers + .iter() + .map(|s| s.as_ref()) + .collect::>() + .join(", "); + builder = builder.header("atproto-accept-labelers", joined); + } + } + for (name, value) in &opts.extra_headers { + builder = builder.header(name, value); + } + + let (parts, _) = builder + .body(()) + .map_err(|e| StreamError::protocol(e.to_string()))? + .into_parts(); + + let body_stream = + jacquard_common::stream::ByteStream::new(stream.0.map_ok(|f| f.buffer).boxed()); + + let response = self + .client + .send_http_bidirectional(parts.clone(), body_stream.into_inner()) + .await + .map_err(StreamError::transport)?; + + let (resp_parts, resp_body) = response.into_parts(); + let status = resp_parts.status; + + // Check if expired + if status == http::StatusCode::UNAUTHORIZED || status == http::StatusCode::BAD_REQUEST { + // Try to refresh + let auth = self.refresh().await.map_err(StreamError::transport)?; + opts.auth = Some(auth); + + // Rebuild request with new auth + let mut builder = http::Request::post(url.to_string()); + if let Some(token) = &opts.auth { + use jacquard_common::AuthorizationToken; + let hv = match token { + AuthorizationToken::Bearer(t) => { + http::HeaderValue::from_str(&format!("Bearer {}", t.as_ref())) + } + AuthorizationToken::Dpop(t) => { + http::HeaderValue::from_str(&format!("DPoP {}", t.as_ref())) + } + } + .map_err(|e| StreamError::protocol(format!("Invalid authorization token: {}", e)))?; + builder = builder.header(http::header::AUTHORIZATION, hv); + } + if let Some(proxy) = &opts.atproto_proxy { + builder = builder.header("atproto-proxy", proxy.as_ref()); + } + if let Some(labelers) = &opts.atproto_accept_labelers { + if !labelers.is_empty() { + let joined = labelers + .iter() + .map(|s| s.as_ref()) + .collect::>() + .join(", "); + builder = builder.header("atproto-accept-labelers", joined); + } + } + for (name, value) in &opts.extra_headers { + builder = builder.header(name, value); + } + + let (parts, _) = builder + .body(()) + .map_err(|e| StreamError::protocol(e.to_string()))? + .into_parts(); + + // Can't retry with the same stream - it's been consumed + // This is a limitation of streaming upload with auth refresh + return Err(StreamError::protocol("Authentication failed on streaming upload and stream cannot be retried".to_string())); + } + + Ok(jacquard_common::xrpc::streaming::XrpcResponseStream::from_typed_parts( + resp_parts, resp_body, + )) + } +} + impl IdentityResolver for CredentialSession where S: SessionStore + Send + Sync + 'static, @@ -596,10 +789,8 @@ where AuthorizationToken::Bearer(t) => format!("Bearer {}", t.as_ref()), AuthorizationToken::Dpop(t) => format!("DPoP {}", t.as_ref()), }; - opts.headers.push(( - CowStr::from("Authorization"), - CowStr::from(auth_value), - )); + opts.headers + .push((CowStr::from("Authorization"), CowStr::from(auth_value))); } opts } diff --git a/crates/jacquard/src/lib.rs b/crates/jacquard/src/lib.rs index 08cd9b96..7e18fadf 100644 --- a/crates/jacquard/src/lib.rs +++ b/crates/jacquard/src/lib.rs @@ -219,6 +219,9 @@ pub mod client; +#[cfg(feature = "streaming")] +pub mod streaming; + pub use common::*; #[cfg(feature = "api")] pub use jacquard_api as api; diff --git a/crates/jacquard/src/streaming.rs b/crates/jacquard/src/streaming.rs new file mode 100644 index 00000000..8fca8b72 --- /dev/null +++ b/crates/jacquard/src/streaming.rs @@ -0,0 +1,3 @@ +pub mod blob; +pub mod repo; +pub mod video; diff --git a/crates/jacquard/src/streaming/blob.rs b/crates/jacquard/src/streaming/blob.rs new file mode 100644 index 00000000..1c33261e --- /dev/null +++ b/crates/jacquard/src/streaming/blob.rs @@ -0,0 +1,61 @@ +//! Streaming support for blob uploads + +use bytes::Bytes; +use jacquard_api::com_atproto::repo::upload_blob::{UploadBlob, UploadBlobOutput}; +use jacquard_common::{ + StreamError, + xrpc::streaming::{XrpcProcedureStream, XrpcStreamResp}, +}; +use serde::{Deserialize, Serialize}; + +/// Streaming implementation for com.atproto.repo.uploadBlob +pub struct UploadBlobStream; + +impl XrpcProcedureStream for UploadBlobStream { + const NSID: &'static str = "com.atproto.repo.uploadBlob"; + const ENCODING: &'static str = "*/*"; + + type Frame<'de> = Bytes; + type Request = UploadBlob; + type Response = UploadBlobStreamResponse; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(data) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(Bytes::copy_from_slice(frame)) + } +} + +/// Response marker for streaming uploadBlob +pub struct UploadBlobStreamResponse; + +impl XrpcStreamResp for UploadBlobStreamResponse { + const NSID: &'static str = "com.atproto.repo.uploadBlob"; + const ENCODING: &'static str = "application/json"; + + type Frame<'de> = UploadBlobOutput<'de>; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(Bytes::from_owner( + serde_json::to_vec(&data).map_err(StreamError::encode)?, + )) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(serde_json::from_slice(frame).map_err(StreamError::decode)?) + } +} diff --git a/crates/jacquard/src/streaming/repo.rs b/crates/jacquard/src/streaming/repo.rs new file mode 100644 index 00000000..63fe5b1e --- /dev/null +++ b/crates/jacquard/src/streaming/repo.rs @@ -0,0 +1,83 @@ +//! Streaming support for repository operations + +use bytes::Bytes; +use jacquard_api::com_atproto::repo::import_repo::ImportRepo; +use jacquard_common::{ + xrpc::streaming::{XrpcProcedureStream, XrpcStreamResp}, + StreamError, +}; +use serde::{Deserialize, Serialize}; + +/// Streaming implementation for com.atproto.repo.importRepo +pub struct ImportRepoStream; + +impl XrpcProcedureStream for ImportRepoStream { + const NSID: &'static str = "com.atproto.repo.importRepo"; + const ENCODING: &'static str = "application/vnd.ipld.car"; + + type Frame<'de> = Bytes; + type Request = ImportRepo; + type Response = ImportRepoStreamResponse; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(data) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(Bytes::copy_from_slice(frame)) + } +} + +/// Response marker for streaming importRepo +pub struct ImportRepoStreamResponse; + +impl XrpcStreamResp for ImportRepoStreamResponse { + const NSID: &'static str = "com.atproto.repo.importRepo"; + const ENCODING: &'static str = "application/json"; + + type Frame<'de> = (); + + fn encode_frame<'de>(_data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(Bytes::new()) + } + + fn decode_frame<'de>(_frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(()) + } +} + +/// Streaming implementation for com.atproto.sync.getRepo +pub struct GetRepoStream; + +impl XrpcStreamResp for GetRepoStream { + const NSID: &'static str = "com.atproto.sync.getRepo"; + const ENCODING: &'static str = "application/vnd.ipld.car"; + + type Frame<'de> = Bytes; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(data) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(Bytes::copy_from_slice(frame)) + } +} diff --git a/crates/jacquard/src/streaming/video.rs b/crates/jacquard/src/streaming/video.rs new file mode 100644 index 00000000..1f7dfb4f --- /dev/null +++ b/crates/jacquard/src/streaming/video.rs @@ -0,0 +1,61 @@ +//! Streaming support for video uploads + +use bytes::Bytes; +use jacquard_api::app_bsky::video::upload_video::{UploadVideo, UploadVideoOutput}; +use jacquard_common::{ + xrpc::streaming::{XrpcProcedureStream, XrpcStreamResp}, + StreamError, +}; +use serde::{Deserialize, Serialize}; + +/// Streaming implementation for app.bsky.video.uploadVideo +pub struct UploadVideoStream; + +impl XrpcProcedureStream for UploadVideoStream { + const NSID: &'static str = "app.bsky.video.uploadVideo"; + const ENCODING: &'static str = "video/mp4"; + + type Frame<'de> = Bytes; + type Request = UploadVideo; + type Response = UploadVideoStreamResponse; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(data) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(Bytes::copy_from_slice(frame)) + } +} + +/// Response marker for streaming uploadVideo +pub struct UploadVideoStreamResponse; + +impl XrpcStreamResp for UploadVideoStreamResponse { + const NSID: &'static str = "app.bsky.video.uploadVideo"; + const ENCODING: &'static str = "application/json"; + + type Frame<'de> = UploadVideoOutput<'de>; + + fn encode_frame<'de>(data: Self::Frame<'de>) -> Result + where + Self::Frame<'de>: Serialize, + { + Ok(Bytes::from_owner( + serde_json::to_vec(&data).map_err(StreamError::encode)?, + )) + } + + fn decode_frame<'de>(frame: &'de [u8]) -> Result, StreamError> + where + Self::Frame<'de>: Deserialize<'de>, + { + Ok(serde_json::from_slice(frame).map_err(StreamError::decode)?) + } +} diff --git a/examples/stream_get_blob.rs b/examples/stream_get_blob.rs new file mode 100644 index 00000000..1f54f736 --- /dev/null +++ b/examples/stream_get_blob.rs @@ -0,0 +1,61 @@ +use clap::Parser; +use jacquard::StreamingResponse; +use jacquard::api::com_atproto::sync::get_blob::GetBlob; +use jacquard::client::Agent; +use jacquard::types::cid::Cid; +use jacquard::types::did::Did; +use jacquard::xrpc::XrpcStreamingClient; +use jacquard_oauth::authstore::MemoryAuthStore; +use jacquard_oauth::client::OAuthClient; +use jacquard_oauth::loopback::LoopbackConfig; +use n0_future::StreamExt; + +#[derive(Parser, Debug)] +#[command( + author, + version, + about = "Download a blob from a PDS and stream the response, then display it, if it's an image" +)] +struct Args { + input: String, + #[arg(short, long)] + did: String, + #[arg(short, long)] + cid: String, +} + +#[tokio::main] +async fn main() -> miette::Result<()> { + let args = Args::parse(); + + let oauth = OAuthClient::with_default_config(MemoryAuthStore::new()); + let session = oauth + .login_with_local_server(args.input, Default::default(), LoopbackConfig::default()) + .await?; + + let agent: Agent<_> = Agent::from(session); + // Use the streaming `.download()` method with the generated API parameter struct + let output: StreamingResponse = agent + .download(GetBlob { + did: Did::new_owned(args.did)?, + cid: Cid::str(&args.cid), + }) + .await?; + + let (parts, body_stream) = output.into_parts(); + + println!("Parts: {:?}", parts); + + let mut buf: Vec = Vec::new(); + let mut stream = body_stream.into_inner(); + + while let Some(Ok(chunk)) = stream.as_mut().next().await { + buf.append(&mut chunk.to_vec()); + } + + if let Ok(img) = image::load_from_memory(&buf) { + viuer::print(&img, &viuer::Config::default()).expect("Image printing failed."); + } + + Ok(()) +}