From fa60458d5a6f07cfc11f11e146c8f6142788cd96 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jakub=20Koz=C5=82owski?= Date: Thu, 1 Oct 2026 18:54:36 +0200 Subject: [PATCH 1/3] Support Cloudflare Workers AI (Clef) as a provider JevConfig gets a provider: TypeSafe (default) or WorkersAI(accountId), which posts to accounts/{id}/ai/run/@cf/cloudflare/{model} and reads the answers from Cloudflare's result envelope. Listing models isn't supported there. Adds JevConfig.workersAI/workersAIFromEnv and the ModelId.clef/clefFlash constants. A Noul without criteria now omits the field instead of sending null, which Clef rejects. The triage example's core now takes a Jev, shared by Main and a new ClefTriage entry point. Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 11 +++ core/src/main/scala/jev4s/Jev.scala | 84 +++++++++++++++++-- core/src/main/scala/jev4s/internal/wire.scala | 10 ++- core/src/main/scala/jev4s/model.scala | 6 ++ core/src/test/scala/example/Triage.scala | 66 +++++++++------ core/src/test/scala/jev4s/JevTests.scala | 47 ++++++++++- 6 files changed, 184 insertions(+), 40 deletions(-) diff --git a/README.md b/README.md index 1ee12c3..760eedf 100644 --- a/README.md +++ b/README.md @@ -130,6 +130,17 @@ You can also construct a `JevConfig` directly. Requests that fail with 408, 429 To use a different model for some requests, use `jev.withModel(ModelId("jev-1.13.0"))`. `jev.models` lists the available models. +### Cloudflare Workers AI (Clef) + +Cloudflare's [Clef models](https://developers.cloudflare.com/workers-ai/models/clef/) speak the same System One format, so the same questions work against them: + +```scala +val config = JevConfig.workersAI(accountId, ApiKey(apiToken)) // or JevConfig.workersAIFromEnv[IO] +val jev = Jev.instance[IO](config, client).withModel(ModelId.clefFlash) +``` + +`JevConfig.workersAIFromEnv` reads `CLOUDFLARE_ACCOUNT_ID`, `CLOUDFLARE_AUTH_TOKEN` and optionally `CLOUDFLARE_MODEL` (`clef` by default). The token needs the "Workers AI - Read" and "Workers AI - Edit" permissions. `jev.models` isn't supported on Workers AI. `example.ClefTriage` runs the triage example against Clef. + ## License Licensed under the [Apache License, Version 2.0](LICENSE). diff --git a/core/src/main/scala/jev4s/Jev.scala b/core/src/main/scala/jev4s/Jev.scala index 919138e..55bb793 100644 --- a/core/src/main/scala/jev4s/Jev.scala +++ b/core/src/main/scala/jev4s/Jev.scala @@ -63,10 +63,36 @@ final case class JevConfig( baseUri: Uri = uri"https://api.typesafe.ai", model: ModelId = ModelId.latest, retry: JevConfig.RetryConfig = JevConfig.RetryConfig.default, + provider: JevConfig.Provider = JevConfig.Provider.TypeSafe, ) object JevConfig { + /** Who serves the System One API. The request and answer formats are the same, the endpoints + * differ. + */ + enum Provider { + + /** TypeSafe's own API: `POST v1/systemone`. */ + case TypeSafe + + /** Cloudflare Workers AI (the Clef models): + * `POST accounts/{accountId}/ai/run/@cf/cloudflare/{model}`. Model listing isn't supported. + */ + case WorkersAI(accountId: String) + } + + /** Clef on Cloudflare Workers AI. `apiToken` needs the "Workers AI - Read" and "Workers AI - + * Edit" permissions. + */ + def workersAI(accountId: String, apiToken: ApiKey, model: ModelId = ModelId.clef): JevConfig = + JevConfig( + apiKey = apiToken, + baseUri = uri"https://api.cloudflare.com/client/v4", + model = model, + provider = Provider.WorkersAI(accountId), + ) + final case class RetryConfig(maxRetries: Int, maxBackoff: FiniteDuration) object RetryConfig { @@ -90,10 +116,29 @@ object JevConfig { ) } + /** Reads `CLOUDFLARE_ACCOUNT_ID` and `CLOUDFLARE_AUTH_TOKEN` (both required) and + * `CLOUDFLARE_MODEL` (`clef` by default). + */ + def workersAIFromEnv[F[_]: Env: Concurrent]: F[JevConfig] = + ( + Env[F] + .get("CLOUDFLARE_ACCOUNT_ID") + .flatMap(_.liftTo[F](JevError.MissingEnv("CLOUDFLARE_ACCOUNT_ID"))), + Env[F] + .get("CLOUDFLARE_AUTH_TOKEN") + .flatMap(_.liftTo[F](JevError.MissingEnv("CLOUDFLARE_AUTH_TOKEN"))), + Env[F].get("CLOUDFLARE_MODEL"), + ).mapN { (accountId, token, model) => + workersAI(accountId, ApiKey(token), model.fold(ModelId.clef)(ModelId(_))) + } + } enum JevError(message: String) extends Exception(message) { case MissingApiKey extends JevError("TYPESAFE_API_KEY is not set") + case MissingEnv(name: String) extends JevError(s"$name is not set") + case Unsupported(operation: String) + extends JevError(s"Not supported by this provider: $operation") case InvalidQuestion(reason: String) extends JevError(s"Invalid question: $reason") case Unauthorized(body: String) extends JevError(s"Unauthorized: $body") case Unprocessable(body: Json) extends JevError(s"Request failed validation: ${body.noSpaces}") @@ -140,8 +185,12 @@ object Jev { .leftMap(JevError.InvalidQuestion(_)) .liftTo[F] response <- client - .run(authorized(Method.POST, "v1/systemone").withEntity(body)) - .use(decodeOrFail[ResponseBody]) + .run(authorized(Method.POST, evaluateUri).withEntity(body)) + .use(r => + decodeOrFail(r)( + using responseDecoder + ) + ) evaluation <- Question .decodeResponse(question, response) .leftMap(JevError.UnexpectedAnswer(_)) @@ -149,15 +198,34 @@ object Jev { } yield evaluation def models: F[List[ModelCard]] = - client - .run(authorized(Method.GET, "v1/models")) - .use(decodeOrFail[ModelList]) - .map(_.models) + config.provider match { + case JevConfig.Provider.TypeSafe => + client + .run(authorized(Method.GET, config.baseUri / "v1" / "models")) + .use(decodeOrFail[ModelList]) + .map(_.models) + case JevConfig.Provider.WorkersAI(_) => JevError.Unsupported("listing models").raiseError + } def withModel(model: ModelId): Jev[F] = JevImpl(config.copy(model = model), client) - private def authorized(method: Method, path: String): Request[F] = - Request[F](method, config.baseUri.addPath(path)) + private def evaluateUri: Uri = + config.provider match { + case JevConfig.Provider.TypeSafe => config.baseUri / "v1" / "systemone" + case JevConfig.Provider.WorkersAI(accountId) => + config.baseUri / "accounts" / accountId / "ai" / "run" / "@cf" / "cloudflare" / + config.model.value + } + + // Workers AI wraps the answers in Cloudflare's {"result": ..., "success": ...} envelope. + private def responseDecoder: Decoder[ResponseBody] = + config.provider match { + case JevConfig.Provider.TypeSafe => Decoder[ResponseBody] + case JevConfig.Provider.WorkersAI(_) => Decoder[ResponseBody].at("result") + } + + private def authorized(method: Method, uri: Uri): Request[F] = + Request[F](method, uri) .putHeaders(Authorization(Credentials.Token(AuthScheme.Bearer, config.apiKey.value))) private def decodeOrFail[A: Decoder](response: Response[F]): F[A] = { diff --git a/core/src/main/scala/jev4s/internal/wire.scala b/core/src/main/scala/jev4s/internal/wire.scala index 49efb46..3637049 100644 --- a/core/src/main/scala/jev4s/internal/wire.scala +++ b/core/src/main/scala/jev4s/internal/wire.scala @@ -67,9 +67,13 @@ private[jev4s] enum QuestionSpec { private[jev4s] object QuestionSpec { - given Encoder[QuestionSpec] = KindlingsEncoder.derived( - using wireConfig - ) + // An absent Noul `criteria` is omitted rather than sent as null: Clef rejects the null. Only the + // top level is filtered, since null is a valid Choice option description. + given Encoder[QuestionSpec] = KindlingsEncoder + .derived[QuestionSpec]( + using wireConfig + ) + .mapJson(_.mapObject(_.filter((_, v) => !v.isNull))) } diff --git a/core/src/main/scala/jev4s/model.scala b/core/src/main/scala/jev4s/model.scala index 1466b16..a78e2c9 100644 --- a/core/src/main/scala/jev4s/model.scala +++ b/core/src/main/scala/jev4s/model.scala @@ -78,6 +78,12 @@ object ModelId { /** The most recent release, official or not. */ val preview: ModelId = "jev-preview" + /** Cloudflare's 27B multimodal Clef model, for [[JevConfig.Provider.WorkersAI]]. */ + val clef: ModelId = "clef" + + /** The faster Clef variant, for [[JevConfig.Provider.WorkersAI]]. */ + val clefFlash: ModelId = "clef-flash" + extension (m: ModelId) { def value: String = m } diff --git a/core/src/test/scala/example/Triage.scala b/core/src/test/scala/example/Triage.scala index 1155fc4..56dd4d0 100644 --- a/core/src/test/scala/example/Triage.scala +++ b/core/src/test/scala/example/Triage.scala @@ -69,30 +69,46 @@ object Triage { object Main extends IOApp.Simple { - val run: IO[Unit] = - EmberClientBuilder.default[IO].withTimeout(10.seconds).build.use { client => - for { - config <- JevConfig.fromEnv[IO] - jev = Jev.instance[IO](config, client) - result <- jev.evaluate( - Ticket("Payouts", "Help! My payouts have been failing for 3 days."), - Triage.question, - ) - triage = result.answers - _ <- IO.println( - s"answered by ${result.model.value}, ${result.usage.inputTokens} input tokens" - ) - _ <- IO.println(s"urgent: ${triage.urgent.yes.value}") - _ <- IO.println( - s"department: ${triage.department.choice} (confidence ${triage.department.confidence.value})" - ) - _ <- IO.println( - s"frustration: ${triage.frustration.score} ~ ${triage.frustration.mostLikely}" - ) - _ <- IO.println( - s"topics: ${triage.topics.filter(_._2.yes.value > 0.5).keys.mkString(", ")}" - ) - } yield () - } + val run: IO[Unit] = JevConfig.fromEnv[IO].flatMap(TriageDemo.run) + +} + +// The same triage, answered by Cloudflare's Clef. Needs CLOUDFLARE_ACCOUNT_ID and CLOUDFLARE_AUTH_TOKEN. +object ClefTriage extends IOApp.Simple { + + val run: IO[Unit] = JevConfig.workersAIFromEnv[IO].flatMap(TriageDemo.run) + +} + +object TriageDemo { + + def run(config: JevConfig): IO[Unit] = + EmberClientBuilder + .default[IO] + .withTimeout(30.seconds) + .build + .use(client => triage(Jev.instance[IO](config, client))) + + def triage(jev: Jev[IO]): IO[Unit] = + for { + result <- jev.evaluate( + Ticket("Payouts", "Help! My payouts have been failing for 3 days."), + Triage.question, + ) + triage = result.answers + _ <- IO.println( + s"answered by ${result.model.value}, ${result.usage.inputTokens} input tokens" + ) + _ <- IO.println(s"urgent: ${triage.urgent.yes.value}") + _ <- IO.println( + s"department: ${triage.department.choice} (confidence ${triage.department.confidence.value})" + ) + _ <- IO.println( + s"frustration: ${triage.frustration.score} ~ ${triage.frustration.mostLikely}" + ) + _ <- IO.println( + s"topics: ${triage.topics.filter(_._2.yes.value > 0.5).keys.mkString(", ")}" + ) + } yield () } diff --git a/core/src/test/scala/jev4s/JevTests.scala b/core/src/test/scala/jev4s/JevTests.scala index 75519e6..17ca93b 100644 --- a/core/src/test/scala/jev4s/JevTests.scala +++ b/core/src/test/scala/jev4s/JevTests.scala @@ -32,7 +32,8 @@ class JevTests extends CatsEffectSuite { private val config = JevConfig(ApiKey("test-key"), retry = JevConfig.RetryConfig.disabled) - private def fake(status: Status, body: Json): IO[(Ref[IO, List[(String, Json)]], Jev[IO])] = + private def fake(status: Status, body: Json, config: JevConfig = config) + : IO[(Ref[IO, List[(String, Json)]], Jev[IO])] = Ref[IO].of(List.empty[(String, Json)]).map { seen => val client = Client.fromHttpApp(HttpApp[IO] { req => req @@ -97,9 +98,7 @@ class JevTests extends CatsEffectSuite { assertEquals( body.hcursor.downField("questions").downField("q3").focus, Some( - json( - """{"type": "noul", "instructions": "Does the ticket ask about refund?", "criteria": null}""" - ) + json("""{"type": "noul", "instructions": "Does the ticket ask about refund?"}""") ), ) } @@ -150,6 +149,46 @@ class JevTests extends CatsEffectSuite { } } + private val workersAI = + JevConfig.workersAI("acc123", ApiKey("cf-token")).copy(retry = JevConfig.RetryConfig.disabled) + + // As returned by the real API, envelope included. + private val noulResponse = json("""{ + "result": { + "model": "clef", + "answers": {"q0": {"type": "noul", "noul": 0.7}}, + "usage": {"input_tokens": 10, "output_tokens": 0} + }, + "success": true, + "errors": [], + "messages": [] + }""") + + test("Workers AI: model goes in the path and the body") { + for { + (seen, jev) <- fake(Status.Ok, noulResponse, workersAI) + result <- jev.withModel(ModelId.clefFlash).evaluate("x", Question.noul("?")) + requests <- seen.get + } yield { + assertEquals(result.answers.yes.value, 0.7) + val List((uri, body)) = requests: @unchecked + assertEquals( + uri, + "https://api.cloudflare.com/client/v4/accounts/acc123/ai/run/@cf/cloudflare/clef-flash", + ) + assertEquals(body.hcursor.downField("model").focus, Some(Json.fromString("clef-flash"))) + } + } + + test("Workers AI: listing models is unsupported") { + fake(Status.Ok, Json.obj(), workersAI) + .flatMap((seen, jev) => jev.models.attempt.product(seen.get)) + .map { (result, requests) => + assertEquals(result, Left(JevError.Unsupported("listing models"))) + assertEquals(requests, Nil) + } + } + test("422 is surfaced with its body") { fake(Status.UnprocessableContent, Json.obj("detail" -> Json.fromString("bad"))) .flatMap((_, jev) => jev.evaluate("x", Question.noul("?")).attempt) From bcaee5fd94d8866e3bd143b5d386cdf711e28b88 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jakub=20Koz=C5=82owski?= Date: Thu, 1 Oct 2026 19:04:17 +0200 Subject: [PATCH 2/3] Make providers pluggable via a Provider[F] trait Everything that differs between System One APIs now lives in a Provider[F]: the default model, endpoints, credentials and unwrapping the response body. The client no longer matches on providers, so a new API is just another implementation. authorize is effectful and runs under the retry middleware, so every attempt can fetch or refresh credentials. The TypeSafe and Workers AI providers read their env vars once, at build time. JevConfig shrinks to a provider plus retry settings. JevConfig.fromEnv and JevConfig.workersAI are replaced by Provider.typeSafeFromEnv and Provider.workersAI(FromEnv), and MissingApiKey by MissingEnv(name). Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 18 ++- core/src/main/scala/jev4s/Jev.scala | 149 ++++++-------------- core/src/main/scala/jev4s/Provider.scala | 130 +++++++++++++++++ core/src/main/scala/jev4s/model.scala | 4 +- core/src/test/scala/example/Triage.scala | 8 +- core/src/test/scala/example/VibeCheck.scala | 4 +- core/src/test/scala/jev4s/JevTests.scala | 53 ++++++- 7 files changed, 239 insertions(+), 127 deletions(-) create mode 100644 core/src/main/scala/jev4s/Provider.scala diff --git a/README.md b/README.md index 760eedf..7023b74 100644 --- a/README.md +++ b/README.md @@ -86,8 +86,8 @@ object Main extends IOApp.Simple { val run: IO[Unit] = EmberClientBuilder.default[IO].build.use { client => for { - config <- JevConfig.fromEnv[IO] - jev = Jev.instance[IO](config, client) + provider <- Provider.typeSafeFromEnv[IO] + jev = Jev.instance[IO](JevConfig(provider), client) result <- jev.evaluate( Ticket("Payouts", "Help! My payouts have been failing for 3 days."), triage, @@ -118,7 +118,7 @@ on topic: 92% ## Configuration -`JevConfig.fromEnv` reads the same environment variables as the official SDKs: +A `Provider` knows where the API lives, how to authenticate with it and which model to use by default. `Provider.typeSafeFromEnv` reads the same environment variables as the official SDKs: | Variable | Required | Default | | --- | --- | --- | @@ -126,7 +126,7 @@ on topic: 92% | `TYPESAFE_BASE_URL` | no | `https://api.typesafe.ai` | | `TYPESAFE_DEFAULT_MODEL` | no | `jev-latest` | -You can also construct a `JevConfig` directly. Requests that fail with 408, 429 or 5xx are retried with exponential backoff (2 retries by default, see `JevConfig.RetryConfig`). Timeouts and connection pooling are up to the `Client` you provide. +You can also construct one directly with `Provider.typeSafe[IO](ApiKey(...))`. Requests that fail with 408, 429 or 5xx are retried with exponential backoff (2 retries by default, see `JevConfig.RetryConfig`). Timeouts and connection pooling are up to the `Client` you provide. To use a different model for some requests, use `jev.withModel(ModelId("jev-1.13.0"))`. `jev.models` lists the available models. @@ -135,11 +135,15 @@ To use a different model for some requests, use `jev.withModel(ModelId("jev-1.13 Cloudflare's [Clef models](https://developers.cloudflare.com/workers-ai/models/clef/) speak the same System One format, so the same questions work against them: ```scala -val config = JevConfig.workersAI(accountId, ApiKey(apiToken)) // or JevConfig.workersAIFromEnv[IO] -val jev = Jev.instance[IO](config, client).withModel(ModelId.clefFlash) +val provider = Provider.workersAI[IO](accountId, ApiKey(apiToken)) // or Provider.workersAIFromEnv[IO] +val jev = Jev.instance[IO](JevConfig(provider), client).withModel(ModelId.clefFlash) ``` -`JevConfig.workersAIFromEnv` reads `CLOUDFLARE_ACCOUNT_ID`, `CLOUDFLARE_AUTH_TOKEN` and optionally `CLOUDFLARE_MODEL` (`clef` by default). The token needs the "Workers AI - Read" and "Workers AI - Edit" permissions. `jev.models` isn't supported on Workers AI. `example.ClefTriage` runs the triage example against Clef. +`Provider.workersAIFromEnv` reads `CLOUDFLARE_ACCOUNT_ID`, `CLOUDFLARE_AUTH_TOKEN` and optionally `CLOUDFLARE_MODEL` (`clef` by default). The token needs the "Workers AI - Read" and "Workers AI - Edit" permissions. `jev.models` isn't supported on Workers AI. `example.ClefTriage` runs the triage example against Clef. + +### Other APIs + +Anything else that serves the System One format can be plugged in by implementing `Provider`: its endpoints, how it adds credentials to a request, and how to find the answers in its response body if they're wrapped. ## License diff --git a/core/src/main/scala/jev4s/Jev.scala b/core/src/main/scala/jev4s/Jev.scala index 55bb793..0c99c0d 100644 --- a/core/src/main/scala/jev4s/Jev.scala +++ b/core/src/main/scala/jev4s/Jev.scala @@ -17,29 +17,24 @@ package jev4s import cats.effect.Concurrent +import cats.effect.MonadCancelThrow +import cats.effect.Resource import cats.effect.Temporal -import cats.effect.std.Env import cats.syntax.all.* import io.circe.Decoder import io.circe.Encoder import io.circe.Json import jev4s.internal.ModelList import jev4s.internal.ResponseBody -import org.http4s.EntityDecoder import org.http4s.Headers import org.http4s.Method import org.http4s.Request import org.http4s.Response import org.http4s.Status -import org.http4s.Uri import org.http4s.circe.* import org.http4s.client.Client import org.http4s.client.middleware.Retry import org.http4s.client.middleware.RetryPolicy -import org.http4s.headers.Authorization -import org.http4s.implicits.* -import org.http4s.AuthScheme -import org.http4s.Credentials import scala.concurrent.duration.* @@ -54,45 +49,17 @@ trait Jev[F[_]] { def models: F[List[ModelCard]] - /** The same client, sending `model` instead of the configured one. */ + /** The same client, sending `model` instead of the provider's default. */ def withModel(model: ModelId): Jev[F] } -final case class JevConfig( - apiKey: ApiKey, - baseUri: Uri = uri"https://api.typesafe.ai", - model: ModelId = ModelId.latest, +final case class JevConfig[F[_]]( + provider: Provider[F], retry: JevConfig.RetryConfig = JevConfig.RetryConfig.default, - provider: JevConfig.Provider = JevConfig.Provider.TypeSafe, ) object JevConfig { - /** Who serves the System One API. The request and answer formats are the same, the endpoints - * differ. - */ - enum Provider { - - /** TypeSafe's own API: `POST v1/systemone`. */ - case TypeSafe - - /** Cloudflare Workers AI (the Clef models): - * `POST accounts/{accountId}/ai/run/@cf/cloudflare/{model}`. Model listing isn't supported. - */ - case WorkersAI(accountId: String) - } - - /** Clef on Cloudflare Workers AI. `apiToken` needs the "Workers AI - Read" and "Workers AI - - * Edit" permissions. - */ - def workersAI(accountId: String, apiToken: ApiKey, model: ModelId = ModelId.clef): JevConfig = - JevConfig( - apiKey = apiToken, - baseUri = uri"https://api.cloudflare.com/client/v4", - model = model, - provider = Provider.WorkersAI(accountId), - ) - final case class RetryConfig(maxRetries: Int, maxBackoff: FiniteDuration) object RetryConfig { @@ -100,42 +67,9 @@ object JevConfig { val disabled: RetryConfig = RetryConfig(maxRetries = 0, maxBackoff = Duration.Zero) } - /** Reads `TYPESAFE_API_KEY` (required), `TYPESAFE_BASE_URL` and `TYPESAFE_DEFAULT_MODEL`, like - * the official SDKs. - */ - def fromEnv[F[_]: Env: Concurrent]: F[JevConfig] = - ( - Env[F].get("TYPESAFE_API_KEY").flatMap(_.liftTo[F](JevError.MissingApiKey)), - Env[F].get("TYPESAFE_BASE_URL").flatMap(_.traverse(Uri.fromString(_).liftTo[F])), - Env[F].get("TYPESAFE_DEFAULT_MODEL"), - ).mapN { (key, baseUri, model) => - val base = JevConfig(ApiKey(key)) - base.copy( - baseUri = baseUri.getOrElse(base.baseUri), - model = model.fold(base.model)(ModelId(_)), - ) - } - - /** Reads `CLOUDFLARE_ACCOUNT_ID` and `CLOUDFLARE_AUTH_TOKEN` (both required) and - * `CLOUDFLARE_MODEL` (`clef` by default). - */ - def workersAIFromEnv[F[_]: Env: Concurrent]: F[JevConfig] = - ( - Env[F] - .get("CLOUDFLARE_ACCOUNT_ID") - .flatMap(_.liftTo[F](JevError.MissingEnv("CLOUDFLARE_ACCOUNT_ID"))), - Env[F] - .get("CLOUDFLARE_AUTH_TOKEN") - .flatMap(_.liftTo[F](JevError.MissingEnv("CLOUDFLARE_AUTH_TOKEN"))), - Env[F].get("CLOUDFLARE_MODEL"), - ).mapN { (accountId, token, model) => - workersAI(accountId, ApiKey(token), model.fold(ModelId.clef)(ModelId(_))) - } - } enum JevError(message: String) extends Exception(message) { - case MissingApiKey extends JevError("TYPESAFE_API_KEY is not set") case MissingEnv(name: String) extends JevError(s"$name is not set") case Unsupported(operation: String) extends JevError(s"Not supported by this provider: $operation") @@ -154,8 +88,16 @@ object Jev { /** Timeouts, connection pooling etc. are up to the `Client` you provide; retries are added on top * of it. */ - def instance[F[_]: Temporal](config: JevConfig, client: Client[F]): Jev[F] = - JevImpl(config, withRetries(config.retry, client)) + def instance[F[_]: Temporal](config: JevConfig[F], client: Client[F]): Jev[F] = + JevImpl( + config.provider, + config.provider.defaultModel, + withRetries(config.retry, authorized(config.provider, client)), + ) + + // Under the retries, so that every attempt asks the provider for credentials. + private def authorized[F[_]: MonadCancelThrow](provider: Provider[F], client: Client[F]) + : Client[F] = Client(request => Resource.eval(provider.authorize(request)).flatMap(client.run)) // Same statuses as the official SDKs: 408, 429, 5xx (529 included). Retry-After is honored by the middleware. private def withRetries[F[_]: Temporal](config: JevConfig.RetryConfig, client: Client[F]) @@ -175,22 +117,23 @@ object Jev { Headers.SensitiveHeaders.contains, )(client) - private final class JevImpl[F[_]: Concurrent] private[Jev] (config: JevConfig, client: Client[F]) - extends Jev[F] { + private final class JevImpl[F[_]: Concurrent] private[Jev] ( + provider: Provider[F], + model: ModelId, + client: Client[F], + ) extends Jev[F] { def evaluate[S: Encoder, A](state: S, question: Question[A]): F[Evaluation[A]] = for { body <- Question - .requestBody(state, question, config.model) + .requestBody(state, question, model) .leftMap(JevError.InvalidQuestion(_)) .liftTo[F] response <- client - .run(authorized(Method.POST, evaluateUri).withEntity(body)) - .use(r => - decodeOrFail(r)( - using responseDecoder - ) + .run( + Request[F](Method.POST, provider.evaluateUri(model)).withEntity(body) ) + .use(decodeOrFail[ResponseBody]) evaluation <- Question .decodeResponse(question, response) .leftMap(JevError.UnexpectedAnswer(_)) @@ -198,40 +141,29 @@ object Jev { } yield evaluation def models: F[List[ModelCard]] = - config.provider match { - case JevConfig.Provider.TypeSafe => + provider.modelsUri match { + case Some(uri) => client - .run(authorized(Method.GET, config.baseUri / "v1" / "models")) + .run(Request[F](Method.GET, uri)) .use(decodeOrFail[ModelList]) .map(_.models) - case JevConfig.Provider.WorkersAI(_) => JevError.Unsupported("listing models").raiseError - } - - def withModel(model: ModelId): Jev[F] = JevImpl(config.copy(model = model), client) - - private def evaluateUri: Uri = - config.provider match { - case JevConfig.Provider.TypeSafe => config.baseUri / "v1" / "systemone" - case JevConfig.Provider.WorkersAI(accountId) => - config.baseUri / "accounts" / accountId / "ai" / "run" / "@cf" / "cloudflare" / - config.model.value + case None => JevError.Unsupported("listing models").raiseError } - // Workers AI wraps the answers in Cloudflare's {"result": ..., "success": ...} envelope. - private def responseDecoder: Decoder[ResponseBody] = - config.provider match { - case JevConfig.Provider.TypeSafe => Decoder[ResponseBody] - case JevConfig.Provider.WorkersAI(_) => Decoder[ResponseBody].at("result") - } + def withModel(model: ModelId): Jev[F] = JevImpl(provider, model, client) - private def authorized(method: Method, uri: Uri): Request[F] = - Request[F](method, uri) - .putHeaders(Authorization(Credentials.Token(AuthScheme.Bearer, config.apiKey.value))) - - private def decodeOrFail[A: Decoder](response: Response[F]): F[A] = { - given EntityDecoder[F, A] = jsonOf[F, A] + private def decodeOrFail[A: Decoder](response: Response[F]): F[A] = response.status match { - case s if s.isSuccess => response.as[A] + case s if s.isSuccess => + response + .as[Json] + .flatMap( + provider + .unwrap(_) + .flatMap(_.as[A].leftMap(_.getMessage)) + .leftMap(JevError.UnexpectedAnswer(_)) + .liftTo[F] + ) case Status.Unauthorized => response.as[String].flatMap(JevError.Unauthorized(_).raiseError) case Status.UnprocessableContent => response.as[Json].flatMap(JevError.Unprocessable(_).raiseError) @@ -240,7 +172,6 @@ object Jev { case s if s.code == 529 => response.as[String].flatMap(JevError.Overloaded(_).raiseError) case s => response.as[String].flatMap(JevError.UnexpectedStatus(s, _).raiseError) } - } } diff --git a/core/src/main/scala/jev4s/Provider.scala b/core/src/main/scala/jev4s/Provider.scala new file mode 100644 index 0000000..e41fb7e --- /dev/null +++ b/core/src/main/scala/jev4s/Provider.scala @@ -0,0 +1,130 @@ +/* + * Copyright 2026 Polyvariant + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package jev4s + +import cats.Applicative +import cats.MonadThrow +import cats.effect.std.Env +import cats.syntax.all.* +import io.circe.Json +import org.http4s.AuthScheme +import org.http4s.Credentials +import org.http4s.Request +import org.http4s.Uri +import org.http4s.headers.Authorization +import org.http4s.implicits.* + +/** An API that serves the System One format. Requests and answers look the same everywhere, but + * endpoints and credentials differ: implement this to point jev4s at a new one. + */ +trait Provider[F[_]] { + + /** The model to use unless `Jev.withModel` says otherwise. */ + def defaultModel: ModelId + + /** Where evaluations with `model` are POSTed. */ + def evaluateUri(model: ModelId): Uri + + /** Where the model list is fetched from, if the API has one. */ + def modelsUri: Option[Uri] + + /** Adds credentials to a request. Runs before every request, so it can fetch or refresh them. */ + def authorize(request: Request[F]): F[Request[F]] + + /** Finds the System One response in a successful response body, e.g. inside an envelope. */ + def unwrap(body: Json): Either[String, Json] = Right(body) +} + +object Provider { + + private val typeSafeUri = uri"https://api.typesafe.ai" + + /** TypeSafe's own API. */ + def typeSafe[F[_]: Applicative]( + apiKey: ApiKey, + baseUri: Uri = typeSafeUri, + defaultModel: ModelId = ModelId.latest, + ): Provider[F] = { + val model = defaultModel + new Provider[F] { + def defaultModel: ModelId = model + def evaluateUri(model: ModelId): Uri = baseUri / "v1" / "systemone" + def modelsUri: Option[Uri] = Some(baseUri / "v1" / "models") + def authorize(request: Request[F]): F[Request[F]] = bearer(request, apiKey).pure[F] + } + } + + /** Reads `TYPESAFE_API_KEY` (required), `TYPESAFE_BASE_URL` and `TYPESAFE_DEFAULT_MODEL`, like + * the official SDKs. They're read once, not on every request. + */ + def typeSafeFromEnv[F[_]: Env: MonadThrow]: F[Provider[F]] = + ( + required[F]("TYPESAFE_API_KEY"), + Env[F].get("TYPESAFE_BASE_URL").flatMap(_.traverse(Uri.fromString(_).liftTo[F])), + Env[F].get("TYPESAFE_DEFAULT_MODEL"), + ).mapN { (key, baseUri, model) => + typeSafe[F]( + ApiKey(key), + baseUri = baseUri.getOrElse(typeSafeUri), + defaultModel = model.fold(ModelId.latest)(ModelId(_)), + ) + } + + /** Cloudflare Workers AI, serving the Clef models. `apiToken` needs the "Workers AI - Read" and + * "Workers AI - Edit" permissions. There's no model list, and answers come wrapped in + * Cloudflare's `{"result": ..., "success": ...}` envelope. + */ + def workersAI[F[_]: Applicative]( + accountId: String, + apiToken: ApiKey, + baseUri: Uri = uri"https://api.cloudflare.com/client/v4", + defaultModel: ModelId = ModelId.clef, + ): Provider[F] = { + val model = defaultModel + new Provider[F] { + def defaultModel: ModelId = model + + def evaluateUri(model: ModelId): Uri = + baseUri / "accounts" / accountId / "ai" / "run" / "@cf" / "cloudflare" / model.value + + def modelsUri: Option[Uri] = None + def authorize(request: Request[F]): F[Request[F]] = bearer(request, apiToken).pure[F] + + override def unwrap(body: Json): Either[String, Json] = + body.hcursor.downField("result").focus.toRight("Missing `result` in the response") + } + } + + /** Reads `CLOUDFLARE_ACCOUNT_ID` and `CLOUDFLARE_AUTH_TOKEN` (both required) and + * `CLOUDFLARE_MODEL` (`clef` by default). They're read once, not on every request. + */ + def workersAIFromEnv[F[_]: Env: MonadThrow]: F[Provider[F]] = + ( + required[F]("CLOUDFLARE_ACCOUNT_ID"), + required[F]("CLOUDFLARE_AUTH_TOKEN"), + Env[F].get("CLOUDFLARE_MODEL"), + ).mapN { (accountId, token, model) => + workersAI[F](accountId, ApiKey(token), defaultModel = model.fold(ModelId.clef)(ModelId(_))) + } + + private def required[F[_]: Env: MonadThrow](name: String): F[String] = + Env[F].get(name).flatMap(_.liftTo[F](JevError.MissingEnv(name))) + + private def bearer[F[_]](request: Request[F], key: ApiKey): Request[F] = + request.putHeaders(Authorization(Credentials.Token(AuthScheme.Bearer, key.value))) + +} diff --git a/core/src/main/scala/jev4s/model.scala b/core/src/main/scala/jev4s/model.scala index a78e2c9..a681046 100644 --- a/core/src/main/scala/jev4s/model.scala +++ b/core/src/main/scala/jev4s/model.scala @@ -78,10 +78,10 @@ object ModelId { /** The most recent release, official or not. */ val preview: ModelId = "jev-preview" - /** Cloudflare's 27B multimodal Clef model, for [[JevConfig.Provider.WorkersAI]]. */ + /** Cloudflare's 27B multimodal Clef model, for [[Provider.workersAI]]. */ val clef: ModelId = "clef" - /** The faster Clef variant, for [[JevConfig.Provider.WorkersAI]]. */ + /** The faster Clef variant, for [[Provider.workersAI]]. */ val clefFlash: ModelId = "clef-flash" extension (m: ModelId) { diff --git a/core/src/test/scala/example/Triage.scala b/core/src/test/scala/example/Triage.scala index 56dd4d0..9079f87 100644 --- a/core/src/test/scala/example/Triage.scala +++ b/core/src/test/scala/example/Triage.scala @@ -69,25 +69,25 @@ object Triage { object Main extends IOApp.Simple { - val run: IO[Unit] = JevConfig.fromEnv[IO].flatMap(TriageDemo.run) + val run: IO[Unit] = Provider.typeSafeFromEnv[IO].flatMap(TriageDemo.run) } // The same triage, answered by Cloudflare's Clef. Needs CLOUDFLARE_ACCOUNT_ID and CLOUDFLARE_AUTH_TOKEN. object ClefTriage extends IOApp.Simple { - val run: IO[Unit] = JevConfig.workersAIFromEnv[IO].flatMap(TriageDemo.run) + val run: IO[Unit] = Provider.workersAIFromEnv[IO].flatMap(TriageDemo.run) } object TriageDemo { - def run(config: JevConfig): IO[Unit] = + def run(provider: Provider[IO]): IO[Unit] = EmberClientBuilder .default[IO] .withTimeout(30.seconds) .build - .use(client => triage(Jev.instance[IO](config, client))) + .use(client => triage(Jev.instance[IO](JevConfig(provider), client))) def triage(jev: Jev[IO]): IO[Unit] = for { diff --git a/core/src/test/scala/example/VibeCheck.scala b/core/src/test/scala/example/VibeCheck.scala index f271451..38f83fd 100644 --- a/core/src/test/scala/example/VibeCheck.scala +++ b/core/src/test/scala/example/VibeCheck.scala @@ -110,9 +110,9 @@ object VibeCheck extends IOApp.Simple { val run: IO[Unit] = EmberClientBuilder.default[IO].withTimeout(30.seconds).build.use { client => for { - config <- JevConfig.fromEnv[IO] + provider <- Provider.typeSafeFromEnv[IO] _ <- IO.println("Say something (empty line to quit).") - _ <- loop(Jev.instance[IO](config, client), Nil) + _ <- loop(Jev.instance[IO](JevConfig(provider), client), Nil) } yield () } diff --git a/core/src/test/scala/jev4s/JevTests.scala b/core/src/test/scala/jev4s/JevTests.scala index 17ca93b..cc0dd03 100644 --- a/core/src/test/scala/jev4s/JevTests.scala +++ b/core/src/test/scala/jev4s/JevTests.scala @@ -22,17 +22,27 @@ import example.* import io.circe.Json import io.circe.parser.parse import munit.CatsEffectSuite +import org.http4s.Header import org.http4s.HttpApp +import org.http4s.Request +import org.http4s.Uri import org.http4s.Response import org.http4s.Status import org.http4s.circe.* import org.http4s.client.Client +import org.http4s.implicits.* +import org.typelevel.ci.* + +import scala.concurrent.duration.* class JevTests extends CatsEffectSuite { - private val config = JevConfig(ApiKey("test-key"), retry = JevConfig.RetryConfig.disabled) + private val config = JevConfig( + Provider.typeSafe[IO](ApiKey("test-key")), + retry = JevConfig.RetryConfig.disabled, + ) - private def fake(status: Status, body: Json, config: JevConfig = config) + private def fake(status: Status, body: Json, config: JevConfig[IO] = config) : IO[(Ref[IO, List[(String, Json)]], Jev[IO])] = Ref[IO].of(List.empty[(String, Json)]).map { seen => val client = Client.fromHttpApp(HttpApp[IO] { req => @@ -150,7 +160,10 @@ class JevTests extends CatsEffectSuite { } private val workersAI = - JevConfig.workersAI("acc123", ApiKey("cf-token")).copy(retry = JevConfig.RetryConfig.disabled) + JevConfig( + Provider.workersAI[IO]("acc123", ApiKey("cf-token")), + retry = JevConfig.RetryConfig.disabled, + ) // As returned by the real API, envelope included. private val noulResponse = json("""{ @@ -189,6 +202,40 @@ class JevTests extends CatsEffectSuite { } } + test("credentials are fetched for every attempt, retries included") { + for { + tokens <- Ref[IO].of(0) + seen <- Ref[IO].of(List.empty[String]) + provider = + new Provider[IO] { + def defaultModel: ModelId = ModelId.latest + def evaluateUri(model: ModelId): Uri = uri"https://example.com/eval" + def modelsUri: Option[Uri] = None + def authorize(request: Request[IO]): IO[Request[IO]] = + tokens.updateAndGet(_ + 1).map(n => request.putHeaders(Header.Raw(ci"X-Token", s"t$n"))) + } + client = Client.fromHttpApp(HttpApp[IO] { req => + val token = req.headers.get(ci"X-Token").fold("none")(_.head.value) + seen.updateAndGet(_ :+ token).map { all => + if (all.sizeIs == 1) + Response[IO](Status.ServiceUnavailable) + else + Response[IO](Status.Ok).withEntity( + json( + """{"model": "m", "answers": {"q0": {"type": "noul", "noul": 0.5}}, "usage": {"input_tokens": 1, "output_tokens": 0}}""" + ) + ) + } + }) + jev = Jev.instance[IO]( + JevConfig(provider, JevConfig.RetryConfig(maxRetries = 1, maxBackoff = 1.milli)), + client, + ) + _ <- jev.evaluate("x", Question.noul("?")) + tokensSent <- seen.get + } yield assertEquals(tokensSent, List("t1", "t2")) + } + test("422 is surfaced with its body") { fake(Status.UnprocessableContent, Json.obj("detail" -> Json.fromString("bad"))) .flatMap((_, jev) => jev.evaluate("x", Question.noul("?")).attempt) From 13e4e43948a69647d1b98ae20b45335288b1a17d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jakub=20Koz=C5=82owski?= Date: Thu, 1 Oct 2026 19:06:15 +0200 Subject: [PATCH 3/3] Bump versions: sbt, http4s --- build.sbt | 2 +- project/build.properties | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/build.sbt b/build.sbt index 684ce85..add5409 100644 --- a/build.sbt +++ b/build.sbt @@ -19,7 +19,7 @@ ThisBuild / mergifyStewardConfig ~= (_.map(_.withMergeMinors(true))) val hearthVersion = "0.4.2" val kindlingsVersion = "0.3.2" -val http4sVersion = "0.23.37" +val http4sVersion = "0.23.38" val circeVersion = "0.14.16" lazy val core = crossProject(JVMPlatform, JSPlatform, NativePlatform) diff --git a/project/build.properties b/project/build.properties index 7c95fc1..e0a1aa0 100644 --- a/project/build.properties +++ b/project/build.properties @@ -1 +1 @@ -sbt.version=1.12.13 +sbt.version=1.13.0