runtime.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326
  1. import { Effect, Option, Schema, Stream } from "effect"
  2. import { HttpClient, HttpClientRequest, HttpClientResponse, type HttpMethod } from "effect/unstable/http"
  3. import { ToolError, toolError } from "../tool-error.js"
  4. import { isRecord, own } from "./spec.js"
  5. import type { AppliedAuth, Credential, Plan, SecurityScheme } from "./types.js"
  6. const decodeJson = Schema.decodeUnknownOption(Schema.UnknownFromJsonString)
  7. const maxErrorBodyChars = 1_024
  8. const maxResponseBodyBytes = 50 * 1024 * 1024
  9. export const invoke = (plan: Plan, input: unknown): Effect.Effect<unknown, unknown, HttpClient.HttpClient> =>
  10. Effect.gen(function* () {
  11. const value = isRecord(input) ? input : {}
  12. let request = yield* buildRequest(plan, value)
  13. const auth = yield* resolveAuth(plan)
  14. for (const [name, item] of Object.entries(auth.query)) {
  15. request = HttpClientRequest.setUrlParam(request, name, item)
  16. }
  17. request = HttpClientRequest.setHeaders(request, auth.headers)
  18. const client = yield* HttpClient.HttpClient
  19. const response = yield* client
  20. .execute(request)
  21. .pipe(
  22. Effect.catch((cause) =>
  23. Effect.fail(toolError(`${plan.operation.method} ${plan.operation.path} failed: transport error`, cause)),
  24. ),
  25. )
  26. const text = yield* readResponseBody(response, plan)
  27. const mediaType = response.headers["content-type"]?.split(";")[0]?.trim().toLowerCase()
  28. const json = mediaType === "application/json" || mediaType?.endsWith("+json") === true
  29. const decoded = text === "" ? Option.some(null) : json ? decodeJson(text) : Option.none()
  30. const parsed = json ? Option.getOrElse(decoded, () => text) : text === "" ? null : text
  31. if (response.status < 200 || response.status >= 300) {
  32. const rendered = typeof parsed === "string" ? parsed : (JSON.stringify(parsed) ?? "")
  33. const summary =
  34. rendered === "" || rendered === "null"
  35. ? "no response body"
  36. : rendered.length > maxErrorBodyChars
  37. ? `${rendered.slice(0, maxErrorBodyChars)}...`
  38. : rendered
  39. return yield* Effect.fail(
  40. toolError(`${plan.operation.method} ${plan.operation.path} failed with HTTP ${response.status}: ${summary}`),
  41. )
  42. }
  43. if (json && Option.isNone(decoded)) {
  44. return yield* Effect.fail(toolError(`${plan.operation.method} ${plan.operation.path} returned malformed JSON.`))
  45. }
  46. return parsed
  47. })
  48. const buildRequest = (
  49. plan: Plan,
  50. input: Readonly<Record<string, unknown>>,
  51. ): Effect.Effect<HttpClientRequest.HttpClientRequest, ToolError> =>
  52. Effect.gen(function* () {
  53. // Validate every model-controlled value before auth resolution, which may refresh tokens.
  54. const url = buildUrl(plan, input)
  55. if (url instanceof ToolError) return yield* Effect.fail(url)
  56. const missing = plan.fields.find(
  57. (field) => field.required && field.location !== "path" && own(input, field.inputName) === undefined,
  58. )
  59. if (missing !== undefined) {
  60. const label = missing.location === "body" ? "body field" : `${missing.location} parameter`
  61. return yield* Effect.fail(toolError(`Missing required ${label} '${missing.inputName}'.`))
  62. }
  63. let request = HttpClientRequest.make(plan.operation.method as HttpMethod.HttpMethod)(url)
  64. for (const field of plan.fields) {
  65. if (field.location !== "query") continue
  66. const item = own(input, field.inputName)
  67. if (item === undefined) continue
  68. const serialized = serializeQuery(request, field, item)
  69. if (serialized instanceof ToolError) return yield* Effect.fail(serialized)
  70. request = serialized
  71. }
  72. // Host headers first, then declared header parameters.
  73. request = HttpClientRequest.setHeaders(request, plan.headers)
  74. for (const field of plan.fields) {
  75. if (field.location !== "header") continue
  76. const item = own(input, field.inputName)
  77. if (item === undefined) continue
  78. const serialized = serializeSimple(field, item, String)
  79. if (serialized instanceof ToolError) return yield* Effect.fail(serialized)
  80. request = HttpClientRequest.setHeader(request, field.name, serialized)
  81. }
  82. const setBody = (value: unknown, mediaType: string) =>
  83. HttpClientRequest.bodyJson(request, value).pipe(
  84. Effect.map((next) => HttpClientRequest.setHeader(next, "content-type", mediaType)),
  85. Effect.mapError((cause) =>
  86. toolError(`Invalid JSON body for ${plan.operation.method} ${plan.operation.path}.`, cause),
  87. ),
  88. )
  89. if (plan.body?.mode === "value") {
  90. const field = plan.fields.find((field) => field.location === "body")
  91. const body = field === undefined ? undefined : own(input, field.inputName)
  92. if (body !== undefined) request = yield* setBody(body, plan.body.mediaType)
  93. }
  94. if (plan.body?.mode === "object") {
  95. const entries = plan.fields.flatMap((field) => {
  96. if (field.location !== "body") return []
  97. const item = own(input, field.inputName)
  98. return item === undefined ? [] : [[field.name, item] as const]
  99. })
  100. if (plan.body.required || entries.length > 0) {
  101. request = yield* setBody(Object.fromEntries(entries), plan.body.mediaType)
  102. }
  103. }
  104. return request
  105. })
  106. const resolveAuth = (plan: Plan): Effect.Effect<AppliedAuth, unknown> =>
  107. Effect.gen(function* () {
  108. const none: AppliedAuth = { headers: {}, query: {} }
  109. if (plan.security.length === 0) return none
  110. const unavailable: Array<string> = []
  111. alternatives: for (const requirement of plan.security) {
  112. const names = Object.keys(requirement)
  113. if (names.length === 0) return none
  114. const credentials: Array<readonly [string, SecurityScheme, Credential]> = []
  115. for (const name of names) {
  116. const scheme = own(plan.schemes, name)
  117. if (scheme === undefined || plan.auth === undefined) {
  118. unavailable.push(name)
  119. continue alternatives
  120. }
  121. const credential = yield* plan.auth.resolve({
  122. name,
  123. definition: scheme,
  124. scopes: requirement[name] ?? [],
  125. operation: plan.operation,
  126. })
  127. if (credential === undefined) {
  128. unavailable.push(name)
  129. continue alternatives
  130. }
  131. credentials.push([name, scheme, credential])
  132. }
  133. const applied = applyCredentials(credentials)
  134. return applied instanceof ToolError ? yield* Effect.fail(applied) : applied
  135. }
  136. return yield* Effect.fail(
  137. toolError(
  138. `${plan.operation.method} ${plan.operation.path} requires authentication; no credential available for: ${[...new Set(unavailable)].join(", ")}.`,
  139. ),
  140. )
  141. })
  142. const applyCredentials = (
  143. credentials: ReadonlyArray<readonly [string, SecurityScheme, Credential]>,
  144. ): AppliedAuth | ToolError => {
  145. const headers = new Map<string, string>()
  146. const query = new Map<string, string>()
  147. const add = (carrier: "header" | "query", name: string, value: string): ToolError | undefined => {
  148. const target = carrier === "header" ? headers : query
  149. if (target.has(name)) return toolError(`Authentication resolves multiple credentials for ${carrier} '${name}'.`)
  150. target.set(name, value)
  151. }
  152. for (const [name, definition, credential] of credentials) {
  153. if (credential.type === "bearer") {
  154. const duplicate = add("header", "authorization", `Bearer ${credential.token}`)
  155. if (duplicate !== undefined) return duplicate
  156. continue
  157. }
  158. if (credential.type === "basic") {
  159. // Buffer instead of btoa: btoa throws on non-Latin-1 credentials.
  160. const duplicate = add(
  161. "header",
  162. "authorization",
  163. `Basic ${Buffer.from(`${credential.username}:${credential.password}`, "utf8").toString("base64")}`,
  164. )
  165. if (duplicate !== undefined) return duplicate
  166. continue
  167. }
  168. if (credential.type === "header") {
  169. const duplicate = add("header", credential.name.toLowerCase(), credential.value)
  170. if (duplicate !== undefined) return duplicate
  171. continue
  172. }
  173. // apiKey: the carrier comes from the scheme declaration.
  174. if (definition.type !== "apiKey") {
  175. return toolError(
  176. `Security scheme '${name}' is not an apiKey scheme; resolve a bearer, basic, or header credential for it.`,
  177. )
  178. }
  179. if (definition.in === "cookie") return toolError(`Cookie authentication '${name}' is not supported.`)
  180. const parameter = definition.in === "header" ? definition.name.toLowerCase() : definition.name
  181. const duplicate = add(definition.in, parameter, credential.value)
  182. if (duplicate !== undefined) return duplicate
  183. }
  184. return { headers: Object.fromEntries(headers), query: Object.fromEntries(query) }
  185. }
  186. const buildUrl = (plan: Plan, input: Readonly<Record<string, unknown>>): string | ToolError => {
  187. let url = plan.url
  188. for (const field of plan.fields) {
  189. if (field.location !== "path") continue
  190. const item = own(input, field.inputName)
  191. if (item === undefined) {
  192. return toolError(`Missing required path parameter '${field.inputName}'.`)
  193. }
  194. const fieldValue = serializeSimple(field, item, (value) =>
  195. encodeURIComponent(value).replace(
  196. /[!'()*]/g,
  197. (character) => `%${character.charCodeAt(0).toString(16).toUpperCase()}`,
  198. ),
  199. )
  200. if (fieldValue instanceof ToolError) return fieldValue
  201. // '.'/'..' survive encoding and URL normalization collapses them, letting a
  202. // model-supplied value retarget the request to a different endpoint.
  203. if (fieldValue === "" || fieldValue === "." || fieldValue === "..") {
  204. return toolError(`Invalid path parameter '${field.inputName}'.`)
  205. }
  206. url = url.replaceAll(`{${field.name}}`, fieldValue)
  207. }
  208. const unresolved = url.match(/\{[^{}]+\}/)
  209. if (unresolved !== null) return toolError(`Unresolved path parameter ${unresolved[0]}.`)
  210. return url
  211. }
  212. const serializeSimple = (
  213. field: Plan["fields"][number],
  214. value: unknown,
  215. encode: (value: string) => string,
  216. ): string | ToolError => {
  217. const scalar = (item: unknown): string | ToolError =>
  218. item !== null && typeof item !== "string" && typeof item !== "number" && typeof item !== "boolean"
  219. ? toolError(`Parameter '${field.inputName}' contains an unsupported nested value.`)
  220. : encode(String(item))
  221. if (Array.isArray(value)) {
  222. const items = value.map(scalar)
  223. const invalid = items.find((item): item is ToolError => item instanceof ToolError)
  224. return invalid ?? items.join(",")
  225. }
  226. if (!isRecord(value)) return scalar(value)
  227. const entries = Object.entries(value).flatMap<string | ToolError>(([name, item]) => {
  228. const rendered = scalar(item)
  229. if (rendered instanceof ToolError) return [rendered]
  230. return field.explode ? [`${encode(name)}=${rendered}`] : [encode(name), rendered]
  231. })
  232. const invalid = entries.find((item): item is ToolError => item instanceof ToolError)
  233. return invalid ?? entries.join(",")
  234. }
  235. const serializeQuery = (
  236. request: HttpClientRequest.HttpClientRequest,
  237. field: Plan["fields"][number],
  238. value: unknown,
  239. ): HttpClientRequest.HttpClientRequest | ToolError => {
  240. if (field.style === "deepObject") {
  241. if (!isRecord(value)) return toolError(`Deep-object parameter '${field.inputName}' must be an object.`)
  242. return Object.entries(value).reduce<HttpClientRequest.HttpClientRequest | ToolError>((current, [name, item]) => {
  243. if (current instanceof ToolError) return current
  244. if (item === undefined || (item !== null && typeof item === "object")) {
  245. return toolError(`Deep-object parameter '${field.inputName}' contains an unsupported nested value.`)
  246. }
  247. return HttpClientRequest.appendUrlParam(current, `${field.name}[${name}]`, String(item))
  248. }, request)
  249. }
  250. if (Array.isArray(value)) {
  251. const rendered = serializeSimple(field, value, String)
  252. if (rendered instanceof ToolError) return rendered
  253. if (!field.explode) return HttpClientRequest.appendUrlParam(request, field.name, rendered)
  254. if (value.some((item) => item === undefined || (item !== null && typeof item === "object"))) {
  255. return toolError(`Query parameter '${field.inputName}' contains an unsupported nested value.`)
  256. }
  257. return value.reduce((current, item) => HttpClientRequest.appendUrlParam(current, field.name, String(item)), request)
  258. }
  259. if (isRecord(value) && field.explode) {
  260. return Object.entries(value).reduce<HttpClientRequest.HttpClientRequest | ToolError>((current, [name, item]) => {
  261. if (current instanceof ToolError) return current
  262. if (item === undefined || (item !== null && typeof item === "object")) {
  263. return toolError(`Query parameter '${field.inputName}' contains an unsupported nested value.`)
  264. }
  265. return HttpClientRequest.appendUrlParam(current, name, String(item))
  266. }, request)
  267. }
  268. const rendered = serializeSimple(field, value, String)
  269. return rendered instanceof ToolError ? rendered : HttpClientRequest.appendUrlParam(request, field.name, rendered)
  270. }
  271. const readResponseBody = (
  272. response: HttpClientResponse.HttpClientResponse,
  273. plan: Plan,
  274. ): Effect.Effect<string, ToolError> =>
  275. Effect.gen(function* () {
  276. const contentLength = response.headers["content-length"]
  277. const parsedSize = contentLength === undefined ? undefined : Number.parseInt(contentLength, 10)
  278. const declaredSize =
  279. parsedSize !== undefined && Number.isSafeInteger(parsedSize) && parsedSize >= 0 ? parsedSize : undefined
  280. if (declaredSize !== undefined && declaredSize > maxResponseBodyBytes) {
  281. return yield* Effect.fail(toolError(`${plan.operation.method} ${plan.operation.path} response exceeds 50 MiB.`))
  282. }
  283. let body = Buffer.allocUnsafe(Math.min(maxResponseBodyBytes, declaredSize ?? 64 * 1024))
  284. let size = 0
  285. yield* Stream.runForEach(response.stream, (chunk) => {
  286. if (size + chunk.byteLength > maxResponseBodyBytes) {
  287. return Effect.fail(toolError(`${plan.operation.method} ${plan.operation.path} response exceeds 50 MiB.`))
  288. }
  289. if (size + chunk.byteLength > body.byteLength) {
  290. const grown = Buffer.allocUnsafe(
  291. Math.min(maxResponseBodyBytes, Math.max(size + chunk.byteLength, body.byteLength * 2)),
  292. )
  293. body.copy(grown, 0, 0, size)
  294. body = grown
  295. }
  296. body.set(chunk, size)
  297. size += chunk.byteLength
  298. return Effect.void
  299. }).pipe(
  300. Effect.catch((cause) => {
  301. if (cause instanceof ToolError) return Effect.fail(cause)
  302. if (cause.reason._tag === "EmptyBodyError") return Effect.void
  303. return Effect.fail(
  304. toolError(`${plan.operation.method} ${plan.operation.path} failed while reading the response body.`, cause),
  305. )
  306. }),
  307. )
  308. return new TextDecoder().decode(body.subarray(0, size))
  309. })