diff --git a/.github/workflows/main.yaml b/.github/workflows/main.yaml index 266224164..b685f3367 100644 --- a/.github/workflows/main.yaml +++ b/.github/workflows/main.yaml @@ -106,7 +106,7 @@ jobs: path: coverage-${{ matrix.folder }}.xml pyodide-e2e: - name: Run Pyodide (WASM) Unit + e2e Tests + name: Run Pyodide Unit + e2e Tests runs-on: ubuntu-latest timeout-minutes: 15 # The interpreter comes from the pyodide pin in ci/pyodide-e2e/package.json. @@ -128,19 +128,13 @@ jobs: with: username: ${{secrets.DOCKER_USERNAME}} password: ${{secrets.DOCKER_PASSWORD}} - - name: Build pure wheels (base client + grpc-web) + - name: Build weaviate-client and weaviate-client-web wheels run: | pip install build python -m build --wheel --outdir dist . python -m build --wheel --outdir dist packages/web - name: Assert both wheels carry the same version - # weaviate-client-web is versioned in lockstep with weaviate-client: both derive - # from the same git tag via setuptools_scm. - run: | - base=$(basename dist/weaviate_client-*.whl); base=${base#weaviate_client-}; base=${base%%-*} - web=$(basename dist/weaviate_client_web-*.whl); web=${web#weaviate_client_web-}; web=${web%%-*} - echo "weaviate-client=$base weaviate-client-web=$web" - test "$base" = "$web" + run: ci/assert-lockstep.sh dist - name: Run the unit suite inside Pyodide under Node # weaviate-client-web imports pyodide, so its unit tests run only here. # JSPI lets pytest run async tests. @@ -354,7 +348,7 @@ jobs: cache: 'pip' # caching pip dependencies - name: Install dependencies run: pip install -r requirements-test.txt -r requirements-devel.txt - - name: Build binary wheels (base client + grpc-web companion) + - name: Build weaviate-client and weaviate-client-web wheels run: | python -m build python -m build --wheel --outdir dist packages/web @@ -428,7 +422,7 @@ jobs: cache: 'pip' # caching pip dependencies - name: Install dependencies run: pip install -r requirements-devel.txt - - name: Build distributions (base client + grpc-web companion) + - name: Build distributions (weaviate-client sdist+wheel, weaviate-client-web wheel) # weaviate-client-web is wheel-only: its version comes from git tags, which an # sdist lacks. run: | @@ -437,11 +431,7 @@ jobs: - name: Assert both packages carry the same version # weaviate-client-web pins weaviate-client==; a mismatch leaves the # [grpc-web] extra unresolvable. - run: | - base=$(basename dist/weaviate_client-*.whl); base=${base#weaviate_client-}; base=${base%%-*} - web=$(basename dist/weaviate_client_web-*.whl); web=${web#weaviate_client_web-}; web=${web%%-*} - echo "weaviate-client=$base weaviate-client-web=$web" - test "$base" = "$web" + run: ci/assert-lockstep.sh dist - name: Publish distributions 📦 to PyPI on new tags if: startsWith(github.ref, 'refs/tags') uses: pypa/gh-action-pypi-publish@cef221092ed1bacb1cc03d23a2d87d1d172e277b # release/v1 diff --git a/ci/assert-lockstep.sh b/ci/assert-lockstep.sh new file mode 100755 index 000000000..c97797599 --- /dev/null +++ b/ci/assert-lockstep.sh @@ -0,0 +1,12 @@ +#!/usr/bin/env bash +# Fails unless the weaviate-client and weaviate-client-web wheels in (default: +# dist) carry the same version; the two packages are released in lockstep. +set -euo pipefail +dist="${1:-dist}" +base=$(basename "$dist"/weaviate_client-*.whl); base=${base#weaviate_client-}; base=${base%%-*} +web=$(basename "$dist"/weaviate_client_web-*.whl); web=${web#weaviate_client_web-}; web=${web%%-*} +echo "weaviate-client=$base weaviate-client-web=$web" +if [ "$base" != "$web" ]; then + echo "version mismatch: weaviate-client $base != weaviate-client-web $web (must release in lockstep)" >&2 + exit 1 +fi diff --git a/ci/pyodide-e2e/package.json b/ci/pyodide-e2e/package.json index fbf78f303..675a1f12c 100644 --- a/ci/pyodide-e2e/package.json +++ b/ci/pyodide-e2e/package.json @@ -1,7 +1,7 @@ { "name": "weaviate-pyodide-e2e", "private": true, - "description": "Runs the weaviate-client e2e suite inside Pyodide (WASM) under Node", + "description": "Runs the weaviate-client-web unit and e2e suites inside Pyodide under Node", "dependencies": { "pyodide": "314.0.4" } diff --git a/ci/pyodide-e2e/run.mjs b/ci/pyodide-e2e/run.mjs index 703d0f05b..56041a736 100644 --- a/ci/pyodide-e2e/run.mjs +++ b/ci/pyodide-e2e/run.mjs @@ -2,32 +2,16 @@ // Usage: node run.mjs (one weaviate_client-*.whl, one weaviate_client_web-*.whl) // Env: WEAVIATE_HOST (default localhost), WEAVIATE_PORT (default 8090). // The pyodide npm pin in package.json fixes the interpreter. -import { readdirSync, readFileSync } from "node:fs"; +import { readFileSync } from "node:fs"; import { dirname, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import { loadPyodide } from "pyodide"; -if (!process.argv[2]) { - console.error("usage: node run.mjs "); - process.exit(2); -} -const wheelsDir = resolve(process.argv[2]); -const here = dirname(fileURLToPath(import.meta.url)); +import { wheelsFromArgv } from "./wheels.mjs"; -const wheels = readdirSync(wheelsDir) - .filter((f) => f.endsWith(".whl")) - .sort(); // installs weaviate_client before weaviate_client_web, which depends on it -const prefixes = ["weaviate_client-", "weaviate_client_web-"]; -if ( - wheels.length !== 2 || - !prefixes.every((p) => wheels.some((w) => w.startsWith(p))) -) { - console.error( - `expected exactly one weaviate_client-*.whl and one weaviate_client_web-*.whl in ${wheelsDir}, found: ${JSON.stringify(wheels)}`, - ); - process.exit(2); -} +const { wheelsDir, wheels } = wheelsFromArgv("run.mjs"); +const here = dirname(fileURLToPath(import.meta.url)); const pyodide = await loadPyodide({ env: { @@ -69,3 +53,5 @@ try { console.error(err); process.exit(1); } +// The interpreter keeps live handles on the Node event loop, so exit explicitly. +process.exit(0); diff --git a/ci/pyodide-e2e/units.mjs b/ci/pyodide-e2e/units.mjs index 469a3ccee..44c43177e 100644 --- a/ci/pyodide-e2e/units.mjs +++ b/ci/pyodide-e2e/units.mjs @@ -2,32 +2,15 @@ // bootstrap check. No Weaviate needed. // Usage: node --experimental-wasm-jspi units.mjs (same wheels as run.mjs) // JSPI is required: pytest runs async tests through run_until_complete (stack switching). -import { readdirSync } from "node:fs"; import { dirname, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import { loadPyodide } from "pyodide"; -if (!process.argv[2]) { - console.error("usage: node units.mjs "); - process.exit(2); -} -const wheelsDir = resolve(process.argv[2]); -const here = dirname(fileURLToPath(import.meta.url)); +import { wheelsFromArgv } from "./wheels.mjs"; -const wheels = readdirSync(wheelsDir) - .filter((f) => f.endsWith(".whl")) - .sort(); // installs weaviate_client before weaviate_client_web, which depends on it -const prefixes = ["weaviate_client-", "weaviate_client_web-"]; -if ( - wheels.length !== 2 || - !prefixes.every((p) => wheels.some((w) => w.startsWith(p))) -) { - console.error( - `expected exactly one weaviate_client-*.whl and one weaviate_client_web-*.whl in ${wheelsDir}, found: ${JSON.stringify(wheels)}`, - ); - process.exit(2); -} +const { wheelsDir, wheels } = wheelsFromArgv("units.mjs"); +const here = dirname(fileURLToPath(import.meta.url)); // Fresh interpreter with micropip ready and the wheels dir mounted. async function freshPyodide() { @@ -50,8 +33,8 @@ async function freshPyodide() { pyodide.runPython(` import sys assert "weaviate_client_web" not in sys.modules -import weaviate # the ONLY weaviate-side import: must bootstrap the companion -assert "weaviate_client_web" in sys.modules, "hook did not import the companion" +import weaviate # first weaviate import: installs weaviate_client_web +assert "weaviate_client_web" in sys.modules, "hook did not import weaviate_client_web" import weaviate_client_web import grpc import httpx @@ -62,9 +45,9 @@ assert getattr( httpx.AsyncHTTPTransport.handle_async_request, "__weaviate_fetch_shim__", False ) is True `); - console.log("OK scenario: bare 'import weaviate' bootstraps the companion"); + console.log("OK scenario: bare 'import weaviate' bootstraps weaviate_client_web"); } catch (err) { - console.error("FAIL scenario: bare 'import weaviate' bootstraps the companion"); + console.error("FAIL scenario: bare 'import weaviate' bootstraps weaviate_client_web"); console.error(err); process.exit(1); } diff --git a/ci/pyodide-e2e/wheels.mjs b/ci/pyodide-e2e/wheels.mjs new file mode 100644 index 000000000..79fd748e1 --- /dev/null +++ b/ci/pyodide-e2e/wheels.mjs @@ -0,0 +1,26 @@ +// Shared by run.mjs and units.mjs: resolves (argv[2]) to exactly the two +// locally-built pure wheels, weaviate_client-*.whl and weaviate_client_web-*.whl. +import { readdirSync } from "node:fs"; +import { resolve } from "node:path"; + +export function wheelsFromArgv(script) { + if (!process.argv[2]) { + console.error(`usage: node ${script} `); + process.exit(2); + } + const wheelsDir = resolve(process.argv[2]); + const wheels = readdirSync(wheelsDir) + .filter((f) => f.endsWith(".whl")) + .sort(); // installs weaviate_client before weaviate_client_web, which depends on it + const prefixes = ["weaviate_client-", "weaviate_client_web-"]; + if ( + wheels.length !== 2 || + !prefixes.every((p) => wheels.some((w) => w.startsWith(p))) + ) { + console.error( + `expected exactly one weaviate_client-*.whl and one weaviate_client_web-*.whl in ${wheelsDir}, found: ${JSON.stringify(wheels)}`, + ); + process.exit(2); + } + return { wheelsDir, wheels }; +} diff --git a/packages/web/README.md b/packages/web/README.md index 2331127b2..ec7089fec 100644 --- a/packages/web/README.md +++ b/packages/web/README.md @@ -79,15 +79,16 @@ Weaviate Cloud: browser use requires **Allow all CORS origins** in the cluster's in the Weaviate Cloud console (takes a few minutes to apply). Until it applies, the client fails at its first REST call with a `Failed to fetch` connection error. -`use_async_with_custom()` still requires `grpc_host`/`grpc_port`/`grpc_secure` — Python -cannot drop required parameters on one platform the way TypeScript drops them from a -type. Pass the HTTP values; anything else is overridden with them and warned about -(`Con006`), so a browser client never silently points somewhere it cannot reach. +`use_async_with_custom()` still requires `grpc_host`/`grpc_port`/`grpc_secure`. Under +Pyodide they are ignored: gRPC always uses the REST endpoint, and values that differ from +the HTTP ones are reported in a `Con006` warning. Code that also runs on CPython should +keep its native gRPC values (there, the same host and port for both is an error) and may +ignore that warning. ```python client = weaviate.use_async_with_custom( http_host="localhost", http_port=8080, http_secure=False, - grpc_host="localhost", grpc_port=8080, grpc_secure=False, # = the HTTP endpoint + grpc_host="localhost", grpc_port=50051, grpc_secure=False, # CPython; ignored under Pyodide ) ``` @@ -106,7 +107,7 @@ Importing `weaviate_client_web` before `weaviate` is equivalent. | Bulk insert: `collection.data.insert_many()` | unary gRPC | Yes, the bulk path under Pyodide | | `batch.stream()` / `batch.experimental()` (BatchStream) | bidi streaming | No: grpc-web has no bidirectional streaming; raises at once, use `insert_many()` | | `batch.dynamic()` / `fixed_size()` / `rate_limit()` | sync-client API | No: sync client only | -| Embedded Weaviate (`use_async_with_embedded`) | subprocess | No: raises "not supported under WebAssembly/Pyodide" | +| Embedded Weaviate (`use_async_with_embedded`) | subprocess | No: raises "not supported under Pyodide" | | Synchronous client | — | No: async only | ## Configuration not honored in the browser @@ -145,7 +146,9 @@ a bad API key) are reported as `INTERNAL: grpc-web response contained no message instead of the real error. In the browser a CORS-blocked request is indistinguishable from a network failure -(`TypeError: Failed to fetch`), and is retried as UNAVAILABLE. +(`TypeError: Failed to fetch`); the error text names CORS as one possible cause. On a +channel that has not yet received any response such a failure is reported as `UNKNOWN` +and fails at once; after the first response it is `UNAVAILABLE` and retried. ## Testing diff --git a/packages/web/pyproject.toml b/packages/web/pyproject.toml index 0f4872e18..ed7cdf4f0 100644 --- a/packages/web/pyproject.toml +++ b/packages/web/pyproject.toml @@ -4,16 +4,13 @@ build-backend = "setuptools.build_meta" [project] name = "weaviate-client-web" -description = "grpc-web / WASM (Pyodide) transport for the Weaviate Python client" +description = "grpc-web and fetch transports for the Weaviate Python client under Pyodide" readme = "README.md" requires-python = ">=3.10" license = { text = "BSD-3-Clause" } authors = [{ name = "Weaviate", email = "hello@weaviate.io" }] keywords = ["weaviate", "grpc-web", "pyodide", "wasm", "emscripten"] -# Version comes from the repository's git tags via setuptools_scm (root below), so every -# build carries the same version as weaviate-client; CI asserts the built wheels match. -# Dependencies are computed in setup.py: the weaviate-client requirement is pinned to -# that same version at build time, so mismatched pairs cannot resolve at install time. +# version: git tags (setuptools_scm); dependencies: setup.py (lockstep weaviate-client pin) dynamic = ["version", "dependencies"] [project.urls] diff --git a/packages/web/src/weaviate_client_web/__init__.py b/packages/web/src/weaviate_client_web/__init__.py index 1ef01eedf..8565df0b7 100644 --- a/packages/web/src/weaviate_client_web/__init__.py +++ b/packages/web/src/weaviate_client_web/__init__.py @@ -8,6 +8,14 @@ import os import sys +try: + import pyodide # noqa: F401 +except ImportError as exc: + raise ImportError( + "weaviate-client-web only works under Pyodide (sys.platform == 'emscripten'); " + "on CPython the weaviate-client package uses native gRPC and does not need it." + ) from exc + from ._channel import GrpcWebChannel, set_sender from ._httpx_fetch import ( install_fetch_transport, diff --git a/packages/web/src/weaviate_client_web/_channel.py b/packages/web/src/weaviate_client_web/_channel.py index e5fbbe852..fda329e36 100644 --- a/packages/web/src/weaviate_client_web/_channel.py +++ b/packages/web/src/weaviate_client_web/_channel.py @@ -50,10 +50,32 @@ def _encode_timeout(seconds: Optional[float]) -> Optional[str]: return None +# gRPC spec: keys are lower-case [0-9a-z_.-], ASCII values are printable (0x20-0x7E). +_METADATA_KEY_CHARS = frozenset("0123456789abcdefghijklmnopqrstuvwxyz_.-") + + +def _is_legal_metadata(name: str, text: str) -> bool: + return ( + bool(name) + and all(c in _METADATA_KEY_CHARS for c in name) + and all(" " <= c <= "~" for c in text) + ) + + +def _normalize_path_prefix(path_prefix: Optional[str]) -> str: + """One leading slash, no trailing or repeated slashes, no surrounding whitespace. + + ``""`` (also for ``None`` or a blank value) means native gRPC paths. + """ + segments = [part for part in (path_prefix or "").strip().split("/") if part] + return "/" + "/".join(segments) if segments else "" + + def _fold_metadata(headers: Dict[str, str], metadata: Any) -> None: """Fold gRPC call metadata (``[(key, value), ...]``) into fetch headers. - Binary ``-bin`` keys are base64-encoded as grpc-web requires. + Binary ``-bin`` keys are base64-encoded as grpc-web requires. Keys and values outside + the gRPC spec raise ``ValueError`` before any I/O, as native grpcio does. """ if not metadata: return @@ -64,8 +86,9 @@ def _fold_metadata(headers: Dict[str, str], metadata: Any) -> None: text = base64.b64encode(raw).decode("ascii") else: text = value if isinstance(value, str) else str(value) - # This path bypasses h11/grpcio's header validation, so keep their defence here. - if any(c in name or c in text for c in ("\r", "\n", "\0")): + # grpcio's metadata validation, redone here: fetch would reject some of these + # values synchronously, as a transport error. + if not _is_legal_metadata(name, text): raise ValueError(f"Illegal character in gRPC metadata {name!r}") headers[name] = text @@ -125,7 +148,7 @@ def __call__(self, *args: Any, **kwargs: Any) -> Any: # batch.dynamic()/fixed_size()/rate_limit() are sync-only, so not suggested here. raise RuntimeError( f"Bidirectional streaming RPC {self._path!r} (server-side batching / " - "BatchStream) is not supported over grpc-web/fetch. Use " + "BatchStream) is not supported over grpc-web. Use " "collection.data.insert_many() instead of batch.stream()." ) @@ -145,10 +168,11 @@ def __init__( raise ValueError("GrpcWebChannel requires a target (host:port)") scheme = "https" if secure else "http" self._base_url = f"{scheme}://{target}" - # Normalize to a single leading slash and no trailing slash; "" == native path. - cleaned = (path_prefix or "").strip("/") - self._path_prefix = f"/{cleaned}" if cleaned else "" + self._path_prefix = _normalize_path_prefix(path_prefix) self._sender: Sender = sender or get_sender() + # Until a first HTTP response arrives, a fetch rejection is most likely + # deterministic (CORS, wrong host/port) and must not enter the UNAVAILABLE retry loop. + self._got_response = False def unary_unary( self, @@ -197,36 +221,35 @@ async def _unary( timeout = None # None / non-finite: no deadline, server- or client-side else: headers["grpc-timeout"] = grpc_timeout + if timeout is not None and timeout <= 0: + raise _deadline_exceeded(path, timeout) url = self._base_url + self._path_prefix + path framed = encode_message(payload) - # Send. Enforce a client-side deadline (the grpc-timeout header is server-side - # only; pyfetch ignores its timeout arg, so without this a stalled request could - # hang forever). Any transport/parse failure is surfaced as AioRpcError; the only - # non-gRPC error a caller can see is the ValueError from metadata validation - # above, raised before any I/O (as native grpcio does). + # The sender enforces the deadline; grpc-timeout binds only the server. Every failure + # below becomes AioRpcError (CancelledError propagates); only metadata validation + # above, before I/O, raises ValueError. try: - send = self._sender(url, headers, framed, timeout) - if timeout is not None: - status, resp_headers, body = await asyncio.wait_for(send, timeout) - else: - status, resp_headers, body = await send + status, resp_headers, body = await self._sender(url, headers, framed, timeout) except AioRpcError: raise - except asyncio.TimeoutError as exc: - raise AioRpcError( - code=StatusCode.DEADLINE_EXCEEDED, - details=f"grpc-web request to {path} timed out after {timeout}s", - ) from exc - except Exception as exc: # network/transport failure -> retryable UNAVAILABLE + except (asyncio.TimeoutError, TimeoutError) as exc: + raise _deadline_exceeded(path, timeout) from exc + except Exception as exc: # network/transport failure # str() of transport errors can be empty (e.g. httpx.ConnectError) — always # include the exception type so failures stay diagnosable detail = f"{type(exc).__name__}: {exc}" if str(exc) else repr(exc) details = f"grpc-web transport error for {path}: {detail}" - if not self._path_prefix and sys.platform == "emscripten": - details += " " + _no_path_prefix_hint() - raise AioRpcError(code=StatusCode.UNAVAILABLE, details=details) from exc + if sys.platform == "emscripten": + if not self._path_prefix: + details += " " + _no_path_prefix_hint() + if isinstance(exc, OSError): # how pyfetch reports every fetch rejection + details += ". " + _cors_hint() + # retryable UNAVAILABLE only once this channel has reached the server + code = StatusCode.UNAVAILABLE if self._got_response else StatusCode.UNKNOWN + raise AioRpcError(code=code, details=details) from exc + self._got_response = True try: return self._handle_response(status, resp_headers, body, deserialize, url) @@ -300,7 +323,7 @@ def _handle_response( # grpc-message headers were stripped by CORS in the browser. details += ( " and no grpc-status was visible. If this is a cross-origin browser " - "request, configure the grpc-web proxy to send " + "request, configure the server or proxy to send " "'Access-Control-Expose-Headers: grpc-status, grpc-message' so " "trailers-only error responses are readable." ) @@ -311,6 +334,13 @@ def _handle_response( _BODY_EXCERPT_LIMIT = 200 +def _deadline_exceeded(path: str, timeout: Optional[float]) -> AioRpcError: + return AioRpcError( + code=StatusCode.DEADLINE_EXCEEDED, + details=f"grpc-web request to {path} timed out after {timeout}s", + ) + + def _body_excerpt(body: bytes, limit: int = _BODY_EXCERPT_LIMIT) -> str: """Render a short, printable, one-line excerpt of a response body for error details. @@ -332,7 +362,7 @@ def _no_path_prefix_hint() -> str: from weaviate.connect.base import GRPC_WEB_MIN_SERVER_VERSION, GRPC_WEB_SERVER_PATH_PREFIX return ( - "(no grpc_path_prefix set — under WebAssembly the connect helpers route gRPC to " + "(no grpc_path_prefix set — under Pyodide the connect helpers route gRPC to " f"the REST endpoint under '{GRPC_WEB_SERVER_PATH_PREFIX}' by themselves, so use " "one of them; hand-built ConnectionParams must set " f"grpc_path_prefix='{GRPC_WEB_SERVER_PATH_PREFIX}' for Weaviate >= " @@ -341,6 +371,12 @@ def _no_path_prefix_hint() -> str: ) +def _cors_hint() -> str: + from weaviate.connect.base import GRPC_WEB_CORS_HINT # lazy: see _no_path_prefix_hint + + return GRPC_WEB_CORS_HINT + + def _frame_error_to_rpc( http_status: int, url: str, body: bytes, frame_error: Optional[BaseException] ) -> AioRpcError: @@ -365,6 +401,8 @@ def _non_grpc_web_error( Details start with "HTTP ", then the URL and a body excerpt. """ + from weaviate.connect.base import GRPC_WEB_MIN_SERVER_VERSION, GRPC_WEB_SERVER_PATH_PREFIX + truncated = isinstance(frame_error, TruncatedFrameError) what = "not a grpc-web response" if http_status == 200 and frame_error is not None: @@ -380,16 +418,16 @@ def _non_grpc_web_error( # either cause is possible; the channel does not know the server version parts.append( "The grpc-web endpoint does not exist at that path: either this Weaviate " - "server predates 1.38.3, the first release to serve grpc-web natively, or " - "the configured grpc-web path prefix is wrong for the proxy in front of it. " - "Weaviate's native prefix is '/v1/grpc-web'." + f"server predates {GRPC_WEB_MIN_SERVER_VERSION}, the first release to serve " + "grpc-web, or the configured grpc-web path prefix is wrong. Weaviate's prefix " + f"is '{GRPC_WEB_SERVER_PATH_PREFIX}'." ) elif http_status == 405: # a 405 comes only from an existing HTTP route: the prefix points at one parts.append( "An HTTP route answered instead of the grpc-web endpoint (method not " - "allowed): the configured grpc-web path prefix is wrong. Weaviate's native " - "prefix is '/v1/grpc-web'." + "allowed): the configured grpc-web path prefix is wrong. Weaviate's prefix " + f"is '{GRPC_WEB_SERVER_PATH_PREFIX}'." ) elif http_status in (502, 503, 504): parts.append("Weaviate or the proxy in front of it is unavailable.") @@ -402,7 +440,7 @@ def _non_grpc_web_error( parts.append( "Something other than a grpc-web endpoint answered — typically a proxy " "error page or a single-page-app catch-all route serving index.html. Check " - "the grpc-web path prefix (Weaviate's native prefix is '/v1/grpc-web')." + f"the grpc-web path prefix (Weaviate's prefix is '{GRPC_WEB_SERVER_PATH_PREFIX}')." ) parts.append(f"Response body: {_body_excerpt(body)}") diff --git a/packages/web/src/weaviate_client_web/_framing.py b/packages/web/src/weaviate_client_web/_framing.py index 8ee2e61b9..bbf1def4c 100644 --- a/packages/web/src/weaviate_client_web/_framing.py +++ b/packages/web/src/weaviate_client_web/_framing.py @@ -11,6 +11,7 @@ _FLAG_COMPRESSED = 0x01 _KNOWN_FLAGS = _FLAG_TRAILER | _FLAG_COMPRESSED _HEADER = struct.Struct(">BI") # 1 flag byte + 4-byte big-endian length +_SINGLE_VALUE_TRAILERS = ("grpc-status", "grpc-message") class FrameError(ValueError): @@ -55,6 +56,7 @@ def parse_trailers(raw: bytes) -> Dict[str, str]: """Parse a trailer payload into a lower-cased dict; accepts CRLF or LF line ends. Undecodable bytes are replaced, so an odd grpc-message never drops grpc-status. + Conflicting repeats of grpc-status or grpc-message raise FrameError. """ out: Dict[str, str] = {} for line in raw.split(b"\n"): @@ -63,7 +65,10 @@ def parse_trailers(raw: bytes) -> Dict[str, str]: continue key, _, value = line.partition(b":") name = key.strip().decode("utf-8", "replace").lower() - out[name] = value.strip().decode("utf-8", "replace") + text = value.strip().decode("utf-8", "replace") + if name in _SINGLE_VALUE_TRAILERS and out.get(name, text) != text: + raise FrameError(f"conflicting {name} values in the grpc-web trailer frame") + out[name] = text return out diff --git a/packages/web/src/weaviate_client_web/_httpx_fetch.py b/packages/web/src/weaviate_client_web/_httpx_fetch.py index dcf337f7f..142b4ac0f 100644 --- a/packages/web/src/weaviate_client_web/_httpx_fetch.py +++ b/packages/web/src/weaviate_client_web/_httpx_fetch.py @@ -24,6 +24,9 @@ "accept-encoding", "content-length", "transfer-encoding", + # httpx's "python-httpx/x.y"; a User-Agent set on fetch (Firefox honours it) needs a + # CORS preflight that Weaviate's default CORS_ALLOW_HEADERS does not allow + "user-agent", } # fetch has already decoded the body; passing content-encoding/length through would make diff --git a/packages/web/src/weaviate_client_web/_sender.py b/packages/web/src/weaviate_client_web/_sender.py index 472ada397..eea5c1546 100644 --- a/packages/web/src/weaviate_client_web/_sender.py +++ b/packages/web/src/weaviate_client_web/_sender.py @@ -1,14 +1,17 @@ """HTTP senders for the grpc-web transport. -A *sender* is ``async def sender(url, headers, body, timeout) -> (status, headers, body)``. -The default uses ``pyodide.http.pyfetch``; tests inject one with +A *sender* is ``async def sender(url, headers, body, timeout) -> (status, headers, body)``; +it enforces ``timeout`` (seconds, ``None`` = no deadline) and raises ``TimeoutError`` when +it expires. The default uses ``pyodide.http.pyfetch``; tests inject one with :func:`weaviate_client_web.set_sender`. """ -from typing import Awaitable, Callable, Dict, Optional, Tuple +from typing import Any, Awaitable, Callable, Dict, Optional, Tuple from pyodide.http import pyfetch # type: ignore[import-not-found] +from ._httpx_fetch import _abort_signal_ms + Sender = Callable[ [str, Dict[str, str], bytes, Optional[float]], Awaitable[Tuple[int, Dict[str, str], bytes]], @@ -18,13 +21,30 @@ async def pyfetch_sender( url: str, headers: Dict[str, str], body: bytes, timeout: Optional[float] ) -> Tuple[int, Dict[str, str], bytes]: - """Default browser sender. + """POST through pyodide.http.pyfetch; raises ``TimeoutError`` when ``timeout`` expires. - ``pyfetch`` has no timeout parameter of its own; the call deadline is enforced by - ``GrpcWebChannel._unary`` via ``asyncio.wait_for``. + The deadline aborts the fetch; its JS timer is cleared when the call ends. """ - response = await pyfetch(url, method="POST", headers=headers, body=body) - data = await response.bytes() + kwargs: Dict[str, Any] = {} + controller = timer = clear_timeout = None + delay_ms = _abort_signal_ms(timeout) + if delay_ms is not None: + from js import AbortController, clearTimeout, setTimeout # type: ignore[import-not-found] + + controller = AbortController.new() + clear_timeout = clearTimeout + timer = setTimeout(controller.abort.bind(controller), delay_ms) + kwargs["signal"] = controller.signal + try: + response = await pyfetch(url, method="POST", headers=headers, body=body, **kwargs) + data = await response.bytes() + except OSError as exc: # pyodide.http.AbortError included + if controller is not None and controller.signal.aborted: + raise TimeoutError(f"deadline of {timeout}s exceeded") from exc + raise + finally: + if timer is not None and clear_timeout is not None: + clear_timeout(timer) try: resp_headers = dict(response.headers) except Exception: # pragma: no cover - header shape varies across Pyodide versions diff --git a/packages/web/src/weaviate_client_web/_shim.py b/packages/web/src/weaviate_client_web/_shim.py index e8c7aa768..6176cb4ce 100644 --- a/packages/web/src/weaviate_client_web/_shim.py +++ b/packages/web/src/weaviate_client_web/_shim.py @@ -132,7 +132,7 @@ def first_version_is_lower(_version: str, _other: str) -> bool: _ASYNC_ONLY_MESSAGE = ( "weaviate-client-web provides an asynchronous-only gRPC transport under " - "WebAssembly/Pyodide. Use an async client (weaviate.use_async_with_local / " + "Pyodide. Use an async client (weaviate.use_async_with_local / " "use_async_with_weaviate_cloud / use_async_with_custom, or WeaviateAsyncClient); " "the synchronous client is not supported in the browser." ) diff --git a/packages/web/tests/test_framing.py b/packages/web/tests/test_framing.py index d2f0e745d..7522bc720 100644 --- a/packages/web/tests/test_framing.py +++ b/packages/web/tests/test_framing.py @@ -131,3 +131,22 @@ def test_compressed_trailer_frame_rejected(): body = _frame(b"grpc-status:0\r\n", 0x81) with pytest.raises(FrameError, match="compressed"): split_response(body) + + +@pytest.mark.parametrize( + "raw", + [ + b"grpc-status:13\r\ngrpc-message:boom\r\ngrpc-status:0\r\n", + b"grpc-status:0\r\ngrpc-message:a\r\ngrpc-message:b\r\n", + ], +) +def test_conflicting_duplicate_status_or_message_rejected(raw): + # last-value-wins would read "13 ... 0" as OK + with pytest.raises(FrameError, match="conflicting grpc-"): + parse_trailers(raw) + + +def test_identical_duplicate_status_accepted(): + parsed = parse_trailers(b"grpc-status:5\r\nGrpc-Status: 5\r\ngrpc-message:gone\r\n") + assert parsed["grpc-status"] == "5" + assert parsed["grpc-message"] == "gone" diff --git a/packages/web/tests/test_httpx_fetch.py b/packages/web/tests/test_httpx_fetch.py index fe6adb2be..3b9c4c7ef 100644 --- a/packages/web/tests/test_httpx_fetch.py +++ b/packages/web/tests/test_httpx_fetch.py @@ -171,6 +171,14 @@ async def test_fetch_managed_request_headers_stripped(fake_pyfetch): assert managed not in sent +async def test_httpx_user_agent_is_not_forwarded(fake_pyfetch): + # httpx.AsyncClient adds "user-agent: python-httpx/..." to every request + async with httpx.AsyncClient(transport=_FetchTransport()) as client: + await client.get("http://h:8080/v1/meta") + sent = fake_pyfetch.calls[0]["headers"] + assert all(key.lower() != "user-agent" for key in sent), sent + + async def test_get_without_body_omits_body_kwarg(fake_pyfetch): # fetch rejects GET/HEAD requests that carry a body, so the kwarg must be absent await _handle(httpx.Request("GET", "http://h:8080/v1/.well-known/ready")) @@ -233,9 +241,7 @@ async def test_read_timeout_maps_to_abort_signal_ms(fake_pyfetch, fake_abort_sig assert fake_pyfetch.calls[0]["signal"] == "signal-30000" -async def test_read_none_means_no_deadline_even_with_pool_and_connect_set( - fake_pyfetch, fake_abort_signal -): +async def test_read_none_is_no_deadline(fake_pyfetch, fake_abort_signal): # a non-finite timeout arrives as read=None with pool set; pool must not become the deadline await _handle(_request_with_timeout({"connect": None, "read": None, "write": None, "pool": 5})) await _handle(_request_with_timeout({"connect": 2.0, "read": None, "write": None, "pool": 9.0})) @@ -339,9 +345,7 @@ async def test_fetch_abort_with_deadline_maps_to_read_timeout(monkeypatch, fake_ ) -async def test_fetch_failure_with_deadline_but_no_timeout_message_stays_connect_error( - monkeypatch, fake_abort_signal -): +async def test_network_failure_with_deadline_is_connect_error(monkeypatch, fake_abort_signal): # nearly every weaviate request sets a read deadline; a plain network failure on # such a request must remain a connection error, not become a timeout _install_raising_pyfetch(monkeypatch, OSError("TypeError: Failed to fetch")) diff --git a/packages/web/tests/test_shim_install.py b/packages/web/tests/test_shim_install.py index 710d4f164..ff141c7ee 100644 --- a/packages/web/tests/test_shim_install.py +++ b/packages/web/tests/test_shim_install.py @@ -20,7 +20,7 @@ def test_import_weaviate_under_shim(): assert weaviate_client_web.is_installed() assert weaviate_client_web.install() is True # idempotent, reports the shim in place assert getattr(grpc, "__weaviate_client_web_shim__", False) is True - assert grpc.__version__ == "1.72.1" + assert grpc.__version__ == FAKE_GRPC_VERSION assert grpc._utilities.first_version_is_lower("1.0.0", "2.0.0") is False # type: ignore[attr-defined] from grpc.aio._typing import ChannelArgumentType # noqa: F401 diff --git a/packages/web/tests/test_transport.py b/packages/web/tests/test_transport.py index a9c85bd0c..787fd63ef 100644 --- a/packages/web/tests/test_transport.py +++ b/packages/web/tests/test_transport.py @@ -1,14 +1,17 @@ -"""GrpcWebChannel tests through fake senders (no network, no Weaviate).""" +"""GrpcWebChannel and pyfetch_sender tests through fake senders and a fake pyfetch (no network).""" import asyncio import struct import sys -from typing import Dict, List, Optional, Tuple +import types +from typing import Any, Dict, List, Optional, Tuple import pytest +import weaviate_client_web._sender as _sender_mod from weaviate_client_web import GrpcWebChannel, set_sender from weaviate_client_web._channel import _body_excerpt, _encode_timeout +from weaviate_client_web._httpx_fetch import _MAX_ABORT_SIGNAL_MS from weaviate_client_web._sender import pyfetch_sender from weaviate_client_web._shim import ( AioChannel, @@ -167,7 +170,7 @@ async def test_weaviate_404_json_names_both_candidate_causes(): assert "malformed grpc-web response" not in details -async def test_nginx_502_maps_to_unavailable_so_the_client_retries(): +async def test_nginx_502_maps_to_unavailable(): # weaviate/retry.py retries only UNAVAILABLE err = await _details_of(502, NGINX_502_HTML) assert err.code() is StatusCode.UNAVAILABLE @@ -202,7 +205,7 @@ async def test_405_names_the_wrong_prefix(): assert "method POST is not allowed" in err.details() -async def test_truncated_grpc_web_body_is_reported_as_truncated_not_as_wrong_prefix(): +async def test_truncated_body_is_reported_as_truncated(): # a valid frame header with a cut-short payload: the endpoint is grpc-web, so no # SPA / path-prefix hint body = _ok_response(b"reply-bytes")[:-6] @@ -231,6 +234,14 @@ async def test_message_frame_after_trailer_is_internal(): assert "single-page-app" not in err.details() +async def test_conflicting_grpc_status_in_one_trailer_is_internal(): + body = _frame(b"r") + _frame(b"grpc-status:13\r\ngrpc-message:x\r\ngrpc-status:0\r\n", 0x80) + err = await _details_of(200, body) + assert err.code() is StatusCode.INTERNAL + assert "malformed grpc-web response" in err.details() + assert "conflicting grpc-status" in err.details() + + async def test_multiple_message_frames_in_unary_response_is_internal(): # a unary RPC has exactly one message; silently taking the first would hide a proxy # or server that streams several @@ -329,19 +340,46 @@ def test_stream_stream_raises_clear_error(): mc(request_iterator=iter([]), timeout=5, metadata=None) -async def test_timeout_maps_to_deadline_exceeded(): - async def slow_sender(url, headers, body, timeout): - await asyncio.sleep(0.5) - return 200, {}, _ok_response(b"x") +@pytest.mark.parametrize("error", [TimeoutError, asyncio.TimeoutError]) +async def test_sender_timeout_maps_to_deadline_exceeded(error): + # senders enforce the deadline and raise TimeoutError when it expires + async def expired(url, headers, body, timeout): + raise error("deadline") - channel = GrpcWebChannel("h:1", secure=False, sender=slow_sender) + channel = GrpcWebChannel("h:1", secure=False, sender=expired) mc = channel.unary_unary("/svc/M", lambda x: x, lambda b: b) with pytest.raises(AioRpcError) as excinfo: await mc(b"q", timeout=0.01) assert excinfo.value.code() is StatusCode.DEADLINE_EXCEEDED + assert excinfo.value.details() == "grpc-web request to /svc/M timed out after 0.01s" + + +@pytest.mark.parametrize("timeout", [0, -1]) +async def test_zero_or_negative_timeout_fails_at_once_without_sending(timeout): + sender = FakeSender(body=_ok_response(b"x")) + mc = _channel(sender).unary_unary("/svc/M", lambda x: x, lambda b: b) + with pytest.raises(AioRpcError) as excinfo: + await mc(b"q", timeout=timeout) + assert excinfo.value.code() is StatusCode.DEADLINE_EXCEEDED + assert excinfo.value.details() == f"grpc-web request to /svc/M timed out after {timeout}s" + assert sender.calls == [] + + +async def test_cancellation_while_sending_propagates_unchanged(): + # CancelledError is a BaseException: it must not become a retried UNAVAILABLE + async def cancelled(url, headers, body, timeout): + raise asyncio.CancelledError() + + mc = GrpcWebChannel("h:1", secure=False, sender=cancelled).unary_unary( + "/svc/M", lambda x: x, lambda b: b + ) + with pytest.raises(asyncio.CancelledError): + await mc(b"q", timeout=5) -async def test_transport_exception_maps_to_unavailable(): +async def test_transport_exception_before_any_response_is_unknown(): + # a first-call fetch rejection is usually deterministic (CORS, wrong host/port): + # UNKNOWN fails at once instead of entering the UNAVAILABLE retry loop async def boom(url, headers, body, timeout): raise ConnectionError("connection refused") @@ -349,10 +387,64 @@ async def boom(url, headers, body, timeout): mc = channel.unary_unary("/svc/M", lambda x: x, lambda b: b) with pytest.raises(AioRpcError) as excinfo: await mc(b"q") - assert excinfo.value.code() is StatusCode.UNAVAILABLE + assert excinfo.value.code() is StatusCode.UNKNOWN assert "ConnectionError: connection refused" in str(excinfo.value.details()) +def _scripted_sender(*outcomes: Any): + remaining = list(outcomes) + + async def sender(url, headers, body, timeout): + outcome = remaining.pop(0) + if isinstance(outcome, BaseException): + raise outcome + return outcome + + return sender + + +async def test_transport_exception_after_a_response_is_unavailable(): + # once the channel has reached the server, a drop is transient and stays retryable + sender = _scripted_sender((200, {}, _ok_response(b"x")), ConnectionError("gone")) + mc = GrpcWebChannel("h:1", secure=False, sender=sender).unary_unary( + "/svc/M", lambda x: x, lambda b: b + ) + assert await mc(b"q") == b"x" + with pytest.raises(AioRpcError) as excinfo: + await mc(b"q") + assert excinfo.value.code() is StatusCode.UNAVAILABLE + + +async def test_an_error_response_also_counts_as_reaching_the_server(): + sender = _scripted_sender((404, {}, b"not found"), ConnectionError("gone")) + mc = GrpcWebChannel("h:1", secure=False, sender=sender).unary_unary( + "/svc/M", lambda x: x, lambda b: b + ) + with pytest.raises(AioRpcError): + await mc(b"q") + with pytest.raises(AioRpcError) as excinfo: + await mc(b"q") + assert excinfo.value.code() is StatusCode.UNAVAILABLE + + +@pytest.mark.parametrize("platform,expect_hint", [("emscripten", True), ("linux", False)]) +async def test_fetch_rejection_mentions_cors_under_emscripten(monkeypatch, platform, expect_hint): + # a CORS block and a dead port reject fetch identically; name CORS as a possibility + async def boom(url, headers, body, timeout): + raise OSError("TypeError: fetch failed") + + monkeypatch.setattr(sys, "platform", platform) + channel = GrpcWebChannel("h:1", secure=False, sender=boom, path_prefix="/v1/grpc-web") + mc = channel.unary_unary("/svc/M", lambda x: x, lambda b: b) + with pytest.raises(AioRpcError) as excinfo: + await mc(b"q") + details = excinfo.value.details() + assert "fetch failed" in details + assert ("CORS_ALLOW_ORIGIN" in details) is expect_hint + assert ("CORS_ALLOW_HEADERS" in details) is expect_hint + assert ("fetch failed. In a browser" in details) is expect_hint + + async def test_transport_exception_with_empty_str_keeps_type(): # httpx transport errors commonly stringify to '' — the detail must still name them async def boom(url, headers, body, timeout): @@ -386,7 +478,7 @@ async def test_empty_ok_response_with_grpc_status_has_no_cors_hint(): assert "Access-Control-Expose-Headers" not in str(excinfo.value.details()) -async def test_message_frame_without_grpc_status_is_internal_not_success(): +async def test_message_frame_without_grpc_status_is_internal(): # a message frame without grpc-status (dropped trailer) is INTERNAL channel = _channel(FakeSender(status=200, headers={}, body=_frame(b"reply-bytes"))) mc = channel.unary_unary("/svc/M", lambda x: x, lambda b: b) @@ -513,7 +605,38 @@ async def test_crlf_in_metadata_rejected(bad): assert sender.calls == [] -async def _unavailable_details(monkeypatch, path_prefix, platform): +@pytest.mark.parametrize( + "key,value", + [ + ("x-name", "caf\u00e9"), # Latin-1: fetch accepts it, gRPC does not + ("x-name", "\u2603"), # not Latin-1: Request.new would throw synchronously + ("x-name", "tab\there"), + ("x key", "v"), + ("x-key:", "v"), + ("", "v"), + ], +) +async def test_metadata_outside_the_grpc_spec_is_a_value_error_before_sending(key, value): + sender = FakeSender(body=_ok_response(b"x")) + mc = _channel(sender).unary_unary("/svc/M", lambda x: x, lambda b: b) + with pytest.raises(ValueError, match="Illegal character"): + await mc(b"q", metadata=[(key, value)]) + assert sender.calls == [] + + +async def test_metadata_with_printable_ascii_and_binary_values_is_sent(): + sender = FakeSender(body=_ok_response(b"x")) + mc = _channel(sender).unary_unary("/svc/M", lambda x: x, lambda b: b) + await mc( + b"q", + metadata=[("X-OpenAI-Api-Key", "sk-A_b.c ~!"), ("name-bin", "caf\u00e9".encode())], + ) + headers = sender.calls[0][1] + assert headers["x-openai-api-key"] == "sk-A_b.c ~!" + assert headers["name-bin"] == "Y2Fmw6k=" + + +async def _transport_error_details(monkeypatch, path_prefix, platform): async def boom(url, headers, body, timeout): raise ConnectionError("Failed to fetch") @@ -522,31 +645,29 @@ async def boom(url, headers, body, timeout): mc = channel.unary_unary("/svc/M", lambda x: x, lambda b: b) with pytest.raises(AioRpcError) as excinfo: await mc(b"q") - assert excinfo.value.code() is StatusCode.UNAVAILABLE + assert excinfo.value.code() is StatusCode.UNKNOWN # fresh channel: no response yet return excinfo.value.details() -async def test_unavailable_without_path_prefix_under_emscripten_hints_at_grpc_path_prefix( - monkeypatch, -): +async def test_transport_error_without_prefix_hints_prefix_under_emscripten(monkeypatch): # the connect helpers always set the prefix under Emscripten, so a prefix-less channel # here means hand-built ConnectionParams; the error must say what to do instead - details = await _unavailable_details(monkeypatch, path_prefix="", platform="emscripten") + details = await _transport_error_details(monkeypatch, path_prefix="", platform="emscripten") assert "grpc_path_prefix='/v1/grpc-web'" in details assert "1.38.3" in details assert "connect helpers" in details -async def test_unavailable_with_path_prefix_has_no_prefix_hint(monkeypatch): - details = await _unavailable_details( +async def test_transport_error_with_prefix_has_no_prefix_hint(monkeypatch): + details = await _transport_error_details( monkeypatch, path_prefix="/v1/grpc-web", platform="emscripten" ) assert "no grpc_path_prefix" not in details -async def test_unavailable_without_path_prefix_off_emscripten_has_no_prefix_hint(monkeypatch): +async def test_transport_error_without_prefix_off_emscripten_has_no_prefix_hint(monkeypatch): # off Emscripten an empty prefix against a transcoder is the normal configuration - details = await _unavailable_details(monkeypatch, path_prefix="", platform="linux") + details = await _transport_error_details(monkeypatch, path_prefix="", platform="linux") assert "no grpc_path_prefix" not in details @@ -572,6 +693,9 @@ async def test_path_prefix_prepended_to_url(): ("/grpc-web/", "http://h:1/grpc-web/svc/M"), ("/a/b", "http://h:1/a/b/svc/M"), ("", "http://h:1/svc/M"), + (" ", "http://h:1/svc/M"), + ("//a//b/", "http://h:1/a/b/svc/M"), + (" grpc-web/ ", "http://h:1/grpc-web/svc/M"), ], ) async def test_path_prefix_normalized_in_url(raw, expected_url): @@ -629,3 +753,154 @@ async def test_metadata_cannot_replace_protocol_headers(): assert headers["x-grpc-web"] == "1" assert headers["x-user-agent"] == "weaviate-client-web" assert headers["x-custom"] == "kept" + + +# --- the default pyfetch sender: deadline via AbortController ---------------------- + + +class _FakeSignal: + def __init__(self) -> None: + self.aborted = False + + +class _FakeController: + def __init__(self) -> None: + self.signal = _FakeSignal() + # JS: controller.abort.bind(controller) + self.abort = types.SimpleNamespace(bind=lambda _this: self._do_abort) + + def _do_abort(self) -> None: + self.signal.aborted = True + + +class _FakeJs: + """Stand-in for the ``js`` module: records timers and fires them on demand.""" + + def __init__(self) -> None: + self.controllers: List[_FakeController] = [] + self.timers: Dict[int, Any] = {} + self.delays: List[int] = [] + self._next = 0 + self.AbortController = types.SimpleNamespace(new=self._new_controller) + + def _new_controller(self) -> _FakeController: + controller = _FakeController() + self.controllers.append(controller) + return controller + + def setTimeout(self, callback: Any, delay: int) -> int: # noqa: N802 - JS name + self._next += 1 + self.timers[self._next] = callback + self.delays.append(delay) + return self._next + + def clearTimeout(self, timer: int) -> None: # noqa: N802 - JS name + self.timers.pop(timer, None) + + def fire_all(self) -> None: + for callback in list(self.timers.values()): + callback() + + +@pytest.fixture +def fake_js(monkeypatch) -> _FakeJs: + js = _FakeJs() + js_mod = types.ModuleType("js") + for name in ("AbortController", "setTimeout", "clearTimeout"): + setattr(js_mod, name, getattr(js, name)) + monkeypatch.setitem(sys.modules, "js", js_mod) + return js + + +class _FakeFetchResponse: + status = 200 + headers: Dict[str, str] = {} + + async def bytes(self) -> bytes: # noqa: A003 - mirrors pyodide's FetchResponse + return _ok_response(b"ok") + + +def _install_pyfetch(monkeypatch, behaviour) -> List[Dict[str, Any]]: + calls: List[Dict[str, Any]] = [] + + async def fake_pyfetch(url: str, **kwargs: Any) -> Any: + calls.append({"url": url, **kwargs}) + return await behaviour() + + monkeypatch.setattr(_sender_mod, "pyfetch", fake_pyfetch) + return calls + + +async def _ok() -> _FakeFetchResponse: + return _FakeFetchResponse() + + +async def test_pyfetch_sender_aborts_via_signal_and_clears_its_timer(monkeypatch, fake_js): + calls = _install_pyfetch(monkeypatch, _ok) + status, _, body = await pyfetch_sender("http://h/svc/M", {}, b"q", 30) + assert status == 200 and body == _ok_response(b"ok") + assert fake_js.delays == [30_000] + assert calls[0]["signal"] is fake_js.controllers[0].signal + assert fake_js.timers == {} # cleared: no JS timer outlives the request + + +async def test_pyfetch_sender_caps_the_timer_at_int32_ms(monkeypatch, fake_js): + # setTimeout delays above 2^31-1 ms overflow and fire at once + _install_pyfetch(monkeypatch, _ok) + await pyfetch_sender("http://h/svc/M", {}, b"q", 1e9) + assert fake_js.delays == [_MAX_ABORT_SIGNAL_MS] + assert fake_js.timers == {} + + +async def test_pyfetch_sender_without_deadline_sets_no_timer(monkeypatch, fake_js): + calls = _install_pyfetch(monkeypatch, _ok) + await pyfetch_sender("http://h/svc/M", {}, b"q", None) + assert fake_js.delays == [] + assert "signal" not in calls[0] + + +async def test_pyfetch_sender_raises_timeout_error_when_its_timer_fires(monkeypatch, fake_js): + async def aborted(): + fake_js.fire_all() # the deadline passes while the fetch is in flight + raise OSError("AbortError: signal is aborted without reason") + + _install_pyfetch(monkeypatch, aborted) + with pytest.raises(TimeoutError): + await pyfetch_sender("http://h/svc/M", {}, b"q", 0.5) + assert fake_js.timers == {} + + +async def test_pyfetch_sender_keeps_other_fetch_failures_and_clears_its_timer(monkeypatch, fake_js): + async def rejected(): + raise OSError("TypeError: fetch failed") + + _install_pyfetch(monkeypatch, rejected) + with pytest.raises(OSError, match="fetch failed") as excinfo: + await pyfetch_sender("http://h/svc/M", {}, b"q", 30) + assert not isinstance(excinfo.value, TimeoutError) + assert fake_js.timers == {} + + +async def test_pyfetch_sender_clears_its_timer_when_cancelled(monkeypatch, fake_js): + async def cancelled(): + raise asyncio.CancelledError() + + _install_pyfetch(monkeypatch, cancelled) + with pytest.raises(asyncio.CancelledError): + await pyfetch_sender("http://h/svc/M", {}, b"q", 30) + assert fake_js.timers == {} + + +async def test_channel_deadline_through_pyfetch_sender_is_deadline_exceeded(monkeypatch, fake_js): + async def aborted(): + fake_js.fire_all() + raise OSError("AbortError: signal is aborted without reason") + + _install_pyfetch(monkeypatch, aborted) + mc = GrpcWebChannel("h:1", secure=False, sender=pyfetch_sender).unary_unary( + "/svc/M", lambda x: x, lambda b: b + ) + with pytest.raises(AioRpcError) as excinfo: + await mc(b"q", timeout=0.5) + assert excinfo.value.code() is StatusCode.DEADLINE_EXCEEDED + assert excinfo.value.details() == "grpc-web request to /svc/M timed out after 0.5s" diff --git a/proto_test/test_proto.py b/proto_test/test_proto.py index e7f1a28b0..7ee7c9cad 100644 --- a/proto_test/test_proto.py +++ b/proto_test/test_proto.py @@ -23,9 +23,7 @@ def _versions_incompatible() -> bool: _skip_if_incompatible = pytest.mark.skipif( _versions_incompatible(), - reason="weaviate.proto.v1 cannot be imported with an incompatible grpcio/protobuf " - "pair (CI version-gate matrix); the gate is covered by test_proto_import and the " - "fallback is exercised in every compatible cell", + reason="incompatible grpcio/protobuf pair: weaviate.proto.v1 does not import", ) @@ -36,9 +34,7 @@ def test_proto_import(): pb_ver >= version.parse("5.26.1") and grpc_ver < version.parse("1.63.0") ): with pytest.raises(Exception) as e: - import weaviate - - assert weaviate.version is not None + importlib.import_module("weaviate") assert "WeaviateProtobufIncompatibility" in str(e.type) else: import weaviate @@ -80,7 +76,7 @@ def raises(pkg: str) -> str: @_skip_if_incompatible -def test_grpcio_fallback_version_passes_every_vendored_stub_gate(): +def test_grpcio_fallback_version_is_at_least_every_stub_generated_version(): """_GRPCIO_FALLBACK_VERSION is at least every vendored stub's GRPC_GENERATED_VERSION.""" try: from grpc._utilities import first_version_is_lower @@ -104,8 +100,27 @@ def test_grpcio_fallback_version_passes_every_vendored_stub_gate(): gated += 1 generated = match.group(1) assert not first_version_is_lower(fallback, generated), ( - f"{stub.relative_to(proto_root)} requires grpcio>={generated} but " - f"_GRPCIO_FALLBACK_VERSION is {fallback}; bump the fallback (and the " - "grpc-web shim's FAKE_GRPC_VERSION) to match the regenerated stubs" + f"{stub.relative_to(proto_root)} was generated for grpcio {generated}, above " + f"_GRPCIO_FALLBACK_VERSION {fallback}; bump it and the grpc shim's " + "FAKE_GRPC_VERSION to match" ) assert gated > 0, "no stub carried a GRPC_GENERATED_VERSION gate; check the extraction regex" + + +def test_shim_fake_grpc_version_matches_the_fallback(): + """The grpc shim's FAKE_GRPC_VERSION equals _GRPCIO_FALLBACK_VERSION. + + Read as text: weaviate_client_web imports pyodide, so it cannot be imported here. + """ + repo = pathlib.Path(__file__).resolve().parents[1] + + def literal(path: pathlib.Path, name: str) -> str: + match = re.search(rf'^{name} = "([^"]+)"', path.read_text(), re.MULTILINE) + assert match is not None, f"{name} not found in {path}" + return match.group(1) + + fallback = literal( + repo / "weaviate" / "proto" / "v1" / "__init__.py", "_GRPCIO_FALLBACK_VERSION" + ) + shim = repo / "packages" / "web" / "src" / "weaviate_client_web" / "_shim.py" + assert literal(shim, "FAKE_GRPC_VERSION") == fallback diff --git a/setup.cfg b/setup.cfg index 26b397613..c2c92c149 100644 --- a/setup.cfg +++ b/setup.cfg @@ -45,15 +45,7 @@ install_requires = packaging>=21.0 python_requires = >=3.10 -[options.extras_require] -agents = - weaviate-agents >=1.0.0, <2.0.0 -grpc-web = - # The Pyodide/WASM grpc-web companion, built from packages/web in this repo. Both - # packages derive their version from the same git tag (setuptools_scm; CI asserts - # the built wheels match). The marker makes the extra a no-op on CPython: the - # companion imports pyodide at module scope and only makes sense under Emscripten. - weaviate-client-web; sys_platform == "emscripten" +# extras_require lives in setup.py: the [grpc-web] extra pins this build's version. [options.package_data] # If any package or subpackage contains *.txt, *.rst or *.md files, include them: diff --git a/setup.py b/setup.py index 7f1a1763c..e0cd39f25 100644 --- a/setup.py +++ b/setup.py @@ -1,4 +1,14 @@ from setuptools import setup +from setuptools_scm import get_version + +# [grpc-web] pins weaviate-client-web to this build's version (lockstep: see +# packages/web/setup.py); micropip does not backtrack to find a matching pair. +version = get_version(root=".", relative_to=__file__) if __name__ == "__main__": - setup() + setup( + extras_require={ + "agents": ["weaviate-agents >=1.0.0, <2.0.0"], + "grpc-web": [f'weaviate-client-web=={version}; sys_platform == "emscripten"'], + } + ) diff --git a/test/test_connection_params.py b/test/test_connection_params.py index 0be4b87cf..3aaabe762 100644 --- a/test/test_connection_params.py +++ b/test/test_connection_params.py @@ -50,6 +50,9 @@ def test_from_url_same_host_port_allowed_with_prefix() -> None: ("/grpc-web", "/grpc-web"), ("grpc-web/", "/grpc-web"), ("/a/b/", "/a/b"), + (" ", ""), # blank means native gRPC, not a whitespace path + (" /grpc-web ", "/grpc-web"), + ("//a//b/", "/a/b"), ], ) def test_path_prefix_normalization(raw, expected) -> None: @@ -123,8 +126,9 @@ def test_async_client_construction_rejects_prefix_without_shim(monkeypatch) -> N from weaviate import WeaviateAsyncClient monkeypatch.delattr(base_mod.grpc, "__weaviate_client_web_shim__", raising=False) - with pytest.raises(WeaviateInvalidInputError, match="weaviate-client-web"): + with pytest.raises(WeaviateInvalidInputError, match="weaviate-client-web") as excinfo: WeaviateAsyncClient(_grpc_web_params()) + assert str(excinfo.value).endswith("use native gRPC instead.") # one period, not two def test_async_client_construction_allows_prefix_with_shim(monkeypatch) -> None: @@ -138,8 +142,13 @@ def test_async_client_construction_allows_prefix_with_shim(monkeypatch) -> None: def test_sync_client_construction_rejects_grpc_web_prefix() -> None: from weaviate import WeaviateClient - with pytest.raises(WeaviateInvalidInputError, match="async"): + with pytest.raises(WeaviateInvalidInputError, match="async") as excinfo: WeaviateClient(_grpc_web_params()) + msg = str(excinfo.value) + # use_async_with_custom() has no grpc_path_prefix; point at what does + assert "WeaviateAsyncClient(ConnectionParams.from_params(" in msg + assert "grpc_path_prefix=" in msg + assert "use_async_with_custom" not in msg @pytest.mark.parametrize( @@ -178,7 +187,7 @@ def test_sync_client_construction_rejects_grpc_web_prefix() -> None: ), ], ) -def test_helper_params_off_emscripten_are_unchanged(call, expected) -> None: +def test_helper_params_off_emscripten_use_native_grpc(call, expected) -> None: # off Emscripten the helpers use native gRPC: grpc_path_prefix stays None import weaviate diff --git a/test/test_wasm_compat.py b/test/test_wasm_compat.py index bcb4b9e8e..78f6413fa 100644 --- a/test/test_wasm_compat.py +++ b/test/test_wasm_compat.py @@ -5,6 +5,7 @@ import subprocess import sys import textwrap +import warnings import grpc import pytest @@ -12,7 +13,7 @@ from weaviate import WeaviateClient from weaviate.collections.batch.async_ import _BatchBaseAsync -from weaviate.connect.base import ConnectionParams +from weaviate.connect.base import GRPC_WEB_SERVER_PATH_PREFIX, ConnectionParams from weaviate.connect.v4 import _ConnectionBase from weaviate.embedded import _EmbeddedBase from weaviate.exceptions import ( @@ -25,7 +26,7 @@ def test_embedded_raises_explicit_error_under_emscripten(monkeypatch) -> None: monkeypatch.setattr(sys, "platform", "emscripten") - with pytest.raises(WeaviateStartUpError, match="WebAssembly/Pyodide"): + with pytest.raises(WeaviateStartUpError, match="not supported under Pyodide"): _EmbeddedBase.check_supported_platform() @@ -47,30 +48,37 @@ def test_batch_stream_fails_fast_when_grpc_web_shim_active(monkeypatch) -> None: # --- grpc-web diagnostics ------------------------------------------------------------- -def _connection(prefix=None) -> _ConnectionBase: +def _connection(prefix=None, grpc_port=None) -> _ConnectionBase: conn = object.__new__(_ConnectionBase) conn._client = None conn._grpc_channel = None conn._weaviate_version = _ServerVersion.from_string("1.36.0") + if grpc_port is None: + grpc_port = 8080 if prefix else 50051 conn._connection_params = ConnectionParams.from_url( "http://localhost:8080", - grpc_port=8080 if prefix else 50051, + grpc_port=grpc_port, grpc_path_prefix=prefix, ) return conn +# the shape of the grpc-web channel's details for a response that is not grpc-web +_CHANNEL_404_DETAILS = ( + "HTTP 404 from http://localhost:8080/grpc-web/grpc.health.v1.Health/Check: not a " + "grpc-web response. The grpc-web endpoint does not exist at that path. " + "Response body: 404 page not found" +) + + def _ping_exception(conn: _ConnectionBase, error: Exception) -> None: getattr(conn, "_ConnectionBase__handle_ping_exception")(error) # noqa: B009 -def test_grpc_web_404_names_the_two_real_causes_and_drops_firewall_advice() -> None: +def test_grpc_web_404_names_both_causes_without_port_advice() -> None: conn = _connection(prefix="/grpc-web") error = AioRpcError( - grpc.StatusCode.UNIMPLEMENTED, - Metadata(), - Metadata(), - details="HTTP 404 for /grpc-web/grpc.health.v1.Health/Check: 404 page not found", + grpc.StatusCode.UNIMPLEMENTED, Metadata(), Metadata(), details=_CHANNEL_404_DETAILS ) with pytest.raises(WeaviateGRPCUnavailableError) as excinfo: _ping_exception(conn, error) @@ -79,11 +87,50 @@ def test_grpc_web_404_names_the_two_real_causes_and_drops_firewall_advice() -> N assert "firewall" not in msg assert "port (localhost:8080) are correct" not in msg assert "UNIMPLEMENTED" in msg # the real code, not swallowed - assert "HTTP 404 for /grpc-web/grpc.health.v1.Health/Check" in msg # ... and details + assert "HTTP 404 from http://localhost:8080/grpc-web/grpc.health" in msg # ... and details assert "/grpc-web" in msg # the prefix that was actually used assert "1.38.3" in msg # candidate 1: server too old ... assert "v1.36.0" in msg # ... shown against the observed server version assert "/v1/grpc-web" in msg # candidate 2: wrong prefix + assert "over the REST endpoint localhost:8080" in msg # gRPC shares the REST address + assert "CORS" not in msg # a routed-path problem, not a blocked request + + +def test_grpc_web_405_gets_the_same_wrong_path_diagnosis() -> None: + # a 405 means an HTTP route answered instead of the grpc-web endpoint + conn = _connection(prefix="/grpc-web") + error = AioRpcError( + grpc.StatusCode.UNIMPLEMENTED, + Metadata(), + Metadata(), + details=( + "HTTP 405 from http://localhost:8080/grpc-web/grpc.health.v1.Health/Check: not a " + "grpc-web response. An HTTP route answered instead of the grpc-web endpoint. " + 'Response body: {"code":405,"message":"method POST is not allowed"}' + ), + ) + with pytest.raises(WeaviateGRPCUnavailableError) as excinfo: + _ping_exception(conn, error) + msg = str(excinfo.value) + + assert "did not route the grpc-web path '/grpc-web'" in msg + assert "1.38.3" in msg + assert "/v1/grpc-web" in msg + assert "HTTP 405" in msg + + +def test_grpc_web_on_a_separate_endpoint_is_not_called_the_rest_endpoint() -> None: + # a hand-built prefix can target a transcoder on another port + conn = _connection(prefix="/grpc-web", grpc_port=50290) + error = AioRpcError( + grpc.StatusCode.UNIMPLEMENTED, Metadata(), Metadata(), details=_CHANNEL_404_DETAILS + ) + with pytest.raises(WeaviateGRPCUnavailableError) as excinfo: + _ping_exception(conn, error) + msg = str(excinfo.value) + + assert "REST endpoint" not in msg + assert "over localhost:50290 (grpc-web)" in msg def test_grpc_web_genuine_unimplemented_is_not_diagnosed_as_a_wrong_path() -> None: @@ -106,7 +153,14 @@ def test_grpc_web_genuine_unimplemented_is_not_diagnosed_as_a_wrong_path() -> No def test_grpc_web_non_404_error_still_omits_the_native_port_advice() -> None: conn = _connection(prefix="/grpc-web") error = AioRpcError( - grpc.StatusCode.UNAVAILABLE, Metadata(), Metadata(), details="HTTP 502 for /grpc-web/..." + grpc.StatusCode.UNAVAILABLE, + Metadata(), + Metadata(), + details=( + "HTTP 502 from http://localhost:8080/grpc-web/grpc.health.v1.Health/Check: not a " + "grpc-web response. Weaviate or the proxy in front of it is unavailable. " + "Response body: 502 Bad Gateway" + ), ) with pytest.raises(WeaviateGRPCUnavailableError) as excinfo: _ping_exception(conn, error) @@ -116,9 +170,12 @@ def test_grpc_web_non_404_error_still_omits_the_native_port_advice() -> None: assert "UNAVAILABLE" in msg assert "HTTP 502" in msg assert "skip_init_checks=True" in msg # the still-useful advice is kept + # a CORS block looks like any other fetch failure, so it is named as a possibility + assert "CORS_ALLOW_ORIGIN" in msg + assert "CORS_ALLOW_HEADERS" in msg -def test_native_grpc_message_keeps_its_advice_and_gains_the_real_status() -> None: +def test_native_grpc_message_includes_status() -> None: conn = _connection() error = AioRpcError( grpc.StatusCode.UNAVAILABLE, Metadata(), Metadata(), details="failed to connect" @@ -135,6 +192,28 @@ def test_native_grpc_message_keeps_its_advice_and_gains_the_real_status() -> Non assert "failed to connect" in msg +def test_prefixless_params_under_emscripten_point_at_the_async_helpers_and_the_prefix( + monkeypatch, +) -> None: + # hand-built ConnectionParams without a prefix under Pyodide: the firewall advice and + # weaviate.connect_to_local (a sync helper, which raises there) would mislead + conn = _connection() + monkeypatch.setattr(sys, "platform", "emscripten") + error = AioRpcError( + grpc.StatusCode.UNKNOWN, Metadata(), Metadata(), details="TypeError: fetch failed" + ) + with pytest.raises(WeaviateGRPCUnavailableError) as excinfo: + _ping_exception(conn, error) + msg = str(excinfo.value) + + assert "firewall" not in msg + assert "connect_to_local" not in msg + assert "use_async_with_local" in msg + assert "grpc_path_prefix='/v1/grpc-web'" in msg + assert "localhost:50051" in msg + assert "fetch failed" in msg # the observed error is kept + + def test_non_grpc_ping_error_is_still_reported() -> None: # not every ping failure is an RpcError; those must not lose the generic advice conn = _connection() @@ -145,7 +224,7 @@ def test_non_grpc_ping_error_is_still_reported() -> None: # --- grpc-web routing under Emscripten ------------------------------------------------ -GRPC_WEB_PREFIX = "/v1/grpc-web" +GRPC_WEB_PREFIX = GRPC_WEB_SERVER_PATH_PREFIX @pytest.fixture @@ -244,7 +323,8 @@ def test_overridden_grpc_arguments_are_warned_about(emscripten) -> None: msg = str(record[0].message) assert "grpc.example.com:50051" in msg # what was discarded ... assert "localhost:8080" in msg # ... and what is used instead - assert "WebAssembly" in msg # ... and why + assert "Pyodide" in msg # ... and why + assert "may ignore this warning" in msg # portable advice: no CPython port collision _assert_grpc_rides_rest(client) @@ -274,6 +354,11 @@ def test_an_explicit_local_grpc_port_is_warned_about_but_the_default_is_not(emsc client = weaviate.use_async_with_local(port=8080, grpc_port=8081) _assert_grpc_rides_rest(client) + with warnings.catch_warnings(): + warnings.simplefilter("error") # any warning fails the default case + client = weaviate.use_async_with_local(port=8080) + _assert_grpc_rides_rest(client) + # --- the single-import hook (weaviate/__init__.py) ------------------------------------ # @@ -302,7 +387,7 @@ def _run_hook_scenario( return subprocess.run([*interp, "-c", script], capture_output=True, text=True) -def test_bare_import_without_companion_raises_clear_import_error() -> None: +def test_bare_import_without_weaviate_client_web_raises_import_error() -> None: # No site-packages, so neither weaviate_client_web nor grpcio is importable; the repo # root goes on sys.path so the weaviate package itself is still found. result = _run_hook_scenario( @@ -313,10 +398,10 @@ def test_bare_import_without_companion_raises_clear_import_error() -> None: except ImportError as e: assert "weaviate-client-web" in str(e), str(e) assert "weaviate-client[grpc-web]" in str(e), str(e) - assert "WebAssembly/Pyodide" in str(e), str(e) + assert "Pyodide" in str(e), str(e) print("OK") else: - raise AssertionError("expected ImportError without the companion") + raise AssertionError("expected ImportError without weaviate-client-web") """, no_site=True, ) @@ -344,7 +429,7 @@ def test_bare_import_with_grpc_present_falls_through_silently() -> None: assert "OK" in result.stdout -def test_bare_import_with_broken_companion_surfaces_its_own_error(tmp_path) -> None: +def test_bare_import_with_broken_weaviate_client_web_surfaces_its_own_error(tmp_path) -> None: # an installed weaviate-client-web that fails to import surfaces its own error, not the # install hint fake_pkg = tmp_path / "weaviate_client_web" @@ -364,9 +449,29 @@ def test_bare_import_with_broken_companion_surfaces_its_own_error(tmp_path) -> N assert "grpc-web" not in str(e), str(e) print("OK") else: - raise AssertionError("expected the companion's own ImportError to surface") + raise AssertionError("expected weaviate_client_web's own ImportError to surface") """, path_entry=str(tmp_path), ) assert result.returncode == 0, result.stderr assert "OK" in result.stdout + + +def test_importing_weaviate_client_web_on_cpython_says_it_needs_pyodide() -> None: + # an accidental install on CPython must explain itself, not fail on 'pyodide' + web_src = str(pathlib.Path(_REPO_ROOT) / "packages" / "web" / "src") + result = _run_hook_scenario( + """ + try: + import weaviate_client_web + except ImportError as e: + assert "Pyodide" in str(e), str(e) + assert e.name != "pyodide", e.name + print("OK") + else: + raise AssertionError("expected ImportError outside Pyodide") + """, + path_entry=web_src, + ) + assert result.returncode == 0, result.stderr + assert "OK" in result.stdout diff --git a/weaviate/__init__.py b/weaviate/__init__.py index 751663920..6a991b4f3 100644 --- a/weaviate/__init__.py +++ b/weaviate/__init__.py @@ -16,12 +16,9 @@ raise if find_spec("grpc") is None: raise ImportError( - "weaviate requires the weaviate-client-web package under " - "WebAssembly/Pyodide: there is no grpcio wheel for Emscripten, and " - "weaviate-client-web provides the grpc-web (fetch) transport in its " - "place. Install it via the extra (e.g. " - "micropip.install('weaviate-client[grpc-web]')) and import weaviate " - "again." + "weaviate needs weaviate-client-web under Pyodide (there is no grpcio " + "wheel). Install it with micropip.install('weaviate-client[grpc-web]'), " + "then import weaviate again." ) from exc import os diff --git a/weaviate/collections/batch/async_.py b/weaviate/collections/batch/async_.py index 07083fe9b..92515e4c6 100644 --- a/weaviate/collections/batch/async_.py +++ b/weaviate/collections/batch/async_.py @@ -138,9 +138,8 @@ async def _start(self): # fail early: over grpc-web the BatchStream RPC would fail inside the background # tasks, which shows up as silently dropped objects or a flush() that never ends raise WeaviateBatchStreamError( - "batch.stream() requires bidirectional gRPC streaming, which is not " - "possible over grpc-web/fetch (WebAssembly/Pyodide). Use " - "collection.data.insert_many() instead." + "batch.stream() needs bidirectional gRPC streaming, which grpc-web does " + "not support (Pyodide). Use collection.data.insert_many() instead" ) self.__number_of_nodes = await self.__cluster.get_number_of_nodes() diff --git a/weaviate/connect/base.py b/weaviate/connect/base.py index 5ac4689d8..a652aa647 100644 --- a/weaviate/connect/base.py +++ b/weaviate/connect/base.py @@ -23,6 +23,12 @@ GRPC_WEB_MIN_SERVER_VERSION = "1.38.3" # Weaviate's grpc-web path prefix GRPC_WEB_SERVER_PATH_PREFIX = "/v1/grpc-web" +# appended to fetch failures under Pyodide; a dead port fails the same way, so it is a hint +GRPC_WEB_CORS_HINT = ( + "In a browser, this also happens when the server's CORS policy blocks the request " + "(Weaviate: CORS_ALLOW_ORIGIN / CORS_ALLOW_HEADERS; Weaviate Cloud: enable 'Allow all " + "CORS origins' in the cluster's settings)" +) def _grpc_web_shim_active() -> bool: @@ -138,11 +144,11 @@ def _grpc_target(self) -> str: def _grpc_web_path_prefix(self) -> str: """Normalized grpc-web path prefix; "" means native gRPC. - One leading slash, no trailing slash ("grpc-web/" -> "/grpc-web"); empty or - None -> "". + One leading slash, no trailing or repeated slashes ("grpc-web/" -> "/grpc-web", + "//a//b/" -> "/a/b"); blank or None -> "". """ - cleaned = (self.grpc_path_prefix or "").strip("/") - return f"/{cleaned}" if cleaned else "" + segments = [part for part in (self.grpc_path_prefix or "").strip().split("/") if part] + return "/" + "/".join(segments) if segments else "" def _check_grpc_web_usable(self, is_async: bool) -> None: """Raise if a grpc-web prefix is set but unusable (sync client, or no grpc shim). @@ -153,18 +159,14 @@ def _check_grpc_web_usable(self, is_async: bool) -> None: return if not is_async: raise WeaviateInvalidInputError( - "grpc_path_prefix (grpc-web) is only supported for async clients; " - "use use_async_with_custom(...) / WeaviateAsyncClient" + "grpc_path_prefix (grpc-web) is only supported for async clients; use " + "WeaviateAsyncClient(ConnectionParams.from_params(..., grpc_path_prefix=...))" ) if not _grpc_web_shim_active(): raise WeaviateInvalidInputError( - "grpc_path_prefix enables grpc-web, which requires the " - "'weaviate-client-web' package (it installs a grpc shim before " - "'import weaviate'); it is not active in this environment. grpc-web is " - "only available under WebAssembly/Pyodide, where a plain `import " - "weaviate` activates it (install the companion with " - "micropip.install('weaviate-client[grpc-web]')); on CPython use native " - "gRPC instead." + "grpc_path_prefix (grpc-web) needs weaviate-client-web, which runs only " + "under Pyodide (micropip.install('weaviate-client[grpc-web]')); on CPython, " + "use native gRPC instead" ) def _grpc_channel( diff --git a/weaviate/connect/helpers.py b/weaviate/connect/helpers.py index c35e23b25..92f1193f6 100644 --- a/weaviate/connect/helpers.py +++ b/weaviate/connect/helpers.py @@ -41,9 +41,7 @@ def _webify( replaced caller endpoint raises warning Con006. """ if sys.platform != "emscripten": - # grpc_path_prefix=None is passed explicitly so the constructor call (and pydantic's - # error output for it) looks exactly as it did before grpc-web existed - return ConnectionParams(http=http, grpc=grpc, grpc_path_prefix=None) + return ConnectionParams(http=http, grpc=grpc) web_grpc = ProtocolParams(host=http.host, port=http.port, secure=http.secure) if grpc_chosen_by_caller and web_grpc != grpc: @@ -417,7 +415,11 @@ def use_async_with_weaviate_cloud( in an `async with` statement, which will automatically open/close the connection when the context is entered/exited. See the examples below for details. Under Pyodide, gRPC runs over grpc-web on the cluster's REST endpoint (443) instead of - the ``grpc-`` host. + the ``grpc-`` host. Browser use requires "Allow all CORS origins" in the cluster's + settings in the Weaviate Cloud console; until it applies, the first REST call fails + with a "Failed to fetch" connection error. + ``AdditionalConfig`` ``proxies`` / ``trust_env`` and ``GrpcConfig.channel_options`` have no + effect under Pyodide: the browser's fetch makes every request. Args: cluster_url: The WCD cluster URL or hostname to connect to. Usually in the form: rAnD0mD1g1t5.something.weaviate.cloud @@ -471,7 +473,7 @@ def use_async_with_weaviate_cloud( def use_async_with_local( host: str = "localhost", port: int = 8080, - grpc_port: int = 50051, + grpc_port: int = _LOCAL_GRPC_PORT_DEFAULT, headers: Optional[Dict[str, str]] = None, additional_config: Optional[AdditionalConfig] = None, skip_init_checks: bool = False, @@ -485,6 +487,8 @@ def use_async_with_local( Under Pyodide, gRPC runs over grpc-web on the REST endpoint; a non-default ``grpc_port`` is ignored with a warning. + ``AdditionalConfig`` ``proxies`` / ``trust_env`` and ``GrpcConfig.channel_options`` have no + effect under Pyodide: the browser's fetch makes every request. Args: host: The host to use for the underlying REST and GraphQL API calls. @@ -641,6 +645,8 @@ def use_async_with_custom( Under Pyodide, gRPC runs over grpc-web on the REST endpoint: ``grpc_host``, ``grpc_port`` and ``grpc_secure`` are replaced by the HTTP values, with a warning if they differ. + ``AdditionalConfig`` ``proxies`` / ``trust_env`` and ``GrpcConfig.channel_options`` have no + effect under Pyodide: the browser's fetch makes every request. Args: http_host: The host to use for the underlying REST and GraphQL API calls. diff --git a/weaviate/connect/v4.py b/weaviate/connect/v4.py index 1cba0ec50..463e2a288 100644 --- a/weaviate/connect/v4.py +++ b/weaviate/connect/v4.py @@ -150,7 +150,7 @@ def __init__( # fail here, at construction, instead of with an unclear ConnectError on the # first REST call; _client/_grpc_channel are already set, so __del__ does not warn raise WeaviateStartUpError( - "The synchronous client is not supported under WebAssembly/Pyodide. " + "The synchronous client is not supported under Pyodide. " "Use an async client (weaviate.use_async_with_local / " "use_async_with_weaviate_cloud / use_async_with_custom, or " "WeaviateAsyncClient) instead." @@ -354,6 +354,7 @@ def __handle_ping_response(self, res: health_weaviate_pb2.WeaviateHealthCheckRes f"v{self.server_version}", self._connection_params._grpc_address, grpc_path_prefix=self.__grpc_web_prefix(), + http_address=self.__http_address(), ) return None @@ -364,12 +365,16 @@ def __handle_ping_exception(self, e: Exception) -> None: self._connection_params._grpc_address, grpc_path_prefix=self.__grpc_web_prefix(), error=e, + http_address=self.__http_address(), ) from e def __grpc_web_prefix(self) -> Optional[str]: """The configured grpc-web path prefix, or None for native gRPC.""" return self._connection_params._grpc_web_path_prefix or None + def __http_address(self) -> Tuple[str, int]: + return (self._connection_params.http.host, self._connection_params.http.port) + @property def grpc_stub(self) -> Optional[weaviate_pb2_grpc.WeaviateStub]: if not self.is_connected(): @@ -811,11 +816,8 @@ async def _execute() -> None: async with AsyncClient() as client: res = await client.get(PYPI_PACKAGE_URL, timeout=self.timeout_config.init) return resp(res) - except (RequestError, OSError): - # ignore any request error, this is a best-effort warning. OSError covers - # fetch failures under Pyodide/WASM, where the page's CSP often blocks - # pypi.org; that must not fail connect(). - pass + except RequestError: + pass # ignore any errors related to requests, it is a best-effort warning return _execute() @@ -823,7 +825,7 @@ async def _execute() -> None: with Client() as client: res = client.get(PYPI_PACKAGE_URL, timeout=self.timeout_config.init) return resp(res) - except (RequestError, OSError): + except RequestError: pass # ignore any errors related to requests, it is a best-effort warning def delete( diff --git a/weaviate/embedded.py b/weaviate/embedded.py index 731560547..e8811d3a1 100644 --- a/weaviate/embedded.py +++ b/weaviate/embedded.py @@ -180,8 +180,7 @@ def check_supported_platform() -> None: # without this check the port probe below "succeeds" under Emscripten's fake # sockets and wrongly reports that Weaviate is already running raise WeaviateStartUpError( - "Embedded Weaviate is not supported under WebAssembly/Pyodide: it spawns a " - "local Weaviate subprocess, and processes are unavailable in the browser. " + "Embedded Weaviate is not supported under Pyodide: it needs a subprocess. " "Connect to a remote Weaviate instance instead." ) if platform.system() in ["Windows"]: diff --git a/weaviate/exceptions.py b/weaviate/exceptions.py index f83e05224..4a51540c0 100644 --- a/weaviate/exceptions.py +++ b/weaviate/exceptions.py @@ -1,5 +1,6 @@ """Weaviate Exceptions.""" +import sys from json.decoder import JSONDecodeError from typing import Optional, Tuple, Union, cast @@ -338,6 +339,7 @@ def __init__( grpc_address: Tuple[str, int] = ("not provided", 0), grpc_path_prefix: Optional[str] = None, error: Optional[BaseException] = None, + http_address: Optional[Tuple[str, int]] = None, ) -> None: code, details = _grpc_status_of(error) observed = "" @@ -351,6 +353,7 @@ def __init__( # local import: weaviate.connect imports this module at import time, so a # module-level import here would be circular from weaviate.connect.base import ( + GRPC_WEB_CORS_HINT, GRPC_WEB_MIN_SERVER_VERSION, GRPC_WEB_SERVER_PATH_PREFIX, ) @@ -358,6 +361,12 @@ def __init__( # no firewall/wrong-port advice: REST already worked, and grpc-web normally # shares its endpoint address = f"{grpc_address[0]}:{grpc_address[1]}" + if http_address is not None and tuple(http_address) == tuple(grpc_address): + transport = ( + f"carries gRPC over the REST endpoint {address}; there is no separate gRPC port" + ) + else: + transport = f"carries gRPC over {address} (grpc-web)" # weaviate_client_web reports an unrouted path as UNIMPLEMENTED with "HTTP 404/405" # in details; a routed endpoint's own UNIMPLEMENTED is not a wrong path if code is StatusCode.UNIMPLEMENTED and any( @@ -365,25 +374,40 @@ def __init__( ): reason = f"""The server did not route the grpc-web path '{grpc_path_prefix}' at {address}. Either: - the server is too old: grpc-web is served from Weaviate {GRPC_WEB_MIN_SERVER_VERSION} onwards, and this server reports {weaviate_version or "an unknown version"}, or -- the grpc-web base path is wrong: Weaviate serves grpc-web at '{GRPC_WEB_SERVER_PATH_PREFIX}'. The connect helpers set it themselves; only hand-built ConnectionParams choose it (grpc_path_prefix). +- the grpc-web path prefix is wrong: Weaviate serves grpc-web at '{GRPC_WEB_SERVER_PATH_PREFIX}' (set by the connect helpers; hand-built ConnectionParams set grpc_path_prefix). """ else: reason = f"""This error could be due to one of several reasons: - grpc-web is not enabled or is incorrectly configured on the server at {address}. +- {GRPC_WEB_CORS_HINT}. - your connection is unstable or has a high latency. In this case you can: - increase init-timeout in `weaviate.use_async_with_custom(additional_config=wvc.init.AdditionalConfig(timeout=wvc.init.Timeout(init=X)))` - disable startup checks by connecting using `skip_init_checks=True` """ msg = f""" Weaviate {weaviate_version} makes use of a high-speed gRPC API as well as a REST API. -Unfortunately, the gRPC health check against Weaviate could not be completed. - -This client speaks grpc-web (base path '{grpc_path_prefix}'), which carries gRPC over the REST endpoint {address}; there is no separate gRPC port. +The gRPC health check over grpc-web (path prefix '{grpc_path_prefix}') failed. This client {transport}. {reason}{observed}""" super().__init__(msg) return + if sys.platform == "emscripten": + # no prefix under Pyodide means hand-built ConnectionParams; the sync helpers and + # the firewall advice below do not apply + from weaviate.connect.base import GRPC_WEB_SERVER_PATH_PREFIX + + msg = f""" +Weaviate {weaviate_version} makes use of a high-speed gRPC API as well as a REST API. +The gRPC health check against Weaviate could not be completed. + +Under Pyodide gRPC travels as grpc-web, and these connection parameters set no grpc-web path prefix, so requests went to {grpc_address[0]}:{grpc_address[1]} without one. Either: +- connect with `weaviate.use_async_with_local`, `use_async_with_weaviate_cloud` or `use_async_with_custom`, which route gRPC over the REST endpoint themselves, or +- set grpc_path_prefix='{GRPC_WEB_SERVER_PATH_PREFIX}' with the gRPC host and port equal to the REST ones on the ConnectionParams you build. +{observed}""" + super().__init__(msg) + return + if grpc_address[0] == "not provided": grpc_msg = "Please check the server address and port." else: diff --git a/weaviate/warnings.py b/weaviate/warnings.py index d69027fbb..3d6d73628 100644 --- a/weaviate/warnings.py +++ b/weaviate/warnings.py @@ -328,12 +328,12 @@ def grpc_max_msg_size_not_found() -> None: @staticmethod def grpc_endpoint_forced_to_grpc_web(requested: str, effective: str) -> None: warnings.warn( - message=f"""Con006: The gRPC endpoint you gave ({requested}) was overridden with {effective}. + message=f"""Con006: The gRPC endpoint you gave ({requested}) was replaced with {effective}. - Under WebAssembly/Pyodide there is no socket and no grpcio wheel, so native gRPC cannot be used at all; - gRPC runs over grpc-web on the REST listener, which is the endpoint above. Pass gRPC arguments matching - the HTTP ones to silence this warning. A grpc-web transcoder on a separate endpoint is not reachable - through these helpers - build weaviate.connect.ConnectionParams yourself if you need one.""", + Under Pyodide, gRPC runs over grpc-web on the REST endpoint, so the gRPC values you pass + are ignored. Code that also runs on CPython should keep its native gRPC values and + may ignore this warning. For a grpc-web endpoint elsewhere, build + weaviate.connect.ConnectionParams yourself.""", category=UserWarning, stacklevel=1, )