Skip to content

Commit ed424fc

Browse files
committed
feat(bigtable): route read_row/mutate_row through the accelerator with native fallback
Change-Id: I363bc3393ff13033e6c37ef9df4be13571837932
1 parent a71d88b commit ed424fc

8 files changed

Lines changed: 892 additions & 44 deletions

File tree

‎packages/google-cloud-bigtable/google/cloud/bigtable/data/_accelerator/_daemon.py‎

Lines changed: 46 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -25,12 +25,18 @@
2525

2626
import os
2727
import shutil
28+
import signal
2829
import socket
2930
import subprocess
3031
import tempfile
3132
import time
3233
from typing import Sequence
3334

35+
# Environment variable that overrides the bundled binary location. Primarily
36+
# for development against a locally-built daemon, and for tests pointing at a
37+
# fake binary.
38+
_BIN_ENV_VAR = "BIGTABLE_ACCELERATOR_BIN"
39+
3440
# Wheels ship the binary at this path relative to the `_accelerator/` package.
3541
_DEFAULT_BIN_RELATIVE_PATH = "bin/accelerator"
3642

@@ -44,12 +50,34 @@
4450
_SIGKILL_GRACE_SECONDS = 2.0
4551

4652

47-
def _resolve_binary_path() -> str:
48-
"""Return the path of the bundled daemon binary, raising if it is absent."""
53+
def _resolve_binary_path(explicit_path: str | None = None) -> str:
54+
"""Resolve the daemon binary path, validating that it is a regular file.
55+
56+
Precedence: an explicit ``binary_path`` argument, then the
57+
``BIGTABLE_ACCELERATOR_BIN`` env var, then the binary bundled in the wheel.
58+
An explicit path or env override that does not point at a regular file is a
59+
hard error (a caller who named a path meant it); a missing bundled binary
60+
reports how to supply one.
61+
"""
62+
if explicit_path is not None:
63+
if not os.path.isfile(explicit_path):
64+
raise FileNotFoundError(
65+
f"binary_path={explicit_path!r} does not point at a regular file"
66+
)
67+
return explicit_path
68+
override = os.environ.get(_BIN_ENV_VAR)
69+
if override:
70+
if not os.path.isfile(override):
71+
raise FileNotFoundError(
72+
f"{_BIN_ENV_VAR}={override!r} does not point at a regular file"
73+
)
74+
return override
4975
bundled = os.path.join(os.path.dirname(__file__), _DEFAULT_BIN_RELATIVE_PATH)
5076
if not os.path.isfile(bundled):
5177
raise FileNotFoundError(
52-
"Accelerator binary not found. Install a wheel that bundles the binary."
78+
"No accelerator binary found. Set the "
79+
f"{_BIN_ENV_VAR} env var to a daemon binary path, or install a "
80+
"wheel that bundles the binary."
5381
)
5482
return bundled
5583

@@ -76,20 +104,25 @@ def __init__(
76104
self,
77105
cli_flags: Sequence[str] = (),
78106
*,
107+
binary_path: str | None = None,
79108
startup_timeout: float = _DEFAULT_STARTUP_TIMEOUT,
80109
):
81110
"""Resolve the binary and pick the UDS path (does not spawn anything).
82111
83112
Args:
84113
cli_flags: extra arguments appended after ``--uds-path`` when
85114
spawning the daemon (e.g. ``--project``/``--instance``).
115+
binary_path: explicit path to the daemon binary. When omitted, the
116+
path is resolved from the ``BIGTABLE_ACCELERATOR_BIN`` env var
117+
and then the binary bundled in the wheel.
86118
startup_timeout: seconds ``start()`` waits for the daemon's UDS to
87119
become connectable before raising.
88120
89121
Raises:
90-
FileNotFoundError: the bundled binary is not present in this wheel.
122+
FileNotFoundError: no binary could be resolved, or an explicit
123+
``binary_path``/env override does not point at a regular file.
91124
"""
92-
self._binary_path = _resolve_binary_path()
125+
self._binary_path = _resolve_binary_path(binary_path)
93126
self._cli_flags = list(cli_flags)
94127
self._startup_timeout = startup_timeout
95128
self._tempdir: str | None = None
@@ -139,10 +172,6 @@ def start(self) -> None:
139172
"""
140173
if self._proc is not None:
141174
raise RuntimeError("AcceleratorDaemon.start() called twice")
142-
if not hasattr(socket, "AF_UNIX"):
143-
raise OSError(
144-
"Unix domain sockets (AF_UNIX) are not supported on this platform."
145-
)
146175
self._tempdir = tempfile.mkdtemp(prefix="bt-accel-")
147176
self._uds_path = os.path.join(self._tempdir, "sock")
148177
# Redirect the daemon's stdout/stderr to a log file rather than
@@ -157,10 +186,9 @@ def start(self) -> None:
157186
# The log file handle only needs to live long enough for Popen to dup
158187
# it into the child, so it stays local to start() rather than being an
159188
# attribute. Startup failures read the tail back from the path.
189+
log_file = open(self._log_path, "wb")
160190
argv = [self._binary_path, "--uds-path", self._uds_path, *self._cli_flags]
161-
log_file = None
162191
try:
163-
log_file = open(self._log_path, "wb")
164192
self._proc = subprocess.Popen(
165193
argv,
166194
stdin=subprocess.PIPE,
@@ -177,11 +205,10 @@ def start(self) -> None:
177205
# Whether or not the spawn succeeded, the parent no longer needs its
178206
# copy of the log fd: on success the child holds its own dup, and on
179207
# failure there is nothing to keep open.
180-
if log_file is not None:
181-
try:
182-
log_file.close()
183-
except OSError:
184-
pass
208+
try:
209+
log_file.close()
210+
except OSError:
211+
pass
185212
try:
186213
self._wait_until_ready(self._startup_timeout)
187214
except BaseException:
@@ -205,17 +232,11 @@ def close(self) -> None:
205232
pass
206233
if not self._wait_for_exit(_STDIN_GRACE_SECONDS):
207234
# Step 2: SIGTERM.
208-
try:
209-
proc.terminate()
210-
except OSError:
211-
pass
235+
proc.terminate()
212236
if not self._wait_for_exit(_SIGTERM_GRACE_SECONDS):
213237
# Step 3: SIGKILL. Bounded wait so teardown can't hang
214238
# forever if the process is stuck unreapable.
215-
try:
216-
proc.kill()
217-
except OSError:
218-
pass
239+
proc.kill()
219240
self._wait_for_exit(_SIGKILL_GRACE_SECONDS)
220241
finally:
221242
self._proc = None
@@ -295,7 +316,7 @@ def _force_kill(self) -> None:
295316
if self._proc is None or self._proc.poll() is not None:
296317
return
297318
try:
298-
self._proc.kill()
319+
self._proc.send_signal(signal.SIGKILL)
299320
except (OSError, ProcessLookupError):
300321
pass
301322
try:
@@ -309,4 +330,3 @@ def _cleanup_tempdir(self) -> None:
309330
shutil.rmtree(self._tempdir, ignore_errors=True)
310331
self._tempdir = None
311332
self._uds_path = None
312-
self._log_path = None
Lines changed: 146 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,146 @@
1+
# Copyright 2026 Google LLC
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
#
15+
"""Client-side fallback policy for accelerator-routed RPCs.
16+
17+
A daemon that cannot open any sessions replies ``UNIMPLEMENTED``, and the routing
18+
layer transparently retries the call on the native client. The first
19+
``UNIMPLEMENTED`` reply trips a sticky breaker so a persistently-degraded daemon
20+
stops being dialed at all. A daemon whose subprocess has died mid-flight trips
21+
the breaker immediately — it will never recover.
22+
23+
Any other gRPC error is a real, daemon-served result the native client would
24+
reproduce (the daemon owns retries, so it has already exhausted them), so it is
25+
translated to the corresponding ``google.api_core`` exception and raised without
26+
falling back.
27+
28+
This module is plain sync-only logic shared verbatim by the async and generated
29+
sync clients; ``grpc.RpcError`` is the common base of both ``grpc.RpcError`` and
30+
``grpc.aio.AioRpcError``, so no CrossSync branching is needed here.
31+
"""
32+
33+
from __future__ import annotations
34+
35+
import logging
36+
import threading
37+
from typing import TYPE_CHECKING
38+
39+
from grpc import RpcError, StatusCode
40+
41+
from google.api_core import exceptions as core_exceptions
42+
43+
if TYPE_CHECKING:
44+
from google.cloud.bigtable.data._accelerator._daemon import AcceleratorDaemon
45+
46+
_LOGGER = logging.getLogger(__name__)
47+
48+
49+
class _AcceleratorFallback(Exception):
50+
"""Internal signal that an accelerator attempt should be retried natively.
51+
52+
Never escapes the Table method that raises it: the method catches it and
53+
falls through to the native code path.
54+
"""
55+
56+
57+
class AcceleratorBreaker:
58+
"""Tracks accelerator health and decides when to stop using it.
59+
60+
One instance per Table. Thread-safe so the generated sync client can share a
61+
Table across threads. Two triggers permanently bypass the accelerator:
62+
63+
* the first ``UNIMPLEMENTED`` reply (the daemon understands the RPC shape but
64+
has no working sessions), and
65+
* an explicit :meth:`trip` when the daemon subprocess is found dead.
66+
"""
67+
68+
def __init__(self):
69+
self._tripped = False
70+
self._lock = threading.Lock()
71+
72+
def bypass(self) -> bool:
73+
"""Whether the accelerator should be skipped entirely from now on."""
74+
return self._tripped
75+
76+
def trip(self) -> None:
77+
"""Permanently bypass the accelerator (e.g. the daemon process died)."""
78+
with self._lock:
79+
self._tripped = True
80+
81+
82+
def _grpc_code(exc: BaseException) -> StatusCode | None:
83+
"""Best-effort extraction of a gRPC status code from an exception.
84+
85+
Handles both ``grpc.RpcError`` (status via a ``code()`` method) and
86+
``google.api_core.exceptions.GoogleAPICallError`` (status stored on the
87+
``grpc_status_code`` attribute), returning ``None`` for anything else.
88+
"""
89+
code = getattr(exc, "code", None)
90+
if callable(code):
91+
try:
92+
return code()
93+
except Exception:
94+
return None
95+
grpc_status = getattr(exc, "grpc_status_code", None)
96+
if isinstance(grpc_status, StatusCode):
97+
return grpc_status
98+
return None
99+
100+
101+
def handle_accelerator_error(
102+
exc: BaseException,
103+
*,
104+
daemon: "AcceleratorDaemon | None",
105+
breaker: AcceleratorBreaker,
106+
) -> None:
107+
"""Classify an exception raised by an accelerator-routed RPC.
108+
109+
Always raises. Either raises :class:`_AcceleratorFallback` to tell the caller
110+
to retry on the native path, or raises the translated ``google.api_core``
111+
exception for the caller to propagate:
112+
113+
* daemon subprocess dead -> trip the breaker, fall back (it will not recover)
114+
* ``UNIMPLEMENTED`` -> trip the breaker, fall back immediately
115+
* any other gRPC error -> translate and raise
116+
* a non-gRPC exception -> re-raise unchanged (never masked as a fallback)
117+
"""
118+
# TODO(accelerator): emit a metric here (e.g. a fallback/error counter keyed
119+
# by reason: dead-daemon / unimplemented / translated-error) once client-side
120+
# accelerator metrics are wired up.
121+
# A dead subprocess can surface as a channel error under any status code, so
122+
# check liveness first: the "daemon died mid-flight" case always wins and is
123+
# never recoverable.
124+
if daemon is not None and not daemon.is_running:
125+
_LOGGER.warning(
126+
"Accelerator daemon is no longer running; permanently falling back "
127+
"to the native Bigtable client for this table.",
128+
exc_info=exc,
129+
)
130+
breaker.trip()
131+
raise _AcceleratorFallback() from exc
132+
if not isinstance(exc, RpcError):
133+
# A bug in our own merge machinery, not a daemon result. Do not mask it
134+
# as a fallback; let it propagate unchanged.
135+
raise exc
136+
if _grpc_code(exc) == StatusCode.UNIMPLEMENTED:
137+
# The daemon only replies UNIMPLEMENTED once it has no working sessions,
138+
# a persistent condition, so trip the breaker and fall back immediately
139+
# rather than re-dialing on every subsequent call.
140+
_LOGGER.warning(
141+
"Accelerator daemon replied UNIMPLEMENTED; permanently falling back "
142+
"to the native Bigtable client for this table."
143+
)
144+
breaker.trip()
145+
raise _AcceleratorFallback() from exc
146+
raise core_exceptions.from_grpc_error(exc) from exc
Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
# Copyright 2026 Google LLC
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
#
15+
"""Single source of truth for which RPCs are routed through the accelerator."""
16+
17+
from __future__ import annotations
18+
19+
# Names mirror the `_DataApiTarget` method names, not the gRPC method names.
20+
# Adding an entry here is not enough on its own: the corresponding method must
21+
# also include a top-of-function branch that dispatches to the accelerator
22+
# service. Keep this set in lockstep with the bundled daemon's capabilities.
23+
_ACCELERATOR_SUPPORTED: frozenset[str] = frozenset({"read_row", "mutate_row"})
24+
25+
26+
def is_supported(method_name: str) -> bool:
27+
return method_name in _ACCELERATOR_SUPPORTED

0 commit comments

Comments
 (0)