mirror of
https://tangled.org/tranquil.farm/tranquil-pds
synced 2026-07-30 08:32:38 +00:00
refactor(api): update repo batch/delete to use repo_ops, clean up remaining repo endpoints
This commit is contained in:
@@ -2,7 +2,6 @@ use axum::body::Body;
|
||||
use axum::{
|
||||
Json,
|
||||
extract::{Query, State},
|
||||
http::StatusCode,
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use bytes::Bytes;
|
||||
@@ -58,10 +57,7 @@ pub async fn upload_blob(
|
||||
}
|
||||
let mime_type_for_check = get_header_str(&headers, http::header::CONTENT_TYPE)
|
||||
.unwrap_or("application/octet-stream");
|
||||
let scope_proof = match user.verify_blob_upload(mime_type_for_check) {
|
||||
Ok(proof) => proof,
|
||||
Err(e) => return Ok(e.into_response()),
|
||||
};
|
||||
let scope_proof = user.verify_blob_upload(mime_type_for_check)?;
|
||||
(
|
||||
scope_proof.principal_did().into_did(),
|
||||
scope_proof.controller_did().map(|c| c.into_did()),
|
||||
@@ -237,7 +233,7 @@ pub async fn list_missing_blobs(
|
||||
State(state): State<AppState>,
|
||||
auth: Auth<NotTakendown>,
|
||||
Query(params): Query<ListMissingBlobsParams>,
|
||||
) -> Result<Response, ApiError> {
|
||||
) -> Result<Json<ListMissingBlobsOutput>, ApiError> {
|
||||
let did = &auth.did;
|
||||
let user = state
|
||||
.user_repo
|
||||
@@ -269,12 +265,8 @@ pub async fn list_missing_blobs(
|
||||
} else {
|
||||
None
|
||||
};
|
||||
Ok((
|
||||
StatusCode::OK,
|
||||
Json(ListMissingBlobsOutput {
|
||||
cursor: next_cursor,
|
||||
blobs,
|
||||
}),
|
||||
)
|
||||
.into_response())
|
||||
Ok(Json(ListMissingBlobsOutput {
|
||||
cursor: next_cursor,
|
||||
blobs,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -1,8 +1,4 @@
|
||||
use axum::{
|
||||
body::Bytes,
|
||||
extract::State,
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use axum::{Json, body::Bytes, extract::State};
|
||||
use jacquard_common::types::{integer::LimitedU32, string::Tid};
|
||||
use jacquard_repo::storage::BlockStore;
|
||||
use k256::ecdsa::SigningKey;
|
||||
@@ -22,7 +18,7 @@ pub async fn import_repo(
|
||||
State(state): State<AppState>,
|
||||
auth: Auth<NotTakendown>,
|
||||
body: Bytes,
|
||||
) -> Result<Response, ApiError> {
|
||||
) -> Result<Json<EmptyResponse>, ApiError> {
|
||||
let accepting_imports = tranquil_config::get().import.accepting;
|
||||
if !accepting_imports {
|
||||
return Err(ApiError::InvalidRequest(
|
||||
@@ -340,7 +336,7 @@ pub async fn import_repo(
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(EmptyResponse::ok().into_response())
|
||||
Ok(Json(EmptyResponse {}))
|
||||
}
|
||||
Err(ImportError::SizeLimitExceeded) => Err(ApiError::PayloadTooLarge(format!(
|
||||
"Import exceeds block limit of {}",
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
use crate::common;
|
||||
use axum::{
|
||||
Json,
|
||||
extract::{Query, State},
|
||||
@@ -5,7 +6,6 @@ use axum::{
|
||||
};
|
||||
use serde::Deserialize;
|
||||
use serde_json::json;
|
||||
use tranquil_pds::api::error::ApiError;
|
||||
use tranquil_pds::state::AppState;
|
||||
use tranquil_pds::types::AtIdentifier;
|
||||
|
||||
@@ -18,57 +18,22 @@ pub async fn describe_repo(
|
||||
State(state): State<AppState>,
|
||||
Query(input): Query<DescribeRepoInput>,
|
||||
) -> Response {
|
||||
let hostname_for_handles = tranquil_config::get().server.hostname_without_port();
|
||||
let user_row = if input.repo.is_did() {
|
||||
let did: tranquil_pds::types::Did = match input.repo.as_str().parse() {
|
||||
Ok(d) => d,
|
||||
Err(_) => return ApiError::InvalidRequest("Invalid DID format".into()).into_response(),
|
||||
};
|
||||
state
|
||||
.user_repo
|
||||
.get_by_did(&did)
|
||||
.await
|
||||
.map(|opt| opt.map(|r| (r.id, r.handle, r.did)))
|
||||
} else {
|
||||
let repo_str = input.repo.as_str();
|
||||
let handle_str = if !repo_str.contains('.') {
|
||||
format!("{}.{}", repo_str, hostname_for_handles)
|
||||
} else {
|
||||
repo_str.to_string()
|
||||
};
|
||||
let handle: tranquil_pds::types::Handle = match handle_str.parse() {
|
||||
Ok(h) => h,
|
||||
Err(_) => {
|
||||
return ApiError::InvalidRequest("Invalid handle format".into()).into_response();
|
||||
}
|
||||
};
|
||||
state
|
||||
.user_repo
|
||||
.get_by_handle(&handle)
|
||||
.await
|
||||
.map(|opt| opt.map(|r| (r.id, r.handle, r.did)))
|
||||
};
|
||||
let (user_id, handle, did) = match user_row {
|
||||
Ok(Some((id, handle, did))) => (id, handle, did),
|
||||
Ok(None) => {
|
||||
return ApiError::RepoNotFound(Some("Repo not found".into())).into_response();
|
||||
}
|
||||
Err(_) => {
|
||||
return ApiError::InternalError(None).into_response();
|
||||
}
|
||||
let resolved = match common::resolve_repo(state.user_repo.as_ref(), &input.repo).await {
|
||||
Ok(r) => r,
|
||||
Err(e) => return e.into_response(),
|
||||
};
|
||||
let collections = state
|
||||
.repo_repo
|
||||
.list_collections(user_id)
|
||||
.list_collections(resolved.user_id)
|
||||
.await
|
||||
.unwrap_or_default();
|
||||
let did_doc = json!({
|
||||
"id": did,
|
||||
"alsoKnownAs": [format!("at://{}", handle)]
|
||||
"id": resolved.did,
|
||||
"alsoKnownAs": [format!("at://{}", resolved.handle)]
|
||||
});
|
||||
Json(json!({
|
||||
"handle": handle,
|
||||
"did": did,
|
||||
"handle": resolved.handle,
|
||||
"did": resolved.did,
|
||||
"didDoc": did_doc,
|
||||
"collections": collections,
|
||||
"handleIsCorrect": true
|
||||
|
||||
@@ -1,27 +1,20 @@
|
||||
use super::validation::validate_record_with_status;
|
||||
use super::validation_mode::{ValidationMode, deserialize_validation_mode};
|
||||
use crate::repo::record::utils::{CommitParams, RecordOp, commit_and_log, extract_blob_cids};
|
||||
use axum::{
|
||||
Json,
|
||||
extract::State,
|
||||
http::StatusCode,
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use cid::Cid;
|
||||
use jacquard_repo::{commit::Commit, mst::Mst, storage::BlockStore};
|
||||
use crate::repo::record::write::CommitInfo;
|
||||
use axum::{Json, extract::State};
|
||||
use jacquard_repo::{mst::Mst, storage::BlockStore};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::json;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use tracing::info;
|
||||
use tranquil_pds::api::error::ApiError;
|
||||
use tranquil_pds::api::error::{ApiError, DbResultExt};
|
||||
use tranquil_pds::auth::{
|
||||
Active, Auth, WriteOpKind, require_not_migrated, require_verified_or_delegated,
|
||||
verify_batch_write_scopes,
|
||||
};
|
||||
use tranquil_pds::cid_types::CommitCid;
|
||||
use tranquil_pds::delegation::DelegationActionType;
|
||||
use tranquil_pds::repo::tracking::TrackingBlockStore;
|
||||
use tranquil_pds::repo_ops::{
|
||||
FinalizeParams, RecordOp, begin_repo_write, extract_blob_cids, finalize_repo_write,
|
||||
};
|
||||
use tranquil_pds::state::AppState;
|
||||
use tranquil_pds::types::{AtIdentifier, AtUri, Did, Nsid, Rkey};
|
||||
use tranquil_pds::validation::ValidationStatus;
|
||||
@@ -42,7 +35,7 @@ async fn process_single_write(
|
||||
did: &Did,
|
||||
validate: ValidationMode,
|
||||
tracking_store: &TrackingBlockStore,
|
||||
) -> Result<WriteAccumulator, Response> {
|
||||
) -> Result<WriteAccumulator, ApiError> {
|
||||
let WriteAccumulator {
|
||||
mst,
|
||||
mut results,
|
||||
@@ -60,32 +53,31 @@ async fn process_single_write(
|
||||
let validation_status = if validate.should_skip() {
|
||||
None
|
||||
} else {
|
||||
match validate_record_with_status(
|
||||
value,
|
||||
collection,
|
||||
rkey.as_ref(),
|
||||
validate.requires_lexicon(),
|
||||
Some(
|
||||
validate_record_with_status(
|
||||
value,
|
||||
collection,
|
||||
rkey.as_ref(),
|
||||
validate.requires_lexicon(),
|
||||
)
|
||||
.await?,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(status) => Some(status),
|
||||
Err(err_response) => return Err(*err_response),
|
||||
}
|
||||
};
|
||||
all_blob_cids.extend(extract_blob_cids(value));
|
||||
let rkey = rkey.clone().unwrap_or_else(Rkey::generate);
|
||||
let record_ipld = tranquil_pds::util::json_to_ipld(value);
|
||||
let record_bytes = serde_ipld_dagcbor::to_vec(&record_ipld).map_err(|_| {
|
||||
ApiError::InvalidRecord("Failed to serialize record".into()).into_response()
|
||||
})?;
|
||||
let record_cid = tracking_store.put(&record_bytes).await.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to store record".into())).into_response()
|
||||
})?;
|
||||
let record_bytes = serde_ipld_dagcbor::to_vec(&record_ipld)
|
||||
.map_err(|_| ApiError::InvalidRecord("Failed to serialize record".into()))?;
|
||||
let record_cid = tracking_store
|
||||
.put(&record_bytes)
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to store record".into())))?;
|
||||
let key = format!("{}/{}", collection, rkey);
|
||||
modified_keys.push(key.clone());
|
||||
let new_mst = mst.add(&key, record_cid).await.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to add to MST".into())).into_response()
|
||||
})?;
|
||||
let new_mst = mst
|
||||
.add(&key, record_cid)
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to add to MST".into())))?;
|
||||
let uri = AtUri::from_parts(did, collection, &rkey);
|
||||
results.push(WriteResult::CreateResult {
|
||||
uri,
|
||||
@@ -113,32 +105,31 @@ async fn process_single_write(
|
||||
let validation_status = if validate.should_skip() {
|
||||
None
|
||||
} else {
|
||||
match validate_record_with_status(
|
||||
value,
|
||||
collection,
|
||||
Some(rkey),
|
||||
validate.requires_lexicon(),
|
||||
Some(
|
||||
validate_record_with_status(
|
||||
value,
|
||||
collection,
|
||||
Some(rkey),
|
||||
validate.requires_lexicon(),
|
||||
)
|
||||
.await?,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(status) => Some(status),
|
||||
Err(err_response) => return Err(*err_response),
|
||||
}
|
||||
};
|
||||
all_blob_cids.extend(extract_blob_cids(value));
|
||||
let record_ipld = tranquil_pds::util::json_to_ipld(value);
|
||||
let record_bytes = serde_ipld_dagcbor::to_vec(&record_ipld).map_err(|_| {
|
||||
ApiError::InvalidRecord("Failed to serialize record".into()).into_response()
|
||||
})?;
|
||||
let record_cid = tracking_store.put(&record_bytes).await.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to store record".into())).into_response()
|
||||
})?;
|
||||
let record_bytes = serde_ipld_dagcbor::to_vec(&record_ipld)
|
||||
.map_err(|_| ApiError::InvalidRecord("Failed to serialize record".into()))?;
|
||||
let record_cid = tracking_store
|
||||
.put(&record_bytes)
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to store record".into())))?;
|
||||
let key = format!("{}/{}", collection, rkey);
|
||||
modified_keys.push(key.clone());
|
||||
let prev_record_cid = mst.get(&key).await.ok().flatten();
|
||||
let new_mst = mst.update(&key, record_cid).await.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to update MST".into())).into_response()
|
||||
})?;
|
||||
let new_mst = mst
|
||||
.update(&key, record_cid)
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to update MST".into())))?;
|
||||
let uri = AtUri::from_parts(did, collection, rkey);
|
||||
results.push(WriteResult::UpdateResult {
|
||||
uri,
|
||||
@@ -163,9 +154,10 @@ async fn process_single_write(
|
||||
let key = format!("{}/{}", collection, rkey);
|
||||
modified_keys.push(key.clone());
|
||||
let prev_record_cid = mst.get(&key).await.ok().flatten();
|
||||
let new_mst = mst.delete(&key).await.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to delete from MST".into())).into_response()
|
||||
})?;
|
||||
let new_mst = mst
|
||||
.delete(&key)
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to delete from MST".into())))?;
|
||||
results.push(WriteResult::DeleteResult {});
|
||||
ops.push(RecordOp::Delete {
|
||||
collection: collection.clone(),
|
||||
@@ -189,7 +181,7 @@ async fn process_writes(
|
||||
did: &Did,
|
||||
validate: ValidationMode,
|
||||
tracking_store: &TrackingBlockStore,
|
||||
) -> Result<WriteAccumulator, Response> {
|
||||
) -> Result<WriteAccumulator, ApiError> {
|
||||
use futures::stream::{self, TryStreamExt};
|
||||
let initial_acc = WriteAccumulator {
|
||||
mst: initial_mst,
|
||||
@@ -198,7 +190,7 @@ async fn process_writes(
|
||||
modified_keys: Vec::new(),
|
||||
all_blob_cids: Vec::new(),
|
||||
};
|
||||
stream::iter(writes.iter().map(Ok::<_, Response>))
|
||||
stream::iter(writes.iter().map(Ok::<_, ApiError>))
|
||||
.try_fold(initial_acc, |acc, write| async move {
|
||||
process_single_write(write, acc, did, validate, tracking_store).await
|
||||
})
|
||||
@@ -261,17 +253,11 @@ pub struct ApplyWritesOutput {
|
||||
pub results: Vec<WriteResult>,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct CommitInfo {
|
||||
pub cid: String,
|
||||
pub rev: String,
|
||||
}
|
||||
|
||||
pub async fn apply_writes(
|
||||
State(state): State<AppState>,
|
||||
auth: Auth<Active>,
|
||||
Json(input): Json<ApplyWritesInput>,
|
||||
) -> Result<Response, ApiError> {
|
||||
) -> Result<Json<ApplyWritesOutput>, ApiError> {
|
||||
info!(
|
||||
"apply_writes called: repo={}, writes={}",
|
||||
input.repo,
|
||||
@@ -288,7 +274,7 @@ pub async fn apply_writes(
|
||||
)));
|
||||
}
|
||||
|
||||
let batch_proof = match verify_batch_write_scopes(
|
||||
let batch_proof = verify_batch_write_scopes(
|
||||
&auth,
|
||||
&auth,
|
||||
&input.writes,
|
||||
@@ -302,10 +288,7 @@ pub async fn apply_writes(
|
||||
WriteOp::Update { .. } => WriteOpKind::Update,
|
||||
WriteOp::Delete { .. } => WriteOpKind::Delete,
|
||||
},
|
||||
) {
|
||||
Ok(proof) => proof,
|
||||
Err(e) => return Ok(e.into_response()),
|
||||
};
|
||||
)?;
|
||||
|
||||
let principal_did = batch_proof.principal_did();
|
||||
let controller_did = batch_proof.controller_did().map(|c| c.into_did());
|
||||
@@ -317,140 +300,35 @@ pub async fn apply_writes(
|
||||
}
|
||||
|
||||
let did = principal_did.into_did();
|
||||
if let Err(e) = require_not_migrated(&state, &did).await {
|
||||
return Ok(e);
|
||||
}
|
||||
if let Err(e) = require_verified_or_delegated(&state, batch_proof.user()).await {
|
||||
return Ok(e);
|
||||
}
|
||||
require_not_migrated(&state, &did).await?;
|
||||
require_verified_or_delegated(&state, batch_proof.user()).await?;
|
||||
|
||||
let user_id: uuid::Uuid = state
|
||||
.user_repo
|
||||
.get_id_by_did(&did)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.ok_or_else(|| ApiError::InternalError(Some("User not found".into())))?;
|
||||
.log_db_err("fetching user for batch write")?
|
||||
.ok_or(ApiError::InternalError(Some("User not found".into())))?;
|
||||
|
||||
let _write_lock = state.repo_write_locks.lock(user_id).await;
|
||||
let (ctx, mst) = begin_repo_write(&state, user_id, input.swap_commit.as_deref()).await?;
|
||||
|
||||
let root_cid_str = state
|
||||
.repo_repo
|
||||
.get_repo_root_cid_by_user_id(user_id)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.ok_or_else(|| ApiError::InternalError(Some("Repo root not found".into())))?;
|
||||
let current_root_cid = CommitCid::from_str(&root_cid_str)
|
||||
.map_err(|_| ApiError::InternalError(Some("Invalid repo root CID".into())))?;
|
||||
if let Some(swap_commit) = &input.swap_commit
|
||||
&& CommitCid::from_str(swap_commit).ok().as_ref() != Some(¤t_root_cid)
|
||||
{
|
||||
return Err(ApiError::InvalidSwap(Some("Repo has been modified".into())));
|
||||
}
|
||||
let tracking_store = TrackingBlockStore::new(state.block_store.clone());
|
||||
let commit_bytes = tracking_store
|
||||
.get(current_root_cid.as_cid())
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.ok_or_else(|| ApiError::InternalError(Some("Commit block not found".into())))?;
|
||||
let commit = Commit::from_cbor(&commit_bytes)
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to parse commit".into())))?;
|
||||
let original_mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None);
|
||||
let initial_mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None);
|
||||
let WriteAccumulator {
|
||||
mst,
|
||||
mst: final_mst,
|
||||
results,
|
||||
ops,
|
||||
modified_keys,
|
||||
all_blob_cids,
|
||||
} = match process_writes(
|
||||
} = process_writes(
|
||||
&input.writes,
|
||||
initial_mst,
|
||||
mst,
|
||||
&did,
|
||||
input.validate,
|
||||
&tracking_store,
|
||||
&ctx.tracking_store,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(acc) => acc,
|
||||
Err(response) => return Ok(response),
|
||||
};
|
||||
let new_mst_root = mst
|
||||
.persist()
|
||||
.await
|
||||
.map_err(|_| ApiError::InternalError(Some("Failed to persist MST".into())))?;
|
||||
let (new_mst_blocks, old_mst_blocks) = {
|
||||
let mut new_blocks = std::collections::BTreeMap::new();
|
||||
let mut old_blocks = std::collections::BTreeMap::new();
|
||||
for key in &modified_keys {
|
||||
mst.blocks_for_path(key, &mut new_blocks)
|
||||
.await
|
||||
.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to get new MST blocks for path".into()))
|
||||
})?;
|
||||
original_mst
|
||||
.blocks_for_path(key, &mut old_blocks)
|
||||
.await
|
||||
.map_err(|_| {
|
||||
ApiError::InternalError(Some("Failed to get old MST blocks for path".into()))
|
||||
})?;
|
||||
}
|
||||
(new_blocks, old_blocks)
|
||||
};
|
||||
let mut relevant_blocks = new_mst_blocks.clone();
|
||||
relevant_blocks.extend(old_mst_blocks.iter().map(|(k, v)| (*k, v.clone())));
|
||||
let written_cids: Vec<Cid> = tracking_store
|
||||
.get_all_relevant_cids()
|
||||
.into_iter()
|
||||
.chain(relevant_blocks.keys().copied())
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
let written_cids_str: Vec<String> = written_cids.iter().map(|c| c.to_string()).collect();
|
||||
let prev_record_cids = ops.iter().filter_map(|op| match op {
|
||||
RecordOp::Update {
|
||||
prev: Some(cid), ..
|
||||
}
|
||||
| RecordOp::Delete {
|
||||
prev: Some(cid), ..
|
||||
} => Some(*cid),
|
||||
_ => None,
|
||||
});
|
||||
let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid.into_cid())
|
||||
.chain(
|
||||
old_mst_blocks
|
||||
.keys()
|
||||
.filter(|cid| !new_mst_blocks.contains_key(*cid))
|
||||
.copied(),
|
||||
)
|
||||
.chain(prev_record_cids)
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
let commit_res = match commit_and_log(
|
||||
&state,
|
||||
CommitParams {
|
||||
did: &did,
|
||||
user_id,
|
||||
current_root_cid: Some(current_root_cid.into_cid()),
|
||||
prev_data_cid: Some(commit.data),
|
||||
new_mst_root,
|
||||
ops,
|
||||
blocks_cids: &written_cids_str,
|
||||
blobs: &all_blob_cids,
|
||||
obsolete_cids,
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(res) => res,
|
||||
Err(e) => return Err(ApiError::from(e)),
|
||||
};
|
||||
.await?;
|
||||
|
||||
if let Some(ref controller) = controller_did {
|
||||
let write_summary: Vec<serde_json::Value> = input
|
||||
let write_summary: Option<serde_json::Value> = controller_did.as_ref().map(|_| {
|
||||
let writes: Vec<serde_json::Value> = input
|
||||
.writes
|
||||
.iter()
|
||||
.map(|w| match w {
|
||||
@@ -475,34 +353,34 @@ pub async fn apply_writes(
|
||||
}),
|
||||
})
|
||||
.collect();
|
||||
json!({
|
||||
"action": "apply_writes",
|
||||
"count": input.writes.len(),
|
||||
"writes": writes
|
||||
})
|
||||
});
|
||||
|
||||
let _ = state
|
||||
.delegation_repo
|
||||
.log_delegation_action(
|
||||
&did,
|
||||
controller,
|
||||
Some(controller),
|
||||
DelegationActionType::RepoWrite,
|
||||
Some(json!({
|
||||
"action": "apply_writes",
|
||||
"count": input.writes.len(),
|
||||
"writes": write_summary
|
||||
})),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
Ok((
|
||||
StatusCode::OK,
|
||||
Json(ApplyWritesOutput {
|
||||
commit: CommitInfo {
|
||||
cid: commit_res.commit_cid.to_string(),
|
||||
rev: commit_res.rev,
|
||||
},
|
||||
results,
|
||||
}),
|
||||
let commit_result = finalize_repo_write(
|
||||
&state,
|
||||
ctx,
|
||||
final_mst,
|
||||
FinalizeParams {
|
||||
did: &did,
|
||||
user_id,
|
||||
controller_did: controller_did.as_ref(),
|
||||
delegation_detail: write_summary,
|
||||
ops,
|
||||
modified_keys: &modified_keys,
|
||||
blob_cids: &all_blob_cids,
|
||||
},
|
||||
)
|
||||
.into_response())
|
||||
.await?;
|
||||
|
||||
Ok(Json(ApplyWritesOutput {
|
||||
commit: CommitInfo {
|
||||
cid: commit_result.commit_cid.to_string(),
|
||||
rev: commit_result.rev,
|
||||
},
|
||||
results,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -1,13 +1,5 @@
|
||||
use crate::repo::record::utils::{
|
||||
CommitError, CommitParams, RecordOp, commit_and_log, get_current_root_cid,
|
||||
};
|
||||
use crate::repo::record::write::{CommitInfo, prepare_repo_write};
|
||||
use axum::{
|
||||
Json,
|
||||
extract::State,
|
||||
http::StatusCode,
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use axum::{Json, extract::State};
|
||||
use cid::Cid;
|
||||
use jacquard_repo::{commit::Commit, mst::Mst, storage::BlockStore};
|
||||
use serde::{Deserialize, Serialize};
|
||||
@@ -17,11 +9,12 @@ use std::sync::Arc;
|
||||
use tracing::error;
|
||||
use tranquil_pds::api::error::ApiError;
|
||||
use tranquil_pds::auth::{Active, Auth, VerifyScope};
|
||||
use tranquil_pds::cid_types::CommitCid;
|
||||
use tranquil_pds::delegation::DelegationActionType;
|
||||
use tranquil_pds::repo::tracking::TrackingBlockStore;
|
||||
use tranquil_pds::repo_ops::{
|
||||
CommitError, FinalizeParams, RecordOp, begin_repo_write, finalize_repo_write,
|
||||
};
|
||||
use tranquil_pds::state::AppState;
|
||||
use tranquil_pds::types::{AtIdentifier, AtUri, Nsid, Rkey};
|
||||
use tranquil_pds::types::{AtIdentifier, AtUri, Did, Nsid, Rkey};
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct DeleteRecordInput {
|
||||
@@ -45,168 +38,66 @@ pub async fn delete_record(
|
||||
State(state): State<AppState>,
|
||||
auth: Auth<Active>,
|
||||
Json(input): Json<DeleteRecordInput>,
|
||||
) -> Result<Response, tranquil_pds::api::error::ApiError> {
|
||||
let scope_proof = match auth.verify_repo_delete(&input.collection) {
|
||||
Ok(proof) => proof,
|
||||
Err(e) => return Ok(e.into_response()),
|
||||
};
|
||||
|
||||
let repo_auth = match prepare_repo_write(&state, &scope_proof, &input.repo).await {
|
||||
Ok(res) => res,
|
||||
Err(err_res) => return Ok(err_res),
|
||||
};
|
||||
|
||||
) -> Result<Json<DeleteRecordOutput>, ApiError> {
|
||||
let scope_proof = auth.verify_repo_delete(&input.collection)?;
|
||||
let repo_auth = prepare_repo_write(&state, &scope_proof, &input.repo).await?;
|
||||
let did = repo_auth.did;
|
||||
let user_id = repo_auth.user_id;
|
||||
let controller_did = repo_auth.controller_did;
|
||||
|
||||
let _write_lock = state.repo_write_locks.lock(user_id).await;
|
||||
let current_root_cid = get_current_root_cid(&state, user_id).await?;
|
||||
let (ctx, mst) = begin_repo_write(&state, user_id, input.swap_commit.as_deref()).await?;
|
||||
|
||||
if let Some(swap_commit) = &input.swap_commit
|
||||
&& CommitCid::from_str(swap_commit).ok().as_ref() != Some(¤t_root_cid)
|
||||
{
|
||||
return Ok(ApiError::InvalidSwap(Some("Repo has been modified".into())).into_response());
|
||||
}
|
||||
let tracking_store = TrackingBlockStore::new(state.block_store.clone());
|
||||
let commit_bytes = match tracking_store.get(current_root_cid.as_cid()).await {
|
||||
Ok(Some(b)) => b,
|
||||
_ => {
|
||||
return Ok(
|
||||
ApiError::InternalError(Some("Commit block not found".into())).into_response(),
|
||||
);
|
||||
}
|
||||
};
|
||||
let commit = match Commit::from_cbor(&commit_bytes) {
|
||||
Ok(c) => c,
|
||||
_ => {
|
||||
return Ok(
|
||||
ApiError::InternalError(Some("Failed to parse commit".into())).into_response(),
|
||||
);
|
||||
}
|
||||
};
|
||||
let mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None);
|
||||
let key = format!("{}/{}", input.collection, input.rkey);
|
||||
|
||||
if let Some(swap_record_str) = &input.swap_record {
|
||||
let expected_cid = Cid::from_str(swap_record_str).ok();
|
||||
let actual_cid = mst.get(&key).await.ok().flatten();
|
||||
if expected_cid != actual_cid {
|
||||
return Ok(ApiError::InvalidSwap(Some(
|
||||
return Err(ApiError::InvalidSwap(Some(
|
||||
"Record has been modified or does not exist".into(),
|
||||
))
|
||||
.into_response());
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
let prev_record_cid = mst.get(&key).await.ok().flatten();
|
||||
if prev_record_cid.is_none() {
|
||||
return Ok((StatusCode::OK, Json(DeleteRecordOutput { commit: None })).into_response());
|
||||
return Ok(Json(DeleteRecordOutput { commit: None }));
|
||||
}
|
||||
let new_mst = match mst.delete(&key).await {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
error!("Failed to delete from MST: {:?}", e);
|
||||
return Ok(ApiError::InternalError(Some(format!(
|
||||
"Failed to delete from MST: {:?}",
|
||||
e
|
||||
)))
|
||||
.into_response());
|
||||
}
|
||||
};
|
||||
let new_mst_root = match new_mst.persist().await {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
error!("Failed to persist MST: {:?}", e);
|
||||
return Ok(
|
||||
ApiError::InternalError(Some("Failed to persist MST".into())).into_response(),
|
||||
);
|
||||
}
|
||||
};
|
||||
let collection_for_audit = input.collection.to_string();
|
||||
let rkey_for_audit = input.rkey.to_string();
|
||||
|
||||
let new_mst = mst.delete(&key).await.map_err(|e| {
|
||||
error!("Failed to delete from MST: {:?}", e);
|
||||
ApiError::InternalError(Some("Failed to delete from MST".into()))
|
||||
})?;
|
||||
|
||||
let op = RecordOp::Delete {
|
||||
collection: input.collection.clone(),
|
||||
rkey: input.rkey.clone(),
|
||||
prev: prev_record_cid,
|
||||
};
|
||||
let mut new_mst_blocks = std::collections::BTreeMap::new();
|
||||
let mut old_mst_blocks = std::collections::BTreeMap::new();
|
||||
if new_mst
|
||||
.blocks_for_path(&key, &mut new_mst_blocks)
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
return Ok(
|
||||
ApiError::InternalError(Some("Failed to get new MST blocks for path".into()))
|
||||
.into_response(),
|
||||
);
|
||||
}
|
||||
if mst
|
||||
.blocks_for_path(&key, &mut old_mst_blocks)
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
return Ok(
|
||||
ApiError::InternalError(Some("Failed to get old MST blocks for path".into()))
|
||||
.into_response(),
|
||||
);
|
||||
}
|
||||
let mut relevant_blocks = new_mst_blocks.clone();
|
||||
relevant_blocks.extend(old_mst_blocks.iter().map(|(k, v)| (*k, v.clone())));
|
||||
let written_cids: Vec<Cid> = tracking_store
|
||||
.get_all_relevant_cids()
|
||||
.into_iter()
|
||||
.chain(relevant_blocks.keys().copied())
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
let written_cids_str: Vec<String> = written_cids.iter().map(|c| c.to_string()).collect();
|
||||
let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid.into_cid())
|
||||
.chain(
|
||||
old_mst_blocks
|
||||
.keys()
|
||||
.filter(|cid| !new_mst_blocks.contains_key(*cid))
|
||||
.copied(),
|
||||
)
|
||||
.chain(prev_record_cid)
|
||||
.collect();
|
||||
let commit_result = match commit_and_log(
|
||||
|
||||
let modified_keys = [key];
|
||||
|
||||
let commit_result = finalize_repo_write(
|
||||
&state,
|
||||
CommitParams {
|
||||
ctx,
|
||||
new_mst,
|
||||
FinalizeParams {
|
||||
did: &did,
|
||||
user_id,
|
||||
current_root_cid: Some(current_root_cid.into_cid()),
|
||||
prev_data_cid: Some(commit.data),
|
||||
new_mst_root,
|
||||
controller_did: controller_did.as_ref(),
|
||||
delegation_detail: controller_did.as_ref().map(|_| {
|
||||
json!({
|
||||
"action": "delete",
|
||||
"collection": input.collection,
|
||||
"rkey": input.rkey
|
||||
})
|
||||
}),
|
||||
ops: vec![op],
|
||||
blocks_cids: &written_cids_str,
|
||||
blobs: &[],
|
||||
obsolete_cids,
|
||||
modified_keys: &modified_keys,
|
||||
blob_cids: &[],
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(res) => res,
|
||||
Err(e) => return Ok(ApiError::from(e).into_response()),
|
||||
};
|
||||
|
||||
if let Some(ref controller) = controller_did {
|
||||
let _ = state
|
||||
.delegation_repo
|
||||
.log_delegation_action(
|
||||
&did,
|
||||
controller,
|
||||
Some(controller),
|
||||
DelegationActionType::RepoWrite,
|
||||
Some(json!({
|
||||
"action": "delete",
|
||||
"collection": collection_for_audit,
|
||||
"rkey": rkey_for_audit
|
||||
})),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
.await?;
|
||||
|
||||
let deleted_uri = AtUri::from_parts(&did, &input.collection, &input.rkey);
|
||||
if let Err(e) = state
|
||||
@@ -217,19 +108,14 @@ pub async fn delete_record(
|
||||
error!("Failed to remove backlinks for {}: {}", deleted_uri, e);
|
||||
}
|
||||
|
||||
Ok((
|
||||
StatusCode::OK,
|
||||
Json(DeleteRecordOutput {
|
||||
commit: Some(CommitInfo {
|
||||
cid: commit_result.commit_cid.to_string(),
|
||||
rev: commit_result.rev,
|
||||
}),
|
||||
Ok(Json(DeleteRecordOutput {
|
||||
commit: Some(CommitInfo {
|
||||
cid: commit_result.commit_cid.to_string(),
|
||||
rev: commit_result.rev,
|
||||
}),
|
||||
)
|
||||
.into_response())
|
||||
}))
|
||||
}
|
||||
|
||||
use tranquil_pds::types::Did;
|
||||
use uuid::Uuid;
|
||||
|
||||
pub async fn delete_record_internal(
|
||||
@@ -239,6 +125,8 @@ pub async fn delete_record_internal(
|
||||
collection: &Nsid,
|
||||
rkey: &Rkey,
|
||||
) -> Result<(), CommitError> {
|
||||
use tranquil_pds::repo_ops::{CommitParams, RecordOp, commit_and_log};
|
||||
|
||||
let _write_lock = state.repo_write_locks.lock(user_id).await;
|
||||
|
||||
let root_cid_str = state
|
||||
@@ -303,19 +191,6 @@ pub async fn delete_record_internal(
|
||||
.await
|
||||
.map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?;
|
||||
|
||||
let mut relevant_blocks = new_mst_blocks.clone();
|
||||
relevant_blocks.extend(old_mst_blocks.iter().map(|(k, v)| (*k, v.clone())));
|
||||
|
||||
let written_cids: Vec<Cid> = tracking_store
|
||||
.get_all_relevant_cids()
|
||||
.into_iter()
|
||||
.chain(relevant_blocks.keys().copied())
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
|
||||
let written_cids_str: Vec<String> = written_cids.iter().map(|c| c.to_string()).collect();
|
||||
|
||||
let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid)
|
||||
.chain(
|
||||
old_mst_blocks
|
||||
@@ -326,6 +201,19 @@ pub async fn delete_record_internal(
|
||||
.chain(std::iter::once(prev_cid))
|
||||
.collect();
|
||||
|
||||
let mut relevant_blocks = new_mst_blocks;
|
||||
relevant_blocks.extend(old_mst_blocks);
|
||||
|
||||
let written_cids: Vec<Cid> = tracking_store
|
||||
.get_all_relevant_cids()
|
||||
.into_iter()
|
||||
.chain(relevant_blocks.keys().copied())
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
|
||||
let written_cids_str: Vec<String> = written_cids.iter().map(ToString::to_string).collect();
|
||||
|
||||
commit_and_log(
|
||||
state,
|
||||
CommitParams {
|
||||
|
||||
@@ -23,9 +23,8 @@ pub use account_status::{
|
||||
};
|
||||
pub use app_password::{create_app_password, list_app_passwords, revoke_app_password};
|
||||
pub use email::{
|
||||
authorize_email_update, check_channel_verified, check_email_in_use,
|
||||
check_email_update_status, check_email_verified, confirm_email, request_email_update,
|
||||
update_email,
|
||||
authorize_email_update, check_channel_verified, check_email_in_use, check_email_update_status,
|
||||
check_email_verified, confirm_email, request_email_update, update_email,
|
||||
};
|
||||
pub use invite::{create_invite_code, create_invite_codes, get_account_invite_codes};
|
||||
pub use logo::get_logo;
|
||||
@@ -44,9 +43,8 @@ pub use password::{
|
||||
set_password,
|
||||
};
|
||||
pub use reauth::{
|
||||
check_legacy_session_mfa, check_reauth_required, get_reauth_status,
|
||||
legacy_mfa_required_response, reauth_passkey_finish, reauth_passkey_start, reauth_password,
|
||||
reauth_required_response, reauth_totp, update_mfa_verified,
|
||||
check_legacy_session_mfa, check_reauth_required, get_reauth_status, reauth_passkey_finish,
|
||||
reauth_passkey_start, reauth_password, reauth_totp, update_mfa_verified,
|
||||
};
|
||||
pub use service_auth::get_service_auth;
|
||||
pub use session::{
|
||||
@@ -64,4 +62,7 @@ pub use trusted_devices::{
|
||||
trust_device, update_trusted_device,
|
||||
};
|
||||
pub use verify_email::{resend_migration_verification, verify_migration_email};
|
||||
pub use verify_token::{VerifyTokenInput, VerifyTokenOutput, verify_token, verify_token_internal};
|
||||
pub use verify_token::{
|
||||
VerifyTokenInput, VerifyTokenOutput, confirm_channel_verification, verify_token,
|
||||
verify_token_internal,
|
||||
};
|
||||
|
||||
@@ -112,7 +112,7 @@ pub async fn get_service_auth(
|
||||
};
|
||||
|
||||
let lxm = params.lxm.as_ref();
|
||||
let lxm_for_token = lxm.map_or("*", |n| n.as_str());
|
||||
let lxm_for_token = lxm.map_or("*", |v| v.as_str());
|
||||
|
||||
if let Some(method) = lxm {
|
||||
if let Err(e) = tranquil_pds::auth::scope_check::check_rpc_scope(
|
||||
@@ -121,7 +121,7 @@ pub async fn get_service_auth(
|
||||
params.aud.as_str(),
|
||||
method.as_str(),
|
||||
) {
|
||||
return e;
|
||||
return e.into_response();
|
||||
}
|
||||
} else if auth.is_oauth() {
|
||||
let permissions = auth.permissions();
|
||||
|
||||
@@ -74,27 +74,14 @@ pub async fn resend_migration_verification(
|
||||
return Ok(Json(ResendMigrationVerificationOutput { sent: true }));
|
||||
}
|
||||
|
||||
let hostname = &tranquil_config::get().server.hostname;
|
||||
let token = tranquil_pds::auth::verification_token::generate_migration_token(
|
||||
crate::identity::provision::enqueue_migration_verification(
|
||||
&state,
|
||||
user.id,
|
||||
&user.did,
|
||||
channel,
|
||||
&identifier,
|
||||
);
|
||||
let formatted_token = tranquil_pds::auth::verification_token::format_token_for_display(&token);
|
||||
|
||||
if let Err(e) = tranquil_pds::comms::comms_repo::enqueue_migration_verification(
|
||||
state.user_repo.as_ref(),
|
||||
state.infra_repo.as_ref(),
|
||||
user.id,
|
||||
channel,
|
||||
&identifier,
|
||||
&formatted_token,
|
||||
hostname,
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!(error = ?e, channel = ?channel, "Failed to enqueue migration verification");
|
||||
}
|
||||
.await;
|
||||
|
||||
info!(did = %user.did, channel = ?channel, "Resent migration verification");
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ use axum::{
|
||||
use cid::Cid;
|
||||
use jacquard_repo::storage::BlockStore;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::Value;
|
||||
use std::str::FromStr;
|
||||
use tranquil_pds::api::error::ApiError;
|
||||
use tranquil_pds::auth::{Active, Auth, Permissive};
|
||||
@@ -51,7 +52,7 @@ pub async fn dereference_scope(
|
||||
State(state): State<AppState>,
|
||||
_auth: Auth<Active>,
|
||||
Json(input): Json<DereferenceScopeInput>,
|
||||
) -> Result<Response, ApiError> {
|
||||
) -> Result<Json<DereferenceScopeOutput>, ApiError> {
|
||||
let scope_parts: Vec<&str> = input.scope.split_whitespace().collect();
|
||||
let mut resolved_scopes: Vec<String> = Vec::new();
|
||||
|
||||
@@ -96,7 +97,7 @@ pub async fn dereference_scope(
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(scope_value) = scope_record.get("scope").and_then(|v| v.as_str()) {
|
||||
if let Some(scope_value) = scope_record.get("scope").and_then(Value::as_str) {
|
||||
let _ = state
|
||||
.cache
|
||||
.set(
|
||||
@@ -118,6 +119,5 @@ pub async fn dereference_scope(
|
||||
|
||||
Ok(Json(DereferenceScopeOutput {
|
||||
scope: resolved_scopes.join(" "),
|
||||
})
|
||||
.into_response())
|
||||
}))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user