aboutsummaryrefslogtreecommitdiffhomepage
path: root/crates/shirabe/src/util/http_downloader.rs
diff options
context:
space:
mode:
authornsfisis <nsfisis@gmail.com>2026-07-17 16:37:58 +0900
committernsfisis <nsfisis@gmail.com>2026-07-17 16:37:58 +0900
commitc5c527d87425ff93a4f983d36e12bd1c4b565e52 (patch)
tree47696883b4f7bf132f6dc6ebc1db777ab57addbc /crates/shirabe/src/util/http_downloader.rs
parent2dc55eec52ee2c6993e64d8bb11c96a4f2ec5c6d (diff)
downloadphp-shirabe-c5c527d87425ff93a4f983d36e12bd1c4b565e52.tar.gz
php-shirabe-c5c527d87425ff93a4f983d36e12bd1c4b565e52.tar.zst
php-shirabe-c5c527d87425ff93a4f983d36e12bd1c4b565e52.zip
refactor(curl-downloader): rewrite as a single async fn, drop Job/tick
Replaces the Job-table + tick()-driven polling loop with one async download() that sends, decides (retry/redirect/fail/succeed via a new decide() extracted from the former run_job), and loops until it resolves — no more resolve/reject callbacks. The client switches from reqwest::blocking::Client to the non-blocking reqwest::Client, with body streaming now via tokio::fs. Because real async I/O needs a live tokio reactor and none runs yet at the process level (sync_executor::block_on is a no-reactor busy-spin executor that only works when awaited futures resolve synchronously), HttpDownloader::start_job drives CurlDownloader::download() through a dedicated temporary current_thread Runtime (curl_runtime(), marked TODO(phase-e)) instead. This keeps concurrency characteristics unchanged for now — start_job still resolves one job at a time — real parallel I/O lands once HttpDownloader/Loop are rearchitected on top of FuturesUnordered. count_active_jobs' curl.tick() polling and the Job.settled/curl_id plumbing are removed as dead weight now that start_job settles curl jobs synchronously, same as the rfs path already did. abort_request is dropped: it had no caller (the PHP Promise-cancellation flow it backs was never ported), and the job table it operated on no longer exists. Verified manually against real network I/O (sandbox disabled): `shirabe show -a` (JSON metadata, in-memory body) and `shirabe create-project` (actual dist zip download + extraction) both complete correctly with no hang. Two unrelated pre-existing bugs surfaced during manual testing (an event-dispatcher subscriber wiring gap during `require`, and a RefCell reentrancy panic in `diagnose`) reproduce identically on the pre-change code and are out of scope here.
Diffstat (limited to 'crates/shirabe/src/util/http_downloader.rs')
-rw-r--r--crates/shirabe/src/util/http_downloader.rs106
1 files changed, 35 insertions, 71 deletions
diff --git a/crates/shirabe/src/util/http_downloader.rs b/crates/shirabe/src/util/http_downloader.rs
index 448c7019..c2c83b85 100644
--- a/crates/shirabe/src/util/http_downloader.rs
+++ b/crates/shirabe/src/util/http_downloader.rs
@@ -97,34 +97,15 @@ impl Default for HttpDownloaderMockHandler {
}
}
+#[derive(Debug)]
struct Job {
id: i64,
status: i64,
request: Request,
sync: bool,
origin: String,
- curl_id: Option<i64>,
response: Option<Response>,
exception: Option<anyhow::Error>,
- /// Completion slot written by the curl resolve/reject closures (driven by `curl.tick()`)
- /// and read by `count_active_jobs`. Uses `Arc<Mutex>` because `CurlDownloader::download`
- /// requires `Send + Sync` callbacks.
- settled: std::sync::Arc<std::sync::Mutex<Option<anyhow::Result<Response>>>>,
-}
-
-impl std::fmt::Debug for Job {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- f.debug_struct("Job")
- .field("id", &self.id)
- .field("status", &self.status)
- .field("request", &self.request)
- .field("sync", &self.sync)
- .field("origin", &self.origin)
- .field("curl_id", &self.curl_id)
- .field("response", &self.response)
- .field("exception", &self.exception)
- .finish()
- }
}
#[derive(Debug, Clone)]
@@ -134,6 +115,24 @@ struct Request {
copy_to: Option<String>,
}
+/// A single-threaded tokio Runtime used only to drive `CurlDownloader::download()` (which needs a
+/// real reactor now that it uses the non-blocking `reqwest::Client`, unlike `sync_executor::block_on`
+/// which assumes every awaited future resolves synchronously). `current_thread` is used because
+/// `block_on` (unlike `spawn`) has no `Send` bound, and `download()`'s future closes over
+/// `Rc<RefCell<...>>` handles that are not `Send`.
+///
+/// TODO(phase-e): remove this once `HttpDownloader::add`/`get` are driven by `Loop::wait`'s
+/// `FuturesUnordered` under a single top-level Runtime (see the async re-architecture design).
+fn curl_runtime() -> &'static tokio::runtime::Runtime {
+ static RT: std::sync::LazyLock<tokio::runtime::Runtime> = std::sync::LazyLock::new(|| {
+ tokio::runtime::Builder::new_current_thread()
+ .enable_all()
+ .build()
+ .expect("failed to build the temporary CurlDownloader bridge runtime")
+ });
+ &RT
+}
+
impl HttpDownloader {
const STATUS_QUEUED: i64 = 1;
const STATUS_STARTED: i64 = 2;
@@ -397,10 +396,8 @@ impl HttpDownloader {
request: request.clone(),
sync,
origin,
- curl_id: None,
response: None,
exception: None,
- settled: std::sync::Arc::new(std::sync::Mutex::new(None)),
};
let can_use_curl = self.can_use_curl(&job);
self.jobs.insert(id, job);
@@ -560,38 +557,20 @@ impl HttpDownloader {
return;
}
- // curl branch: register the request with the curl multi handle. Completion is delivered
- // by curl.tick() into the job's `settled` slot (read by count_active_jobs). The resolve
- // callback stores the Response, reject stores the error; this mirrors PHP's promise
- // resolve/reject firing during tick(). PHP catches any exception from download() and
- // rejects the job.
- let settled = self.jobs.get(&id).unwrap().settled.clone();
- let settled_for_reject = settled.clone();
- let resolve: Box<dyn Fn(Response) + Send + Sync> = Box::new(move |response: Response| {
- *settled.lock().unwrap() = Some(Ok(response));
- });
- let reject: Box<dyn Fn(anyhow::Error) + Send + Sync> =
- Box::new(move |error: anyhow::Error| {
- *settled_for_reject.lock().unwrap() = Some(Err(error));
- });
-
- let download_result = {
- let curl = self.curl.as_mut().unwrap();
- curl.download(resolve, reject, &origin, &url, options, copy_to.as_deref())
- };
- match download_result {
- Ok(Ok(curl_id)) => {
- if let Some(job) = self.jobs.get_mut(&id) {
- job.curl_id = Some(curl_id);
- }
- }
- Ok(Err(e)) => {
- self.settle_job(id, Err(e.into()));
- }
- Err(e) => {
- self.settle_job(id, Err(e));
+ // curl branch: `CurlDownloader::download` now runs the whole redirect/retry/auth state
+ // machine to completion itself and returns the settled result directly, so this drives
+ // it through the temporary bridge runtime instead of PHP's promise resolve/reject firing
+ // during tick(). PHP catches any exception from download() and rejects the job.
+ let result: anyhow::Result<Response> = {
+ let curl = self.curl.as_ref().unwrap();
+ match curl_runtime().block_on(curl.download(&origin, &url, options, copy_to.as_deref()))
+ {
+ Ok(Ok(response)) => Ok(response),
+ Ok(Err(transport_exception)) => Err(transport_exception.into()),
+ Err(e) => Err(e),
}
- }
+ };
+ self.settle_job(id, result);
}
fn mark_job_done(&mut self) {
@@ -637,24 +616,9 @@ impl HttpDownloader {
}
}
- if let Some(curl) = self.curl.as_mut() {
- curl.tick()?;
- }
-
- // Apply completions delivered by curl.tick() into each started job's `settled` slot.
- // This reproduces the effect of PHP's resolve/reject callbacks firing during tick().
- let started_ids: Vec<i64> = self
- .jobs
- .values()
- .filter(|j| j.status == Self::STATUS_STARTED)
- .map(|j| j.id)
- .collect();
- for id in started_ids {
- let settled = self.jobs.get(&id).unwrap().settled.lock().unwrap().take();
- if let Some(result) = settled {
- self.settle_job(id, result);
- }
- }
+ // Unlike the old tick()-driven curl path, `start_job` now settles curl jobs synchronously
+ // (via the temporary bridge runtime), so no job is ever left lingering in STATUS_STARTED
+ // by the time we get here — nothing left to poll or collect.
if let Some(index) = index {
return Ok(