From 96c8a95094c91317f191805b6f069f741b413a27 Mon Sep 17 00:00:00 2001 From: Gregory Schier Date: Sun, 16 Aug 2026 09:20:29 -0700 Subject: [PATCH] Add a plugin API for reading HTTP response bodies Plugins read bodies by response id through ctx.httpResponse.body() instead of opening HttpResponse.bodyPath themselves. The accessors are named after fetch's, minus the single-use semantics, since the bytes are durable and re-reading should work. Underneath is a chunked pull over the existing plugin protocol, so the host can move bodies off the filesystem without plugins noticing. text() now decodes with the response's charset rather than assuming UTF-8, and the buffering accessors refuse past 32 MiB and point at chunks(). --- Cargo.lock | 1 + crates-cli/yaak-cli/src/plugin_events.rs | 2 + .../yaak-app-client/src/plugin_events.rs | 2 + crates/yaak-plugins/bindings/gen_events.ts | 43 +++- crates/yaak-plugins/src/events.rs | 62 ++++++ crates/yaak/Cargo.toml | 1 + crates/yaak/src/lib.rs | 1 + crates/yaak/src/plugin_events.rs | 188 +++++++++++++++-- crates/yaak/src/response_body.rs | 197 ++++++++++++++++++ package-lock.json | 2 +- packages/plugin-runtime-types/package.json | 2 +- .../src/bindings/gen_events.ts | 43 +++- .../src/plugins/Context.ts | 57 +++++ .../plugin-runtime-types/src/plugins/index.ts | 7 +- packages/plugin-runtime/src/PluginInstance.ts | 36 +++- packages/plugin-runtime/src/responseBody.ts | 151 ++++++++++++++ .../plugin-runtime/tests/responseBody.test.ts | 122 +++++++++++ .../template-function-response/src/index.ts | 44 ++-- 18 files changed, 915 insertions(+), 46 deletions(-) create mode 100644 crates/yaak/src/response_body.rs create mode 100644 packages/plugin-runtime/src/responseBody.ts create mode 100644 packages/plugin-runtime/tests/responseBody.test.ts diff --git a/Cargo.lock b/Cargo.lock index 61734d63..305f1dcd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11196,6 +11196,7 @@ name = "yaak" version = "0.1.0" dependencies = [ "async-trait", + "base64 0.22.1", "log 0.4.29", "md5 0.8.0", "serde_json", diff --git a/crates-cli/yaak-cli/src/plugin_events.rs b/crates-cli/yaak-cli/src/plugin_events.rs index d232a88a..583d1a9a 100644 --- a/crates-cli/yaak-cli/src/plugin_events.rs +++ b/crates-cli/yaak-cli/src/plugin_events.rs @@ -12,6 +12,7 @@ use yaak::plugin_events::{ GroupedPluginEvent, HostRequest, SharedPluginEventContext, handle_shared_plugin_event, }; use yaak::render::{render_grpc_request, render_http_request}; +use yaak::response_body::FileResponseBodyStore; use yaak::send::{SendHttpRequestWithPluginsParams, send_http_request_with_plugins}; use yaak_crypto::manager::EncryptionManager; use yaak_http::cookies::get_cookie_value_from_jar; @@ -131,6 +132,7 @@ async fn build_plugin_reply( match handle_shared_plugin_event( &host_context.query_manager, + &FileResponseBodyStore::new(&host_context.query_manager), &event.payload, SharedPluginEventContext { plugin_name, workspace_id: shared_workspace_id }, ) { diff --git a/crates-tauri/yaak-app-client/src/plugin_events.rs b/crates-tauri/yaak-app-client/src/plugin_events.rs index 95fb3297..461f3b82 100644 --- a/crates-tauri/yaak-app-client/src/plugin_events.rs +++ b/crates-tauri/yaak-app-client/src/plugin_events.rs @@ -16,6 +16,7 @@ use tauri_plugin_opener::OpenerExt; use yaak::plugin_events::{ GroupedPluginEvent, HostRequest, SharedPluginEventContext, handle_shared_plugin_event, }; +use yaak::response_body::FileResponseBodyStore; use yaak_crypto::manager::EncryptionManager; use yaak_http::cookies::get_cookie_value_from_jar; use yaak_models::models::{HttpResponse, Plugin}; @@ -54,6 +55,7 @@ pub(crate) async fn handle_plugin_event( match handle_shared_plugin_event( app_handle.db_manager().inner(), + &FileResponseBodyStore::new(app_handle.db_manager().inner()), &event.payload, SharedPluginEventContext { plugin_name: &plugin_name, diff --git a/crates/yaak-plugins/bindings/gen_events.ts b/crates/yaak-plugins/bindings/gen_events.ts index 0dbd3212..3b25831a 100644 --- a/crates/yaak-plugins/bindings/gen_events.ts +++ b/crates/yaak-plugins/bindings/gen_events.ts @@ -416,6 +416,26 @@ export type GetHttpRequestByIdRequest = { id: string, }; export type GetHttpRequestByIdResponse = { httpRequest: HttpRequest | null, }; +/** + * Ask what a response's body is, before deciding whether to pull it. + * + * Bodies are addressed by response id and never by path, so where the host + * keeps the bytes is its own business. + */ +export type GetHttpResponseBodyInfoRequest = { responseId: string, }; + +export type GetHttpResponseBodyInfoResponse = { +/** + * How many bytes are actually stored, which is not necessarily what the + * `Content-Length` header claimed. Zero when the response has no body. + */ +contentLength: number, +/** + * The response's `Content-Type` header, verbatim, so the reader can pick a + * charset. + */ +contentType?: string | null, }; + export type GetKeyValueRequest = { key: string, }; export type GetKeyValueResponse = { value?: string, }; @@ -452,7 +472,7 @@ export type ImportResponse = { resources: ImportResources, }; export type InternalEvent = { id: string, pluginRefId: string, pluginName: string, replyId: string | null, context: PluginContext, payload: InternalEventPayload, }; -export type InternalEventPayload = { "type": "boot_request" } & BootRequest | { "type": "boot_response" } | { "type": "reload_response" } & ReloadResponse | { "type": "terminate_request" } | { "type": "terminate_response" } | { "type": "import_request" } & ImportRequest | { "type": "import_response" } & ImportResponse | { "type": "filter_request" } & FilterRequest | { "type": "filter_response" } & FilterResponse | { "type": "export_http_request_request" } & ExportHttpRequestRequest | { "type": "export_http_request_response" } & ExportHttpRequestResponse | { "type": "send_http_request_request" } & SendHttpRequestRequest | { "type": "send_http_request_response" } & SendHttpRequestResponse | { "type": "list_cookie_names_request" } & ListCookieNamesRequest | { "type": "list_cookie_names_response" } & ListCookieNamesResponse | { "type": "get_cookie_value_request" } & GetCookieValueRequest | { "type": "get_cookie_value_response" } & GetCookieValueResponse | { "type": "get_http_request_actions_request" } & EmptyPayload | { "type": "get_http_request_actions_response" } & GetHttpRequestActionsResponse | { "type": "call_http_request_action_request" } & CallHttpRequestActionRequest | { "type": "get_websocket_request_actions_request" } & EmptyPayload | { "type": "get_websocket_request_actions_response" } & GetWebsocketRequestActionsResponse | { "type": "call_websocket_request_action_request" } & CallWebsocketRequestActionRequest | { "type": "get_workspace_actions_request" } & EmptyPayload | { "type": "get_workspace_actions_response" } & GetWorkspaceActionsResponse | { "type": "call_workspace_action_request" } & CallWorkspaceActionRequest | { "type": "get_folder_actions_request" } & EmptyPayload | { "type": "get_folder_actions_response" } & GetFolderActionsResponse | { "type": "call_folder_action_request" } & CallFolderActionRequest | { "type": "get_grpc_request_actions_request" } & EmptyPayload | { "type": "get_grpc_request_actions_response" } & GetGrpcRequestActionsResponse | { "type": "call_grpc_request_action_request" } & CallGrpcRequestActionRequest | { "type": "get_template_function_summary_request" } & EmptyPayload | { "type": "get_template_function_summary_response" } & GetTemplateFunctionSummaryResponse | { "type": "get_template_function_config_request" } & GetTemplateFunctionConfigRequest | { "type": "get_template_function_config_response" } & GetTemplateFunctionConfigResponse | { "type": "call_template_function_request" } & CallTemplateFunctionRequest | { "type": "call_template_function_response" } & CallTemplateFunctionResponse | { "type": "get_http_authentication_summary_request" } & EmptyPayload | { "type": "get_http_authentication_summary_response" } & GetHttpAuthenticationSummaryResponse | { "type": "get_http_authentication_config_request" } & GetHttpAuthenticationConfigRequest | { "type": "get_http_authentication_config_response" } & GetHttpAuthenticationConfigResponse | { "type": "call_http_authentication_request" } & CallHttpAuthenticationRequest | { "type": "call_http_authentication_response" } & CallHttpAuthenticationResponse | { "type": "call_http_authentication_action_request" } & CallHttpAuthenticationActionRequest | { "type": "call_http_authentication_action_response" } & EmptyPayload | { "type": "copy_text_request" } & CopyTextRequest | { "type": "copy_text_response" } & EmptyPayload | { "type": "render_http_request_request" } & RenderHttpRequestRequest | { "type": "render_http_request_response" } & RenderHttpRequestResponse | { "type": "render_grpc_request_request" } & RenderGrpcRequestRequest | { "type": "render_grpc_request_response" } & RenderGrpcRequestResponse | { "type": "template_render_request" } & TemplateRenderRequest | { "type": "template_render_response" } & TemplateRenderResponse | { "type": "get_key_value_request" } & GetKeyValueRequest | { "type": "get_key_value_response" } & GetKeyValueResponse | { "type": "set_key_value_request" } & SetKeyValueRequest | { "type": "set_key_value_response" } & SetKeyValueResponse | { "type": "delete_key_value_request" } & DeleteKeyValueRequest | { "type": "delete_key_value_response" } & DeleteKeyValueResponse | { "type": "open_window_request" } & OpenWindowRequest | { "type": "window_navigate_event" } & WindowNavigateEvent | { "type": "window_close_event" } | { "type": "close_window_request" } & CloseWindowRequest | { "type": "open_external_url_request" } & OpenExternalUrlRequest | { "type": "open_external_url_response" } & EmptyPayload | { "type": "show_toast_request" } & ShowToastRequest | { "type": "show_toast_response" } & EmptyPayload | { "type": "prompt_text_request" } & PromptTextRequest | { "type": "prompt_text_response" } & PromptTextResponse | { "type": "prompt_form_request" } & PromptFormRequest | { "type": "prompt_form_response" } & PromptFormResponse | { "type": "window_info_request" } & WindowInfoRequest | { "type": "window_info_response" } & WindowInfoResponse | { "type": "list_open_workspaces_request" } & ListOpenWorkspacesRequest | { "type": "list_open_workspaces_response" } & ListOpenWorkspacesResponse | { "type": "get_http_request_by_id_request" } & GetHttpRequestByIdRequest | { "type": "get_http_request_by_id_response" } & GetHttpRequestByIdResponse | { "type": "find_http_responses_request" } & FindHttpResponsesRequest | { "type": "find_http_responses_response" } & FindHttpResponsesResponse | { "type": "list_http_requests_request" } & ListHttpRequestsRequest | { "type": "list_http_requests_response" } & ListHttpRequestsResponse | { "type": "list_folders_request" } & ListFoldersRequest | { "type": "list_folders_response" } & ListFoldersResponse | { "type": "upsert_model_request" } & UpsertModelRequest | { "type": "upsert_model_response" } & UpsertModelResponse | { "type": "delete_model_request" } & DeleteModelRequest | { "type": "delete_model_response" } & DeleteModelResponse | { "type": "get_themes_request" } & GetThemesRequest | { "type": "get_themes_response" } & GetThemesResponse | { "type": "empty_response" } & EmptyPayload | { "type": "error_response" } & ErrorResponse; +export type InternalEventPayload = { "type": "boot_request" } & BootRequest | { "type": "boot_response" } | { "type": "reload_response" } & ReloadResponse | { "type": "terminate_request" } | { "type": "terminate_response" } | { "type": "import_request" } & ImportRequest | { "type": "import_response" } & ImportResponse | { "type": "filter_request" } & FilterRequest | { "type": "filter_response" } & FilterResponse | { "type": "export_http_request_request" } & ExportHttpRequestRequest | { "type": "export_http_request_response" } & ExportHttpRequestResponse | { "type": "send_http_request_request" } & SendHttpRequestRequest | { "type": "send_http_request_response" } & SendHttpRequestResponse | { "type": "list_cookie_names_request" } & ListCookieNamesRequest | { "type": "list_cookie_names_response" } & ListCookieNamesResponse | { "type": "get_cookie_value_request" } & GetCookieValueRequest | { "type": "get_cookie_value_response" } & GetCookieValueResponse | { "type": "get_http_request_actions_request" } & EmptyPayload | { "type": "get_http_request_actions_response" } & GetHttpRequestActionsResponse | { "type": "call_http_request_action_request" } & CallHttpRequestActionRequest | { "type": "get_websocket_request_actions_request" } & EmptyPayload | { "type": "get_websocket_request_actions_response" } & GetWebsocketRequestActionsResponse | { "type": "call_websocket_request_action_request" } & CallWebsocketRequestActionRequest | { "type": "get_workspace_actions_request" } & EmptyPayload | { "type": "get_workspace_actions_response" } & GetWorkspaceActionsResponse | { "type": "call_workspace_action_request" } & CallWorkspaceActionRequest | { "type": "get_folder_actions_request" } & EmptyPayload | { "type": "get_folder_actions_response" } & GetFolderActionsResponse | { "type": "call_folder_action_request" } & CallFolderActionRequest | { "type": "get_grpc_request_actions_request" } & EmptyPayload | { "type": "get_grpc_request_actions_response" } & GetGrpcRequestActionsResponse | { "type": "call_grpc_request_action_request" } & CallGrpcRequestActionRequest | { "type": "get_template_function_summary_request" } & EmptyPayload | { "type": "get_template_function_summary_response" } & GetTemplateFunctionSummaryResponse | { "type": "get_template_function_config_request" } & GetTemplateFunctionConfigRequest | { "type": "get_template_function_config_response" } & GetTemplateFunctionConfigResponse | { "type": "call_template_function_request" } & CallTemplateFunctionRequest | { "type": "call_template_function_response" } & CallTemplateFunctionResponse | { "type": "get_http_authentication_summary_request" } & EmptyPayload | { "type": "get_http_authentication_summary_response" } & GetHttpAuthenticationSummaryResponse | { "type": "get_http_authentication_config_request" } & GetHttpAuthenticationConfigRequest | { "type": "get_http_authentication_config_response" } & GetHttpAuthenticationConfigResponse | { "type": "call_http_authentication_request" } & CallHttpAuthenticationRequest | { "type": "call_http_authentication_response" } & CallHttpAuthenticationResponse | { "type": "call_http_authentication_action_request" } & CallHttpAuthenticationActionRequest | { "type": "call_http_authentication_action_response" } & EmptyPayload | { "type": "copy_text_request" } & CopyTextRequest | { "type": "copy_text_response" } & EmptyPayload | { "type": "render_http_request_request" } & RenderHttpRequestRequest | { "type": "render_http_request_response" } & RenderHttpRequestResponse | { "type": "render_grpc_request_request" } & RenderGrpcRequestRequest | { "type": "render_grpc_request_response" } & RenderGrpcRequestResponse | { "type": "template_render_request" } & TemplateRenderRequest | { "type": "template_render_response" } & TemplateRenderResponse | { "type": "get_key_value_request" } & GetKeyValueRequest | { "type": "get_key_value_response" } & GetKeyValueResponse | { "type": "set_key_value_request" } & SetKeyValueRequest | { "type": "set_key_value_response" } & SetKeyValueResponse | { "type": "delete_key_value_request" } & DeleteKeyValueRequest | { "type": "delete_key_value_response" } & DeleteKeyValueResponse | { "type": "open_window_request" } & OpenWindowRequest | { "type": "window_navigate_event" } & WindowNavigateEvent | { "type": "window_close_event" } | { "type": "close_window_request" } & CloseWindowRequest | { "type": "open_external_url_request" } & OpenExternalUrlRequest | { "type": "open_external_url_response" } & EmptyPayload | { "type": "show_toast_request" } & ShowToastRequest | { "type": "show_toast_response" } & EmptyPayload | { "type": "prompt_text_request" } & PromptTextRequest | { "type": "prompt_text_response" } & PromptTextResponse | { "type": "prompt_form_request" } & PromptFormRequest | { "type": "prompt_form_response" } & PromptFormResponse | { "type": "window_info_request" } & WindowInfoRequest | { "type": "window_info_response" } & WindowInfoResponse | { "type": "list_open_workspaces_request" } & ListOpenWorkspacesRequest | { "type": "list_open_workspaces_response" } & ListOpenWorkspacesResponse | { "type": "get_http_request_by_id_request" } & GetHttpRequestByIdRequest | { "type": "get_http_request_by_id_response" } & GetHttpRequestByIdResponse | { "type": "find_http_responses_request" } & FindHttpResponsesRequest | { "type": "find_http_responses_response" } & FindHttpResponsesResponse | { "type": "get_http_response_body_info_request" } & GetHttpResponseBodyInfoRequest | { "type": "get_http_response_body_info_response" } & GetHttpResponseBodyInfoResponse | { "type": "read_http_response_body_chunk_request" } & ReadHttpResponseBodyChunkRequest | { "type": "read_http_response_body_chunk_response" } & ReadHttpResponseBodyChunkResponse | { "type": "list_http_requests_request" } & ListHttpRequestsRequest | { "type": "list_http_requests_response" } & ListHttpRequestsResponse | { "type": "list_folders_request" } & ListFoldersRequest | { "type": "list_folders_response" } & ListFoldersResponse | { "type": "upsert_model_request" } & UpsertModelRequest | { "type": "upsert_model_response" } & UpsertModelResponse | { "type": "delete_model_request" } & DeleteModelRequest | { "type": "delete_model_response" } & DeleteModelResponse | { "type": "get_themes_request" } & GetThemesRequest | { "type": "get_themes_response" } & GetThemesResponse | { "type": "empty_response" } & EmptyPayload | { "type": "error_response" } & ErrorResponse; export type JsonPrimitive = string | number | boolean | null; @@ -502,6 +522,27 @@ required?: boolean, }; export type PromptTextResponse = { value: string | null, }; +/** + * Pull one window of a response body. + * + * Reads are idempotent: the bytes live in durable storage, so the same window + * can be asked for as many times as the plugin likes. + */ +export type ReadHttpResponseBodyChunkRequest = { responseId: string, offset: number, length: number, }; + +export type ReadHttpResponseBodyChunkResponse = { +/** + * Base64, because the desktop transport is a WebSocket that only sends + * text frames today. A host that can carry binary sends the bytes as they + * are and fills this in from them. + */ +data: string, +/** + * Bytes decoded from `data`. Short of the requested length means the body + * ended here. + */ +length: number, }; + export type ReloadResponse = { silent: boolean, }; export type RenderGrpcRequestRequest = { grpcRequest: GrpcRequest, purpose: RenderPurpose, }; diff --git a/crates/yaak-plugins/src/events.rs b/crates/yaak-plugins/src/events.rs index 5fd480e5..5faeb3c0 100644 --- a/crates/yaak-plugins/src/events.rs +++ b/crates/yaak-plugins/src/events.rs @@ -171,6 +171,12 @@ pub enum InternalEventPayload { FindHttpResponsesRequest(FindHttpResponsesRequest), FindHttpResponsesResponse(FindHttpResponsesResponse), + + GetHttpResponseBodyInfoRequest(GetHttpResponseBodyInfoRequest), + GetHttpResponseBodyInfoResponse(GetHttpResponseBodyInfoResponse), + ReadHttpResponseBodyChunkRequest(ReadHttpResponseBodyChunkRequest), + ReadHttpResponseBodyChunkResponse(ReadHttpResponseBodyChunkResponse), + ListHttpRequestsRequest(ListHttpRequestsRequest), ListHttpRequestsResponse(ListHttpRequestsResponse), ListFoldersRequest(ListFoldersRequest), @@ -1413,6 +1419,62 @@ pub struct FindHttpResponsesResponse { pub http_responses: Vec, } +/// Ask what a response's body is, before deciding whether to pull it. +/// +/// Bodies are addressed by response id and never by path, so where the host +/// keeps the bytes is its own business. +#[derive(Debug, Clone, Default, Serialize, Deserialize, TS)] +#[serde(default, rename_all = "camelCase")] +#[ts(export, export_to = "gen_events.ts")] +pub struct GetHttpResponseBodyInfoRequest { + pub response_id: String, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, TS)] +#[serde(default, rename_all = "camelCase")] +#[ts(export, export_to = "gen_events.ts")] +pub struct GetHttpResponseBodyInfoResponse { + /// How many bytes are actually stored, which is not necessarily what the + /// `Content-Length` header claimed. Zero when the response has no body. + #[ts(type = "number")] + pub content_length: u64, + + /// The response's `Content-Type` header, verbatim, so the reader can pick a + /// charset. + #[ts(optional = nullable)] + pub content_type: Option, +} + +/// Pull one window of a response body. +/// +/// Reads are idempotent: the bytes live in durable storage, so the same window +/// can be asked for as many times as the plugin likes. +#[derive(Debug, Clone, Default, Serialize, Deserialize, TS)] +#[serde(default, rename_all = "camelCase")] +#[ts(export, export_to = "gen_events.ts")] +pub struct ReadHttpResponseBodyChunkRequest { + pub response_id: String, + #[ts(type = "number")] + pub offset: u64, + #[ts(type = "number")] + pub length: u64, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, TS)] +#[serde(default, rename_all = "camelCase")] +#[ts(export, export_to = "gen_events.ts")] +pub struct ReadHttpResponseBodyChunkResponse { + /// Base64, because the desktop transport is a WebSocket that only sends + /// text frames today. A host that can carry binary sends the bytes as they + /// are and fills this in from them. + pub data: String, + + /// Bytes decoded from `data`. Short of the requested length means the body + /// ended here. + #[ts(type = "number")] + pub length: u64, +} + #[derive(Debug, Clone, Default, Serialize, Deserialize, TS)] #[serde(default, rename_all = "camelCase")] #[ts(export, export_to = "gen_events.ts")] diff --git a/crates/yaak/Cargo.toml b/crates/yaak/Cargo.toml index b0af1b57..03b8e161 100644 --- a/crates/yaak/Cargo.toml +++ b/crates/yaak/Cargo.toml @@ -6,6 +6,7 @@ publish = false [dependencies] async-trait = "0.1" +base64 = "0.22.1" # For carrying body chunks over a text-only plugin transport log = { workspace = true } md5 = "0.8.0" serde_json = { workspace = true } diff --git a/crates/yaak/src/lib.rs b/crates/yaak/src/lib.rs index 7bc790e8..e21be7e3 100644 --- a/crates/yaak/src/lib.rs +++ b/crates/yaak/src/lib.rs @@ -3,6 +3,7 @@ pub mod export; pub mod import; pub mod plugin_events; pub mod render; +pub mod response_body; pub mod send; pub use error::Error; diff --git a/crates/yaak/src/plugin_events.rs b/crates/yaak/src/plugin_events.rs index 96e2e4df..ef367ada 100644 --- a/crates/yaak/src/plugin_events.rs +++ b/crates/yaak/src/plugin_events.rs @@ -1,3 +1,6 @@ +use crate::response_body::ResponseBodyStore; +use base64::Engine; +use base64::prelude::BASE64_STANDARD; use yaak_models::models::AnyModel; use yaak_models::query_manager::QueryManager; use yaak_models::util::UpdateSource; @@ -5,12 +8,14 @@ use yaak_plugins::events::{ CloseWindowRequest, CopyTextRequest, DeleteKeyValueRequest, DeleteKeyValueResponse, DeleteModelRequest, DeleteModelResponse, ErrorResponse, FindHttpResponsesRequest, FindHttpResponsesResponse, GetCookieValueRequest, GetHttpRequestByIdRequest, - GetHttpRequestByIdResponse, GetKeyValueRequest, GetKeyValueResponse, InternalEventPayload, - ListCookieNamesRequest, ListFoldersRequest, ListFoldersResponse, ListHttpRequestsRequest, - ListHttpRequestsResponse, ListOpenWorkspacesRequest, OpenExternalUrlRequest, OpenWindowRequest, - PromptFormRequest, PromptTextRequest, ReloadResponse, RenderGrpcRequestRequest, - RenderHttpRequestRequest, SendHttpRequestRequest, SetKeyValueRequest, ShowToastRequest, - TemplateRenderRequest, UpsertModelRequest, UpsertModelResponse, WindowInfoRequest, + GetHttpRequestByIdResponse, GetHttpResponseBodyInfoRequest, GetHttpResponseBodyInfoResponse, + GetKeyValueRequest, GetKeyValueResponse, InternalEventPayload, ListCookieNamesRequest, + ListFoldersRequest, ListFoldersResponse, ListHttpRequestsRequest, ListHttpRequestsResponse, + ListOpenWorkspacesRequest, OpenExternalUrlRequest, OpenWindowRequest, PromptFormRequest, + PromptTextRequest, ReadHttpResponseBodyChunkRequest, ReadHttpResponseBodyChunkResponse, + ReloadResponse, RenderGrpcRequestRequest, RenderHttpRequestRequest, SendHttpRequestRequest, + SetKeyValueRequest, ShowToastRequest, TemplateRenderRequest, UpsertModelRequest, + UpsertModelResponse, WindowInfoRequest, }; pub struct SharedPluginEventContext<'a> { @@ -40,6 +45,8 @@ pub enum SharedRequest<'a> { ListFolders(&'a ListFoldersRequest), ListHttpRequests(&'a ListHttpRequestsRequest), FindHttpResponses(&'a FindHttpResponsesRequest), + GetHttpResponseBodyInfo(&'a GetHttpResponseBodyInfoRequest), + ReadHttpResponseBodyChunk(&'a ReadHttpResponseBodyChunkRequest), UpsertModel(&'a UpsertModelRequest), DeleteModel(&'a DeleteModelRequest), } @@ -136,6 +143,12 @@ impl<'a> From<&'a InternalEventPayload> for GroupedPluginRequest<'a> { InternalEventPayload::FindHttpResponsesRequest(req) => { GroupedPluginRequest::Shared(SharedRequest::FindHttpResponses(req)) } + InternalEventPayload::GetHttpResponseBodyInfoRequest(req) => { + GroupedPluginRequest::Shared(SharedRequest::GetHttpResponseBodyInfo(req)) + } + InternalEventPayload::ReadHttpResponseBodyChunkRequest(req) => { + GroupedPluginRequest::Shared(SharedRequest::ReadHttpResponseBodyChunk(req)) + } InternalEventPayload::UpsertModelRequest(req) => { GroupedPluginRequest::Shared(SharedRequest::UpsertModel(req)) } @@ -182,13 +195,17 @@ impl<'a> From<&'a InternalEventPayload> for GroupedPluginRequest<'a> { pub fn handle_shared_plugin_event<'a>( query_manager: &QueryManager, + body_store: &dyn ResponseBodyStore, payload: &'a InternalEventPayload, context: SharedPluginEventContext<'_>, ) -> GroupedPluginEvent<'a> { match GroupedPluginRequest::from(payload) { - GroupedPluginRequest::Shared(req) => { - GroupedPluginEvent::Handled(Some(build_shared_reply(query_manager, req, context))) - } + GroupedPluginRequest::Shared(req) => GroupedPluginEvent::Handled(Some(build_shared_reply( + query_manager, + body_store, + req, + context, + ))), GroupedPluginRequest::Host(req) => GroupedPluginEvent::ToHandle(req), GroupedPluginRequest::Ignore => GroupedPluginEvent::Handled(None), } @@ -196,6 +213,7 @@ pub fn handle_shared_plugin_event<'a>( fn build_shared_reply( query_manager: &QueryManager, + body_store: &dyn ResponseBodyStore, request: SharedRequest<'_>, context: SharedPluginEventContext<'_>, ) -> InternalEventPayload { @@ -283,6 +301,30 @@ fn build_shared_reply( http_responses, }) } + SharedRequest::GetHttpResponseBodyInfo(req) => match body_store.info(&req.response_id) { + Ok(info) => InternalEventPayload::GetHttpResponseBodyInfoResponse( + GetHttpResponseBodyInfoResponse { + content_length: info.content_length, + content_type: info.content_type, + }, + ), + Err(err) => InternalEventPayload::ErrorResponse(ErrorResponse { + error: format!("Failed to read body of response {}: {err}", req.response_id), + }), + }, + SharedRequest::ReadHttpResponseBodyChunk(req) => { + match body_store.read_chunk(&req.response_id, req.offset, req.length) { + Ok(bytes) => InternalEventPayload::ReadHttpResponseBodyChunkResponse( + ReadHttpResponseBodyChunkResponse { + length: bytes.len() as u64, + data: BASE64_STANDARD.encode(bytes), + }, + ), + Err(err) => InternalEventPayload::ErrorResponse(ErrorResponse { + error: format!("Failed to read body of response {}: {err}", req.response_id), + }), + } + } SharedRequest::UpsertModel(req) => { use AnyModel::*; @@ -437,10 +479,26 @@ fn build_shared_reply( #[cfg(test)] mod tests { use super::*; + use crate::response_body::{FileResponseBodyStore, ResponseBodyInfo}; + use std::cell::RefCell; use tempfile::TempDir; use yaak_models::models::{AnyModel, Folder, HttpRequest, Workspace}; use yaak_models::util::UpdateSource; + /// The real dispatch, with the store the desktop and CLI hand it. + fn dispatch<'a>( + query_manager: &QueryManager, + payload: &'a InternalEventPayload, + context: SharedPluginEventContext<'_>, + ) -> GroupedPluginEvent<'a> { + handle_shared_plugin_event( + query_manager, + &FileResponseBodyStore::new(query_manager), + payload, + context, + ) + } + fn seed_query_manager() -> (QueryManager, TempDir) { let temp_dir = TempDir::new().expect("Failed to create temp dir"); let db_path = temp_dir.path().join("db.sqlite"); @@ -498,7 +556,7 @@ mod tests { let payload = InternalEventPayload::ListHttpRequestsRequest( yaak_plugins::events::ListHttpRequestsRequest { folder_id: None }, ); - let result = handle_shared_plugin_event( + let result = dispatch( &query_manager, &payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, @@ -517,7 +575,7 @@ mod tests { let by_workspace_payload = InternalEventPayload::ListHttpRequestsRequest( yaak_plugins::events::ListHttpRequestsRequest { folder_id: None }, ); - let by_workspace = handle_shared_plugin_event( + let by_workspace = dispatch( &query_manager, &by_workspace_payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: Some("wk_test") }, @@ -536,7 +594,7 @@ mod tests { folder_id: Some("fl_test".to_string()), }, ); - let by_folder = handle_shared_plugin_event( + let by_folder = dispatch( &query_manager, &by_folder_payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, @@ -559,7 +617,7 @@ mod tests { limit: Some(1), }); - let result = handle_shared_plugin_event( + let result = dispatch( &query_manager, &payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: Some("wk_test") }, @@ -575,6 +633,104 @@ mod tests { } } + /// A store that answers from memory, standing in for whatever holds the + /// bytes — the point being that the dispatch below never learns which. + struct FakeBodyStore { + body: Vec, + reads: RefCell>, + } + + impl ResponseBodyStore for FakeBodyStore { + fn info(&self, _response_id: &str) -> crate::error::Result { + Ok(ResponseBodyInfo { + content_length: self.body.len() as u64, + content_type: Some("text/plain; charset=utf-8".to_string()), + }) + } + + fn read_chunk( + &self, + _response_id: &str, + offset: u64, + length: u64, + ) -> crate::error::Result> { + self.reads.borrow_mut().push((offset, length)); + let start = (offset as usize).min(self.body.len()); + let end = (start + length as usize).min(self.body.len()); + Ok(self.body[start..end].to_vec()) + } + } + + #[test] + fn response_body_is_read_by_id_through_the_store() { + let (query_manager, _temp_dir) = seed_query_manager(); + let store = FakeBodyStore { body: b"hello".to_vec(), reads: RefCell::new(Vec::new()) }; + + let info_payload = InternalEventPayload::GetHttpResponseBodyInfoRequest( + GetHttpResponseBodyInfoRequest { response_id: "rs_test".to_string() }, + ); + let info = handle_shared_plugin_event( + &query_manager, + &store, + &info_payload, + SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, + ); + match info { + GroupedPluginEvent::Handled(Some( + InternalEventPayload::GetHttpResponseBodyInfoResponse(resp), + )) => { + assert_eq!(resp.content_length, 5); + assert_eq!(resp.content_type.as_deref(), Some("text/plain; charset=utf-8")); + } + other => panic!("unexpected body info result: {other:?}"), + } + + let chunk_payload = InternalEventPayload::ReadHttpResponseBodyChunkRequest( + ReadHttpResponseBodyChunkRequest { + response_id: "rs_test".to_string(), + offset: 1, + length: 3, + }, + ); + let chunk = handle_shared_plugin_event( + &query_manager, + &store, + &chunk_payload, + SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, + ); + match chunk { + GroupedPluginEvent::Handled(Some( + InternalEventPayload::ReadHttpResponseBodyChunkResponse(resp), + )) => { + assert_eq!(resp.length, 3); + assert_eq!(BASE64_STANDARD.decode(resp.data).unwrap(), b"ell"); + } + other => panic!("unexpected body chunk result: {other:?}"), + } + + assert_eq!(*store.reads.borrow(), vec![(1, 3)]); + } + + #[test] + fn an_unreadable_response_body_becomes_an_error_reply() { + let (query_manager, _temp_dir) = seed_query_manager(); + let payload = InternalEventPayload::GetHttpResponseBodyInfoRequest( + GetHttpResponseBodyInfoRequest { response_id: "rs_never_persisted".to_string() }, + ); + let result = dispatch( + &query_manager, + &payload, + SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, + ); + + match result { + GroupedPluginEvent::Handled(Some(InternalEventPayload::ErrorResponse(resp))) => { + assert!(resp.error.contains("rs_never_persisted"), "unhelpful error: {}", resp.error) + } + other => panic!("unexpected missing-response result: {other:?}"), + } + } + #[test] fn upsert_and_delete_model_are_shared_handled() { let (query_manager, _temp_dir) = seed_query_manager(); @@ -590,7 +746,7 @@ mod tests { }), }); - let upsert_result = handle_shared_plugin_event( + let upsert_result = dispatch( &query_manager, &upsert_payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: Some("wk_test") }, @@ -609,7 +765,7 @@ mod tests { model: "http_request".to_string(), id: "rq_test".to_string(), }); - let delete_result = handle_shared_plugin_event( + let delete_result = dispatch( &query_manager, &delete_payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: Some("wk_test") }, @@ -631,7 +787,7 @@ mod tests { let payload = InternalEventPayload::WindowInfoRequest(WindowInfoRequest { label: "main".to_string(), }); - let result = handle_shared_plugin_event( + let result = dispatch( &query_manager, &payload, SharedPluginEventContext { plugin_name: "@yaak/test", workspace_id: None }, diff --git a/crates/yaak/src/response_body.rs b/crates/yaak/src/response_body.rs new file mode 100644 index 00000000..f0d7945f --- /dev/null +++ b/crates/yaak/src/response_body.rs @@ -0,0 +1,197 @@ +//! Reading response bodies back out, by response id. +//! +//! Plugins only ever name a response. Where its bytes actually live — files the +//! engine wrote under `/responses/` today, blob rows later — is +//! behind [`ResponseBodyStore`], so moving the bytes is a change to this file +//! and nothing a plugin can see. + +use crate::error::Result; +use std::fs::File; +use std::io::{Read, Seek, SeekFrom}; +use yaak_models::query_manager::QueryManager; + +/// The most bytes one read will hand back, however much was asked for. +/// +/// A chunk is buffered whole and, on the desktop transport, base64'd into a +/// single WebSocket frame, so an unbounded request is a way to make the host +/// allocate on a plugin's say-so. +pub const MAX_CHUNK_BYTES: u64 = 8 * 1024 * 1024; + +/// What a stored body is, without reading any of it. +#[derive(Debug, Clone, Default)] +pub struct ResponseBodyInfo { + /// Bytes actually stored, which is not necessarily what `Content-Length` + /// claimed. Zero when the response has no body. + pub content_length: u64, + /// The response's `Content-Type` header, verbatim. + pub content_type: Option, +} + +/// Somewhere response bodies can be read from, a window at a time. +/// +/// Reads are repeatable — the bytes are durable, so nothing is consumed by +/// looking at it. +pub trait ResponseBodyStore { + fn info(&self, response_id: &str) -> Result; + + /// Bytes `[offset, offset + length)`, clamped to what is there. A short + /// read means the body ended. + fn read_chunk(&self, response_id: &str, offset: u64, length: u64) -> Result>; +} + +/// The desktop and CLI store: the database says where the file is, and the +/// filesystem holds it. +pub struct FileResponseBodyStore<'a> { + query_manager: &'a QueryManager, +} + +impl<'a> FileResponseBodyStore<'a> { + pub fn new(query_manager: &'a QueryManager) -> Self { + Self { query_manager } + } + + /// The file backing a response, or `None` when the response stored no body. + /// + /// A response that was never persisted (GraphQL introspection and other + /// ephemeral sends) has no row here at all, so it fails as "not found" + /// rather than reading as an empty body. + fn body_path(&self, response_id: &str) -> Result> { + Ok(self.query_manager.connect().get_http_response(response_id)?.body_path) + } +} + +impl ResponseBodyStore for FileResponseBodyStore<'_> { + fn info(&self, response_id: &str) -> Result { + let response = self.query_manager.connect().get_http_response(response_id)?; + + let content_type = response + .headers + .iter() + .find(|h| h.name.eq_ignore_ascii_case("content-type")) + .map(|h| h.value.clone()); + + let content_length = match response.body_path { + Some(path) => std::fs::metadata(path)?.len(), + None => 0, + }; + + Ok(ResponseBodyInfo { content_length, content_type }) + } + + fn read_chunk(&self, response_id: &str, offset: u64, length: u64) -> Result> { + let Some(path) = self.body_path(response_id)? else { + return Ok(Vec::new()); + }; + + let length = length.min(MAX_CHUNK_BYTES); + if length == 0 { + return Ok(Vec::new()); + } + + let mut file = File::open(path)?; + file.seek(SeekFrom::Start(offset))?; + + let mut buf = Vec::new(); + file.take(length).read_to_end(&mut buf)?; + Ok(buf) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Write; + use tempfile::TempDir; + use yaak_models::models::{HttpRequest, HttpResponse, HttpResponseHeader, Workspace}; + use yaak_models::util::UpdateSource; + + fn seed(body: Option<&[u8]>) -> (QueryManager, TempDir, String) { + let temp_dir = TempDir::new().unwrap(); + let (query_manager, blob_manager, _rx) = yaak_models::init_standalone( + &temp_dir.path().join("db.sqlite"), + &temp_dir.path().join("blobs.sqlite"), + ) + .unwrap(); + + query_manager + .connect() + .upsert_workspace( + &Workspace { id: "wk_test".to_string(), ..Default::default() }, + &UpdateSource::Sync, + ) + .unwrap(); + + query_manager + .connect() + .upsert_http_request( + &HttpRequest { + id: "rq_test".to_string(), + workspace_id: "wk_test".to_string(), + ..Default::default() + }, + &UpdateSource::Sync, + ) + .unwrap(); + + let body_path = body.map(|bytes| { + let path = temp_dir.path().join("body"); + let mut f = std::fs::File::create(&path).unwrap(); + f.write_all(bytes).unwrap(); + path.to_string_lossy().to_string() + }); + + let response = query_manager + .connect() + .upsert_http_response( + &HttpResponse { + workspace_id: "wk_test".to_string(), + request_id: "rq_test".to_string(), + body_path, + headers: vec![HttpResponseHeader { + name: "Content-Type".to_string(), + value: "application/json; charset=utf-8".to_string(), + }], + ..Default::default() + }, + &UpdateSource::Sync, + &blob_manager, + ) + .unwrap(); + + let id = response.id.clone(); + (query_manager, temp_dir, id) + } + + #[test] + fn info_reports_stored_size_and_content_type() { + let (qm, _tmp, id) = seed(Some(b"hello world")); + let info = FileResponseBodyStore::new(&qm).info(&id).unwrap(); + assert_eq!(info.content_length, 11); + assert_eq!(info.content_type.as_deref(), Some("application/json; charset=utf-8")); + } + + #[test] + fn chunks_cover_the_body_and_stop_short_at_the_end() { + let (qm, _tmp, id) = seed(Some(b"hello world")); + let store = FileResponseBodyStore::new(&qm); + assert_eq!(store.read_chunk(&id, 0, 5).unwrap(), b"hello"); + assert_eq!(store.read_chunk(&id, 6, 100).unwrap(), b"world"); + assert!(store.read_chunk(&id, 11, 100).unwrap().is_empty()); + // Reading the same window twice gives the same bytes; nothing is consumed. + assert_eq!(store.read_chunk(&id, 0, 5).unwrap(), b"hello"); + } + + #[test] + fn a_response_with_no_body_is_empty_not_an_error() { + let (qm, _tmp, id) = seed(None); + let store = FileResponseBodyStore::new(&qm); + assert_eq!(store.info(&id).unwrap().content_length, 0); + assert!(store.read_chunk(&id, 0, 100).unwrap().is_empty()); + } + + #[test] + fn an_unknown_response_fails() { + let (qm, _tmp, _id) = seed(Some(b"hi")); + assert!(FileResponseBodyStore::new(&qm).info("rs_nope").is_err()); + } +} diff --git a/package-lock.json b/package-lock.json index 2187134b..a2d37125 100644 --- a/package-lock.json +++ b/package-lock.json @@ -17226,7 +17226,7 @@ }, "packages/plugin-runtime-types": { "name": "@yaakapp/api", - "version": "0.8.0", + "version": "0.9.0", "dependencies": { "@types/node": "^24.0.13" }, diff --git a/packages/plugin-runtime-types/package.json b/packages/plugin-runtime-types/package.json index d4dd862d..f282c3f0 100644 --- a/packages/plugin-runtime-types/package.json +++ b/packages/plugin-runtime-types/package.json @@ -1,6 +1,6 @@ { "name": "@yaakapp/api", - "version": "0.8.0", + "version": "0.9.0", "keywords": [ "api-client", "bruno-alternative", diff --git a/packages/plugin-runtime-types/src/bindings/gen_events.ts b/packages/plugin-runtime-types/src/bindings/gen_events.ts index 0dbd3212..3b25831a 100644 --- a/packages/plugin-runtime-types/src/bindings/gen_events.ts +++ b/packages/plugin-runtime-types/src/bindings/gen_events.ts @@ -416,6 +416,26 @@ export type GetHttpRequestByIdRequest = { id: string, }; export type GetHttpRequestByIdResponse = { httpRequest: HttpRequest | null, }; +/** + * Ask what a response's body is, before deciding whether to pull it. + * + * Bodies are addressed by response id and never by path, so where the host + * keeps the bytes is its own business. + */ +export type GetHttpResponseBodyInfoRequest = { responseId: string, }; + +export type GetHttpResponseBodyInfoResponse = { +/** + * How many bytes are actually stored, which is not necessarily what the + * `Content-Length` header claimed. Zero when the response has no body. + */ +contentLength: number, +/** + * The response's `Content-Type` header, verbatim, so the reader can pick a + * charset. + */ +contentType?: string | null, }; + export type GetKeyValueRequest = { key: string, }; export type GetKeyValueResponse = { value?: string, }; @@ -452,7 +472,7 @@ export type ImportResponse = { resources: ImportResources, }; export type InternalEvent = { id: string, pluginRefId: string, pluginName: string, replyId: string | null, context: PluginContext, payload: InternalEventPayload, }; -export type InternalEventPayload = { "type": "boot_request" } & BootRequest | { "type": "boot_response" } | { "type": "reload_response" } & ReloadResponse | { "type": "terminate_request" } | { "type": "terminate_response" } | { "type": "import_request" } & ImportRequest | { "type": "import_response" } & ImportResponse | { "type": "filter_request" } & FilterRequest | { "type": "filter_response" } & FilterResponse | { "type": "export_http_request_request" } & ExportHttpRequestRequest | { "type": "export_http_request_response" } & ExportHttpRequestResponse | { "type": "send_http_request_request" } & SendHttpRequestRequest | { "type": "send_http_request_response" } & SendHttpRequestResponse | { "type": "list_cookie_names_request" } & ListCookieNamesRequest | { "type": "list_cookie_names_response" } & ListCookieNamesResponse | { "type": "get_cookie_value_request" } & GetCookieValueRequest | { "type": "get_cookie_value_response" } & GetCookieValueResponse | { "type": "get_http_request_actions_request" } & EmptyPayload | { "type": "get_http_request_actions_response" } & GetHttpRequestActionsResponse | { "type": "call_http_request_action_request" } & CallHttpRequestActionRequest | { "type": "get_websocket_request_actions_request" } & EmptyPayload | { "type": "get_websocket_request_actions_response" } & GetWebsocketRequestActionsResponse | { "type": "call_websocket_request_action_request" } & CallWebsocketRequestActionRequest | { "type": "get_workspace_actions_request" } & EmptyPayload | { "type": "get_workspace_actions_response" } & GetWorkspaceActionsResponse | { "type": "call_workspace_action_request" } & CallWorkspaceActionRequest | { "type": "get_folder_actions_request" } & EmptyPayload | { "type": "get_folder_actions_response" } & GetFolderActionsResponse | { "type": "call_folder_action_request" } & CallFolderActionRequest | { "type": "get_grpc_request_actions_request" } & EmptyPayload | { "type": "get_grpc_request_actions_response" } & GetGrpcRequestActionsResponse | { "type": "call_grpc_request_action_request" } & CallGrpcRequestActionRequest | { "type": "get_template_function_summary_request" } & EmptyPayload | { "type": "get_template_function_summary_response" } & GetTemplateFunctionSummaryResponse | { "type": "get_template_function_config_request" } & GetTemplateFunctionConfigRequest | { "type": "get_template_function_config_response" } & GetTemplateFunctionConfigResponse | { "type": "call_template_function_request" } & CallTemplateFunctionRequest | { "type": "call_template_function_response" } & CallTemplateFunctionResponse | { "type": "get_http_authentication_summary_request" } & EmptyPayload | { "type": "get_http_authentication_summary_response" } & GetHttpAuthenticationSummaryResponse | { "type": "get_http_authentication_config_request" } & GetHttpAuthenticationConfigRequest | { "type": "get_http_authentication_config_response" } & GetHttpAuthenticationConfigResponse | { "type": "call_http_authentication_request" } & CallHttpAuthenticationRequest | { "type": "call_http_authentication_response" } & CallHttpAuthenticationResponse | { "type": "call_http_authentication_action_request" } & CallHttpAuthenticationActionRequest | { "type": "call_http_authentication_action_response" } & EmptyPayload | { "type": "copy_text_request" } & CopyTextRequest | { "type": "copy_text_response" } & EmptyPayload | { "type": "render_http_request_request" } & RenderHttpRequestRequest | { "type": "render_http_request_response" } & RenderHttpRequestResponse | { "type": "render_grpc_request_request" } & RenderGrpcRequestRequest | { "type": "render_grpc_request_response" } & RenderGrpcRequestResponse | { "type": "template_render_request" } & TemplateRenderRequest | { "type": "template_render_response" } & TemplateRenderResponse | { "type": "get_key_value_request" } & GetKeyValueRequest | { "type": "get_key_value_response" } & GetKeyValueResponse | { "type": "set_key_value_request" } & SetKeyValueRequest | { "type": "set_key_value_response" } & SetKeyValueResponse | { "type": "delete_key_value_request" } & DeleteKeyValueRequest | { "type": "delete_key_value_response" } & DeleteKeyValueResponse | { "type": "open_window_request" } & OpenWindowRequest | { "type": "window_navigate_event" } & WindowNavigateEvent | { "type": "window_close_event" } | { "type": "close_window_request" } & CloseWindowRequest | { "type": "open_external_url_request" } & OpenExternalUrlRequest | { "type": "open_external_url_response" } & EmptyPayload | { "type": "show_toast_request" } & ShowToastRequest | { "type": "show_toast_response" } & EmptyPayload | { "type": "prompt_text_request" } & PromptTextRequest | { "type": "prompt_text_response" } & PromptTextResponse | { "type": "prompt_form_request" } & PromptFormRequest | { "type": "prompt_form_response" } & PromptFormResponse | { "type": "window_info_request" } & WindowInfoRequest | { "type": "window_info_response" } & WindowInfoResponse | { "type": "list_open_workspaces_request" } & ListOpenWorkspacesRequest | { "type": "list_open_workspaces_response" } & ListOpenWorkspacesResponse | { "type": "get_http_request_by_id_request" } & GetHttpRequestByIdRequest | { "type": "get_http_request_by_id_response" } & GetHttpRequestByIdResponse | { "type": "find_http_responses_request" } & FindHttpResponsesRequest | { "type": "find_http_responses_response" } & FindHttpResponsesResponse | { "type": "list_http_requests_request" } & ListHttpRequestsRequest | { "type": "list_http_requests_response" } & ListHttpRequestsResponse | { "type": "list_folders_request" } & ListFoldersRequest | { "type": "list_folders_response" } & ListFoldersResponse | { "type": "upsert_model_request" } & UpsertModelRequest | { "type": "upsert_model_response" } & UpsertModelResponse | { "type": "delete_model_request" } & DeleteModelRequest | { "type": "delete_model_response" } & DeleteModelResponse | { "type": "get_themes_request" } & GetThemesRequest | { "type": "get_themes_response" } & GetThemesResponse | { "type": "empty_response" } & EmptyPayload | { "type": "error_response" } & ErrorResponse; +export type InternalEventPayload = { "type": "boot_request" } & BootRequest | { "type": "boot_response" } | { "type": "reload_response" } & ReloadResponse | { "type": "terminate_request" } | { "type": "terminate_response" } | { "type": "import_request" } & ImportRequest | { "type": "import_response" } & ImportResponse | { "type": "filter_request" } & FilterRequest | { "type": "filter_response" } & FilterResponse | { "type": "export_http_request_request" } & ExportHttpRequestRequest | { "type": "export_http_request_response" } & ExportHttpRequestResponse | { "type": "send_http_request_request" } & SendHttpRequestRequest | { "type": "send_http_request_response" } & SendHttpRequestResponse | { "type": "list_cookie_names_request" } & ListCookieNamesRequest | { "type": "list_cookie_names_response" } & ListCookieNamesResponse | { "type": "get_cookie_value_request" } & GetCookieValueRequest | { "type": "get_cookie_value_response" } & GetCookieValueResponse | { "type": "get_http_request_actions_request" } & EmptyPayload | { "type": "get_http_request_actions_response" } & GetHttpRequestActionsResponse | { "type": "call_http_request_action_request" } & CallHttpRequestActionRequest | { "type": "get_websocket_request_actions_request" } & EmptyPayload | { "type": "get_websocket_request_actions_response" } & GetWebsocketRequestActionsResponse | { "type": "call_websocket_request_action_request" } & CallWebsocketRequestActionRequest | { "type": "get_workspace_actions_request" } & EmptyPayload | { "type": "get_workspace_actions_response" } & GetWorkspaceActionsResponse | { "type": "call_workspace_action_request" } & CallWorkspaceActionRequest | { "type": "get_folder_actions_request" } & EmptyPayload | { "type": "get_folder_actions_response" } & GetFolderActionsResponse | { "type": "call_folder_action_request" } & CallFolderActionRequest | { "type": "get_grpc_request_actions_request" } & EmptyPayload | { "type": "get_grpc_request_actions_response" } & GetGrpcRequestActionsResponse | { "type": "call_grpc_request_action_request" } & CallGrpcRequestActionRequest | { "type": "get_template_function_summary_request" } & EmptyPayload | { "type": "get_template_function_summary_response" } & GetTemplateFunctionSummaryResponse | { "type": "get_template_function_config_request" } & GetTemplateFunctionConfigRequest | { "type": "get_template_function_config_response" } & GetTemplateFunctionConfigResponse | { "type": "call_template_function_request" } & CallTemplateFunctionRequest | { "type": "call_template_function_response" } & CallTemplateFunctionResponse | { "type": "get_http_authentication_summary_request" } & EmptyPayload | { "type": "get_http_authentication_summary_response" } & GetHttpAuthenticationSummaryResponse | { "type": "get_http_authentication_config_request" } & GetHttpAuthenticationConfigRequest | { "type": "get_http_authentication_config_response" } & GetHttpAuthenticationConfigResponse | { "type": "call_http_authentication_request" } & CallHttpAuthenticationRequest | { "type": "call_http_authentication_response" } & CallHttpAuthenticationResponse | { "type": "call_http_authentication_action_request" } & CallHttpAuthenticationActionRequest | { "type": "call_http_authentication_action_response" } & EmptyPayload | { "type": "copy_text_request" } & CopyTextRequest | { "type": "copy_text_response" } & EmptyPayload | { "type": "render_http_request_request" } & RenderHttpRequestRequest | { "type": "render_http_request_response" } & RenderHttpRequestResponse | { "type": "render_grpc_request_request" } & RenderGrpcRequestRequest | { "type": "render_grpc_request_response" } & RenderGrpcRequestResponse | { "type": "template_render_request" } & TemplateRenderRequest | { "type": "template_render_response" } & TemplateRenderResponse | { "type": "get_key_value_request" } & GetKeyValueRequest | { "type": "get_key_value_response" } & GetKeyValueResponse | { "type": "set_key_value_request" } & SetKeyValueRequest | { "type": "set_key_value_response" } & SetKeyValueResponse | { "type": "delete_key_value_request" } & DeleteKeyValueRequest | { "type": "delete_key_value_response" } & DeleteKeyValueResponse | { "type": "open_window_request" } & OpenWindowRequest | { "type": "window_navigate_event" } & WindowNavigateEvent | { "type": "window_close_event" } | { "type": "close_window_request" } & CloseWindowRequest | { "type": "open_external_url_request" } & OpenExternalUrlRequest | { "type": "open_external_url_response" } & EmptyPayload | { "type": "show_toast_request" } & ShowToastRequest | { "type": "show_toast_response" } & EmptyPayload | { "type": "prompt_text_request" } & PromptTextRequest | { "type": "prompt_text_response" } & PromptTextResponse | { "type": "prompt_form_request" } & PromptFormRequest | { "type": "prompt_form_response" } & PromptFormResponse | { "type": "window_info_request" } & WindowInfoRequest | { "type": "window_info_response" } & WindowInfoResponse | { "type": "list_open_workspaces_request" } & ListOpenWorkspacesRequest | { "type": "list_open_workspaces_response" } & ListOpenWorkspacesResponse | { "type": "get_http_request_by_id_request" } & GetHttpRequestByIdRequest | { "type": "get_http_request_by_id_response" } & GetHttpRequestByIdResponse | { "type": "find_http_responses_request" } & FindHttpResponsesRequest | { "type": "find_http_responses_response" } & FindHttpResponsesResponse | { "type": "get_http_response_body_info_request" } & GetHttpResponseBodyInfoRequest | { "type": "get_http_response_body_info_response" } & GetHttpResponseBodyInfoResponse | { "type": "read_http_response_body_chunk_request" } & ReadHttpResponseBodyChunkRequest | { "type": "read_http_response_body_chunk_response" } & ReadHttpResponseBodyChunkResponse | { "type": "list_http_requests_request" } & ListHttpRequestsRequest | { "type": "list_http_requests_response" } & ListHttpRequestsResponse | { "type": "list_folders_request" } & ListFoldersRequest | { "type": "list_folders_response" } & ListFoldersResponse | { "type": "upsert_model_request" } & UpsertModelRequest | { "type": "upsert_model_response" } & UpsertModelResponse | { "type": "delete_model_request" } & DeleteModelRequest | { "type": "delete_model_response" } & DeleteModelResponse | { "type": "get_themes_request" } & GetThemesRequest | { "type": "get_themes_response" } & GetThemesResponse | { "type": "empty_response" } & EmptyPayload | { "type": "error_response" } & ErrorResponse; export type JsonPrimitive = string | number | boolean | null; @@ -502,6 +522,27 @@ required?: boolean, }; export type PromptTextResponse = { value: string | null, }; +/** + * Pull one window of a response body. + * + * Reads are idempotent: the bytes live in durable storage, so the same window + * can be asked for as many times as the plugin likes. + */ +export type ReadHttpResponseBodyChunkRequest = { responseId: string, offset: number, length: number, }; + +export type ReadHttpResponseBodyChunkResponse = { +/** + * Base64, because the desktop transport is a WebSocket that only sends + * text frames today. A host that can carry binary sends the bytes as they + * are and fills this in from them. + */ +data: string, +/** + * Bytes decoded from `data`. Short of the requested length means the body + * ended here. + */ +length: number, }; + export type ReloadResponse = { silent: boolean, }; export type RenderGrpcRequestRequest = { grpcRequest: GrpcRequest, purpose: RenderPurpose, }; diff --git a/packages/plugin-runtime-types/src/plugins/Context.ts b/packages/plugin-runtime-types/src/plugins/Context.ts index 0a1afa49..7de049e2 100644 --- a/packages/plugin-runtime-types/src/plugins/Context.ts +++ b/packages/plugin-runtime-types/src/plugins/Context.ts @@ -6,6 +6,7 @@ import type { GetCookieValueResponse, GetHttpRequestByIdRequest, GetHttpRequestByIdResponse, + GetHttpResponseBodyInfoRequest, JsonPrimitive, ListCookieNamesResponse, ListFoldersRequest, @@ -65,6 +66,56 @@ type DynamicPromptFormRequest = Omit & { export type WorkspaceHandle = Pick; +export interface ReadHttpResponseBodyOptions { + /** + * Refuse to buffer a body larger than this, in bytes. Defaults to 32 MiB. + * Pass `Infinity` to read whatever is there. + */ + maxBytes?: number; + + /** Bytes to pull from the host at a time. Defaults to 1 MiB. */ + chunkSize?: number; +} + +/** + * A response body, read back from wherever the host stored it. + * + * The accessors are named after `fetch`'s, but unlike `fetch` the body is not + * used up by reading it: these bytes are in durable storage, so every accessor + * can be called as many times as you like, in any order. + */ +export interface HttpResponseBody { + /** The response these bytes belong to. */ + readonly responseId: string; + + /** + * How many bytes are stored, which is not necessarily what the + * `Content-Length` header claimed. Zero when the response has no body. + */ + readonly contentLength: number; + + /** The response's `Content-Type` header, verbatim, or null if it had none. */ + readonly contentType: string | null; + + /** + * The body decoded to a string, using the charset from `contentType` and + * falling back to UTF-8. Throws if the body is over `maxBytes`. + */ + text(options?: ReadHttpResponseBodyOptions): Promise; + + /** `text()`, parsed as JSON. */ + json(options?: ReadHttpResponseBodyOptions): Promise; + + /** The raw bytes. Throws if the body is over `maxBytes`. */ + arrayBuffer(options?: ReadHttpResponseBodyOptions): Promise; + + /** + * The raw bytes, a chunk at a time, so a body of any size can be read + * without holding all of it at once. Not subject to `maxBytes`. + */ + chunks(options?: Pick): AsyncIterable; +} + export interface Context { clipboard: { copyText(text: string): Promise; @@ -129,6 +180,12 @@ export interface Context { }; httpResponse: { find(args: FindHttpResponsesRequest): Promise; + /** + * Read a response's body by id. Where the host keeps the bytes — files on + * a desktop, rows in a database, somewhere else later — is not something a + * plugin sees or should depend on. + */ + body(args: GetHttpResponseBodyInfoRequest): Promise; }; templates: { render(args: TemplateRenderRequest & { data: T }): Promise; diff --git a/packages/plugin-runtime-types/src/plugins/index.ts b/packages/plugin-runtime-types/src/plugins/index.ts index 86c2b67d..92e8391f 100644 --- a/packages/plugin-runtime-types/src/plugins/index.ts +++ b/packages/plugin-runtime-types/src/plugins/index.ts @@ -13,7 +13,12 @@ import type { WorkspaceActionPlugin } from "./WorkspaceActionPlugin"; export type { Context }; export type { DynamicAuthenticationArg } from "./AuthenticationPlugin"; -export type { CallPromptFormDynamicArgs, DynamicPromptFormArg } from "./Context"; +export type { + CallPromptFormDynamicArgs, + DynamicPromptFormArg, + HttpResponseBody, + ReadHttpResponseBodyOptions, +} from "./Context"; export type { DynamicTemplateFunctionArg } from "./TemplateFunctionPlugin"; export type { TemplateFunctionPlugin }; export type { FolderActionPlugin } from "./FolderActionPlugin"; diff --git a/packages/plugin-runtime/src/PluginInstance.ts b/packages/plugin-runtime/src/PluginInstance.ts index 031a6d0f..ea6ad752 100644 --- a/packages/plugin-runtime/src/PluginInstance.ts +++ b/packages/plugin-runtime/src/PluginInstance.ts @@ -21,6 +21,7 @@ import type { GetCookieValueRequest, GetCookieValueResponse, GetHttpRequestByIdResponse, + GetHttpResponseBodyInfoResponse, GetKeyValueResponse, GrpcRequestAction, HttpAuthenticationAction, @@ -37,6 +38,7 @@ import type { PluginContext, PromptFormResponse, PromptTextResponse, + ReadHttpResponseBodyChunkResponse, RenderGrpcRequestResponse, RenderHttpRequestResponse, SendHttpRequestResponse, @@ -49,6 +51,7 @@ import type { import { applyDynamicFormInput } from "./common"; import { EventChannel } from "./EventChannel"; import { migrateTemplateFunctionSelectOptions } from "./migrations"; +import { createResponseBody, decodeBase64Chunk } from "./responseBody"; export interface PluginWorkerData { bootRequest: BootRequest; @@ -555,16 +558,24 @@ export class PluginInstance { #sendForReply>( context: PluginContext, payload: InternalEventPayload, + // Off by default because a reply-shaped object with none of the expected + // fields is what every existing caller already copes with; turning it on + // for a new call is how that stops spreading. + { throwOnError = false }: { throwOnError?: boolean } = {}, ): Promise { // 1. Build event to send const eventToSend = this.#buildEventToSend(context, payload, null); // 2. Spawn listener in background - const promise = new Promise((resolve) => { + const promise = new Promise((resolve, reject) => { const cb = (event: InternalEvent) => { if (event.replyId === eventToSend.id) { this.#appToPluginEvents.unlisten(cb); // Unlisten, now that we're done const { type: _, ...payload } = event.payload; + if (throwOnError && event.payload.type === "error_response") { + reject(new Error(String((payload as { error?: string }).error ?? "Unknown error"))); + return; + } resolve(payload as T); } }; @@ -751,6 +762,29 @@ export class PluginInstance { ); return httpResponses; }, + body: async ({ responseId }) => { + const info = await this.#sendForReply( + context, + { type: "get_http_response_body_info_request", responseId }, + { throwOnError: true }, + ); + + return createResponseBody( + { + responseId, + contentLength: info.contentLength, + contentType: info.contentType ?? null, + }, + async (offset, length) => { + const chunk = await this.#sendForReply( + context, + { type: "read_http_response_body_chunk_request", responseId, offset, length }, + { throwOnError: true }, + ); + return decodeBase64Chunk(chunk.data); + }, + ); + }, }, grpcRequest: { render: async (args) => { diff --git a/packages/plugin-runtime/src/responseBody.ts b/packages/plugin-runtime/src/responseBody.ts new file mode 100644 index 00000000..dc46b8e9 --- /dev/null +++ b/packages/plugin-runtime/src/responseBody.ts @@ -0,0 +1,151 @@ +import type { HttpResponseBody, ReadHttpResponseBodyOptions } from "@yaakapp/api"; + +/** Bytes pulled from the host per round trip, when the caller doesn't say. */ +const DEFAULT_CHUNK_SIZE = 1024 * 1024; + +/** + * The most a plugin buffers by default. + * + * Reading a body used to be unbounded, so any ceiling is an improvement; this + * one is set well above what an API returns and well below what makes the + * plugin runtime fall over. `chunks()` has no ceiling, and any caller that + * really wants the whole thing can raise `maxBytes`. + */ +const DEFAULT_MAX_BYTES = 32 * 1024 * 1024; + +/** Fetch one window of body bytes from the host. */ +export type ReadResponseBodyChunk = (offset: number, length: number) => Promise; + +export interface ResponseBodyInfo { + responseId: string; + contentLength: number; + contentType: string | null; +} + +export function createResponseBody( + info: ResponseBodyInfo, + readChunk: ReadResponseBodyChunk, +): HttpResponseBody { + const { responseId, contentLength, contentType } = info; + + async function* chunks( + options?: Pick, + ): AsyncIterable { + const chunkSize = Math.max(1, Math.floor(options?.chunkSize ?? DEFAULT_CHUNK_SIZE)); + let offset = 0; + // Bounded by the length the host reported, but a short read still ends it: + // the body may have been rewritten between the two calls. + while (offset < contentLength) { + const chunk = await readChunk(offset, Math.min(chunkSize, contentLength - offset)); + if (chunk.byteLength === 0) return; + yield chunk; + offset += chunk.byteLength; + } + } + + async function readAll(accessor: string, options?: ReadHttpResponseBodyOptions) { + const maxBytes = options?.maxBytes ?? DEFAULT_MAX_BYTES; + refuseIfTooBig(accessor, contentLength, maxBytes); + + const parts: Uint8Array[] = []; + let total = 0; + for await (const chunk of chunks(options)) { + total += chunk.byteLength; + // The size the host reported is a claim about a moment ago, so check the + // bytes actually arriving too. + refuseIfTooBig(accessor, total, maxBytes); + parts.push(chunk); + } + + const bytes = new Uint8Array(total); + let offset = 0; + for (const part of parts) { + bytes.set(part, offset); + offset += part.byteLength; + } + return bytes; + } + + return { + responseId, + contentLength, + contentType, + chunks, + async arrayBuffer(options) { + const bytes = await readAll("arrayBuffer", options); + return bytes.buffer as ArrayBuffer; + }, + async text(options) { + return decodeBody(await readAll("text", options), contentType); + }, + async json(options?: ReadHttpResponseBodyOptions) { + return JSON.parse(decodeBody(await readAll("json", options), contentType)) as T; + }, + }; +} + +function refuseIfTooBig(accessor: string, bytes: number, maxBytes: number) { + if (bytes <= maxBytes) return; + throw new Error( + `Response body is ${formatBytes(bytes)}, over the ${formatBytes(maxBytes)} limit for ` + + `${accessor}(). Read it with chunks() instead, or pass a larger maxBytes.`, + ); +} + +/** + * Decode using the charset the response declared. + * + * Assuming UTF-8 mangles every response that isn't, and the header is right + * there. An unknown label is the one case worth guessing on, since the + * alternative is refusing to read a body we can very likely still read. + */ +function decodeBody(bytes: Uint8Array, contentType: string | null): string { + const charset = parseCharset(contentType); + if (charset != null) { + try { + return new TextDecoder(charset).decode(bytes); + } catch { + // Not a label this runtime knows. + } + } + // TextDecoder drops a leading BOM on its own. + return new TextDecoder("utf-8").decode(bytes); +} + +function parseCharset(contentType: string | null): string | null { + const match = contentType?.match(/;\s*charset\s*=\s*"?([^";]+)"?/i); + return match?.[1]?.trim() || null; +} + +function formatBytes(bytes: number): string { + if (bytes === Infinity) return "unlimited"; + if (bytes < 1024) return `${bytes} B`; + const units = ["KB", "MB", "GB"]; + let value = bytes / 1024; + let unit = 0; + while (value >= 1024 && unit < units.length - 1) { + value /= 1024; + unit++; + } + return `${value.toFixed(1)} ${units[unit]}`; +} + +/** + * Decode a chunk that arrived as base64. + * + * The desktop transport is a WebSocket carrying JSON text frames, so bytes + * have to be spelled out. A host that can pass an ArrayBuffer along skips this. + */ +export function decodeBase64Chunk(data: string): Uint8Array { + if (typeof Buffer !== "undefined") { + const buf = Buffer.from(data, "base64"); + return new Uint8Array(buf.buffer, buf.byteOffset, buf.byteLength); + } + + const binary = atob(data); + const bytes = new Uint8Array(binary.length); + for (let i = 0; i < binary.length; i++) { + bytes[i] = binary.charCodeAt(i); + } + return bytes; +} diff --git a/packages/plugin-runtime/tests/responseBody.test.ts b/packages/plugin-runtime/tests/responseBody.test.ts new file mode 100644 index 00000000..4e0effa8 --- /dev/null +++ b/packages/plugin-runtime/tests/responseBody.test.ts @@ -0,0 +1,122 @@ +import { describe, expect, test } from "vite-plus/test"; +import { createResponseBody, decodeBase64Chunk } from "../src/responseBody"; + +/** A store of bytes that records every window it was asked for. */ +function fakeBody(bytes: Uint8Array, contentType: string | null) { + const reads: Array<[number, number]> = []; + const body = createResponseBody( + { responseId: "rs_test", contentLength: bytes.byteLength, contentType }, + async (offset, length) => { + reads.push([offset, length]); + return bytes.slice(offset, offset + length); + }, + ); + return { body, reads }; +} + +function utf8(text: string) { + return new TextEncoder().encode(text); +} + +describe("response body", () => { + test("pulls a body in chunks and reassembles it", async () => { + const { body, reads } = fakeBody(utf8("abcdefghij"), "text/plain"); + + expect(await body.text({ chunkSize: 4 })).toEqual("abcdefghij"); + expect(reads).toEqual([ + [0, 4], + [4, 4], + [8, 2], + ]); + }); + + test("can be read more than once, unlike fetch", async () => { + const { body } = fakeBody(utf8('{"a":1}'), "application/json"); + + expect(await body.text()).toEqual('{"a":1}'); + expect(await body.json()).toEqual({ a: 1 }); + expect(new Uint8Array(await body.arrayBuffer())).toEqual(utf8('{"a":1}')); + }); + + test("decodes using the charset the response declared", async () => { + // "café naïve" as Latin-1, which is mojibake if read as UTF-8. + const latin1 = new Uint8Array([0x63, 0x61, 0x66, 0xe9, 0x20, 0x6e, 0x61, 0xef, 0x76, 0x65]); + + const declared = fakeBody(latin1, "text/plain; charset=iso-8859-1"); + expect(await declared.body.text()).toEqual("café naïve"); + + const undeclared = fakeBody(latin1, "text/plain"); + expect(await undeclared.body.text()).not.toEqual("café naïve"); + }); + + test("falls back to UTF-8 for a charset the runtime doesn't know", async () => { + const { body } = fakeBody(utf8("hello"), "text/plain; charset=not-a-real-charset"); + expect(await body.text()).toEqual("hello"); + }); + + test("drops a UTF-8 BOM", async () => { + const withBom = new Uint8Array([0xef, 0xbb, 0xbf, ...utf8('{"a":1}')]); + const { body } = fakeBody(withBom, "application/json"); + + expect(await body.text()).toEqual('{"a":1}'); + expect(await body.json()).toEqual({ a: 1 }); + }); + + test("refuses to buffer past maxBytes, and says what to do instead", async () => { + const { body } = fakeBody(utf8("x".repeat(100)), "text/plain"); + + await expect(body.text({ maxBytes: 50 })).rejects.toThrow(/chunks\(\)/); + await expect(body.json({ maxBytes: 50 })).rejects.toThrow(/over the/); + await expect(body.arrayBuffer({ maxBytes: 50 })).rejects.toThrow(/arrayBuffer\(\)/); + + // The ceiling is the caller's to raise. + expect(await body.text({ maxBytes: 100 })).toHaveLength(100); + }); + + test("streams past maxBytes through chunks()", async () => { + const { body } = fakeBody(utf8("x".repeat(100)), "text/plain"); + + let total = 0; + for await (const chunk of body.chunks({ chunkSize: 10 })) { + total += chunk.byteLength; + } + expect(total).toEqual(100); + }); + + test("stops early when the host runs out of bytes sooner than it claimed", async () => { + // contentLength says 100; the store only ever hands back 10. + const body = createResponseBody( + { responseId: "rs_test", contentLength: 100, contentType: "text/plain" }, + async (offset) => (offset === 0 ? utf8("0123456789") : new Uint8Array()), + ); + + expect(await body.text()).toEqual("0123456789"); + }); + + test("a response with no body reads as empty", async () => { + const { body, reads } = fakeBody(new Uint8Array(), "application/json"); + + expect(body.contentLength).toEqual(0); + expect(await body.text()).toEqual(""); + expect(reads).toEqual([]); + }); + + test("keeps binary bytes intact", async () => { + const png = new Uint8Array([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 0x00, 0xff]); + const { body } = fakeBody(png, "image/png"); + + expect(new Uint8Array(await body.arrayBuffer({ chunkSize: 3 }))).toEqual(png); + }); +}); + +describe("decodeBase64Chunk", () => { + test("round-trips arbitrary bytes", () => { + const bytes = new Uint8Array([0, 1, 127, 128, 254, 255]); + const base64 = Buffer.from(bytes).toString("base64"); + expect(decodeBase64Chunk(base64)).toEqual(bytes); + }); + + test("decodes an empty chunk", () => { + expect(decodeBase64Chunk("")).toEqual(new Uint8Array()); + }); +}); diff --git a/plugins/template-function-response/src/index.ts b/plugins/template-function-response/src/index.ts index 9de85ffb..541482f1 100644 --- a/plugins/template-function-response/src/index.ts +++ b/plugins/template-function-response/src/index.ts @@ -1,4 +1,3 @@ -import { readFileSync } from "node:fs"; import type { CallTemplateFunctionArgs, Context, @@ -196,17 +195,8 @@ export const plugin: PluginDefinition = { }); if (response == null) return null; - if (response.bodyPath == null) { - return null; - } - - const BOM = "\ufeff"; - let body: string; - try { - body = readFileSync(response.bodyPath, "utf-8").replace(BOM, ""); - } catch { - return null; - } + const body = await readResponseBody(ctx, response); + if (body == null) return null; try { const result: JSONPathResult = @@ -261,23 +251,29 @@ export const plugin: PluginDefinition = { }); if (response == null) return null; - if (response.bodyPath == null) { - return null; - } - - let body: string; - try { - body = readFileSync(response.bodyPath, "utf-8"); - } catch { - return null; - } - - return body; + return await readResponseBody(ctx, response); }, }, ], }; +/** + * The response's body as text, or null when there is nothing to read. + * + * The host is asked for it by response id, so this works wherever the bytes + * happen to live. A body over the runtime's size limit throws rather than + * coming back empty, since a template silently rendering to nothing is worse + * than one that says why. + */ +async function readResponseBody(ctx: Context, response: HttpResponse): Promise { + if (!response.id) return null; + + const body = await ctx.httpResponse.body({ responseId: response.id }); + if (body.contentLength === 0) return null; + + return await body.text(); +} + async function getResponse( ctx: Context, {