ensureModulesLoaded loads all pending modules (from compat parsing, the implicit CWD module, and -m flags). Called from serveQuery after ensureWorkspaceLoaded. Uses a mutex+flag instead of sync.Once so that transient failures (e.g. session not yet registered) can be retried.
(ctx context.Context, client *daggerClient)
| 655 | // ensureWorkspaceLoaded. Uses a mutex+flag instead of sync.Once so that |
| 656 | // transient failures (e.g. session not yet registered) can be retried. |
| 657 | func (srv *Server) ensureModulesLoaded(ctx context.Context, client *daggerClient) error { |
| 658 | if len(client.pendingModules) == 0 && len(client.pendingExtraModules) == 0 { |
| 659 | return nil |
| 660 | } |
| 661 | |
| 662 | client.modulesMu.Lock() |
| 663 | defer client.modulesMu.Unlock() |
| 664 | |
| 665 | if client.modulesLoaded { |
| 666 | return client.modulesErr |
| 667 | } |
| 668 | |
| 669 | // Wait for the client's session attachables to be available. |
| 670 | // Don't mark as loaded on failure — allow retry on next request. |
| 671 | if _, err := client.getClientCaller(ctx, client.clientID); err != nil { |
| 672 | return fmt.Errorf("waiting for client session attachables: %w", err) |
| 673 | } |
| 674 | |
| 675 | loads := gatherModuleLoadRequests(client.pendingModules, client.pendingExtraModules) |
| 676 | resolvedLoads := make([]resolvedModuleLoad, len(loads)) |
| 677 | resolveErrs := make([]error, len(loads)) |
| 678 | |
| 679 | // Resolve modules in parallel, then apply to client state in deterministic order. |
| 680 | jobs := parallel.New(). |
| 681 | WithContextualTracer(true). |
| 682 | WithLimit(moduleResolveParallelism(len(loads))) |
| 683 | for i, load := range loads { |
| 684 | i := i |
| 685 | load := load |
| 686 | jobs = jobs.WithJob(moduleLoadJobName(load), func(ctx context.Context) error { |
| 687 | resolved, err := srv.resolveModuleLoad(ctx, client.dag, load) |
| 688 | if err != nil { |
| 689 | resolveErrs[i] = err |
| 690 | return nil //nolint:nilerr // errors collected for deterministic ordering |
| 691 | } |
| 692 | resolvedLoads[i] = resolved |
| 693 | return nil |
| 694 | }) |
| 695 | } |
| 696 | if err := jobs.Run(ctx); err != nil { |
| 697 | client.modulesErr = fmt.Errorf("resolving modules: %w", err) |
| 698 | client.modulesLoaded = true |
| 699 | return client.modulesErr |
| 700 | } |
| 701 | |
| 702 | for i, load := range loads { |
| 703 | if resolveErrs[i] != nil { |
| 704 | client.modulesErr = moduleLoadErr(load, resolveErrs[i]) |
| 705 | client.modulesLoaded = true |
| 706 | return client.modulesErr |
| 707 | } |
| 708 | } |
| 709 | |
| 710 | client.stateMu.Lock() |
| 711 | defer client.stateMu.Unlock() |
| 712 | if err := srv.serveAllResolvedModuleLoads(client, loads, resolvedLoads); err != nil { |
| 713 | client.modulesErr = err |
| 714 | client.modulesLoaded = true |
no test coverage detected