messqueue.cpp 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198
  1. #include "key_value_printer.h"
  2. #include "mb_lwtrace.h"
  3. #include "remote_client_session.h"
  4. #include "remote_server_session.h"
  5. #include "ybus.h"
  6. #include <util/generic/singleton.h>
  7. using namespace NBus;
  8. using namespace NBus::NPrivate;
  9. using namespace NActor;
  10. TBusMessageQueuePtr NBus::CreateMessageQueue(const TBusQueueConfig& config, TExecutorPtr executor, TBusLocator* locator, const char* name) {
  11. return new TBusMessageQueue(config, executor, locator, name);
  12. }
  13. TBusMessageQueuePtr NBus::CreateMessageQueue(const TBusQueueConfig& config, TBusLocator* locator, const char* name) {
  14. TExecutor::TConfig executorConfig;
  15. executorConfig.WorkerCount = config.NumWorkers;
  16. executorConfig.Name = name;
  17. TExecutorPtr executor = new TExecutor(executorConfig);
  18. return CreateMessageQueue(config, executor, locator, name);
  19. }
  20. TBusMessageQueuePtr NBus::CreateMessageQueue(const TBusQueueConfig& config, const char* name) {
  21. return CreateMessageQueue(config, new TBusLocator, name);
  22. }
  23. TBusMessageQueuePtr NBus::CreateMessageQueue(TExecutorPtr executor, const char* name) {
  24. return CreateMessageQueue(TBusQueueConfig(), executor, new TBusLocator, name);
  25. }
  26. TBusMessageQueuePtr NBus::CreateMessageQueue(const char* name) {
  27. TBusQueueConfig config;
  28. return CreateMessageQueue(config, name);
  29. }
  30. namespace {
  31. TBusQueueConfig QueueConfigFillDefaults(const TBusQueueConfig& orig, const TString& name) {
  32. TBusQueueConfig patched = orig;
  33. if (!patched.Name) {
  34. patched.Name = name;
  35. }
  36. return patched;
  37. }
  38. }
  39. TBusMessageQueue::TBusMessageQueue(const TBusQueueConfig& config, TExecutorPtr executor, TBusLocator* locator, const char* name)
  40. : Config(QueueConfigFillDefaults(config, name))
  41. , Locator(locator)
  42. , WorkQueue(executor)
  43. , Running(1)
  44. {
  45. InitBusLwtrace();
  46. InitNetworkSubSystem();
  47. }
  48. TBusMessageQueue::~TBusMessageQueue() {
  49. Stop();
  50. }
  51. void TBusMessageQueue::Stop() {
  52. if (!AtomicCas(&Running, 0, 1)) {
  53. ShutdownComplete.WaitI();
  54. return;
  55. }
  56. Scheduler.Stop();
  57. DestroyAllSessions();
  58. WorkQueue->Stop();
  59. ShutdownComplete.Signal();
  60. }
  61. bool TBusMessageQueue::IsRunning() {
  62. return AtomicGet(Running);
  63. }
  64. TBusMessageQueueStatus TBusMessageQueue::GetStatusRecordInternal() const {
  65. TBusMessageQueueStatus r;
  66. r.ExecutorStatus = WorkQueue->GetStatusRecordInternal();
  67. r.Config = Config;
  68. return r;
  69. }
  70. TString TBusMessageQueue::GetStatusSelf() const {
  71. return GetStatusRecordInternal().PrintToString();
  72. }
  73. TString TBusMessageQueue::GetStatusSingleLine() const {
  74. return WorkQueue->GetStatusSingleLine();
  75. }
  76. TString TBusMessageQueue::GetStatus(ui16 flags) const {
  77. TStringStream ss;
  78. ss << GetStatusSelf();
  79. TList<TIntrusivePtr<TBusSessionImpl>> sessions;
  80. {
  81. TGuard<TMutex> scope(Lock);
  82. sessions = Sessions;
  83. }
  84. for (TList<TIntrusivePtr<TBusSessionImpl>>::const_iterator session = sessions.begin();
  85. session != sessions.end(); ++session) {
  86. ss << Endl;
  87. ss << (*session)->GetStatus(flags);
  88. }
  89. ss << Endl;
  90. ss << "object counts (not necessarily owned by this message queue):" << Endl;
  91. TKeyValuePrinter p;
  92. p.AddRow("TRemoteClientConnection", TObjectCounter<TRemoteClientConnection>::ObjectCount(), false);
  93. p.AddRow("TRemoteServerConnection", TObjectCounter<TRemoteServerConnection>::ObjectCount(), false);
  94. p.AddRow("TRemoteClientSession", TObjectCounter<TRemoteClientSession>::ObjectCount(), false);
  95. p.AddRow("TRemoteServerSession", TObjectCounter<TRemoteServerSession>::ObjectCount(), false);
  96. p.AddRow("NEventLoop::TEventLoop", TObjectCounter<NEventLoop::TEventLoop>::ObjectCount(), false);
  97. p.AddRow("NEventLoop::TChannel", TObjectCounter<NEventLoop::TChannel>::ObjectCount(), false);
  98. ss << p.PrintToString();
  99. return ss.Str();
  100. }
  101. TBusClientSessionPtr TBusMessageQueue::CreateSource(TBusProtocol* proto, IBusClientHandler* handler, const TBusClientSessionConfig& config, const TString& name) {
  102. TRemoteClientSessionPtr session(new TRemoteClientSession(this, proto, handler, config, name));
  103. Add(session.Get());
  104. return session.Get();
  105. }
  106. TBusServerSessionPtr TBusMessageQueue::CreateDestination(TBusProtocol* proto, IBusServerHandler* handler, const TBusClientSessionConfig& config, const TString& name) {
  107. TRemoteServerSessionPtr session(new TRemoteServerSession(this, proto, handler, config, name));
  108. try {
  109. int port = config.ListenPort;
  110. if (port == 0) {
  111. port = Locator->GetLocalPort(proto->GetService());
  112. }
  113. if (port == 0) {
  114. port = proto->GetPort();
  115. }
  116. session->Listen(port, this);
  117. Add(session.Get());
  118. return session.Release();
  119. } catch (...) {
  120. Y_ABORT("create destination failure: %s", CurrentExceptionMessage().c_str());
  121. }
  122. }
  123. TBusServerSessionPtr TBusMessageQueue::CreateDestination(TBusProtocol* proto, IBusServerHandler* handler, const TBusServerSessionConfig& config, const TVector<TBindResult>& bindTo, const TString& name) {
  124. TRemoteServerSessionPtr session(new TRemoteServerSession(this, proto, handler, config, name));
  125. try {
  126. session->Listen(bindTo, this);
  127. Add(session.Get());
  128. return session.Release();
  129. } catch (...) {
  130. Y_ABORT("create destination failure: %s", CurrentExceptionMessage().c_str());
  131. }
  132. }
  133. void TBusMessageQueue::Add(TIntrusivePtr<TBusSessionImpl> session) {
  134. TGuard<TMutex> scope(Lock);
  135. Sessions.push_back(session);
  136. }
  137. void TBusMessageQueue::Remove(TBusSession* session) {
  138. TGuard<TMutex> scope(Lock);
  139. TList<TIntrusivePtr<TBusSessionImpl>>::iterator it = std::find(Sessions.begin(), Sessions.end(), session);
  140. Y_ABORT_UNLESS(it != Sessions.end(), "do not destroy session twice");
  141. Sessions.erase(it);
  142. }
  143. void TBusMessageQueue::Destroy(TBusSession* session) {
  144. session->Shutdown();
  145. }
  146. void TBusMessageQueue::DestroyAllSessions() {
  147. TList<TIntrusivePtr<TBusSessionImpl>> sessions;
  148. {
  149. TGuard<TMutex> scope(Lock);
  150. sessions = Sessions;
  151. }
  152. for (auto& session : sessions) {
  153. Y_ABORT_UNLESS(session->IsDown(), "Session must be shut down prior to queue shutdown");
  154. }
  155. }
  156. void TBusMessageQueue::Schedule(IScheduleItemAutoPtr i) {
  157. Scheduler.Schedule(i);
  158. }
  159. TString TBusMessageQueue::GetNameInternal() const {
  160. return Config.Name;
  161. }