Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
import com.google.bigtable.v2.SessionCheckAndMutateRowResponse;
import com.google.bigtable.v2.SessionMutateRowRequest;
import com.google.bigtable.v2.SessionMutateRowResponse;
import com.google.bigtable.v2.SessionReadModifyWriteRowRequest;
import com.google.bigtable.v2.SessionReadModifyWriteRowResponse;
import com.google.bigtable.v2.SessionReadRowRequest;
import com.google.bigtable.v2.SessionReadRowResponse;
import com.google.cloud.bigtable.data.v2.internal.channels.ChannelPool;
Expand Down Expand Up @@ -77,6 +79,7 @@ static AuthorizedViewAsync createAndStart(
VRpcDescriptor.READ_ROW_AUTH_VIEW,
VRpcDescriptor.MUTATE_ROW_AUTH_VIEW,
VRpcDescriptor.CHECK_AND_MUTATE_ROW_AUTH_VIEW,
VRpcDescriptor.READ_MODIFY_WRITE_ROW_AUTH_VIEW,
featureFlags,
clientInfo,
configManager,
Expand Down Expand Up @@ -120,6 +123,13 @@ public CompletableFuture<SessionCheckAndMutateRowResponse> checkAndMutateRow(
return f;
}

public CompletableFuture<SessionReadModifyWriteRowResponse> readModifyWriteRow(
SessionReadModifyWriteRowRequest req, Deadline deadline) {
UnaryResponseFuture<SessionReadModifyWriteRowResponse> f = new UnaryResponseFuture<>();
base.readModifyWriteRow(req, f, deadline);
return f;
}

@Override
public void close() {
this.base.close();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,9 +68,9 @@ public static MaterializedViewAsync createAndStart(
openReq,
VRpcDescriptor.MATERIALIZED_VIEW_SESSION,
VRpcDescriptor.READ_ROW_MAT_VIEW,
// Materialized views are read-only, so mutateRow and checkAndMutateRow don't apply.
null,
null,
/* mutateRowDescriptor= */ null, // Materialized views are read-only
/* checkAndMutateRowDescriptor= */ null,
/* readModifyWriteRowDescriptor= */ null,
featureFlags,
clientInfo,
configManager,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
import com.google.bigtable.v2.SessionCheckAndMutateRowResponse;
import com.google.bigtable.v2.SessionMutateRowRequest;
import com.google.bigtable.v2.SessionMutateRowResponse;
import com.google.bigtable.v2.SessionReadModifyWriteRowRequest;
import com.google.bigtable.v2.SessionReadModifyWriteRowResponse;
import com.google.bigtable.v2.SessionReadRowRequest;
import com.google.bigtable.v2.SessionReadRowResponse;
import com.google.cloud.bigtable.data.v2.internal.channels.ChannelPool;
Expand Down Expand Up @@ -75,6 +77,7 @@ public static TableAsync createAndStart(
VRpcDescriptor.READ_ROW,
VRpcDescriptor.MUTATE_ROW,
VRpcDescriptor.CHECK_AND_MUTATE_ROW,
VRpcDescriptor.READ_MODIFY_WRITE_ROW,
featureFlags,
clientInfo,
configManager,
Expand Down Expand Up @@ -102,7 +105,6 @@ public SessionPool<?> getSessionPool() {
return base.getSessionPool();
}

// TODO: get deadline from compatibility layer
// Currently these are the deadlines from gax:
// ApiCallContext#timeout (attempt timeout)
// ApiCallContext#RetrySettings#totalTimeout
Expand All @@ -117,19 +119,24 @@ public CompletableFuture<SessionReadRowResponse> readRow(
return f;
}

// TODO: get deadline from compatibility layer
public CompletableFuture<SessionMutateRowResponse> mutateRow(
SessionMutateRowRequest req, Deadline deadline) {
UnaryResponseFuture<SessionMutateRowResponse> f = new UnaryResponseFuture<>();
base.mutateRow(req, f, deadline);
return f;
}

// TODO: get deadline from compatibility layer
public CompletableFuture<SessionCheckAndMutateRowResponse> checkAndMutateRow(
SessionCheckAndMutateRowRequest req, Deadline deadline) {
UnaryResponseFuture<SessionCheckAndMutateRowResponse> f = new UnaryResponseFuture<>();
base.checkAndMutateRow(req, f, deadline);
return f;
}

public CompletableFuture<SessionReadModifyWriteRowResponse> readModifyWriteRow(
SessionReadModifyWriteRowRequest req, Deadline deadline) {
UnaryResponseFuture<SessionReadModifyWriteRowResponse> f = new UnaryResponseFuture<>();
base.readModifyWriteRow(req, f, deadline);
return f;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
import com.google.bigtable.v2.SessionCheckAndMutateRowResponse;
import com.google.bigtable.v2.SessionMutateRowRequest;
import com.google.bigtable.v2.SessionMutateRowResponse;
import com.google.bigtable.v2.SessionReadModifyWriteRowRequest;
import com.google.bigtable.v2.SessionReadModifyWriteRowResponse;
import com.google.bigtable.v2.SessionReadRowRequest;
import com.google.bigtable.v2.SessionReadRowResponse;
import com.google.cloud.bigtable.data.v2.internal.channels.ChannelPool;
Expand Down Expand Up @@ -55,6 +57,9 @@ class TableBase implements AutoCloseable {
mutateRowDescriptor;
private final VRpcDescriptor<?, SessionCheckAndMutateRowRequest, SessionCheckAndMutateRowResponse>
checkAndMutateRowDescriptor;
private final VRpcDescriptor<
?, SessionReadModifyWriteRowRequest, SessionReadModifyWriteRowResponse>
readModifyWriteRowDescriptor;

static <ReqT extends Message> TableBase createAndStart(
ReqT openReq,
Expand All @@ -63,6 +68,8 @@ static <ReqT extends Message> TableBase createAndStart(
VRpcDescriptor<?, SessionMutateRowRequest, SessionMutateRowResponse> mutateRowDescriptor,
VRpcDescriptor<?, SessionCheckAndMutateRowRequest, SessionCheckAndMutateRowResponse>
checkAndMutateRowDescriptor,
VRpcDescriptor<?, SessionReadModifyWriteRowRequest, SessionReadModifyWriteRowResponse>
readModifyWriteRowDescriptor,
FeatureFlags featureFlags,
ClientInfo clientInfo,
ClientConfigurationManager configManager,
Expand Down Expand Up @@ -100,6 +107,7 @@ static <ReqT extends Message> TableBase createAndStart(
readRowDescriptor,
mutateRowDescriptor,
checkAndMutateRowDescriptor,
readModifyWriteRowDescriptor,
metrics,
timer,
userCallbackExecutor);
Expand All @@ -112,13 +120,16 @@ static <ReqT extends Message> TableBase createAndStart(
VRpcDescriptor<?, SessionMutateRowRequest, SessionMutateRowResponse> mutateRowDescriptor,
VRpcDescriptor<?, SessionCheckAndMutateRowRequest, SessionCheckAndMutateRowResponse>
checkAndMutateRowDescriptor,
VRpcDescriptor<?, SessionReadModifyWriteRowRequest, SessionReadModifyWriteRowResponse>
readModifyWriteRowDescriptor,
Metrics metrics,
BigtableTimer timer,
Executor userCallbackExecutor) {
this.sessionPool = sessionPool;
this.readRowDescriptor = readRowDescriptor;
this.mutateRowDescriptor = mutateRowDescriptor;
this.checkAndMutateRowDescriptor = checkAndMutateRowDescriptor;
this.readModifyWriteRowDescriptor = readModifyWriteRowDescriptor;
this.metrics = metrics;
this.timer = timer;
this.userCallbackExecutor = userCallbackExecutor;
Expand Down Expand Up @@ -166,11 +177,28 @@ public void checkAndMutateRow(
SessionCheckAndMutateRowRequest req,
VRpcListener<SessionCheckAndMutateRowResponse> listener,
Deadline deadline) {
// RetryingVRpc is still needed even for non-idempotent ops: the idempotent=false flag prevents
// client-initiated retries on application-level errors, but the server may still direct a retry
// via response status.
RetryingVRpc<SessionCheckAndMutateRowRequest, SessionCheckAndMutateRowResponse> retry =
new RetryingVRpc<>(() -> sessionPool.newCall(checkAndMutateRowDescriptor), timer);
VRpcTracer tracer =
metrics.newTableTracer(sessionPool.getInfo(), checkAndMutateRowDescriptor, deadline);
// CheckAndMutateRow is not idempotent and must never be retried.
new VOperationImpl<>(retry, Context.current(), userCallbackExecutor, tracer, deadline, false)
.start(req, listener);
}

public void readModifyWriteRow(
SessionReadModifyWriteRowRequest req,
VRpcListener<SessionReadModifyWriteRowResponse> listener,
Deadline deadline) {
// RetryingVRpc is still needed even for non-idempotent ops: the idempotent=false flag prevents
// client-initiated retries on application-level errors, but the server may still direct a retry
// via response status.
RetryingVRpc<SessionReadModifyWriteRowRequest, SessionReadModifyWriteRowResponse> retry =
Comment thread
mutianf marked this conversation as resolved.
new RetryingVRpc<>(() -> sessionPool.newCall(readModifyWriteRowDescriptor), timer);
VRpcTracer tracer =
metrics.newTableTracer(sessionPool.getInfo(), readModifyWriteRowDescriptor, deadline);
new VOperationImpl<>(retry, Context.current(), userCallbackExecutor, tracer, deadline, false)
.start(req, listener);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import com.google.api.gax.rpc.UnaryCallable;
import com.google.cloud.bigtable.data.v2.models.ConditionalRowMutation;
import com.google.cloud.bigtable.data.v2.models.Query;
import com.google.cloud.bigtable.data.v2.models.ReadModifyWriteRow;
import com.google.cloud.bigtable.data.v2.models.RowAdapter;
import com.google.cloud.bigtable.data.v2.models.RowMutation;

Expand Down Expand Up @@ -46,4 +47,12 @@ public UnaryCallable<ConditionalRowMutation, Boolean> decorateCheckAndMutateRow(
UnaryCallable<ConditionalRowMutation, Boolean> classic, UnaryCallSettings<?, ?> settings) {
return classic;
}

@Override
public <RowT> UnaryCallable<ReadModifyWriteRow, RowT> decorateReadModifyWriteRow(
UnaryCallable<ReadModifyWriteRow, RowT> classic,
RowAdapter<RowT> rowAdapter,
UnaryCallSettings<?, ?> settings) {
return classic;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import com.google.api.gax.rpc.UnaryCallable;
import com.google.cloud.bigtable.data.v2.models.ConditionalRowMutation;
import com.google.cloud.bigtable.data.v2.models.Query;
import com.google.cloud.bigtable.data.v2.models.ReadModifyWriteRow;
import com.google.cloud.bigtable.data.v2.models.RowAdapter;
import com.google.cloud.bigtable.data.v2.models.RowMutation;

Expand All @@ -36,4 +37,9 @@ UnaryCallable<RowMutation, Void> decorateMutateRow(

UnaryCallable<ConditionalRowMutation, Boolean> decorateCheckAndMutateRow(
UnaryCallable<ConditionalRowMutation, Boolean> classic, UnaryCallSettings<?, ?> settings);

<RowT> UnaryCallable<ReadModifyWriteRow, RowT> decorateReadModifyWriteRow(
UnaryCallable<ReadModifyWriteRow, RowT> classic,
RowAdapter<RowT> rowAdapter,
UnaryCallSettings<?, ?> settings);
}
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,10 @@
import com.google.cloud.bigtable.data.v2.internal.compat.ops.CheckAndMutateRowShim;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.DivertingUnaryCallable;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.MutateRowShim;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.ReadModifyWriteRowShim;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.ReadRowShim;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.ReadRowShimInner;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.ReadWriteSessionPools;
import com.google.cloud.bigtable.data.v2.internal.compat.ops.RowBuilderShim;
import com.google.cloud.bigtable.data.v2.internal.csm.Metrics;
import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo;
import com.google.cloud.bigtable.data.v2.internal.csm.tracers.DebugTagTracer;
Expand All @@ -47,6 +49,7 @@
import com.google.cloud.bigtable.data.v2.internal.util.ClientConfigurationManager;
import com.google.cloud.bigtable.data.v2.models.ConditionalRowMutation;
import com.google.cloud.bigtable.data.v2.models.Query;
import com.google.cloud.bigtable.data.v2.models.ReadModifyWriteRow;
import com.google.cloud.bigtable.data.v2.models.RowAdapter;
import com.google.cloud.bigtable.data.v2.models.RowMutation;
import com.google.cloud.bigtable.data.v2.stub.MetadataExtractorInterceptor;
Expand Down Expand Up @@ -77,18 +80,17 @@
public class ShimImpl implements Shim {
private static final Logger logger = Logger.getLogger(ShimImpl.class.getName());

// TODO: this should be a client config
public static final int MAX_CONSECUTIVE_UNIMPLEMENTED_FAILURES = 30;
private static final Duration DA_CHECK_TIMEOUT = Duration.ofSeconds(5);

private final ClientConfigurationManager configManager;
private final Resource<ClientConfigurationManager> configManagerResource;
private final Client client;
private final DebugTagTracer debugTagTracer;

private final ReadRowShimInner readRowShimInner;
private final ReadRowShim readRowShim;
private final MutateRowShim mutateRowShim;
private final CheckAndMutateRowShim checkAndMutateRowShim;
private final ReadModifyWriteRowShim readModifyWriteRowShim;

public static Shim create(
ClientInfo clientInfo,
Expand Down Expand Up @@ -202,9 +204,11 @@ public ShimImpl(
this.client = client;
this.debugTagTracer = debugTagTracer;

this.readRowShimInner = new ReadRowShimInner(client);
this.readRowShim = new ReadRowShim(client);
this.mutateRowShim = new MutateRowShim(client);
this.checkAndMutateRowShim = new CheckAndMutateRowShim(client);
ReadWriteSessionPools rwPools = new ReadWriteSessionPools(client);
this.checkAndMutateRowShim = new CheckAndMutateRowShim(rwPools);
this.readModifyWriteRowShim = new ReadModifyWriteRowShim(rwPools);
}

/**
Expand Down Expand Up @@ -368,7 +372,7 @@ public <RowT> UnaryCallable<Query, RowT> decorateReadRow(
return new DivertingUnaryCallable<>(
configManager,
classic,
new ReadRowShim<>(readRowShimInner, rowAdapter),
new RowBuilderShim<>(readRowShim, rowAdapter, r -> r.hasRow() ? r.getRow() : null),
Util.extractTimeout(settings),
debugTagTracer);
}
Expand All @@ -383,11 +387,31 @@ public UnaryCallable<RowMutation, Void> decorateMutateRow(
@Override
public UnaryCallable<ConditionalRowMutation, Boolean> decorateCheckAndMutateRow(
UnaryCallable<ConditionalRowMutation, Boolean> classic, UnaryCallSettings<?, ?> settings) {
if (!WipFeatures.CHECK_AND_MUTATE_ROW_ENABLED) {
return classic;
}
return new DivertingUnaryCallable<>(
configManager,
classic,
checkAndMutateRowShim,
Util.extractTimeout(settings),
debugTagTracer);
}

@Override
public <RowT> UnaryCallable<ReadModifyWriteRow, RowT> decorateReadModifyWriteRow(
UnaryCallable<ReadModifyWriteRow, RowT> classic,
RowAdapter<RowT> rowAdapter,
UnaryCallSettings<?, ?> settings) {
if (!WipFeatures.READ_MODIFY_WRITE_ROW_ENABLED) {
return classic;
}
return new DivertingUnaryCallable<>(
configManager,
classic,
new RowBuilderShim<>(
readModifyWriteRowShim, rowAdapter, r -> r.hasRow() ? r.getRow() : null),
Util.extractTimeout(settings),
debugTagTracer);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
/*
* Copyright 2026 Google LLC
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.google.cloud.bigtable.data.v2.internal.compat;

/** Tracks work-in-progress session-path features that are not yet enabled in production. */
public final class WipFeatures {
// Enable session-path diversion for CheckAndMutateRow once per-method diversion is supported.
public static final boolean CHECK_AND_MUTATE_ROW_ENABLED = false;

// Enable session-path diversion for ReadModifyWriteRow once per-method diversion is supported.
public static final boolean READ_MODIFY_WRITE_ROW_ENABLED = false;

private WipFeatures() {}
}
Loading
Loading