aboutsummaryrefslogtreecommitdiffhomepage
path: root/crates/shirabe/src/util/http/curl_downloader.rs
diff options
context:
space:
mode:
Diffstat (limited to 'crates/shirabe/src/util/http/curl_downloader.rs')
-rw-r--r--crates/shirabe/src/util/http/curl_downloader.rs173
1 files changed, 101 insertions, 72 deletions
diff --git a/crates/shirabe/src/util/http/curl_downloader.rs b/crates/shirabe/src/util/http/curl_downloader.rs
index 0519f26..94e703d 100644
--- a/crates/shirabe/src/util/http/curl_downloader.rs
+++ b/crates/shirabe/src/util/http/curl_downloader.rs
@@ -14,9 +14,9 @@ use shirabe_php_shim::{
CURLOPT_FILE, CURLOPT_FOLLOWLOCATION, CURLOPT_HTTP_VERSION, CURLOPT_IPRESOLVE,
CURLOPT_PROTOCOLS, CURLOPT_SHARE, CURLOPT_TIMEOUT, CURLOPT_URL, CURLOPT_WRITEHEADER,
CURLPROTO_HTTP, CURLPROTO_HTTPS, CURLSHOPT_SHARE, CurlMultiHandle, CurlShareHandle,
- LogicException, PHP_VERSION_ID, PhpMixed, RuntimeException, array_diff, array_diff_key,
- array_merge, count, curl_errno, curl_error, curl_getinfo, curl_handle_id, curl_init,
- curl_multi_add_handle, curl_multi_exec, curl_multi_info_read, curl_multi_init,
+ LogicException, PHP_VERSION_ID, PhpMixed, PhpResource, RuntimeException, array_diff,
+ array_diff_key, array_merge, count, curl_errno, curl_error, curl_getinfo, curl_handle_id,
+ curl_init, curl_multi_add_handle, curl_multi_exec, curl_multi_info_read, curl_multi_init,
curl_multi_select, curl_multi_setopt, curl_setopt, curl_setopt_array, curl_share_init,
curl_share_setopt, curl_strerror, curl_version, defined, explode, fclose, fopen,
function_exists, implode, in_array, ini_get, is_resource, json_decode, parse_url, preg_quote,
@@ -47,6 +47,10 @@ pub struct CurlDownloader {
share_handle: Option<CurlShareHandle>,
/// @var Job[]
jobs: IndexMap<i64, IndexMap<String, PhpMixed>>,
+ /// The `headerHandle`/`bodyHandle` stream resources of each job, keyed by job id. PHP keeps
+ /// them inside the Job array, but a PhpResource cannot live in the PhpMixed-typed job map.
+ header_handles: IndexMap<i64, PhpResource>,
+ body_handles: IndexMap<i64, PhpResource>,
/// @var IOInterface
io: std::rc::Rc<std::cell::RefCell<dyn IOInterface>>,
/// @var Config
@@ -232,6 +236,8 @@ impl CurlDownloader {
multi_handle,
share_handle,
jobs: IndexMap::new(),
+ header_handles: IndexMap::new(),
+ body_handles: IndexMap::new(),
io,
config,
auth_helper,
@@ -329,16 +335,18 @@ impl CurlDownloader {
}
let curl_handle = curl_init();
- let header_handle = fopen("php://temp/maxmemory:32768", "w+b");
- if matches!(header_handle, PhpMixed::Bool(false)) {
- anyhow::bail!(
- RuntimeException {
- message: "Failed to open a temp stream to store curl headers".to_string(),
- code: 0,
- }
- .message
- );
- }
+ let header_handle = match fopen("php://temp/maxmemory:32768", "w+b") {
+ Ok(header_handle) => header_handle,
+ Err(_) => {
+ anyhow::bail!(
+ RuntimeException {
+ message: "Failed to open a temp stream to store curl headers".to_string(),
+ code: 0,
+ }
+ .message
+ );
+ }
+ };
let body_target: String = if let Some(copy_to) = copy_to {
format!("{}~", copy_to)
@@ -350,18 +358,20 @@ impl CurlDownloader {
// (stripping the "fopen(...): " prefix) into the message. Rust I/O reports failures through
// return values, not warnings, so error_message stays empty until fopen surfaces its reason.
let error_message = String::new();
- let body_handle = fopen(&body_target, "w+b");
- if matches!(body_handle, PhpMixed::Bool(false)) {
- return Ok(Err(TransportException::new(
- format!(
- "The \"{}\" file could not be written to {}: {}",
- url,
- copy_to.unwrap_or("a temporary file"),
- error_message
- ),
- 0,
- )));
- }
+ let body_handle = match fopen(&body_target, "w+b") {
+ Ok(body_handle) => body_handle,
+ Err(_) => {
+ return Ok(Err(TransportException::new(
+ format!(
+ "The \"{}\" file could not be written to {}: {}",
+ url,
+ copy_to.unwrap_or("a temporary file"),
+ error_message
+ ),
+ 0,
+ )));
+ }
+ };
curl_setopt(&curl_handle, CURLOPT_URL, PhpMixed::String(url.to_string()));
curl_setopt(&curl_handle, CURLOPT_FOLLOWLOCATION, PhpMixed::Bool(false));
@@ -378,8 +388,10 @@ impl CurlDownloader {
.max(300),
),
);
- curl_setopt(&curl_handle, CURLOPT_WRITEHEADER, header_handle.clone());
- curl_setopt(&curl_handle, CURLOPT_FILE, body_handle.clone());
+ // TODO(phase-d): CURLOPT_WRITEHEADER/CURLOPT_FILE hand the header/body stream resources to
+ // curl, but curl_setopt takes a PhpMixed and a PhpResource cannot be passed through it.
+ // curl_setopt(&curl_handle, CURLOPT_WRITEHEADER, header_handle);
+ // curl_setopt(&curl_handle, CURLOPT_FILE, body_handle);
curl_setopt(
&curl_handle,
CURLOPT_ENCODING,
@@ -601,8 +613,11 @@ impl CurlDownloader {
.map(|s| PhpMixed::String(s.to_string()))
.unwrap_or(PhpMixed::Null),
);
- job.insert("headerHandle".to_string(), header_handle.clone());
- job.insert("bodyHandle".to_string(), body_handle.clone());
+ // headerHandle/bodyHandle are PHP stream resources; they live in the Rust-side
+ // header_handles/body_handles maps (keyed by job id) since a PhpResource cannot be stored
+ // in the PhpMixed job map.
+ job.insert("headerHandle".to_string(), PhpMixed::Null);
+ job.insert("bodyHandle".to_string(), PhpMixed::Null);
job.insert("resolve".to_string(), PhpMixed::Null);
job.insert("reject".to_string(), PhpMixed::Null);
job.insert("primaryIp".to_string(), PhpMixed::String(String::new()));
@@ -611,7 +626,10 @@ impl CurlDownloader {
// above) — blocked on the typed-Job refactor and the React\Promise model.
let _ = (resolve, reject);
- self.jobs.insert(curl_handle_id(&curl_handle), job);
+ let job_id = curl_handle_id(&curl_handle);
+ self.header_handles.insert(job_id, header_handle);
+ self.body_handles.insert(job_id, body_handle);
+ self.jobs.insert(job_id, job);
let using_proxy = proxy
.get_status(Some(" using proxy (%s)"))
@@ -677,15 +695,17 @@ impl CurlDownloader {
if PHP_VERSION_ID < 80000 {
// curl_close($job['curlHandle']);
}
- if is_resource(job.get("headerHandle").unwrap_or(&PhpMixed::Null)) {
- fclose(job.get("headerHandle").cloned().unwrap_or(PhpMixed::Null));
+ if let Some(handle) = self.header_handles.get(&id) {
+ fclose(handle);
}
- if is_resource(job.get("bodyHandle").unwrap_or(&PhpMixed::Null)) {
- fclose(job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null));
+ if let Some(handle) = self.body_handles.get(&id) {
+ fclose(handle);
}
if let Some(PhpMixed::String(filename)) = job.get("filename") {
unlink_silent(&format!("{}~", filename));
}
+ self.header_handles.shift_remove(&id);
+ self.body_handles.shift_remove(&id);
self.jobs.shift_remove(&id);
}
}
@@ -939,18 +959,19 @@ impl CurlDownloader {
)));
}
status_code = progress.get("http_code").and_then(|b| b.as_int());
- rewind(job.get("headerHandle").cloned().unwrap_or(PhpMixed::Null));
- headers = Some(explode(
- "\r\n",
- &rtrim(
- &stream_get_contents(
- job.get("headerHandle").cloned().unwrap_or(PhpMixed::Null),
- )
- .unwrap_or_default(),
- None,
- ),
- ));
- fclose(job.get("headerHandle").cloned().unwrap_or(PhpMixed::Null));
+ let header_handle = self.header_handles.get(&i).cloned();
+ let header_contents = match &header_handle {
+ Some(handle) => {
+ rewind(handle);
+ stream_get_contents(handle).unwrap_or_default()
+ }
+ None => String::new(),
+ };
+ headers = Some(explode("\r\n", &rtrim(&header_contents, None)));
+ if let Some(handle) = &header_handle {
+ fclose(handle);
+ }
+ self.header_handles.shift_remove(&i);
if status_code == Some(0) {
anyhow::bail!(
@@ -984,17 +1005,15 @@ impl CurlDownloader {
}
// prepare response object
+ let body_handle = self.body_handles.get(&i).cloned();
let contents: PhpMixed;
if let Some(PhpMixed::String(filename)) = job.get("filename") {
let mut c: PhpMixed = PhpMixed::String(format!("{}~", filename));
if status_code.unwrap_or(0) >= 300 {
- rewind(job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null));
- c = PhpMixed::String(
- stream_get_contents(
- job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null),
- )
- .unwrap_or_default(),
- );
+ if let Some(handle) = &body_handle {
+ rewind(handle);
+ c = PhpMixed::String(stream_get_contents(handle).unwrap_or_default());
+ }
}
contents = c;
response = Some(CurlResponse::new(
@@ -1027,12 +1046,16 @@ impl CurlDownloader {
.and_then(|v| v.as_array())
.and_then(|a| a.get("max_file_size"))
.and_then(|b| b.as_int());
- rewind(job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null));
+ if let Some(handle) = &body_handle {
+ rewind(handle);
+ }
if let Some(max_file_size) = max_file_size {
- let c = stream_get_contents_with_max(
- job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null),
- Some(max_file_size),
- );
+ let c = match &body_handle {
+ Some(handle) => {
+ stream_get_contents_with_max(handle, Some(max_file_size))
+ }
+ None => None,
+ };
// Gzipped responses with missing Content-Length header cannot be detected during the file download
// because $progress['size_download'] refers to the gzipped size downloaded, not the actual file size
if let Some(c_str) = c.as_deref()
@@ -1053,12 +1076,10 @@ impl CurlDownloader {
}
contents = PhpMixed::String(c.unwrap_or_default());
} else {
- contents = PhpMixed::String(
- stream_get_contents(
- job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null),
- )
- .unwrap_or_default(),
- );
+ contents = PhpMixed::String(match &body_handle {
+ Some(handle) => stream_get_contents(handle).unwrap_or_default(),
+ None => String::new(),
+ });
}
response = Some(CurlResponse::new(
@@ -1086,7 +1107,10 @@ impl CurlDownloader {
crate::io::DEBUG,
);
}
- fclose(job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null));
+ if let Some(handle) = &body_handle {
+ fclose(handle);
+ }
+ self.body_handles.shift_remove(&i);
let response_ref = response.as_ref().unwrap();
if response_ref.inner.get_status_code() >= 300
@@ -1310,11 +1334,11 @@ impl CurlDownloader {
e.set_response(r.inner.get_body().map(|s| s.to_string()));
}
e.set_response_info(progress.values().cloned().collect::<Vec<_>>());
- self.reject_job(&job, anyhow::anyhow!(e.message));
+ self.reject_job(i, &job, anyhow::anyhow!(e.message));
}
Err(e) => {
// Non-TransportException fatal error: pass through reject_job to mirror PHP catch (\Exception $e).
- self.reject_job(&job, e);
+ self.reject_job(i, &job, e);
}
}
}
@@ -1369,6 +1393,7 @@ impl CurlDownloader {
if max_file_size < download_content_length {
let job = self.jobs.get(&i).cloned().unwrap_or_default();
self.reject_job(
+ i,
&job,
anyhow::anyhow!(
MaxFileSizeExceededException(TransportException::new(
@@ -1392,6 +1417,7 @@ impl CurlDownloader {
if max_file_size < size_download {
let job = self.jobs.get(&i).cloned().unwrap_or_default();
self.reject_job(
+ i,
&job,
anyhow::anyhow!(
MaxFileSizeExceededException(TransportException::new(
@@ -1438,6 +1464,7 @@ impl CurlDownloader {
if blocked {
let job = self.jobs.get(&i).cloned().unwrap_or_default();
self.reject_job(
+ i,
&job,
anyhow::anyhow!(
TransportException::new(
@@ -1801,13 +1828,15 @@ impl CurlDownloader {
}
/// @param Job $job
- fn reject_job(&mut self, job: &IndexMap<String, PhpMixed>, _e: anyhow::Error) {
- if is_resource(job.get("headerHandle").unwrap_or(&PhpMixed::Null)) {
- fclose(job.get("headerHandle").cloned().unwrap_or(PhpMixed::Null));
+ fn reject_job(&mut self, id: i64, job: &IndexMap<String, PhpMixed>, _e: anyhow::Error) {
+ if let Some(handle) = self.header_handles.get(&id) {
+ fclose(handle);
}
- if is_resource(job.get("bodyHandle").unwrap_or(&PhpMixed::Null)) {
- fclose(job.get("bodyHandle").cloned().unwrap_or(PhpMixed::Null));
+ if let Some(handle) = self.body_handles.get(&id) {
+ fclose(handle);
}
+ self.header_handles.shift_remove(&id);
+ self.body_handles.shift_remove(&id);
if let Some(PhpMixed::String(filename)) = job.get("filename") {
unlink_silent(&format!("{}~", filename));
}