diff --git a/google/cloud/bigtable/data_connection.cc b/google/cloud/bigtable/data_connection.cc index bac6ccf317c70..55120ce7d2139 100644 --- a/google/cloud/bigtable/data_connection.cc +++ b/google/cloud/bigtable/data_connection.cc @@ -209,30 +209,34 @@ std::shared_ptr MakeInstanceAffinityDataConnection( std::move(components.limiter), std::move(components.options)); } -std::shared_ptr MakeDirectPathDataConnection( - DataConnectionComponents components, - std::vector const& instances) { - // Step 1: Configure separate Options for DirectPath and CloudPath (fallback). - Options directpath_options = components.options; - directpath_options - .set<::google::cloud::bigtable_internal::DataEndpointOption>( - bigtable::internal::DefaultDirectPathDataEndpoint()); - directpath_options.set( +Options MakeDirectPathOptions(Options options) { + options.set<::google::cloud::bigtable_internal::DataEndpointOption>( + bigtable::internal::DefaultDirectPathDataEndpoint()); + options.set( bigtable::internal::DefaultDirectPathDataEndpoint()); - directpath_options.set( + options.set( bigtable::internal::DefaultDirectPathAuthority()); - directpath_options.set( + options.set( bigtable::experimental::DirectPathMode::kEnabled); + return options; +} - Options cloudpath_options = components.options; - cloudpath_options.set<::google::cloud::bigtable_internal::DataEndpointOption>( +Options MakeCloudPathOptions(Options options) { + options.set<::google::cloud::bigtable_internal::DataEndpointOption>( "bigtable.googleapis.com"); - cloudpath_options.set("bigtable.googleapis.com"); - cloudpath_options.set("bigtable.googleapis.com"); - cloudpath_options.set( + options.set("bigtable.googleapis.com"); + options.set("bigtable.googleapis.com"); + options.set( bigtable::experimental::DirectPathMode::kDisabled); + return options; +} - // Step 2: Concurrently create StubManagers for both DirectPath and CloudPath +std::shared_ptr +MakeSpeculativeDirectPathDataConnection( + DataConnectionComponents components, + std::vector const& instances, + Options directpath_options, Options cloudpath_options) { + // Step 1: Concurrently create StubManagers for both DirectPath and CloudPath // on background threads so that stubs are warming up while probing proceeds. promise> directpath_promise; future> directpath_future = @@ -279,14 +283,14 @@ std::shared_ptr MakeDirectPathDataConnection( std::move(affinity_stubs), stub_creation_fn)); }); - // Step 3: Probe DirectPath connectivity and ALTS negotiation on the primary + // Step 2: Probe DirectPath connectivity and ALTS negotiation on the primary // instance. StatusOr const probe_result = bigtable_internal::DirectPathProber::Probe( components.auth, instances.front(), directpath_options, components.background->cq()); - // Step 4a: If DirectPath probing succeeds, use DirectPath stubs, discard the + // Step 3a: If DirectPath probing succeeds, use DirectPath stubs, discard the // fallback CloudPath future asynchronously, and record successful metrics. if (probe_result.ok() && probe_result->success) { std::unique_ptr stub_manager = @@ -309,7 +313,7 @@ std::shared_ptr MakeDirectPathDataConnection( std::move(components.limiter), std::move(directpath_options)); } - // Step 4b: If DirectPath probing fails, fall back to the CloudPath + // Step 3b: If DirectPath probing fails, fall back to the CloudPath // StubManager, discard the unused DirectPath future asynchronously, and run // environmental diagnostics asynchronously to diagnose and record metrics for // the failure reason. @@ -330,6 +334,88 @@ std::shared_ptr MakeDirectPathDataConnection( std::move(components.limiter), std::move(cloudpath_options)); } +std::shared_ptr MakeBlockingDirectPathDataConnection( + DataConnectionComponents components, + std::vector const& instances, + Options directpath_options, Options cloudpath_options) { + // Step 1: Probe DirectPath connectivity and ALTS negotiation on the primary + // instance first. + StatusOr const probe_result = + bigtable_internal::DirectPathProber::Probe( + components.auth, instances.front(), directpath_options, + components.background->cq()); + + auto create_stub_manager = [&](Options const& options) { + auto stub_creation_fn = + [auth = components.auth, cq = components.background->cq(), options]( + std::string_view instance_name, + bigtable_internal::StubManager::Priming priming) { + return bigtable_internal::CreateBigtableStub(auth, cq, instance_name, + priming, options); + }; + absl::flat_hash_map> + affinity_stubs = bigtable_internal::CreateBigtableAffinityStubs( + instances, stub_creation_fn); + return std::make_unique( + std::move(affinity_stubs), std::move(stub_creation_fn)); + }; + + // Step 2a: If DirectPath probing succeeds, create the DirectPath StubManager + // and record successful metrics. + if (probe_result.ok() && probe_result->success) { + auto stub_manager = create_stub_manager(directpath_options); + +#ifdef GOOGLE_CLOUD_CPP_BIGTABLE_WITH_OTEL_METRICS + if (components.direct_access_compatibility != nullptr) { + components.direct_access_compatibility->Record( + opentelemetry::context::RuntimeContext::GetCurrent(), 1, + bigtable_internal::DirectAccessCompatibilityLabels{ + bigtable_internal::ToString(probe_result->ip_preference), ""}); + } +#endif + return std::make_shared( + std::move(components.background), std::move(stub_manager), + std::move(components.operation_context_factory), + std::move(components.limiter), std::move(directpath_options)); + } + + // Step 2b: If DirectPath probing fails, create the CloudPath StubManager, + // and run environmental diagnostics asynchronously. + auto stub_manager = create_stub_manager(cloudpath_options); + +#ifdef GOOGLE_CLOUD_CPP_BIGTABLE_WITH_OTEL_METRICS + bigtable_internal::DirectPathDiagnostics::RunAsync( + components.background->cq(), directpath_options, + components.direct_access_compatibility); +#endif + return std::make_shared( + std::move(components.background), std::move(stub_manager), + std::move(components.operation_context_factory), + std::move(components.limiter), std::move(cloudpath_options)); +} + +std::shared_ptr MakeDirectPathDataConnection( + DataConnectionComponents components, + std::vector const& instances) { + experimental::DirectPathInitializationMode const initialization_mode = + components.options + .get(); + Options directpath_options = MakeDirectPathOptions(components.options); + Options cloudpath_options = + MakeCloudPathOptions(std::move(components.options)); + + if (initialization_mode == + bigtable::experimental::DirectPathInitializationMode::kBlocking) { + return MakeBlockingDirectPathDataConnection( + std::move(components), instances, std::move(directpath_options), + std::move(cloudpath_options)); + } + return MakeSpeculativeDirectPathDataConnection( + std::move(components), instances, std::move(directpath_options), + std::move(cloudpath_options)); +} + std::shared_ptr WrapDataTracing( std::shared_ptr conn) { if (google::cloud::internal::TracingEnabled(conn->options())) { diff --git a/google/cloud/bigtable/data_connection_test.cc b/google/cloud/bigtable/data_connection_test.cc index ca4c710426ffc..c9b0a3a676d2d 100644 --- a/google/cloud/bigtable/data_connection_test.cc +++ b/google/cloud/bigtable/data_connection_test.cc @@ -149,6 +149,30 @@ TEST(MakeDataConnection, DirectPathModeOptionDisabledNoGrpcMetrics) { experimental::DirectPathMetricsMode::kEnabled)); EXPECT_NE(conn, nullptr); } + +TEST(MakeDataConnection, DirectPathInitializationModeBlocking) { + InstanceResource instance_a{Project("my-project"), "instance-a"}; + auto conn = MakeDataConnection( + {instance_a}, + TestOptions() + .set( + experimental::DirectPathMode::kEnabled) + .set( + experimental::DirectPathInitializationMode::kBlocking)); + EXPECT_NE(conn, nullptr); +} + +TEST(MakeDataConnection, DirectPathInitializationModeSpeculative) { + InstanceResource instance_a{Project("my-project"), "instance-a"}; + auto conn = MakeDataConnection( + {instance_a}, TestOptions() + .set( + experimental::DirectPathMode::kEnabled) + .set( + experimental::DirectPathInitializationMode:: + kAsynchronousSpeculative)); + EXPECT_NE(conn, nullptr); +} #endif } // namespace diff --git a/google/cloud/bigtable/internal/defaults.cc b/google/cloud/bigtable/internal/defaults.cc index 6ad585cd9aecb..a35da7b76f6c5 100644 --- a/google/cloud/bigtable/internal/defaults.cc +++ b/google/cloud/bigtable/internal/defaults.cc @@ -260,6 +260,10 @@ Options DefaultOptions(Options opts) { opts.set( DefaultDirectPathDiagnosticsTimeout()); } + if (!opts.has()) { + opts.set( + DefaultDirectPathInitializationMode()); + } return opts; } diff --git a/google/cloud/bigtable/internal/defaults.h b/google/cloud/bigtable/internal/defaults.h index 117ede67482f8..7cffb697c6192 100644 --- a/google/cloud/bigtable/internal/defaults.h +++ b/google/cloud/bigtable/internal/defaults.h @@ -52,6 +52,11 @@ constexpr std::chrono::milliseconds DefaultDirectPathDiagnosticsTimeout() { return std::chrono::seconds(60); } +constexpr experimental::DirectPathInitializationMode +DefaultDirectPathInitializationMode() { + return experimental::DirectPathInitializationMode::kBlocking; +} + /** * Returns true if Direct Path is enabled for Bigtable. */ diff --git a/google/cloud/bigtable/internal/defaults_test.cc b/google/cloud/bigtable/internal/defaults_test.cc index 4942faf0b6309..38d22c639ab2c 100644 --- a/google/cloud/bigtable/internal/defaults_test.cc +++ b/google/cloud/bigtable/internal/defaults_test.cc @@ -636,6 +636,8 @@ TEST(EndpointEnvTest, DirectPathTimeoutDefaults) { Eq(DefaultDirectPathProbeTimeout())); EXPECT_THAT(opts.get(), Eq(DefaultDirectPathDiagnosticsTimeout())); + EXPECT_THAT(opts.get(), + Eq(DefaultDirectPathInitializationMode())); } TEST(EndpointEnvTest, EmulatorOverridesCloudDirectPath) { diff --git a/google/cloud/bigtable/options.h b/google/cloud/bigtable/options.h index 400baf6c9f046..a9a3dfd4d2ffa 100644 --- a/google/cloud/bigtable/options.h +++ b/google/cloud/bigtable/options.h @@ -265,6 +265,31 @@ struct DirectPathDiagnosticsTimeoutOption { using Type = std::chrono::milliseconds; }; +/** + * Option to configure the initialization behavior for DirectPath. + * + * - kBlocking - Channel Pool initialization begins only after + * DirectPath or CloudPath is determined. + * - kAsynchronousSpeculative - Both DirectPath and CloudPath channel pools + * are initialized on background threads while the probe runs. The unused + * channel pool is discarded. + * + * @ingroup google-cloud-bigtable-options + */ +enum class DirectPathInitializationMode { + kBlocking, // Default. + kAsynchronousSpeculative, +}; + +/** + * Option to configure the initialization behavior for DirectPath. + * + * @ingroup google-cloud-bigtable-options + */ +struct DirectPathInitializationModeOption { + using Type = DirectPathInitializationMode; +}; + } // namespace experimental /// The complete list of options accepted by `bigtable::*Client` @@ -355,7 +380,8 @@ using DataPolicyOptionList = OptionList; + experimental::DynamicChannelPoolSizingPolicyOption, + experimental::DirectPathInitializationModeOption>; GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace bigtable diff --git a/google/cloud/bigtable/tests/observability_integration_test.cc b/google/cloud/bigtable/tests/observability_integration_test.cc index cc87d470193aa..9cc00305f101e 100644 --- a/google/cloud/bigtable/tests/observability_integration_test.cc +++ b/google/cloud/bigtable/tests/observability_integration_test.cc @@ -159,6 +159,12 @@ class ObservabilityIntegrationTest void TearDown() override { TableIntegrationTest::TearDown(); } + static void VerifyDirectAccessCompatibleMetricOnDirectPath( + experimental::DirectPathInitializationMode mode); + + static void VerifyDirectAccessFallbackAndDiagnosticsMetric( + experimental::DirectPathInitializationMode mode); + static google::cloud::testing_util::OtelCollectorServer collector_service_; static std::unique_ptr server_; static std::string server_address_; @@ -453,8 +459,9 @@ TEST_F(ObservabilityIntegrationTest, VerifyOutstandingRpcsMetric) { HasMetricLabel("streaming", Not(IsEmpty())))))); } -TEST_F(ObservabilityIntegrationTest, - VerifyDirectAccessCompatibleMetricOnDirectPath) { +void ObservabilityIntegrationTest:: + VerifyDirectAccessCompatibleMetricOnDirectPath( + experimental::DirectPathInitializationMode mode) { if (UsingCloudBigtableEmulator()) { GTEST_SKIP() << "Metrics export integration test runs against production"; } @@ -481,13 +488,15 @@ TEST_F(ObservabilityIntegrationTest, server_address_); ScopedEnvironment env_otel("GOOGLE_CLOUD_CPP_TESTING_OTEL_COLLECTOR", "1"); - Options options = Options{} - .set(true) - .set( - experimental::DirectPathMode::kEnabled) - .set(std::chrono::seconds(5)) - .set(std::chrono::hours(1)) - .set(std::chrono::hours(1)); + Options options = + Options{} + .set(true) + .set( + experimental::DirectPathMode::kEnabled) + .set(mode) + .set(std::chrono::seconds(5)) + .set(std::chrono::hours(1)) + .set(std::chrono::hours(1)); std::string const& table_id = TableTestEnvironment::table_id(); @@ -498,7 +507,12 @@ TEST_F(ObservabilityIntegrationTest, Table table(std::move(conn), TableResource(project_id(), instance_id(), table_id)); - std::string const row_key = "observability-directpath-compatible-row-1"; + std::string const mode_str = + mode == experimental::DirectPathInitializationMode::kBlocking + ? "blocking" + : "speculative"; + std::string const row_key = absl::StrCat( + "observability-directpath-compatible-", mode_str, "-row-1"); std::vector expected{ {row_key, "family4", "c0", 1000, "v1000"}, {row_key, "family4", "c1", 2000, "v2000"}, @@ -532,8 +546,9 @@ TEST_F(ObservabilityIntegrationTest, HasMetricLabel("ip_preference", AnyOf(Eq("ipv4"), Eq("ipv6"))))))); } -TEST_F(ObservabilityIntegrationTest, - VerifyDirectAccessFallbackAndDiagnosticsMetric) { +void ObservabilityIntegrationTest:: + VerifyDirectAccessFallbackAndDiagnosticsMetric( + experimental::DirectPathInitializationMode mode) { if (UsingCloudBigtableEmulator()) { GTEST_SKIP() << "Metrics export integration test runs against production"; } @@ -558,13 +573,15 @@ TEST_F(ObservabilityIntegrationTest, server_address_); ScopedEnvironment env_otel("GOOGLE_CLOUD_CPP_TESTING_OTEL_COLLECTOR", "1"); - Options options = Options{} - .set(true) - .set( - experimental::DirectPathMode::kEnabled) - .set(std::chrono::seconds(5)) - .set(std::chrono::hours(1)) - .set(std::chrono::hours(1)); + Options options = + Options{} + .set(true) + .set( + experimental::DirectPathMode::kEnabled) + .set(mode) + .set(std::chrono::seconds(5)) + .set(std::chrono::hours(1)) + .set(std::chrono::hours(1)); std::string const& table_id = TableTestEnvironment::table_id(); @@ -575,7 +592,12 @@ TEST_F(ObservabilityIntegrationTest, Table table(std::move(conn), TableResource(project_id(), instance_id(), table_id)); - std::string const row_key = "observability-directpath-fallback-row-1"; + std::string const mode_str = + mode == experimental::DirectPathInitializationMode::kBlocking + ? "blocking" + : "speculative"; + std::string const row_key = + absl::StrCat("observability-directpath-fallback-", mode_str, "-row-1"); std::vector expected{ {row_key, "family4", "c0", 1000, "v1000"}, {row_key, "family4", "c1", 2000, "v2000"}, @@ -607,6 +629,30 @@ TEST_F(ObservabilityIntegrationTest, HasMetricLabel("ip_preference", IsEmpty()))))); } +TEST_F(ObservabilityIntegrationTest, + VerifyDirectAccessCompatibleMetricOnDirectPathBlocking) { + VerifyDirectAccessCompatibleMetricOnDirectPath( + experimental::DirectPathInitializationMode::kBlocking); +} + +TEST_F(ObservabilityIntegrationTest, + VerifyDirectAccessCompatibleMetricOnDirectPathSpeculative) { + VerifyDirectAccessCompatibleMetricOnDirectPath( + experimental::DirectPathInitializationMode::kAsynchronousSpeculative); +} + +TEST_F(ObservabilityIntegrationTest, + VerifyDirectAccessFallbackAndDiagnosticsMetricBlocking) { + VerifyDirectAccessFallbackAndDiagnosticsMetric( + experimental::DirectPathInitializationMode::kBlocking); +} + +TEST_F(ObservabilityIntegrationTest, + VerifyDirectAccessFallbackAndDiagnosticsMetricSpeculative) { + VerifyDirectAccessFallbackAndDiagnosticsMetric( + experimental::DirectPathInitializationMode::kAsynchronousSpeculative); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace bigtable