| name | sdk-patterns-rust |
| summary | Couchbase SDK patterns for Rust — CAS optimistic locking, retry on CasMismatch, tokio bulk operations with join_all/FuturesUnordered, atomic counters, sub-document (MutateInSpec/LookupInSpec), array operations, exists, replica reads, touch/getAndTouch, preserve_expiry, error handling |
| description | Couchbase SDK patterns for Rust — CAS optimistic locking, retry on CasMismatch, tokio bulk operations with join_all/FuturesUnordered, atomic counters, sub-document (MutateInSpec/LookupInSpec), array operations, exists, replica reads, touch/getAndTouch, preserve_expiry, error handling |
| compatibility | Rust SDK 1.0+. Requires tokio async runtime. No KV range scan or transactions yet. |
| metadata | {"last_verified":"2026-05","min_server_version":"7.0","handoff":[{"condition":"user asks about SQL++ queries","skill":"server-querying-rust"},{"condition":"user asks about full-text search or vector search","skill":"search-rust"},{"condition":"user asks about error handling, retries, or exception types","skill":"error-handling"}]} |
SDK Patterns — Rust
CAS Optimistic Locking
Read-modify-write without pessimistic locks. Retry on CasMismatch.
use couchbase::error::ErrorKind;
use couchbase::options::kv_options::ReplaceOptions;
use std::future::Future;
async fn update_with_cas(
collection: &couchbase::collection::Collection,
key: &str,
) -> Result<(), couchbase::error::Error> {
loop {
let get_result = collection.get(key, None).await?;
let mut doc: serde_json::Value = get_result.content_as()?;
doc["updated_at"] = serde_json::json!(chrono::Utc::now().timestamp());
match collection.replace(
key,
doc,
ReplaceOptions::new().cas(get_result.cas()),
).await {
Ok(_) => return Ok(()),
Err(e) if matches!(e.kind(), ErrorKind::CasMismatch) => continue,
Err(e) => return Err(e),
}
}
}
Bulk Operations
Use tokio::join! for a fixed set, or futures::future::join_all / FuturesUnordered for dynamic sets.
use futures::future::join_all;
let keys = vec!["airline::1", "airline::2", "airline::3"];
let futures: Vec<_> = keys.iter()
.map(|k| collection.get(*k, None))
.collect();
let results = join_all(futures).await;
for (key, result) in keys.iter().zip(results) {
match result {
Ok(r) => println!("{key}: {}", r.content_as::<serde_json::Value>().unwrap()),
Err(e) if matches!(e.kind(), ErrorKind::DocumentNotFound) => println!("{key}: not found"),
Err(e) => eprintln!("{key}: error {e}"),
}
}
Bulk Upsert with concurrency limit
use futures::stream::{self, StreamExt};
let docs: Vec<(&str, serde_json::Value)> = vec![
("airline::10", serde_json::json!({"name": "Alpha Air"})),
("airline::11", serde_json::json!({"name": "Beta Air"})),
];
stream::iter(docs)
.map(|(key, doc)| collection.upsert(key, doc, None))
.buffer_unordered(16)
.for_each(|result| async move {
if let Err(e) = result {
eprintln!("Upsert error: {e}");
}
})
.await;
Atomic Counters
use couchbase::options::kv_options::{IncrementOptions, DecrementOptions};
use couchbase::kv::CounterDelta;
let result = collection.increment(
"counter::page_views",
IncrementOptions::new()
.delta(CounterDelta::from(1u64))
.initial(Some(0)),
).await?;
println!("New value: {}", result.content());
let result = collection.decrement(
"counter::stock::item42",
DecrementOptions::new().delta(CounterDelta::from(1u64)),
).await?;
Sub-Document Operations
Fetch or mutate specific fields without transferring the full document.
use couchbase::kv::{LookupInSpec, MutateInSpec};
let result = collection.lookup_in("order::1001", &[
LookupInSpec::get("status", None)?,
LookupInSpec::get("items[0].sku", None)?,
LookupInSpec::exists("metadata.archived", None)?,
LookupInSpec::count("items", None)?,
], None).await?;
let status: String = result.content_as(0)?;
let first_sku: String = result.content_as(1)?;
let archived = result.exists(2);
let item_count: u32 = result.content_as(3)?;
collection.mutate_in("order::1001", &[
MutateInSpec::upsert("status", "shipped", None)?,
MutateInSpec::upsert("tracking.carrier", "FedEx", None)?,
MutateInSpec::increment("metrics.view_count", 1, None)?,
], None).await?;
Array Operations
collection.mutate_in("order::1001", &[
MutateInSpec::array_append("events",
serde_json::json!({"ts": 1700000000, "action": "shipped"}), None)?,
], None).await?;
collection.mutate_in("order::1001", &[
MutateInSpec::array_prepend("events",
serde_json::json!({"ts": 1700000001, "action": "delivered"}), None)?,
], None).await?;
collection.mutate_in("user::42", &[
MutateInSpec::array_add_unique("tags", "premium", None)?,
], None).await?;
Full sub-document API reference (all languages): shared/server/subdocument.md
Replica Reads
let result = collection.get_any_replica("airline::1001", None).await?;
let doc: serde_json::Value = result.content_as()?;
Touch (Reset TTL)
use couchbase::options::kv_options::TouchOptions;
use std::time::Duration;
collection.touch("session::abc", Duration::from_secs(3600), None).await?;
let result = collection.get_and_touch(
"session::abc",
Duration::from_secs(3600),
None,
).await?;
let doc: serde_json::Value = result.content_as()?;
Preserve Expiry on Mutate
use couchbase::options::kv_options::ReplaceOptions;
let get_result = collection.get("session::abc", None).await?;
let mut doc: serde_json::Value = get_result.content_as()?;
doc["last_seen"] = serde_json::json!(chrono::Utc::now().timestamp());
collection.replace(
"session::abc",
doc,
ReplaceOptions::new()
.cas(get_result.cas())
.preserve_expiry(true),
).await?;
Error Handling Patterns
use couchbase::error::ErrorKind;
fn handle_kv_error(e: couchbase::error::Error) {
match e.kind() {
ErrorKind::DocumentNotFound => eprintln!("Document does not exist"),
ErrorKind::DocumentExists => eprintln!("Document already exists (insert conflict)"),
ErrorKind::CasMismatch => eprintln!("Concurrent modification, retry"),
ErrorKind::ServerTimeout => eprintln!("Operation timed out"),
ErrorKind::ValueTooLarge => eprintln!("Document exceeds 20MB limit"),
_ => eprintln!("Unexpected error: {e}"),
}
}