Skip to content

Commit 8b842bd

Browse files
mpartipiloclaude
andcommitted
Remove max_workers and alive_nodes_checking_frequency from async replication config
Both fields are no-ops on servers >= 1.37.3 (the server silently ignores them), making round-trip assertions unreliable. Remove from the Pydantic create/update models, the read dataclass, both factory methods (Configure/Reconfigure.Replication.async_config), and the parser. Integration and unit tests are updated to use propagation_concurrency and hashtree_height, which do round-trip correctly on all supported server versions. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent f26bee0 commit 8b842bd

5 files changed

Lines changed: 20 additions & 49 deletions

File tree

integration/test_collection_config.py

Lines changed: 13 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1598,8 +1598,8 @@ def test_replication_config_with_async_config(collection_factory: CollectionFact
15981598
factor=1,
15991599
async_enabled=True,
16001600
async_config=Configure.Replication.async_config(
1601-
max_workers=8,
16021601
hashtree_height=20,
1602+
propagation_concurrency=4,
16031603
),
16041604
),
16051605
)
@@ -1608,8 +1608,8 @@ def test_replication_config_with_async_config(collection_factory: CollectionFact
16081608
assert config.replication_config.async_enabled is True
16091609
assert config.replication_config.async_config is not None
16101610
ac = config.replication_config.async_config
1611-
assert ac.max_workers == 8
16121611
assert ac.hashtree_height == 20
1612+
assert ac.propagation_concurrency == 4
16131613

16141614

16151615
def test_replication_config_remove_async_config_by_disabling_async_replication(
@@ -1624,14 +1624,13 @@ def test_replication_config_remove_async_config_by_disabling_async_replication(
16241624
factor=1,
16251625
async_enabled=True,
16261626
async_config=Configure.Replication.async_config(
1627-
max_workers=8,
1628-
hashtree_height=20,
1627+
propagation_concurrency=4,
16291628
),
16301629
),
16311630
)
16321631
config = collection.config.get()
16331632
assert config.replication_config.async_config is not None
1634-
assert config.replication_config.async_config.max_workers == 8
1633+
assert config.replication_config.async_config.propagation_concurrency == 4
16351634

16361635
collection.config.update(
16371636
replication_config=Reconfigure.replication(
@@ -1653,14 +1652,13 @@ def test_replication_config_remove_async_config(collection_factory: CollectionFa
16531652
factor=1,
16541653
async_enabled=True,
16551654
async_config=Configure.Replication.async_config(
1656-
max_workers=8,
1657-
hashtree_height=20,
1655+
propagation_concurrency=4,
16581656
),
16591657
),
16601658
)
16611659
config = collection.config.get()
16621660
assert config.replication_config.async_config is not None
1663-
assert config.replication_config.async_config.max_workers == 8
1661+
assert config.replication_config.async_config.propagation_concurrency == 4
16641662

16651663
collection.config.update(
16661664
replication_config=Reconfigure.replication(
@@ -1685,29 +1683,29 @@ def test_replication_config_unset_single_async_field(
16851683
factor=1,
16861684
async_enabled=True,
16871685
async_config=Configure.Replication.async_config(
1688-
max_workers=8,
16891686
hashtree_height=20,
1687+
propagation_concurrency=4,
16901688
),
16911689
),
16921690
)
16931691
config = collection.config.get()
16941692
ac = config.replication_config.async_config
16951693
assert ac is not None
1696-
assert ac.max_workers == 8
16971694
assert ac.hashtree_height == 20
1695+
assert ac.propagation_concurrency == 4
16981696

1699-
# Update with only max_workers — hashtree_height reverts to server default
1697+
# Update with only propagation_concurrency — hashtree_height reverts to server default
17001698
collection.config.update(
17011699
replication_config=Reconfigure.replication(
17021700
async_config=Reconfigure.Replication.async_config(
1703-
max_workers=8,
1701+
propagation_concurrency=4,
17041702
),
17051703
),
17061704
)
17071705
config = collection.config.get()
17081706
ac = config.replication_config.async_config
17091707
assert ac is not None
1710-
assert ac.max_workers == 8
1708+
assert ac.propagation_concurrency == 4
17111709
assert ac.hashtree_height != 20
17121710

17131711

@@ -1734,15 +1732,15 @@ def test_replication_config_add_async_config_to_existing_collection(
17341732
collection.config.update(
17351733
replication_config=Reconfigure.replication(
17361734
async_config=Reconfigure.Replication.async_config(
1737-
max_workers=8,
1735+
hashtree_height=20,
17381736
propagation_concurrency=4,
17391737
),
17401738
),
17411739
)
17421740
config = collection.config.get()
17431741
assert config.replication_config.async_config is not None
17441742
ac = config.replication_config.async_config
1745-
assert ac.max_workers == 8
1743+
assert ac.hashtree_height == 20
17461744
assert ac.propagation_concurrency == 4
17471745

17481746

test/collection/test_config.py

Lines changed: 4 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -2853,11 +2853,9 @@ def test_config_with_vectors(vector_config: List[_VectorConfigCreate], expected:
28532853
(
28542854
Configure.replication(
28552855
async_config=Configure.Replication.async_config(
2856-
max_workers=10,
28572856
hashtree_height=5,
28582857
frequency=60,
28592858
frequency_while_propagating=30,
2860-
alive_nodes_checking_frequency=120,
28612859
logging_frequency=15,
28622860
diff_batch_size=100,
28632861
diff_per_node_timeout=10,
@@ -2871,11 +2869,9 @@ def test_config_with_vectors(vector_config: List[_VectorConfigCreate], expected:
28712869
),
28722870
{
28732871
"asyncConfig": {
2874-
"maxWorkers": 10,
28752872
"hashtreeHeight": 5,
28762873
"frequency": 60,
28772874
"frequencyWhilePropagating": 30,
2878-
"aliveNodesCheckingFrequency": 120,
28792875
"loggingFrequency": 15,
28802876
"diffBatchSize": 100,
28812877
"diffPerNodeTimeout": 10,
@@ -2923,11 +2919,9 @@ def test_configure_with_replication(config: _ReplicationConfigCreate, expected:
29232919
(
29242920
Reconfigure.replication(
29252921
async_config=Reconfigure.Replication.async_config(
2926-
max_workers=10,
29272922
hashtree_height=5,
29282923
frequency=60,
29292924
frequency_while_propagating=30,
2930-
alive_nodes_checking_frequency=120,
29312925
logging_frequency=15,
29322926
diff_batch_size=100,
29332927
diff_per_node_timeout=10,
@@ -2944,11 +2938,9 @@ def test_configure_with_replication(config: _ReplicationConfigCreate, expected:
29442938
"asyncEnabled": None,
29452939
"deletionStrategy": None,
29462940
"asyncConfig": {
2947-
"maxWorkers": 10,
29482941
"hashtreeHeight": 5,
29492942
"frequency": 60,
29502943
"frequencyWhilePropagating": 30,
2951-
"aliveNodesCheckingFrequency": 120,
29522944
"loggingFrequency": 15,
29532945
"diffBatchSize": 100,
29542946
"diffPerNodeTimeout": 10,
@@ -2977,29 +2969,26 @@ def test_replication_config_to_dict_with_async_config() -> None:
29772969
async_enabled=True,
29782970
deletion_strategy=ReplicationDeletionStrategy.TIME_BASED_RESOLUTION,
29792971
async_config=_AsyncReplicationConfig(
2980-
max_workers=8,
29812972
hashtree_height=20,
29822973
frequency=None,
29832974
frequency_while_propagating=None,
2984-
alive_nodes_checking_frequency=3,
29852975
logging_frequency=None,
29862976
diff_batch_size=None,
29872977
diff_per_node_timeout=None,
29882978
pre_propagation_timeout=None,
29892979
propagation_timeout=None,
29902980
propagation_limit=None,
29912981
propagation_delay=None,
2992-
propagation_concurrency=None,
2982+
propagation_concurrency=4,
29932983
propagation_batch_size=None,
29942984
),
29952985
)
29962986
d = config.to_dict()
29972987
assert d["factor"] == 3
29982988
assert d["asyncEnabled"] is True
29992989
assert d["deletionStrategy"] == "TimeBasedResolution"
3000-
assert d["asyncConfig"]["maxWorkers"] == 8
30012990
assert d["asyncConfig"]["hashtreeHeight"] == 20
3002-
assert d["asyncConfig"]["aliveNodesCheckingFrequency"] == 3
2991+
assert d["asyncConfig"]["propagationConcurrency"] == 4
30032992

30042993

30052994
def test_replication_config_to_dict_without_async_config() -> None:
@@ -3025,7 +3014,7 @@ def test_replication_config_update_merge_with_missing_async_config() -> None:
30253014
"""
30263015
update = Reconfigure.replication(
30273016
async_config=Reconfigure.Replication.async_config(
3028-
max_workers=12,
3017+
hashtree_height=5,
30293018
propagation_concurrency=4,
30303019
),
30313020
)
@@ -3036,7 +3025,7 @@ def test_replication_config_update_merge_with_missing_async_config() -> None:
30363025
"deletionStrategy": "NoAutomatedResolution",
30373026
}
30383027
result = update.merge_with_existing(existing_schema)
3039-
assert result["asyncConfig"]["maxWorkers"] == 12
3028+
assert result["asyncConfig"]["hashtreeHeight"] == 5
30403029
assert result["asyncConfig"]["propagationConcurrency"] == 4
30413030
assert result["factor"] == 3
30423031

test/collection/test_config_update.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -112,13 +112,13 @@ def test_replication_async_config_replace_on_update() -> None:
112112
schema = {
113113
"factor": 1,
114114
"asyncEnabled": True,
115-
"asyncConfig": {"maxWorkers": 8, "hashtreeHeight": 20},
115+
"asyncConfig": {"hashtreeHeight": 20, "propagationConcurrency": 4},
116116
}
117117
update = Reconfigure.replication(
118-
async_config=Reconfigure.Replication.async_config(max_workers=16),
118+
async_config=Reconfigure.Replication.async_config(propagation_concurrency=8),
119119
)
120120
result = update.merge_with_existing(schema)
121-
assert result["asyncConfig"] == {"maxWorkers": 16}
121+
assert result["asyncConfig"] == {"propagationConcurrency": 8}
122122
assert "hashtreeHeight" not in result["asyncConfig"]
123123

124124

weaviate/collections/classes/config.py

Lines changed: 0 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -298,11 +298,9 @@ class _ShardingConfigCreate(_ConfigCreateModel):
298298

299299

300300
class _AsyncReplicationConfigCreate(_ConfigCreateModel):
301-
maxWorkers: Optional[int]
302301
hashtreeHeight: Optional[int]
303302
frequency: Optional[int]
304303
frequencyWhilePropagating: Optional[int]
305-
aliveNodesCheckingFrequency: Optional[int]
306304
loggingFrequency: Optional[int]
307305
diffBatchSize: Optional[int]
308306
diffPerNodeTimeout: Optional[int]
@@ -315,11 +313,9 @@ class _AsyncReplicationConfigCreate(_ConfigCreateModel):
315313

316314

317315
class _AsyncReplicationConfigUpdate(_ConfigUpdateModel):
318-
maxWorkers: Optional[int]
319316
hashtreeHeight: Optional[int]
320317
frequency: Optional[int]
321318
frequencyWhilePropagating: Optional[int]
322-
aliveNodesCheckingFrequency: Optional[int]
323319
loggingFrequency: Optional[int]
324320
diffBatchSize: Optional[int]
325321
diffPerNodeTimeout: Optional[int]
@@ -1809,11 +1805,9 @@ def to_dict(self) -> Dict[str, Any]:
18091805

18101806
@dataclass
18111807
class _AsyncReplicationConfig(_ConfigBase):
1812-
max_workers: Optional[int]
18131808
hashtree_height: Optional[int]
18141809
frequency: Optional[int]
18151810
frequency_while_propagating: Optional[int]
1816-
alive_nodes_checking_frequency: Optional[int]
18171811
logging_frequency: Optional[int]
18181812
diff_batch_size: Optional[int]
18191813
diff_per_node_timeout: Optional[int]
@@ -2565,11 +2559,9 @@ class _Replication:
25652559
@staticmethod
25662560
def async_config(
25672561
*,
2568-
max_workers: Optional[int] = None,
25692562
hashtree_height: Optional[int] = None,
25702563
frequency: Optional[int] = None,
25712564
frequency_while_propagating: Optional[int] = None,
2572-
alive_nodes_checking_frequency: Optional[int] = None,
25732565
logging_frequency: Optional[int] = None,
25742566
diff_batch_size: Optional[int] = None,
25752567
diff_per_node_timeout: Optional[int] = None,
@@ -2585,11 +2577,9 @@ def async_config(
25852577
This is only available with WeaviateDB `>=v1.36.0`.
25862578
"""
25872579
return _AsyncReplicationConfigCreate(
2588-
maxWorkers=max_workers,
25892580
hashtreeHeight=hashtree_height,
25902581
frequency=frequency,
25912582
frequencyWhilePropagating=frequency_while_propagating,
2592-
aliveNodesCheckingFrequency=alive_nodes_checking_frequency,
25932583
loggingFrequency=logging_frequency,
25942584
diffBatchSize=diff_batch_size,
25952585
diffPerNodeTimeout=diff_per_node_timeout,
@@ -2606,11 +2596,9 @@ class _ReplicationUpdate:
26062596
@staticmethod
26072597
def async_config(
26082598
*,
2609-
max_workers: Optional[int] = None,
26102599
hashtree_height: Optional[int] = None,
26112600
frequency: Optional[int] = None,
26122601
frequency_while_propagating: Optional[int] = None,
2613-
alive_nodes_checking_frequency: Optional[int] = None,
26142602
logging_frequency: Optional[int] = None,
26152603
diff_batch_size: Optional[int] = None,
26162604
diff_per_node_timeout: Optional[int] = None,
@@ -2626,11 +2614,9 @@ def async_config(
26262614
This is only available with WeaviateDB `>=v1.36.0`.
26272615
"""
26282616
return _AsyncReplicationConfigUpdate(
2629-
maxWorkers=max_workers,
26302617
hashtreeHeight=hashtree_height,
26312618
frequency=frequency,
26322619
frequencyWhilePropagating=frequency_while_propagating,
2633-
aliveNodesCheckingFrequency=alive_nodes_checking_frequency,
26342620
loggingFrequency=logging_frequency,
26352621
diffBatchSize=diff_batch_size,
26362622
diffPerNodeTimeout=diff_per_node_timeout,

weaviate/collections/classes/config_methods.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -385,11 +385,9 @@ def _collection_config_from_json(schema: Dict[str, Any]) -> _CollectionConfig:
385385
),
386386
async_config=(
387387
_AsyncReplicationConfig(
388-
max_workers=async_cfg.get("maxWorkers"),
389388
hashtree_height=async_cfg.get("hashtreeHeight"),
390389
frequency=async_cfg.get("frequency"),
391390
frequency_while_propagating=async_cfg.get("frequencyWhilePropagating"),
392-
alive_nodes_checking_frequency=async_cfg.get("aliveNodesCheckingFrequency"),
393391
logging_frequency=async_cfg.get("loggingFrequency"),
394392
diff_batch_size=async_cfg.get("diffBatchSize"),
395393
diff_per_node_timeout=async_cfg.get("diffPerNodeTimeout"),

0 commit comments

Comments
 (0)