diff options
| author | nsfisis <nsfisis@gmail.com> | 2026-07-18 18:03:57 +0900 |
|---|---|---|
| committer | nsfisis <nsfisis@gmail.com> | 2026-07-18 18:03:57 +0900 |
| commit | 9b291d454b059977639db26c847af4e835cda0c3 (patch) | |
| tree | e54770570a80e67681f1614841156b864452713c /crates/shirabe/src/util/process_executor.rs | |
| parent | f8e385f7a1bd752d39f4d4f87b0439596839dde9 (diff) | |
| download | php-shirabe-9b291d454b059977639db26c847af4e835cda0c3.tar.gz php-shirabe-9b291d454b059977639db26c847af4e835cda0c3.tar.zst php-shirabe-9b291d454b059977639db26c847af4e835cda0c3.zip | |
refactor(downloader): take &self across the downloader hierarchy
Concurrent package operations call into the same downloader instances
through Rc<RefCell<dyn DownloaderInterface>>; with &mut self methods
every call holds a RefMut across its awaits, which panics with
'already mutably borrowed' the moment two operations overlap. This is
groundwork for fanning out InstallationManager's download/install
loops (same rework HttpDownloader/CurlDownloader already got).
- DownloaderInterface/ChangeReportInterface/ArchiveDownloader/
VcsDownloader methods now take &self; as_change_report_interface
returns &dyn instead of &mut dyn.
- Implementors move their genuinely mutable state behind cells:
FileDownloader.additional_cleanup_paths, the archive downloaders'
cleanup_executed, ZipDownloader.zip_archive_object,
VcsDownloaderBase.has_cleaned_changes, GitDownloader's stash/discard/
cache maps and GitUtil, SvnDownloader.cache_credentials,
PerforceDownloader.perforce. FileDownloader.io gains a RefCell layer
so get_local_changes can keep PHP's NullIO swap under &self.
- ProcessExecutor::execute_async now returns a future that captures
everything up front instead of borrowing the executor, and call
sites build the future before awaiting, so no borrow on the shared
executor is held while a subprocess runs.
- Filesystem::remove_directory_async becomes remove_directory_async_via
taking the Rc handle: the Filesystem is only borrowed for the sync
head/tail, never across the rm subprocess await (sync borrow_mut
users like rename/ensure_directory_exists would otherwise collide).
- DownloadManager async call sites hold shared borrows only.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Diffstat (limited to 'crates/shirabe/src/util/process_executor.rs')
| -rw-r--r-- | crates/shirabe/src/util/process_executor.rs | 136 |
1 files changed, 80 insertions, 56 deletions
diff --git a/crates/shirabe/src/util/process_executor.rs b/crates/shirabe/src/util/process_executor.rs index 767e6a00..ab4768a4 100644 --- a/crates/shirabe/src/util/process_executor.rs +++ b/crates/shirabe/src/util/process_executor.rs @@ -548,9 +548,16 @@ impl ProcessExecutor { /// starts a process on the commandline in async mode /// - /// `&self` so that concurrent calls through the same `Rc<RefCell<ProcessExecutor>>` can - /// coexist (shared borrows); the max_jobs throttle is enforced by the semaphore. - pub async fn execute_async<C>(&self, command: C, cwd: Option<&str>) -> anyhow::Result<Process> + /// Returns a future that does NOT borrow the executor: everything it needs is captured up + /// front, so callers can drop their `Ref`/`RefMut` on the shared `Rc<RefCell<ProcessExecutor>>` + /// before awaiting (`let fut = pe.borrow().execute_async(...); fut.await`). Holding a borrow + /// across the await would panic as soon as a sibling future or a sync `execute()` call touches + /// the same executor. The max_jobs throttle is enforced by the semaphore. + pub fn execute_async<C>( + &self, + command: C, + cwd: Option<&str>, + ) -> std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<Process>>>> where C: IntoExecCommand, { @@ -563,64 +570,70 @@ impl ProcessExecutor { // returning a misleading Process. todo!("ProcessExecutorMock async path needs a Process mock seam in external-packages"); } - if !self.allow_async { - return Err(LogicException { - message: "You must use the ProcessExecutor instance which is part of a Composer\\Loop instance to be able to run async processes".to_string(), - code: 0, + let allow_async = self.allow_async; + let semaphore = self.semaphore.clone(); + let io = self.io.clone(); + let cwd = cwd.map(ToOwned::to_owned); + + Box::pin(async move { + if !allow_async { + return Err(LogicException { + message: "You must use the ProcessExecutor instance which is part of a Composer\\Loop instance to be able to run async processes".to_string(), + code: 0, + } + .into()); } - .into()); - } - // PHP queues the job and only startJob()s it once runningJobs < maxJobs; the permit is the - // equivalent gate, so everything below (including the "Executing async command" debug - // line PHP prints from startJob) happens only once a slot is free. - let semaphore = self.semaphore.clone(); - let _permit = semaphore - .acquire() - .await - .expect("the semaphore is never closed"); + // PHP queues the job and only startJob()s it once runningJobs < maxJobs; the permit is + // the equivalent gate, so everything below (including the "Executing async command" + // debug line PHP prints from startJob) happens only once a slot is free. + let _permit = semaphore + .acquire() + .await + .expect("the semaphore is never closed"); - self.output_command_run(&command, cwd, true); + Self::output_command_run_with(&io, &command, cwd.as_deref(), true); - // PHP: $job['reject']($e) on process construction/start failure — surfaced as Err here. - let mut process = if is_string(&command) { - Process::from_shell_commandline( - command.as_string().unwrap_or(""), - cwd, - None, - PhpMixed::Null, - Some(Self::get_timeout() as f64), - )? - } else if let PhpMixed::List(ref list) = command { - Process::new( - list.iter() - .map(|v| v.as_string().unwrap_or("").to_string()) - .collect(), - cwd.map(ToOwned::to_owned), - None, - PhpMixed::Null, - Some(Self::get_timeout() as f64), - )? - } else { - return Err(LogicException { - message: "Invalid command type".to_string(), - code: 0, - } - .into()); - }; + // PHP: $job['reject']($e) on process construction/start failure — surfaced as Err here. + let mut process = if is_string(&command) { + Process::from_shell_commandline( + command.as_string().unwrap_or(""), + cwd.as_deref(), + None, + PhpMixed::Null, + Some(Self::get_timeout() as f64), + )? + } else if let PhpMixed::List(ref list) = command { + Process::new( + list.iter() + .map(|v| v.as_string().unwrap_or("").to_string()) + .collect(), + cwd.clone(), + None, + PhpMixed::Null, + Some(Self::get_timeout() as f64), + )? + } else { + return Err(LogicException { + message: "Invalid command type".to_string(), + code: 0, + } + .into()); + }; - process.start(None, IndexMap::new())?; + process.start(None, IndexMap::new())?; - // PHP's countActiveJobs tick: pump the process until it exits, checking the timeout each - // round. The async sleep yields to the reactor so sibling jobs genuinely overlap. - while process.is_running() { - process.check_timeout()?; - tokio::time::sleep(std::time::Duration::from_millis(1)).await; - } + // PHP's countActiveJobs tick: pump the process until it exits, checking the timeout + // each round. The async sleep yields to the reactor so sibling jobs genuinely overlap. + while process.is_running() { + process.check_timeout()?; + tokio::time::sleep(std::time::Duration::from_millis(1)).await; + } - // PHP resolves the promise with the Process regardless of its exit status; callers - // inspect is_successful() themselves. - Ok(process) + // PHP resolves the promise with the Process regardless of its exit status; callers + // inspect is_successful() themselves. + Ok(process) + }) } fn output_handler( @@ -712,7 +725,18 @@ impl ProcessExecutor { /// @param string|list<string> $command fn output_command_run(&self, command: &PhpMixed, cwd: Option<&str>, r#async: bool) { - if self.io.is_none() || !self.io.as_ref().unwrap().is_debug() { + Self::output_command_run_with(&self.io, command, cwd, r#async); + } + + /// `output_command_run` body as an associated fn so the `execute_async` future can carry a + /// clone of the io handle instead of borrowing the executor. + fn output_command_run_with( + io: &Option<std::rc::Rc<std::cell::RefCell<dyn IOInterface>>>, + command: &PhpMixed, + cwd: Option<&str>, + r#async: bool, + ) { + if io.is_none() || !io.as_ref().unwrap().is_debug() { return; } @@ -754,7 +778,7 @@ impl ProcessExecutor { "--password '***' ", &safe_command, ); - self.io.as_ref().unwrap().write_error(&format!( + io.as_ref().unwrap().write_error(&format!( "Executing{} command ({}): {}", if r#async { " async" } else { "" }, cwd.unwrap_or("CWD"), |
