Drain on net thread.
| 420 | |
| 421 | // Drain on net thread. |
| 422 | void DrainOnNetThread() { |
| 423 | if (!MetadataSync::SchemaFetchEnabled()) return; |
| 424 | if (t_draining) return; |
| 425 | if (g_shuttingDown.load(std::memory_order_acquire)) return; |
| 426 | |
| 427 | // Only log/work when the queue actually has something, to avoid spamming |
| 428 | // every BYieldingSend tick. |
| 429 | { |
| 430 | std::lock_guard<std::mutex> lock(g_sendMutex); |
| 431 | if (g_sendQueue.empty()) return; |
| 432 | } |
| 433 | |
| 434 | if (!g_sessionCaptured.load(std::memory_order_relaxed)) { |
| 435 | LOG("[SchemaFetch] Drain: queue non-empty but session not captured, deferring"); |
| 436 | return; |
| 437 | } |
| 438 | |
| 439 | uint32_t conn = g_statsConnHandle.load(std::memory_order_relaxed); |
| 440 | if (conn == 0) conn = g_connHandle.load(std::memory_order_relaxed); |
| 441 | if (conn == 0) { |
| 442 | LOG("[SchemaFetch] Drain: queue non-empty but conn==0, deferring"); |
| 443 | return; |
| 444 | } |
| 445 | |
| 446 | void* cmInterface = g_cmInterface.load(std::memory_order_relaxed); |
| 447 | if (!cmInterface) { |
| 448 | LOG("[SchemaFetch] Drain: queue non-empty but cmInterface==null, deferring"); |
| 449 | return; |
| 450 | } |
| 451 | |
| 452 | size_t qsize; |
| 453 | { |
| 454 | std::lock_guard<std::mutex> lock(g_sendMutex); |
| 455 | qsize = g_sendQueue.size(); |
| 456 | } |
| 457 | LOG("[SchemaFetch] Drain: on net thread, conn=%u, %zu queued", conn, qsize); |
| 458 | |
| 459 | t_draining = true; |
| 460 | constexpr int kMaxPerTick = 2; |
| 461 | for (int i = 0; i < kMaxPerTick; ++i) { |
| 462 | SchemaSendItem item; |
| 463 | { |
| 464 | std::lock_guard<std::mutex> lock(g_sendMutex); |
| 465 | if (g_sendQueue.empty()) break; |
| 466 | item = g_sendQueue.front(); |
| 467 | g_sendQueue.pop(); |
| 468 | } |
| 469 | SendSchemaRequest(item.appId, item.owner, cmInterface, conn); |
| 470 | } |
| 471 | t_draining = false; |
| 472 | } |
| 473 | |
| 474 | // Inbound 819 capture (mirror of Windows TryHandleSchemaResponse). |
| 475 |
no test coverage detected