hubcio commented on code in PR #4017:
URL: https://github.com/apache/iggy/pull/4017#discussion_r3940655008


##########
core/metadata/src/stm/stream.rs:
##########
@@ -561,12 +561,25 @@ impl StatsRegistry {
     }
 
     fn remove_partitions_from(&self, stream_id: usize, topic_id: usize, 
first_removed: usize) {
-        self.partitions
+        let mut partitions = self
+            .partitions
             .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, pid), _| {
-                !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed)
-            });
+            .expect("stats registry mutex poisoned");
+        let mut removed = Vec::new();
+        partitions.retain(|(sid, tid, pid), entry| {
+            let keep = !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed);
+            if !keep {
+                removed.push(entry.stats.clone());
+            }
+            keep
+        });
+        drop(partitions);
+
+        // Roll counters out of the parents after releasing the registry lock.
+        // A second left-right apply finds no entries and is a no-op.
+        for stats in removed {
+            stats.zero_out_all();

Review Comment:
   warning: this zero runs while the owning shard can still touch the partition 
- a late segment-cleaner decrement wraps the partition counter and 
double-subtracts from topic and stream totals, a late append leaks into them. 
move the roll-out to the `ConfirmRemove` arm, where writes are fenced.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -561,12 +561,25 @@ impl StatsRegistry {
     }

Review Comment:
   warning: pre-existing, not blocking - `remove_topic` above has the same leak 
one level up: `TopicStats` is dropped without rolling it out of `StreamStats`, 
so `get_stream` and `get_stats` keep a deleted topic's totals until restart. 
zeroing partitions at `ConfirmRemove` cascades through the topic and covers 
this too.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -561,12 +561,25 @@ impl StatsRegistry {
     }
 
     fn remove_partitions_from(&self, stream_id: usize, topic_id: usize, 
first_removed: usize) {
-        self.partitions
+        let mut partitions = self
+            .partitions
             .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, pid), _| {
-                !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed)
-            });
+            .expect("stats registry mutex poisoned");
+        let mut removed = Vec::new();
+        partitions.retain(|(sid, tid, pid), entry| {
+            let keep = !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed);
+            if !keep {
+                removed.push(entry.stats.clone());
+            }
+            keep
+        });
+        drop(partitions);
+
+        // Roll counters out of the parents after releasing the registry lock.

Review Comment:
   nit: not true - the second absorb runs after the next publish, and the 
reconciler's get-or-create can re-insert the key in between, so the second 
apply does find entries. replay is safe only because `zero_out_all` swaps, say 
that instead.



##########
foreign/python/src/client.rs:
##########
@@ -777,6 +781,82 @@ impl IggyClient {
         })
     }
 
+    /// Create partitions for a topic. New partitions use consecutive, 
zero-based IDs

Review Comment:
   nit: "zero-based" only holds for an empty topic, ids start at max + 1. say 
they continue from one past the current highest id and that freed ids are 
reused.



##########
foreign/python/src/client.rs:
##########
@@ -57,6 +57,10 @@ pub struct IggyClient {
     inner: Arc<RustIggyClient>,
 }
 
+fn to_runtime_error(error: impl ToString) -> PyErr {

Review Comment:
   nit: this helper has 2 callers while 38 siblings in this file still inline 
`PyErr::new::<PyRuntimeError, _>(e.to_string())`. inline these two like the 
rest, or move everything over in one sweep.



##########
foreign/python/src/client.rs:
##########
@@ -777,6 +781,82 @@ impl IggyClient {
         })
     }
 
+    /// Create partitions for a topic. New partitions use consecutive, 
zero-based IDs
+    /// after the current maximum; IDs removed by deletion can be reused. 
Existing
+    /// consumer groups leave them unassigned until their next rebalance.

Review Comment:
   warning: this is backwards - `CreatePartitions` rebalances every group in 
the same commit (`rebalance_consumer_groups` in the STM apply): generation 
bump, all partitions redistributed round-robin, pending revocations dropped. 
say that, and regenerate the stub.
   
   also at line 823 - the delete side understates it the same way.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -3069,6 +3082,62 @@ mod tests {
         }
     }
 
+    #[test]
+    fn given_counted_partitions_when_deleted_should_roll_back_parent_stats() {
+        let mut inner = StreamsInner::new();
+        create_stream(&mut inner, "stream");
+        let create_topic = CreateTopicWithAssignmentsRequest {
+            created_view: 0,
+            request: make_topic_request(0, 3, "topic"),
+            derived_options: WireOptions::empty(),
+            partitions: (0..3)
+                .map(|partition_id| CreatedPartitionAssignment {
+                    partition_id,
+                    consensus_group_id: 1,
+                })
+                .collect(),
+        };
+        let _ = StateHandler::apply(&create_topic, &mut inner, 
IggyTimestamp::now());
+
+        let topic_stats = inner.items[0].topics[0].stats.clone();
+        let partition_stats: Vec<_> = (0..3)
+            .map(|partition_id| {
+                inner
+                    .stats_registry
+                    .partition(0, 0, partition_id, topic_stats.clone())
+            })
+            .collect();
+        for (index, stats) in partition_stats.iter().enumerate() {
+            stats.increment_messages_count((index + 1) as u64);
+            stats.increment_size_bytes(((index + 1) * 100) as u64);
+            stats.increment_segments_count(1);
+        }
+
+        let delete = DeletePartitionsRequest {
+            stream_id: WireIdentifier::numeric(0),
+            topic_id: WireIdentifier::numeric(0),
+            partitions_count: 2,
+        };
+        let apply = StateHandler::apply(&delete, &mut inner, 
IggyTimestamp::now());
+        assert_eq!(apply.code, 0);
+
+        assert_eq!(partition_stats[0].messages_count_inconsistent(), 1);
+        assert_eq!(partition_stats[0].size_bytes_inconsistent(), 100);
+        assert_eq!(partition_stats[0].segments_count_inconsistent(), 1);
+        for stats in &partition_stats[1..] {
+            assert_eq!(stats.messages_count_inconsistent(), 0);
+            assert_eq!(stats.size_bytes_inconsistent(), 0);
+            assert_eq!(stats.segments_count_inconsistent(), 0);
+        }
+        assert_eq!(topic_stats.messages_count_inconsistent(), 1);
+        assert_eq!(topic_stats.size_bytes_inconsistent(), 100);
+        assert_eq!(inner.items[0].stats.messages_count_inconsistent(), 1);
+        assert_eq!(inner.items[0].stats.size_bytes_inconsistent(), 100);

Review Comment:
   nit: if this test stays, also assert `segments_count_inconsistent()` on 
topic and stream, and run the delete through both left-right buffers - clone 
`inner` before the delete and apply to both copies.



##########
foreign/python/tests/test_partition.py:
##########
@@ -0,0 +1,149 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you 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.
+
+import pytest
+
+from apache_iggy import IggyClient, SendMessage
+
+
+async def _create_topic(iggy_client: IggyClient, unique_name):
+    stream_name = unique_name()
+    topic_name = unique_name()
+
+    await iggy_client.create_stream(stream_name)
+    await iggy_client.create_topic(
+        stream=stream_name, name=topic_name, partitions_count=2
+    )
+    return stream_name, topic_name
+
+
+class TestPartitionManagement:
+    @pytest.mark.asyncio
+    async def test_create_and_delete_partitions(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+
+        await iggy_client.create_partitions(stream_name, topic_name, 2)
+        created = await iggy_client.get_topic(stream_name, topic_name)
+        assert created is not None
+        assert created.partitions_count == 4
+        assert [partition.id for partition in created.partitions] == [0, 1, 2, 
3]
+
+        await iggy_client.send_messages(
+            stream_name, topic_name, 3, [SendMessage("partition payload")]
+        )
+        await iggy_client.delete_partitions(stream_name, topic_name, 2)
+        deleted = await iggy_client.get_topic(stream_name, topic_name)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]

Review Comment:
   warning: nothing end-to-end checks the stats roll-out - the only message 
goes to partition 3, which gets deleted, so an over-eviction bug would pass 
too. send to partition 0 as well, then poll `get_topic` until `messages_count 
== 1` with partition 0 still holding its message, since stats fold 
asynchronously.



##########
foreign/python/tests/test_partition.py:
##########
@@ -0,0 +1,149 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you 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.
+
+import pytest
+
+from apache_iggy import IggyClient, SendMessage
+
+
+async def _create_topic(iggy_client: IggyClient, unique_name):
+    stream_name = unique_name()
+    topic_name = unique_name()
+
+    await iggy_client.create_stream(stream_name)
+    await iggy_client.create_topic(
+        stream=stream_name, name=topic_name, partitions_count=2
+    )
+    return stream_name, topic_name
+
+
+class TestPartitionManagement:
+    @pytest.mark.asyncio
+    async def test_create_and_delete_partitions(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+
+        await iggy_client.create_partitions(stream_name, topic_name, 2)
+        created = await iggy_client.get_topic(stream_name, topic_name)
+        assert created is not None
+        assert created.partitions_count == 4
+        assert [partition.id for partition in created.partitions] == [0, 1, 2, 
3]
+
+        await iggy_client.send_messages(
+            stream_name, topic_name, 3, [SendMessage("partition payload")]
+        )
+        await iggy_client.delete_partitions(stream_name, topic_name, 2)
+        deleted = await iggy_client.get_topic(stream_name, topic_name)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]
+
+    @pytest.mark.asyncio
+    async def test_partition_management_accepts_numeric_ids(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+        stream = await iggy_client.get_stream(stream_name)
+        assert stream is not None
+        topic = await iggy_client.get_topic(stream.id, topic_name)
+        assert topic is not None
+
+        await iggy_client.create_partitions(stream.id, topic.id, 1)
+        created = await iggy_client.get_topic(stream.id, topic.id)
+        assert created is not None
+        assert created.partitions_count == 3
+        assert [partition.id for partition in created.partitions] == [0, 1, 2]
+
+        await iggy_client.delete_partitions(stream.id, topic.id, 1)
+        deleted = await iggy_client.get_topic(stream.id, topic.id)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]
+
+    @pytest.mark.asyncio
+    async def test_partition_management_rejects_zero_count(

Review Comment:
   nit: the docstrings promise `1..=1000` but only 0 is tested. parametrize 
over `[0, 1001]`, both hit `TooManyPartitions` before consensus.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -561,12 +561,25 @@ impl StatsRegistry {
     }
 
     fn remove_partitions_from(&self, stream_id: usize, topic_id: usize, 
first_removed: usize) {
-        self.partitions
+        let mut partitions = self
+            .partitions
             .lock()
-            .expect("stats registry mutex poisoned")
-            .retain(|(sid, tid, pid), _| {
-                !(*sid == stream_id && *tid == topic_id && *pid >= 
first_removed)
-            });
+            .expect("stats registry mutex poisoned");
+        let mut removed = Vec::new();
+        partitions.retain(|(sid, tid, pid), entry| {

Review Comment:
   simplification: if the zero stays here, do it under the guard - 
`zero_out_all` only touches atomics, so releasing the lock first buys nothing. 
`partitions.extract_if(|(sid, tid, pid), _| *sid == stream_id && *tid == 
topic_id && *pid >= first_removed).for_each(|(_, entry)| 
entry.stats.zero_out_all())` replaces the vec, the clones and the `drop`.



##########
foreign/python/src/client.rs:
##########
@@ -777,6 +781,82 @@ impl IggyClient {
         })
     }
 
+    /// Create partitions for a topic. New partitions use consecutive, 
zero-based IDs
+    /// after the current maximum; IDs removed by deletion can be reused. 
Existing
+    /// consumer groups leave them unassigned until their next rebalance.
+    ///
+    /// Args:
+    ///     stream_id: Stream identifier as `str | int`.
+    ///     topic_id: Topic identifier as `str | int`.
+    ///     partitions_count: Number of partitions to create as `int`; must be
+    ///         1..=1000.
+    ///
+    /// Returns:
+    ///     An awaitable that resolves to `None` when the partitions are 
created.

Review Comment:
   nit: create is just as async as delete - the reply comes after the STM apply 
and the reconciler materializes storage later, sends park until then. say "when 
the partitions are committed" here, same shape as the delete docstring.



##########
foreign/python/tests/test_partition.py:
##########
@@ -0,0 +1,149 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you 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.
+
+import pytest
+
+from apache_iggy import IggyClient, SendMessage
+
+
+async def _create_topic(iggy_client: IggyClient, unique_name):
+    stream_name = unique_name()
+    topic_name = unique_name()
+
+    await iggy_client.create_stream(stream_name)
+    await iggy_client.create_topic(
+        stream=stream_name, name=topic_name, partitions_count=2
+    )
+    return stream_name, topic_name
+
+
+class TestPartitionManagement:
+    @pytest.mark.asyncio
+    async def test_create_and_delete_partitions(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+
+        await iggy_client.create_partitions(stream_name, topic_name, 2)
+        created = await iggy_client.get_topic(stream_name, topic_name)
+        assert created is not None
+        assert created.partitions_count == 4
+        assert [partition.id for partition in created.partitions] == [0, 1, 2, 
3]
+
+        await iggy_client.send_messages(
+            stream_name, topic_name, 3, [SendMessage("partition payload")]
+        )
+        await iggy_client.delete_partitions(stream_name, topic_name, 2)
+        deleted = await iggy_client.get_topic(stream_name, topic_name)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]
+
+    @pytest.mark.asyncio
+    async def test_partition_management_accepts_numeric_ids(

Review Comment:
   simplification: this repeats the create/delete flow above with int ids. 
parametrize the plain flow over name vs numeric ids and keep the 
message-bearing delete as its own test.



##########
foreign/python/src/client.rs:
##########
@@ -777,6 +781,82 @@ impl IggyClient {
         })
     }
 
+    /// Create partitions for a topic. New partitions use consecutive, 
zero-based IDs
+    /// after the current maximum; IDs removed by deletion can be reused. 
Existing
+    /// consumer groups leave them unassigned until their next rebalance.
+    ///
+    /// Args:
+    ///     stream_id: Stream identifier as `str | int`.
+    ///     topic_id: Topic identifier as `str | int`.
+    ///     partitions_count: Number of partitions to create as `int`; must be
+    ///         1..=1000.

Review Comment:
   nit: the `1..=1000` literal duplicates `MAX_PARTITIONS_COUNT`, and 
`create_topic` documents no cap at all. either drop the number or add it there 
too.
   
   also at line 829.



##########
foreign/python/tests/test_partition.py:
##########
@@ -0,0 +1,149 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you 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.
+
+import pytest
+
+from apache_iggy import IggyClient, SendMessage
+
+
+async def _create_topic(iggy_client: IggyClient, unique_name):
+    stream_name = unique_name()
+    topic_name = unique_name()
+
+    await iggy_client.create_stream(stream_name)
+    await iggy_client.create_topic(
+        stream=stream_name, name=topic_name, partitions_count=2
+    )
+    return stream_name, topic_name
+
+
+class TestPartitionManagement:
+    @pytest.mark.asyncio
+    async def test_create_and_delete_partitions(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+
+        await iggy_client.create_partitions(stream_name, topic_name, 2)
+        created = await iggy_client.get_topic(stream_name, topic_name)
+        assert created is not None
+        assert created.partitions_count == 4
+        assert [partition.id for partition in created.partitions] == [0, 1, 2, 
3]
+
+        await iggy_client.send_messages(
+            stream_name, topic_name, 3, [SendMessage("partition payload")]
+        )
+        await iggy_client.delete_partitions(stream_name, topic_name, 2)
+        deleted = await iggy_client.get_topic(stream_name, topic_name)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]
+
+    @pytest.mark.asyncio
+    async def test_partition_management_accepts_numeric_ids(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+        stream = await iggy_client.get_stream(stream_name)
+        assert stream is not None
+        topic = await iggy_client.get_topic(stream.id, topic_name)
+        assert topic is not None
+
+        await iggy_client.create_partitions(stream.id, topic.id, 1)
+        created = await iggy_client.get_topic(stream.id, topic.id)
+        assert created is not None
+        assert created.partitions_count == 3
+        assert [partition.id for partition in created.partitions] == [0, 1, 2]
+
+        await iggy_client.delete_partitions(stream.id, topic.id, 1)
+        deleted = await iggy_client.get_topic(stream.id, topic.id)
+        assert deleted is not None
+        assert deleted.partitions_count == 2
+        assert [partition.id for partition in deleted.partitions] == [0, 1]
+
+    @pytest.mark.asyncio
+    async def test_partition_management_rejects_zero_count(
+        self, iggy_client: IggyClient, unique_name
+    ):
+        stream_name, topic_name = await _create_topic(iggy_client, unique_name)
+
+        with pytest.raises(RuntimeError, match="Too many partitions"):

Review Comment:
   nit: zero hitting "Too many partitions" is a deliberate legacy quirk of 
`validate_partitions_change_count` - leave a one-line comment so nobody "fixes" 
the error code and breaks this test.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to