diff --git a/.github/workflows/cql-conformance.yml b/.github/workflows/cql-conformance.yml new file mode 100644 index 0000000..a1c8211 --- /dev/null +++ b/.github/workflows/cql-conformance.yml @@ -0,0 +1,63 @@ +name: CQL Conformance + +on: + workflow_dispatch: + inputs: + test_filter: + description: 'pytest -k expression to select tests (empty = all)' + required: false + default: '' + scylla_ref: + description: 'ScyllaDB git ref for test content (default: master)' + required: false + default: 'master' + push: + branches: [master, main] + pull_request: + branches: [master, main] + +jobs: + cql-conformance: + runs-on: ubuntu-24.04 + + steps: + - uses: actions/checkout@v4 + + - name: Install system dependencies + run: | + sudo apt-get update + sudo apt-get install -y ninja-build liburing-dev clang-19 clang-tools-19 + + - name: Configure and build server + run: | + cmake -B build -G Ninja \ + -DCMAKE_BUILD_TYPE=Debug \ + -DCMAKE_CXX_COMPILER=clang++-19 \ + -DCMAKE_C_COMPILER=clang-19 + ninja -C build objstore_server + + - name: Install uv + run: | + curl -Ls https://astral.sh/uv/install.sh | sh + echo "$HOME/.cargo/bin" >> $GITHUB_PATH + + - name: Run CQL conformance tests + run: | + ARGS="-P 9042" + ARGS="$ARGS -r ${{ inputs.scylla_ref || 'master' }}" + if [ -n "${{ inputs.test_filter }}" ]; then + ARGS="$ARGS -t '${{ inputs.test_filter }}'" + fi + bash extra/cql_tests/run.sh $ARGS || true + + - name: Generate test summary + if: always() + run: | + echo "### CQL Conformance Test Results" >> "$GITHUB_STEP_SUMMARY" + echo "" >> "$GITHUB_STEP_SUMMARY" + echo "Tests sourced from [scylladb/scylladb](https://github.com/scylladb/scylladb) \`test/cqlpy/cassandra_tests/\`" >> "$GITHUB_STEP_SUMMARY" + echo "(originally ported from Apache Cassandra)" >> "$GITHUB_STEP_SUMMARY" + echo "" >> "$GITHUB_STEP_SUMMARY" + echo "Use \`workflow_dispatch\` to run with custom filters:" >> "$GITHUB_STEP_SUMMARY" + echo "- \`test_filter\`: pytest \`-k\` expression (e.g. \`testInsert\`, \`insert_test\`)" >> "$GITHUB_STEP_SUMMARY" + echo "- \`scylla_ref\`: git ref for test content" >> "$GITHUB_STEP_SUMMARY" diff --git a/extra/cql_tests/.gitignore b/extra/cql_tests/.gitignore new file mode 100644 index 0000000..11196d5 --- /dev/null +++ b/extra/cql_tests/.gitignore @@ -0,0 +1,5 @@ +# Fetched upstream tests and Python artifacts +scylla_cql_tests/ +.venv/ +__pycache__/ +*.pyc diff --git a/extra/cql_tests/conftest.py b/extra/cql_tests/conftest.py new file mode 100644 index 0000000..e0cf484 --- /dev/null +++ b/extra/cql_tests/conftest.py @@ -0,0 +1,270 @@ +# Minimal conftest.py shim for running ScyllaDB's cassandra_tests against objstore. +# +# This replaces ScyllaDB's own conftest.py with a lightweight version that +# provides only the pytest fixtures the upstream test files actually need: +# cql — a cassandra-driver Session connected to objstore +# test_keyspace — a freshly-created keyspace, dropped after the test +# this_dc — the data-center name reported by system.local +# scylla_only — skips the test (objstore is neither Scylla nor Cassandra) +# cassandra_bug — no-op (we don't need to xfail Cassandra bugs) +# + +import os +import socket +import subprocess +import tempfile +import time +import pytest +from cassandra.cluster import Cluster, ConsistencyLevel, ExecutionProfile, EXEC_PROFILE_DEFAULT, NoHostAvailable +from cassandra.policies import RoundRobinPolicy + +# ── CLI options ─────────────────────────────────────────────────────────── +def pytest_addoption(parser): + parser.addoption("--host", action="store", default="127.0.0.1", + help="CQL contact point") + parser.addoption("--port", action="store", default="9042", + help="CQL native transport port") + parser.addoption("--binary", action="store", default=None, + help="Path to objstore_server binary; enables auto-restart on crash") + +# ── Server manager ──────────────────────────────────────────────────────── + +def _wait_for_port(host, port, attempts=50, delay=0.2): + for _ in range(attempts): + try: + s = socket.create_connection((host, port), timeout=1) + s.close() + return + except OSError: + time.sleep(delay) + raise RuntimeError(f"Server did not become ready on {host}:{port}") + +class ServerManager: + """Manages the objstore_server process lifecycle.""" + + def __init__(self, binary, host, port): + self.binary = binary + self.host = host + self.port = port + self._proc = None + self._db = None + + def _spawn(self): + # mkstemp gives us a unique path; remove it so the server creates the file fresh. + fd, self._db = tempfile.mkstemp(prefix="cql_conformance_", suffix=".db", dir="/tmp") + os.close(fd) + os.remove(self._db) + self._proc = subprocess.Popen( + [self.binary, self._db, "--port", str(self.port)], + ) + _wait_for_port(self.host, self.port) + + def start(self): + if not os.path.isfile(self.binary): + pytest.exit( + f"error: objstore binary not found at '{self.binary}'\n" + f"Build with: cmake -B build -G Ninja && ninja -C build objstore_server", + returncode=1, + ) + self._spawn() + print("Server ready.") + + def restart(self): + if self._proc is not None: + self._proc.kill() + self._proc.wait() + if self._db and os.path.exists(self._db): + os.remove(self._db) + self._spawn() + + def stop(self): + if self._proc is not None: + self._proc.kill() + self._proc.wait() + self._proc = None + if self._db and os.path.exists(self._db): + os.remove(self._db) + self._db = None + +@pytest.fixture(scope="session") +def server_manager(request): + binary = request.config.getoption("--binary") + if binary is None: + yield None + return + host = request.config.getoption("--host") + port = int(request.config.getoption("--port")) + mgr = ServerManager(binary, host, port) + mgr.start() + yield mgr + mgr.stop() + +# ── CQL session proxy ───────────────────────────────────────────────────── + +def _make_cluster(host, port): + profile = ExecutionProfile( + load_balancing_policy=RoundRobinPolicy(), + consistency_level=ConsistencyLevel.ONE, + request_timeout=30, + ) + return Cluster( + contact_points=[host], + port=port, + protocol_version=4, + execution_profiles={EXEC_PROFILE_DEFAULT: profile}, + connect_timeout=30, + control_connection_timeout=30, + ) + +class CqlSession: + """Proxy around a cassandra Session; supports disconnect/reconnect between tests.""" + + def __init__(self): + self._cluster = None + self._session = None + + @property + def connected(self): + return self._session is not None + + def connect(self, host, port): + """Open a new connection, closing any existing one. Raises NoHostAvailable on failure.""" + if self._cluster is not None: + try: + self._cluster.shutdown() + except Exception: + pass + self._cluster = None + self._session = None + self._cluster = _make_cluster(host, port) + self._session = self._cluster.connect() # raises NoHostAvailable on failure + + def disconnect(self): + self._session = None + + def __getattr__(self, name): + return getattr(self._session, name) + + def shutdown(self): + if self._cluster is not None: + try: + self._cluster.shutdown() + except Exception: + pass + +# ── Fixtures ────────────────────────────────────────────────────────────── + +@pytest.fixture(scope="session") +def cql(request, server_manager): + host = request.config.getoption("--host") + port = int(request.config.getoption("--port")) + cql_session = CqlSession() + try: + cql_session.connect(host, port) + except NoHostAvailable: + if server_manager is None: + pytest.exit(f"Cannot connect to objstore at {host}:{port}", + returncode=pytest.ExitCode.INTERNAL_ERROR) + # else: cql_test_connection will restart+reconnect before each test + yield cql_session + cql_session.shutdown() + +# Autouse fixture that runs around every test: +# setup — ensure a live CQL connection; restart server + reconnect if needed +# teardown — detect a crash; restart + reconnect so the next test can try +@pytest.fixture(scope="function", autouse=True) +def cql_test_connection(cql, server_manager, request): + host = request.config.getoption("--host") + port = int(request.config.getoption("--port")) + + # ── Pre-test: ensure connected ──────────────────────────────────────── + if not cql.connected: + if server_manager is None: + pytest.skip("No connection to server") + try: + server_manager.restart() + cql.connect(host, port) + except Exception: + pytest.skip("Cannot connect after server restart") + + yield + + # ── Post-test: detect crash; restart for next test ──────────────────── + try: + cql.execute("SELECT * FROM system.local") + except Exception: + cql.disconnect() + if server_manager is not None: + print(f"\nServer crashed during {request.node.nodeid}; restarting...") + try: + server_manager.restart() + cql.connect(host, port) + except Exception: + pass # next test's pre-test check will retry + else: + pytest.exit( + f"objstore crashed during {request.node.nodeid}", + returncode=pytest.ExitCode.INTERNAL_ERROR, + ) + +@pytest.fixture(scope="session") +def this_dc(cql): + if not cql.connected: + yield "datacenter1" + return + row = cql.execute("SELECT data_center FROM system.local").one() + yield row[0] if row else "datacenter1" + +def _unique_ks(): + ms = int(round(time.time() * 1000)) + if _unique_ks.last >= ms: + ms = _unique_ks.last + 1 + _unique_ks.last = ms + return f"cqltest{ms}" +_unique_ks.last = 0 + +@pytest.fixture(scope="function") +def test_keyspace(cql, this_dc): + ks = _unique_ks() + cql.execute( + f"CREATE KEYSPACE {ks} WITH replication = " + f"{{'class': 'SimpleStrategy', 'replication_factor': 1}}" + ) + yield ks + try: + cql.execute(f"DROP KEYSPACE {ks}") + except Exception: + pass + +# objstore is neither Scylla nor Cassandra — skip Scylla-only tests. +@pytest.fixture(scope="session") +def scylla_only(cql): + pytest.skip("Scylla-only test skipped (running on objstore)") + +# cassandra_bug: on objstore we just run the test normally. +@pytest.fixture(scope="session") +def cassandra_bug(): + pass + +@pytest.fixture(scope="session") +def has_tablets(): + return False + +@pytest.fixture(scope="function") +def skip_without_tablets(): + pytest.skip("Tablets not supported") + +@pytest.fixture(scope="function") +def driver_bug_1(): + pass + +@pytest.fixture(scope="function") +def random_seed(): + import random + seed = time.time() + random.seed(seed) + yield seed + +@pytest.fixture(scope="function") +def compact_storage(): + pytest.skip("Compact storage not supported") diff --git a/extra/cql_tests/nodetool.py b/extra/cql_tests/nodetool.py new file mode 100644 index 0000000..39a52ed --- /dev/null +++ b/extra/cql_tests/nodetool.py @@ -0,0 +1,11 @@ +# Minimal nodetool shim — objstore does not support nodetool operations. +# These are no-ops so that porting.py functions like flush() don't crash. + +def flush(cql, table): + pass + +def compact(cql, table): + pass + +def scylla_log(cql, msg, level='info'): + pass diff --git a/extra/cql_tests/pytest.ini b/extra/cql_tests/pytest.ini new file mode 100644 index 0000000..d0f2657 --- /dev/null +++ b/extra/cql_tests/pytest.ini @@ -0,0 +1,5 @@ +[pytest] +# Minimal pytest config for running ScyllaDB's CQL tests against objstore. +# Tests are expected to mostly fail/xfail since objstore implements a subset +# of CQL. Progress is tracked by counting passes over time. +addopts = --no-header diff --git a/extra/cql_tests/run.sh b/extra/cql_tests/run.sh new file mode 100755 index 0000000..90f363f --- /dev/null +++ b/extra/cql_tests/run.sh @@ -0,0 +1,187 @@ +#!/usr/bin/env bash +# run.sh — fetch ScyllaDB CQL tests and run them against a local objstore instance. +# +# The tests come from scylladb/scylladb's test/cqlpy/cassandra_tests/ directory, +# which contains Apache Cassandra CQL tests ported to Python+pytest. Only a +# compatibility shim (conftest.py, porting.py, util.py) is local; the actual +# test content is fetched verbatim from the upstream repository. +# +# Usage: +# ./extra/cql_tests/run.sh [options] +# -b objstore_server binary (default: ./build/objstore/objstore_server) +# -P native protocol port (default: 9042) +# -r scylladb git ref/tag (default: master) +# -t pytest -k expression (default: run all) +# -x stop on first failure (pytest -x) +# -l list tests only (pytest --collect-only) +# -h show this help +# +# Environment: +# CQL_TEST_HOST override contact point (default: 127.0.0.1) +# CQL_TEST_PORT override port (default: value of -p) +# SCYLLA_REF same as -r + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" +ROOT_DIR="$(cd "$SCRIPT_DIR/../.." && pwd)" + +BINARY="${BINARY:-$ROOT_DIR/build/objstore/objstore_server}" +PORT="${PORT:-9042}" +SCYLLA_REF="${SCYLLA_REF:-master}" +PYTEST_K="" +PYTEST_EXTRA="" +LIST_ONLY=false + +while getopts "b:P:r:t:xlh" opt; do + case "$opt" in + b) BINARY="$OPTARG" ;; + P) PORT="$OPTARG" ;; + r) SCYLLA_REF="$OPTARG" ;; + t) PYTEST_K="$OPTARG" ;; + x) PYTEST_EXTRA="$PYTEST_EXTRA -x" ;; + l) LIST_ONLY=true ;; + h) + sed -n '2,22p' "$0" | sed 's/^# \?//' + exit 0 + ;; + *) exit 1 ;; + esac +done + +HOST="${CQL_TEST_HOST:-127.0.0.1}" +PORT="${CQL_TEST_PORT:-$PORT}" + +# ── Fetch ScyllaDB tests ────────────────────────────────────────────────── +TESTS_DIR="$SCRIPT_DIR/scylla_cql_tests" +STAMP="$TESTS_DIR/.ref" + +fetch_tests() { + local ref="$1" + local base_url="https://raw.githubusercontent.com/scylladb/scylladb/${ref}" + + echo "Fetching ScyllaDB CQL tests (ref: ${ref})..." + rm -rf "$TESTS_DIR" + + # The upstream import structure is: + # validation/operations/insert_test.py → from ...porting import * + # cassandra_tests/porting.py → from ..util import unique_name + # from .. import nodetool + # We mirror the hierarchy so relative imports resolve to our shims: + # .scylla_tests/ ← parent package ("test.cqlpy" equivalent) + # util.py ← our shim + # nodetool.py ← our shim + # conftest.py ← our pytest fixtures + # cassandra_tests/ + # porting.py ← UPSTREAM (unmodified) + # validation/operations/ ← UPSTREAM test files + + mkdir -p "$TESTS_DIR/cassandra_tests/validation/operations" + mkdir -p "$TESTS_DIR/cassandra_tests/validation/entities" + + # Operation tests (from Apache Cassandra, ported to Python by ScyllaDB) + local ops=( + create_test.py + insert_test.py + select_test.py + update_test.py + delete_test.py + drop_test.py + truncate_test.py + use_test.py + alter_test.py + batch_test.py + select_limit_test.py + select_order_by_test.py + ) + for f in "${ops[@]}"; do + curl -sfL "${base_url}/test/cqlpy/cassandra_tests/validation/operations/${f}" \ + -o "$TESTS_DIR/cassandra_tests/validation/operations/${f}" || true + done + + # Entity tests + local ents=( + collections_test.py + counters_test.py + timestamp_test.py + type_test.py + ) + for f in "${ents[@]}"; do + curl -sfL "${base_url}/test/cqlpy/cassandra_tests/validation/entities/${f}" \ + -o "$TESTS_DIR/cassandra_tests/validation/entities/${f}" || true + done + + # Upstream porting layer (the test helper library) + curl -sfL "${base_url}/test/cqlpy/cassandra_tests/porting.py" \ + -o "$TESTS_DIR/cassandra_tests/porting.py" + + # __init__.py files for package structure + touch "$TESTS_DIR/__init__.py" + touch "$TESTS_DIR/cassandra_tests/__init__.py" + touch "$TESTS_DIR/cassandra_tests/validation/__init__.py" + touch "$TESTS_DIR/cassandra_tests/validation/operations/__init__.py" + touch "$TESTS_DIR/cassandra_tests/validation/entities/__init__.py" + + echo "$ref" > "$STAMP" + echo "Tests fetched." +} + +# Fetch if ref changed or not present +if [ ! -f "$STAMP" ] || [ "$(cat "$STAMP")" != "$SCYLLA_REF" ]; then + fetch_tests "$SCYLLA_REF" +else + echo "Using cached ScyllaDB tests (ref: $SCYLLA_REF)" +fi + +cp "$SCRIPT_DIR/conftest.py" "$TESTS_DIR/conftest.py" +cp "$SCRIPT_DIR/util.py" "$TESTS_DIR/util.py" +cp "$SCRIPT_DIR/nodetool.py" "$TESTS_DIR/nodetool.py" +echo "Installed shims." + +# ── Python venv ─────────────────────────────────────────────────────────── +VENV="$SCRIPT_DIR/.venv" +if [ ! -d "$VENV" ]; then + echo "Creating venv with Python 3.12 via uv..." + uv venv --python 3.12 "$VENV" +fi +# shellcheck disable=SC1091 +source "$VENV/bin/activate" +uv pip install -U pip setuptools wheel +uv pip install --no-build-isolation cassandra-driver +uv pip install pytest + +# ── List mode ───────────────────────────────────────────────────────────── +# The parent of .scylla_tests/ must be on sys.path so that the package +# (whose __init__.py lives at .scylla_tests/__init__.py) is importable. +PARENT_OF_TESTS="$(dirname "$TESTS_DIR")" + +if $LIST_ONLY; then + echo "" + python3 -m pytest "$TESTS_DIR/cassandra_tests/" \ + --collect-only -q \ + --override-ini="pythonpath=$PARENT_OF_TESTS" \ + --rootdir="$TESTS_DIR" \ + --confcutdir="$TESTS_DIR" \ + -c "$SCRIPT_DIR/pytest.ini" \ + ${PYTEST_K:+-k "$PYTEST_K"} \ + 2>&1 || true + exit 0 +fi + +# ── Run tests ───────────────────────────────────────────────────────────── +echo "" +echo "════════════════════════════════════════════════════════════" +echo " CQL Conformance Tests (scylladb ref: ${SCYLLA_REF})" +echo " Server: ${HOST}:${PORT}" +echo "════════════════════════════════════════════════════════════" +echo "" + +python3 -m pytest "$TESTS_DIR/cassandra_tests/" \ + --host "$HOST" --port "$PORT" --binary "$BINARY" \ + --override-ini="pythonpath=$PARENT_OF_TESTS" \ + --rootdir="$TESTS_DIR" \ + --confcutdir="$TESTS_DIR" \ + -c "$SCRIPT_DIR/pytest.ini" \ + -v --tb=short \ + ${PYTEST_K:+-k "$PYTEST_K"} \ + $PYTEST_EXTRA diff --git a/extra/cql_tests/util.py b/extra/cql_tests/util.py new file mode 100644 index 0000000..774242f --- /dev/null +++ b/extra/cql_tests/util.py @@ -0,0 +1,155 @@ +# Minimal util.py shim — provides helper functions expected by ScyllaDB's +# cassandra_tests porting layer and test files. + +import time +import string +import random +import collections +from contextlib import contextmanager + +def random_string(length=10, chars=string.ascii_uppercase + string.digits): + return ''.join(random.choice(chars) for x in range(length)) + +def random_bytes(length=10): + return bytearray(random.getrandbits(8) for _ in range(length)) + +# Unique names that are valid unquoted CQL identifiers. +unique_name_prefix = 'cqltest' +def unique_name(): + current_ms = int(round(time.time() * 1000)) + if unique_name.last_ms >= current_ms: + current_ms = unique_name.last_ms + 1 + unique_name.last_ms = current_ms + return unique_name_prefix + str(current_ms) +unique_name.last_ms = 0 + +def unique_key_string(): + unique_key_string.i += 1 + return 's' + str(unique_key_string.i) +unique_key_string.i = 0 + +def unique_key_int(): + unique_key_int.i += 1 + return unique_key_int.i +unique_key_int.i = 0 + +def is_scylla(cql): + return False + +def keyspace_has_tablets(cql, keyspace): + return False + +@contextmanager +def new_test_keyspace(cql, opts): + keyspace = unique_name() + cql.execute("CREATE KEYSPACE " + keyspace + " " + opts) + try: + yield keyspace + finally: + cql.execute("DROP KEYSPACE " + keyspace) + +previously_used_table_names = [] +@contextmanager +def new_test_table(cql, keyspace, schema, extra=""): + global previously_used_table_names + if not previously_used_table_names: + previously_used_table_names.append(unique_name()) + table_name = previously_used_table_names.pop() + table = keyspace + "." + table_name + cql.execute("CREATE TABLE " + table + "(" + schema + ")" + extra) + try: + yield table + finally: + cql.execute("DROP TABLE " + table) + previously_used_table_names.append(table_name) + +@contextmanager +def new_type(cql, keyspace, cmd, name=None): + type_name = keyspace + "." + (name or unique_name()) + cql.execute("CREATE TYPE " + type_name + " " + cmd) + try: + yield type_name + finally: + cql.execute("DROP TYPE " + type_name) + +@contextmanager +def new_function(cql, keyspace, body, name=None, args=None): + fun = name if name else unique_name() + cql.execute(f"CREATE FUNCTION {keyspace}.{fun} {body}") + try: + yield fun + finally: + if args: + cql.execute(f"DROP FUNCTION {keyspace}.{fun}({args})") + else: + cql.execute(f"DROP FUNCTION {keyspace}.{fun}") + +@contextmanager +def new_materialized_view(cql, table, select, pk, where, extra=""): + keyspace = table.split('.')[0] + mv = keyspace + "." + unique_name() + cql.execute(f"CREATE MATERIALIZED VIEW {mv} AS SELECT {select} FROM {table} WHERE {where} PRIMARY KEY ({pk}) {extra}") + try: + yield mv + finally: + cql.execute(f"DROP MATERIALIZED VIEW {mv}") + +@contextmanager +def new_secondary_index(cql, table, column, name='', extra=''): + keyspace = table.split('.')[0] + if not name: + name = unique_name() + cql.execute(f"CREATE INDEX {name} ON {table} ({column}) {extra}") + try: + yield f"{keyspace}.{name}" + finally: + cql.execute(f"DROP INDEX {keyspace}.{name}") + +@contextmanager +def new_cql(cql): + session = cql.cluster.connect() + try: + yield session + finally: + session.shutdown() + +def user_type(*args): + return collections.namedtuple('user_type', args[::2])(*args[1::2]) + +@contextmanager +def cql_session(host, port, is_ssl, username, password, request_timeout=120, protocol_version=4): + from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT + from cassandra.policies import RoundRobinPolicy + from cassandra.auth import PlainTextAuthProvider + from cassandra import ConsistencyLevel + profile = ExecutionProfile( + load_balancing_policy=RoundRobinPolicy(), + consistency_level=ConsistencyLevel.ONE, + request_timeout=request_timeout) + cluster = Cluster( + contact_points=[host], + port=int(port), + protocol_version=protocol_version, + execution_profiles={EXEC_PROFILE_DEFAULT: profile}, + auth_provider=PlainTextAuthProvider(username=username, password=password), + connect_timeout=60, + control_connection_timeout=60, + ) + try: + yield cluster.connect() + finally: + cluster.shutdown() + +def project(column_name_string, rows): + return [getattr(r, column_name_string) for r in rows] + +def local_process_id(cql): + return None + +class config_value_context: + def __init__(self, cql, key, value): + pass + def __enter__(self): + pass + def __exit__(self, *args): + pass diff --git a/objstore/engine/engine.cpp b/objstore/engine/engine.cpp index 5ab1648..b54fb31 100644 --- a/objstore/engine/engine.cpp +++ b/objstore/engine/engine.cpp @@ -297,6 +297,17 @@ namespace objstore::engine { if (table == "keyspaces") return make_schema_keyspaces(engine.schema); if (table == "tables") return make_schema_tables(engine.schema); if (table == "columns") return make_schema_columns(engine.schema); + // Return empty result sets for system_schema tables queried by + // drivers during metadata sync that we do not populate. + if (table == "indexes" || table == "triggers" || table == "types" || + table == "functions" || table == "aggregates" || table == "views" || + table == "vertices" || table == "edges") { + VirtualRows vr; + vr.keyspace = "system_schema"; + vr.table = table; + push_back(vr.columns, VirtualColumn{"keyspace_name", types::make_native(NativeType::text)}); + return vr; + } } return {}; } diff --git a/objstore/native/native.cppm b/objstore/native/native.cppm index d534caf..c4946cb 100644 --- a/objstore/native/native.cppm +++ b/objstore/native/native.cppm @@ -473,10 +473,12 @@ namespace objstore::native { case opcode::OPTIONS: { auto frame = make_native_frame(conn, &chunk, write, opcode::SUPPORTED, stream); // [string multimap]: n pairs of ([string] key, [string list] values) - append_be_u16(frame, 1); + append_be_u16(frame, 2); append_cql_string(frame, "CQL_VERSION"); append_be_u16(frame, 1); append_cql_string(frame, "3.0.0"); + append_cql_string(frame, "COMPRESSION"); + append_be_u16(frame, 0); // no compression algorithms supported } break; case opcode::QUERY: { @@ -583,21 +585,22 @@ export namespace objstore::native { push_back(state.recv_buf, chunk.data.ptr[i]); } - // Process one complete frame (if available) - if (state.recv_buf.length >= FRAME_HEADER_SIZE) { + // Process all complete frames in the buffer + while (state.recv_buf.length >= FRAME_HEADER_SIZE) { const U8* hdr = state.recv_buf.ptr; S32 body_len = read_be_s32(hdr + 5); - if (body_len >= 0 && state.recv_buf.length >= FRAME_HEADER_SIZE + U64(body_len)) { - handle_frame(state, engine, req.connection, req.write, hdr, hdr + FRAME_HEADER_SIZE, body_len); + if (body_len < 0 || state.recv_buf.length < FRAME_HEADER_SIZE + U64(body_len)) + break; + + handle_frame(state, engine, req.connection, req.write, hdr, hdr + FRAME_HEADER_SIZE, body_len); - // Remove the consumed frame from the front of recv_buf - U64 frame_size = FRAME_HEADER_SIZE + U64(body_len); - U64 remaining = state.recv_buf.length - frame_size; - if (remaining > 0) - os::memory_move(state.recv_buf.ptr, state.recv_buf.ptr + frame_size, remaining); - state.recv_buf.length = remaining; - } + // Remove the consumed frame from the front of recv_buf + U64 frame_size = FRAME_HEADER_SIZE + U64(body_len); + U64 remaining = state.recv_buf.length - frame_size; + if (remaining > 0) + os::memory_move(state.recv_buf.ptr, state.recv_buf.ptr + frame_size, remaining); + state.recv_buf.length = remaining; } // Always release the chunk chain; unprocessed data is in recv_buf diff --git a/objstore/native/native.test.cpp b/objstore/native/native.test.cpp index 0f44dca..d2a634d 100644 --- a/objstore/native/native.test.cpp +++ b/objstore/native/native.test.cpp @@ -227,6 +227,7 @@ TEST_CASE("Native protocol OPTIONS returns SUPPORTED", "[objstore.native]") { CHECK(resp.body_len > 0); CHECK(body_contains(resp, "CQL_VERSION")); CHECK(body_contains(resp, "3.0.0")); + CHECK(body_contains(resp, "COMPRESSION")); exit_signal = true; os::signal_notify_safe(signal_pipe);