prolly-store-spanner 0.2.0

Cloud Spanner store adapter for prolly-map.
Documentation
use std::error::Error;
use std::time::{SystemTime, UNIX_EPOCH};

use google_cloud_spanner::client::ClientConfig;
use prolly::{AsyncProlly, Config, Error as ProllyError, Mutation, RemoteProllyStore};
use prolly::{Resolution, Resolver};
use prolly_store_spanner::SpannerBackend;

fn main() -> Result<(), Box<dyn Error>> {
    runtime().block_on(run())
}

async fn run() -> Result<(), Box<dyn Error>> {
    let database = required_env("PROLLY_STORE_SPANNER_DATABASE")?;
    let mut config = ClientConfig::default();
    if std::env::var("PROLLY_STORE_SPANNER_AUTH").is_ok() {
        config = config.with_auth().await?;
    }
    let backend = SpannerBackend::connect(&database, config).await?;
    let client = backend.client().clone();

    let prolly = AsyncProlly::new(RemoteProllyStore::new(backend), Config::default());
    let base = seed_tree(&prolly).await?;
    let left = prolly
        .batch(
            &base,
            vec![
                upsert("account/001", "suspended"),
                upsert("account/003", "active"),
            ],
        )
        .await?;
    let right = prolly
        .batch(&base, vec![upsert("account/002", "closed")])
        .await?;

    let diffs = prolly.diff(&base, &left).await?;
    assert_eq!(diffs.len(), 2);

    let merged = prolly.merge(&base, &left, &right, None).await?;
    assert_value(&prolly, &merged, "account/001", "suspended").await?;
    assert_value(&prolly, &merged, "account/002", "closed").await?;
    assert_value(&prolly, &merged, "account/003", "active").await?;

    let root_name = format!("examples/spanner/{}/main", now_nanos());
    prolly
        .publish_named_root(root_name.as_bytes(), &merged)
        .await?;
    let loaded = prolly
        .load_named_root(root_name.as_bytes())
        .await?
        .expect("named root");
    assert_value(&prolly, &loaded, "account/002", "closed").await?;

    let conflict_left = prolly
        .batch(&base, vec![upsert("account/001", "left")])
        .await?;
    let conflict_right = prolly
        .batch(&base, vec![upsert("account/001", "right")])
        .await?;
    assert!(matches!(
        prolly
            .merge(&base, &conflict_left, &conflict_right, None)
            .await,
        Err(ProllyError::Conflict(_))
    ));

    let resolver: Resolver = Box::new(|conflict| {
        let mut value = conflict.left.clone().unwrap_or_default();
        value.extend_from_slice(b"+");
        value.extend_from_slice(
            conflict
                .right
                .as_ref()
                .map(Vec::as_slice)
                .unwrap_or_default(),
        );
        Resolution::value(value)
    });
    let resolved = prolly
        .merge(&base, &conflict_left, &conflict_right, Some(resolver))
        .await?;
    assert_value(&prolly, &resolved, "account/001", "left+right").await?;

    let roots = prolly.list_named_roots().await?;
    println!("spanner example ok; named_roots={}", roots.len());
    client.close().await;
    Ok(())
}

async fn seed_tree<S>(prolly: &AsyncProlly<S>) -> Result<prolly::Tree, prolly::Error>
where
    S: prolly::AsyncStore,
    S::Error: Send + Sync,
{
    prolly
        .batch(
            &prolly.create(),
            vec![
                upsert("account/001", "active"),
                upsert("account/002", "active"),
            ],
        )
        .await
}

async fn assert_value<S>(
    prolly: &AsyncProlly<S>,
    tree: &prolly::Tree,
    key: &str,
    expected: &str,
) -> Result<(), prolly::Error>
where
    S: prolly::AsyncStore,
    S::Error: Send + Sync,
{
    assert_eq!(
        prolly.get(tree, key.as_bytes()).await?,
        Some(expected.as_bytes().to_vec())
    );
    Ok(())
}

fn upsert(key: &str, value: &str) -> Mutation {
    Mutation::Upsert {
        key: key.as_bytes().to_vec(),
        val: value.as_bytes().to_vec(),
    }
}

fn required_env(name: &str) -> Result<String, Box<dyn Error>> {
    std::env::var(name).map_err(|_| format!("set {name} to run this example").into())
}

fn now_nanos() -> u128 {
    SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .unwrap()
        .as_nanos()
}

fn runtime() -> tokio::runtime::Runtime {
    tokio::runtime::Builder::new_multi_thread()
        .worker_threads(2)
        .enable_all()
        .build()
        .unwrap()
}