use std::collections::HashMap;
use affinidi_tdk::secrets_resolver::secrets::Secret;
use chrono::Utc;
use didwebvh_rs::log_entry::LogEntryMethods;
use serde_json::Value;
use vta_sdk::keys::{KeyOrigin, KeyRecord, KeyStatus, KeyType};
use super::errors::UpdateDidWebvhError;
use super::options::{RotateDidWebvhKeysOptions, UpdateDidWebvhOptions, UpdateDidWebvhResult};
use super::orchestrator::{Commit, OnCommit, update_did_webvh_tracked};
use super::state::{find_record_by_scid, state_from_jsonl};
use crate::auth::AuthClaims;
use crate::error::AppError;
use crate::keys::derivation::Bip32Extension;
use crate::keys::seeds::{get_active_seed_id, load_seed_bytes};
use crate::keys::{encode_public_multibase, store_key};
use crate::operations::did_webvh::concurrency::DID_UPDATE_LOCKS;
use crate::store::KeyspaceHandle;
use crate::webvh_store;
const RELATIONSHIPS: [&str; 5] = [
"authentication",
"assertionMethod",
"keyAgreement",
"capabilityInvocation",
"capabilityDelegation",
];
#[derive(Clone)]
struct Rotated {
vm_id: String,
key_type: KeyType,
path: String,
public_key: String,
previous: Option<KeyRecord>,
}
pub async fn rotate_did_webvh_keys(
deps: &super::super::WebvhDeps<'_>,
auth: &AuthClaims,
scid: &str,
opts: RotateDidWebvhKeysOptions,
vta_did: Option<&str>,
identity_secrets: Option<&affinidi_tdk::secrets_resolver::ThreadedSecretsResolver>,
channel: &str,
) -> Result<UpdateDidWebvhResult, UpdateDidWebvhError> {
let record = find_record_by_scid(deps.webvh_ks, scid)
.await?
.ok_or_else(|| UpdateDidWebvhError::NotFound(format!("SCID {scid} not found")))?;
auth.require_admin()
.map_err(|e| UpdateDidWebvhError::Forbidden(format!("admin required: {e}")))?;
auth.require_context(&record.context_id).map_err(|_| {
UpdateDidWebvhError::Forbidden(format!(
"caller has no admin role in context `{}`",
record.context_id
))
})?;
let did_log = webvh_store::get_did_log(deps.webvh_ks, &record.did)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("get_did_log: {e}")))?
.ok_or_else(|| {
UpdateDidWebvhError::Library(format!("DID log missing for {}", record.did))
})?;
let state = state_from_jsonl(&did_log)?;
let last = state.log_entries().last().ok_or_else(|| {
UpdateDidWebvhError::Library(format!("DID {} has no log entries", record.did))
})?;
let prior_version_id = last.get_version_id().to_string();
let current_doc = last.log_entry.get_did_document().map_err(|e| {
UpdateDidWebvhError::Library(format!("extract document from last entry: {e}"))
})?;
let context = crate::contexts::get_context(deps.contexts_ks, &record.context_id)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("get_context: {e}")))?
.ok_or_else(|| {
UpdateDidWebvhError::Library(format!(
"context `{}` referenced by DID is missing",
record.context_id
))
})?;
let mut new_doc = current_doc.clone();
let seed_id = get_active_seed_id(deps.keys_ks).await.map_err(|e| {
UpdateDidWebvhError::Persistence(format!("could not load active seed id: {e}"))
})?;
let seed = load_seed_bytes(deps.keys_ks, deps.seed_store, Some(seed_id))
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("could not load seed: {e}")))?;
let root = vti_common::slip10::ExtendedSigningKey::from_seed(&seed)
.map_err(|e| UpdateDidWebvhError::Persistence(format!("BIP-32 root derivation: {e}")))?;
let mut rotated: Vec<Rotated> = Vec::new();
{
let obj = new_doc.as_object_mut().ok_or_else(|| {
UpdateDidWebvhError::Library("current document is not a JSON object".into())
})?;
let mut methods: Vec<&mut Value> = Vec::new();
for (field, value) in obj.iter_mut() {
let Value::Array(entries) = value else {
continue;
};
if field == "verificationMethod" {
methods.extend(entries.iter_mut());
} else if RELATIONSHIPS.contains(&field.as_str()) {
methods.extend(entries.iter_mut().filter(|e| e.is_object()));
}
}
if methods.is_empty() {
return Err(UpdateDidWebvhError::Library(
"current document declares no verification methods to rotate".into(),
));
}
for (i, vm) in methods.into_iter().enumerate() {
let Some(entry) = rotate_one(
deps.keys_ks,
&root,
&context.base_path,
&record.did,
&record.context_id,
i,
vm,
)
.await?
else {
continue;
};
if rotated.iter().any(|r| r.vm_id == entry.vm_id) {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"verification method `{}` is declared more than once",
entry.vm_id
)));
}
rotated.push(entry);
}
}
if rotated.is_empty() {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"{} declares no verification method whose key this VTA holds; there is nothing \
it can rotate",
record.did
)));
}
let rotation_id = uuid::Uuid::new_v4();
for r in &rotated {
let id = staging_id(&r.vm_id, &rotation_id);
deps.keys_ks
.insert(
store_key(&id),
&new_record(r, &id, &record.context_id, seed_id, false),
)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("stage rotated key: {e}")))?;
}
let own_identity = if vta_did == Some(record.did.as_str()) {
match identity_secrets {
Some(resolver) => Some((resolver.clone(), own_identity_secrets(&root, &rotated))),
None => {
tracing::warn!(
did = %record.did,
"rotating this VTA's own DID on a path with no live secrets resolver; \
restart the VTA to load the new keys"
);
None
}
}
} else {
None
};
let on_commit: OnCommit = {
let keys_ks = deps.keys_ks.clone();
let rotated = rotated.clone();
let prior_version_id = prior_version_id.clone();
let context_id = record.context_id.clone();
Box::new(move || {
Box::pin(async move {
promote(
&keys_ks,
&rotated,
&prior_version_id,
&rotation_id,
&context_id,
seed_id,
)
.await?;
if let Some((resolver, secrets)) = own_identity {
install_own_identity(&resolver, secrets).await;
}
Ok(())
})
})
};
let label = opts
.label
.or_else(|| Some(format!("rotate-keys for {}", record.did)));
let result = update_did_webvh_tracked(
deps,
auth,
scid,
UpdateDidWebvhOptions {
document: Some(new_doc),
pre_rotation_count: opts.pre_rotation_count,
witnesses: None,
watchers: None,
ttl: None,
label,
expected_version_id: Some(prior_version_id.clone()),
},
vta_did,
channel,
Some(on_commit),
)
.await;
let result = match result {
Ok(result) => result,
Err((e, Commit::NotCommitted)) => {
for r in &rotated {
let _ = deps
.keys_ks
.remove(store_key(&staging_id(&r.vm_id, &rotation_id)))
.await;
}
return Err(e);
}
Err((e, Commit::Committed)) => {
tracing::warn!(
did = %record.did, error = %e,
"rotation committed locally but a later step failed"
);
return Err(e);
}
};
crate::audit::record_best_effort(
deps.audit,
"did.webvh.rotate_keys",
&auth.did,
Some(&record.did),
"success",
Some(channel),
Some(&record.context_id),
)
.await;
tracing::info!(
channel,
did = %record.did,
scid = %scid,
methods = rotated.len(),
version = %result.new_version_id,
"did:webvh keys rotated"
);
Ok(result)
}
fn own_identity_secrets(
root: &vti_common::slip10::ExtendedSigningKey,
rotated: &[Rotated],
) -> Vec<Secret> {
let mut secrets = Vec::new();
for r in rotated {
let secret = match r.key_type {
KeyType::Ed25519 => root.derive_ed25519(&r.path),
KeyType::X25519 => root.derive_x25519(&r.path),
KeyType::MlDsa44 => root.derive_ml_dsa_44(&r.path),
KeyType::MlDsa65 => root.derive_ml_dsa_65(&r.path),
_ => {
tracing::warn!(
vm_id = %r.vm_id,
"rotated a {:?} key on this VTA's own DID; it is not held by the \
messaging resolver and is loaded on its next use",
r.key_type
);
continue;
}
};
match secret {
Ok(mut secret) => {
secret.id = r.vm_id.clone();
secrets.push(secret);
}
Err(e) => tracing::error!(
vm_id = %r.vm_id, error = %e,
"could not derive the rotated key for the live resolver; restart the VTA"
),
}
}
secrets
}
async fn install_own_identity(
resolver: &affinidi_tdk::secrets_resolver::ThreadedSecretsResolver,
secrets: Vec<Secret>,
) {
use affinidi_tdk::secrets_resolver::SecretsResolver;
for secret in secrets {
let vm_id = secret.id.clone();
resolver.insert(secret).await;
tracing::info!(vm_id = %vm_id, "live secret replaced after rotation");
}
}
fn staging_id(vm_id: &str, rotation_id: &uuid::Uuid) -> String {
format!("{vm_id}{STAGING_MARK}{rotation_id}")
}
const STAGING_MARK: &str = "@rotating-";
#[derive(Debug, Default, PartialEq, Eq)]
pub struct StagedRotationRecovery {
pub promoted: Vec<String>,
pub retired: Vec<String>,
pub removed: Vec<String>,
}
pub async fn recover_staged_rotations(
keys_ks: &KeyspaceHandle,
webvh_ks: &KeyspaceHandle,
) -> Result<StagedRotationRecovery, AppError> {
let mut report = StagedRotationRecovery::default();
let mut logs: HashMap<String, Option<Vec<(String, Value)>>> = HashMap::new();
for (_, value) in keys_ks.prefix_iter_raw("key:").await? {
let Ok(staged) = serde_json::from_slice::<KeyRecord>(&value) else {
continue;
};
let Some((vm_id, _)) = staged.key_id.split_once(STAGING_MARK) else {
continue;
};
let vm_id = vm_id.to_string();
let Some((did, _)) = vm_id.split_once('#') else {
keys_ks.remove(store_key(&staged.key_id)).await?;
report.removed.push(staged.key_id.clone());
continue;
};
let _guard = DID_UPDATE_LOCKS.acquire(did).await;
if !logs.contains_key(did) {
let entries = match webvh_store::get_did_log(webvh_ks, did).await? {
None => None,
Some(jsonl) => match state_from_jsonl(&jsonl) {
Ok(state) => Some(
state
.log_entries()
.iter()
.filter_map(|e| {
Some((
e.get_version_id().to_string(),
e.log_entry.get_did_document().ok()?,
))
})
.collect(),
),
Err(e) => {
tracing::error!(
did, error = %e,
"cannot read the DID log to recover an interrupted key rotation; \
leaving its staging records"
);
continue;
}
},
};
logs.insert(did.to_string(), entries);
}
let history: Vec<(&str, Option<String>)> = logs[did]
.iter()
.flatten()
.map(|(version, doc)| (version.as_str(), published_key(doc, did, &vm_id)))
.collect();
let publishes = |k: &Option<String>| k.as_deref() == Some(staged.public_key.as_str());
if history.last().is_some_and(|(_, k)| publishes(k)) {
let active: Option<KeyRecord> = keys_ks.get(store_key(&vm_id)).await?;
let already = active
.as_ref()
.is_some_and(|a| a.public_key == staged.public_key);
if !already {
if let Some(active) = &active {
if active.context_id != staged.context_id {
tracing::error!(
vm_id = %vm_id,
"the active record of a method an interrupted rotation staged a \
key for belongs to another context; leaving it for an operator"
);
continue;
}
let first = history
.iter()
.rposition(|(_, k)| !publishes(k))
.map(|i| i + 1)
.unwrap_or(0);
if let Some(prior) = first.checked_sub(1).map(|i| history[i].0) {
let retired_id = format!("{vm_id}@{prior}");
let mut retired = active.clone();
retired.key_id = retired_id.clone();
retired.status = KeyStatus::Revoked;
retired.updated_at = Utc::now();
keys_ks.insert(store_key(&retired_id), &retired).await?;
}
}
let mut promoted = staged.clone();
promoted.key_id = vm_id.clone();
promoted.status = KeyStatus::Active;
promoted.updated_at = Utc::now();
keys_ks.insert(store_key(&vm_id), &promoted).await?;
}
keys_ks.remove(store_key(&staged.key_id)).await?;
tracing::warn!(vm_id = %vm_id, "promoted the key of an interrupted rotation");
report.promoted.push(vm_id);
} else if let Some(last) = history.iter().rposition(|(_, k)| publishes(k)) {
let retired_id = format!("{vm_id}@{}", history[last].0);
if keys_ks
.get::<KeyRecord>(store_key(&retired_id))
.await?
.is_none()
{
let mut retired = staged.clone();
retired.key_id = retired_id.clone();
retired.status = KeyStatus::Revoked;
retired.updated_at = Utc::now();
keys_ks.insert(store_key(&retired_id), &retired).await?;
}
keys_ks.remove(store_key(&staged.key_id)).await?;
report.retired.push(retired_id);
} else {
keys_ks.remove(store_key(&staged.key_id)).await?;
report.removed.push(staged.key_id.clone());
}
}
Ok(report)
}
fn published_key(doc: &Value, did: &str, vm_id: &str) -> Option<String> {
let obj = doc.as_object()?;
std::iter::once("verificationMethod")
.chain(RELATIONSHIPS)
.filter_map(|field| obj.get(field)?.as_array())
.flatten()
.filter_map(Value::as_object)
.find(|m| {
m.get("id").and_then(Value::as_str).is_some_and(|id| {
id == vm_id || (id.starts_with('#') && format!("{did}{id}") == vm_id)
})
})
.and_then(|m| m.get("publicKeyMultibase")?.as_str().map(str::to_string))
}
async fn promote(
keys_ks: &KeyspaceHandle,
rotated: &[Rotated],
prior_version_id: &str,
rotation_id: &uuid::Uuid,
context_id: &str,
seed_id: u32,
) -> Result<(), UpdateDidWebvhError> {
let now = Utc::now();
for r in rotated {
if let Some(previous) = &r.previous {
let retired_id = format!("{}@{prior_version_id}", r.vm_id);
let mut retired = previous.clone();
retired.key_id = retired_id.clone();
retired.status = KeyStatus::Revoked;
retired.updated_at = now;
keys_ks
.insert(store_key(&retired_id), &retired)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("retire key: {e}")))?;
}
keys_ks
.insert(
store_key(&r.vm_id),
&new_record(r, &r.vm_id, context_id, seed_id, true),
)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("install rotated key: {e}")))?;
let _ = keys_ks
.remove(store_key(&staging_id(&r.vm_id, rotation_id)))
.await;
}
Ok(())
}
fn new_record(
r: &Rotated,
key_id: &str,
context_id: &str,
seed_id: u32,
active: bool,
) -> KeyRecord {
let now = Utc::now();
KeyRecord {
key_id: key_id.to_string(),
derivation_path: r.path.clone(),
key_type: r.key_type.clone(),
status: if active {
KeyStatus::Active
} else {
KeyStatus::Revoked
},
public_key: r.public_key.clone(),
label: Some(
r.previous
.as_ref()
.and_then(|p| p.label.clone())
.unwrap_or_else(|| r.vm_id.clone()),
),
context_id: Some(context_id.to_string()),
seed_id: Some(seed_id),
exportable: r.previous.as_ref().and_then(|p| p.exportable),
origin: KeyOrigin::Derived,
created_at: now,
updated_at: now,
}
}
async fn rotate_one(
keys_ks: &KeyspaceHandle,
root: &vti_common::slip10::ExtendedSigningKey,
base_path: &str,
did: &str,
context_id: &str,
index: usize,
vm: &mut Value,
) -> Result<Option<Rotated>, UpdateDidWebvhError> {
let obj = vm.as_object_mut().ok_or_else(|| {
UpdateDidWebvhError::InvalidDocument(format!(
"verification method {index} is not an object"
))
})?;
let raw_id = obj
.get("id")
.and_then(Value::as_str)
.ok_or_else(|| {
UpdateDidWebvhError::InvalidDocument(format!("verification method {index} has no id"))
})?
.to_string();
let vm_id = if raw_id.starts_with('#') {
format!("{did}{raw_id}")
} else {
raw_id
};
if !vm_id.starts_with(&format!("{did}#")) {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"verification method `{vm_id}` is not a method of {did}; rotation replaces only \
this DID's own keys"
)));
}
let Some(current_public) = obj
.get("publicKeyMultibase")
.and_then(Value::as_str)
.map(str::to_string)
else {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"verification method `{vm_id}` carries no publicKeyMultibase; rotation replaces \
multibase keys only"
)));
};
let previous: Option<KeyRecord> = keys_ks
.get(store_key(&vm_id))
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("load key record: {e}")))?;
let Some(previous) = previous else {
tracing::info!(vm_id = %vm_id, "rotation skipped a method whose key this VTA does not hold");
return Ok(None);
};
if previous.key_id != vm_id {
return Err(UpdateDidWebvhError::Forbidden(format!(
"the key record stored for `{vm_id}` names another key (`{}`); refusing to \
rotate it",
previous.key_id
)));
}
if previous.context_id.as_deref() != Some(context_id) {
return Err(UpdateDidWebvhError::Forbidden(format!(
"the key record for `{vm_id}` does not belong to this DID's context; refusing to \
rotate it"
)));
}
if previous.public_key != current_public {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"the key record for `{vm_id}` holds a different key than the document publishes; \
realign the records (webvh/dids/realign-keys) before rotating"
)));
}
if previous.origin == KeyOrigin::Internal {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"verification method `{vm_id}` is backed by an internal (non-extractable) key; \
rotating it here would replace it with a key the seed can reproduce, which is \
weaker. Mint a new internal key and publish it with vta/webvh/dids/update instead"
)));
}
let key_type = previous.key_type.clone();
let path = crate::keys::paths::allocate_path(keys_ks, base_path)
.await
.map_err(|e| UpdateDidWebvhError::Persistence(format!("allocate_path: {e}")))?;
let public_key = derive_public(root, &key_type, &path)?;
obj.insert(
"publicKeyMultibase".into(),
Value::String(public_key.clone()),
);
Ok(Some(Rotated {
vm_id,
key_type,
path,
public_key,
previous: Some(previous),
}))
}
fn derive_public(
root: &vti_common::slip10::ExtendedSigningKey,
key_type: &KeyType,
path: &str,
) -> Result<String, UpdateDidWebvhError> {
let err = |e: crate::error::AppError| {
UpdateDidWebvhError::Persistence(format!("derive {key_type:?} at `{path}`: {e}"))
};
let secret = match key_type {
KeyType::Ed25519 => root.derive_ed25519(path).map_err(err)?,
KeyType::X25519 => root.derive_x25519(path).map_err(err)?,
KeyType::MlDsa44 => root.derive_ml_dsa_44(path).map_err(err)?,
KeyType::MlDsa65 => root.derive_ml_dsa_65(path).map_err(err)?,
KeyType::P256 => {
use p256::elliptic_curve::sec1::ToSec1Point;
let secret = root.derive_p256(path).map_err(err)?;
let point = secret.secret_key.public_key().to_sec1_point(true);
return Ok(encode_public_multibase(&KeyType::P256, point.as_bytes()));
}
#[allow(unreachable_patterns)]
other => {
return Err(UpdateDidWebvhError::InvalidDocument(format!(
"rotation cannot mint a {other:?} key"
)));
}
};
secret
.get_public_keymultibase()
.map_err(|e| UpdateDidWebvhError::Persistence(format!("public key encoding: {e}")))
}