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