Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
52 commits
Select commit Hold shift + click to select a range
d39dc31
wip: serverside AEL parsing
gagan405 Jun 4, 2026
6a2c646
fix tests against dsl server
gagan405 Jun 5, 2026
3aec6e8
fixed unit tests to run against dsl on server
gagan405 Jun 6, 2026
b10deee
updated tests
gagan405 Jun 8, 2026
a9686f8
cleanups
gagan405 Jun 8, 2026
35c5963
add gating logic
gagan405 Jun 8, 2026
a160e16
refactors
gagan405 Jun 8, 2026
cd10701
more refactors
gagan405 Jun 8, 2026
5cbe19e
update tests: running version against serverside ael parsing
gagan405 Jun 8, 2026
1a9ce9f
fixed tests to run with both client and server side ael parsing
gagan405 Jun 8, 2026
3b89012
merge dev
gagan405 Jun 9, 2026
aec778b
fixed failing tests
gagan405 Jun 9, 2026
734eebd
updated tests to reflect latest server changes
gagan405 Jun 10, 2026
460981a
fixed bit scan assertion
gagan405 Jun 10, 2026
43b476c
fixed and cleaned up tests
gagan405 Jun 10, 2026
b8eeb70
refactored more test
gagan405 Jun 10, 2026
9cc66ac
[broken] merge from dev
gagan405 Jun 11, 2026
59a65a9
merged changes from dev; updated tests; marked upsert tests as failin…
gagan405 Jun 11, 2026
0483b03
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jun 22, 2026
9d4335b
use sdk to use specific branch of pac
gagan405 Jun 24, 2026
159f05f
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jun 29, 2026
410ba59
add index selection support on server side
gagan405 Jul 6, 2026
1302536
minor refactors
gagan405 Jul 7, 2026
c2a7223
add integ tests
gagan405 Jul 7, 2026
9f7fd09
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jul 7, 2026
074ad0d
Merge branch 'dev' into CLIENT-5011-server-side-idx-selection
gagan405 Jul 7, 2026
07d0d60
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jul 9, 2026
1ab1156
add wire protocol for index hints
gagan405 Jul 9, 2026
ffc6820
merge ael branch with index selection
gagan405 Jul 9, 2026
54bc0be
add integ tests
gagan405 Jul 9, 2026
be780c1
add blob index type
gagan405 Jul 9, 2026
9b379a7
merge dev
gagan405 Jul 13, 2026
005560d
merge dev
gagan405 Jul 13, 2026
08818c0
fix tests
gagan405 Jul 13, 2026
7822403
merge ael parsing
gagan405 Jul 13, 2026
dc6a49b
merged dev
gagan405 Jul 27, 2026
41512ab
restore git url for PAC dependency
gagan405 Jul 27, 2026
d542ef6
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jul 29, 2026
8221d4b
pin PAC version
gagan405 Jul 29, 2026
df156c7
Merge branch 'dev' into CLIENT-4878-serverside-ael-parsing
gagan405 Jul 30, 2026
666c859
merge dev
gagan405 Jul 30, 2026
8400667
Bump PAC pin to pick up multi-byte Field 44 WHERE decode.
gagan405 Jul 30, 2026
1424032
Merge branch 'dev' into CLIENT-5011-server-side-idx-selection
gagan405 Jul 30, 2026
82e0212
fix tests, add feature gate
gagan405 Jul 30, 2026
94786d1
use conditional import till PAC exposes the QueryAPI
gagan405 Jul 30, 2026
be4cae1
applied changes for PR comments
gagan405 Jul 31, 2026
2b8518e
apply review changes
gagan405 Jul 31, 2026
a3bf736
avoid hot path checks
gagan405 Jul 31, 2026
10a5f2a
fix test assertions
gagan405 Aug 1, 2026
5aca5ec
some more caching of cluster version and capabilities
gagan405 Aug 1, 2026
ab9a8e6
fixed the refactors
gagan405 Aug 1, 2026
3ab5e3c
optimizations to shared ops
gagan405 Aug 1, 2026
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
45 changes: 45 additions & 0 deletions aerospike_sdk/ael/server_filter.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
# Copyright 2025-2026 Aerospike, Inc.
#
# Portions may be licensed to Aerospike, Inc. under one or more contributor
# license agreements WHICH ARE COMPATIBLE WITH THE APACHE LICENSE, VERSION 2.0.
#
# 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.

"""Pick client-parsed vs server-compiled filter wire form for AEL strings."""

from __future__ import annotations

from aerospike_async import FilterExpression

from aerospike_sdk.ael.parser import parse_ael

# Resolved once at import — the PAC factory does not change at runtime.
_SERVER_COMPILED_FACTORY = getattr(
FilterExpression, "from_server_compiled_ael", None,
)
_PAC_EXPOSES_SERVER_COMPILED: bool = callable(_SERVER_COMPILED_FACTORY)


def filter_expression_from_ael_string(
ael: str,
*,
supports_server_compiled_ael: bool,
) -> FilterExpression:
"""Return a ``FilterExpression`` for *ael*, using server-compiled wire when allowed.

When ``supports_server_compiled_ael`` is true and PAC exposes the factory,
returns field **43** MessagePack ``[128, "<utf-8 ael>"]`` via
:meth:`~aerospike_async.FilterExpression.from_server_compiled_ael`.
Otherwise parses on the client via :func:`~aerospike_sdk.ael.parser.parse_ael`.
"""
if supports_server_compiled_ael and _PAC_EXPOSES_SERVER_COMPILED:
return _SERVER_COMPILED_FACTORY(ael) # type: ignore[misc]
return parse_ael(ael)
18 changes: 15 additions & 3 deletions aerospike_sdk/aio/background.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
reject_unsupported_background_write_ops,
)
from aerospike_sdk.dataset import DataSet
from aerospike_sdk.ael.parser import parse_ael
from aerospike_sdk.ael.server_filter import filter_expression_from_ael_string
from aerospike_sdk.exceptions import _convert_pac_exception
from aerospike_sdk.operations_shared import _seconds_from_timedelta, _seconds_until

Expand Down Expand Up @@ -201,6 +201,9 @@ def __init__(
self._records_per_second: Optional[int] = None
self._durable_delete_command_default: Optional[bool] = None
self._durable_delete_override: Optional[bool] = None
self._supports_server_compiled_ael = bool(
getattr(session.client, "_cached_supports_server_compiled_ael", False),
)

def default_with_durable_delete(self) -> BackgroundOperationBuilder:
"""Prefer durable deletes when resolving policy defaults (SC namespaces)."""
Expand Down Expand Up @@ -246,7 +249,10 @@ def where(
"use one narrowing mechanism.",
)
if isinstance(expression, str):
self._filter_expression = parse_ael(expression)
self._filter_expression = filter_expression_from_ael_string(
expression,
supports_server_compiled_ael=self._supports_server_compiled_ael,
)
else:
self._filter_expression = expression
return self
Expand Down Expand Up @@ -542,6 +548,9 @@ def __init__(
self._records_per_second: Optional[int] = None
self._durable_delete_command_default: Optional[bool] = None
self._durable_delete_override: Optional[bool] = None
self._supports_server_compiled_ael = bool(
getattr(session.client, "_cached_supports_server_compiled_ael", False),
)

def default_with_durable_delete(self) -> BackgroundUdfBuilder:
"""Prefer durable deletes when resolving policy defaults (SC namespaces)."""
Expand Down Expand Up @@ -587,7 +596,10 @@ def where(
) -> BackgroundUdfBuilder:
"""Optional predicate limiting which records invoke the UDF."""
if isinstance(expression, str):
self._filter_expression = parse_ael(expression)
self._filter_expression = filter_expression_from_ael_string(
expression,
supports_server_compiled_ael=self._supports_server_compiled_ael,
)
else:
self._filter_expression = expression
return self
Expand Down
79 changes: 79 additions & 0 deletions aerospike_sdk/aio/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,20 @@
from aerospike_sdk.policy.behavior_settings import Mode
from aerospike_sdk.policy.sdk_config_loader import fill_hard_defaults
from aerospike_sdk.policy.system_settings import SystemSettings
from aerospike_sdk.feature_gates import (
PSDK_ENABLE_QUERY_SELECTION,
PSDK_ENABLE_SERVER_COMPILED_AEL,
cached_ael_capability_kwargs,
)
from aerospike_sdk.query_selection import (
compute_query_selection_support,
compute_query_selection_support_blocking,
)
from aerospike_sdk.sdk_config_monitor import AsyncSdkConfigMonitor, SdkConfigSource
from aerospike_sdk.server_compiled_ael import (
compute_server_compiled_ael_support,
compute_server_compiled_ael_support_blocking,
)

if typing.TYPE_CHECKING:
from aerospike_sdk.aio.session import Session
Expand Down Expand Up @@ -147,6 +160,8 @@ def __init__(
# Shared by all Session instances from this client; avoids repeated
# namespace/<ns> info probes when callers use multiple sessions.
self._namespace_mode_cache: Dict[str, Mode] = {}
self._cached_supports_query_selection: Optional[bool] = None
self._cached_supports_server_compiled_ael: Optional[bool] = None
# Resolved SDK-level settings (file over programmatic over defaults).
# A frozen snapshot swapped wholesale by the config monitor, so the
# operation path reads it lock-free.
Expand Down Expand Up @@ -220,6 +235,18 @@ async def connect(self) -> None:
log.debug("Connecting to cluster seeds=%r", self._seeds)
self._client = await new_client(self._policy, self._seeds)
self._connected = True
if PSDK_ENABLE_QUERY_SELECTION:
self._cached_supports_query_selection = await compute_query_selection_support(
self._client,
)
else:
self._cached_supports_query_selection = False
if PSDK_ENABLE_SERVER_COMPILED_AEL:
self._cached_supports_server_compiled_ael = (
await compute_server_compiled_ael_support(self._client)
)
else:
self._cached_supports_server_compiled_ael = False
log.info(
"Connected seeds=%r", self._seeds,
extra={"aerospike.cluster": self._policy.cluster_name},
Expand Down Expand Up @@ -261,9 +288,33 @@ async def close(self) -> None:
self._client = None
self._connected = False
log.info("Client closed")
self._cached_supports_query_selection = None
self._cached_supports_server_compiled_ael = None
self._namespace_mode_cache.clear()
self._supports_mrt_cache = None

@property
def supports_query_selection(self) -> bool:
"""``True`` when all cluster nodes support field ``44`` query selection (>= 8.1.3).

Computed at :meth:`connect` / :meth:`connect_blocking` from PAC
``Version.supports_query_selection()`` on every node.
"""
if not self._connected or self._client is None:
return False
return bool(self._cached_supports_query_selection)

@property
def supports_server_compiled_ael(self) -> bool:
"""``True`` when server-compiled AEL filters are usable on this connection.

Requires all nodes >= 8.1.3 (PAC ``Version.supports_server_compiled_ael``)
and PAC ``FilterExpression.from_server_compiled_ael``. Cached at connect.
"""
if not self._connected or self._client is None:
return False
return bool(self._cached_supports_server_compiled_ael)

def connect_blocking(self) -> None:
"""Synchronously open a connection without requiring an asyncio loop.

Expand Down Expand Up @@ -295,6 +346,18 @@ def connect_blocking(self) -> None:
log.debug("Connecting (blocking) to cluster seeds=%r", self._seeds)
self._client = new_client_blocking(self._policy, self._seeds)
self._connected = True
if PSDK_ENABLE_QUERY_SELECTION:
self._cached_supports_query_selection = compute_query_selection_support_blocking(
self._client,
)
else:
self._cached_supports_query_selection = False
if PSDK_ENABLE_SERVER_COMPILED_AEL:
self._cached_supports_server_compiled_ael = (
compute_server_compiled_ael_support_blocking(self._client)
)
else:
self._cached_supports_server_compiled_ael = False
log.info(
"Connected seeds=%r", self._seeds,
extra={"aerospike.cluster": self._policy.cluster_name},
Expand All @@ -313,6 +376,8 @@ def close_blocking(self) -> None:
self._client = None
self._connected = False
log.info("Client closed")
self._cached_supports_query_selection = None
self._cached_supports_server_compiled_ael = None
self._namespace_mode_cache.clear()
self._supports_mrt_cache = None

Expand Down Expand Up @@ -487,7 +552,12 @@ def _query(
behavior=behavior,
indexes_monitor=self._indexes_monitor,
namespace_mode_resolver=namespace_mode_resolver,
namespace_mode_resolver_blocking=namespace_mode_resolver_blocking,
sdk_client=self,
**cached_ael_capability_kwargs(
self._cached_supports_server_compiled_ael,
self._cached_supports_query_selection,
),
)
builder._single_key = key
return builder
Expand All @@ -505,7 +575,12 @@ def _query(
behavior=behavior,
indexes_monitor=self._indexes_monitor,
namespace_mode_resolver=namespace_mode_resolver,
namespace_mode_resolver_blocking=namespace_mode_resolver_blocking,
sdk_client=self,
**cached_ael_capability_kwargs(
self._cached_supports_server_compiled_ael,
self._cached_supports_query_selection,
),
)
builder._keys = keys
return builder
Expand Down Expand Up @@ -535,6 +610,10 @@ def _query(
namespace_mode_resolver=namespace_mode_resolver,
namespace_mode_resolver_blocking=namespace_mode_resolver_blocking,
sdk_client=self,
**cached_ael_capability_kwargs(
self._cached_supports_server_compiled_ael,
self._cached_supports_query_selection,
),
)

@overload
Expand Down
53 changes: 36 additions & 17 deletions aerospike_sdk/aio/operations/query.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@
AerospikeError,
_convert_pac_exception,
)
from aerospike_sdk.feature_gates import cached_ael_capability_kwargs
from aerospike_sdk.policy.behavior_settings import Mode, OpKind, OpShape
from aerospike_sdk.record_result import RecordResult
from aerospike_sdk.record_stream import RecordStream
Expand Down Expand Up @@ -965,7 +966,9 @@ async def _execute_dataset_query(self) -> RecordStream:
log.debug(
"dataset query: %s.%s filter=%s chunk=%s hint=%s",
self._namespace, self._set_name,
self._filter_expression is not None or bool(self._filter_records),
self._filter_expression is not None
or self._where_ael is not None
or bool(self._filter_records),
self._chunk_size,
self._query_hint is not None,
extra={"aerospike.cluster": _cmd_cluster(self._client)},
Expand All @@ -980,40 +983,52 @@ async def _execute_dataset_query(self) -> RecordStream:
policy = self._apply_txn(QueryPolicy())
if self._chunk_size is not None and self._chunk_size > 0:
policy.max_records = self._chunk_size
if self._filter_expression is not None:
policy.filter_expression = self._filter_expression

hint = self._query_hint
use_server_query_selection = self._use_server_query_selection(hint)
self._apply_dataset_query_policy_filter(
policy, use_server_query_selection=use_server_query_selection,
)

if hint is not None and hint.query_duration is not None:
policy.expected_duration = hint.query_duration

if self._where_ael is not None and self._indexes_monitor is not None:
# Lazy start: the monitor's daemon thread only spins up on the
# first AEL ``where()`` query. ``start()`` is idempotent.
self._indexes_monitor.start(self._client)
# Offload the readiness wait so the event loop isn't pinned for
# the first-fetch case (subsequent calls return immediately).
await asyncio.to_thread(self._indexes_monitor.wait_until_ready)
self._prepare_dataset_query_index_context(
use_server_query_selection=use_server_query_selection,
)
await self._wait_for_dataset_query_index_context(
use_server_query_selection=use_server_query_selection,
)

self._resolve_index_context()
if not use_server_query_selection and self._where_ael is not None:
self._resolve_index_context()

partition_filter = self._partition_filter or PartitionFilter.all()

if self._where_ael is not None and self._index_context is not None:
self._auto_generate_filters(hint, policy)
self._maybe_auto_generate_filters(
hint, policy, use_server_query_selection=use_server_query_selection,
)

statement = self._build_statement()

try:
recordset = await self._client.query(statement, partition_filter, policy=policy)
recordset, plan = await self._run_dataset_query_async(
policy, partition_filter, hint, statement,
use_server_query_selection=use_server_query_selection,
)
except Exception as e:
raise _convert_pac_exception(e) from e

if self._chunk_size is not None and self._chunk_size > 0:
client = self._client

async def _reexecute(pf: PartitionFilter) -> Any:
return await client.query(statement, pf, policy=policy)
if plan is not None:
async def _reexecute(pf: PartitionFilter) -> Any:
return await client.query_with_plan(
statement, pf, plan, policy=policy,
)
else:
async def _reexecute(pf: PartitionFilter) -> Any:
return await client.query(statement, pf, policy=policy)

return RecordStream._from_chunked_pac_recordset(
recordset,
Expand Down Expand Up @@ -1148,6 +1163,10 @@ def _promote(self) -> None:
txn=self._txn,
namespace_mode_resolver=self._namespace_mode_resolver,
namespace_mode_resolver_blocking=self._namespace_mode_resolver_blocking,
**cached_ael_capability_kwargs(
getattr(self._sdk_client_fast, "_cached_supports_server_compiled_ael", None),
getattr(self._sdk_client_fast, "_cached_supports_query_selection", None),
),
)
qb._op_type = self._op_type_fast
qb._single_key = self._key
Expand Down
13 changes: 13 additions & 0 deletions aerospike_sdk/aio/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
)
from aerospike_sdk.aio.operations.udf import UdfFunctionBuilder
from aerospike_sdk.dataset import DataSet
from aerospike_sdk.feature_gates import cached_ael_capability_kwargs
from aerospike_sdk.policy.behavior import Behavior, OpKind, OpShape
from aerospike_sdk.policy.behavior_settings import Mode
from aerospike_sdk.policy.policy_mapper import to_read_policy, to_write_policy
Expand Down Expand Up @@ -610,6 +611,10 @@ def execute_udf(self, *keys: Key) -> "UdfFunctionBuilder":
namespace_mode_resolver=self._resolve_namespace_mode,
namespace_mode_resolver_blocking=self._resolve_namespace_mode_blocking,
sdk_client=self._client,
**cached_ael_capability_kwargs(
self._client._cached_supports_server_compiled_ael,
self._client._cached_supports_query_selection,
),
)
qb._set_current_keys_from_varargs(keys)
return UdfFunctionBuilder(qb)
Expand Down Expand Up @@ -691,6 +696,10 @@ def _build_write_segment(
namespace_mode_resolver=self._resolve_namespace_mode,
namespace_mode_resolver_blocking=self._resolve_namespace_mode_blocking,
sdk_client=self._client,
**cached_ael_capability_kwargs(
self._client._cached_supports_server_compiled_ael,
self._client._cached_supports_query_selection,
),
)
target: Union[Key, List[Key]] = all_keys[0] if len(all_keys) == 1 else all_keys
return qb._start_write_verb(op_type, target)
Expand Down Expand Up @@ -746,6 +755,10 @@ def _fast_query_builder(self, key: Key, behavior: Behavior) -> QueryBuilder:
self._resolve_namespace_mode,
self._resolve_namespace_mode_blocking,
self._client,
**cached_ael_capability_kwargs(
self._client._cached_supports_server_compiled_ael,
self._client._cached_supports_query_selection,
),
)
builder._single_key = key
return builder
Expand Down
3 changes: 3 additions & 0 deletions aerospike_sdk/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -661,6 +661,9 @@ def _convert_pac_exception(exc: Exception) -> AerospikeError:
should use ``raise convert_pac_exception(e) from e``.
:func:`_result_code_to_exception`
"""
if isinstance(exc, AerospikeError):
return exc

if isinstance(exc, PacServerError):
return _result_code_to_exception(
exc.result_code,
Expand Down
Loading
Loading