|
|
@@ -15,7 +15,10 @@
|
|
|
#include <nlohmann/json.hpp>
|
|
|
|
|
|
#include <algorithm>
|
|
|
+#include <chrono>
|
|
|
#include <functional>
|
|
|
+#include <memory>
|
|
|
+#include <mutex>
|
|
|
#include <queue>
|
|
|
#include <utility>
|
|
|
#include <vector>
|
|
|
@@ -2430,13 +2433,45 @@ grpc::Status DatabaseGrpcImpl::Subscribe(
|
|
|
std::vector<std::string> patterns(request->patterns().begin(), request->patterns().end());
|
|
|
bool includeData = request->include_data();
|
|
|
|
|
|
+ // v2.11.2 — WRITER LIFETIME. `writer` and `context` belong to this call and
|
|
|
+ // die when this function returns, but the event callback runs on whatever
|
|
|
+ // thread published the event: EventManager::publish() copies the matching
|
|
|
+ // callbacks out under its lock and then invokes them OUTSIDE it, so
|
|
|
+ // unsubscribing does not stop a dispatch already in progress. A callback
|
|
|
+ // could therefore still be writing through `writer` after this handler had
|
|
|
+ // returned and gRPC had freed it.
|
|
|
+ //
|
|
|
+ // That window was always there (a client disconnecting while an event was
|
|
|
+ // being published) but it stayed narrow. Making every subscriber unwind at
|
|
|
+ // once on shutdown - the point of this change - is exactly the situation
|
|
|
+ // that widens it, so the lifetime has to be made safe first.
|
|
|
+ //
|
|
|
+ // This shared state closes it in both directions: the callback takes the
|
|
|
+ // mutex and gives up if the stream is closed, and the handler closes it
|
|
|
+ // under the same mutex, so on return either the write finished or it never
|
|
|
+ // started. The mutex also serialises concurrent writes from several
|
|
|
+ // publishing threads, which grpc::ServerWriter does not permit and which
|
|
|
+ // nothing prevented before.
|
|
|
+ struct StreamState {
|
|
|
+ std::mutex mu;
|
|
|
+ bool open = true;
|
|
|
+ grpc::ServerWriter<pb::DatabaseEvent>* writer = nullptr;
|
|
|
+ };
|
|
|
+ auto state = std::make_shared<StreamState>();
|
|
|
+ state->writer = writer;
|
|
|
+
|
|
|
// v2.8.0 — Subscribe is gated PER EVENT, not once at entry. An empty
|
|
|
// `collections` list means "every collection" and `patterns` accepts globs,
|
|
|
// so there is no single name to authorise up front. Filtering in the
|
|
|
// callback is also what makes a partial grant work: a principal subscribed
|
|
|
// to everything receives only the collections it may read.
|
|
|
uint64_t subId = events_.subscribe(collections, patterns,
|
|
|
- [this, writer, includeData, context](const DatabaseEvent& event) {
|
|
|
+ [this, state, includeData, context](const DatabaseEvent& event) {
|
|
|
+ // Held for the whole callback: the writer must not be touched once the
|
|
|
+ // handler has closed the stream, and two publishers must not write
|
|
|
+ // concurrently.
|
|
|
+ std::lock_guard<std::mutex> streamLock(state->mu);
|
|
|
+ if (!state->open) return;
|
|
|
if (context->IsCancelled()) return;
|
|
|
|
|
|
smartbotic::database::Decision edec;
|
|
|
@@ -2467,19 +2502,68 @@ grpc::Status DatabaseGrpcImpl::Subscribe(
|
|
|
}
|
|
|
}
|
|
|
|
|
|
- writer->Write(protoEvent);
|
|
|
+ state->writer->Write(protoEvent);
|
|
|
});
|
|
|
|
|
|
- // Block until cancelled
|
|
|
- while (!context->IsCancelled()) {
|
|
|
- std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
|
|
+ // v2.11.2 — park until the client goes away OR the service starts stopping.
|
|
|
+ //
|
|
|
+ // This used to be a 100 ms sleep loop over context->IsCancelled() alone,
|
|
|
+ // which meant a stopping server was noticed only once the shutdown deadline
|
|
|
+ // expired and gRPC cancelled the call - so every stop waited out the full
|
|
|
+ // grace period with nothing to show for it. Waiting on the condition
|
|
|
+ // variable returns the moment beginShutdown() fires; the 100 ms timeout
|
|
|
+ // remains because the sync API offers no way to be woken when a client
|
|
|
+ // disconnects, so that half still has to be polled.
|
|
|
+ bool stopping = false;
|
|
|
+ {
|
|
|
+ std::unique_lock<std::mutex> lock(shutdownMutex_);
|
|
|
+ while (!context->IsCancelled()) {
|
|
|
+ if (shutdownCv_.wait_for(lock, std::chrono::milliseconds(100), [this] {
|
|
|
+ return shuttingDown_.load(std::memory_order_acquire);
|
|
|
+ })) {
|
|
|
+ stopping = true;
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ }
|
|
|
}
|
|
|
|
|
|
+ // Order matters. unsubscribe() first so no NEW dispatch picks up this
|
|
|
+ // callback, then close the stream under its own mutex so any dispatch
|
|
|
+ // already running has either finished or will see a closed stream. Only
|
|
|
+ // after both is it safe to return and let gRPC free `writer`.
|
|
|
events_.unsubscribe(subId);
|
|
|
-
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> streamLock(state->mu);
|
|
|
+ state->open = false;
|
|
|
+ }
|
|
|
+
|
|
|
+ // A client that sees OK would reasonably conclude its subscription ended
|
|
|
+ // normally. It did not - the server is going away - and UNAVAILABLE is the
|
|
|
+ // status a gRPC client already treats as "retry against a new connection".
|
|
|
+ if (stopping) {
|
|
|
+ // Logged at INFO deliberately. It is bounded (at most
|
|
|
+ // maxSubscribeStreams_ lines, 50 by default, once per process
|
|
|
+ // lifetime), it tells an operator why consumers saw UNAVAILABLE, and it
|
|
|
+ // is the only externally visible evidence that the stream unwound on
|
|
|
+ // the shutdown SIGNAL rather than by waiting out the Shutdown()
|
|
|
+ // deadline - which is what tests/load_test/test_shutdown_with_subscriber.sh
|
|
|
+ // measures.
|
|
|
+ spdlog::info("Subscribe: closing stream, server is shutting down");
|
|
|
+ return grpc::Status(grpc::StatusCode::UNAVAILABLE, "server is shutting down");
|
|
|
+ }
|
|
|
return grpc::Status::OK;
|
|
|
}
|
|
|
|
|
|
+void DatabaseGrpcImpl::beginShutdown() {
|
|
|
+ spdlog::info("Signalling shutdown to {} active Subscribe stream(s)",
|
|
|
+ activeSubscribeStreams_.load());
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(shutdownMutex_);
|
|
|
+ shuttingDown_.store(true, std::memory_order_release);
|
|
|
+ }
|
|
|
+ shutdownCv_.notify_all();
|
|
|
+}
|
|
|
+
|
|
|
// ===== Health and Stats =====
|
|
|
|
|
|
grpc::Status DatabaseGrpcImpl::HealthCheck(
|