| 776 | } |
| 777 | |
| 778 | Status Statestore::Init(int32_t state_store_port) { |
| 779 | #ifndef NDEBUG |
| 780 | if (FLAGS_stress_statestore_startup_delay_ms > 0) { |
| 781 | LOG(INFO) << "Stress statestore startup delay: " |
| 782 | << FLAGS_stress_statestore_startup_delay_ms << " ms"; |
| 783 | SleepForMs(FLAGS_stress_statestore_startup_delay_ms); |
| 784 | } |
| 785 | #endif |
| 786 | std::shared_ptr<TProcessor> processor(new StatestoreServiceProcessor(thrift_iface())); |
| 787 | std::shared_ptr<TProcessorEventHandler> event_handler( |
| 788 | new RpcEventHandler("statestore", metrics_)); |
| 789 | processor->setEventHandler(event_handler); |
| 790 | ThriftServerBuilder builder("StatestoreService", processor, state_store_port); |
| 791 | // Mark this as an internal service to use a more permissive Thrift max message size |
| 792 | builder.is_external_facing(false); |
| 793 | if (IsInternalTlsConfigured()) { |
| 794 | SSLProtocol ssl_version; |
| 795 | RETURN_IF_ERROR( |
| 796 | SSLProtoVersions::StringToProtocol(FLAGS_ssl_minimum_version, &ssl_version)); |
| 797 | LOG(INFO) << "Enabling SSL for Statestore"; |
| 798 | builder.ssl(FLAGS_ssl_server_certificate, FLAGS_ssl_private_key) |
| 799 | .pem_password_cmd(FLAGS_ssl_private_key_password_cmd) |
| 800 | .ssl_version(ssl_version) |
| 801 | .cipher_list(FLAGS_ssl_cipher_list) |
| 802 | .tls_ciphersuites(FLAGS_tls_ciphersuites); |
| 803 | } |
| 804 | RETURN_IF_ERROR(subscriber_topic_update_threadpool_.Init()); |
| 805 | RETURN_IF_ERROR(subscriber_priority_topic_update_threadpool_.Init()); |
| 806 | RETURN_IF_ERROR(subscriber_heartbeat_threadpool_.Init()); |
| 807 | |
| 808 | ThriftServer* server; |
| 809 | RETURN_IF_ERROR(builder.metrics(metrics_).Build(&server)); |
| 810 | thrift_server_.reset(server); |
| 811 | RETURN_IF_ERROR(thrift_server_->Start()); |
| 812 | |
| 813 | RETURN_IF_ERROR(Thread::Create("statestore-heartbeat", "heartbeat-monitoring-thread", |
| 814 | &Statestore::MonitorSubscriberHeartbeat, this, &heartbeat_monitoring_thread_)); |
| 815 | RETURN_IF_ERROR(Thread::Create("statestore-update-catalogd", "update-catalogd-thread", |
| 816 | &Statestore::MonitorUpdateCatalogd, this, &update_catalogd_thread_)); |
| 817 | service_started_ = true; |
| 818 | service_started_metric_->SetValue(true); |
| 819 | return Status::OK(); |
| 820 | } |
| 821 | |
| 822 | void Statestore::RegisterWebpages(Webserver* webserver, bool metrics_only) { |
| 823 | Webserver::RawUrlCallback healthz_callback = |