forked from ravendb/ravendb-python-client
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathquery.py
More file actions
3008 lines (2503 loc) · 115 KB
/
Copy pathquery.py
File metadata and controls
3008 lines (2503 loc) · 115 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
from __future__ import annotations
import datetime
import enum
import os
import warnings
from copy import copy
from typing import (
BinaryIO,
Generic,
TypeVar,
List,
Optional,
Callable,
Dict,
Union,
Set,
Tuple,
Collection,
Type,
Iterator,
TYPE_CHECKING,
)
from ravendb.documents.indexes.vector.embedding import VectorEmbeddingType
from ravendb.documents.queries.time_series import TimeSeriesQueryBuilder
from ravendb.documents.session.time_series import TimeSeriesRange, ITimeSeriesValuesBindable
from ravendb.primitives import constants
from ravendb.documents.conventions import DocumentConventions
from ravendb.documents.indexes.spatial.configuration import SpatialUnits, SpatialRelation
from ravendb.documents.queries.explanation import Explanations, ExplanationOptions
from ravendb.documents.queries.facets.builders import FacetBuilder
from ravendb.documents.queries.facets.definitions import FacetBase
from ravendb.documents.queries.facets.queries import AggregationDocumentQuery, AggregationRawDocumentQuery
from ravendb.documents.queries.group_by import GroupBy
from ravendb.documents.queries.highlighting import QueryHighlightings, HighlightingOptions, Highlightings
from ravendb.documents.queries.index_query import IndexQuery, Parameters
from ravendb.documents.queries.misc import SearchOperator
from ravendb.documents.queries.more_like_this import (
MoreLikeThisScope,
MoreLikeThisBase,
MoreLikeThisBuilder,
MoreLikeThisUsingDocument,
MoreLikeThisUsingDocumentForDocumentQuery,
)
from ravendb.documents.queries.query import QueryOperator, QueryResult, QueryTimings, ProjectionBehavior, QueryData
from ravendb.documents.queries.spatial import DynamicSpatialField, SpatialCriteria, SpatialCriteriaFactory
from ravendb.documents.queries.suggestions import (
SuggestionBase,
SuggestionWithTerm,
SuggestionOptions,
SuggestionWithTerms,
SuggestionDocumentQuery,
SuggestionBuilder,
)
from ravendb.documents.queries.utils import QueryFieldUtil
from ravendb.exceptions.exceptions import InvalidOperationException
from ravendb.documents.session.event_args import BeforeQueryEventArgs
from ravendb.documents.session.loaders.include import IncludeBuilderBase, QueryIncludeBuilder
from ravendb.documents.session.misc import (
MethodCall,
CmpXchg,
OrderingType,
NullsOrdering,
DocumentQueryCustomization,
)
from ravendb.documents.session.operations.lazy import LazyQueryOperation
from ravendb.documents.session.operations.query import QueryOperation
from ravendb.documents.session.query_group_by import GroupByDocumentQuery
from ravendb.documents.session.tokens.misc import WhereOperator
from ravendb.documents.session.tokens.query_tokens.facets import FacetToken
from ravendb.documents.session.tokens.query_tokens.query_token import QueryToken
from ravendb.documents.session.tokens.query_tokens.definitions import (
DeclareToken,
LoadToken,
FromToken,
FieldsToFetchToken,
OrderByToken,
GroupByToken,
GroupByKeyToken,
GroupBySumToken,
GroupByCountToken,
MoreLikeThisToken,
CloseSubclauseToken,
WhereToken,
OpenSubclauseToken,
TrueToken,
NegateToken,
QueryOperatorToken,
CompareExchangeValueIncludesToken,
HighlightingToken,
ExplanationToken,
TimingsToken,
IntersectMarkerToken,
DistinctToken,
ShapeToken,
SuggestToken,
CounterIncludesToken,
TimeSeriesIncludesToken,
VectorSearchToken,
)
from ravendb.documents.session.utils.document_query import DocumentQueryHelper
from ravendb.documents.session.utils.includes_util import IncludesUtil
from ravendb.primitives.constants import VectorSearch
from ravendb.tools.utils import Utils
_T = TypeVar("_T")
_TResult = TypeVar("_TResult")
_TProjection = TypeVar("_TProjection")
_T_TS_Bindable = TypeVar("_T_TS_Bindable", bound=ITimeSeriesValuesBindable)
if TYPE_CHECKING:
from ravendb.documents.store.definition import Lazy
from ravendb.documents.session.document_session import DocumentSession
from ravendb.documents.session import InMemoryDocumentSessionOperations
class AbstractDocumentQuery(Generic[_T]):
def __init__(
self,
object_type: Type[_T],
session: "InMemoryDocumentSessionOperations",
index_name: Optional[str],
collection_name: Optional[str],
is_group_by: bool,
declare_tokens: Optional[List[DeclareToken]],
load_tokens: Optional[List[LoadToken]],
from_alias: Optional[str] = None,
is_project_into: Optional[bool] = False,
):
self._object_type = object_type
self._root_types = {object_type}
self._is_group_by = is_group_by
self.__index_name = index_name
self.__collection_name = collection_name
self._from_token = FromToken.create(index_name, collection_name, from_alias)
self._declare_tokens = declare_tokens
self._load_tokens = load_tokens
self._the_session: "InMemoryDocumentSessionOperations" = session
self.__conventions = DocumentConventions() if session is None else session.conventions
self.is_project_into = is_project_into if is_project_into is not None else False
self._is_intersect: Optional[bool] = None
self.__is_in_more_like_this: Optional[bool] = None
self._includes_alias: Optional[str] = None
self._negate: Optional[bool] = None
self._query_raw: Optional[str] = None
self._query_parameters: Dict[str, object] = {}
self.__alias_to_group_by_field_name: Dict[str, str] = {}
self._document_includes: Set[str] = set()
self._page_size: Optional[int] = None
self._start: Optional[int] = None
self.__current_clause_depth: int = 0
self._select_tokens: List[QueryToken] = []
self._fields_to_fetch_token: Optional[FieldsToFetchToken] = None
self._where_tokens: List[QueryToken] = []
self._group_by_tokens: List[QueryToken] = []
self._order_by_tokens: List[QueryToken] = []
self._with_tokens: List[QueryToken] = []
self._compare_exchange_includes_tokens: List[CompareExchangeValueIncludesToken] = []
self._highlighting_tokens: List[HighlightingToken] = []
self._query_highlightings: QueryHighlightings = QueryHighlightings()
self._explanation_token: Optional[ExplanationToken] = None
self._explanations: Optional[Explanations] = None
self._counter_includes_tokens: Optional[List[CounterIncludesToken]] = None
self._time_series_includes_tokens: Optional[List[TimeSeriesIncludesToken]] = None
self._before_query_executed_callback: List[Callable[[IndexQuery], None]] = []
self._after_query_executed_callback: List[Callable[[QueryResult], None]] = [
self.__update_stats_highlightings_and_explanations
]
self._after_stream_executed_callback: List[Callable[[Dict], None]] = []
self._query_stats = QueryStatistics()
self._disable_entities_tracking: Optional[bool] = None
self._disable_caching: Optional[bool] = None
self._query_tag: Optional[str] = None
self._projection_behavior: Optional[ProjectionBehavior] = None
self.parameter_prefix = "p"
self._query_timings: Optional[QueryTimings] = None
self._default_operator = QueryOperator.AND
self._the_wait_for_non_stale_results: Optional[bool] = None
self._timeout: Optional[datetime.timedelta] = None
self._query_operation: Optional[QueryOperation] = None
@property
def conventions(self) -> DocumentConventions:
return self.__conventions
@property
def query_operation(self) -> QueryOperation:
return self._query_operation
@property
def query_class(self) -> Type[_T]:
return self._object_type
@property
def index_name(self) -> str:
return self.__index_name
@property
def collection_name(self) -> str:
return self.__collection_name
@property
def is_distinct(self) -> bool:
return len(self._select_tokens) > 0 and isinstance(self._select_tokens[0], DistinctToken)
@property
def session(self) -> "DocumentSession":
self._the_session: "DocumentSession"
return self._the_session
@property
def __default_timeout(self) -> datetime.timedelta:
return self.conventions.wait_for_non_stale_results_timeout
def _to_string(self) -> str:
"""
Should be used only for testing
Overrides of __str__ cause side effects during debugging
@return: str
"""
return self.__to_string(False)
def _using_default_operator(self, operator: QueryOperator) -> None:
if self._where_tokens:
raise RuntimeError("Default operator can only be set before any where clause is added")
self._default_operator = operator
def _wait_for_non_stale_results(self, wait_timeout: Optional[datetime.timedelta]) -> None:
# Graph queries may set this property multiple times
if self._the_wait_for_non_stale_results:
if self._timeout is None or (wait_timeout is not None and self._timeout.seconds < wait_timeout.seconds):
self._timeout = wait_timeout
return
self._the_wait_for_non_stale_results = True
self._timeout = wait_timeout or self.__default_timeout
@property
def lazy_query_operation(self) -> LazyQueryOperation:
if self._query_operation is None:
self._query_operation = self.initialize_query_operation()
return LazyQueryOperation(
self._object_type, self._the_session, self._query_operation, self._after_query_executed_callback
)
def initialize_query_operation(self) -> QueryOperation:
self._the_session.before_query_invoke(
BeforeQueryEventArgs(self._the_session, DocumentQueryCustomizationDelegate(self))
)
index_query = self.index_query
return QueryOperation(
self._the_session,
self.__index_name,
index_query,
self._fields_to_fetch_token,
self._disable_entities_tracking,
False,
False,
self.is_project_into,
)
@property
def index_query(self) -> IndexQuery:
server_version = None
if self._the_session is not None and self._the_session.advanced.request_executor is not None:
server_version = self._the_session.advanced.request_executor.last_server_version
compability_mode = server_version is not None and server_version < "4.2"
query = self.__to_string(compability_mode)
index_query = self._generate_index_query(query)
self.invoke_before_query_executed(index_query)
return index_query
def _random_ordering(self, seed: Optional[str] = None) -> None:
self.__assert_no_raw_query()
self._no_caching()
self._order_by_tokens.append(OrderByToken.random() if not seed else OrderByToken.create_random(seed))
def _projection(self, projection_behavior):
self._projection_behavior = projection_behavior
def _add_group_by_alias(self, field_name: str, projected_name: str) -> None:
self.__alias_to_group_by_field_name[projected_name] = field_name
def __assert_no_raw_query(self) -> None:
if self._query_raw is not None:
raise RuntimeError(
"raw_query was called, cannot modify this query by calling on operations "
"that would modify the query (such as where, select, order_by, group_by, etc)"
)
def _add_parameter(self, name: str, value: object) -> None:
name = name.lstrip("$")
if name in self._query_parameters:
raise ValueError(f"The parameter {name} was already added")
self._query_parameters[name] = value
def _group_by(self, field_or_field_name: Union[GroupBy, str], *fields_or_field_names: Union[GroupBy, str]) -> None:
field = field_or_field_name if isinstance(field_or_field_name, GroupBy) else GroupBy.field(field_or_field_name)
fields = (
(
list(map(GroupBy.field, fields_or_field_names))
if isinstance(fields_or_field_names[0], str)
else list(fields_or_field_names)
)
if fields_or_field_names
else []
)
fields.insert(0, field)
if not self._from_token.dynamic:
raise RuntimeError("group_by only works with dynamic queries")
self.__assert_no_raw_query()
self._is_group_by = True
for item in fields:
field_name = self._ensure_valid_field_name(item.field, False)
self._group_by_tokens.append(GroupByToken.create(field_name, item.method))
def _group_by_key(self, field_name: str, projected_name: str = None) -> None:
self.__assert_no_raw_query()
self._is_group_by = True
if projected_name is not None and projected_name in self.__alias_to_group_by_field_name:
aliased_field_name = self.__alias_to_group_by_field_name[projected_name]
if field_name is None or field_name.lower() == projected_name.lower():
field_name = aliased_field_name
elif field_name is not None and field_name in self.__alias_to_group_by_field_name:
aliased_field_name = self.__alias_to_group_by_field_name[field_name]
field_name = aliased_field_name
self._select_tokens.append(GroupByKeyToken.create(field_name, projected_name))
def _group_by_sum(self, field_name: str, projected_name: Optional[str] = None):
self.__assert_no_raw_query()
self._is_group_by = True
field_name = self._ensure_valid_field_name(field_name, False)
self._select_tokens.append(GroupBySumToken.create(field_name, projected_name))
def _group_by_count(self, projected_name: Optional[str] = None):
self.__assert_no_raw_query()
self._is_group_by = True
self._select_tokens.append(GroupByCountToken.create(projected_name))
def _where_true(self) -> None:
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, None)
tokens.append(TrueToken.instance())
def _more_like_this(self) -> MoreLikeThisScope:
self.__append_operator_if_needed(self._where_tokens)
token = MoreLikeThisToken()
self._where_tokens.append(token)
self.__is_in_more_like_this = True
def __action():
self.__is_in_more_like_this = False
return MoreLikeThisScope(token, self.__add_query_parameter, __action)
def _include(self, path_or_include_builder: Union[str, IncludeBuilderBase]) -> None:
if self._the_session is not None and self._the_session.no_tracking:
raise InvalidOperationException(
"Cannot register includes when no_tracking is enabled. "
"Included documents are not tracked, so subsequent load operations for that data will still trigger additional server requests. "
"To avoid confusion, include operations are disallowed when tracking is disabled on the session or query."
)
if isinstance(path_or_include_builder, str):
self._document_includes.add(path_or_include_builder)
elif isinstance(path_or_include_builder, IncludeBuilderBase):
includes = path_or_include_builder
if includes is None:
return
if includes.documents_to_include is not None:
self._document_includes.update(includes.documents_to_include)
self._include_counters(includes.alias, includes.counters_to_include_by_source_path)
if includes.time_series_to_include_by_source_alias is not None:
self._include_time_series(includes.alias, includes.time_series_to_include_by_source_alias)
if includes.compare_exchange_values_to_include is not None:
self._compare_exchange_includes_tokens = []
for compare_exchange_value in includes.compare_exchange_values_to_include:
self._compare_exchange_includes_tokens.append(
CompareExchangeValueIncludesToken.create(compare_exchange_value)
)
def _take(self, count: int) -> None:
self._page_size = count
def _skip(self, count: int) -> None:
self._start = count
def _where_lucene(self, field_name: str, where_clause: str, exact: bool) -> None:
field_name = self._ensure_valid_field_name(field_name, False)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, field_name)
options = WhereToken.WhereOptions(exact__from__to=(exact, None, None)) if exact else None
where_token = WhereToken.create(
WhereOperator.LUCENE, field_name, self.__add_query_parameter(where_clause), options
)
tokens.append(where_token)
def _open_subclause(self) -> None:
self.__current_clause_depth += 1
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, None)
tokens.append(OpenSubclauseToken.create())
def _close_subclause(self) -> None:
self.__current_clause_depth -= 1
tokens = self.__get_current_where_tokens()
tokens.append(CloseSubclauseToken.create())
def _where_equals(
self,
field_name__value_or_method__exact: Optional[Tuple[str, Union[MethodCall, object], Optional[bool]]] = None,
where_params: Optional[WhereParams] = None,
) -> None:
if not ((field_name__value_or_method__exact is not None) ^ (where_params is not None)):
raise ValueError("Expected only one argument, got both or any")
if field_name__value_or_method__exact:
params = WhereParams()
params.field_name = field_name__value_or_method__exact[0]
params.value = field_name__value_or_method__exact[1]
params.exact = (
field_name__value_or_method__exact[2] if field_name__value_or_method__exact[2] is not None else False
)
elif where_params:
params = where_params
else:
raise ValueError("Unexpected None value of the argument.")
if self._negate:
self._negate = False
self._where_not_equals(where_params=params)
return
params.field_name = self._ensure_valid_field_name(params.field_name, params.nested_path)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
if self._if_value_is_method(WhereOperator.EQUALS, params, tokens):
return
transform_to_equal_value = self.__transform_value(params)
add_query_parameter = self.__add_query_parameter(transform_to_equal_value)
where_token = WhereToken.create(
WhereOperator.EQUALS,
params.field_name,
add_query_parameter,
WhereToken.WhereOptions(exact__from__to=(params.exact, None, None)),
)
tokens.append(where_token)
def _if_value_is_method(self, op: WhereOperator, where_params: WhereParams, tokens: List[QueryToken]) -> bool:
# MethodCall values (RavenDocumentQuery.now/today, CmpXchg) emit a
# method-flavored WhereToken instead of binding as a parameter.
# Returns True if a token was appended (caller should short-circuit).
if isinstance(where_params.value, MethodCall):
from ravendb.documents.queries.raven_document_query import RavenDocumentQuery
mc = where_params.value
args = []
for arg in mc.args:
args.append(self.__add_query_parameter(arg))
token: Optional[WhereToken] = None
if isinstance(mc, CmpXchg):
token = WhereToken.create(
op,
where_params.field_name,
None,
WhereToken.WhereOptions(
method_type__parameters__property__exact=(
WhereToken.MethodsType.CMP_X_CHG,
args,
mc.access_path,
where_params.exact,
)
),
)
elif isinstance(mc, RavenDocumentQuery.Time):
token = WhereToken.create(
op,
where_params.field_name,
None,
WhereToken.WhereOptions(
method_type__parameters__property__exact=(
mc.method_type,
args,
mc.access_path,
where_params.exact,
)
),
)
else:
raise TypeError(f"Unknown method {type(mc)}")
tokens.append(token)
return True
return False
def _where_not_equals(
self,
field_name__value_or_method__exact: Optional[Tuple[str, Union[MethodCall, object], Optional[bool]]] = None,
where_params: Optional[WhereParams] = None,
):
is_where_params = where_params is not None
if not ((field_name__value_or_method__exact is not None) ^ is_where_params):
raise ValueError("Expected only one argument, got both or any")
if not is_where_params:
where_params = WhereParams()
where_params.field_name = field_name__value_or_method__exact[0]
where_params.value = field_name__value_or_method__exact[1]
exact = field_name__value_or_method__exact[2] if len(field_name__value_or_method__exact) == 3 else False
exact = False if exact is None else exact
where_params.exact = exact
if self._negate:
self._negate = False
self._where_equals(where_params)
return
transform_to_equal_value = self.__transform_value(where_params)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
where_params.field_name = self._ensure_valid_field_name(where_params.field_name, where_params.nested_path)
if self._if_value_is_method(WhereOperator.NOT_EQUALS, where_params, tokens):
return
where_token = WhereToken.create(
WhereOperator.NOT_EQUALS,
where_params.field_name,
self.__add_query_parameter(transform_to_equal_value),
WhereToken.WhereOptions(exact__from__to=(where_params.exact, None, None)),
)
tokens.append(where_token)
def _negate_next(self) -> None:
self._negate = not self._negate
def _where_in(self, field_name: str, values: Collection, exact: Optional[bool] = False) -> None:
field_name = self._ensure_valid_field_name(field_name, False)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, field_name)
where_token = WhereToken.create(
WhereOperator.IN,
field_name,
self.__add_query_parameter(self.__transform_collection(field_name, Utils.unpack_collection(values))),
)
tokens.append(where_token)
def _where_starts_with(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
where_params = WhereParams()
where_params.field_name = field_name
where_params.value = value
where_params.allow_wildcards = True
transform_to_equal_value = self.__transform_value(where_params)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
where_params.field_name = self._ensure_valid_field_name(where_params.field_name, where_params.nested_path)
self.__negate_if_needed(tokens, where_params.field_name)
where_token = WhereToken.create(
WhereOperator.STARTS_WITH,
where_params.field_name,
self.__add_query_parameter(transform_to_equal_value),
WhereToken.WhereOptions(exact__from__to=(exact, None, None)),
)
tokens.append(where_token)
def _where_ends_with(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
where_params = WhereParams()
where_params.field_name = field_name
where_params.value = value
where_params.allow_wildcards = True
transform_to_equal_value = self.__transform_value(where_params)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
where_params.field_name = self._ensure_valid_field_name(where_params.field_name, where_params.nested_path)
self.__negate_if_needed(tokens, where_params.field_name)
where_token = WhereToken.create(
WhereOperator.ENDS_WITH,
where_params.field_name,
self.__add_query_parameter(transform_to_equal_value),
WhereToken.WhereOptions(exact__from__to=(exact, None, None)),
)
tokens.append(where_token)
def _where_between(self, field_name: str, start: object, end: object, exact: Optional[bool] = False) -> None:
field_name = self._ensure_valid_field_name(field_name, False)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, field_name)
start_params = WhereParams()
start_params.value = start
start_params.field_name = field_name
end_params = WhereParams()
end_params.value = end
end_params.field_name = field_name
from_parameter_name = self.__add_query_parameter(
"*" if start is None else self.__transform_value(start_params, True)
)
to_parameter_name = self.__add_query_parameter(
"NULL" if end is None else self.__transform_value(end_params, True)
)
where_token = WhereToken.create(
WhereOperator.BETWEEN,
field_name,
None,
WhereToken.WhereOptions(exact__from__to=(exact, from_parameter_name, to_parameter_name)),
)
tokens.append(where_token)
def _where_compare(
self,
op: WhereOperator,
field_name: str,
value: object,
exact: Optional[bool],
null_sentinel: str,
) -> None:
# Shared body for >/>=/</<=. Routes through _if_value_is_method first
# so MethodCall values (now/today/cmpxchg) are emitted as RQL calls.
field_name = self._ensure_valid_field_name(field_name, False)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, field_name)
where_params = WhereParams()
where_params.value = value
where_params.field_name = field_name
where_params.exact = exact
if self._if_value_is_method(op, where_params, tokens):
return
parameter = self.__add_query_parameter(
null_sentinel if value is None else self.__transform_value(where_params, True)
)
tokens.append(
WhereToken.create(op, field_name, parameter, WhereToken.WhereOptions(exact__from__to=(exact, None, None)))
)
def _where_greater_than(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
self._where_compare(WhereOperator.GREATER_THAN, field_name, value, exact, "*")
def _where_greater_than_or_equal(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
self._where_compare(WhereOperator.GREATER_THAN_OR_EQUAL, field_name, value, exact, "*")
def _where_less_than(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
self._where_compare(WhereOperator.LESS_THAN, field_name, value, exact, "NULL")
def _where_less_than_or_equal(self, field_name: str, value: object, exact: Optional[bool] = False) -> None:
self._where_compare(WhereOperator.LESS_THAN_OR_EQUAL, field_name, value, exact, "NULL")
def _where_regex(self, field_name: str, pattern: str) -> None:
field_name = self._ensure_valid_field_name(field_name, False)
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
self.__negate_if_needed(tokens, field_name)
where_params = WhereParams()
where_params.value = pattern
where_params.field_name = field_name
parameter = self.__add_query_parameter(self.__transform_value(where_params))
where_token = WhereToken.create(WhereOperator.REGEX, field_name, parameter)
tokens.append(where_token)
def _and_also(self, wrap_previous_query_clauses: bool = False) -> None:
tokens = self.__get_current_where_tokens()
if not tokens:
return
if isinstance(tokens[-1], QueryOperatorToken):
raise TypeError("Cannot add AND, previous token was already an operator token")
if wrap_previous_query_clauses:
tokens.insert(0, OpenSubclauseToken.create())
tokens.append(CloseSubclauseToken.create())
tokens.append(QueryOperatorToken.AND())
def _or_else(self) -> None:
tokens = self.__get_current_where_tokens()
if len(tokens) == 0:
return
if isinstance(tokens[-1], QueryOperatorToken):
raise RuntimeError("Cannot add OR, previous token was already an operator token")
tokens.append(QueryOperatorToken.OR())
def _boost(self, boost: float) -> None:
if boost == 1.0:
return
if boost < 0.0:
raise ValueError("Boost factor must be a non-negative number")
tokens = self.__get_current_where_tokens()
last = tokens[-1] if tokens else None
if isinstance(last, WhereToken):
where_token = last
where_token.options.boost = boost
elif isinstance(last, CloseSubclauseToken):
close = last
parameter = self.__add_query_parameter(boost)
index = tokens.index(last)
while last is not None and index > 0:
index -= 1
last = tokens[index] # find the previous option
if isinstance(last, OpenSubclauseToken):
open_token = last
open_token.boost_parameter_name = parameter
close.boost_parameter_name = parameter
return
else:
raise RuntimeError("Cannot apply boost")
def _fuzzy(self, fuzzy: float) -> None:
tokens = self.__get_current_where_tokens()
if not tokens:
raise RuntimeError("Fuzzy can only be used right after where clause")
where_token = tokens[-1]
if not isinstance(where_token, WhereToken):
raise RuntimeError("Fuzzy can only be used right after where clause")
if where_token.where_operator != WhereOperator.EQUALS:
raise RuntimeError("Fuzzy can only be used right after where clause with equals operator")
if fuzzy < 0.0 or fuzzy > 1.0:
raise ValueError("Fuzzy distance must be between 0.0 and 1")
where_token.options.fuzzy = fuzzy
def _proximity(self, proximity: int) -> None:
tokens = self.__get_current_where_tokens()
if not tokens:
raise RuntimeError("Proximity can only be used right after search clause")
where_token = tokens[-1]
if not isinstance(where_token, WhereToken):
raise RuntimeError("Proximity can only be used right after search clause")
if where_token.where_operator != WhereOperator.SEARCH:
raise RuntimeError("Proximity can only be used right after search clause")
if proximity < 1:
raise ValueError("Proximity distance must be a positive number")
where_token.options.proximity = proximity
def _order_by(
self,
field: str,
sorter_name_or_ordering_type: Optional[Union[str, OrderingType]] = OrderingType.STRING,
nulls: NullsOrdering = NullsOrdering.DEFAULT,
) -> None:
is_ordering_type = isinstance(sorter_name_or_ordering_type, OrderingType)
if not is_ordering_type and sorter_name_or_ordering_type.isspace():
raise ValueError("Sorter name cannot be None or whitespace")
self.__assert_no_raw_query()
f = self._ensure_valid_field_name(field, False)
self._order_by_tokens.append(OrderByToken.create_ascending(f, sorter_name_or_ordering_type, nulls))
def _order_by_descending(
self,
field: str,
sorter_name_or_ordering_type: Optional[Union[str, OrderingType]] = OrderingType.STRING,
nulls: NullsOrdering = NullsOrdering.DEFAULT,
) -> None:
is_ordering_type = isinstance(sorter_name_or_ordering_type, OrderingType)
if not is_ordering_type and sorter_name_or_ordering_type.isspace():
raise ValueError("Sorter name cannot be None or whitespace")
self.__assert_no_raw_query()
f = self._ensure_valid_field_name(field, False)
self._order_by_tokens.append(OrderByToken.create_descending(f, sorter_name_or_ordering_type, nulls))
def _order_by_score(self) -> None:
self.__assert_no_raw_query()
self._order_by_tokens.append(OrderByToken.score_ascending())
def _order_by_score_descending(self) -> None:
self.__assert_no_raw_query()
self._order_by_tokens.append(OrderByToken.score_descending())
def _statistics(self, stats_callback: Callable[[QueryStatistics], None]) -> None:
stats_callback(self._query_stats)
def invoke_after_query_executed(self, result: QueryResult) -> None:
for callback in self._after_query_executed_callback:
callback(result)
def invoke_before_query_executed(self, query: IndexQuery) -> None:
for callback in self._before_query_executed_callback:
callback(query)
def invoke_after_stream_executed(self, result: dict) -> None:
for callback in self._after_stream_executed_callback:
callback(result)
def _generate_index_query(self, query: str) -> IndexQuery:
index_query = IndexQuery()
index_query.query = query
index_query.start = self._start
index_query.wait_for_non_stale_results = self._the_wait_for_non_stale_results
index_query.wait_for_non_stale_results_timeout = self._timeout
index_query.query_parameters = self._query_parameters
index_query.disable_caching = self._disable_caching
index_query.tag = self._query_tag
index_query.projection_behavior = self._projection_behavior
if self._page_size is not None:
index_query.page_size = self._page_size
return index_query
def _search(self, field_name: str, search_terms: str, operator: SearchOperator) -> None:
tokens = self.__get_current_where_tokens()
self.__append_operator_if_needed(tokens)
field_name = self._ensure_valid_field_name(field_name, False)
self.__negate_if_needed(tokens, field_name)
where_token = WhereToken.create(
WhereOperator.SEARCH,
field_name,
self.__add_query_parameter(search_terms),
WhereToken.WhereOptions(search=operator),
)
tokens.append(where_token)
def __to_string(self, compatibility_mode: bool) -> str:
if self._query_raw is not None:
return self._query_raw
if self.__current_clause_depth != 0:
raise RuntimeError(
f"A clause was not closed correctly within this query, "
f"current clause depth = {self.__current_clause_depth}"
)
query_text = []
self.__build_declare(query_text)
self.__build_from(query_text)
self.__build_group_by(query_text)
self.__build_where(query_text)
self.__build_order_by(query_text)
self.__build_load(query_text)
self.__build_select(query_text)
self.__build_include(query_text)
if not compatibility_mode:
self.__build_pagination(query_text)
return "".join(query_text)
def __build_with(self, query_text: List[str]) -> None:
for with_token in self._with_tokens:
with_token.write_to(query_text)
query_text.append(os.linesep)
def __build_pagination(self, query_text: List[str]) -> None:
if (self._start is not None and self._start > 0) or self._page_size is not None:
query_text.append(" limit $")
query_text.append(self.__add_query_parameter(self._start or 0))
query_text.append(", $")
query_text.append(self.__add_query_parameter(self._page_size))
def __build_include(self, query_text: List[str]) -> None:
if (
not self._document_includes
and not self._highlighting_tokens
and self._explanation_token is None
and self._query_timings is None
and not self._compare_exchange_includes_tokens
and self._counter_includes_tokens is None
and self._time_series_includes_tokens is None
):
return
query_text.append(" include ")
first = True
for include in self._document_includes:
if not first:
query_text.append(",")
first = False
required, escaped_include = IncludesUtil.requires_quotes(include)
if required:
query_text.append("'")
query_text.append(escaped_include)
query_text.append("'")
else:
query_text.append(include)
first = self.__write_include_tokens(self._counter_includes_tokens, first, query_text)
first = self.__write_include_tokens(self._time_series_includes_tokens, first, query_text)
first = self.__write_include_tokens(self._compare_exchange_includes_tokens, first, query_text)
first = self.__write_include_tokens(self._highlighting_tokens, first, query_text)
if self._explanation_token is not None:
if not first:
query_text.append(",")
first = False
self._explanation_token.write_to(query_text)
if self._query_timings is not None:
if not first:
query_text.append(",")
first = False
TimingsToken.instance().write_to(query_text)
def __write_include_tokens(self, tokens: Collection[QueryToken], first: bool, query_text: List[str]) -> bool:
if tokens is None:
return first
for token in tokens:
if not first:
query_text.append(",")
first = False
token.write_to(query_text)
return first
def _intersect(self) -> None:
tokens = self.__get_current_where_tokens()
if tokens:
last = tokens[-1]
if isinstance(last, (WhereToken, CloseSubclauseToken)):
self._is_intersect = True