1 //===--- TUScheduler.cpp -----------------------------------------*-C++-*-===//
2 //
3 //                     The LLVM Compiler Infrastructure
4 //
5 // This file is distributed under the University of Illinois Open Source
6 // License. See LICENSE.TXT for details.
7 //
8 //===----------------------------------------------------------------------===//
9 // For each file, managed by TUScheduler, we create a single ASTWorker that
10 // manages an AST for that file. All operations that modify or read the AST are
11 // run on a separate dedicated thread asynchronously in FIFO order.
12 //
13 // We start processing each update immediately after we receive it. If two or
14 // more updates come subsequently without reads in-between, we attempt to drop
15 // an older one to not waste time building the ASTs we don't need.
16 //
17 // The processing thread of the ASTWorker is also responsible for building the
18 // preamble. However, unlike AST, the same preamble can be read concurrently, so
19 // we run each of async preamble reads on its own thread.
20 //
21 // To limit the concurrent load that clangd produces we mantain a semaphore that
22 // keeps more than a fixed number of threads from running concurrently.
23 //
24 // Rationale for cancelling updates.
25 // LSP clients can send updates to clangd on each keystroke. Some files take
26 // significant time to parse (e.g. a few seconds) and clangd can get starved by
27 // the updates to those files. Therefore we try to process only the last update,
28 // if possible.
29 // Our current strategy to do that is the following:
30 // - For each update we immediately schedule rebuild of the AST.
31 // - Rebuild of the AST checks if it was cancelled before doing any actual work.
32 //   If it was, it does not do an actual rebuild, only reports llvm::None to the
33 //   callback
34 // - When adding an update, we cancel the last update in the queue if it didn't
35 //   have any reads.
36 // There is probably a optimal ways to do that. One approach we might take is
37 // the following:
38 // - For each update we remember the pending inputs, but delay rebuild of the
39 //   AST for some timeout.
40 // - If subsequent updates come before rebuild was started, we replace the
41 //   pending inputs and reset the timer.
42 // - If any reads of the AST are scheduled, we start building the AST
43 //   immediately.
44 
45 #include "TUScheduler.h"
46 #include "Logger.h"
47 #include "Trace.h"
48 #include "clang/Frontend/PCHContainerOperations.h"
49 #include "llvm/Support/Errc.h"
50 #include "llvm/Support/Path.h"
51 #include <memory>
52 #include <queue>
53 #include <thread>
54 
55 namespace clang {
56 namespace clangd {
57 namespace {
58 class ASTWorkerHandle;
59 
60 /// Owns one instance of the AST, schedules updates and reads of it.
61 /// Also responsible for building and providing access to the preamble.
62 /// Each ASTWorker processes the async requests sent to it on a separate
63 /// dedicated thread.
64 /// The ASTWorker that manages the AST is shared by both the processing thread
65 /// and the TUScheduler. The TUScheduler should discard an ASTWorker when
66 /// remove() is called, but its thread may be busy and we don't want to block.
67 /// So the workers are accessed via an ASTWorkerHandle. Destroying the handle
68 /// signals the worker to exit its run loop and gives up shared ownership of the
69 /// worker.
70 class ASTWorker {
71   friend class ASTWorkerHandle;
72   ASTWorker(llvm::StringRef File, Semaphore &Barrier, CppFile AST,
73             bool RunSync);
74 
75 public:
76   /// Create a new ASTWorker and return a handle to it.
77   /// The processing thread is spawned using \p Tasks. However, when \p Tasks
78   /// is null, all requests will be processed on the calling thread
79   /// synchronously instead. \p Barrier is acquired when processing each
80   /// request, it is be used to limit the number of actively running threads.
81   static ASTWorkerHandle Create(llvm::StringRef File, AsyncTaskRunner *Tasks,
82                                 Semaphore &Barrier, CppFile AST);
83   ~ASTWorker();
84 
85   void update(ParseInputs Inputs, WantDiagnostics,
86               UniqueFunction<void(std::vector<DiagWithFixIts>)> OnUpdated);
87   void runWithAST(llvm::StringRef Name,
88                   UniqueFunction<void(llvm::Expected<InputsAndAST>)> Action);
89   bool blockUntilIdle(Deadline Timeout) const;
90 
91   std::shared_ptr<const PreambleData> getPossiblyStalePreamble() const;
92   std::size_t getUsedBytes() const;
93 
94 private:
95   // Must be called exactly once on processing thread. Will return after
96   // stop() is called on a separate thread and all pending requests are
97   // processed.
98   void run();
99   /// Signal that run() should finish processing pending requests and exit.
100   void stop();
101   /// Adds a new task to the end of the request queue.
102   void startTask(llvm::StringRef Name, UniqueFunction<void()> Task,
103                  llvm::Optional<WantDiagnostics> UpdateType);
104   /// Should the first task in the queue be skipped instead of run?
105   bool shouldSkipHeadLocked() const;
106 
107   struct Request {
108     UniqueFunction<void()> Action;
109     std::string Name;
110     Context Ctx;
111     llvm::Optional<WantDiagnostics> UpdateType;
112   };
113 
114   std::string File;
115   const bool RunSync;
116   Semaphore &Barrier;
117   // AST and FileInputs are only accessed on the processing thread from run().
118   CppFile AST;
119   // Inputs, corresponding to the current state of AST.
120   ParseInputs FileInputs;
121   // Guards members used by both TUScheduler and the worker thread.
122   mutable std::mutex Mutex;
123   std::shared_ptr<const PreambleData> LastBuiltPreamble; /* GUARDED_BY(Mutex) */
124   // Result of getUsedBytes() after the last rebuild or read of AST.
125   std::size_t LastASTSize; /* GUARDED_BY(Mutex) */
126   // Set to true to signal run() to finish processing.
127   bool Done;                           /* GUARDED_BY(Mutex) */
128   std::deque<Request> Requests;        /* GUARDED_BY(Mutex) */
129   mutable std::condition_variable RequestsCV;
130 };
131 
132 /// A smart-pointer-like class that points to an active ASTWorker.
133 /// In destructor, signals to the underlying ASTWorker that no new requests will
134 /// be sent and the processing loop may exit (after running all pending
135 /// requests).
136 class ASTWorkerHandle {
137   friend class ASTWorker;
138   ASTWorkerHandle(std::shared_ptr<ASTWorker> Worker)
139       : Worker(std::move(Worker)) {
140     assert(this->Worker);
141   }
142 
143 public:
144   ASTWorkerHandle(const ASTWorkerHandle &) = delete;
145   ASTWorkerHandle &operator=(const ASTWorkerHandle &) = delete;
146   ASTWorkerHandle(ASTWorkerHandle &&) = default;
147   ASTWorkerHandle &operator=(ASTWorkerHandle &&) = default;
148 
149   ~ASTWorkerHandle() {
150     if (Worker)
151       Worker->stop();
152   }
153 
154   ASTWorker &operator*() {
155     assert(Worker && "Handle was moved from");
156     return *Worker;
157   }
158 
159   ASTWorker *operator->() {
160     assert(Worker && "Handle was moved from");
161     return Worker.get();
162   }
163 
164   /// Returns an owning reference to the underlying ASTWorker that can outlive
165   /// the ASTWorkerHandle. However, no new requests to an active ASTWorker can
166   /// be schedule via the returned reference, i.e. only reads of the preamble
167   /// are possible.
168   std::shared_ptr<const ASTWorker> lock() { return Worker; }
169 
170 private:
171   std::shared_ptr<ASTWorker> Worker;
172 };
173 
174 ASTWorkerHandle ASTWorker::Create(llvm::StringRef File, AsyncTaskRunner *Tasks,
175                                   Semaphore &Barrier, CppFile AST) {
176   std::shared_ptr<ASTWorker> Worker(
177       new ASTWorker(File, Barrier, std::move(AST), /*RunSync=*/!Tasks));
178   if (Tasks)
179     Tasks->runAsync("worker:" + llvm::sys::path::filename(File),
180                     [Worker]() { Worker->run(); });
181 
182   return ASTWorkerHandle(std::move(Worker));
183 }
184 
185 ASTWorker::ASTWorker(llvm::StringRef File, Semaphore &Barrier, CppFile AST,
186                      bool RunSync)
187     : File(File), RunSync(RunSync), Barrier(Barrier), AST(std::move(AST)),
188       Done(false) {
189   if (RunSync)
190     return;
191 }
192 
193 ASTWorker::~ASTWorker() {
194 #ifndef NDEBUG
195   std::lock_guard<std::mutex> Lock(Mutex);
196   assert(Done && "handle was not destroyed");
197   assert(Requests.empty() && "unprocessed requests when destroying ASTWorker");
198 #endif
199 }
200 
201 void ASTWorker::update(
202     ParseInputs Inputs, WantDiagnostics WantDiags,
203     UniqueFunction<void(std::vector<DiagWithFixIts>)> OnUpdated) {
204   auto Task = [=](decltype(OnUpdated) OnUpdated) mutable {
205     FileInputs = Inputs;
206     auto Diags = AST.rebuild(std::move(Inputs));
207 
208     {
209       std::lock_guard<std::mutex> Lock(Mutex);
210       if (AST.getPreamble())
211         LastBuiltPreamble = AST.getPreamble();
212       LastASTSize = AST.getUsedBytes();
213     }
214     // We want to report the diagnostics even if this update was cancelled.
215     // It seems more useful than making the clients wait indefinitely if they
216     // spam us with updates.
217     if (Diags && WantDiags != WantDiagnostics::No)
218       OnUpdated(std::move(*Diags));
219   };
220 
221   startTask("Update", Bind(Task, std::move(OnUpdated)), WantDiags);
222 }
223 
224 void ASTWorker::runWithAST(
225     llvm::StringRef Name,
226     UniqueFunction<void(llvm::Expected<InputsAndAST>)> Action) {
227   auto Task = [=](decltype(Action) Action) {
228     ParsedAST *ActualAST = AST.getAST();
229     if (!ActualAST) {
230       Action(llvm::make_error<llvm::StringError>("invalid AST",
231                                                  llvm::errc::invalid_argument));
232       return;
233     }
234     Action(InputsAndAST{FileInputs, *ActualAST});
235 
236     // Size of the AST might have changed after reads too, e.g. if some decls
237     // were deserialized from preamble.
238     std::lock_guard<std::mutex> Lock(Mutex);
239     LastASTSize = ActualAST->getUsedBytes();
240   };
241 
242   startTask(Name, Bind(Task, std::move(Action)),
243             /*UpdateType=*/llvm::None);
244 }
245 
246 std::shared_ptr<const PreambleData>
247 ASTWorker::getPossiblyStalePreamble() const {
248   std::lock_guard<std::mutex> Lock(Mutex);
249   return LastBuiltPreamble;
250 }
251 
252 std::size_t ASTWorker::getUsedBytes() const {
253   std::lock_guard<std::mutex> Lock(Mutex);
254   return LastASTSize;
255 }
256 
257 void ASTWorker::stop() {
258   {
259     std::lock_guard<std::mutex> Lock(Mutex);
260     assert(!Done && "stop() called twice");
261     Done = true;
262   }
263   RequestsCV.notify_all();
264 }
265 
266 void ASTWorker::startTask(llvm::StringRef Name, UniqueFunction<void()> Task,
267                           llvm::Optional<WantDiagnostics> UpdateType) {
268   if (RunSync) {
269     assert(!Done && "running a task after stop()");
270     trace::Span Tracer(Name + ":" + llvm::sys::path::filename(File));
271     Task();
272     return;
273   }
274 
275   {
276     std::lock_guard<std::mutex> Lock(Mutex);
277     assert(!Done && "running a task after stop()");
278     Requests.push_back(
279         {std::move(Task), Name, Context::current().clone(), UpdateType});
280   }
281   RequestsCV.notify_all();
282 }
283 
284 void ASTWorker::run() {
285   while (true) {
286     Request Req;
287     {
288       std::unique_lock<std::mutex> Lock(Mutex);
289       RequestsCV.wait(Lock, [&]() { return Done || !Requests.empty(); });
290       if (Requests.empty()) {
291         assert(Done);
292         return;
293       }
294       // Even when Done is true, we finish processing all pending requests
295       // before exiting the processing loop.
296 
297       while (shouldSkipHeadLocked())
298         Requests.pop_front();
299       assert(!Requests.empty() && "skipped the whole queue");
300       Req = std::move(Requests.front());
301       // Leave it on the queue for now, so waiters don't see an empty queue.
302     } // unlock Mutex
303 
304     {
305       std::lock_guard<Semaphore> BarrierLock(Barrier);
306       WithContext Guard(std::move(Req.Ctx));
307       trace::Span Tracer(Req.Name);
308       Req.Action();
309     }
310 
311     {
312       std::lock_guard<std::mutex> Lock(Mutex);
313       Requests.pop_front();
314     }
315     RequestsCV.notify_all();
316   }
317 }
318 
319 // Returns true if Requests.front() is a dead update that can be skipped.
320 bool ASTWorker::shouldSkipHeadLocked() const {
321   assert(!Requests.empty());
322   auto Next = Requests.begin();
323   auto UpdateType = Next->UpdateType;
324   if (!UpdateType) // Only skip updates.
325     return false;
326   ++Next;
327   // An update is live if its AST might still be read.
328   // That is, if it's not immediately followed by another update.
329   if (Next == Requests.end() || !Next->UpdateType)
330     return false;
331   // The other way an update can be live is if its diagnostics might be used.
332   switch (*UpdateType) {
333   case WantDiagnostics::Yes:
334     return false; // Always used.
335   case WantDiagnostics::No:
336     return true; // Always dead.
337   case WantDiagnostics::Auto:
338     // Used unless followed by an update that generates diagnostics.
339     for (; Next != Requests.end(); ++Next)
340       if (Next->UpdateType == WantDiagnostics::Yes ||
341           Next->UpdateType == WantDiagnostics::Auto)
342         return true; // Prefer later diagnostics.
343     return false;
344   }
345   llvm_unreachable("Unknown WantDiagnostics");
346 }
347 
348 bool ASTWorker::blockUntilIdle(Deadline Timeout) const {
349   std::unique_lock<std::mutex> Lock(Mutex);
350   return wait(Lock, RequestsCV, Timeout, [&] { return Requests.empty(); });
351 }
352 
353 } // namespace
354 
355 unsigned getDefaultAsyncThreadsCount() {
356   unsigned HardwareConcurrency = std::thread::hardware_concurrency();
357   // C++ standard says that hardware_concurrency()
358   // may return 0, fallback to 1 worker thread in
359   // that case.
360   if (HardwareConcurrency == 0)
361     return 1;
362   return HardwareConcurrency;
363 }
364 
365 struct TUScheduler::FileData {
366   /// Latest inputs, passed to TUScheduler::update().
367   ParseInputs Inputs;
368   ASTWorkerHandle Worker;
369 };
370 
371 TUScheduler::TUScheduler(unsigned AsyncThreadsCount,
372                          bool StorePreamblesInMemory,
373                          ASTParsedCallback ASTCallback)
374     : StorePreamblesInMemory(StorePreamblesInMemory),
375       PCHOps(std::make_shared<PCHContainerOperations>()),
376       ASTCallback(std::move(ASTCallback)), Barrier(AsyncThreadsCount) {
377   if (0 < AsyncThreadsCount) {
378     PreambleTasks.emplace();
379     WorkerThreads.emplace();
380   }
381 }
382 
383 TUScheduler::~TUScheduler() {
384   // Notify all workers that they need to stop.
385   Files.clear();
386 
387   // Wait for all in-flight tasks to finish.
388   if (PreambleTasks)
389     PreambleTasks->wait();
390   if (WorkerThreads)
391     WorkerThreads->wait();
392 }
393 
394 bool TUScheduler::blockUntilIdle(Deadline D) const {
395   for (auto &File : Files)
396     if (!File.getValue()->Worker->blockUntilIdle(D))
397       return false;
398   if (PreambleTasks)
399     if (!PreambleTasks->wait(D))
400       return false;
401   return true;
402 }
403 
404 void TUScheduler::update(
405     PathRef File, ParseInputs Inputs, WantDiagnostics WantDiags,
406     UniqueFunction<void(std::vector<DiagWithFixIts>)> OnUpdated) {
407   std::unique_ptr<FileData> &FD = Files[File];
408   if (!FD) {
409     // Create a new worker to process the AST-related tasks.
410     ASTWorkerHandle Worker = ASTWorker::Create(
411         File, WorkerThreads ? WorkerThreads.getPointer() : nullptr, Barrier,
412         CppFile(File, StorePreamblesInMemory, PCHOps, ASTCallback));
413     FD = std::unique_ptr<FileData>(new FileData{Inputs, std::move(Worker)});
414   } else {
415     FD->Inputs = Inputs;
416   }
417   FD->Worker->update(std::move(Inputs), WantDiags, std::move(OnUpdated));
418 }
419 
420 void TUScheduler::remove(PathRef File) {
421   bool Removed = Files.erase(File);
422   if (!Removed)
423     log("Trying to remove file from TUScheduler that is not tracked. File:" +
424         File);
425 }
426 
427 void TUScheduler::runWithAST(
428     llvm::StringRef Name, PathRef File,
429     UniqueFunction<void(llvm::Expected<InputsAndAST>)> Action) {
430   auto It = Files.find(File);
431   if (It == Files.end()) {
432     Action(llvm::make_error<llvm::StringError>(
433         "trying to get AST for non-added document",
434         llvm::errc::invalid_argument));
435     return;
436   }
437 
438   It->second->Worker->runWithAST(Name, std::move(Action));
439 }
440 
441 void TUScheduler::runWithPreamble(
442     llvm::StringRef Name, PathRef File,
443     UniqueFunction<void(llvm::Expected<InputsAndPreamble>)> Action) {
444   auto It = Files.find(File);
445   if (It == Files.end()) {
446     Action(llvm::make_error<llvm::StringError>(
447         "trying to get preamble for non-added document",
448         llvm::errc::invalid_argument));
449     return;
450   }
451 
452   if (!PreambleTasks) {
453     trace::Span Tracer(Name);
454     SPAN_ATTACH(Tracer, "file", File);
455     std::shared_ptr<const PreambleData> Preamble =
456         It->second->Worker->getPossiblyStalePreamble();
457     Action(InputsAndPreamble{It->second->Inputs, Preamble.get()});
458     return;
459   }
460 
461   ParseInputs InputsCopy = It->second->Inputs;
462   std::shared_ptr<const ASTWorker> Worker = It->second->Worker.lock();
463   auto Task = [InputsCopy, Worker, this](std::string Name, std::string File,
464                                          Context Ctx,
465                                          decltype(Action) Action) mutable {
466     std::lock_guard<Semaphore> BarrierLock(Barrier);
467     WithContext Guard(std::move(Ctx));
468     trace::Span Tracer(Name);
469     SPAN_ATTACH(Tracer, "file", File);
470     std::shared_ptr<const PreambleData> Preamble =
471         Worker->getPossiblyStalePreamble();
472     Action(InputsAndPreamble{InputsCopy, Preamble.get()});
473   };
474 
475   PreambleTasks->runAsync("task:" + llvm::sys::path::filename(File),
476                           Bind(Task, std::string(Name), std::string(File),
477                                Context::current().clone(), std::move(Action)));
478 }
479 
480 std::vector<std::pair<Path, std::size_t>>
481 TUScheduler::getUsedBytesPerFile() const {
482   std::vector<std::pair<Path, std::size_t>> Result;
483   Result.reserve(Files.size());
484   for (auto &&PathAndFile : Files)
485     Result.push_back(
486         {PathAndFile.first(), PathAndFile.second->Worker->getUsedBytes()});
487   return Result;
488 }
489 
490 } // namespace clangd
491 } // namespace clang
492