use std::path::Path;
use serde_json::json;
use sqlx::{Pool, Sqlite};
use time::OffsetDateTime;
use time::format_description::FormatItem;
use time::macros::format_description;
use super::error::CliError;
use super::operator_session::OperatorSession;
use super::pds::PdsClient;
use crate::config::{Config, LabelerConfigToml};
use crate::service_record::{self, RECORD_COLLECTION, RECORD_RKEY};
const CTS_FORMAT: &[FormatItem<'_>] =
format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]");
pub const AUDIT_ACTION_SERVICE_RECORD_PUBLISHED: &str = "service_record_published";
pub const AUDIT_REASON_SERVICE_RECORD: &str =
"service_record_published: { cid, content_hash_hex, content_changed }";
#[derive(Debug)]
pub enum PublishOutcome {
NoChange,
Published {
cid: String,
created_at: String,
},
}
pub async fn publish(
pool: &Pool<Sqlite>,
config: &Config,
session_path: &Path,
) -> Result<PublishOutcome, CliError> {
let labeler_cfg = config
.labeler
.as_ref()
.ok_or_else(|| CliError::Config("missing [labeler] section in config".into()))?;
let operator_cfg = config
.operator
.as_ref()
.ok_or_else(|| CliError::Config("missing [operator] section in config".into()))?;
let session = OperatorSession::load(session_path)
.map_err(|e| CliError::Config(format!("operator session: {e}")))?
.ok_or_else(|| {
CliError::Config(format!(
"no operator session at {}; run `cairn operator-login`",
session_path.display()
))
})?;
let prior = load_prior_state(pool).await?;
let created_at = prior
.as_ref()
.map(|p| p.created_at.clone())
.unwrap_or_else(now_rfc3339);
let record = service_record::render(labeler_cfg, &created_at)
.map_err(|e| CliError::Config(format!("render service record: {e}")))?;
let new_hash = service_record::content_hash(&record);
let new_hash_hex = hex::encode(new_hash);
if let Some(p) = &prior
&& p.content_hash_hex == new_hash_hex
{
record_skip_audit(pool, &session.operator_did, p).await?;
return Ok(PublishOutcome::NoChange);
}
let swap = prior.as_ref().map(|p| p.cid.as_str());
let pds = PdsClient::new(&operator_cfg.pds_url)?;
let record_json = serde_json::to_value(&record).expect("record serializes");
let put = pds
.put_record(
&session.access_jwt,
&session.operator_did,
RECORD_COLLECTION,
RECORD_RKEY,
&record_json,
swap,
)
.await?;
commit_publish(
pool,
&session.operator_did,
&put.cid,
&new_hash_hex,
&created_at,
labeler_cfg,
)
.await?;
Ok(PublishOutcome::Published {
cid: put.cid,
created_at,
})
}
#[derive(Debug, Clone)]
struct PriorState {
cid: String,
content_hash_hex: String,
created_at: String,
}
async fn load_prior_state(pool: &Pool<Sqlite>) -> Result<Option<PriorState>, CliError> {
let cid = get_labeler_config(pool, "service_record_cid").await?;
let hash = get_labeler_config(pool, "service_record_content_hash").await?;
let created = get_labeler_config(pool, "service_record_created_at").await?;
match (cid, hash, created) {
(Some(cid), Some(content_hash_hex), Some(created_at)) => Ok(Some(PriorState {
cid,
content_hash_hex,
created_at,
})),
(None, None, None) => Ok(None),
_ => Ok(None),
}
}
async fn get_labeler_config(pool: &Pool<Sqlite>, key: &str) -> Result<Option<String>, CliError> {
sqlx::query_scalar!("SELECT value FROM labeler_config WHERE key = ?1", key)
.fetch_optional(pool)
.await
.map_err(|e| CliError::Startup(format!("read labeler_config: {e}")))
}
async fn commit_publish(
pool: &Pool<Sqlite>,
actor_did: &str,
cid: &str,
hash_hex: &str,
created_at: &str,
_labeler: &LabelerConfigToml,
) -> Result<(), CliError> {
let now_ms = crate::writer::epoch_ms_now();
let reason = audit_reason_json(cid, hash_hex, true);
let mut tx = pool
.begin()
.await
.map_err(|e| CliError::Startup(format!("begin tx: {e}")))?;
for (key, val) in [
("service_record_cid", cid),
("service_record_content_hash", hash_hex),
("service_record_created_at", created_at),
] {
sqlx::query!(
"INSERT INTO labeler_config (key, value, updated_at) VALUES (?1, ?2, ?3)
ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.updated_at",
key,
val,
now_ms,
)
.execute(&mut *tx)
.await
.map_err(|e| CliError::Startup(format!("upsert labeler_config {key}: {e}")))?;
}
crate::audit::append::append_in_tx(
&mut tx,
&crate::audit::append::AuditRowForAppend {
created_at: now_ms,
action: AUDIT_ACTION_SERVICE_RECORD_PUBLISHED.into(),
actor_did: actor_did.into(),
target: None,
target_cid: Some(cid.into()),
outcome: "success".into(),
reason: Some(reason),
},
)
.await
.map_err(|e| CliError::Startup(format!("audit insert: {e}")))?;
tx.commit()
.await
.map_err(|e| CliError::Startup(format!("commit publish: {e}")))?;
Ok(())
}
async fn record_skip_audit(
pool: &Pool<Sqlite>,
actor_did: &str,
prior: &PriorState,
) -> Result<(), CliError> {
let now_ms = crate::writer::epoch_ms_now();
let reason = audit_reason_json(&prior.cid, &prior.content_hash_hex, false);
crate::audit::append::append_via_pool(
pool,
&crate::audit::append::AuditRowForAppend {
created_at: now_ms,
action: AUDIT_ACTION_SERVICE_RECORD_PUBLISHED.into(),
actor_did: actor_did.into(),
target: None,
target_cid: Some(prior.cid.clone()),
outcome: "success".into(),
reason: Some(reason),
},
)
.await
.map_err(|e| CliError::Startup(format!("audit insert: {e}")))?;
Ok(())
}
fn audit_reason_json(cid: &str, hash_hex: &str, content_changed: bool) -> String {
json!({
"cid": cid,
"content_hash_hex": hash_hex,
"content_changed": content_changed,
})
.to_string()
}
fn now_rfc3339() -> String {
let dt = OffsetDateTime::now_utc();
let formatted = dt.format(&CTS_FORMAT).expect("format rfc3339");
format!("{formatted}Z")
}