Testing with Iggy
Run the Iggy Docker image from your own test suite, with examples for Rust, Python and Node.js.
Iggy has no embedded, in-process test server. The server ships as a binary, not a library, so a test suite can't start it in-process. To test against a real server, start the published Docker image when your tests start and stop it when they finish. The examples below use each language's testcontainers library to do that.
What the container needs
These settings apply whichever language your tests are in:
seccomp=unconfinedas a security option. Docker's default seccomp profile blocks theio_uringcalls the server needs.IGGY_ROOT_USERNAMEandIGGY_ROOT_PASSWORD. Without them the server generates a random root password and only logs it once, which a test can't easily pick up.IGGY_TCP_ADDRESS=0.0.0.0:8090andIGGY_NODE_ADVERTISED_ADDRESS=localhost. By default the server listens on127.0.0.1inside the container, which a mapped port can't reach. Listening on all addresses also needs an advertised address, or the server won't start.IGGY_SHARDING_CPU_ALLOCATION=allon Docker Desktop, WSL2 and other hosts with a single NUMA node. Without it, 0.9.0 exits withfailed to bind shard 0 memory to its NUMA node. It does no harm elsewhere.
Wait for the server client listeners started log line before connecting. The server logs server running; waiting on shard threads earlier, before the port is open, and connecting on that line fails some of the time.
If the server exits with io_uring was denied locked memory, the container's locked-memory limit is too low for it. Raise it, as the --ulimit memlock=-1:-1 flag does in Docker & Helm.
Rust
This example uses testcontainers 0.27 and the iggy crate 0.11, which goes with server 0.9.0.
[dev-dependencies]
iggy = "0.11"
tokio = { version = "1", features = ["full"] }
testcontainers = "0.27"use iggy::prelude::*;
use std::str::FromStr;
use testcontainers::core::{IntoContainerPort, WaitFor};
use testcontainers::runners::AsyncRunner;
use testcontainers::{GenericImage, ImageExt};
#[tokio::test]
async fn send_and_poll_roundtrip() {
let container = GenericImage::new("apache/iggy", "0.9.0")
.with_exposed_port(8090.tcp())
.with_wait_for(WaitFor::message_on_stdout(
"server client listeners started",
))
.with_env_var("IGGY_ROOT_USERNAME", "iggy")
.with_env_var("IGGY_ROOT_PASSWORD", "iggy")
.with_env_var("IGGY_TCP_ADDRESS", "0.0.0.0:8090")
.with_env_var("IGGY_NODE_ADVERTISED_ADDRESS", "localhost")
.with_env_var("IGGY_SHARDING_CPU_ALLOCATION", "all")
.with_security_opt("seccomp=unconfined")
.start()
.await
.expect("failed to start apache/iggy container");
let host = container.get_host().await.unwrap();
let port = container.get_host_port_ipv4(8090).await.unwrap();
let client = IggyClientBuilder::new()
.with_tcp()
.with_server_address(format!("{host}:{port}"))
.build()
.unwrap();
client.connect().await.unwrap();
client
.login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
.await
.unwrap();
let stream = client.create_stream("test-stream").await.unwrap();
let topic = client
.create_topic(
&Identifier::named("test-stream").unwrap(),
"test-topic",
&TopicCreateOptions {
partitions_count: Some(1),
message_expiry: Some(IggyExpiry::NeverExpire),
..TopicCreateOptions::default()
},
)
.await
.unwrap();
let mut messages = vec![IggyMessage::from_str("hello from cargo test").unwrap()];
client
.send_messages(
&Identifier::numeric(stream.id).unwrap(),
&Identifier::numeric(topic.id).unwrap(),
&Partitioning::partition_id(0),
&mut messages,
)
.await
.unwrap();
let consumer = Consumer::default();
let polled = client
.poll_messages(
&Identifier::numeric(stream.id).unwrap(),
&Identifier::numeric(topic.id).unwrap(),
Some(0),
&consumer,
&PollingStrategy::offset(0),
1,
true,
)
.await
.unwrap();
assert_eq!(polled.messages.len(), 1);
assert_eq!(polled.messages[0].payload.as_ref(), b"hello from cargo test");
}Run it with cargo test.
Python
This example uses testcontainers, pytest and pytest-asyncio with apache-iggy 0.9.0:
pip install apache-iggy==0.9.0 testcontainers pytest pytest-asyncioimport asyncio
from datetime import timedelta
import pytest
from testcontainers.core.container import DockerContainer
from testcontainers.core.wait_strategies import LogMessageWaitStrategy
from apache_iggy import (
AutoLogin,
Consumer,
IggyClient,
PollingStrategy,
SendMessage,
TcpConfig,
TcpReconnectionConfig,
)
@pytest.fixture(scope="module")
def iggy_container():
container = (
DockerContainer("apache/iggy:0.9.0")
.with_exposed_ports(8090)
.with_env("IGGY_ROOT_USERNAME", "iggy")
.with_env("IGGY_ROOT_PASSWORD", "iggy")
.with_env("IGGY_TCP_ADDRESS", "0.0.0.0:8090")
.with_env("IGGY_NODE_ADVERTISED_ADDRESS", "localhost")
.with_env("IGGY_SHARDING_CPU_ALLOCATION", "all")
.with_kwargs(security_opt=["seccomp=unconfined"])
.waiting_for(LogMessageWaitStrategy("server client listeners started"))
)
with container:
yield container
@pytest.mark.asyncio
async def test_send_and_poll_roundtrip(iggy_container):
host = iggy_container.get_container_host_ip()
port = iggy_container.get_exposed_port(8090)
client = IggyClient(
TcpConfig(
server_address=f"{host}:{port}",
auto_login=AutoLogin.username_password("iggy", "iggy"),
reconnection=TcpReconnectionConfig(
enabled=True, interval=timedelta(seconds=1)
),
)
)
await client.connect()
await client.create_stream(name="test-stream")
await client.create_topic(
stream="test-stream", partitions_count=1, name="test-topic"
)
await client.send_messages(
stream="test-stream",
topic="test-topic",
partitioning=0,
messages=[SendMessage("hello from pytest")],
)
for _ in range(10):
polled = await client.poll_messages(
stream="test-stream",
topic="test-topic",
consumer=Consumer.Single("test-consumer"),
partition_id=0,
polling_strategy=PollingStrategy.Next(),
count=1,
auto_commit=True,
)
if polled:
break
await asyncio.sleep(0.5)
assert len(polled) == 1
assert polled[0].payload().decode("utf-8") == "hello from pytest"Run it with pytest -v.
Node.js
This example uses testcontainers with apache-iggy 0.10.0, the Node.js SDK that goes with server 0.9.0, and the built-in node:test runner. testcontainers needs Node.js 22.19 or newer.
npm install apache-iggy testcontainersimport assert from 'node:assert/strict';
import { test, before, after } from 'node:test';
import { GenericContainer, Wait } from 'testcontainers';
import { Client, Partitioning, Consumer, PollingStrategy } from 'apache-iggy';
let container;
let client;
before(async () => {
container = await new GenericContainer('apache/iggy:0.9.0')
.withExposedPorts(8090)
.withEnvironment({
IGGY_ROOT_USERNAME: 'iggy',
IGGY_ROOT_PASSWORD: 'iggy',
IGGY_TCP_ADDRESS: '0.0.0.0:8090',
IGGY_NODE_ADVERTISED_ADDRESS: 'localhost',
IGGY_SHARDING_CPU_ALLOCATION: 'all',
})
.withSecurityOpt('seccomp=unconfined')
.withWaitStrategy(Wait.forLogMessage('server client listeners started'))
.start();
client = new Client({
transport: 'TCP',
options: {
host: container.getHost(),
port: container.getMappedPort(8090),
keepAlive: true,
},
reconnect: { enabled: true, interval: 5000, maxRetries: 5 },
heartbeatInterval: 5000,
credentials: { username: 'iggy', password: 'iggy' },
});
});
after(async () => {
await client?.destroy();
await container?.stop();
});
test('send and poll roundtrip', async () => {
const stream = await client.stream.create({ name: 'test-stream' });
const topic = await client.topic.create({
streamId: stream.id,
name: 'test-topic',
partitionCount: 1,
compressionAlgorithm: 1,
});
await client.message.send({
streamId: stream.id,
topicId: topic.id,
messages: [{ id: 1, headers: [], payload: 'hello from node:test' }],
partition: Partitioning.PartitionId(topic.partitions[0].id),
});
let polled;
for (let i = 0; i < 10; i++) {
polled = await client.message.poll({
streamId: stream.id,
topicId: topic.id,
consumer: Consumer.Single,
partitionId: topic.partitions[0].id,
pollingStrategy: PollingStrategy.Offset(0n),
count: 1,
autocommit: false,
});
if (polled?.messages?.length) break;
await new Promise((resolve) => setTimeout(resolve, 500));
}
assert.equal(polled.messages.length, 1);
assert.equal(polled.messages[0].payload.toString('utf-8'), 'hello from node:test');
});Save it as iggy.test.mjs and run it with node --test.
