use std::time::Duration;
use async_nats::jetstream;
use async_nats::jetstream::object_store::ObjectInfo;
use futures::StreamExt;
use kanade_shared::kv::{
OBJECT_AGENT_RELEASES, OBJECT_APP_PACKAGES, OBJECT_COLLECTIONS, OBJECT_SCRIPTS,
};
use sqlx::SqlitePool;
use tracing::{info, warn};
const REOPEN_BACKOFF: Duration = Duration::from_secs(5);
pub const BUCKETS: [&str; 4] = [
OBJECT_COLLECTIONS,
OBJECT_AGENT_RELEASES,
OBJECT_APP_PACKAGES,
OBJECT_SCRIPTS,
];
pub async fn run(pool: SqlitePool, jetstream: jetstream::Context, bucket: &'static str) {
loop {
let store = match jetstream.get_object_store(bucket).await {
Ok(s) => s,
Err(e) => {
warn!(error = %e, bucket, "object_meta: get_object_store failed; retrying");
tokio::time::sleep(REOPEN_BACKOFF).await;
continue;
}
};
let mut watcher = match store.watch_with_history().await {
Ok(w) => w,
Err(e) => {
warn!(error = %e, bucket, "object_meta: watch failed; retrying");
tokio::time::sleep(REOPEN_BACKOFF).await;
continue;
}
};
info!(bucket, "object_meta projector (re)attached");
while let Some(item) = watcher.next().await {
match item {
Ok(meta) => {
if let Err(e) = apply(&pool, bucket, &meta).await {
warn!(error = %e, bucket, key = %meta.name, "object_meta: apply failed");
}
}
Err(e) => {
warn!(error = %e, bucket, "object_meta watch: entry error; continuing");
}
}
}
warn!(bucket, "object_meta watch ended; reopening");
tokio::time::sleep(REOPEN_BACKOFF).await;
}
}
pub async fn apply(pool: &SqlitePool, bucket: &str, meta: &ObjectInfo) -> Result<(), sqlx::Error> {
if meta.deleted {
return delete_key(pool, bucket, &meta.name).await;
}
let modified = meta.modified.and_then(|t| {
chrono::DateTime::from_timestamp(t.unix_timestamp(), t.nanosecond()).map(|d| d.to_rfc3339())
});
sqlx::query(
"INSERT INTO object_store_meta (bucket, key, size, digest, modified)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT (bucket, key) DO UPDATE SET
size = excluded.size,
digest = excluded.digest,
modified = excluded.modified",
)
.bind(bucket)
.bind(&meta.name)
.bind(meta.size as i64)
.bind(&meta.digest)
.bind(modified)
.execute(pool)
.await?;
Ok(())
}
pub async fn delete_key(pool: &SqlitePool, bucket: &str, key: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM object_store_meta WHERE bucket = ? AND key = ?")
.bind(bucket)
.bind(key)
.execute(pool)
.await?;
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MetaRow {
pub key: String,
pub size: i64,
pub digest: Option<String>,
pub modified: Option<String>,
}
pub async fn list_bucket(pool: &SqlitePool, bucket: &str) -> Result<Vec<MetaRow>, sqlx::Error> {
let rows = sqlx::query_as::<_, (String, i64, Option<String>, Option<String>)>(
"SELECT key, size, digest, modified FROM object_store_meta WHERE bucket = ?",
)
.bind(bucket)
.fetch_all(pool)
.await?;
Ok(rows
.into_iter()
.map(|(key, size, digest, modified)| MetaRow {
key,
size,
digest,
modified,
})
.collect())
}
#[cfg(test)]
mod tests {
use super::*;
use sqlx::sqlite::SqlitePoolOptions;
async fn pool() -> SqlitePool {
let pool = SqlitePoolOptions::new()
.max_connections(1)
.connect("sqlite::memory:")
.await
.unwrap();
sqlx::migrate!("./migrations").run(&pool).await.unwrap();
pool
}
fn meta(name: &str, size: usize, deleted: bool) -> ObjectInfo {
ObjectInfo {
name: name.to_string(),
description: None,
metadata: Default::default(),
headers: None,
options: None,
bucket: OBJECT_SCRIPTS.to_string(),
nuid: format!("nuid-{name}"),
size,
chunks: 1,
modified: None,
digest: Some("SHA-256=abc".to_string()),
deleted,
}
}
#[tokio::test]
async fn put_then_list_returns_the_object() {
let pool = pool().await;
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 42, false))
.await
.unwrap();
let rows = list_bucket(&pool, OBJECT_SCRIPTS).await.unwrap();
assert_eq!(
rows,
vec![MetaRow {
key: "a/1.0.0".to_string(),
size: 42,
digest: Some("SHA-256=abc".to_string()),
modified: None,
}],
);
}
#[tokio::test]
async fn re_put_overwrites_in_place() {
let pool = pool().await;
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 42, false))
.await
.unwrap();
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 99, false))
.await
.unwrap();
let rows = list_bucket(&pool, OBJECT_SCRIPTS).await.unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].size, 99);
}
#[tokio::test]
async fn tombstone_removes_the_row() {
let pool = pool().await;
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 42, false))
.await
.unwrap();
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 0, true))
.await
.unwrap();
assert!(list_bucket(&pool, OBJECT_SCRIPTS).await.unwrap().is_empty());
}
#[tokio::test]
async fn buckets_are_isolated() {
let pool = pool().await;
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 1, false))
.await
.unwrap();
apply(&pool, OBJECT_APP_PACKAGES, &meta("a/1.0.0", 2, false))
.await
.unwrap();
assert_eq!(list_bucket(&pool, OBJECT_SCRIPTS).await.unwrap()[0].size, 1);
assert_eq!(
list_bucket(&pool, OBJECT_APP_PACKAGES).await.unwrap()[0].size,
2,
);
apply(&pool, OBJECT_SCRIPTS, &meta("a/1.0.0", 0, true))
.await
.unwrap();
assert!(list_bucket(&pool, OBJECT_SCRIPTS).await.unwrap().is_empty());
assert_eq!(
list_bucket(&pool, OBJECT_APP_PACKAGES).await.unwrap().len(),
1
);
}
}