1use std::{
2 collections::{HashMap, HashSet},
3 io,
4 sync::Arc,
5};
6
7use crate::{
8 build::{Build, BuildBehaviour, BuildError, RemotePackageSourceSpec, SrcRockSource},
9 config::Config,
10 lockfile::{
11 FlushLockfileError, LocalPackage, LocalPackageId, LockConstraint, Lockfile, OptState,
12 PinnedState, ReadOnly, ReadWrite,
13 },
14 lua_installation::{LuaInstallation, LuaInstallationError},
15 lua_rockspec::BuildBackendSpec,
16 lua_version::LuaVersionUnset,
17 luarocks::{
18 install_binary_rock::{BinaryRockInstall, InstallBinaryRockError},
19 luarocks_installation::{LuaRocksError, LuaRocksInstallError, LuaRocksInstallation},
20 },
21 operations::resolve::{
22 build_dependencies_to_install, PackageInstallData, Resolve, ResolveDependenciesError,
23 },
24 package::{PackageName, PackageNameList, PackageReq},
25 remote_package_db::{RemotePackageDB, RemotePackageDBError, RemotePackageDbIntegrityError},
26 rockspec::Rockspec,
27 tree::{self, InstallTree, Tree, TreeError},
28 workspace::{Workspace, WorkspaceTreeError},
29};
30
31pub use crate::operations::install::spec::PackageInstallSpec;
32
33use super::{DownloadedRockspec, RemoteRockDownload};
34use bon::Builder;
35use bytes::Bytes;
36use futures::stream::FuturesUnordered;
37use futures::StreamExt;
38use itertools::Itertools;
39use miette::Diagnostic;
40use thiserror::Error;
41use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
42use tokio::task::{JoinError, JoinHandle};
43
44use tracing::{info_span, span, Instrument};
45pub mod spec;
46
47#[derive(Builder)]
51#[builder(start_fn = new, finish_fn(name = _build, vis = ""))]
52pub struct Install<'a, T>
53where
54 T: InstallTree + Clone + Send + Sync,
55{
56 #[builder(start_fn)]
57 config: &'a Config,
58 #[builder(field)]
59 packages: Vec<PackageInstallSpec>,
60 #[builder(setters(name = "_tree", vis = ""))]
61 tree: T,
62 package_db: Option<RemotePackageDB>,
63}
64
65impl<'a, State> InstallBuilder<'a, Tree, State>
66where
67 State: install_builder::State,
68{
69 pub fn workspace(
70 self,
71 workspace: &'a Workspace,
72 ) -> Result<InstallBuilder<'a, Tree, install_builder::SetTree<State>>, WorkspaceTreeError>
73 where
74 State::Tree: install_builder::IsUnset,
75 {
76 let config = self.config;
77 Ok(self._tree(workspace.tree(config)?))
78 }
79}
80
81impl<'a, T, State> InstallBuilder<'a, T, State>
82where
83 State: install_builder::State,
84 T: InstallTree + Clone + Send + Sync,
85{
86 pub fn tree(self, tree: T) -> InstallBuilder<'a, T, install_builder::SetTree<State>>
87 where
88 State::Tree: install_builder::IsUnset,
89 {
90 self._tree(tree)
91 }
92
93 pub fn packages(self, packages: Vec<PackageInstallSpec>) -> Self {
94 Self { packages, ..self }
95 }
96
97 pub fn package(self, package: PackageInstallSpec) -> Self {
98 Self {
99 packages: self
100 .packages
101 .into_iter()
102 .chain(std::iter::once(package))
103 .collect(),
104 ..self
105 }
106 }
107}
108
109impl<State, T> InstallBuilder<'_, T, State>
110where
111 State: install_builder::State + install_builder::IsComplete,
112 T: InstallTree + Clone + Send + Sync + 'static,
113{
114 pub async fn install(self) -> Result<Vec<LocalPackage>, InstallError> {
116 let install_built = self._build();
117 if install_built.packages.is_empty() {
118 return Ok(Vec::default());
119 }
120 let count = install_built.packages.len();
121 let span = if count > 1 {
122 info_span!("Installing", count,)
123 } else {
124 let install_spec = &install_built.packages[0];
125 info_span!("Installing", package = install_spec.package.to_string(),)
126 };
127 let _enter = span.enter();
128 let package_db = match install_built.package_db {
129 Some(db) => db,
130 None => RemotePackageDB::from_config(install_built.config).await?,
131 };
132
133 let duplicate_entrypoints = install_built
134 .packages
135 .iter()
136 .filter(|pkg| pkg.entry_type == tree::EntryType::Entrypoint)
137 .map(|pkg| pkg.package.name())
138 .duplicates()
139 .cloned()
140 .collect_vec();
141
142 if !duplicate_entrypoints.is_empty() {
143 return Err(InstallError::DuplicateEntrypoints(PackageNameList::new(
144 duplicate_entrypoints,
145 )));
146 }
147
148 install_impl(
149 install_built.packages,
150 Arc::new(package_db),
151 install_built.config,
152 &install_built.tree,
153 )
154 .await
155 }
156}
157
158type InstallWorkerOutput = Result<(LocalPackageId, (LocalPackage, tree::EntryType)), InstallError>;
159
160#[derive(Error, Debug, Diagnostic)]
161pub enum InstallError {
162 #[error("unable to resolve dependencies:\n{0}")]
163 #[diagnostic(forward(0))]
164 ResolveDependencies(#[from] ResolveDependenciesError),
165 #[error(transparent)]
166 #[diagnostic(transparent)]
167 LuaVersionUnset(#[from] LuaVersionUnset),
168 #[error(transparent)]
169 #[diagnostic(transparent)]
170 LuaInstallation(#[from] LuaInstallationError),
171 #[error(transparent)]
172 #[diagnostic(transparent)]
173 FlushLockfile(#[from] FlushLockfileError),
174 #[error(transparent)]
175 #[diagnostic(transparent)]
176 Tree(#[from] TreeError),
177 #[error(transparent)]
178 #[diagnostic(transparent)]
179 WorkspaceTree(#[from] WorkspaceTreeError),
180 #[error("error instantiating LuaRocks compatibility layer:\n{0}")]
181 #[diagnostic(forward(0))]
182 LuaRocks(#[from] LuaRocksError),
183 #[error("error installing LuaRocks compatibility layer:\n{0}")]
184 #[diagnostic(forward(0))]
185 LuaRocksInstall(#[from] LuaRocksInstallError),
186 #[error("failed to build {0}: {1}")]
187 Build(PackageName, BuildError),
188 #[error("failed to install build depencency {0}:\n{1}")]
189 BuildDependency(PackageName, BuildError),
190 #[error("error initialising remote package DB:\n{0}")]
191 #[diagnostic(forward(0))]
192 RemotePackageDB(#[from] RemotePackageDBError),
193 #[error("failed to install pre-built rock {0}:\n{1}")]
194 InstallBinaryRock(PackageName, InstallBinaryRockError),
195 #[error("integrity error for package '{package}'")]
196 Integrity {
197 package: PackageName,
198 #[diagnostic_source]
199 err: RemotePackageDbIntegrityError,
200 },
201 #[error("cannot install duplicate entrypoints:\n{0}")]
202 DuplicateEntrypoints(PackageNameList),
203 #[error("install worker panicked")]
204 #[diagnostic(help(
205 r#"this is a bug in Lux, please report it, ideally with `RUST_BACKTRACE=1`.
206retrying with fewer parallel jobs (`--max-jobs`) may avoid the panic in the meantime"#
207 ))]
208 Join(#[from] tokio::task::JoinError),
209}
210
211async fn install_impl<T>(
212 packages: Vec<PackageInstallSpec>,
213 package_db: Arc<RemotePackageDB>,
214 config: &Config,
215 tree: &T,
216) -> Result<Vec<LocalPackage>, InstallError>
217where
218 T: InstallTree + Clone + Send + Sync + 'static,
219{
220 let (dep_tx, mut dep_rx) = tokio::sync::mpsc::unbounded_channel();
221 let (build_dep_tx, build_dep_rx) = tokio::sync::mpsc::unbounded_channel();
222 let (build_dep_install_done_tx, mut build_dep_install_done_rx) =
223 tokio::sync::mpsc::unbounded_channel::<PackageName>();
224
225 let lockfile = tree.lockfile()?;
226 let build_lockfile = tree.build_tree(config)?.lockfile()?;
227
228 let lua = Arc::new(LuaInstallation::new_from_config(config).await?);
229
230 let mut resolve_worker = spawn_resolve_worker(
231 config,
232 packages,
233 package_db,
234 lockfile.clone(),
235 build_lockfile.clone(),
236 dep_tx,
237 build_dep_tx,
238 );
239 let mut build_deps_worker = spawn_build_deps_worker(
240 config,
241 tree,
242 lua.clone(),
243 build_dep_rx,
244 build_dep_install_done_tx,
245 );
246
247 let mut all_packages: HashMap<LocalPackageId, PackageInstallData> = HashMap::new();
248 let mut scheduled_packages: HashSet<LocalPackageId> = HashSet::new();
249 let mut installed_packages: HashMap<LocalPackageId, (LocalPackage, tree::EntryType)> =
250 HashMap::new();
251 let mut installed_build_deps: HashSet<PackageName> = HashSet::new();
252 let mut ongoing_installs: FuturesUnordered<
253 tracing::instrument::Instrumented<tokio::task::JoinHandle<InstallWorkerOutput>>,
254 > = FuturesUnordered::new();
255 let mut resolve_done = false;
256 let mut build_deps_done = false;
257 let mut dep_rx_drained = false;
258 let mut build_dep_rx_drained = false;
259 let mut install_loop_result: Result<(), InstallError> = Ok(());
260 let max_jobs = config.max_jobs();
261
262 'install: loop {
263 if resolve_done && build_deps_done && dep_rx_drained && build_dep_rx_drained {
264 break;
265 }
266 tokio::select! {
267 resolve_result = &mut resolve_worker, if !resolve_done => {
268 if let Err(err) = worker_result(resolve_result) {
269 install_loop_result = Err(*err);
270 break 'install;
271 }
272 resolve_done = true;
273 }
274 build_deps_result = &mut build_deps_worker, if !build_deps_done => {
275 if let Err(err) = worker_result(build_deps_result) {
276 install_loop_result = Err(*err);
277 break 'install;
278 }
279 build_deps_done = true;
280 }
281 name = build_dep_install_done_rx.recv(), if !build_dep_rx_drained => {
282 if let Some(name) = name {
283 installed_build_deps.insert(name);
284 } else {
285 build_dep_rx_drained = true;
286 }
287 }
288 dep = dep_rx.recv(), if !dep_rx_drained => {
289 if let Some(dep) = dep {
290 all_packages.insert(dep.spec.id(), dep);
291 } else {
292 dep_rx_drained = true;
293 }
294 }
295 }
296
297 for (package_id, package_install_data) in ready_to_install(
298 &all_packages,
299 &scheduled_packages,
300 &build_lockfile,
301 &installed_build_deps,
302 ) {
303 if max_jobs > 0 && ongoing_installs.len() >= max_jobs {
304 if let Err(err) =
305 wait_for_next_install(&mut ongoing_installs, &mut installed_packages).await
306 {
307 install_loop_result = Err(err);
308 break 'install;
309 }
310 }
311 scheduled_packages.insert(package_id);
312 ongoing_installs.push(spawn_install_worker(
313 package_install_data,
314 &lua,
315 tree,
316 config,
317 ));
318 }
319 }
320
321 match install_loop_result {
322 Ok(_) => {
323 while wait_for_next_install(&mut ongoing_installs, &mut installed_packages).await? {}
324
325 lockfile.map_then_flush(|lockfile| {
326 for (package_id, (package, is_entrypoint)) in installed_packages.iter() {
327 lockfile.add_dependencies(
328 package_id,
329 package,
330 *is_entrypoint,
331 &all_packages,
332 &installed_packages,
333 )?;
334 }
335 Ok::<_, io::Error>(())
336 })?;
337
338 Ok(installed_packages
339 .into_values()
340 .map(|(pkg, _)| pkg)
341 .collect_vec())
342 }
343 Err(err) => {
344 resolve_worker.into_inner().abort();
345 build_deps_worker.into_inner().abort();
346 for install in ongoing_installs {
347 install.into_inner().abort();
348 }
349 Err(err)
350 }
351 }
352}
353
354fn spawn_resolve_worker(
355 config: &Config,
356 packages: Vec<PackageInstallSpec>,
357 package_db: Arc<RemotePackageDB>,
358 lockfile: Lockfile<ReadOnly>,
359 build_lockfile: Lockfile<ReadOnly>,
360 dep_tx: UnboundedSender<PackageInstallData>,
361 build_dep_tx: UnboundedSender<PackageInstallData>,
362) -> tracing::instrument::Instrumented<JoinHandle<Result<(), InstallError>>> {
363 tokio::spawn({
364 let config = config.clone();
365 let lockfile = Arc::new(lockfile);
366 let build_lockfile = Arc::new(build_lockfile);
367 async move {
368 Resolve::new()
369 .dependencies_tx(dep_tx)
370 .build_dependencies_tx(build_dep_tx)
371 .packages(packages)
372 .package_db(package_db)
373 .lockfile(lockfile)
374 .build_lockfile(build_lockfile)
375 .config(&config)
376 .get_all_dependencies()
377 .await?;
378 Ok::<(), InstallError>(())
379 }
380 })
381 .instrument(tracing::trace_span!("resolve_worker"))
382}
383
384fn spawn_build_deps_worker<T>(
385 config: &Config,
386 tree: &T,
387 lua: Arc<LuaInstallation>,
388 mut build_dep_rx: UnboundedReceiver<PackageInstallData>,
389 build_dep_install_done_tx: UnboundedSender<PackageName>,
390) -> tracing::instrument::Instrumented<JoinHandle<Result<(), InstallError>>>
391where
392 T: InstallTree + Clone + Send + Sync + 'static,
393{
394 tokio::spawn({
395 let config = config.clone();
396 let tree = tree.clone();
397 let lua = lua.clone();
398 async move {
399 while let Some(build_dep_spec) = build_dep_rx.recv().await {
400 let rockspec = build_dep_spec.downloaded_rock.rockspec();
401 let package = rockspec.package().clone();
402 let span = info_span!(
403 "Installing build dependency",
404 package = package.to_string(),
405 version = rockspec.version().to_string()
406 );
407 async {
408 let build_tree = tree.build_tree(&config)?;
409 let mut build_lockfile = build_tree.lockfile()?.write_guard();
410 let pkg = Build::new()
411 .rockspec(rockspec)
412 .lua(&lua)
413 .tree(&build_tree)
414 .entry_type(tree::EntryType::Entrypoint)
415 .config(&config)
416 .constraint(build_dep_spec.spec.constraint())
417 .behaviour(build_dep_spec.build_behaviour)
418 .build()
419 .await
420 .map_err(|err| InstallError::BuildDependency(package.clone(), err))?;
421 build_lockfile.add_entrypoint(&pkg);
422 Ok::<_, InstallError>(())
423 }
424 .instrument(span)
425 .await?;
426 let _ = build_dep_install_done_tx.send(package);
427 }
428 Ok::<(), InstallError>(())
429 }
430 })
431 .instrument(tracing::trace_span!("build_deps_worker"))
432}
433
434fn ready_to_install(
435 all_packages: &HashMap<LocalPackageId, PackageInstallData>,
436 scheduled: &HashSet<LocalPackageId>,
437 build_lockfile: &Lockfile<ReadOnly>,
438 installed_build_deps: &HashSet<PackageName>,
439) -> Vec<(LocalPackageId, PackageInstallData)> {
440 all_packages
441 .iter()
442 .filter(|(id, data)| {
443 !scheduled.contains(*id)
444 && match &data.downloaded_rock {
445 RemoteRockDownload::BinaryRock { .. } => true,
446 _ => build_dependencies_ready(
447 data.downloaded_rock.rockspec(),
448 data.build_behaviour,
449 build_lockfile,
450 installed_build_deps,
451 ),
452 }
453 })
454 .map(|(id, data)| (id.clone(), data.clone()))
455 .collect()
456}
457
458fn spawn_install_worker<T>(
459 data: PackageInstallData,
460 lua: &Arc<LuaInstallation>,
461 tree: &T,
462 config: &Config,
463) -> tracing::instrument::Instrumented<JoinHandle<InstallWorkerOutput>>
464where
465 T: InstallTree + Clone + Send + Sync + 'static,
466{
467 let config = config.clone();
468 let tree = tree.clone();
469 let lua = lua.clone();
470 let entry_type = data.entry_type;
471 tokio::spawn(async move {
472 let pkg = install_package(data, &lua, &tree, &config).await?;
473 Ok::<_, InstallError>((pkg.id(), (pkg, entry_type)))
474 })
475 .instrument(tracing::trace_span!("install_worker"))
476}
477
478#[tracing::instrument(level = "trace", skip_all)]
479async fn install_package<T>(
480 data: PackageInstallData,
481 lua: &Arc<LuaInstallation>,
482 tree: &T,
483 config: &Config,
484) -> Result<LocalPackage, InstallError>
485where
486 T: InstallTree + Sync,
487{
488 match data.downloaded_rock {
489 RemoteRockDownload::RockspecOnly { rockspec_download } => {
490 install_rockspec(
491 rockspec_download,
492 None,
493 data.spec.constraint(),
494 data.build_behaviour,
495 data.pin,
496 data.opt,
497 data.entry_type,
498 lua,
499 tree,
500 config,
501 )
502 .await
503 }
504 RemoteRockDownload::BinaryRock {
505 rockspec_download,
506 packed_rock,
507 } => {
508 install_binary_rock(
509 rockspec_download,
510 packed_rock,
511 data.spec.constraint(),
512 data.build_behaviour,
513 data.pin,
514 data.opt,
515 data.entry_type,
516 config,
517 tree,
518 )
519 .await
520 }
521 RemoteRockDownload::SrcRock {
522 rockspec_download,
523 src_rock,
524 source_url,
525 } => {
526 let src_rock_source = SrcRockSource {
527 bytes: src_rock,
528 source_url,
529 };
530 install_rockspec(
531 rockspec_download,
532 Some(src_rock_source),
533 data.spec.constraint(),
534 data.build_behaviour,
535 data.pin,
536 data.opt,
537 data.entry_type,
538 lua,
539 tree,
540 config,
541 )
542 .await
543 }
544 }
545}
546
547fn worker_result(
548 result: Result<Result<(), InstallError>, JoinError>,
549) -> Result<(), Box<InstallError>> {
550 match result {
551 Ok(Ok(())) => Ok(()),
552 Ok(Err(err)) => Err(err.into()),
553 Err(join) => Err(InstallError::from(join).into()),
554 }
555}
556
557async fn wait_for_next_install(
558 ongoing_installs: &mut FuturesUnordered<
559 tracing::instrument::Instrumented<tokio::task::JoinHandle<InstallWorkerOutput>>,
560 >,
561 installed_packages: &mut HashMap<LocalPackageId, (LocalPackage, tree::EntryType)>,
562) -> Result<bool, InstallError> {
563 if let Some(result) = ongoing_installs.next().await {
564 match result {
565 Ok(Ok((id, installed))) => {
566 installed_packages.insert(id, installed);
567 Ok(true)
568 }
569 Ok(Err(err)) => Err(err),
570 Err(join) => Err(InstallError::from(join)),
571 }
572 } else {
573 Ok(false)
574 }
575}
576trait LockfileExt {
577 fn add_dependencies(
578 self,
579 id: &LocalPackageId,
580 pkg: &LocalPackage,
581 entry_type: tree::EntryType,
582 all_packages: &HashMap<LocalPackageId, PackageInstallData>,
583 installed_packages: &HashMap<LocalPackageId, (LocalPackage, tree::EntryType)>,
584 ) -> io::Result<()>;
585}
586
587impl LockfileExt for &mut Lockfile<ReadWrite> {
588 fn add_dependencies(
589 self,
590 id: &LocalPackageId,
591 pkg: &LocalPackage,
592 entry_type: tree::EntryType,
593 all_packages: &HashMap<LocalPackageId, PackageInstallData>,
594 installed_packages: &HashMap<LocalPackageId, (LocalPackage, tree::EntryType)>,
595 ) -> io::Result<()> {
596 if entry_type == tree::EntryType::Entrypoint {
597 self.add_entrypoint(pkg);
598 }
599
600 for dependency_id in all_packages
601 .get(id)
602 .map(|pkg| pkg.spec.dependencies())
603 .unwrap_or_default()
604 .into_iter()
605 {
606 self.add_dependency(
607 pkg,
608 installed_packages
609 .get(dependency_id)
610 .map(|(pkg, _)| pkg)
611 .ok_or(io::Error::other(
612 r#"
613error writing dependencies to the lockfile.
614A required dependency was not installed correctly.
615This is likely because an install thread panicked and was interrupted unexpectedly.
616
617[THIS IS A BUG!]
618"#,
619 ))?,
620 );
621 }
622 Ok(())
623 }
624}
625
626fn build_dependencies_ready(
633 rockspec: &impl Rockspec,
634 behaviour: BuildBehaviour,
635 build_lockfile: &Lockfile<ReadOnly>,
636 installed: &HashSet<PackageName>,
637) -> bool {
638 let build_deps = rockspec.build_dependencies().current_platform();
639 build_dependencies_to_install(rockspec).iter().all(|name| {
640 installed.contains(name)
641 || (behaviour != BuildBehaviour::Force
642 && build_lockfile
643 .has_rock(
644 &build_deps
645 .iter()
646 .find(|dep| dep.name() == name)
647 .map(|dep| dep.package_req().clone())
648 .unwrap_or_else(|| PackageReq::from(name.clone())),
649 None,
650 )
651 .is_some())
652 })
653}
654
655#[allow(clippy::too_many_arguments)]
656async fn install_rockspec<T>(
657 rockspec_download: DownloadedRockspec,
658 src_rock_source: Option<SrcRockSource>,
659 constraint: LockConstraint,
660 behaviour: BuildBehaviour,
661 pin: PinnedState,
662 opt: OptState,
663 entry_type: tree::EntryType,
664 lua: &LuaInstallation,
665 tree: &T,
666 config: &Config,
667) -> Result<LocalPackage, InstallError>
668where
669 T: InstallTree + Sync,
670{
671 let package = rockspec_download.rockspec.package().clone();
672 let rockspec = rockspec_download.rockspec;
673 let span = info_span!(
674 "Installing",
675 package = package.to_string(),
676 version = rockspec.version().to_string(),
677 );
678 let _enter = span.enter();
679 let source = rockspec_download.source;
680
681 if let Some(BuildBackendSpec::LuaRock(_)) = &rockspec.build().current_platform().build_backend {
682 let luarocks_tree = tree.build_tree(config)?;
683 let luarocks = LuaRocksInstallation::new(config, luarocks_tree)?;
684 luarocks.ensure_installed(lua).await?;
685 }
686
687 let source_spec = match src_rock_source {
688 Some(src_rock_source) => RemotePackageSourceSpec::SrcRock(src_rock_source),
689 None => RemotePackageSourceSpec::RockSpec(rockspec_download.source_url),
690 };
691
692 let pkg = Build::new()
693 .rockspec(&rockspec)
694 .lua(lua)
695 .tree(tree)
696 .entry_type(entry_type)
697 .config(config)
698 .pin(pin)
699 .opt(opt)
700 .constraint(constraint)
701 .behaviour(behaviour)
702 .source(source)
703 .source_spec(source_spec)
704 .build()
705 .await
706 .map_err(|err| InstallError::Build(package, err))?;
707 Ok(pkg)
708}
709
710#[allow(clippy::too_many_arguments)]
711async fn install_binary_rock(
712 rockspec_download: DownloadedRockspec,
713 packed_rock: Bytes,
714 constraint: LockConstraint,
715 behaviour: BuildBehaviour,
716 pin: PinnedState,
717 opt: OptState,
718 entry_type: tree::EntryType,
719 config: &Config,
720 tree: &impl InstallTree,
721) -> Result<LocalPackage, InstallError> {
722 let rockspec = rockspec_download.rockspec;
723 let package = rockspec.package().clone();
724 let span = span!(
725 tracing::Level::INFO,
726 "Installing (pre-built)",
727 package = package.to_string(),
728 version = rockspec.version().to_string(),
729 );
730 let _enter = span.enter();
731 let pkg = BinaryRockInstall::new(
732 &rockspec,
733 rockspec_download.source,
734 packed_rock,
735 entry_type,
736 config,
737 tree,
738 )
739 .pin(pin)
740 .opt(opt)
741 .constraint(constraint)
742 .behaviour(behaviour)
743 .install()
744 .await
745 .map_err(|err| InstallError::InstallBinaryRock(package, err))?;
746 Ok(pkg)
747}