mirror of
https://github.com/mountain-loop/yaak.git
synced 2026-08-25 21:04:04 +02:00
Make the HTTP send path runnable without a database (#545)
This commit is contained in:
@@ -593,7 +593,7 @@ impl UpsertModelInfo for WorkspaceMeta {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, TS)]
|
#[derive(Debug, Clone, Serialize, Deserialize, TS, PartialEq)]
|
||||||
#[ts(export, export_to = "gen_models.ts")]
|
#[ts(export, export_to = "gen_models.ts")]
|
||||||
pub enum CookieDomain {
|
pub enum CookieDomain {
|
||||||
HostOnly(String),
|
HostOnly(String),
|
||||||
@@ -602,14 +602,14 @@ pub enum CookieDomain {
|
|||||||
Empty,
|
Empty,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, TS)]
|
#[derive(Debug, Clone, Serialize, Deserialize, TS, PartialEq)]
|
||||||
#[ts(export, export_to = "gen_models.ts")]
|
#[ts(export, export_to = "gen_models.ts")]
|
||||||
pub enum CookieExpires {
|
pub enum CookieExpires {
|
||||||
AtUtc(String),
|
AtUtc(String),
|
||||||
SessionEnd,
|
SessionEnd,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize, TS)]
|
#[derive(Debug, Clone, Serialize, Deserialize, TS, PartialEq)]
|
||||||
#[ts(export, export_to = "gen_models.ts")]
|
#[ts(export, export_to = "gen_models.ts")]
|
||||||
pub enum CookieSameSite {
|
pub enum CookieSameSite {
|
||||||
Strict,
|
Strict,
|
||||||
@@ -617,7 +617,7 @@ pub enum CookieSameSite {
|
|||||||
None,
|
None,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, TS)]
|
#[derive(Debug, Clone, Serialize, TS, PartialEq)]
|
||||||
#[serde(rename_all = "camelCase")]
|
#[serde(rename_all = "camelCase")]
|
||||||
#[ts(export, export_to = "gen_models.ts")]
|
#[ts(export, export_to = "gen_models.ts")]
|
||||||
pub struct Cookie {
|
pub struct Cookie {
|
||||||
|
|||||||
@@ -21,3 +21,4 @@ yaak-tls = { workspace = true }
|
|||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
tempfile = "3"
|
tempfile = "3"
|
||||||
|
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }
|
||||||
|
|||||||
+491
-186
@@ -24,9 +24,9 @@ use yaak_http::types::{
|
|||||||
};
|
};
|
||||||
use yaak_models::blob_manager::{BlobManager, BodyChunk};
|
use yaak_models::blob_manager::{BlobManager, BodyChunk};
|
||||||
use yaak_models::models::{
|
use yaak_models::models::{
|
||||||
ClientCertificate, CookieJar, DnsOverride, Environment, HttpRequest, HttpResponse,
|
ClientCertificate, Cookie, CookieJar, DnsOverride, Environment, HttpRequest, HttpResponse,
|
||||||
HttpResponseEvent, HttpResponseHeader, HttpResponseState, ProxySetting, ProxySettingAuth,
|
HttpResponseEvent, HttpResponseHeader, HttpResponseState, ProxySetting, ProxySettingAuth,
|
||||||
ResolvedSetting,
|
ResolvedHttpRequestSettings, ResolvedSetting,
|
||||||
};
|
};
|
||||||
use yaak_models::query_manager::QueryManager;
|
use yaak_models::query_manager::QueryManager;
|
||||||
use yaak_models::util::{UpdateSource, generate_prefixed_id};
|
use yaak_models::util::{UpdateSource, generate_prefixed_id};
|
||||||
@@ -170,8 +170,7 @@ impl PrepareSendableRequest for PluginPrepareSendableRequest {
|
|||||||
struct ConnectionManagerSendRequestExecutor<'a> {
|
struct ConnectionManagerSendRequestExecutor<'a> {
|
||||||
connection_manager: &'a HttpConnectionManager,
|
connection_manager: &'a HttpConnectionManager,
|
||||||
plugin_context_id: String,
|
plugin_context_id: String,
|
||||||
query_manager: QueryManager,
|
runtime_config: HttpSendRuntimeConfig,
|
||||||
request: HttpRequest,
|
|
||||||
cancelled_rx: Option<watch::Receiver<bool>>,
|
cancelled_rx: Option<watch::Receiver<bool>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -183,18 +182,17 @@ impl SendRequestExecutor for ConnectionManagerSendRequestExecutor<'_> {
|
|||||||
event_tx: mpsc::Sender<SenderHttpResponseEvent>,
|
event_tx: mpsc::Sender<SenderHttpResponseEvent>,
|
||||||
cookie_behavior: CookieBehavior,
|
cookie_behavior: CookieBehavior,
|
||||||
) -> yaak_http::error::Result<yaak_http::sender::HttpResponse> {
|
) -> yaak_http::error::Result<yaak_http::sender::HttpResponse> {
|
||||||
let runtime_config = resolve_http_send_runtime_config(&self.query_manager, &self.request)
|
let runtime_config = &self.runtime_config;
|
||||||
.map_err(|e| yaak_http::error::Error::RequestError(e.to_string()))?;
|
|
||||||
let client_certificate =
|
let client_certificate =
|
||||||
find_client_certificate(&sendable_request.url, &runtime_config.client_certificates);
|
find_client_certificate(&sendable_request.url, &runtime_config.client_certificates);
|
||||||
let cached_client = self
|
let cached_client = self
|
||||||
.connection_manager
|
.connection_manager
|
||||||
.get_client(&HttpConnectionOptions {
|
.get_client(&HttpConnectionOptions {
|
||||||
id: self.plugin_context_id.clone(),
|
id: self.plugin_context_id.clone(),
|
||||||
validate_certificates: runtime_config.validate_certificates,
|
validate_certificates: runtime_config.settings.validate_certificates.value,
|
||||||
proxy: runtime_config.proxy,
|
proxy: runtime_config.proxy.clone(),
|
||||||
client_certificate,
|
client_certificate,
|
||||||
dns_overrides: runtime_config.dns_overrides,
|
dns_overrides: runtime_config.dns_overrides.clone(),
|
||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
@@ -238,20 +236,69 @@ pub struct SendHttpRequestByIdParams<'a, T: TemplateCallback> {
|
|||||||
pub executor: &'a dyn SendRequestExecutor,
|
pub executor: &'a dyn SendRequestExecutor,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct SendHttpRequestParams<'a, T: TemplateCallback> {
|
/// An [`HttpRequest`] carrying the authentication and headers it inherits from its folder or
|
||||||
|
/// workspace, alongside the id of the model that authentication came from.
|
||||||
|
///
|
||||||
|
/// Sending requires inheritance to be applied first, and this type is the only way to say it has
|
||||||
|
/// been. There is no public constructor beyond [`resolve_inherited_request`], which resolves it
|
||||||
|
/// against a database, and [`ResolvedHttpRequest::assume_resolved`], which a caller without a
|
||||||
|
/// database must name explicitly — so skipping inheritance is a deliberate, greppable act rather
|
||||||
|
/// than something a new call site forgets.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct ResolvedHttpRequest {
|
||||||
|
request: HttpRequest,
|
||||||
|
auth_context_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ResolvedHttpRequest {
|
||||||
|
/// Declare a request already resolved, for callers with no database to resolve against. The
|
||||||
|
/// caller owns the promise that inherited authentication and headers are applied.
|
||||||
|
pub fn assume_resolved(request: HttpRequest, auth_context_id: String) -> Self {
|
||||||
|
Self { request, auth_context_id }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn request(&self) -> &HttpRequest {
|
||||||
|
&self.request
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn auth_context_id(&self) -> &str {
|
||||||
|
&self.auth_context_id
|
||||||
|
}
|
||||||
|
|
||||||
|
fn into_parts(self) -> (HttpRequest, String) {
|
||||||
|
(self.request, self.auth_context_id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Everything a send needs that would otherwise be read from the database.
|
||||||
|
///
|
||||||
|
/// Callers backed by a database build this with [`resolve_send_inputs`]. A stateless caller
|
||||||
|
/// constructs it directly, which is what lets [`send_http_request`] run with no database at all.
|
||||||
|
pub struct HttpSendInputs {
|
||||||
|
pub request: ResolvedHttpRequest,
|
||||||
|
pub environment_chain: Vec<Environment>,
|
||||||
|
pub runtime_config: HttpSendRuntimeConfig,
|
||||||
|
/// Cookies the send starts with. The store is shared, so reading it back after the send
|
||||||
|
/// returns (or fails) yields the cookies the transaction collected.
|
||||||
|
pub cookie_store: Option<CookieStore>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Where a send writes its response. Without it, the send keeps everything in memory: no
|
||||||
|
/// response body file, no model writes, and no timeline rows.
|
||||||
|
pub struct ResponseStorage<'a> {
|
||||||
pub query_manager: &'a QueryManager,
|
pub query_manager: &'a QueryManager,
|
||||||
pub blob_manager: &'a BlobManager,
|
pub blob_manager: &'a BlobManager,
|
||||||
pub request: HttpRequest,
|
|
||||||
pub environment_id: Option<&'a str>,
|
|
||||||
pub template_callback: &'a T,
|
|
||||||
pub send_options: Option<SendableHttpRequestOptions>,
|
|
||||||
pub update_source: UpdateSource,
|
pub update_source: UpdateSource,
|
||||||
pub cookie_jar_id: Option<String>,
|
|
||||||
pub response_dir: &'a Path,
|
pub response_dir: &'a Path,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub struct SendHttpRequestParams<'a, T: TemplateCallback> {
|
||||||
|
pub inputs: HttpSendInputs,
|
||||||
|
pub template_callback: &'a T,
|
||||||
|
pub storage: Option<ResponseStorage<'a>>,
|
||||||
pub emit_events_to: Option<mpsc::Sender<SenderHttpResponseEvent>>,
|
pub emit_events_to: Option<mpsc::Sender<SenderHttpResponseEvent>>,
|
||||||
pub emit_response_body_chunks_to: Option<mpsc::UnboundedSender<Vec<u8>>>,
|
pub emit_response_body_chunks_to: Option<mpsc::UnboundedSender<Vec<u8>>>,
|
||||||
pub cancelled_rx: Option<watch::Receiver<bool>>,
|
pub cancelled_rx: Option<watch::Receiver<bool>>,
|
||||||
pub auth_context_id: Option<String>,
|
|
||||||
pub existing_response: Option<HttpResponse>,
|
pub existing_response: Option<HttpResponse>,
|
||||||
pub prepare_sendable_request: Option<&'a dyn PrepareSendableRequest>,
|
pub prepare_sendable_request: Option<&'a dyn PrepareSendableRequest>,
|
||||||
pub executor: &'a dyn SendRequestExecutor,
|
pub executor: &'a dyn SendRequestExecutor,
|
||||||
@@ -296,43 +343,84 @@ pub struct SendHttpRequestResult {
|
|||||||
pub rendered_request: HttpRequest,
|
pub rendered_request: HttpRequest,
|
||||||
pub response: HttpResponse,
|
pub response: HttpResponse,
|
||||||
pub response_body: Vec<u8>,
|
pub response_body: Vec<u8>,
|
||||||
|
/// The cookies held by the jar after the send, for callers that persist one.
|
||||||
|
pub cookies: Option<Vec<Cookie>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
pub struct HttpSendRuntimeConfig {
|
pub struct HttpSendRuntimeConfig {
|
||||||
pub send_options: SendableHttpRequestOptions,
|
pub settings: ResolvedHttpRequestSettings,
|
||||||
pub validate_certificates: bool,
|
|
||||||
pub proxy: HttpConnectionProxySetting,
|
pub proxy: HttpConnectionProxySetting,
|
||||||
pub dns_overrides: Vec<DnsOverride>,
|
pub dns_overrides: Vec<DnsOverride>,
|
||||||
pub client_certificates: Vec<ClientCertificate>,
|
pub client_certificates: Vec<ClientCertificate>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn resolve_http_send_runtime_config(
|
impl HttpSendRuntimeConfig {
|
||||||
query_manager: &QueryManager,
|
pub fn send_options(&self) -> SendableHttpRequestOptions {
|
||||||
request: &HttpRequest,
|
SendableHttpRequestOptions {
|
||||||
) -> Result<HttpSendRuntimeConfig> {
|
follow_redirects: self.settings.follow_redirects.value,
|
||||||
let db = query_manager.connect();
|
timeout: if self.settings.request_timeout.value > 0 {
|
||||||
let workspace =
|
Some(Duration::from_millis(
|
||||||
db.get_workspace(&request.workspace_id).map_err(SendHttpRequestError::LoadWorkspace)?;
|
self.settings.request_timeout.value.unsigned_abs() as u64
|
||||||
let resolved_settings = db
|
|
||||||
.resolve_settings_for_http_request(request)
|
|
||||||
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
|
||||||
let settings = db.get_settings();
|
|
||||||
|
|
||||||
Ok(HttpSendRuntimeConfig {
|
|
||||||
send_options: SendableHttpRequestOptions {
|
|
||||||
follow_redirects: resolved_settings.follow_redirects.value,
|
|
||||||
timeout: if resolved_settings.request_timeout.value > 0 {
|
|
||||||
Some(std::time::Duration::from_millis(
|
|
||||||
resolved_settings.request_timeout.value.unsigned_abs() as u64,
|
|
||||||
))
|
))
|
||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
},
|
},
|
||||||
},
|
}
|
||||||
validate_certificates: resolved_settings.validate_certificates.value,
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Resolve every database-backed input a send needs, in one pass.
|
||||||
|
pub fn resolve_send_inputs(
|
||||||
|
query_manager: &QueryManager,
|
||||||
|
request: &HttpRequest,
|
||||||
|
environment_id: Option<&str>,
|
||||||
|
cookies: Option<Vec<Cookie>>,
|
||||||
|
) -> Result<HttpSendInputs> {
|
||||||
|
let db = query_manager.connect();
|
||||||
|
|
||||||
|
let environment_chain = db
|
||||||
|
.resolve_environments(&request.workspace_id, request.folder_id.as_deref(), environment_id)
|
||||||
|
.map_err(SendHttpRequestError::ResolveEnvironments)?;
|
||||||
|
|
||||||
|
let resolved_request = resolve_inherited_request(query_manager, request)?;
|
||||||
|
|
||||||
|
let workspace =
|
||||||
|
db.get_workspace(&request.workspace_id).map_err(SendHttpRequestError::LoadWorkspace)?;
|
||||||
|
let settings = db.get_settings();
|
||||||
|
let resolved_settings = db
|
||||||
|
.resolve_settings_for_http_request(request)
|
||||||
|
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
||||||
|
|
||||||
|
Ok(HttpSendInputs {
|
||||||
|
request: resolved_request,
|
||||||
|
environment_chain,
|
||||||
|
runtime_config: HttpSendRuntimeConfig {
|
||||||
|
settings: resolved_settings,
|
||||||
proxy: proxy_setting_from_settings(settings.proxy),
|
proxy: proxy_setting_from_settings(settings.proxy),
|
||||||
dns_overrides: workspace.setting_dns_overrides,
|
dns_overrides: workspace.setting_dns_overrides,
|
||||||
client_certificates: settings.client_certificates,
|
client_certificates: settings.client_certificates,
|
||||||
|
},
|
||||||
|
cookie_store: cookies.map(CookieStore::from_cookies),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Apply the authentication and headers a request inherits from its folder or workspace.
|
||||||
|
pub fn resolve_inherited_request(
|
||||||
|
query_manager: &QueryManager,
|
||||||
|
request: &HttpRequest,
|
||||||
|
) -> Result<ResolvedHttpRequest> {
|
||||||
|
let db = query_manager.connect();
|
||||||
|
let (authentication_type, authentication, auth_context_id) = db
|
||||||
|
.resolve_auth_for_http_request(request)
|
||||||
|
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
||||||
|
let headers = db
|
||||||
|
.resolve_headers_for_http_request(request)
|
||||||
|
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
||||||
|
|
||||||
|
Ok(ResolvedHttpRequest {
|
||||||
|
request: HttpRequest { authentication_type, authentication, headers, ..request.clone() },
|
||||||
|
auth_context_id,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -368,6 +456,14 @@ pub async fn send_http_request_by_id_with_plugins(
|
|||||||
pub async fn send_http_request_with_plugins(
|
pub async fn send_http_request_with_plugins(
|
||||||
params: SendHttpRequestWithPluginsParams<'_>,
|
params: SendHttpRequestWithPluginsParams<'_>,
|
||||||
) -> Result<SendHttpRequestResult> {
|
) -> Result<SendHttpRequestResult> {
|
||||||
|
let mut cookie_jar = load_cookie_jar(params.query_manager, params.cookie_jar_id.as_deref())?;
|
||||||
|
let inputs = resolve_send_inputs(
|
||||||
|
params.query_manager,
|
||||||
|
¶ms.request,
|
||||||
|
params.environment_id,
|
||||||
|
cookie_jar.as_ref().map(|jar| jar.cookies.clone()),
|
||||||
|
)?;
|
||||||
|
|
||||||
let template_callback = PluginTemplateCallback::new(
|
let template_callback = PluginTemplateCallback::new(
|
||||||
params.plugin_manager.clone(),
|
params.plugin_manager.clone(),
|
||||||
params.encryption_manager.clone(),
|
params.encryption_manager.clone(),
|
||||||
@@ -382,30 +478,32 @@ pub async fn send_http_request_with_plugins(
|
|||||||
let executor = ConnectionManagerSendRequestExecutor {
|
let executor = ConnectionManagerSendRequestExecutor {
|
||||||
connection_manager: params.connection_manager,
|
connection_manager: params.connection_manager,
|
||||||
plugin_context_id: params.plugin_context.id.clone(),
|
plugin_context_id: params.plugin_context.id.clone(),
|
||||||
query_manager: params.query_manager.clone(),
|
runtime_config: inputs.runtime_config.clone(),
|
||||||
request: params.request.clone(),
|
|
||||||
cancelled_rx: params.cancelled_rx.clone(),
|
cancelled_rx: params.cancelled_rx.clone(),
|
||||||
};
|
};
|
||||||
|
let cookie_store = inputs.cookie_store.clone();
|
||||||
|
|
||||||
send_http_request(SendHttpRequestParams {
|
let result = send_http_request(SendHttpRequestParams {
|
||||||
|
inputs,
|
||||||
|
template_callback: &template_callback,
|
||||||
|
storage: Some(ResponseStorage {
|
||||||
query_manager: params.query_manager,
|
query_manager: params.query_manager,
|
||||||
blob_manager: params.blob_manager,
|
blob_manager: params.blob_manager,
|
||||||
request: params.request,
|
|
||||||
environment_id: params.environment_id,
|
|
||||||
template_callback: &template_callback,
|
|
||||||
send_options: None,
|
|
||||||
update_source: params.update_source,
|
update_source: params.update_source,
|
||||||
cookie_jar_id: params.cookie_jar_id,
|
|
||||||
response_dir: params.response_dir,
|
response_dir: params.response_dir,
|
||||||
|
}),
|
||||||
emit_events_to: params.emit_events_to,
|
emit_events_to: params.emit_events_to,
|
||||||
emit_response_body_chunks_to: params.emit_response_body_chunks_to,
|
emit_response_body_chunks_to: params.emit_response_body_chunks_to,
|
||||||
cancelled_rx: params.cancelled_rx,
|
cancelled_rx: params.cancelled_rx,
|
||||||
auth_context_id: None,
|
|
||||||
existing_response: params.existing_response,
|
existing_response: params.existing_response,
|
||||||
prepare_sendable_request: Some(&auth_hook),
|
prepare_sendable_request: Some(&auth_hook),
|
||||||
executor: &executor,
|
executor: &executor,
|
||||||
})
|
})
|
||||||
.await
|
.await;
|
||||||
|
|
||||||
|
persist_cookies_after_send(params.query_manager, cookie_jar.as_mut(), cookie_store.as_ref())?;
|
||||||
|
|
||||||
|
result
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send_http_request_by_id<T: TemplateCallback>(
|
pub async fn send_http_request_by_id<T: TemplateCallback>(
|
||||||
@@ -416,50 +514,46 @@ pub async fn send_http_request_by_id<T: TemplateCallback>(
|
|||||||
.connect()
|
.connect()
|
||||||
.get_http_request(params.request_id)
|
.get_http_request(params.request_id)
|
||||||
.map_err(SendHttpRequestError::LoadRequest)?;
|
.map_err(SendHttpRequestError::LoadRequest)?;
|
||||||
let (request, auth_context_id) = resolve_inherited_request(params.query_manager, &request)?;
|
let mut cookie_jar = load_cookie_jar(params.query_manager, params.cookie_jar_id.as_deref())?;
|
||||||
|
let inputs = resolve_send_inputs(
|
||||||
|
params.query_manager,
|
||||||
|
&request,
|
||||||
|
params.environment_id,
|
||||||
|
cookie_jar.as_ref().map(|jar| jar.cookies.clone()),
|
||||||
|
)?;
|
||||||
|
let cookie_store = inputs.cookie_store.clone();
|
||||||
|
|
||||||
send_http_request(SendHttpRequestParams {
|
let result = send_http_request(SendHttpRequestParams {
|
||||||
|
inputs,
|
||||||
|
template_callback: params.template_callback,
|
||||||
|
storage: Some(ResponseStorage {
|
||||||
query_manager: params.query_manager,
|
query_manager: params.query_manager,
|
||||||
blob_manager: params.blob_manager,
|
blob_manager: params.blob_manager,
|
||||||
request,
|
|
||||||
environment_id: params.environment_id,
|
|
||||||
template_callback: params.template_callback,
|
|
||||||
send_options: None,
|
|
||||||
update_source: params.update_source,
|
update_source: params.update_source,
|
||||||
cookie_jar_id: params.cookie_jar_id,
|
|
||||||
response_dir: params.response_dir,
|
response_dir: params.response_dir,
|
||||||
|
}),
|
||||||
emit_events_to: params.emit_events_to,
|
emit_events_to: params.emit_events_to,
|
||||||
emit_response_body_chunks_to: params.emit_response_body_chunks_to,
|
emit_response_body_chunks_to: params.emit_response_body_chunks_to,
|
||||||
cancelled_rx: params.cancelled_rx,
|
cancelled_rx: params.cancelled_rx,
|
||||||
existing_response: None,
|
existing_response: None,
|
||||||
prepare_sendable_request: params.prepare_sendable_request,
|
prepare_sendable_request: params.prepare_sendable_request,
|
||||||
executor: params.executor,
|
executor: params.executor,
|
||||||
auth_context_id: Some(auth_context_id),
|
|
||||||
})
|
})
|
||||||
.await
|
.await;
|
||||||
|
|
||||||
|
persist_cookies_after_send(params.query_manager, cookie_jar.as_mut(), cookie_store.as_ref())?;
|
||||||
|
|
||||||
|
result
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send_http_request<T: TemplateCallback>(
|
pub async fn send_http_request<T: TemplateCallback>(
|
||||||
params: SendHttpRequestParams<'_, T>,
|
params: SendHttpRequestParams<'_, T>,
|
||||||
) -> Result<SendHttpRequestResult> {
|
) -> Result<SendHttpRequestResult> {
|
||||||
let environment_chain =
|
let HttpSendInputs { request, environment_chain, runtime_config, cookie_store } = params.inputs;
|
||||||
resolve_environment_chain(params.query_manager, ¶ms.request, params.environment_id)?;
|
let (request, auth_context_id) = request.into_parts();
|
||||||
let (resolved_request, auth_context_id) =
|
let storage = params.storage;
|
||||||
if let Some(auth_context_id) = params.auth_context_id.clone() {
|
let send_options = runtime_config.send_options();
|
||||||
(params.request.clone(), auth_context_id)
|
let resolved_settings = &runtime_config.settings;
|
||||||
} else {
|
|
||||||
resolve_inherited_request(params.query_manager, ¶ms.request)?
|
|
||||||
};
|
|
||||||
let runtime_config = resolve_http_send_runtime_config(params.query_manager, ¶ms.request)?;
|
|
||||||
let send_options = params.send_options.unwrap_or(runtime_config.send_options);
|
|
||||||
let resolved_settings = params
|
|
||||||
.query_manager
|
|
||||||
.connect()
|
|
||||||
.resolve_settings_for_http_request(¶ms.request)
|
|
||||||
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
|
||||||
let mut cookie_jar = load_cookie_jar(params.query_manager, params.cookie_jar_id.as_deref())?;
|
|
||||||
let cookie_store =
|
|
||||||
cookie_jar.as_ref().map(|jar| CookieStore::from_cookies(jar.cookies.clone()));
|
|
||||||
let cookie_behavior = CookieBehavior {
|
let cookie_behavior = CookieBehavior {
|
||||||
store: cookie_store,
|
store: cookie_store,
|
||||||
send_cookies: resolved_settings.send_cookies.value,
|
send_cookies: resolved_settings.send_cookies.value,
|
||||||
@@ -467,7 +561,7 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
};
|
};
|
||||||
|
|
||||||
let rendered_request = render_http_request(
|
let rendered_request = render_http_request(
|
||||||
&resolved_request,
|
&request,
|
||||||
environment_chain,
|
environment_chain,
|
||||||
params.template_callback,
|
params.template_callback,
|
||||||
&RenderOptions::throw(),
|
&RenderOptions::throw(),
|
||||||
@@ -491,8 +585,8 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
Some(SendableBody::Stream { .. }) | None => None,
|
Some(SendableBody::Stream { .. }) | None => None,
|
||||||
};
|
};
|
||||||
let mut response = params.existing_response.unwrap_or_default();
|
let mut response = params.existing_response.unwrap_or_default();
|
||||||
response.request_id = params.request.id.clone();
|
response.request_id = request.id.clone();
|
||||||
response.workspace_id = params.request.workspace_id.clone();
|
response.workspace_id = request.workspace_id.clone();
|
||||||
response.request_content_length = request_content_length;
|
response.request_content_length = request_content_length;
|
||||||
response.request_headers = sendable_request
|
response.request_headers = sendable_request
|
||||||
.headers
|
.headers
|
||||||
@@ -513,12 +607,15 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
response.elapsed = 0;
|
response.elapsed = 0;
|
||||||
response.elapsed_headers = 0;
|
response.elapsed_headers = 0;
|
||||||
response.elapsed_dns = 0;
|
response.elapsed_dns = 0;
|
||||||
let persist_response = !response.request_id.is_empty();
|
// Responses with no request behind them are ephemeral: they belong to whoever called this
|
||||||
if persist_response {
|
// function and never reach the model store.
|
||||||
response = params
|
let store = storage.as_ref().filter(|_| !response.request_id.is_empty());
|
||||||
|
let persist_response = store.is_some();
|
||||||
|
if let Some(store) = store {
|
||||||
|
response = store
|
||||||
.query_manager
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(&response, ¶ms.update_source, params.blob_manager)
|
.upsert_http_response(&response, &store.update_source, store.blob_manager)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)?;
|
.map_err(SendHttpRequestError::PersistResponse)?;
|
||||||
} else if response.id.is_empty() {
|
} else if response.id.is_empty() {
|
||||||
response.id = generate_prefixed_id("rs");
|
response.id = generate_prefixed_id("rs");
|
||||||
@@ -527,14 +624,12 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
let request_body_id = format!("{}.request", response.id);
|
let request_body_id = format!("{}.request", response.id);
|
||||||
let mut request_body_capture_task = None;
|
let mut request_body_capture_task = None;
|
||||||
let mut request_body_capture_error = None;
|
let mut request_body_capture_error = None;
|
||||||
if persist_response {
|
if let Some(store) = store {
|
||||||
match sendable_request.body.as_mut() {
|
match sendable_request.body.as_mut() {
|
||||||
Some(SendableBody::Bytes(bytes)) => {
|
Some(SendableBody::Bytes(bytes)) => {
|
||||||
if let Err(err) = persist_request_body_bytes(
|
if let Err(err) =
|
||||||
params.blob_manager,
|
persist_request_body_bytes(store.blob_manager, &request_body_id, bytes.as_ref())
|
||||||
&request_body_id,
|
{
|
||||||
bytes.as_ref(),
|
|
||||||
) {
|
|
||||||
request_body_capture_error = Some(err);
|
request_body_capture_error = Some(err);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -543,7 +638,7 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
let inner = std::mem::replace(data, Box::pin(tokio::io::empty()));
|
let inner = std::mem::replace(data, Box::pin(tokio::io::empty()));
|
||||||
let tee_reader = TeeReader::new(inner, tx);
|
let tee_reader = TeeReader::new(inner, tx);
|
||||||
*data = Box::pin(tee_reader);
|
*data = Box::pin(tee_reader);
|
||||||
let blob_manager = params.blob_manager.clone();
|
let blob_manager = store.blob_manager.clone();
|
||||||
let body_id = request_body_id.clone();
|
let body_id = request_body_id.clone();
|
||||||
request_body_capture_task = Some(tokio::spawn(async move {
|
request_body_capture_task = Some(tokio::spawn(async move {
|
||||||
persist_request_body_stream(blob_manager, body_id, rx).await
|
persist_request_body_stream(blob_manager, body_id, rx).await
|
||||||
@@ -555,10 +650,9 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
|
|
||||||
let (event_tx, mut event_rx) =
|
let (event_tx, mut event_rx) =
|
||||||
mpsc::channel::<SenderHttpResponseEvent>(HTTP_EVENT_CHANNEL_CAPACITY);
|
mpsc::channel::<SenderHttpResponseEvent>(HTTP_EVENT_CHANNEL_CAPACITY);
|
||||||
let event_query_manager = params.query_manager.clone();
|
let event_store = store.map(|store| (store.query_manager.clone(), store.update_source.clone()));
|
||||||
let event_response_id = response.id.clone();
|
let event_response_id = response.id.clone();
|
||||||
let event_workspace_id = params.request.workspace_id.clone();
|
let event_workspace_id = request.workspace_id.clone();
|
||||||
let event_update_source = params.update_source.clone();
|
|
||||||
let emit_events_to = params.emit_events_to.clone();
|
let emit_events_to = params.emit_events_to.clone();
|
||||||
let dns_elapsed = Arc::new(AtomicI32::new(0));
|
let dns_elapsed = Arc::new(AtomicI32::new(0));
|
||||||
let event_dns_elapsed = dns_elapsed.clone();
|
let event_dns_elapsed = dns_elapsed.clone();
|
||||||
@@ -568,15 +662,14 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
event_dns_elapsed.store(u64_to_i32(*duration), Ordering::Relaxed);
|
event_dns_elapsed.store(u64_to_i32(*duration), Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
|
|
||||||
if persist_response {
|
if let Some((query_manager, update_source)) = event_store.as_ref() {
|
||||||
let db_event = HttpResponseEvent::new(
|
let db_event = HttpResponseEvent::new(
|
||||||
&event_response_id,
|
&event_response_id,
|
||||||
&event_workspace_id,
|
&event_workspace_id,
|
||||||
event.clone().into(),
|
event.clone().into(),
|
||||||
);
|
);
|
||||||
if let Err(err) = event_query_manager
|
if let Err(err) =
|
||||||
.connect()
|
query_manager.connect().upsert_http_response_event(&db_event, update_source)
|
||||||
.upsert_http_response_event(&db_event, &event_update_source)
|
|
||||||
{
|
{
|
||||||
warn!("Failed to persist HTTP response event: {}", err);
|
warn!("Failed to persist HTTP response event: {}", err);
|
||||||
}
|
}
|
||||||
@@ -595,7 +688,7 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
send_setting_event(
|
send_setting_event(
|
||||||
&event_tx,
|
&event_tx,
|
||||||
"validate_certificates",
|
"validate_certificates",
|
||||||
runtime_config.validate_certificates.to_string(),
|
resolved_settings.validate_certificates.value.to_string(),
|
||||||
&resolved_settings.validate_certificates,
|
&resolved_settings.validate_certificates,
|
||||||
);
|
);
|
||||||
send_setting_event(
|
send_setting_event(
|
||||||
@@ -627,16 +720,9 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
match executor.send(sendable_request, event_tx, cookie_behavior.clone()).await {
|
match executor.send(sendable_request, event_tx, cookie_behavior.clone()).await {
|
||||||
Ok(response) => response,
|
Ok(response) => response,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
persist_cookie_jar(
|
if let Some(store) = store {
|
||||||
params.query_manager,
|
|
||||||
cookie_jar.as_mut(),
|
|
||||||
cookie_behavior.store.as_ref(),
|
|
||||||
)?;
|
|
||||||
if persist_response {
|
|
||||||
let _ = persist_response_error(
|
let _ = persist_response_error(
|
||||||
params.query_manager,
|
store,
|
||||||
params.blob_manager,
|
|
||||||
¶ms.update_source,
|
|
||||||
&response,
|
&response,
|
||||||
started_at,
|
started_at,
|
||||||
err.to_string(),
|
err.to_string(),
|
||||||
@@ -654,14 +740,19 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
};
|
};
|
||||||
|
|
||||||
let headers_elapsed = duration_to_i32(started_at.elapsed());
|
let headers_elapsed = duration_to_i32(started_at.elapsed());
|
||||||
std::fs::create_dir_all(params.response_dir).map_err(|source| {
|
let body_path = match storage.as_ref() {
|
||||||
|
Some(storage) => {
|
||||||
|
std::fs::create_dir_all(storage.response_dir).map_err(|source| {
|
||||||
SendHttpRequestError::CreateResponseDirectory {
|
SendHttpRequestError::CreateResponseDirectory {
|
||||||
path: params.response_dir.to_path_buf(),
|
path: storage.response_dir.to_path_buf(),
|
||||||
source,
|
source,
|
||||||
}
|
}
|
||||||
})?;
|
})?;
|
||||||
let body_path = params.response_dir.join(&response.id);
|
Some(storage.response_dir.join(&response.id))
|
||||||
let response_body_path = body_path.to_string_lossy().to_string();
|
}
|
||||||
|
None => None,
|
||||||
|
};
|
||||||
|
let response_body_path = body_path.as_ref().map(|p| p.to_string_lossy().to_string());
|
||||||
let connected_response = HttpResponse {
|
let connected_response = HttpResponse {
|
||||||
state: HttpResponseState::Connected,
|
state: HttpResponseState::Connected,
|
||||||
elapsed_headers: headers_elapsed,
|
elapsed_headers: headers_elapsed,
|
||||||
@@ -671,7 +762,7 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
remote_addr: http_response.remote_addr.clone(),
|
remote_addr: http_response.remote_addr.clone(),
|
||||||
version: http_response.version.clone(),
|
version: http_response.version.clone(),
|
||||||
elapsed_dns: dns_elapsed.load(Ordering::Relaxed),
|
elapsed_dns: dns_elapsed.load(Ordering::Relaxed),
|
||||||
body_path: Some(response_body_path.clone()),
|
body_path: response_body_path.clone(),
|
||||||
content_length: http_response.content_length.map(u64_to_i32),
|
content_length: http_response.content_length.map(u64_to_i32),
|
||||||
headers: http_response
|
headers: http_response
|
||||||
.headers
|
.headers
|
||||||
@@ -685,20 +776,26 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
.collect(),
|
.collect(),
|
||||||
..response
|
..response
|
||||||
};
|
};
|
||||||
if persist_response {
|
if let Some(store) = store {
|
||||||
response = params
|
response = store
|
||||||
.query_manager
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(&connected_response, ¶ms.update_source, params.blob_manager)
|
.upsert_http_response(&connected_response, &store.update_source, store.blob_manager)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)?;
|
.map_err(SendHttpRequestError::PersistResponse)?;
|
||||||
} else {
|
} else {
|
||||||
response = connected_response;
|
response = connected_response;
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut file =
|
let mut body_file = match body_path {
|
||||||
File::options().create(true).truncate(true).write(true).open(&body_path).await.map_err(
|
Some(path) => {
|
||||||
|source| SendHttpRequestError::WriteResponseBody { path: body_path.clone(), source },
|
let file =
|
||||||
|
File::options().create(true).truncate(true).write(true).open(&path).await.map_err(
|
||||||
|
|source| SendHttpRequestError::WriteResponseBody { path: path.clone(), source },
|
||||||
)?;
|
)?;
|
||||||
|
Some((file, path))
|
||||||
|
}
|
||||||
|
None => None,
|
||||||
|
};
|
||||||
let mut body_stream =
|
let mut body_stream =
|
||||||
http_response.into_body_stream().map_err(SendHttpRequestError::ReadResponseBody)?;
|
http_response.into_body_stream().map_err(SendHttpRequestError::ReadResponseBody)?;
|
||||||
let mut response_body = Vec::new();
|
let mut response_body = Vec::new();
|
||||||
@@ -737,9 +834,11 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
Ok(n) => {
|
Ok(n) => {
|
||||||
written_bytes += n;
|
written_bytes += n;
|
||||||
let chunk = &read_buf[..n];
|
let chunk = &read_buf[..n];
|
||||||
|
if let Some((file, path)) = body_file.as_mut() {
|
||||||
file.write_all(chunk).await.map_err(|source| {
|
file.write_all(chunk).await.map_err(|source| {
|
||||||
SendHttpRequestError::WriteResponseBody { path: body_path.clone(), source }
|
SendHttpRequestError::WriteResponseBody { path: path.clone(), source }
|
||||||
})?;
|
})?;
|
||||||
|
}
|
||||||
if let Some(tx) = params.emit_response_body_chunks_to.as_ref() {
|
if let Some(tx) = params.emit_response_body_chunks_to.as_ref() {
|
||||||
let _ = tx.send(chunk.to_vec());
|
let _ = tx.send(chunk.to_vec());
|
||||||
} else if collect_response_body {
|
} else if collect_response_body {
|
||||||
@@ -757,14 +856,14 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
elapsed_dns: dns_elapsed.load(Ordering::Relaxed),
|
elapsed_dns: dns_elapsed.load(Ordering::Relaxed),
|
||||||
..response.clone()
|
..response.clone()
|
||||||
};
|
};
|
||||||
if persist_response {
|
if let Some(store) = store {
|
||||||
response = params
|
response = store
|
||||||
.query_manager
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(
|
.upsert_http_response(
|
||||||
&progress_response,
|
&progress_response,
|
||||||
¶ms.update_source,
|
&store.update_source,
|
||||||
params.blob_manager,
|
store.blob_manager,
|
||||||
)
|
)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)?;
|
.map_err(SendHttpRequestError::PersistResponse)?;
|
||||||
} else {
|
} else {
|
||||||
@@ -782,10 +881,12 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if let Some((file, path)) = body_file.as_mut() {
|
||||||
file.flush().await.map_err(|source| SendHttpRequestError::WriteResponseBody {
|
file.flush().await.map_err(|source| SendHttpRequestError::WriteResponseBody {
|
||||||
path: body_path.clone(),
|
path: path.clone(),
|
||||||
source,
|
source,
|
||||||
})?;
|
})?;
|
||||||
|
}
|
||||||
drop(body_stream);
|
drop(body_stream);
|
||||||
|
|
||||||
if let Some(err) = request_body_capture_error.take() {
|
if let Some(err) = request_body_capture_error.take() {
|
||||||
@@ -796,22 +897,15 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if let Some(err) = body_read_error {
|
if let Some(err) = body_read_error {
|
||||||
if persist_response {
|
if let Some(store) = store {
|
||||||
let _ = persist_response_error(
|
let _ = persist_response_error(
|
||||||
params.query_manager,
|
store,
|
||||||
params.blob_manager,
|
|
||||||
¶ms.update_source,
|
|
||||||
&response,
|
&response,
|
||||||
started_at,
|
started_at,
|
||||||
err.to_string(),
|
err.to_string(),
|
||||||
request_started_url,
|
request_started_url,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
persist_cookie_jar(
|
|
||||||
params.query_manager,
|
|
||||||
cookie_jar.as_mut(),
|
|
||||||
cookie_behavior.store.as_ref(),
|
|
||||||
)?;
|
|
||||||
if let Some(task) = request_body_capture_task.take() {
|
if let Some(task) = request_body_capture_task.take() {
|
||||||
match task.await {
|
match task.await {
|
||||||
Ok(Ok(_)) => {}
|
Ok(Ok(_)) => {}
|
||||||
@@ -827,7 +921,7 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
|
|
||||||
let compressed_length = http_response.content_length.unwrap_or(written_bytes as u64);
|
let compressed_length = http_response.content_length.unwrap_or(written_bytes as u64);
|
||||||
let final_response = HttpResponse {
|
let final_response = HttpResponse {
|
||||||
body_path: Some(response_body_path),
|
body_path: response_body_path,
|
||||||
content_length: Some(usize_to_i32(written_bytes)),
|
content_length: Some(usize_to_i32(written_bytes)),
|
||||||
content_length_compressed: Some(u64_to_i32(compressed_length)),
|
content_length_compressed: Some(u64_to_i32(compressed_length)),
|
||||||
elapsed: duration_to_i32(started_at.elapsed()),
|
elapsed: duration_to_i32(started_at.elapsed()),
|
||||||
@@ -836,18 +930,16 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
state: HttpResponseState::Closed,
|
state: HttpResponseState::Closed,
|
||||||
..response
|
..response
|
||||||
};
|
};
|
||||||
if persist_response {
|
if let Some(store) = store {
|
||||||
response = params
|
response = store
|
||||||
.query_manager
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(&final_response, ¶ms.update_source, params.blob_manager)
|
.upsert_http_response(&final_response, &store.update_source, store.blob_manager)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)?;
|
.map_err(SendHttpRequestError::PersistResponse)?;
|
||||||
} else {
|
} else {
|
||||||
response = final_response;
|
response = final_response;
|
||||||
}
|
}
|
||||||
|
|
||||||
persist_cookie_jar(params.query_manager, cookie_jar.as_mut(), cookie_behavior.store.as_ref())?;
|
|
||||||
|
|
||||||
// Request-body history can be much larger than the response. It should not keep the
|
// Request-body history can be much larger than the response. It should not keep the
|
||||||
// response in a loading state after the network/response-body work has completed.
|
// response in a loading state after the network/response-body work has completed.
|
||||||
if let Some(task) = request_body_capture_task.take() {
|
if let Some(task) = request_body_capture_task.take() {
|
||||||
@@ -876,11 +968,11 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if update_response && persist_response {
|
if update_response && let Some(store) = store {
|
||||||
response = params
|
response = store
|
||||||
.query_manager
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(&response, ¶ms.update_source, params.blob_manager)
|
.upsert_http_response(&response, &store.update_source, store.blob_manager)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)?;
|
.map_err(SendHttpRequestError::PersistResponse)?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -891,7 +983,12 @@ pub async fn send_http_request<T: TemplateCallback>(
|
|||||||
warn!("Failed to join response event task: {}", join_err);
|
warn!("Failed to join response event task: {}", join_err);
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(SendHttpRequestResult { rendered_request, response, response_body })
|
Ok(SendHttpRequestResult {
|
||||||
|
rendered_request,
|
||||||
|
response,
|
||||||
|
response_body,
|
||||||
|
cookies: cookie_behavior.store.as_ref().map(|store| store.get_all_cookies()),
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
fn persist_request_body_bytes(
|
fn persist_request_body_bytes(
|
||||||
@@ -953,37 +1050,7 @@ fn append_error_message(existing_error: Option<String>, message: String) -> Stri
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn resolve_environment_chain(
|
pub fn load_cookie_jar(
|
||||||
query_manager: &QueryManager,
|
|
||||||
request: &HttpRequest,
|
|
||||||
environment_id: Option<&str>,
|
|
||||||
) -> Result<Vec<Environment>> {
|
|
||||||
let db = query_manager.connect();
|
|
||||||
db.resolve_environments(&request.workspace_id, request.folder_id.as_deref(), environment_id)
|
|
||||||
.map_err(SendHttpRequestError::ResolveEnvironments)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn resolve_inherited_request(
|
|
||||||
query_manager: &QueryManager,
|
|
||||||
request: &HttpRequest,
|
|
||||||
) -> Result<(HttpRequest, String)> {
|
|
||||||
let db = query_manager.connect();
|
|
||||||
let (authentication_type, authentication, auth_context_id) = db
|
|
||||||
.resolve_auth_for_http_request(request)
|
|
||||||
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
|
||||||
let resolved_headers = db
|
|
||||||
.resolve_headers_for_http_request(request)
|
|
||||||
.map_err(SendHttpRequestError::ResolveRequestInheritance)?;
|
|
||||||
|
|
||||||
let mut request = request.clone();
|
|
||||||
request.authentication_type = authentication_type;
|
|
||||||
request.authentication = authentication;
|
|
||||||
request.headers = resolved_headers;
|
|
||||||
|
|
||||||
Ok((request, auth_context_id))
|
|
||||||
}
|
|
||||||
|
|
||||||
fn load_cookie_jar(
|
|
||||||
query_manager: &QueryManager,
|
query_manager: &QueryManager,
|
||||||
cookie_jar_id: Option<&str>,
|
cookie_jar_id: Option<&str>,
|
||||||
) -> Result<Option<CookieJar>> {
|
) -> Result<Option<CookieJar>> {
|
||||||
@@ -998,22 +1065,32 @@ fn load_cookie_jar(
|
|||||||
.map_err(SendHttpRequestError::LoadCookieJar)
|
.map_err(SendHttpRequestError::LoadCookieJar)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn persist_cookie_jar(
|
/// Write the cookies a send collected back to its jar.
|
||||||
|
///
|
||||||
|
/// The store is shared with the HTTP transaction, so it holds every cookie picked up along the
|
||||||
|
/// way no matter how the send ended — including when the response arrived but the storage work
|
||||||
|
/// after it failed. Call this whatever the send returned; a send that failed before the
|
||||||
|
/// transaction started leaves the store untouched, which compares equal and writes nothing.
|
||||||
|
pub fn persist_cookies_after_send(
|
||||||
query_manager: &QueryManager,
|
query_manager: &QueryManager,
|
||||||
cookie_jar: Option<&mut CookieJar>,
|
cookie_jar: Option<&mut CookieJar>,
|
||||||
cookie_store: Option<&CookieStore>,
|
cookie_store: Option<&CookieStore>,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
match (cookie_jar, cookie_store) {
|
let (Some(cookie_jar), Some(cookie_store)) = (cookie_jar, cookie_store) else {
|
||||||
(Some(cookie_jar), Some(cookie_store)) => {
|
return Ok(());
|
||||||
cookie_jar.cookies = cookie_store.get_all_cookies();
|
};
|
||||||
|
|
||||||
|
let cookies = cookie_store.get_all_cookies();
|
||||||
|
if cookies == cookie_jar.cookies {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
cookie_jar.cookies = cookies;
|
||||||
query_manager
|
query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_cookie_jar(cookie_jar, &UpdateSource::Background)
|
.upsert_cookie_jar(cookie_jar, &UpdateSource::Background)
|
||||||
.map_err(SendHttpRequestError::PersistCookieJar)?;
|
.map_err(SendHttpRequestError::PersistCookieJar)?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
|
||||||
_ => Ok(()),
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn send_setting_event<T>(
|
fn send_setting_event<T>(
|
||||||
@@ -1118,16 +1195,15 @@ pub async fn apply_plugin_authentication(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn persist_response_error(
|
fn persist_response_error(
|
||||||
query_manager: &QueryManager,
|
store: &ResponseStorage,
|
||||||
blob_manager: &BlobManager,
|
|
||||||
update_source: &UpdateSource,
|
|
||||||
response: &HttpResponse,
|
response: &HttpResponse,
|
||||||
started_at: Instant,
|
started_at: Instant,
|
||||||
error: String,
|
error: String,
|
||||||
fallback_url: String,
|
fallback_url: String,
|
||||||
) -> Result<HttpResponse> {
|
) -> Result<HttpResponse> {
|
||||||
let elapsed = duration_to_i32(started_at.elapsed());
|
let elapsed = duration_to_i32(started_at.elapsed());
|
||||||
query_manager
|
store
|
||||||
|
.query_manager
|
||||||
.connect()
|
.connect()
|
||||||
.upsert_http_response(
|
.upsert_http_response(
|
||||||
&HttpResponse {
|
&HttpResponse {
|
||||||
@@ -1142,8 +1218,8 @@ fn persist_response_error(
|
|||||||
url: if response.url.is_empty() { fallback_url } else { response.url.clone() },
|
url: if response.url.is_empty() { fallback_url } else { response.url.clone() },
|
||||||
..response.clone()
|
..response.clone()
|
||||||
},
|
},
|
||||||
update_source,
|
&store.update_source,
|
||||||
blob_manager,
|
store.blob_manager,
|
||||||
)
|
)
|
||||||
.map_err(SendHttpRequestError::PersistResponse)
|
.map_err(SendHttpRequestError::PersistResponse)
|
||||||
}
|
}
|
||||||
@@ -1173,3 +1249,232 @@ fn u64_to_i32(value: u64) -> i32 {
|
|||||||
fn u128_to_i32(value: u128) -> i32 {
|
fn u128_to_i32(value: u128) -> i32 {
|
||||||
if value > i32::MAX as u128 { i32::MAX } else { value as i32 }
|
if value > i32::MAX as u128 { i32::MAX } else { value as i32 }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::pin::Pin;
|
||||||
|
use tempfile::TempDir;
|
||||||
|
use tokio::io::AsyncRead;
|
||||||
|
use yaak_http::decompress::ContentEncoding;
|
||||||
|
use yaak_models::models::{CookieDomain, CookieExpires, Workspace};
|
||||||
|
|
||||||
|
struct NoopTemplateCallback;
|
||||||
|
|
||||||
|
impl TemplateCallback for NoopTemplateCallback {
|
||||||
|
async fn run(
|
||||||
|
&self,
|
||||||
|
_fn_name: &str,
|
||||||
|
_args: HashMap<String, serde_json::Value>,
|
||||||
|
) -> yaak_templates::error::Result<String> {
|
||||||
|
Ok(String::new())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn transform_arg(
|
||||||
|
&self,
|
||||||
|
_fn_name: &str,
|
||||||
|
_arg_name: &str,
|
||||||
|
arg_value: &str,
|
||||||
|
) -> yaak_templates::error::Result<String> {
|
||||||
|
Ok(arg_value.to_string())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
struct StubExecutor {
|
||||||
|
body: &'static [u8],
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl SendRequestExecutor for StubExecutor {
|
||||||
|
async fn send(
|
||||||
|
&self,
|
||||||
|
sendable_request: SendableHttpRequest,
|
||||||
|
event_tx: mpsc::Sender<SenderHttpResponseEvent>,
|
||||||
|
_cookie_behavior: CookieBehavior,
|
||||||
|
) -> yaak_http::error::Result<yaak_http::sender::HttpResponse> {
|
||||||
|
let _ = event_tx.try_send(SenderHttpResponseEvent::HeaderDown(
|
||||||
|
"content-type".to_string(),
|
||||||
|
"text/plain".to_string(),
|
||||||
|
));
|
||||||
|
let body: Pin<Box<dyn AsyncRead + Send>> =
|
||||||
|
Box::pin(std::io::Cursor::new(self.body.to_vec()));
|
||||||
|
Ok(yaak_http::sender::HttpResponse::new(
|
||||||
|
200,
|
||||||
|
Some("OK".to_string()),
|
||||||
|
vec![("content-type".to_string(), "text/plain".to_string())],
|
||||||
|
Vec::new(),
|
||||||
|
Some(self.body.len() as u64),
|
||||||
|
sendable_request.url.clone(),
|
||||||
|
None,
|
||||||
|
Some("HTTP/1.1".to_string()),
|
||||||
|
body,
|
||||||
|
ContentEncoding::Identity,
|
||||||
|
))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The hosted sender runs with no query manager, blob manager, or response directory. Nothing
|
||||||
|
/// here touches storage, so the caller's only view of the response is what gets streamed out.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn sends_without_a_database() {
|
||||||
|
let (event_tx, mut event_rx) = mpsc::channel(HTTP_EVENT_CHANNEL_CAPACITY);
|
||||||
|
let (chunk_tx, mut chunk_rx) = mpsc::unbounded_channel();
|
||||||
|
let executor = StubExecutor { body: b"hello world" };
|
||||||
|
|
||||||
|
let result = send_http_request(SendHttpRequestParams {
|
||||||
|
inputs: HttpSendInputs {
|
||||||
|
request: ResolvedHttpRequest::assume_resolved(
|
||||||
|
HttpRequest {
|
||||||
|
workspace_id: "wk_test".to_string(),
|
||||||
|
url: "http://localhost/test".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
String::new(),
|
||||||
|
),
|
||||||
|
environment_chain: Vec::new(),
|
||||||
|
runtime_config: HttpSendRuntimeConfig {
|
||||||
|
settings: ResolvedHttpRequestSettings::default(),
|
||||||
|
proxy: HttpConnectionProxySetting::System,
|
||||||
|
dns_overrides: Vec::new(),
|
||||||
|
client_certificates: Vec::new(),
|
||||||
|
},
|
||||||
|
cookie_store: Some(CookieStore::new()),
|
||||||
|
},
|
||||||
|
template_callback: &NoopTemplateCallback,
|
||||||
|
storage: None,
|
||||||
|
emit_events_to: Some(event_tx),
|
||||||
|
emit_response_body_chunks_to: Some(chunk_tx),
|
||||||
|
cancelled_rx: None,
|
||||||
|
existing_response: None,
|
||||||
|
prepare_sendable_request: None,
|
||||||
|
executor: &executor,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("send should succeed without a database");
|
||||||
|
|
||||||
|
assert_eq!(result.response.status, 200);
|
||||||
|
assert!(matches!(result.response.state, HttpResponseState::Closed));
|
||||||
|
assert_eq!(result.response.content_length, Some(11));
|
||||||
|
assert!(result.cookies.is_some());
|
||||||
|
|
||||||
|
// Nothing was written, so the response carries no body file and only ever lived in memory.
|
||||||
|
assert_eq!(result.response.body_path, None);
|
||||||
|
assert!(!result.response.id.is_empty());
|
||||||
|
|
||||||
|
let mut body = Vec::new();
|
||||||
|
while let Some(chunk) = chunk_rx.recv().await {
|
||||||
|
body.extend_from_slice(&chunk);
|
||||||
|
}
|
||||||
|
assert_eq!(body, b"hello world");
|
||||||
|
|
||||||
|
let mut events = Vec::new();
|
||||||
|
while let Some(event) = event_rx.recv().await {
|
||||||
|
events.push(event);
|
||||||
|
}
|
||||||
|
assert!(
|
||||||
|
events.iter().any(
|
||||||
|
|e| matches!(e, SenderHttpResponseEvent::Setting { name, .. } if name == "redirects")
|
||||||
|
),
|
||||||
|
"expected timeline settings events, got {events:?}"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
events.iter().any(|e| matches!(e, SenderHttpResponseEvent::HeaderDown(..))),
|
||||||
|
"expected timeline events from the executor, got {events:?}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn seed_cookie_jar() -> (QueryManager, CookieJar, TempDir) {
|
||||||
|
let temp_dir = TempDir::new().expect("Failed to create temp dir");
|
||||||
|
let (query_manager, _blob_manager, _rx) = yaak_models::init_standalone(
|
||||||
|
&temp_dir.path().join("db.sqlite"),
|
||||||
|
&temp_dir.path().join("blobs.sqlite"),
|
||||||
|
)
|
||||||
|
.expect("Failed to initialize DB");
|
||||||
|
|
||||||
|
query_manager
|
||||||
|
.connect()
|
||||||
|
.upsert_workspace(
|
||||||
|
&Workspace { id: "wk_test".to_string(), ..Default::default() },
|
||||||
|
&UpdateSource::Sync,
|
||||||
|
)
|
||||||
|
.expect("Failed to seed workspace");
|
||||||
|
let cookie_jar = query_manager
|
||||||
|
.connect()
|
||||||
|
.upsert_cookie_jar(
|
||||||
|
&CookieJar {
|
||||||
|
id: "cj_test".to_string(),
|
||||||
|
workspace_id: "wk_test".to_string(),
|
||||||
|
name: "Default".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
&UpdateSource::Sync,
|
||||||
|
)
|
||||||
|
.expect("Failed to seed cookie jar");
|
||||||
|
|
||||||
|
(query_manager, cookie_jar, temp_dir)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn cookie(name: &str) -> Cookie {
|
||||||
|
Cookie {
|
||||||
|
name: name.to_string(),
|
||||||
|
value: "value".to_string(),
|
||||||
|
domain: CookieDomain::HostOnly("localhost".to_string()),
|
||||||
|
expires: CookieExpires::SessionEnd,
|
||||||
|
path: "/".to_string(),
|
||||||
|
secure: false,
|
||||||
|
http_only: false,
|
||||||
|
same_site: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Cookies the transaction collected must survive no matter how the send ended, including the
|
||||||
|
/// storage work that runs after a response has already arrived.
|
||||||
|
#[test]
|
||||||
|
fn persists_cookies_collected_before_a_failure() {
|
||||||
|
let (query_manager, mut cookie_jar, _temp_dir) = seed_cookie_jar();
|
||||||
|
let store = CookieStore::from_cookies(cookie_jar.cookies.clone());
|
||||||
|
store.store_cookies_from_response(
|
||||||
|
&"http://localhost/test".parse().expect("valid url"),
|
||||||
|
&["session=abc123; Path=/".to_string()],
|
||||||
|
);
|
||||||
|
|
||||||
|
persist_cookies_after_send(&query_manager, Some(&mut cookie_jar), Some(&store))
|
||||||
|
.expect("Failed to persist cookies");
|
||||||
|
|
||||||
|
let stored =
|
||||||
|
query_manager.connect().get_cookie_jar("cj_test").expect("Failed to load cookie jar");
|
||||||
|
assert_eq!(stored.cookies.len(), 1);
|
||||||
|
assert_eq!(stored.cookies[0].name, "session");
|
||||||
|
assert_eq!(stored.cookies[0].value, "abc123");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A send that failed before the transaction started leaves the store untouched, so there is
|
||||||
|
/// nothing to write and a concurrent update to the jar must not be clobbered.
|
||||||
|
#[test]
|
||||||
|
fn leaves_the_jar_alone_when_no_cookies_changed() {
|
||||||
|
let (query_manager, mut cookie_jar, _temp_dir) = seed_cookie_jar();
|
||||||
|
cookie_jar.cookies = vec![cookie("original")];
|
||||||
|
cookie_jar = query_manager
|
||||||
|
.connect()
|
||||||
|
.upsert_cookie_jar(&cookie_jar, &UpdateSource::Sync)
|
||||||
|
.expect("Failed to seed cookies");
|
||||||
|
let store = CookieStore::from_cookies(cookie_jar.cookies.clone());
|
||||||
|
|
||||||
|
// Someone else updates the jar while the send is in flight.
|
||||||
|
query_manager
|
||||||
|
.connect()
|
||||||
|
.upsert_cookie_jar(
|
||||||
|
&CookieJar { cookies: vec![cookie("newer")], ..cookie_jar.clone() },
|
||||||
|
&UpdateSource::Sync,
|
||||||
|
)
|
||||||
|
.expect("Failed to update cookie jar");
|
||||||
|
|
||||||
|
persist_cookies_after_send(&query_manager, Some(&mut cookie_jar), Some(&store))
|
||||||
|
.expect("Failed to persist cookies");
|
||||||
|
|
||||||
|
let stored =
|
||||||
|
query_manager.connect().get_cookie_jar("cj_test").expect("Failed to load cookie jar");
|
||||||
|
assert_eq!(stored.cookies, vec![cookie("newer")]);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user