Skip to content

Commit 5ab3087

Browse files
committed
RDBC-1059 Add OptimisticConcurrencyMode with WritesAndReads tracking
1 parent 9ce37a3 commit 5ab3087

9 files changed

Lines changed: 549 additions & 18 deletions

File tree

ravendb/__init__.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -239,6 +239,7 @@
239239
DocumentsChanges,
240240
ForceRevisionStrategy,
241241
MethodCall,
242+
OptimisticConcurrencyMode,
242243
OrderingType,
243244
JavaScriptMap,
244245
DocumentQueryCustomization,

ravendb/documents/commands/batches.py

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44
import json
55
from abc import abstractmethod
66
from enum import Enum
7-
from typing import Callable, Union, Optional, TYPE_CHECKING, List, Set, Dict
7+
from typing import Callable, IO, Union, Optional, TYPE_CHECKING, List, Set, Dict
88

99
import requests
1010

@@ -49,6 +49,7 @@ class CommandType(Enum):
4949
TIME_SERIES_BULK_INSERT = "TIME_SERIES_BULK_INSERT"
5050
TIME_SERIES_COPY = "TIME_SERIES_COPY"
5151
BATCH_PATCH = "BatchPATCH"
52+
BATCH_TRACK_CHANGES = "BatchTrackChanges"
5253
CLIENT_ANY_COMMAND = "CLIENT_ANY_COMMAND"
5354
CLIENT_MODIFY_DOCUMENT_COMMAND = "CLIENT_MODIFY_DOCUMENT_COMMAND"
5455

@@ -81,6 +82,8 @@ def from_csharp_value_str(cls, value: str) -> CommandType:
8182
return cls.COUNTERS
8283
elif value == "BatchPATCH":
8384
return cls.BATCH_PATCH
85+
elif value == "BatchTrackChanges":
86+
return cls.BATCH_TRACK_CHANGES
8487
elif value == "ForceRevisionCreation":
8588
return cls.FORCE_REVISION_CREATION
8689
elif value == "TimeSeries":
@@ -119,7 +122,7 @@ def __init__(
119122
self.__attachment_streams = []
120123
stream = command.stream
121124
if stream is None:
122-
continue # remote-only attachment no stream to track
125+
continue # remote-only attachment has no local stream
123126
if stream in self.__attachment_streams:
124127
raise RuntimeError(
125128
"It is forbidden to re-use the same stream for more than one attachment. "
@@ -270,6 +273,27 @@ def serialize(self, conventions: DocumentConventions) -> dict:
270273
pass
271274

272275

276+
class BatchTrackChangesCommandData(CommandData):
277+
# Emitted in OptimisticConcurrencyMode.WRITES_AND_READS: carries the change
278+
# vector of every tracked entity not already covered by a PUT/DELETE, so
279+
# the server can verify none of them changed underneath us.
280+
def __init__(self, tracked_entities: Dict[str, str], ids_to_skip: Set[str]):
281+
super().__init__(command_type=CommandType.BATCH_TRACK_CHANGES)
282+
self.tracked_entities = tracked_entities
283+
self._ids_to_skip = ids_to_skip
284+
285+
def serialize(self, conventions: DocumentConventions) -> dict:
286+
tracked = {
287+
entity_id: change_vector
288+
for entity_id, change_vector in self.tracked_entities.items()
289+
if entity_id not in self._ids_to_skip
290+
}
291+
return {
292+
"Type": str(CommandType.BATCH_TRACK_CHANGES),
293+
"TrackedEntities": tracked,
294+
}
295+
296+
273297
class DeleteCommandData(CommandData):
274298
def __init__(self, key: str, change_vector: str, original_change_vector: str = None):
275299
super(DeleteCommandData, self).__init__(key=key, command_type=CommandType.DELETE, change_vector=change_vector)
@@ -534,7 +558,7 @@ def __init__(
534558
self,
535559
document_id: str,
536560
name: str,
537-
stream: bytes,
561+
stream: Union[bytes, IO[bytes]],
538562
content_type: str,
539563
change_vector: str,
540564
remote_parameters: Optional["RemoteAttachmentParameters"] = None,
@@ -554,11 +578,11 @@ def __init__(
554578
self.__size_in_bytes = size_in_bytes
555579

556580
@property
557-
def stream(self):
581+
def stream(self) -> Union[bytes, IO[bytes]]:
558582
return self.__stream
559583

560584
@property
561-
def content_type(self):
585+
def content_type(self) -> str:
562586
return self.__content_type
563587

564588
@property

ravendb/documents/conventions.py

Lines changed: 40 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,10 @@ def __init__(self):
5151

5252
# Flags
5353
self.disable_topology_updates = False
54-
self.use_optimistic_concurrency = False
54+
self._optimistic_concurrency_mode = None
55+
# Track which setter the user touched so we can reject mixing them.
56+
self._use_optimistic_concurrency_was_set = False
57+
self._optimistic_concurrency_mode_was_set = False
5558
self.throw_if_query_page_size_is_not_set = False
5659
self._send_application_identifier = True
5760
self._save_enums_as_integers: Optional[bool] = None
@@ -376,6 +379,39 @@ def _assert_not_frozen(self) -> None:
376379
"Conventions has been frozen after documentStore.initialize()" " and no changes can be applied to them"
377380
)
378381

382+
@property
383+
def optimistic_concurrency_mode(self):
384+
from ravendb.documents.session.misc import OptimisticConcurrencyMode
385+
386+
return self._optimistic_concurrency_mode or OptimisticConcurrencyMode.NONE
387+
388+
@optimistic_concurrency_mode.setter
389+
def optimistic_concurrency_mode(self, value) -> None:
390+
self._assert_not_frozen()
391+
if self._use_optimistic_concurrency_was_set:
392+
raise RuntimeError("optimistic_concurrency_mode cannot be combined with use_optimistic_concurrency.")
393+
self._optimistic_concurrency_mode_was_set = True
394+
self._optimistic_concurrency_mode = value
395+
396+
@property
397+
def use_optimistic_concurrency(self) -> bool:
398+
from ravendb.documents.session.misc import OptimisticConcurrencyMode
399+
400+
return self._optimistic_concurrency_mode not in (None, OptimisticConcurrencyMode.NONE)
401+
402+
@use_optimistic_concurrency.setter
403+
def use_optimistic_concurrency(self, value: bool) -> None:
404+
# Legacy bool view: True <-> WRITES, False <-> NONE.
405+
from ravendb.documents.session.misc import OptimisticConcurrencyMode
406+
407+
self._assert_not_frozen()
408+
if self._optimistic_concurrency_mode_was_set:
409+
raise RuntimeError("use_optimistic_concurrency cannot be combined with optimistic_concurrency_mode.")
410+
self._use_optimistic_concurrency_was_set = True
411+
self._optimistic_concurrency_mode = (
412+
OptimisticConcurrencyMode.WRITES if value else OptimisticConcurrencyMode.NONE
413+
)
414+
379415
def clone(self) -> DocumentConventions:
380416
cloned = DocumentConventions()
381417
cloned._list_of_registered_id_conventions = [*self._list_of_registered_id_conventions]
@@ -392,7 +428,9 @@ def clone(self) -> DocumentConventions:
392428
cloned._find_collection_name = self._find_collection_name
393429
cloned._find_python_class_name = self.find_python_class_name
394430

395-
cloned.use_optimistic_concurrency = self.use_optimistic_concurrency
431+
cloned._optimistic_concurrency_mode = self._optimistic_concurrency_mode
432+
cloned._use_optimistic_concurrency_was_set = self._use_optimistic_concurrency_was_set
433+
cloned._optimistic_concurrency_mode_was_set = self._optimistic_concurrency_mode_was_set
396434
cloned.throw_if_query_page_size_is_not_set = self.throw_if_query_page_size_is_not_set
397435
cloned.max_number_of_requests_per_session = self.max_number_of_requests_per_session
398436

ravendb/documents/operations/batch.py

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,7 @@ def get_command_type(obj_node: dict) -> CommandType:
7878
"it. So it was executed ONLY on the requested node on " + self._session.request_executor.url
7979
)
8080

81+
skip = 0
8182
for i in range(self._session_commands_count):
8283
batch_result = result.results[i]
8384
if batch_result is None:
@@ -86,7 +87,7 @@ def get_command_type(obj_node: dict) -> CommandType:
8687
command_type = get_command_type(batch_result)
8788

8889
if command_type == CommandType.PUT:
89-
self._handle_put(i, batch_result, False)
90+
self._handle_put(i - skip, batch_result, False)
9091
elif command_type == CommandType.FORCE_REVISION_CREATION:
9192
self._handle_force_revision_creation(batch_result)
9293
elif command_type == CommandType.DELETE:
@@ -95,6 +96,10 @@ def get_command_type(obj_node: dict) -> CommandType:
9596
self._handle_compare_exchange_put(batch_result)
9697
elif command_type == CommandType.COMPARE_EXCHANGE_DELETE:
9798
self._handle_compare_exchange_delete(batch_result)
99+
elif command_type == CommandType.BATCH_TRACK_CHANGES:
100+
# No client-side state to update; bump skip so PUT indices
101+
# remain aligned with the SaveChangesData.entities array.
102+
skip += 1
98103
else:
99104
raise ValueError(f"Command {command_type} is not supported")
100105

@@ -130,6 +135,8 @@ def get_command_type(obj_node: dict) -> CommandType:
130135
continue # todo: RavenDB-13474 add to time series cache
131136
elif command_type == CommandType.TIME_SERIES_COPY or command_type == CommandType.BATCH_PATCH:
132137
continue
138+
elif command_type == CommandType.BATCH_TRACK_CHANGES:
139+
continue
133140
else:
134141
raise ValueError(f"Command {command_type} is not supported")
135142

@@ -229,6 +236,7 @@ def _handle_delete_internal(self, batch_result: dict, command_type: CommandType)
229236
return
230237

231238
self._session.documents_by_id.pop(key, None)
239+
self._session._tracked_entities.try_remove(key)
232240

233241
if document_info.entity is not None:
234242
self._session.documents_by_entity.pop(document_info.entity, None)
@@ -306,6 +314,8 @@ def _handle_metadata_modifications(
306314
document_info.key = key
307315
document_info.change_vector = change_vector
308316

317+
self._session._tracked_entities.try_update(key, change_vector)
318+
309319
self._apply_metadata_modifications(key, document_info)
310320

311321
def _handle_counters(self, batch_result: Dict) -> None:

ravendb/documents/session/document_session.py

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -667,13 +667,24 @@ def transaction_mode(self) -> TransactionMode:
667667
def transaction_mode(self, value: TransactionMode):
668668
self._session.transaction_mode = value
669669

670+
@property
671+
def optimistic_concurrency_mode(self):
672+
return self._session._optimistic_concurrency_mode
673+
674+
@optimistic_concurrency_mode.setter
675+
def optimistic_concurrency_mode(self, value):
676+
self._session._set_optimistic_concurrency_mode(value)
677+
670678
@property
671679
def use_optimistic_concurrency(self) -> bool:
672-
return self._session._use_optimistic_concurrency
680+
# Derived view; setter routes through optimistic_concurrency_mode.
681+
from ravendb.documents.session.misc import OptimisticConcurrencyMode
682+
683+
return self._session._optimistic_concurrency_mode not in (None, OptimisticConcurrencyMode.NONE)
673684

674685
@use_optimistic_concurrency.setter
675686
def use_optimistic_concurrency(self, value: bool):
676-
self._session._use_optimistic_concurrency = value
687+
self._session._set_use_optimistic_concurrency(value)
677688

678689
def is_loaded(self, key: str) -> bool:
679690
return self._session.is_loaded_or_deleted(key)
@@ -753,6 +764,7 @@ def evict(self, entity: object) -> None:
753764
self._session._counters_by_doc_id.pop(document_info.key, None)
754765
if self._session.time_series_by_doc_id:
755766
self._session.time_series_by_doc_id.pop(document_info.key, None)
767+
self._session._tracked_entities.try_remove(document_info.key)
756768

757769
self._session._deleted_entities.evict(entity)
758770
self._session.entity_to_json.remove_from_missing(entity)
@@ -772,6 +784,7 @@ def clear(self) -> None:
772784
self._session._clear_cluster_session()
773785
self._session._pending_lazy_operations.clear()
774786
self._session.entity_to_json.clear()
787+
self._session._tracked_entities.clear()
775788

776789
def document_query(
777790
self,

0 commit comments

Comments
 (0)