tree: 532a704481716b4931d233995f2b0af718663989
  1. Benchmarks/
  2. Iggy_SDK/
  3. Iggy_SDK.Tests.BDD/
  4. Iggy_SDK.Tests.Integration/
  5. Iggy_SDK_Tests/
  6. scripts/
  7. Shared/
  8. .dockerignore
  9. .editorconfig
  10. Directory.Build.props
  11. Directory.Packages.props
  12. global.json
  13. Iggy_SDK.sln
  14. Iggy_SDK.sln.DotSettings
  15. LICENSE
  16. NOTICE
  17. README.md
foreign/csharp/README.md

C# SDK for Iggy Nuget (with prereleases)

Overview

The Apache Iggy C# SDK provides a comprehensive client library for interacting with Iggy message streaming servers. It offers a modern, async-first API with support for multiple transport protocols and comprehensive message streaming capabilities.

Getting Started

Installation

Install the NuGet package:

dotnet add package Apache.Iggy

Supported Protocols

The SDK supports two transport protocols:

  • TCP - Binary protocol for optimal performance and lower latency (recommended)
  • HTTP - RESTful JSON API for stateless operations

Over TCP the SDK speaks the VSR consensus framing, which is the only wire protocol the server accepts.

See Viewstamped Replication (VSR) for what that means for the client API.

Creating a Client

The SDK is built around the IIggyClient interface. To create a client instance:

var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp
});

await client.ConnectAsync();

Optionally, you can provide an ILoggerFactory for diagnostics and debugging (defaults to NullLoggerFactory.Instance):

var loggerFactory = LoggerFactory.Create(builder =>
{
    builder
        .AddFilter("Apache.Iggy", LogLevel.Information)
        .AddConsole();
});

var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp,
    LoggerFactory = loggerFactory
});

await client.ConnectAsync();

Configuration

The IggyClientConfigurator provides comprehensive configuration options:

var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp,

    // Buffer sizes (optional, default: 4096)
    ReceiveBufferSize = 4096,
    SendBufferSize = 4096,

    // TLS/SSL configuration
    TlsSettings = new TlsSettings
    {
        Enabled = true,
        Hostname = "iggy",
        CertificatePath = "/path/to/cert"
    },

    // Idle ping keeping the session alive under server-side heartbeat verification (TCP only).
    // Default 5 seconds, must be between 1 millisecond and about 49 days. Init-only.
    HeartbeatInterval = TimeSpan.FromSeconds(5),

    // Automatic reconnection with exponential backoff (enabled by default, infinite retries)
    ReconnectionSettings = new ReconnectionSettings
    {
        Enabled = true,
        MaxRetries = 0,              // 0 = infinite retries
        InitialDelay = TimeSpan.FromSeconds(5),
        MaxDelay = TimeSpan.FromSeconds(30),
        WaitAfterReconnect = TimeSpan.FromSeconds(1),
        UseExponentialBackoff = true,
        BackoffMultiplier = 2.0
    },

    // Auto-login after connection. Optional for reconnection: a client that signs in with
    // LoginUserAsync has that sign-in replayed on a reconnect too. Without either, a reconnect
    // cannot restore the session and a lost connection fails the request
    AutoLoginSettings = AutoLoginSettings.For("your_username", "your_password"),
    // or AutoLoginSettings.ForPersonalAccessToken("your_token")

    // Optional: logging
    LoggerFactory = loggerFactory
});

await client.ConnectAsync();

Viewstamped Replication (VSR)

Over TCP every request is wrapped in a 256-byte consensus header, the client registers a consensus session at login, and writes are replicated before they are acknowledged. The IIggyClient surface is unchanged, with the few exceptions listed under Limitations.

var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp,

    // Upper bound on a reply frame the server announces, 64 MiB by default.
    MaxResponseFrameSize = 64 * 1024 * 1024,

    AutoLoginSettings = AutoLoginSettings.For("iggy", "iggy")
});

await client.ConnectAsync();

What changes under VSR

  • Login binds a session. LoginUserAsync / LoginWithPersonalAccessTokenAsync run the register handshake at connect time, and the session lives for as long as the connection. Logging out, being evicted or losing the connection ends it, and the next login registers a fresh one.
  • Leader redirection is automatic. The client reads the cluster roster, follows the current leader and re-checks it when a request is refused because the node stopped being primary.
  • The client picks partitions. The broker never routes: balanced and message-key partitioning are resolved client-side (the message-key hash matches the Rust SDK byte for byte), and consumer-group polls round-robin over the partitions the coordinator assigned to this client.
  • Consumer groups are assignment-based. JoinConsumerGroupAsync makes this client a member; the assignment is synced on demand and refreshed on every PingAsync. Partition counts are cached for 30 seconds, so a topic another client widens is picked up without waiting for a ping.
  • Credentials are bounds-checked locally. A username outside 3-50 bytes, a password outside 3-100 bytes or a personal access token outside 1-255 bytes is rejected before the register body is framed.
  • PingAsync costs more than a ping. Besides the ping it re-syncs the assignment of every consumer group this client has joined, so it makes one extra round trip per joined group. The TCP client pings on its own every HeartbeatInterval (5 seconds by default) while connected, so an idle session survives the server's heartbeat verification and assignments stay fresh; a lost session is repaired by the regular reconnect and auto login.

Retries and failed requests

The SDK replays a request whenever the server says it never admitted it. Two cases surface to the caller:

  • IggyInvalidStatusCodeException carries the server status code, with FromServer telling apart a verdict the cluster reported from a failure the client raised itself.
  • VsrRequestOutcomeUnknownException means no server verdict arrived after the request was written - the connection was lost, the call was cancelled, or the server evicted the session while the request was in flight - so the cluster may or may not have committed it. The SDK will not replay it on a new session, because that would bypass server-side deduplication - re-issuing it is the caller‘s decision. IggyPublisher will not retry it either: it reports the batch through the message-batch-failed event, and IggyConsumer rethrows it rather than swallowing it, because an auto-committing poll may have advanced the offset already. Rethrowing ends the consumer’s polling loop: catch it around the enumeration, decide whether the operation is safe to re-issue, and start consuming again.

Limitations

  • VSR requires Protocol.Tcp; configuring it with Protocol.Http throws at client creation.
  • StoreOffsetAsync / DeleteOffsetAsync need an explicit partition id under VSR: the broker does not resolve a null partition for a consumer-offset request, so passing one throws client-side.
  • FlushUnsavedBufferAsync is not available under VSR; the server refuses it.
  • Polling a topic that does not exist returns an empty poll rather than throwing. The server answers an unresolved topic with the empty-poll reply shape, so the client cannot tell it apart from a topic with no messages. Check the topic exists first if the distinction matters.

Behaviour changes for existing clients

  • MaxResponseFrameSize bounds the reply frames the VSR reader accepts. A reply larger than the 64 MiB default is refused and the connection is dropped, so raise it if a single response legitimately exceeds that
    • a large GetSnapshotAsync is the usual case.
  • Clients built with IggyConsumerBuilder / IggyPublisherBuilder now auto-login with the credentials passed to WithConnection. Before, a builder-created client came back from a reconnect unauthenticated; now the credentials are held for the lifetime of the connection and replayed.
  • The TCP client now pings the server every HeartbeatInterval (5 seconds, always on) on its own, and reconnection is on by default (it was off before) with unlimited retries, like the Rust client. Bringing an endpoint up is what gets another attempt, and one retry is a full pass over every endpoint the client knows - where it is, the configured address, and every node the roster named - rather than one dial of the first. Bad credentials or a missing leader is thrown right away, and so is a TLS fault no retry can fix (an unreadable CA file, a certificate this client will never accept), once the pass has given the other endpoints their turn. A dropped connection fails every in-flight request at once; they share a single reconnect and are replayed on the connection it establishes. With the default MaxRetries = 0 an unreachable server is retried forever, so a request that passes no CancellationToken waits for as long as the server stays down - set MaxRetries or pass a token to bound it. A reconnect restores the session from AutoLoginSettings or from the sign-in a LoginUserAsync call succeeded with, so a client that logged in by hand reconnects too; without either there is nothing to restore and the request fails. A server-side eviction (the heartbeat verifier reacting to silence) is recovered from the same way; only LogoutUserAsync or Dispose ends a session for good. Set ReconnectionSettings.Enabled = false to opt out of reconnection.
  • AutoLoginSettings properties are now init-only, as is IggyClientConfigurator.HeartbeatInterval. Build them with an object initializer or the AutoLoginSettings.For / AutoLoginSettings.ForPersonalAccessToken factories instead of assigning after construction.
  • IggyConsumerBuilder / IggyPublisherBuilder accept a personal access token through the WithConnection(protocol, address, personalAccessToken, ...) overload, as an alternative to a username and password.
  • The SDK now ships a dependency on System.IO.Hashing, used for the client-side message-key partitioner.
  • TCP sockets are opened with NoDelay. The protocol is request/reply, so a write is always the last one before the client waits for the answer and Nagle has nothing to coalesce it with - it only held back the trailing segment of a large request until the previous one was acked.

Authentication

User Login

Begin by using the root account (note: the root account cannot be removed or updated):

var response = await client.LoginUserAsync("iggy", "iggy");

Creating Users

Create new users with customizable permissions:

var permissions = new Permissions
{
    Global = new GlobalPermissions
    {
        ManageServers = true,
        ManageUsers = true,
        ManageStreams = true,
        ManageTopics = true,
        PollMessages = true,
        ReadServers = true,
        ReadStreams = true,
        ReadTopics = true,
        ReadUsers = true,
        SendMessages = true
    }
};

await client.CreateUserAsync("test_user", "secure_password", UserStatus.Active, permissions);

// Login with the new user
var loginResponse = await client.LoginUserAsync("test_user", "secure_password");

Personal Access Tokens

Create and use Personal Access Tokens (PAT) for programmatic access:

// Create a PAT
var patResponse = await client.CreatePersonalAccessTokenAsync("api-token", 3600);

// Login with PAT
await client.LoginWithPersonalAccessTokenAsync(patResponse.Token);

Streams and Topics

Creating Streams

await client.CreateStreamAsync("my-stream");

You can reference streams by either numeric ID or name:

var streamById = Identifier.Numeric(0);
var streamByName = Identifier.String("my-stream");

Creating Topics

Every stream contains topics for organizing messages:

var streamId = Identifier.String("my-stream");

await client.CreateTopicAsync(
    streamId,
    name: "my-topic",
    partitionsCount: 3,
    compressionAlgorithm: CompressionAlgorithm.None,
    messageExpiry: 0,  // 0 = never expire
    maxTopicSize: 0    // 0 = unlimited
);

Note: Stream and topic names use hyphens instead of spaces. Iggy automatically replaces spaces with hyphens.

Publishing Messages

Sending Messages

Send messages using the publisher interface:

var streamId = Identifier.String("my-stream");
var topicId = Identifier.String("my-topic");

var messages = new List<Message>
{
    new(Guid.NewGuid(), "Hello, Iggy!"u8.ToArray()),
    new(1, "Another message"u8.ToArray())
};

await client.SendMessagesAsync(
    streamId,
    topicId,
    Partitioning.None(),  // balanced partitioning
    messages
);

Partitioning Strategies

Control which partition receives each message:

// Balanced partitioning (default)
Partitioning.None()

// Send to specific partition
Partitioning.PartitionId(1)

// Key-based partitioning (string)
Partitioning.EntityIdString("user-123")

// Key-based partitioning (integer)
Partitioning.EntityIdInt(12345)

// Key-based partitioning (GUID)
Partitioning.EntityIdGuid(Guid.NewGuid())

User-Defined Headers

Add custom headers to messages with typed values:

var headers = new Dictionary<HeaderKey, HeaderValue>
{
    { new HeaderKey { Value = "correlation_id" }, HeaderValue.FromString("req-123") },
    { new HeaderKey { Value = "priority" }, HeaderValue.FromInt32(1) },
    { new HeaderKey { Value = "timeout" }, HeaderValue.FromInt64(5000) },
    { new HeaderKey { Value = "confidence" }, HeaderValue.FromFloat(0.95f) },
    { new HeaderKey { Value = "is_urgent" }, HeaderValue.FromBool(true) },
    { new HeaderKey { Value = "request_id" }, HeaderValue.FromGuid(Guid.NewGuid()) }
};

var messages = new List<Message>
{
    new(Guid.NewGuid(), "Message with headers"u8.ToArray(), headers)
};

await client.SendMessagesAsync(
    streamId,
    topicId,
    Partitioning.PartitionId(1),
    messages
);

Flushing Unsaved Buffer

Force a flush of the in-memory buffer to disk for a specific partition:

await client.FlushUnsavedBufferAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    partitionId: 1,
    fsync: true
);

Consumer Groups

Creating Consumer Groups

Coordinate message consumption across multiple consumers:

var groupResponse = await client.CreateConsumerGroupAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    "my-consumer-group"
);

Joining and Leaving Groups

Note: Join/Leave operations are only supported on TCP protocol and will throw FeatureUnavailableException on HTTP.

// Join a consumer group
await client.JoinConsumerGroupAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Identifier.Numeric(1)  // consumer ID
);

// Leave a consumer group
await client.LeaveConsumerGroupAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Identifier.Numeric(1)  // consumer ID
);

Consuming Messages

Fetching Messages

Fetch a batch of messages:

var polledMessages = await client.PollMessagesAsync(new MessageFetchRequest
{
    StreamId = streamId,
    TopicId = topicId,
    Consumer = Consumer.New(1), // or Consumer.Group("my-consumer-group") for consumer group
    Count = 10,
    PartitionId = 0, // optional, null for consumer group
    PollingStrategy = PollingStrategy.Next(),
    AutoCommit = true
});

foreach (var message in polledMessages.Messages)
{
    Console.WriteLine($"Message: {Encoding.UTF8.GetString(message.Payload)}");
}

Polling Strategies

Control where message consumption starts:

// Start from a specific offset
PollingStrategy.Offset(1000)

// Start from a specific timestamp (microseconds since epoch)
PollingStrategy.Timestamp(1699564800000000)

// Start from the first message
PollingStrategy.First()

// Start from the last message
PollingStrategy.Last()

// Start from the next unread message
PollingStrategy.Next()

Offset Management

Storing Offsets

Store the current consumer position:

await client.StoreOffsetAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Identifier.Numeric(1),  // consumer ID
    0,                      // partition ID
    42                      // offset value
);

Retrieving Offsets

Get the current stored offset:

var offsetInfo = await client.GetOffsetAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Identifier.Numeric(1),  // consumer ID
    0                       // partition ID
);

Console.WriteLine($"Current offset: {offsetInfo.Offset}");

Deleting Offsets

Clear stored offsets:

await client.DeleteOffsetAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Identifier.Numeric(1),  // consumer ID
    0                       // partition ID
);

System Operations

Cluster Information

Get cluster metadata and node information:

var metadata = await client.GetClusterMetadataAsync();

Server Statistics

Retrieve server performance metrics:

var stats = await client.GetStatsAsync();

Health Checks

Verify server connectivity:

await client.PingAsync();

Client Information

Get information about connected clients:

var clients = await client.GetClientsAsync();
var currentClient = await client.GetMeAsync();

Snapshots

Capture a system snapshot as a compressed ZIP archive:

var snapshotBytes = await client.GetSnapshotAsync(
    SnapshotCompression.Zstd,
    new List<SystemSnapshotType>
    {
        SystemSnapshotType.ServerLogs,
        SystemSnapshotType.ServerConfig,
        SystemSnapshotType.ResourceUsage
    }
);

// Or capture everything
var fullSnapshot = await client.GetSnapshotAsync(
    SnapshotCompression.Deflated,
    new List<SystemSnapshotType> { SystemSnapshotType.All }
);

Available compression methods: Stored, Deflated, Bzip2, Zstd, Lzma, Xz.

Available snapshot types: FilesystemOverview, ProcessList, ResourceUsage, Test, ServerLogs, ServerConfig, All.

Segment Management

Delete the last N segments from a partition:

await client.DeleteSegmentsAsync(
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    partitionId: 1,
    segmentsCount: 2
);

Event Subscription

Subscribe to connection events:

// Subscribe to connection events
client.SubscribeConnectionEvents(async connectionState =>
{
    Console.WriteLine($"Current connection state: {connectionState.CurrentState}");

    await SaveConnectionStateLog(connectionState.CurrentState);
});

// Unsubscribe
client.UnsubscribeConnectionEvents(handler);

Advanced: IggyPublisher

High-level publisher with background sending, retries, and encryption:

var publisher = IggyPublisherBuilder.Create(
    client,
    Identifier.String("my-stream"),
    Identifier.String("my-topic")
)
.WithBackgroundSending(enabled: true, batchSize: 100)
.WithRetry(maxAttempts: 3)
.Build();

await publisher.InitAsync();

var messages = new List<Message>
{
    new(Guid.NewGuid(), "Message 1"u8.ToArray()),
    new(0, "Message 2"u8.ToArray())
};

await publisher.SendMessages(messages);

// Wait for all messages to be sent
await publisher.WaitUntilAllSends();
await publisher.DisposeAsync();

For automatic object serialization, use the typed variant:

class OrderSerializer : ISerializer<Order>
{
    public byte[] Serialize(Order data) =>
        Encoding.UTF8.GetBytes(JsonSerializer.Serialize(data));
}

var publisher = IggyPublisherBuilder<Order>.Create(
    client,
    Identifier.String("orders-stream"),
    Identifier.String("orders-topic"),
    new OrderSerializer()
).Build();

await publisher.InitAsync();
await publisher.SendAsync(new List<Order> { /* ... */ });

Advanced: IggyConsumer

High-level consumer with automatic offset management and consumer groups:

var consumer = IggyConsumerBuilder.Create(
    client,
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Consumer.New(1)
)
.WithPollingStrategy(PollingStrategy.Next())
.WithBatchSize(10)
.WithAutoCommitMode(AutoCommitMode.Auto)
.Build();

await consumer.InitAsync();

await foreach (var message in consumer.ReceiveAsync())
{
    var payload = Encoding.UTF8.GetString(message.Message.Payload);
    Console.WriteLine($"Offset {message.CurrentOffset}: {payload}");
}

For consumer groups (load-balanced across multiple consumers):

var consumer = IggyConsumerBuilder.Create(
    client,
    Identifier.String("my-stream"),
    Identifier.String("my-topic"),
    Consumer.Group("my-group")
)
.WithConsumerGroup("my-group", createIfNotExists: true)
.WithPollingStrategy(PollingStrategy.Next())
.WithAutoCommitMode(AutoCommitMode.AfterReceive)
.Build();

await consumer.InitAsync();

await foreach (var message in consumer.ReceiveAsync())
{
    Console.WriteLine($"Partition {message.PartitionId}: {message.Message.Payload}");
}

await consumer.DisposeAsync();

For automatic deserialization:

class OrderDeserializer : IDeserializer<OrderEvent>
{
    public OrderEvent Deserialize(byte[] data) =>
        JsonSerializer.Deserialize<OrderEvent>(Encoding.UTF8.GetString(data))!;
}

var consumer = IggyConsumerBuilder<OrderEvent>.Create(
    client,
    Identifier.String("orders-stream"),
    Identifier.String("orders-topic"),
    Consumer.Group("order-processors"),
    new OrderDeserializer()
)
.WithAutoCommitMode(AutoCommitMode.Auto)
.Build();

await consumer.InitAsync();

await foreach (var message in consumer.ReceiveDeserializedAsync())
{
    if (message.Status == MessageStatus.Success)
    {
        Console.WriteLine($"Order: {message.Data?.OrderId}");
    }
}

API Reference

The SDK provides the following main interfaces:

  • IIggyClient - Main client interface (aggregates all features)
  • IIggyPublisher - High-level message publishing interface
  • IIggyConsumer - High-level message consumption interface
  • IIggyStream - Stream management
  • IIggyTopic - Topic management
  • IIggyOffset - Offset management
  • IIggyConsumerGroup - Consumer group operations
  • IIggyPartition - Partition operations
  • IIggySegment - Segment management
  • IIggyUsers - User and authentication management
  • IIggySystem - System and cluster operations
  • IIggyPersonalAccessToken - Personal access token management

Additionally, builder-based APIs are available:

  • IggyPublisherBuilder / IggyPublisherBuilder - Fluent publisher configuration
  • IggyConsumerBuilder / IggyConsumerBuilder - Fluent consumer configuration

Running Examples

Examples are located in examples/csharp/src/ in root iggy directory. Available examples:

  • GettingStarted - Basic producer/consumer setup
  • Basic - Simple message publishing and consuming
  • MessageHeaders - Using custom message headers
  • MessageEnvelope - Envelope pattern for message serialization
  • NewSdk - High-level IggyPublisher/IggyConsumer API

Start the Iggy server:

cargo run --bin iggy-server

Run an example (from the examples/csharp/ directory):

dotnet run -c Release --project src/GettingStarted/Iggy_SDK.Examples.GettingStarted.Producer
dotnet run -c Release --project src/GettingStarted/Iggy_SDK.Examples.GettingStarted.Consumer

Integration Tests

Integration tests are located in Iggy_SDK.Tests.Integration/. Tests can run against:

  • A dockerized Iggy server with TestContainers
  • A local Iggy server (set IGGY_SERVER_HOST environment variable)

Requirements

  • .NET 10 SDK
  • Docker (for TestContainers tests)

Running Integration Tests Locally

1. Dockerization

The suite runs against iggy-server. TCP only: the SDK frames TCP with the VSR wire protocol, the cluster serves reads from the primary, and the HTTP surface has no equivalent path to route them through.

cargo build --bin iggy-server --bin iggy

docker build --no-cache -f core/server/Dockerfile --platform linux/amd64 --target runtime-prebuilt --build-arg PREBUILT_IGGY_SERVER=target/debug/iggy-server --build-arg PREBUILT_IGGY_CLI=target/debug/iggy -t iggy-server:test .

2. Build the Test Project

dotnet build foreign/csharp/Iggy_SDK.Tests.Integration

3. Run test

cd foreign/csharp
export IGGY_SERVER_DOCKER_IMAGE=iggy-server:test
dotnet test -f net10.0 --project Iggy_SDK.Tests.Integration --no-build --verbosity diagnostic

IGGY_SERVER_DOCKER_IMAGE defaults to iggy-server:test, so the export above is only needed to point at a different image. Rider and Visual Studio need nothing configured.

Useful Resources

ROADMAP - TODO

  • [ ] Error handling with status codes and descriptions
  • [ ] Add support for ASP.NET Core Dependency Injection