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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 22 additions & 5 deletions cassandra/cqlengine/management.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
import warnings
from itertools import product

from cassandra import metadata
from cassandra import DriverException, metadata
from cassandra.cqlengine import CQLEngineException
from cassandra.cqlengine import columns, query
from cassandra.cqlengine.connection import execute, get_cluster, format_log_context
Expand Down Expand Up @@ -270,7 +270,7 @@ def _sync_table(model, connection=None):

_update_options(model, connection=connection)

table = cluster.metadata.keyspaces[ks_name].tables[raw_cf_name]
table = _get_table_metadata(model, connection)

indexes = [c for n, c in model._columns.items() if c.index]

Expand Down Expand Up @@ -431,9 +431,26 @@ def _get_table_metadata(model, connection=None):
# returns the table as provided by the native driver for a given model
cluster = get_cluster(connection)
ks = model._get_keyspace()
table = model._raw_column_family_name()
table = cluster.metadata.keyspaces[ks].tables[table]
return table
raw_cf_name = model._raw_column_family_name()
try:
return cluster.metadata.keyspaces[ks].tables[raw_cf_name]
except KeyError:
# Metadata may be stale; force a targeted refresh and retry once.
try:
cluster.refresh_table_metadata(ks, raw_cf_name)
except DriverException as exc:
msg = format_log_context(
"Failed to refresh table metadata for '{0}'.'{1}': {2}",
keyspace=ks, connection=connection)
raise CQLEngineException(msg.format(ks, raw_cf_name, exc)) from exc
try:
return cluster.metadata.keyspaces[ks].tables[raw_cf_name]
except KeyError as exc:
msg = format_log_context(
"Table metadata for '{0}'.'{1}' is not available after refresh. "
"Check schema agreement and cluster health.",
keyspace=ks, connection=connection)
raise CQLEngineException(msg.format(ks, raw_cf_name)) from exc


def _options_map_from_strings(option_strings):
Expand Down
223 changes: 223 additions & 0 deletions tests/unit/cqlengine/test_management.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,223 @@
# Copyright 2025 ScyllaDB, Inc.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""
Unit tests for cassandra.cqlengine.management module.
"""

import unittest
from unittest.mock import patch, MagicMock, PropertyMock

from cassandra import DriverException
from cassandra.cqlengine import CQLEngineException
from cassandra.cqlengine.management import _get_table_metadata, _sync_table


class MockTableMeta:
"""Minimal stand-in for TableMetadata."""

def __init__(self):
self.columns = {}
self.options = {}
self.partition_key = []
self.clustering_key = []


class TestGetTableMetadataRetry(unittest.TestCase):
"""Tests for _get_table_metadata retry on KeyError."""

def _make_model(self, ks="test_ks", table="test_table"):
model = MagicMock()
model._get_keyspace.return_value = ks
model._raw_column_family_name.return_value = table
return model

@patch("cassandra.cqlengine.management.get_cluster")
def test_returns_table_when_present(self, mock_get_cluster):
"""Table metadata is found on first lookup -- no refresh needed."""
table_meta = MockTableMeta()
cluster = MagicMock()
cluster.metadata.keyspaces = {
"test_ks": MagicMock(tables={"test_table": table_meta})
}
mock_get_cluster.return_value = cluster
model = self._make_model()

result = _get_table_metadata(model)
self.assertIs(result, table_meta)
cluster.refresh_table_metadata.assert_not_called()

@patch("cassandra.cqlengine.management.get_cluster")
def test_retries_after_refresh_on_missing_table(self, mock_get_cluster):
"""Table missing initially, but available after refresh."""
table_meta = MockTableMeta()
cluster = MagicMock()

# Table is only in `tables_after`; refresh_table_metadata() itself is
# what flips which dict subsequent `.tables` accesses return, so a
# regression that looks up metadata before refreshing can't pass.
tables_first = {}
tables_after = {"test_table": table_meta}
refreshed = {"done": False}

def refresh_side_effect(ks, table):
refreshed["done"] = True

cluster.refresh_table_metadata.side_effect = refresh_side_effect
ks_meta = MagicMock()
type(ks_meta).tables = PropertyMock(
side_effect=lambda: tables_after if refreshed["done"] else tables_first
)
cluster.metadata.keyspaces = {"test_ks": ks_meta}
mock_get_cluster.return_value = cluster

model = self._make_model()
result = _get_table_metadata(model)

self.assertIs(result, table_meta)
cluster.refresh_table_metadata.assert_called_once_with("test_ks", "test_table")

@patch("cassandra.cqlengine.management.get_cluster")
def test_raises_after_failed_refresh(self, mock_get_cluster):
"""Table missing even after refresh -- raises CQLEngineException."""
cluster = MagicMock()
ks_meta = MagicMock()
type(ks_meta).tables = PropertyMock(return_value={})
cluster.metadata.keyspaces = {"test_ks": ks_meta}
mock_get_cluster.return_value = cluster

model = self._make_model()

with self.assertRaises(CQLEngineException) as ctx:
_get_table_metadata(model)

self.assertIn("not available after refresh", str(ctx.exception))
cluster.refresh_table_metadata.assert_called_once_with("test_ks", "test_table")

@patch("cassandra.cqlengine.management.get_cluster")
def test_raises_cqlengineexception_when_refresh_raises_driver_exception(
self, mock_get_cluster
):
"""Table missing initially and refresh_table_metadata() itself raises
a DriverException -- a contextual CQLEngineException is raised
instead of letting the driver error escape uncaught."""
cluster = MagicMock()
ks_meta = MagicMock()
type(ks_meta).tables = PropertyMock(return_value={})
cluster.metadata.keyspaces = {"test_ks": ks_meta}
cluster.refresh_table_metadata.side_effect = DriverException("refresh failed")
mock_get_cluster.return_value = cluster

model = self._make_model()

with self.assertRaises(CQLEngineException) as ctx:
_get_table_metadata(model)

self.assertIsInstance(ctx.exception.__cause__, DriverException)
cluster.refresh_table_metadata.assert_called_once_with("test_ks", "test_table")


class TestSyncTableMetadataLookup(unittest.TestCase):
"""Tests that _sync_table delegates metadata lookup to _get_table_metadata."""

def _make_model(self, ks="test_ks", table="test_table"):
"""Create a mock model that passes _sync_table's precondition checks."""
model = MagicMock()
model.__abstract__ = False
model.column_family_name.return_value = '"test_ks"."test_table"'
model._raw_column_family_name.return_value = table
model._get_keyspace.return_value = ks
model._get_connection.return_value = None
model._columns = {}
return model

@patch("cassandra.cqlengine.management._get_table_metadata")
@patch(
"cassandra.cqlengine.management._get_create_table",
return_value="CREATE TABLE test",
)
@patch("cassandra.cqlengine.management.execute")
@patch("cassandra.cqlengine.management.get_cluster")
@patch(
"cassandra.cqlengine.management._allow_schema_modification", return_value=True
)
@patch("cassandra.cqlengine.management.issubclass", return_value=True)
def test_calls_get_table_metadata_after_create(
self,
mock_issubclass,
mock_allow,
mock_get_cluster,
mock_execute,
mock_create,
mock_get_meta,
):
"""After creating a new table, _sync_table calls _get_table_metadata
-- and only after execute() has run the CREATE TABLE statement."""
table_meta = MockTableMeta()
call_order = []
mock_execute.side_effect = lambda *a, **k: call_order.append("execute")

def get_meta_side_effect(*a, **k):
call_order.append("get_table_metadata")
return table_meta

mock_get_meta.side_effect = get_meta_side_effect

cluster = MagicMock()
ks_meta = MagicMock()
ks_meta.tables = {} # table not in tables -> triggers CREATE TABLE
cluster.metadata.keyspaces = {"test_ks": ks_meta}
mock_get_cluster.return_value = cluster

model = self._make_model()
_sync_table(model)

mock_get_meta.assert_called_once_with(model, None)
self.assertEqual(call_order, ["execute", "get_table_metadata"])

@patch("cassandra.cqlengine.management._get_table_metadata")
@patch(
"cassandra.cqlengine.management._get_create_table",
return_value="CREATE TABLE test",
)
@patch("cassandra.cqlengine.management.execute")
@patch("cassandra.cqlengine.management.get_cluster")
@patch(
"cassandra.cqlengine.management._allow_schema_modification", return_value=True
)
@patch("cassandra.cqlengine.management.issubclass", return_value=True)
def test_propagates_exception_from_get_table_metadata(
self,
mock_issubclass,
mock_allow,
mock_get_cluster,
mock_execute,
mock_create,
mock_get_meta,
):
"""CQLEngineException from _get_table_metadata propagates out of _sync_table."""
mock_get_meta.side_effect = CQLEngineException("Table metadata not available")

cluster = MagicMock()
ks_meta = MagicMock()
ks_meta.tables = {}
cluster.metadata.keyspaces = {"test_ks": ks_meta}
mock_get_cluster.return_value = cluster

model = self._make_model()

with self.assertRaises(CQLEngineException) as ctx:
_sync_table(model)

self.assertIn("not available", str(ctx.exception))
Loading