diff --git a/foreign/python/apache_iggy.pyi b/foreign/python/apache_iggy.pyi index 5b53669384..58b1d09364 100644 --- a/foreign/python/apache_iggy.pyi +++ b/foreign/python/apache_iggy.pyi @@ -36,6 +36,7 @@ __all__ = [ "PollingStrategy", "ReceiveMessage", "SendMessage", + "Stream", "StreamDetails", "Topic", "TopicDetails", @@ -352,6 +353,32 @@ class IggyClient: Gets stream by id. Returns Option of stream details or a PyRuntimeError on failure. """ + def get_streams(self) -> collections.abc.Awaitable[list[Stream]]: + r""" + Gets all streams. + Returns a list of streams or a PyRuntimeError on failure. + """ + def update_stream( + self, stream_id: builtins.str | builtins.int, name: builtins.str + ) -> collections.abc.Awaitable[None]: + r""" + Updates a stream's name by id. + Returns Ok(()) on successful update or a PyRuntimeError on failure. + """ + def delete_stream( + self, stream_id: builtins.str | builtins.int + ) -> collections.abc.Awaitable[None]: + r""" + Deletes a stream by id. + Returns Ok(()) on successful deletion or a PyRuntimeError on failure. + """ + def purge_stream( + self, stream_id: builtins.str | builtins.int + ) -> collections.abc.Awaitable[None]: + r""" + Purges all messages from a stream by id. + Returns Ok(()) on successful purge or a PyRuntimeError on failure. + """ def create_topic( self, stream: builtins.str | builtins.int, @@ -796,6 +823,17 @@ class SendMessage: directly from Python using the provided string or bytes data. """ +@typing.final +class Stream: + @property + def id(self) -> builtins.int: ... + @property + def name(self) -> builtins.str: ... + @property + def messages_count(self) -> builtins.int: ... + @property + def topics_count(self) -> builtins.int: ... + @typing.final class StreamDetails: @property diff --git a/foreign/python/src/client.rs b/foreign/python/src/client.rs index 2a36ef3335..78532b4bef 100644 --- a/foreign/python/src/client.rs +++ b/foreign/python/src/client.rs @@ -35,7 +35,7 @@ use crate::consumer::{ use crate::identifier::PyIdentifier; use crate::receive_message::{PollingStrategy, ReceiveMessage}; use crate::send_message::SendMessage; -use crate::stream::StreamDetails; +use crate::stream::{Stream, StreamDetails}; use crate::topic::{Topic, TopicDetails}; use tokio::sync::Mutex; @@ -168,6 +168,78 @@ impl IggyClient { }) } + /// Gets all streams. + /// Returns a list of streams or a PyRuntimeError on failure. + #[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[list[Stream]]", imports=("collections.abc")))] + fn get_streams<'a>(&self, py: Python<'a>) -> PyResult> { + let inner = self.inner.clone(); + future_into_py(py, async move { + let streams = inner + .get_streams() + .await + .map_err(|e| PyErr::new::(e.to_string()))?; + Ok(streams.into_iter().map(Stream::from).collect::>()) + }) + } + + /// Updates a stream's name by id. + /// Returns Ok(()) on successful update or a PyRuntimeError on failure. + #[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[None]", imports=("collections.abc")))] + fn update_stream<'a>( + &self, + py: Python<'a>, + stream_id: PyIdentifier, + name: String, + ) -> PyResult> { + let stream_id = Identifier::try_from(stream_id)?; + let inner = self.inner.clone(); + future_into_py(py, async move { + inner + .update_stream(&stream_id, &name) + .await + .map_err(|e| PyErr::new::(e.to_string()))?; + Ok(()) + }) + } + + /// Deletes a stream by id. + /// Returns Ok(()) on successful deletion or a PyRuntimeError on failure. + #[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[None]", imports=("collections.abc")))] + fn delete_stream<'a>( + &self, + py: Python<'a>, + stream_id: PyIdentifier, + ) -> PyResult> { + let stream_id = Identifier::try_from(stream_id)?; + let inner = self.inner.clone(); + future_into_py(py, async move { + inner + .delete_stream(&stream_id) + .await + .map_err(|e| PyErr::new::(e.to_string()))?; + Ok(()) + }) + } + + /// Purges all messages from a stream by id. + /// Returns Ok(()) on successful purge or a PyRuntimeError on failure. + #[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[None]", imports=("collections.abc")))] + fn purge_stream<'a>( + &self, + py: Python<'a>, + stream_id: PyIdentifier, + ) -> PyResult> { + let stream_id = Identifier::try_from(stream_id)?; + let inner = self.inner.clone(); + future_into_py(py, async move { + inner + .purge_stream(&stream_id) + .await + .map_err(|e| PyErr::new::(e.to_string()))?; + Ok(()) + }) + } + /// Creates a new topic with the given parameters. /// Returns Ok(()) on successful topic creation or a PyRuntimeError on failure. #[pyo3( diff --git a/foreign/python/src/lib.rs b/foreign/python/src/lib.rs index 8c78456062..e5632d7896 100644 --- a/foreign/python/src/lib.rs +++ b/foreign/python/src/lib.rs @@ -31,7 +31,7 @@ use consumer::{ use pyo3::prelude::*; use receive_message::{PollingStrategy, ReceiveMessage}; use send_message::SendMessage; -use stream::StreamDetails; +use stream::{Stream, StreamDetails}; use topic::{Topic, TopicDetails}; /// A Python module implemented in Rust. @@ -41,6 +41,7 @@ fn apache_iggy(_py: Python, m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; m.add_class::()?; m.add_class::()?; m.add_class::()?; diff --git a/foreign/python/src/stream.rs b/foreign/python/src/stream.rs index 7343d2c53f..c4c7f5846b 100644 --- a/foreign/python/src/stream.rs +++ b/foreign/python/src/stream.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -use iggy::prelude::StreamDetails as RustStreamDetails; +use iggy::prelude::{Stream as RustStream, StreamDetails as RustStreamDetails}; use pyo3::prelude::*; use pyo3_stub_gen::derive::{gen_stub_pyclass, gen_stub_pymethods}; @@ -56,3 +56,39 @@ impl StreamDetails { self.inner.topics_count } } + +#[gen_stub_pyclass] +#[pyclass] +pub struct Stream { + pub(crate) inner: RustStream, +} + +impl From for Stream { + fn from(stream: RustStream) -> Self { + Self { inner: stream } + } +} + +#[gen_stub_pymethods] +#[pymethods] +impl Stream { + #[getter] + pub fn id(&self) -> u32 { + self.inner.id + } + + #[getter] + pub fn name(&self) -> String { + self.inner.name.to_string() + } + + #[getter] + pub fn messages_count(&self) -> u64 { + self.inner.messages_count + } + + #[getter] + pub fn topics_count(&self) -> u32 { + self.inner.topics_count + } +} diff --git a/foreign/python/tests/test_stream.py b/foreign/python/tests/test_stream.py index 5e71a543a4..8f757fab11 100644 --- a/foreign/python/tests/test_stream.py +++ b/foreign/python/tests/test_stream.py @@ -17,7 +17,7 @@ import pytest -from apache_iggy import IggyClient +from apache_iggy import IggyClient, SendMessage from .utils import get_server_config, wait_for_ping, wait_for_server @@ -203,3 +203,381 @@ async def test_create_stream_before_login_fails(self, unique_name): with pytest.raises(RuntimeError): await client.create_stream(unique_name()) + + +class TestGetStreams: + """Test listing streams via get_streams.""" + + @pytest.mark.asyncio + async def test_get_streams_returns_created_streams( + self, iggy_client: IggyClient, unique_name + ): + """Test get_streams returns every stream created during the test.""" + # Reverse-alphabetical names so that id-ascending (creation) order + # and name order disagree, proving the list isn't accidentally + # name-sorted. The client fixture is session-scoped, so other tests + # may have created streams; assert on the ones created here instead + # of the full server view. + created = [unique_name(f"z{index}") for index in range(3, 0, -1)] + for name in created: + await iggy_client.create_stream(name) + + streams = await iggy_client.get_streams() + by_name = {stream.name: stream for stream in streams} + assert set(created).issubset(by_name) + + mine = [by_name[name] for name in created] + assert [stream.id for stream in mine] == sorted(stream.id for stream in mine) + assert all(stream.messages_count == 0 for stream in mine) + assert all(stream.topics_count == 0 for stream in mine) + + @pytest.mark.asyncio + async def test_get_streams_reflects_topic_count( + self, iggy_client: IggyClient, unique_name + ): + """Test get_streams reports the topic count for a listed stream.""" + 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=1 + ) + + streams = await iggy_client.get_streams() + listed = next( + (stream for stream in streams if stream.name == stream_name), None + ) + assert listed is not None + assert listed.topics_count == 1 + + @pytest.mark.asyncio + async def test_get_streams_returns_same_result_when_called_repeatedly( + self, iggy_client: IggyClient, unique_name + ): + """Test repeated get_streams calls return an identically ordered view.""" + await iggy_client.create_stream(unique_name()) + + first = await iggy_client.get_streams() + second = await iggy_client.get_streams() + assert [stream.id for stream in first] == [stream.id for stream in second] + assert [stream.name for stream in first] == [stream.name for stream in second] + + @pytest.mark.asyncio + async def test_get_streams_before_connect_fails(self): + """Test get_streams requires an established connection.""" + host, port = get_server_config() + client = IggyClient(f"{host}:{port}") + + with pytest.raises(RuntimeError): + await client.get_streams() + + @pytest.mark.asyncio + async def test_get_streams_before_login_fails(self): + """Test get_streams requires authentication.""" + host, port = get_server_config() + wait_for_server(host, port) + + client = IggyClient(f"{host}:{port}") + await client.connect() + + with pytest.raises(RuntimeError): + await client.get_streams() + + +class TestUpdateStream: + """Test updating streams via update_stream.""" + + @pytest.mark.asyncio + async def test_update_stream_renames_stream( + self, iggy_client: IggyClient, unique_name + ): + """Test update_stream renames a stream; old name no longer resolves.""" + stream_name = unique_name() + new_name = unique_name() + + await iggy_client.create_stream(stream_name) + + await iggy_client.update_stream(stream_id=stream_name, name=new_name) + + renamed = await iggy_client.get_stream(new_name) + assert renamed is not None + assert renamed.name == new_name + + old = await iggy_client.get_stream(stream_name) + assert old is None + + @pytest.mark.asyncio + async def test_update_stream_preserves_id( + self, iggy_client: IggyClient, unique_name + ): + """Test update_stream keeps the same numeric id after a rename.""" + stream_name = unique_name() + new_name = unique_name() + + await iggy_client.create_stream(stream_name) + before = await iggy_client.get_stream(stream_name) + assert before is not None + + await iggy_client.update_stream(stream_id=stream_name, name=new_name) + + after = await iggy_client.get_stream(new_name) + assert after is not None + assert after.id == before.id + + @pytest.mark.asyncio + async def test_update_stream_by_numeric_id( + self, iggy_client: IggyClient, unique_name + ): + """Test update_stream accepts a numeric stream id.""" + stream_name = unique_name() + new_name = unique_name() + + await iggy_client.create_stream(stream_name) + stream = await iggy_client.get_stream(stream_name) + assert stream is not None + + await iggy_client.update_stream(stream_id=stream.id, name=new_name) + + renamed = await iggy_client.get_stream(new_name) + assert renamed is not None + assert renamed.name == new_name + + @pytest.mark.asyncio + async def test_update_stream_applies_repeated_updates( + self, iggy_client: IggyClient, unique_name + ): + """Test successive update_stream calls each take effect.""" + stream_name = unique_name() + first_rename = unique_name() + second_rename = unique_name() + + await iggy_client.create_stream(stream_name) + + await iggy_client.update_stream(stream_id=stream_name, name=first_rename) + after_first = await iggy_client.get_stream(first_rename) + assert after_first is not None + assert after_first.name == first_rename + assert await iggy_client.get_stream(stream_name) is None + + await iggy_client.update_stream(stream_id=first_rename, name=second_rename) + after_second = await iggy_client.get_stream(second_rename) + assert after_second is not None + assert after_second.name == second_rename + assert await iggy_client.get_stream(first_rename) is None + + @pytest.mark.asyncio + async def test_update_nonexistent_stream_fails( + self, iggy_client: IggyClient, unique_name + ): + """Test update_stream raises for a non-existent stream.""" + with pytest.raises(RuntimeError): + await iggy_client.update_stream(stream_id=unique_name(), name=unique_name()) + + @pytest.mark.asyncio + async def test_update_stream_to_existing_name_fails( + self, iggy_client: IggyClient, unique_name + ): + """Test update_stream rejects renaming a stream to a name already in use.""" + first_stream = unique_name() + second_stream = unique_name() + + await iggy_client.create_stream(first_stream) + await iggy_client.create_stream(second_stream) + + with pytest.raises(RuntimeError): + await iggy_client.update_stream(stream_id=second_stream, name=first_stream) + + @pytest.mark.asyncio + async def test_update_stream_requires_connection_and_auth(self, unique_name): + """Test update_stream fails both before connecting and before logging in.""" + host, port = get_server_config() + wait_for_server(host, port) + + client = IggyClient(f"{host}:{port}") + with pytest.raises(RuntimeError): + await client.update_stream(stream_id=unique_name(), name=unique_name()) + + await client.connect() + with pytest.raises(RuntimeError): + await client.update_stream(stream_id=unique_name(), name=unique_name()) + + +class TestDeleteStream: + """Test deleting streams via delete_stream.""" + + @pytest.mark.asyncio + async def test_delete_stream_removes_stream( + self, iggy_client: IggyClient, unique_name + ): + """Test delete_stream removes the stream so it no longer resolves.""" + stream_name = unique_name() + + await iggy_client.create_stream(stream_name) + + await iggy_client.delete_stream(stream_name) + + assert await iggy_client.get_stream(stream_name) is None + + @pytest.mark.asyncio + async def test_delete_stream_by_numeric_id( + self, iggy_client: IggyClient, unique_name + ): + """Test delete_stream accepts a numeric stream id.""" + stream_name = unique_name() + + await iggy_client.create_stream(stream_name) + stream = await iggy_client.get_stream(stream_name) + assert stream is not None + + await iggy_client.delete_stream(stream.id) + + assert await iggy_client.get_stream(stream_name) is None + + @pytest.mark.asyncio + async def test_delete_stream_leaves_other_streams( + self, iggy_client: IggyClient, unique_name + ): + """Test delete_stream removes only the targeted stream.""" + stream_to_delete = unique_name() + stream_to_keep = unique_name() + + await iggy_client.create_stream(stream_to_delete) + await iggy_client.create_stream(stream_to_keep) + + await iggy_client.delete_stream(stream_to_delete) + + assert await iggy_client.get_stream(stream_to_delete) is None + kept = await iggy_client.get_stream(stream_to_keep) + assert kept is not None + assert kept.name == stream_to_keep + + @pytest.mark.asyncio + async def test_delete_nonexistent_stream_fails( + self, iggy_client: IggyClient, unique_name + ): + """Test delete_stream raises for a non-existent stream.""" + with pytest.raises(RuntimeError): + await iggy_client.delete_stream(unique_name()) + + @pytest.mark.asyncio + async def test_delete_stream_twice_fails_second_time( + self, iggy_client: IggyClient, unique_name + ): + """Test deleting an already-deleted stream raises on the second call.""" + stream_name = unique_name() + + await iggy_client.create_stream(stream_name) + + await iggy_client.delete_stream(stream_name) + with pytest.raises(RuntimeError): + await iggy_client.delete_stream(stream_name) + + @pytest.mark.asyncio + async def test_delete_stream_requires_connection_and_auth(self, unique_name): + """Test delete_stream fails both before connecting and before logging in.""" + host, port = get_server_config() + wait_for_server(host, port) + + client = IggyClient(f"{host}:{port}") + with pytest.raises(RuntimeError): + await client.delete_stream(unique_name()) + + await client.connect() + with pytest.raises(RuntimeError): + await client.delete_stream(unique_name()) + + +class TestPurgeStream: + """Test purging stream messages via purge_stream.""" + + @pytest.mark.asyncio + async def test_purge_stream_clears_messages_but_keeps_stream( + self, iggy_client: IggyClient, unique_name + ): + """Test purge_stream empties the stream while leaving it in place.""" + 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=1 + ) + + messages = [SendMessage(f"payload-{index}") for index in range(5)] + await iggy_client.send_messages(stream_name, topic_name, 0, messages) + + before = await iggy_client.get_stream(stream_name) + assert before is not None + assert before.messages_count == 5 + + await iggy_client.purge_stream(stream_name) + + after = await iggy_client.get_stream(stream_name) + # Purging clears messages only; the stream itself survives (purge is + # not delete) and keeps its identity and topics. + assert after is not None + assert after.messages_count == 0 + assert after.id == before.id + assert after.name == before.name + assert after.topics_count == before.topics_count + + @pytest.mark.asyncio + async def test_purge_empty_stream_succeeds( + self, iggy_client: IggyClient, unique_name + ): + """Test purge_stream is a no-op on a stream with no messages.""" + stream_name = unique_name() + + await iggy_client.create_stream(stream_name) + + await iggy_client.purge_stream(stream_name) + + stream = await iggy_client.get_stream(stream_name) + assert stream is not None + assert stream.messages_count == 0 + + @pytest.mark.asyncio + async def test_purge_stream_is_idempotent_when_called_repeatedly( + self, iggy_client: IggyClient, unique_name + ): + """Test purge_stream succeeds when called repeatedly on the same stream.""" + 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=1 + ) + + messages = [SendMessage(f"payload-{index}") for index in range(5)] + await iggy_client.send_messages(stream_name, topic_name, 0, messages) + + await iggy_client.purge_stream(stream_name) + await iggy_client.purge_stream(stream_name) + + stream = await iggy_client.get_stream(stream_name) + assert stream is not None + assert stream.messages_count == 0 + + @pytest.mark.asyncio + async def test_purge_nonexistent_stream_fails( + self, iggy_client: IggyClient, unique_name + ): + """Test purge_stream raises for a non-existent stream.""" + with pytest.raises(RuntimeError): + await iggy_client.purge_stream(unique_name()) + + @pytest.mark.asyncio + async def test_purge_stream_requires_connection_and_auth(self, unique_name): + """Test purge_stream fails both before connecting and before logging in.""" + host, port = get_server_config() + wait_for_server(host, port) + + client = IggyClient(f"{host}:{port}") + with pytest.raises(RuntimeError): + await client.purge_stream(unique_name()) + + await client.connect() + with pytest.raises(RuntimeError): + await client.purge_stream(unique_name())