From bbac60d744caa9a5b7156575f8d65343c24479fc Mon Sep 17 00:00:00 2001 From: Erik Simon Date: Wed, 8 Jul 2026 17:25:35 +0200 Subject: [PATCH] harness-providers: Anthropic codec + provider /v1/messages request builder (cache_control breakpoints on the first 2 system blocks + last 2 messages, tool schema -> input_schema, extended thinking budget) and an SSE decoder built on eventsource-stream, mapping content_block_start/delta/stop and message_delta into our normalized LlmEvent stream (text, thinking+signature, streamed tool-call JSON accumulated and parsed at content_block_stop, usage merged from message_start + message_delta). AnthropicProvider wires this to reqwest with x-api-key/anthropic-version/anthropic-beta headers and classifies HTTP errors into RateLimited/Auth/Overloaded/Http. ProviderRegistry resolves "provider/model" strings. Also tightened processor::process_step's cancellation: the event loop now selects the stream poll against ctx.cancel instead of only checking at the top of the loop, so a blocked provider stream is actually interrupted by abort (matches docs/02-engine.md's cancellation semantics). 127 tests passing, clippy clean. --- Cargo.lock | 1025 ++++++++++++++++- Cargo.toml | 3 +- crates/harness-core/src/engine/processor.rs | 24 +- crates/harness-providers/Cargo.toml | 15 + crates/harness-providers/src/anthropic.rs | 157 +++ .../harness-providers/src/codec/anthropic.rs | 499 ++++++++ crates/harness-providers/src/codec/mod.rs | 1 + crates/harness-providers/src/lib.rs | 7 +- crates/harness-providers/src/registry.rs | 57 + 9 files changed, 1761 insertions(+), 27 deletions(-) create mode 100644 crates/harness-providers/src/anthropic.rs create mode 100644 crates/harness-providers/src/codec/anthropic.rs create mode 100644 crates/harness-providers/src/codec/mod.rs create mode 100644 crates/harness-providers/src/registry.rs diff --git a/Cargo.lock b/Cargo.lock index 7bed947..0c5bdc5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,12 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "adler2" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa" + [[package]] name = "ahash" version = "0.8.12" @@ -23,6 +29,40 @@ dependencies = [ "memchr", ] +[[package]] +name = "async-compression" +version = "0.4.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e79b3f8a79cccc2898f31920fc69f304859b3bd567490f75ebf51ae1c792a9ac" +dependencies = [ + "compression-codecs", + "compression-core", + "pin-project-lite", + "tokio", +] + +[[package]] +name = "async-stream" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b5a71a6f37880a80d1d7f19efd781e4b5de42c88f0722cc13bcb6cc2cfe8476" +dependencies = [ + "async-stream-impl", + "futures-core", + "pin-project-lite", +] + +[[package]] +name = "async-stream-impl" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7c24de15d275a1ecfd47a380fb4d5ec9bfe0933f309ed5e705b775596a3574d" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "async-trait" version = "0.1.89" @@ -34,6 +74,18 @@ dependencies = [ "syn", ] +[[package]] +name = "atomic-waker" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" + +[[package]] +name = "base64" +version = "0.22.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" + [[package]] name = "bitflags" version = "2.13.0" @@ -79,6 +131,58 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +[[package]] +name = "cfg_aliases" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" + +[[package]] +name = "chacha20" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" +dependencies = [ + "cfg-if", + "cpufeatures", + "rand_core 0.10.1", +] + +[[package]] +name = "compression-codecs" +version = "0.4.38" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce2548391e9c1929c21bf6aa2680af86fe4c1b33e6cea9ac1cfeec0bd11218cf" +dependencies = [ + "compression-core", + "flate2", + "memchr", +] + +[[package]] +name = "compression-core" +version = "0.4.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc14f565cf027a105f7a44ccf9e5b424348421a1d8952a8fc9d499d313107789" + +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + +[[package]] +name = "crc32fast" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9481c1c90cbf2ac953f07c8d4a58aa3945c425b7185c9154d67a65e4230da511" +dependencies = [ + "cfg-if", +] + [[package]] name = "crossbeam-deque" version = "0.8.7" @@ -125,6 +229,17 @@ dependencies = [ "windows-sys 0.48.0", ] +[[package]] +name = "displaydoc" +version = "0.2.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ac70aa55017e108007fbaf5aa0f54b021c98f92ff8af59d42eda9da96e3dd4f" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "dyn-clone" version = "1.0.20" @@ -159,6 +274,17 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "eventsource-stream" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74fef4569247a5f429d9156b9d0a2599914385dd189c539334c625d8099d90ab" +dependencies = [ + "futures-core", + "nom", + "pin-project-lite", +] + [[package]] name = "fallible-iterator" version = "0.3.0" @@ -183,6 +309,25 @@ version = "0.1.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582" +[[package]] +name = "flate2" +version = "1.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "843fba2746e448b37e26a819579957415c8cef339bf08564fe8b7ddbd959573c" +dependencies = [ + "crc32fast", + "miniz_oxide", +] + +[[package]] +name = "form_urlencoded" +version = "1.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb4cb245038516f5f85277875cdaa4f7d2c9a0fa0468de06ed190163b1581fcf" +dependencies = [ + "percent-encoding", +] + [[package]] name = "futures" version = "0.3.32" @@ -278,8 +423,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", + "js-sys", "libc", "wasi", + "wasm-bindgen", ] [[package]] @@ -290,10 +437,24 @@ checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" dependencies = [ "cfg-if", "libc", - "r-efi", + "r-efi 5.3.0", "wasip2", ] +[[package]] +name = "getrandom" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" +dependencies = [ + "cfg-if", + "js-sys", + "libc", + "r-efi 6.0.0", + "rand_core 0.10.1", + "wasm-bindgen", +] + [[package]] name = "globset" version = "0.4.18" @@ -394,7 +555,19 @@ dependencies = [ name = "harness-providers" version = "0.1.0" dependencies = [ + "async-stream", + "async-trait", + "bytes", + "eventsource-stream", + "futures", "harness-core", + "reqwest", + "serde", + "serde_json", + "thiserror 2.0.18", + "tokio", + "tokio-util", + "tracing", ] [[package]] @@ -445,6 +618,207 @@ dependencies = [ "hashbrown", ] +[[package]] +name = "http" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6970f50e31d6fc17d3fa27329444bfa74e196cf62e95052a3f6fee181dba6425" +dependencies = [ + "bytes", + "itoa", +] + +[[package]] +name = "http-body" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1efedce1fb8e6913f23e0c92de8e62cd5b772a67e7b3946df930a62566c93184" +dependencies = [ + "bytes", + "http", +] + +[[package]] +name = "http-body-util" +version = "0.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b021d93e26becf5dc7e1b75b1bed1fd93124b374ceb73f43d4d4eafec896a64a" +dependencies = [ + "bytes", + "futures-core", + "http", + "http-body", + "pin-project-lite", +] + +[[package]] +name = "httparse" +version = "1.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" + +[[package]] +name = "hyper" +version = "1.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "55281c53a1894c864990125767da440a4e630446785086f52523b20033b74498" +dependencies = [ + "atomic-waker", + "bytes", + "futures-channel", + "futures-core", + "http", + "http-body", + "httparse", + "itoa", + "pin-project-lite", + "smallvec", + "tokio", + "want", +] + +[[package]] +name = "hyper-rustls" +version = "0.27.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" +dependencies = [ + "http", + "hyper", + "hyper-util", + "rustls", + "tokio", + "tokio-rustls", + "tower-service", + "webpki-roots", +] + +[[package]] +name = "hyper-util" +version = "0.1.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" +dependencies = [ + "base64", + "bytes", + "futures-channel", + "futures-util", + "http", + "http-body", + "hyper", + "ipnet", + "libc", + "percent-encoding", + "pin-project-lite", + "socket2", + "tokio", + "tower-service", + "tracing", +] + +[[package]] +name = "icu_collections" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2984d1cd16c883d7935b9e07e44071dca8d917fd52ecc02c04d5fa0b5a3f191c" +dependencies = [ + "displaydoc", + "potential_utf", + "utf8_iter", + "yoke", + "zerofrom", + "zerovec", +] + +[[package]] +name = "icu_locale_core" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92219b62b3e2b4d88ac5119f8904c10f8f61bf7e95b640d25ba3075e6cac2c29" +dependencies = [ + "displaydoc", + "litemap", + "tinystr", + "writeable", + "zerovec", +] + +[[package]] +name = "icu_normalizer" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c56e5ee99d6e3d33bd91c5d85458b6005a22140021cc324cea84dd0e72cff3b4" +dependencies = [ + "icu_collections", + "icu_normalizer_data", + "icu_properties", + "icu_provider", + "smallvec", + "zerovec", +] + +[[package]] +name = "icu_normalizer_data" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da3be0ae77ea334f4da67c12f149704f19f81d1adf7c51cf482943e84a2bad38" + +[[package]] +name = "icu_properties" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bee3b67d0ea5c2cca5003417989af8996f8604e34fb9ddf96208a033901e70de" +dependencies = [ + "icu_collections", + "icu_locale_core", + "icu_properties_data", + "icu_provider", + "zerotrie", + "zerovec", +] + +[[package]] +name = "icu_properties_data" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e2bbb201e0c04f7b4b3e14382af113e17ba4f63e2c9d2ee626b720cbce54a14" + +[[package]] +name = "icu_provider" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "139c4cf31c8b5f33d7e199446eff9c1e02decfc2f0eec2c8d71f65befa45b421" +dependencies = [ + "displaydoc", + "icu_locale_core", + "writeable", + "yoke", + "zerofrom", + "zerotrie", + "zerovec", +] + +[[package]] +name = "idna" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3b0875f23caa03898994f6ddc501886a45c7d3d62d04d2d90788d47be1b1e4de" +dependencies = [ + "idna_adapter", + "smallvec", + "utf8_iter", +] + +[[package]] +name = "idna_adapter" +version = "1.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb68373c0d6620ef8105e855e7745e18b0d00d3bdb07fb532e434244cdb9a714" +dependencies = [ + "icu_normalizer", + "icu_properties", +] + [[package]] name = "ignore" version = "0.4.27" @@ -461,6 +835,12 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "ipnet" +version = "2.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2" + [[package]] name = "itoa" version = "1.0.18" @@ -519,6 +899,12 @@ version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" +[[package]] +name = "litemap" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92daf443525c4cce67b150400bc2316076100ce0b3686209eb8cf3c31612e6f0" + [[package]] name = "lock_api" version = "0.4.14" @@ -534,6 +920,12 @@ version = "0.4.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" +[[package]] +name = "lru-slab" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" + [[package]] name = "memchr" version = "2.8.3" @@ -549,6 +941,22 @@ dependencies = [ "libc", ] +[[package]] +name = "minimal-lexical" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68354c5c6bd36d73ff3feceb05efa59b6acb7626617f4962be322a825e61f79a" + +[[package]] +name = "miniz_oxide" +version = "0.8.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fa76a2c86f704bdb222d66965fb3d63269ce38518b83cb0575fca855ebb6316" +dependencies = [ + "adler2", + "simd-adler32", +] + [[package]] name = "mio" version = "1.2.1" @@ -560,6 +968,16 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "nom" +version = "7.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d273983c5a657a70a3e8f2a01329822f3b8c8172b73826411a55751e404a0a4a" +dependencies = [ + "memchr", + "minimal-lexical", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -595,6 +1013,12 @@ dependencies = [ "windows-link", ] +[[package]] +name = "percent-encoding" +version = "2.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -607,6 +1031,15 @@ version = "0.3.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" +[[package]] +name = "potential_utf" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0103b1cef7ec0cf76490e969665504990193874ea05c85ff9bab8b911d0a0564" +dependencies = [ + "zerovec", +] + [[package]] name = "ppv-lite86" version = "0.2.21" @@ -625,6 +1058,62 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "quinn" +version = "0.11.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +dependencies = [ + "bytes", + "cfg_aliases", + "pin-project-lite", + "quinn-proto", + "quinn-udp", + "rustc-hash", + "rustls", + "socket2", + "thiserror 2.0.18", + "tokio", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-proto" +version = "0.11.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" +dependencies = [ + "bytes", + "getrandom 0.4.3", + "lru-slab", + "rand 0.10.2", + "rand_pcg", + "ring", + "rustc-hash", + "rustls", + "rustls-pki-types", + "slab", + "thiserror 2.0.18", + "tinyvec", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-udp" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35a133f956daabe89a61a685c2649f13d82d5aa4bd5d12d1277e1072a21c0694" +dependencies = [ + "cfg_aliases", + "libc", + "once_cell", + "socket2", + "tracing", + "windows-sys 0.61.2", +] + [[package]] name = "quote" version = "1.0.46" @@ -640,6 +1129,12 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" +[[package]] +name = "r-efi" +version = "6.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" + [[package]] name = "rand" version = "0.9.4" @@ -647,7 +1142,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "44c5af06bb1b7d3216d91932aed5265164bf384dc89cd6ba05cf59a35f5f76ea" dependencies = [ "rand_chacha", - "rand_core", + "rand_core 0.9.5", +] + +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "chacha20", + "getrandom 0.4.3", + "rand_core 0.10.1", ] [[package]] @@ -657,7 +1163,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.9.5", ] [[package]] @@ -669,6 +1175,21 @@ dependencies = [ "getrandom 0.3.4", ] +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + +[[package]] +name = "rand_pcg" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" +dependencies = [ + "rand_core 0.10.1", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -706,6 +1227,61 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" +[[package]] +name = "reqwest" +version = "0.12.28" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" +dependencies = [ + "base64", + "bytes", + "futures-core", + "futures-util", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-rustls", + "hyper-util", + "js-sys", + "log", + "percent-encoding", + "pin-project-lite", + "quinn", + "rustls", + "rustls-pki-types", + "serde", + "serde_json", + "serde_urlencoded", + "sync_wrapper", + "tokio", + "tokio-rustls", + "tokio-util", + "tower", + "tower-http", + "tower-service", + "url", + "wasm-bindgen", + "wasm-bindgen-futures", + "wasm-streams", + "web-sys", + "webpki-roots", +] + +[[package]] +name = "ring" +version = "0.17.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" +dependencies = [ + "cc", + "cfg-if", + "getrandom 0.2.17", + "libc", + "untrusted", + "windows-sys 0.52.0", +] + [[package]] name = "rusqlite" version = "0.32.1" @@ -720,6 +1296,12 @@ dependencies = [ "smallvec", ] +[[package]] +name = "rustc-hash" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b1e7f9a428571be2dc5bc0505c13fb6bf936822b894ec87abf8a08a4e51742d" + [[package]] name = "rustix" version = "1.1.4" @@ -733,12 +1315,53 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "rustls" +version = "0.23.41" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b92b125634d9b795e7beca796cc790df15a7fb38323bf3196fda83292d06b1f" +dependencies = [ + "once_cell", + "ring", + "rustls-pki-types", + "rustls-webpki", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "764899a24af3980067ee14bc143654f297b22eaebfe3c7b6b211920a5a59b046" +dependencies = [ + "web-time", + "zeroize", +] + +[[package]] +name = "rustls-webpki" +version = "0.103.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +dependencies = [ + "ring", + "rustls-pki-types", + "untrusted", +] + [[package]] name = "rustversion" version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f" +[[package]] +name = "ryu" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" + [[package]] name = "same-file" version = "1.0.6" @@ -832,6 +1455,18 @@ dependencies = [ "zmij", ] +[[package]] +name = "serde_urlencoded" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3491c14715ca2294c4d6a88f15e84739788c1d030eed8c110436aafdaa2f3fd" +dependencies = [ + "form_urlencoded", + "itoa", + "ryu", + "serde", +] + [[package]] name = "shell-words" version = "1.1.1" @@ -854,6 +1489,12 @@ dependencies = [ "libc", ] +[[package]] +name = "simd-adler32" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "703d5c7ef118737c72f1af64ad2f6f8c5e1921f818cdcb97b8fe6fc69bf66214" + [[package]] name = "similar" version = "2.7.0" @@ -882,6 +1523,18 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "stable_deref_trait" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" + +[[package]] +name = "subtle" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" + [[package]] name = "syn" version = "2.0.118" @@ -893,6 +1546,26 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "sync_wrapper" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263" +dependencies = [ + "futures-core", +] + +[[package]] +name = "synstructure" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "728a70f3dbaf5bab7f0c4b1ac8d7ae5ea60a4b5549c8a5914361c99147a709d2" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "tempfile" version = "3.27.0" @@ -946,6 +1619,31 @@ dependencies = [ "syn", ] +[[package]] +name = "tinystr" +version = "0.8.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c8323304221c2a851516f22236c5722a72eaa19749016521d6dff0824447d96d" +dependencies = [ + "displaydoc", + "zerovec", +] + +[[package]] +name = "tinyvec" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3e61e67053d25a4e82c844e8424039d9745781b3fc4f32b8d55ed50f5f667ef3" +dependencies = [ + "tinyvec_macros", +] + +[[package]] +name = "tinyvec_macros" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" + [[package]] name = "tokio" version = "1.52.3" @@ -974,6 +1672,16 @@ dependencies = [ "syn", ] +[[package]] +name = "tokio-rustls" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" +dependencies = [ + "rustls", + "tokio", +] + [[package]] name = "tokio-util" version = "0.7.18" @@ -987,6 +1695,56 @@ dependencies = [ "tokio", ] +[[package]] +name = "tower" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" +dependencies = [ + "futures-core", + "futures-util", + "pin-project-lite", + "sync_wrapper", + "tokio", + "tower-layer", + "tower-service", +] + +[[package]] +name = "tower-http" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" +dependencies = [ + "async-compression", + "bitflags", + "bytes", + "futures-core", + "futures-util", + "http", + "http-body", + "http-body-util", + "pin-project-lite", + "tokio", + "tokio-util", + "tower", + "tower-layer", + "tower-service", + "url", +] + +[[package]] +name = "tower-layer" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "121c2a6cda46980bb0fcd1647ffaf6cd3fc79a013de288782836f6df9c48780e" + +[[package]] +name = "tower-service" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8df9b6e13f2d32c91b9bd719c00d1958837bc7dec474d94952798cc8e69eeec3" + [[package]] name = "tracing" version = "0.1.44" @@ -1018,13 +1776,19 @@ dependencies = [ "once_cell", ] +[[package]] +name = "try-lock" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" + [[package]] name = "ulid" version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "470dbf6591da1b39d43c14523b2b469c86879a53e8b758c8e090a470fe7b1fbe" dependencies = [ - "rand", + "rand 0.9.4", "serde", "web-time", ] @@ -1035,6 +1799,30 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "untrusted" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" + +[[package]] +name = "url" +version = "2.5.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff67a8a4397373c3ef660812acab3268222035010ab8680ec4215f38ba3d0eed" +dependencies = [ + "form_urlencoded", + "idna", + "percent-encoding", + "serde", +] + +[[package]] +name = "utf8_iter" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" + [[package]] name = "vcpkg" version = "0.2.15" @@ -1057,6 +1845,15 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "want" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bfa7760aed19e106de2c7c0b581b509f2f25d3dacaf737cb82ac61bc6d760b0e" +dependencies = [ + "try-lock", +] + [[package]] name = "wasi" version = "0.11.1+wasi-snapshot-preview1" @@ -1085,6 +1882,16 @@ dependencies = [ "wasm-bindgen-shared", ] +[[package]] +name = "wasm-bindgen-futures" +version = "0.4.76" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c62df1340f32221cb9c54d6a27b030e3dba64361d4a95bed55f9aacb44da291d" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "wasm-bindgen-macro" version = "0.2.126" @@ -1117,6 +1924,29 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "wasm-streams" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15053d8d85c7eccdbefef60f06769760a563c7f0a9d6902a13d35c7800b0ad65" +dependencies = [ + "futures-util", + "js-sys", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", +] + +[[package]] +name = "web-sys" +version = "0.3.103" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8622dcb61c0bcc9fffa6938bed81210af2da9a7e4a1a834b2e37a59b6dfb6141" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "web-time" version = "1.1.0" @@ -1127,6 +1957,15 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "webpki-roots" +version = "1.0.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf85cb06032201fa7c6f829d7db5a7e5aa45bcc0655327713065f6f0576731bf" +dependencies = [ + "rustls-pki-types", +] + [[package]] name = "winapi-util" version = "0.1.11" @@ -1148,7 +1987,16 @@ version = "0.48.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "677d2418bec65e3338edb076e806bc1ec15693c5d0104683f2efe857f61056a9" dependencies = [ - "windows-targets", + "windows-targets 0.48.5", +] + +[[package]] +name = "windows-sys" +version = "0.52.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "282be5f36a8ce781fad8c8ae18fa3f9beff57ec1b52cb3de0789201425d9a33d" +dependencies = [ + "windows-targets 0.52.6", ] [[package]] @@ -1166,13 +2014,29 @@ version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9a2fa6e2155d7247be68c096456083145c183cbbbc2764150dda45a87197940c" dependencies = [ - "windows_aarch64_gnullvm", - "windows_aarch64_msvc", - "windows_i686_gnu", - "windows_i686_msvc", - "windows_x86_64_gnu", - "windows_x86_64_gnullvm", - "windows_x86_64_msvc", + "windows_aarch64_gnullvm 0.48.5", + "windows_aarch64_msvc 0.48.5", + "windows_i686_gnu 0.48.5", + "windows_i686_msvc 0.48.5", + "windows_x86_64_gnu 0.48.5", + "windows_x86_64_gnullvm 0.48.5", + "windows_x86_64_msvc 0.48.5", +] + +[[package]] +name = "windows-targets" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b724f72796e036ab90c1021d4780d4d3d648aca59e491e6b98e725b84e99973" +dependencies = [ + "windows_aarch64_gnullvm 0.52.6", + "windows_aarch64_msvc 0.52.6", + "windows_i686_gnu 0.52.6", + "windows_i686_gnullvm", + "windows_i686_msvc 0.52.6", + "windows_x86_64_gnu 0.52.6", + "windows_x86_64_gnullvm 0.52.6", + "windows_x86_64_msvc 0.52.6", ] [[package]] @@ -1181,48 +2045,125 @@ version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2b38e32f0abccf9987a4e3079dfb67dcd799fb61361e53e2882c3cbaf0d905d8" +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" + [[package]] name = "windows_aarch64_msvc" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc" +[[package]] +name = "windows_aarch64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" + [[package]] name = "windows_i686_gnu" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e" +[[package]] +name = "windows_i686_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e9b5ad5ab802e97eb8e295ac6720e509ee4c243f69d781394014ebfe8bbfa0b" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" + [[package]] name = "windows_i686_msvc" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406" +[[package]] +name = "windows_i686_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" + [[package]] name = "windows_x86_64_gnu" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "53d40abd2583d23e4718fddf1ebec84dbff8381c07cae67ff7768bbf19c6718e" +[[package]] +name = "windows_x86_64_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" + [[package]] name = "windows_x86_64_gnullvm" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b7b52767868a23d5bab768e390dc5f5c55825b6d30b86c844ff2dc7414044cc" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" + [[package]] name = "windows_x86_64_msvc" version = "0.48.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed94fce61571a4006852b7389a063ab983c02eb1bb37b47f8272ce92d06d9538" +[[package]] +name = "windows_x86_64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" + [[package]] name = "wit-bindgen" version = "0.57.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" +[[package]] +name = "writeable" +version = "0.6.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ffae5123b2d3fc086436f8834ae3ab053a283cfac8fe0a0b8eaae044768a4c4" + +[[package]] +name = "yoke" +version = "0.8.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "709fe23a0424b6a435d82152b1bd3fdfb0833487d5fa90d05d42762a9891fef5" +dependencies = [ + "stable_deref_trait", + "yoke-derive", + "zerofrom", +] + +[[package]] +name = "yoke-derive" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "de844c262c8848816172cef550288e7dc6c7b7814b4ee56b3e1553f275f1858e" +dependencies = [ + "proc-macro2", + "quote", + "syn", + "synstructure", +] + [[package]] name = "zerocopy" version = "0.8.53" @@ -1243,6 +2184,66 @@ dependencies = [ "syn", ] +[[package]] +name = "zerofrom" +version = "0.1.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ec05a11813ea801ff6d75110ad09cd0824ddba17dfe17128ea0d5f68e6c5272" +dependencies = [ + "zerofrom-derive", +] + +[[package]] +name = "zerofrom-derive" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "11532158c46691caf0f2593ea8358fed6bbf68a0315e80aae9bd41fbade684a1" +dependencies = [ + "proc-macro2", + "quote", + "syn", + "synstructure", +] + +[[package]] +name = "zeroize" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e" + +[[package]] +name = "zerotrie" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0f9152d31db0792fa83f70fb2f83148effb5c1f5b8c7686c3459e361d9bc20bf" +dependencies = [ + "displaydoc", + "yoke", + "zerofrom", +] + +[[package]] +name = "zerovec" +version = "0.11.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "90f911cbc359ab6af17377d242225f4d75119aec87ea711a880987b18cd7b239" +dependencies = [ + "yoke", + "zerofrom", + "zerovec-derive", +] + +[[package]] +name = "zerovec-derive" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "625dc425cab0dca6dc3c3319506e6593dcb08a9f387ea3b284dbd52a92c40555" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "zmij" version = "1.0.21" diff --git a/Cargo.toml b/Cargo.toml index e901c8c..bbd1443 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -59,8 +59,7 @@ dirs = "5" shell-words = "1" libc = "0.2" tempfile = "3" -insta = "1" -wiremock = "0.6" +bytes = "1" [workspace.lints.rust] unsafe_code = "deny" diff --git a/crates/harness-core/src/engine/processor.rs b/crates/harness-core/src/engine/processor.rs index e66d81c..a4ff597 100644 --- a/crates/harness-core/src/engine/processor.rs +++ b/crates/harness-core/src/engine/processor.rs @@ -574,18 +574,18 @@ pub async fn process_step( let mut result = StepResult::Stop; loop { - let cancelled = ctx.cancel.is_cancelled(); - if cancelled { - run.mark_errored(&ProviderError::Cancelled).await; - return Ok(StepOutcome { - result: StepResult::Stop, - message_id: run.assistant.as_ref().map(|m| m.id.clone()), - usage, - aborted: true, - }); - } - - let item = stream.next().await; + let item = tokio::select! { + item = stream.next() => item, + _ = ctx.cancel.cancelled() => { + run.mark_errored(&ProviderError::Cancelled).await; + return Ok(StepOutcome { + result: StepResult::Stop, + message_id: run.assistant.as_ref().map(|m| m.id.clone()), + usage, + aborted: true, + }); + } + }; let Some(item) = item else { break }; let event = match item { diff --git a/crates/harness-providers/Cargo.toml b/crates/harness-providers/Cargo.toml index 5a673d1..9972f4c 100644 --- a/crates/harness-providers/Cargo.toml +++ b/crates/harness-providers/Cargo.toml @@ -6,6 +6,21 @@ license.workspace = true [dependencies] harness-core = { workspace = true } +tokio = { workspace = true } +tokio-util = { workspace = true } +futures = { workspace = true } +async-stream = { workspace = true } +async-trait = { workspace = true } +reqwest = { workspace = true } +eventsource-stream = { workspace = true } +bytes = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +thiserror = { workspace = true } +tracing = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["test-util", "macros"] } [lints] workspace = true diff --git a/crates/harness-providers/src/anthropic.rs b/crates/harness-providers/src/anthropic.rs new file mode 100644 index 0000000..0512ee7 --- /dev/null +++ b/crates/harness-providers/src/anthropic.rs @@ -0,0 +1,157 @@ +use std::time::Duration; + +use async_trait::async_trait; +use harness_core::llm::{LlmEventStream, LlmRequest, Provider, ProviderError}; +use harness_core::types::ModelInfo; +use tokio_util::sync::CancellationToken; + +use crate::codec::anthropic::{build_request, decode, initiator_header}; + +const DEFAULT_BASE_URL: &str = "https://api.anthropic.com"; +const ANTHROPIC_VERSION: &str = "2023-06-01"; +const ANTHROPIC_BETA: &str = + "interleaved-thinking-2025-05-14,fine-grained-tool-streaming-2025-05-14"; + +pub struct AnthropicProvider { + api_key: String, + base_url: String, + client: reqwest::Client, +} + +impl AnthropicProvider { + pub fn new(api_key: impl Into) -> Self { + Self { + api_key: api_key.into(), + base_url: DEFAULT_BASE_URL.to_string(), + client: reqwest::Client::new(), + } + } + + #[cfg(test)] + fn with_base_url(api_key: impl Into, base_url: impl Into) -> Self { + Self { + api_key: api_key.into(), + base_url: base_url.into(), + client: reqwest::Client::new(), + } + } + + fn classify_error( + status: reqwest::StatusCode, + body: String, + retry_after: Option, + ) -> ProviderError { + match status.as_u16() { + 401 | 403 => ProviderError::Auth(body), + 429 => ProviderError::RateLimited { retry_after }, + 529 => ProviderError::Overloaded, + s if (500..600).contains(&s) => ProviderError::Overloaded, + s => ProviderError::Http { status: s, body }, + } + } +} + +#[async_trait] +impl Provider for AnthropicProvider { + fn id(&self) -> &str { + "anthropic" + } + + async fn list_models(&self) -> Result, ProviderError> { + // Full models.dev metadata integration lands in M3; callers fall back to + // config-specified model strings until then. + Ok(Vec::new()) + } + + async fn stream( + &self, + req: LlmRequest, + cancel: CancellationToken, + ) -> Result { + let body = build_request(&req); + let url = format!("{}/v1/messages", self.base_url); + let initiator = initiator_header(req.initiator); + + let send = self + .client + .post(&url) + .header("x-api-key", &self.api_key) + .header("anthropic-version", ANTHROPIC_VERSION) + .header("anthropic-beta", ANTHROPIC_BETA) + .header("x-initiator", initiator) + .json(&body) + .send(); + + let response = tokio::select! { + result = send => result.map_err(|e| ProviderError::Network(e.to_string()))?, + _ = cancel.cancelled() => return Err(ProviderError::Cancelled), + }; + + if !response.status().is_success() { + let status = response.status(); + let retry_after = response + .headers() + .get("retry-after") + .and_then(|v| v.to_str().ok()) + .and_then(|s| s.parse::().ok()) + .map(Duration::from_secs); + let body = response.text().await.unwrap_or_default(); + return Err(Self::classify_error(status, body, retry_after)); + } + + Ok(decode(response.bytes_stream())) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn classifies_status_codes() { + assert!(matches!( + AnthropicProvider::classify_error( + reqwest::StatusCode::TOO_MANY_REQUESTS, + String::new(), + None + ), + ProviderError::RateLimited { .. } + )); + assert!(matches!( + AnthropicProvider::classify_error( + reqwest::StatusCode::UNAUTHORIZED, + String::new(), + None + ), + ProviderError::Auth(_) + )); + assert!(matches!( + AnthropicProvider::classify_error( + reqwest::StatusCode::SERVICE_UNAVAILABLE, + String::new(), + None + ), + ProviderError::Overloaded + )); + assert!(matches!( + AnthropicProvider::classify_error( + reqwest::StatusCode::BAD_REQUEST, + String::new(), + None + ), + ProviderError::Http { status: 400, .. } + )); + } + + #[test] + fn id_is_anthropic() { + let provider = AnthropicProvider::new("test-key"); + assert_eq!(provider.id(), "anthropic"); + } + + #[test] + fn constructs_with_custom_base_url() { + let provider = AnthropicProvider::with_base_url("test-key", "http://localhost:1234"); + assert_eq!(provider.base_url, "http://localhost:1234"); + } +} diff --git a/crates/harness-providers/src/codec/anthropic.rs b/crates/harness-providers/src/codec/anthropic.rs new file mode 100644 index 0000000..ebcb04c --- /dev/null +++ b/crates/harness-providers/src/codec/anthropic.rs @@ -0,0 +1,499 @@ +//! Request builder + SSE decoder for Anthropic's `/v1/messages` streaming API. + +use std::collections::HashMap; + +use async_stream::try_stream; +use eventsource_stream::Eventsource; +use futures::Stream; +use harness_core::llm::{ + FinishReason, Initiator, LlmEvent, LlmEventStream, LlmRequest, ProviderError, Role, WireContent, +}; +use harness_core::types::TokenUsage; +use serde_json::{json, Value}; + +const DEFAULT_MAX_TOKENS: u32 = 8192; +/// opencode's `applyCaching`: cache breakpoints on the first 2 system blocks + last 2 +/// non-system messages. +const CACHE_BREAKPOINTS: usize = 2; + +fn cache_control() -> Value { + json!({"type": "ephemeral"}) +} + +fn build_system(system: &[String]) -> Vec { + system + .iter() + .enumerate() + .map(|(i, block)| { + let mut b = json!({"type": "text", "text": block}); + if i < CACHE_BREAKPOINTS { + b["cache_control"] = cache_control(); + } + b + }) + .collect() +} + +fn wire_role_to_anthropic(role: Role) -> &'static str { + match role { + Role::User | Role::Tool => "user", + Role::Assistant => "assistant", + Role::System => "user", // system blocks are carried separately in `system`, not here + } +} + +fn build_messages(messages: &[harness_core::llm::WireMessage]) -> Vec { + let mut built: Vec = messages + .iter() + .map(|m| { + let content: Vec = m + .content + .iter() + .map(|c| match c { + WireContent::Text { text } => json!({"type": "text", "text": text}), + WireContent::ToolCall { call_id, name, input } => { + json!({"type": "tool_use", "id": call_id, "name": name, "input": input}) + } + WireContent::ToolResult { call_id, output, is_error } => { + json!({"type": "tool_result", "tool_use_id": call_id, "content": output, "is_error": is_error}) + } + WireContent::Image { mime_type, data } => { + json!({"type": "image", "source": {"type": "base64", "media_type": mime_type, "data": data}}) + } + }) + .collect(); + json!({"role": wire_role_to_anthropic(m.role), "content": content}) + }) + .collect(); + + let n = built.len(); + for msg in built.iter_mut().skip(n.saturating_sub(CACHE_BREAKPOINTS)) { + if let Some(content) = msg["content"].as_array_mut() { + if let Some(last) = content.last_mut() { + last["cache_control"] = cache_control(); + } + } + } + built +} + +pub fn build_request(req: &LlmRequest) -> Value { + let mut body = json!({ + "model": req.model, + "system": build_system(&req.system), + "messages": build_messages(&req.messages), + "max_tokens": req.max_tokens.unwrap_or(DEFAULT_MAX_TOKENS), + "stream": true, + }); + + if !req.tools.is_empty() { + body["tools"] = Value::Array( + req.tools + .iter() + .map(|t| json!({"name": t.name, "description": t.description, "input_schema": t.parameters})) + .collect(), + ); + } + if let Some(temp) = req.temperature { + body["temperature"] = json!(temp); + } + if let Some(reasoning) = &req.reasoning { + if let Some(budget) = reasoning.budget_tokens { + body["thinking"] = json!({"type": "enabled", "budget_tokens": budget}); + } + } + body +} + +pub fn initiator_header(initiator: Initiator) -> &'static str { + match initiator { + Initiator::User => "user", + Initiator::Agent => "agent", + } +} + +fn map_stop_reason(reason: Option<&str>) -> FinishReason { + match reason { + Some("end_turn") | Some("stop_sequence") => FinishReason::Stop, + Some("tool_use") => FinishReason::ToolCalls, + Some("max_tokens") => FinishReason::Length, + Some(other) => FinishReason::Unknown(other.to_string()), + None => FinishReason::Unknown("none".to_string()), + } +} + +#[derive(Default)] +struct BlockState { + kind: BlockKind, + call_id: String, + name: String, + partial_json: String, + signature: Option, +} + +#[derive(Default, PartialEq, Eq)] +enum BlockKind { + #[default] + Text, + Thinking, + ToolUse, +} + +/// Decodes a raw SSE byte stream into our normalized `LlmEvent` stream. Errors from the +/// underlying HTTP stream and any Anthropic `error` event both surface as `Err`. +pub fn decode(byte_stream: S) -> LlmEventStream +where + S: Stream> + Send + 'static, + E: std::error::Error + Send + Sync + 'static, +{ + let events = byte_stream.eventsource(); + Box::pin(try_stream! { + futures::pin_mut!(events); + let mut blocks: HashMap = HashMap::new(); + let mut usage = TokenUsage::default(); + + while let Some(item) = futures::StreamExt::next(&mut events).await { + let event = item.map_err(|e| ProviderError::Decode(e.to_string()))?; + if event.data.is_empty() { + continue; + } + let value: Value = serde_json::from_str(&event.data) + .map_err(|e| ProviderError::Decode(format!("{e}: {}", event.data)))?; + let kind = value["type"].as_str().unwrap_or_default(); + + match kind { + "message_start" => { + let u = &value["message"]["usage"]; + usage.input = u["input_tokens"].as_u64().unwrap_or(0); + usage.cache_write = u["cache_creation_input_tokens"].as_u64().unwrap_or(0); + usage.cache_read = u["cache_read_input_tokens"].as_u64().unwrap_or(0); + } + "content_block_start" => { + let index = value["index"].as_u64().unwrap_or(0); + let block = &value["content_block"]; + match block["type"].as_str().unwrap_or_default() { + "text" => { + blocks.insert(index, BlockState { kind: BlockKind::Text, ..Default::default() }); + yield LlmEvent::TextStart { id: index.to_string() }; + } + "thinking" => { + blocks.insert(index, BlockState { kind: BlockKind::Thinking, ..Default::default() }); + yield LlmEvent::ReasoningStart { id: index.to_string() }; + } + "tool_use" => { + let call_id = block["id"].as_str().unwrap_or_default().to_string(); + let name = block["name"].as_str().unwrap_or_default().to_string(); + blocks.insert(index, BlockState { + kind: BlockKind::ToolUse, + call_id: call_id.clone(), + name: name.clone(), + ..Default::default() + }); + yield LlmEvent::ToolInputStart { call_id, name }; + } + _ => {} + } + } + "content_block_delta" => { + let index = value["index"].as_u64().unwrap_or(0); + let delta = &value["delta"]; + match delta["type"].as_str().unwrap_or_default() { + "text_delta" => { + let text = delta["text"].as_str().unwrap_or_default().to_string(); + yield LlmEvent::TextDelta { id: index.to_string(), text }; + } + "thinking_delta" => { + let text = delta["thinking"].as_str().unwrap_or_default().to_string(); + yield LlmEvent::ReasoningDelta { id: index.to_string(), text }; + } + "signature_delta" => { + if let Some(state) = blocks.get_mut(&index) { + state.signature = Some(delta["signature"].as_str().unwrap_or_default().to_string()); + } + } + "input_json_delta" => { + let partial = delta["partial_json"].as_str().unwrap_or_default(); + if let Some(state) = blocks.get_mut(&index) { + state.partial_json.push_str(partial); + yield LlmEvent::ToolInputDelta { call_id: state.call_id.clone(), json: partial.to_string() }; + } + } + _ => {} + } + } + "content_block_stop" => { + let index = value["index"].as_u64().unwrap_or(0); + if let Some(state) = blocks.remove(&index) { + match state.kind { + BlockKind::Text => yield LlmEvent::TextEnd { id: index.to_string() }, + BlockKind::Thinking => { + yield LlmEvent::ReasoningEnd { id: index.to_string(), signature: state.signature }; + } + BlockKind::ToolUse => { + let input: Value = if state.partial_json.trim().is_empty() { + json!({}) + } else { + serde_json::from_str(&state.partial_json).unwrap_or(Value::Null) + }; + yield LlmEvent::ToolCall { call_id: state.call_id, name: state.name, input }; + } + } + } + } + "message_delta" => { + if let Some(out) = value["usage"]["output_tokens"].as_u64() { + usage.output = out; + } + let stop_reason = value["delta"]["stop_reason"].as_str(); + yield LlmEvent::Finish { reason: map_stop_reason(stop_reason), usage }; + } + "error" => { + let message = value["error"]["message"].as_str().unwrap_or("unknown error").to_string(); + let err_type = value["error"]["type"].as_str().unwrap_or(""); + Err(match err_type { + "overloaded_error" => ProviderError::Overloaded, + "rate_limit_error" => ProviderError::RateLimited { retry_after: None }, + "authentication_error" | "permission_error" => ProviderError::Auth(message), + _ => ProviderError::Http { status: 0, body: message }, + })?; + } + _ => {} // ping, message_stop: nothing to emit + } + } + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use futures::StreamExt; + use harness_core::llm::{ReasoningOpts, ToolSchema, WireMessage}; + + fn sse_stream(raw: &'static str) -> LlmEventStream { + let chunks: Vec> = + vec![Ok(bytes::Bytes::from_static(raw.as_bytes()))]; + decode(futures::stream::iter(chunks)) + } + + #[tokio::test] + async fn decodes_text_only_response() { + let raw = concat!( + "event: message_start\n", + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10}}}\n\n", + "event: content_block_start\n", + "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n", + "event: content_block_delta\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n\n", + "event: content_block_stop\n", + "data: {\"type\":\"content_block_stop\",\"index\":0}\n\n", + "event: message_delta\n", + "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"output_tokens\":5}}\n\n", + "event: message_stop\n", + "data: {\"type\":\"message_stop\"}\n\n", + ); + let events: Vec = sse_stream(raw).map(|e| e.unwrap()).collect().await; + assert_eq!( + events, + vec![ + LlmEvent::TextStart { id: "0".into() }, + LlmEvent::TextDelta { + id: "0".into(), + text: "Hello".into() + }, + LlmEvent::TextEnd { id: "0".into() }, + LlmEvent::Finish { + reason: FinishReason::Stop, + usage: TokenUsage { + input: 10, + output: 5, + ..Default::default() + }, + }, + ] + ); + } + + #[tokio::test] + async fn decodes_tool_call_with_streamed_json_input() { + let raw = concat!( + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":1}}}\n\n", + "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"file\\\"\"}}\n\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\":\\\"a.txt\\\"}\"}}\n\n", + "data: {\"type\":\"content_block_stop\",\"index\":0}\n\n", + "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":8}}\n\n", + ); + let events: Vec = sse_stream(raw).map(|e| e.unwrap()).collect().await; + assert_eq!( + events, + vec![ + LlmEvent::ToolInputStart { + call_id: "call_1".into(), + name: "read".into() + }, + LlmEvent::ToolInputDelta { + call_id: "call_1".into(), + json: "{\"file\"".into() + }, + LlmEvent::ToolInputDelta { + call_id: "call_1".into(), + json: ":\"a.txt\"}".into() + }, + LlmEvent::ToolCall { + call_id: "call_1".into(), + name: "read".into(), + input: json!({"file": "a.txt"}), + }, + LlmEvent::Finish { + reason: FinishReason::ToolCalls, + usage: TokenUsage { + input: 1, + output: 8, + ..Default::default() + }, + }, + ] + ); + } + + #[tokio::test] + async fn decodes_thinking_block_with_signature() { + let raw = concat!( + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":1}}}\n\n", + "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\"}}\n\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"pondering\"}}\n\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"signature_delta\",\"signature\":\"sig123\"}}\n\n", + "data: {\"type\":\"content_block_stop\",\"index\":0}\n\n", + "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"output_tokens\":2}}\n\n", + ); + let events: Vec = sse_stream(raw).map(|e| e.unwrap()).collect().await; + assert_eq!( + events, + vec![ + LlmEvent::ReasoningStart { id: "0".into() }, + LlmEvent::ReasoningDelta { + id: "0".into(), + text: "pondering".into() + }, + LlmEvent::ReasoningEnd { + id: "0".into(), + signature: Some("sig123".into()) + }, + LlmEvent::Finish { + reason: FinishReason::Stop, + usage: TokenUsage { + input: 1, + output: 2, + ..Default::default() + }, + }, + ] + ); + } + + #[tokio::test] + async fn mid_stream_error_event_surfaces_as_err() { + let raw = concat!( + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":1}}}\n\n", + "data: {\"type\":\"error\",\"error\":{\"type\":\"overloaded_error\",\"message\":\"overloaded\"}}\n\n", + ); + let events: Vec> = sse_stream(raw).collect().await; + assert!(matches!( + events.last(), + Some(Err(ProviderError::Overloaded)) + )); + } + + #[test] + fn build_request_applies_cache_control_to_first_two_system_blocks() { + let req = LlmRequest { + model: "claude-sonnet".into(), + system: vec!["env".into(), "agent".into(), "instructions".into()], + messages: vec![], + tools: vec![], + temperature: None, + max_tokens: None, + reasoning: None, + initiator: Initiator::User, + }; + let body = build_request(&req); + let system = body["system"].as_array().unwrap(); + assert_eq!(system.len(), 3); + assert!(system[0]["cache_control"].is_object()); + assert!(system[1]["cache_control"].is_object()); + assert!(system[2].get("cache_control").is_none()); + } + + #[test] + fn build_request_applies_cache_control_to_last_two_messages() { + let messages = vec![ + WireMessage { + role: Role::User, + content: vec![WireContent::Text { text: "1".into() }], + }, + WireMessage { + role: Role::Assistant, + content: vec![WireContent::Text { text: "2".into() }], + }, + WireMessage { + role: Role::User, + content: vec![WireContent::Text { text: "3".into() }], + }, + ]; + let req = LlmRequest { + model: "claude-sonnet".into(), + system: vec![], + messages, + tools: vec![], + temperature: None, + max_tokens: None, + reasoning: None, + initiator: Initiator::User, + }; + let body = build_request(&req); + let msgs = body["messages"].as_array().unwrap(); + assert!(msgs[0]["content"][0].get("cache_control").is_none()); + assert!(msgs[1]["content"][0]["cache_control"].is_object()); + assert!(msgs[2]["content"][0]["cache_control"].is_object()); + } + + #[test] + fn build_request_maps_tool_schema_to_input_schema_key() { + let req = LlmRequest { + model: "claude-sonnet".into(), + system: vec![], + messages: vec![], + tools: vec![ToolSchema { + name: "read".into(), + description: "reads a file".into(), + parameters: json!({"type": "object"}), + }], + temperature: None, + max_tokens: None, + reasoning: None, + initiator: Initiator::User, + }; + let body = build_request(&req); + assert_eq!(body["tools"][0]["name"], "read"); + assert_eq!(body["tools"][0]["input_schema"], json!({"type": "object"})); + } + + #[test] + fn build_request_includes_thinking_budget_when_reasoning_set() { + let req = LlmRequest { + model: "claude-sonnet".into(), + system: vec![], + messages: vec![], + tools: vec![], + temperature: None, + max_tokens: None, + reasoning: Some(ReasoningOpts { + effort: None, + budget_tokens: Some(2048), + }), + initiator: Initiator::User, + }; + let body = build_request(&req); + assert_eq!(body["thinking"]["budget_tokens"], 2048); + } +} diff --git a/crates/harness-providers/src/codec/mod.rs b/crates/harness-providers/src/codec/mod.rs new file mode 100644 index 0000000..e529997 --- /dev/null +++ b/crates/harness-providers/src/codec/mod.rs @@ -0,0 +1 @@ +pub mod anthropic; diff --git a/crates/harness-providers/src/lib.rs b/crates/harness-providers/src/lib.rs index 8ff030f..63c8f2d 100644 --- a/crates/harness-providers/src/lib.rs +++ b/crates/harness-providers/src/lib.rs @@ -1 +1,6 @@ -// Provider implementations (Copilot, OpenAI, Anthropic) land here in M1/M3. +pub mod anthropic; +pub mod codec; +pub mod registry; + +pub use anthropic::AnthropicProvider; +pub use registry::ProviderRegistry; diff --git a/crates/harness-providers/src/registry.rs b/crates/harness-providers/src/registry.rs new file mode 100644 index 0000000..5e3db43 --- /dev/null +++ b/crates/harness-providers/src/registry.rs @@ -0,0 +1,57 @@ +use std::collections::HashMap; +use std::sync::Arc; + +use harness_core::llm::Provider; + +#[derive(Default, Clone)] +pub struct ProviderRegistry { + providers: HashMap>, +} + +impl ProviderRegistry { + pub fn new() -> Self { + Self::default() + } + + pub fn register(&mut self, provider: Arc) { + self.providers.insert(provider.id().to_string(), provider); + } + + pub fn get(&self, id: &str) -> Option> { + self.providers.get(id).cloned() + } + + /// Resolves a `"provider/model"` string to its provider and bare model id. + pub fn resolve(&self, model_ref: &str) -> Option<(Arc, String)> { + let (provider_id, model_id) = model_ref.split_once('/')?; + self.get(provider_id).map(|p| (p, model_id.to_string())) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::anthropic::AnthropicProvider; + + #[test] + fn resolves_provider_slash_model() { + let mut registry = ProviderRegistry::new(); + registry.register(Arc::new(AnthropicProvider::new("key"))); + let (provider, model) = registry.resolve("anthropic/claude-sonnet-4-5").unwrap(); + assert_eq!(provider.id(), "anthropic"); + assert_eq!(model, "claude-sonnet-4-5"); + } + + #[test] + fn unknown_provider_resolves_to_none() { + let registry = ProviderRegistry::new(); + assert!(registry.resolve("openai/gpt-4o").is_none()); + } + + #[test] + fn malformed_model_ref_resolves_to_none() { + let mut registry = ProviderRegistry::new(); + registry.register(Arc::new(AnthropicProvider::new("key"))); + assert!(registry.resolve("no-slash-here").is_none()); + } +}