Skip to main content

Documentation Index

Fetch the complete documentation index at: https://mintlify.com/iii-hq/sdk/llms.txt

Use this file to discover all available pages before exploring further.

Overview

The Streams API provides atomic operations on stream data. Streams are key-value stores that support atomic updates with multiple operations (increment, set, merge, etc.).

Creating a Streams Instance

Streams::new

Create a new Streams instance for performing atomic updates.
pub fn new(iii: III) -> Self
iii
III
required
III client instance
Example:
use iii_sdk::{III, Streams};

let iii = III::new("ws://localhost:49134");
iii.connect().await?;

let streams = Streams::new(iii);

Stream Keys

Stream keys follow the format: stream_name::group_id::item_id
  • stream_name: The stream namespace
  • group_id: Group identifier (e.g., user ID, session ID)
  • item_id: Item identifier within the group
Examples:
  • orders::user-123::order-456
  • counters::daily::page-views
  • sessions::session-abc::state

Atomic Operations

update

Perform atomic updates with multiple operations.
pub async fn update(
    &self,
    key: impl Into<String>,
    ops: Vec<UpdateOp>,
) -> Result<UpdateResult, IIIError>
key
impl Into<String>
required
Stream key in the format stream::group::item
ops
Vec<UpdateOp>
required
List of operations to apply atomically
UpdateResult
Result<UpdateResult, IIIError>
Result containing old and new values
Example:
use iii_sdk::{Streams, UpdateOp};
use serde_json::json;

let streams = Streams::new(iii);

let result = streams.update(
    "orders::user-123::order-456",
    vec![
        UpdateOp::increment("total", 100),
        UpdateOp::set("status", json!("processing")),
        UpdateOp::set("updated_at", json!("2024-01-15T10:30:00Z")),
    ],
).await?;

println!("Old value: {:?}", result.old_value);
println!("New value: {:?}", result.new_value);
All operations in an update call are applied atomically. Either all succeed or none are applied.

Convenience Methods

increment

Atomically increment a numeric field.
pub async fn increment(
    &self,
    key: impl Into<String>,
    field: impl Into<String>,
    by: i64,
) -> Result<UpdateResult, IIIError>
Example:
// Increment page view counter
streams.increment("counters::daily::page-views", "count", 1).await?;

// Increment by larger amount
streams.increment("wallet::user-123::balance", "amount", 1000).await?;

decrement

Atomically decrement a numeric field.
pub async fn decrement(
    &self,
    key: impl Into<String>,
    field: impl Into<String>,
    by: i64,
) -> Result<UpdateResult, IIIError>
Example:
// Decrement inventory
streams.decrement("inventory::warehouse-1::item-abc", "quantity", 5).await?;

set_field

Atomically set a field value.
pub async fn set_field(
    &self,
    key: impl Into<String>,
    field: impl Into<String>,
    value: impl Into<serde_json::Value>,
) -> Result<UpdateResult, IIIError>
Example:
use serde_json::json;

// Set user status
streams.set_field(
    "users::active::user-123",
    "status",
    json!("online")
).await?;

// Set nested field
streams.set_field(
    "profiles::user-123::data",
    "settings.theme",
    json!("dark")
).await?;

remove_field

Atomically remove a field.
pub async fn remove_field(
    &self,
    key: impl Into<String>,
    field: impl Into<String>,
) -> Result<UpdateResult, IIIError>
Example:
// Remove a field
streams.remove_field("cache::session-abc::data", "expired_token").await?;

merge

Atomically merge an object into the existing value.
pub async fn merge(
    &self,
    key: impl Into<String>,
    value: impl Into<serde_json::Value>,
) -> Result<UpdateResult, IIIError>
Example:
use serde_json::json;

// Merge settings
streams.merge(
    "settings::user-123::preferences",
    json!({
        "theme": "dark",
        "language": "en",
        "notifications": true
    })
).await?;

Update Operations

UpdateOp

Operations that can be performed atomically.
pub enum UpdateOp {
    Set { path: FieldPath, value: Option<Value> },
    Merge { path: Option<FieldPath>, value: Value },
    Increment { path: FieldPath, by: i64 },
    Decrement { path: FieldPath, by: i64 },
    Remove { path: FieldPath },
}

Creating UpdateOps

use iii_sdk::UpdateOp;
use serde_json::json;

// Set operation
let op = UpdateOp::set("name", json!("Alice"));

// Increment operation
let op = UpdateOp::increment("counter", 1);

// Decrement operation
let op = UpdateOp::decrement("stock", 5);

// Remove operation
let op = UpdateOp::remove("temp_data");

// Merge operation (root level)
let op = UpdateOp::merge(json!({ "new_field": "value" }));

// Merge operation (at path)
let op = UpdateOp::merge_at("settings", json!({ "theme": "dark" }));

UpdateBuilder

Build complex update operations fluently.
use iii_sdk::UpdateBuilder;
use serde_json::json;

let ops = UpdateBuilder::new()
    .increment("views", 1)
    .set("last_viewed", json!("2024-01-15T10:30:00Z"))
    .merge(json!({
        "metadata": {
            "source": "web"
        }
    }))
    .build();

let result = streams.update("analytics::page-123::stats", ops).await?;

Result Types

UpdateResult

Result of an atomic update operation.
pub struct UpdateResult {
    pub old_value: Option<Value>,
    pub new_value: Value,
}
Example:
let result = streams.increment("counter::daily::visits", "count", 1).await?;

if let Some(old) = result.old_value {
    println!("Previous count: {}", old.get("count").unwrap());
}
println!("New count: {}", result.new_value.get("count").unwrap());

Field Paths

Field paths use dot notation to access nested fields:
// Top-level field
UpdateOp::set("status", json!("active"))

// Nested field
UpdateOp::set("user.profile.name", json!("Alice"))

// Array index
UpdateOp::set("items[0].quantity", json!(5))

Complete Example

Here’s a complete example showing stream operations:
use iii_sdk::{III, Streams, UpdateOp, UpdateBuilder};
use serde_json::json;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let iii = III::new("ws://localhost:49134");
    iii.connect().await?;
    
    let streams = Streams::new(iii);
    
    // Initialize a counter
    streams.set_field(
        "analytics::app::daily-stats",
        "views",
        json!(0)
    ).await?;
    
    // Increment counter atomically
    for _ in 0..10 {
        streams.increment("analytics::app::daily-stats", "views", 1).await?;
    }
    
    // Update multiple fields atomically
    let result = streams.update(
        "analytics::app::daily-stats",
        vec![
            UpdateOp::increment("unique_visitors", 1),
            UpdateOp::set("last_updated", json!("2024-01-15T10:30:00Z")),
            UpdateOp::merge(json!({
                "metadata": {
                    "version": "1.0",
                    "platform": "web"
                }
            })),
        ],
    ).await?;
    
    println!("Updated stats: {:?}", result.new_value);
    
    // Use builder for complex updates
    let ops = UpdateBuilder::new()
        .increment("views", 5)
        .increment("clicks", 2)
        .set("status", json!("active"))
        .build();
    
    streams.update("analytics::app::daily-stats", ops).await?;
    
    Ok(())
}

Use Cases

Counters and Analytics

// Track page views
streams.increment("analytics::page-123::stats", "views", 1).await?;

// Track unique visitors
streams.update(
    "analytics::page-123::stats",
    vec![
        UpdateOp::increment("views", 1),
        UpdateOp::increment("unique_visitors", 1),
    ],
).await?;

User Sessions

// Update session state
streams.merge(
    "sessions::user-123::current",
    json!({
        "last_activity": "2024-01-15T10:30:00Z",
        "page": "/dashboard",
        "online": true
    })
).await?;

Inventory Management

// Decrement stock atomically
let result = streams.decrement(
    "inventory::warehouse-1::item-abc",
    "quantity",
    5
).await?;

if let Some(old) = result.old_value {
    if old.get("quantity").unwrap().as_i64().unwrap() < 5 {
        println!("Warning: Low stock!");
    }
}

See Also

Build docs developers (and LLMs) love