use serde::{Deserialize, Serialize};
use serde_json::Value;
use trust_tasks_rs::TrustTask;
use vti_common::acl::Capability;
use crate::audit;
use crate::auth::AuthClaims;
use crate::operations::{room_groups, room_invitation};
use crate::server::AppState;
use super::helpers::{
TRANSPORT_TRUST_TASK, TrustTaskOutcome, app_error_to_reject, parse_payload, success_response,
};
const KEY_PACKAGE_LIFETIME_SECS: u64 = 7 * 24 * 60 * 60;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct KeyPackageResponse {
pub key_package: String,
pub expires_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct WelcomeResponse {
pub epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct CommitResponse {
pub epoch: u64,
}
pub(super) async fn handle_key_package(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: trust_tasks_rs::specs::rooms::keys::key_package::v0_1::Payload =
match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let Err(e) =
require_invitation(state, req.invitation.as_deref(), &req.room_id, &auth.did).await
{
return app_error_to_reject(&doc, e);
}
let minted = match room_groups::mint_key_package(
&state.room_groups_ks,
&req.room_id,
&auth.did,
KEY_PACKAGE_LIFETIME_SECS,
now(),
)
.await
{
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.key-package", auth, &req.room_id).await;
success_response(
&doc,
KeyPackageResponse {
key_package: minted.key_package,
expires_at: rfc3339(minted.expires_at),
},
)
}
pub(super) async fn handle_welcome(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: trust_tasks_rs::specs::rooms::keys::welcome::v0_1::Payload = match parse_payload(&doc)
{
Ok(r) => r,
Err(resp) => return resp,
};
let invitation =
match require_invitation(state, req.invitation.as_deref(), &req.room_id, &auth.did).await {
Ok(i) => i,
Err(e) => return app_error_to_reject(&doc, e),
};
let welcome = match decode_b64(&req.welcome, "welcome") {
Ok(b) => b,
Err(e) => return app_error_to_reject(&doc, e),
};
let epoch = match room_groups::join(
&state.room_groups_ks,
&req.room_id,
invitation.subject(),
&welcome,
now(),
)
.await
{
Ok(e) => e,
Err(e) => return app_error_to_reject(&doc, e),
};
if let Err(e) = room_groups::consume_invitation(
&state.room_invitations_ks,
invitation.credential_id(),
&req.room_id,
now(),
)
.await
{
return app_error_to_reject(&doc, e);
}
record(state, "rooms.keys.welcome", auth, &req.room_id).await;
success_response(&doc, WelcomeResponse { epoch })
}
pub(super) async fn handle_commit(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: trust_tasks_rs::specs::rooms::keys::commit::v0_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let commit = match decode_b64(&req.commit, "commit") {
Ok(b) => b,
Err(e) => return app_error_to_reject(&doc, e),
};
let epoch = match room_groups::apply_commit(
&state.room_groups_ks,
&req.room_id,
&commit,
u64::from(req.epoch),
now(),
)
.await
{
Ok(e) => e,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.commit", auth, &req.room_id).await;
success_response(&doc, CommitResponse { epoch })
}
pub(super) async fn handle_seal(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(r) = super::helpers::require_capability(
state,
auth,
&doc,
Capability::RoomOpen,
"sealing a room record",
)
.await
{
return r;
}
let req: trust_tasks_rs::specs::rooms::keys::seal::v0_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let plaintext = match base64::Engine::decode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
&req.plaintext,
) {
Ok(b) => b,
Err(e) => {
return app_error_to_reject(
&doc,
vti_common::error::AppError::Validation(format!("plaintext is not base64url: {e}")),
);
}
};
let sealed = match room_groups::seal_record(
&state.room_groups_ks,
&req.room_id,
&req.key,
req.version,
&plaintext,
)
.await
{
Ok(s) => s,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.seal", auth, &req.room_id).await;
success_response(
&doc,
serde_json::json!({
"sealed": {
"ciphertext": sealed.ciphertext,
"nonce": sealed.nonce,
"epoch": sealed.epoch,
}
}),
)
}
pub(super) async fn handle_list(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(r) = super::helpers::require_capability(
state,
auth,
&doc,
Capability::RoomOpen,
"listing the rooms this VTA holds keys for",
)
.await
{
return r;
}
let rooms = match room_groups::list_rooms(&state.room_groups_ks).await {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.list", auth, "").await;
success_response(&doc, serde_json::json!({ "rooms": rooms }))
}
pub(super) async fn handle_chain(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(r) = super::helpers::require_capability(
state,
auth,
&doc,
Capability::RoomOpen,
"extending a room's readable history",
)
.await
{
return r;
}
let req: trust_tasks_rs::specs::rooms::keys::chain::v0_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let links: Result<Vec<vti_rooms::wire::EpochLink>, _> = req
.links
.iter()
.map(|l| {
u32::try_from(l.epoch).map(|epoch| vti_rooms::wire::EpochLink {
epoch,
wrapped: l.wrapped.clone(),
nonce: l.nonce.clone(),
})
})
.collect();
let links = match links {
Ok(l) => l,
Err(_) => {
return app_error_to_reject(
&doc,
vti_common::error::AppError::Validation(
"an epoch link names an epoch outside the representable range".into(),
),
);
}
};
let (earliest, stored) =
match room_groups::store_links(&state.room_groups_ks, &req.room_id, links, now()).await {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.chain", auth, &req.room_id).await;
success_response(
&doc,
serde_json::json!({
"roomId": req.room_id,
"earliestReadableEpoch": earliest,
"stored": stored,
}),
)
}
fn signing_context<'a>(
state: &'a AppState,
auth: &'a AuthClaims,
) -> crate::operations::room_issuance::SigningContext<'a> {
crate::operations::room_issuance::SigningContext {
keys_ks: &state.keys_ks,
imported_ks: &state.imported_ks,
internal_ks: &state.internal_ks,
contexts_ks: &state.contexts_ks,
acl_ks: &state.acl_ks,
seed_store: &state.seed_store,
auth,
}
}
pub(super) async fn handle_backfill(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(r) = super::helpers::require_capability(
state,
auth,
&doc,
Capability::RoomOpen,
"fetching a room's readable history",
)
.await
{
return r;
}
let req: trust_tasks_rs::specs::rooms::keys::backfill::v0_1::Payload = match parse_payload(&doc)
{
Ok(r) => r,
Err(resp) => return resp,
};
let vta_did = match state.config.read().await.vta_did.clone() {
Some(d) => d,
None => {
return app_error_to_reject(
&doc,
vti_common::error::AppError::Validation(
"this agent has no DID of its own, so it cannot present to a host as itself"
.into(),
),
);
}
};
let Some(resolver) = state.did_resolver.clone() else {
return app_error_to_reject(
&doc,
vti_common::error::AppError::Validation(
"this agent has no DID resolver configured, so it cannot find the host".into(),
),
);
};
let minted = match crate::operations::room_oracle::present(
state,
auth,
&vta_did,
&req.room_id,
"read",
Some(req.host.as_str()),
None,
)
.await
{
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
let mut payload = serde_json::json!({
"roomId": req.room_id,
"presentation": minted.presentation,
});
if let Some(from) = req.from_epoch {
payload["fromEpoch"] = serde_json::json!(u64::from(from));
}
if let Some(limit) = req.limit {
payload["limit"] = serde_json::json!(u64::from(limit));
}
let key = format!("{vta_did}#key-0");
let reply = match crate::operations::room_host::send_room_task(
signing_context(state, auth),
&resolver,
&req.host,
&key,
&vta_did,
&key,
vti_rooms::wire::ROOMS_EPOCH_CHAIN_TYPE,
&format!("{}#response", vti_rooms::wire::ROOMS_EPOCH_CHAIN_TYPE),
payload,
)
.await
{
Ok(v) => v,
Err(e) => return app_error_to_reject(&doc, e),
};
let links: Vec<vti_rooms::wire::EpochLink> =
match serde_json::from_value(reply.get("links").cloned().unwrap_or(Value::Array(vec![]))) {
Ok(l) => l,
Err(e) => {
return app_error_to_reject(
&doc,
vti_common::error::AppError::Internal(format!(
"room host `{}` served rungs this agent cannot read: {e}",
req.host
)),
);
}
};
let fetched = links.len();
let (earliest, stored) =
match room_groups::store_links(&state.room_groups_ks, &req.room_id, links, now()).await {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.backfill", auth, &req.room_id).await;
success_response(
&doc,
serde_json::json!({
"roomId": req.room_id,
"earliestReadableEpoch": earliest,
"fetched": fetched,
"stored": stored,
}),
)
}
pub(super) async fn handle_open(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(r) = super::helpers::require_capability(
state,
auth,
&doc,
Capability::RoomOpen,
"opening a room record",
)
.await
{
return r;
}
let req: trust_tasks_rs::specs::rooms::keys::open::v0_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let plaintext = match room_groups::open_record(
&state.room_groups_ks,
&req.room_id,
&req.key,
u64::from(req.version),
&req.sealed.ciphertext,
&req.sealed.nonce,
u32::try_from(u64::from(req.sealed.epoch)).unwrap_or(u32::MAX),
)
.await
{
Ok(p) => p,
Err(e) => return app_error_to_reject(&doc, e),
};
record(state, "rooms.keys.open", auth, &req.room_id).await;
success_response(
&doc,
serde_json::json!({
"plaintext": base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
&plaintext,
)
}),
)
}
async fn require_invitation(
state: &AppState,
encoded: Option<&str>,
room_id: &str,
member_did: &str,
) -> Result<room_invitation::VerifiedInvitation, vti_common::error::AppError> {
let encoded = encoded.ok_or_else(|| {
vti_common::error::AppError::Validation(format!(
"no invitation presented for room `{room_id}`; joining a room is a two-party \
act and the invitation is the other party's half"
))
})?;
let keys = vti_rooms_dtg::DataIntegrityKeys(state.trust_task_vm_resolver());
let invitation = room_invitation::verify(encoded, room_id, member_did, &keys).await?;
if room_invitation::is_consumed(&state.room_invitations_ks, invitation.credential_id()).await? {
return Err(vti_common::error::AppError::Conflict(format!(
"invitation `{}` has already been used",
invitation.credential_id()
)));
}
Ok(invitation)
}
fn decode_b64(s: &str, what: &str) -> Result<Vec<u8>, vti_common::error::AppError> {
use base64::Engine as _;
base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(s.trim())
.map_err(|e| vti_common::error::AppError::Validation(format!("decode the {what}: {e}")))
}
fn now() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
fn rfc3339(unix_seconds: u64) -> String {
chrono::DateTime::from_timestamp(unix_seconds as i64, 0)
.unwrap_or_else(|| chrono::DateTime::from_timestamp(0, 0).expect("epoch is in range"))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
}
pub(super) async fn record(state: &AppState, action: &str, auth: &AuthClaims, room_id: &str) {
if let Err(e) = audit::record(
&state.audit_sink,
action,
&auth.did,
Some(room_id),
"success",
Some(TRANSPORT_TRUST_TASK),
None,
)
.await
{
tracing::error!(error = %e, action, "failed to record a room-group audit entry");
}
}