mirror of
https://tangled.org/tranquil.farm/tranquil-pds
synced 2026-09-28 13:05:33 +00:00
Idk. Code quality in general?
This commit is contained in:
+22
-2
@@ -15,11 +15,21 @@ use sha2::{Digest, Sha256};
|
||||
use std::str::FromStr;
|
||||
use tracing::error;
|
||||
|
||||
const MAX_BLOB_SIZE: usize = 1_000_000;
|
||||
|
||||
pub async fn upload_blob(
|
||||
State(state): State<AppState>,
|
||||
headers: axum::http::HeaderMap,
|
||||
body: Bytes,
|
||||
) -> Response {
|
||||
if body.len() > MAX_BLOB_SIZE {
|
||||
return (
|
||||
StatusCode::PAYLOAD_TOO_LARGE,
|
||||
Json(json!({"error": "BlobTooLarge", "message": format!("Blob size {} exceeds maximum of {} bytes", body.len(), MAX_BLOB_SIZE)})),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
|
||||
let token = match crate::auth::extract_bearer_token_from_header(
|
||||
headers.get("Authorization").and_then(|h| h.to_str().ok())
|
||||
) {
|
||||
@@ -57,7 +67,17 @@ pub async fn upload_blob(
|
||||
let mut hasher = Sha256::new();
|
||||
hasher.update(&data);
|
||||
let hash = hasher.finalize();
|
||||
let multihash = Multihash::wrap(0x12, &hash).unwrap();
|
||||
let multihash = match Multihash::wrap(0x12, &hash) {
|
||||
Ok(mh) => mh,
|
||||
Err(e) => {
|
||||
error!("Failed to create multihash for blob: {:?}", e);
|
||||
return (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(json!({"error": "InternalError", "message": "Failed to hash blob"})),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
};
|
||||
let cid = Cid::new_v1(0x55, multihash);
|
||||
let cid_str = cid.to_string();
|
||||
|
||||
@@ -207,7 +227,7 @@ pub async fn list_missing_blobs(
|
||||
}
|
||||
};
|
||||
|
||||
let limit = params.limit.unwrap_or(500).min(1000);
|
||||
let limit = params.limit.unwrap_or(500).clamp(1, 1000);
|
||||
let cursor_str = params.cursor.unwrap_or_default();
|
||||
let (cursor_collection, cursor_rkey) = if cursor_str.contains('|') {
|
||||
let parts: Vec<&str> = cursor_str.split('|').collect();
|
||||
|
||||
+3
-14
@@ -1,3 +1,4 @@
|
||||
use crate::api::ApiError;
|
||||
use crate::state::AppState;
|
||||
use crate::sync::import::{apply_import, parse_car, ImportError};
|
||||
use crate::sync::verify::CarVerifier;
|
||||
@@ -54,24 +55,12 @@ pub async fn import_repo(
|
||||
headers.get("Authorization").and_then(|h| h.to_str().ok()),
|
||||
) {
|
||||
Some(t) => t,
|
||||
None => {
|
||||
return (
|
||||
StatusCode::UNAUTHORIZED,
|
||||
Json(json!({"error": "AuthenticationRequired"})),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
None => return ApiError::AuthenticationRequired.into_response(),
|
||||
};
|
||||
|
||||
let auth_user = match crate::auth::validate_bearer_token(&state.db, &token).await {
|
||||
Ok(user) => user,
|
||||
Err(e) => {
|
||||
return (
|
||||
StatusCode::UNAUTHORIZED,
|
||||
Json(json!({"error": "AuthenticationFailed", "message": e})),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
Err(e) => return ApiError::from(e).into_response(),
|
||||
};
|
||||
|
||||
let did = &auth_user.did;
|
||||
|
||||
@@ -17,6 +17,8 @@ use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use tracing::error;
|
||||
|
||||
const MAX_BATCH_WRITES: usize = 200;
|
||||
|
||||
#[derive(Deserialize)]
|
||||
#[serde(tag = "$type")]
|
||||
pub enum WriteOp {
|
||||
@@ -115,10 +117,10 @@ pub async fn apply_writes(
|
||||
.into_response();
|
||||
}
|
||||
|
||||
if input.writes.len() > 200 {
|
||||
if input.writes.len() > MAX_BATCH_WRITES {
|
||||
return (
|
||||
StatusCode::BAD_REQUEST,
|
||||
Json(json!({"error": "InvalidRequest", "message": "Too many writes (max 200)"})),
|
||||
Json(json!({"error": "InvalidRequest", "message": format!("Too many writes (max {})", MAX_BATCH_WRITES)})),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
@@ -213,11 +215,23 @@ pub async fn apply_writes(
|
||||
.clone()
|
||||
.unwrap_or_else(|| Utc::now().format("%Y%m%d%H%M%S%f").to_string());
|
||||
let mut record_bytes = Vec::new();
|
||||
serde_ipld_dagcbor::to_writer(&mut record_bytes, value).unwrap();
|
||||
let record_cid = tracking_store.put(&record_bytes).await.unwrap();
|
||||
if serde_ipld_dagcbor::to_writer(&mut record_bytes, value).is_err() {
|
||||
return (StatusCode::BAD_REQUEST, Json(json!({"error": "InvalidRecord", "message": "Failed to serialize record"}))).into_response();
|
||||
}
|
||||
let record_cid = match tracking_store.put(&record_bytes).await {
|
||||
Ok(c) => c,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to store record"}))).into_response(),
|
||||
};
|
||||
|
||||
let key = format!("{}/{}", collection.parse::<Nsid>().unwrap(), rkey);
|
||||
mst = mst.add(&key, record_cid).await.unwrap();
|
||||
let collection_nsid = match collection.parse::<Nsid>() {
|
||||
Ok(n) => n,
|
||||
Err(_) => return (StatusCode::BAD_REQUEST, Json(json!({"error": "InvalidCollection", "message": "Invalid collection NSID"}))).into_response(),
|
||||
};
|
||||
let key = format!("{}/{}", collection_nsid, rkey);
|
||||
mst = match mst.add(&key, record_cid).await {
|
||||
Ok(m) => m,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to add to MST"}))).into_response(),
|
||||
};
|
||||
|
||||
let uri = format!("at://{}/{}/{}", did, collection, rkey);
|
||||
results.push(WriteResult::CreateResult {
|
||||
@@ -236,11 +250,23 @@ pub async fn apply_writes(
|
||||
value,
|
||||
} => {
|
||||
let mut record_bytes = Vec::new();
|
||||
serde_ipld_dagcbor::to_writer(&mut record_bytes, value).unwrap();
|
||||
let record_cid = tracking_store.put(&record_bytes).await.unwrap();
|
||||
if serde_ipld_dagcbor::to_writer(&mut record_bytes, value).is_err() {
|
||||
return (StatusCode::BAD_REQUEST, Json(json!({"error": "InvalidRecord", "message": "Failed to serialize record"}))).into_response();
|
||||
}
|
||||
let record_cid = match tracking_store.put(&record_bytes).await {
|
||||
Ok(c) => c,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to store record"}))).into_response(),
|
||||
};
|
||||
|
||||
let key = format!("{}/{}", collection.parse::<Nsid>().unwrap(), rkey);
|
||||
mst = mst.update(&key, record_cid).await.unwrap();
|
||||
let collection_nsid = match collection.parse::<Nsid>() {
|
||||
Ok(n) => n,
|
||||
Err(_) => return (StatusCode::BAD_REQUEST, Json(json!({"error": "InvalidCollection", "message": "Invalid collection NSID"}))).into_response(),
|
||||
};
|
||||
let key = format!("{}/{}", collection_nsid, rkey);
|
||||
mst = match mst.update(&key, record_cid).await {
|
||||
Ok(m) => m,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to update MST"}))).into_response(),
|
||||
};
|
||||
|
||||
let uri = format!("at://{}/{}/{}", did, collection, rkey);
|
||||
results.push(WriteResult::UpdateResult {
|
||||
@@ -254,8 +280,15 @@ pub async fn apply_writes(
|
||||
});
|
||||
}
|
||||
WriteOp::Delete { collection, rkey } => {
|
||||
let key = format!("{}/{}", collection.parse::<Nsid>().unwrap(), rkey);
|
||||
mst = mst.delete(&key).await.unwrap();
|
||||
let collection_nsid = match collection.parse::<Nsid>() {
|
||||
Ok(n) => n,
|
||||
Err(_) => return (StatusCode::BAD_REQUEST, Json(json!({"error": "InvalidCollection", "message": "Invalid collection NSID"}))).into_response(),
|
||||
};
|
||||
let key = format!("{}/{}", collection_nsid, rkey);
|
||||
mst = match mst.delete(&key).await {
|
||||
Ok(m) => m,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to delete from MST"}))).into_response(),
|
||||
};
|
||||
|
||||
results.push(WriteResult::DeleteResult {});
|
||||
ops.push(RecordOp::Delete {
|
||||
@@ -266,7 +299,10 @@ pub async fn apply_writes(
|
||||
}
|
||||
}
|
||||
|
||||
let new_mst_root = mst.persist().await.unwrap();
|
||||
let new_mst_root = match mst.persist().await {
|
||||
Ok(c) => c,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to persist MST"}))).into_response(),
|
||||
};
|
||||
let written_cids = tracking_store.get_written_cids();
|
||||
let written_cids_str = written_cids
|
||||
.iter()
|
||||
|
||||
@@ -55,8 +55,11 @@ pub async fn commit_and_log(
|
||||
let new_root_cid = state.block_store.put(&new_commit_bytes).await
|
||||
.map_err(|e| format!("Failed to save commit block: {:?}", e))?;
|
||||
|
||||
let mut tx = state.db.begin().await
|
||||
.map_err(|e| format!("Failed to begin transaction: {}", e))?;
|
||||
|
||||
sqlx::query!("UPDATE repos SET repo_root_cid = $1 WHERE user_id = $2", new_root_cid.to_string(), user_id)
|
||||
.execute(&state.db)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| format!("DB Error (repos): {}", e))?;
|
||||
|
||||
@@ -71,7 +74,7 @@ pub async fn commit_and_log(
|
||||
rkey,
|
||||
cid.to_string()
|
||||
)
|
||||
.execute(&state.db)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| format!("DB Error (records): {}", e))?;
|
||||
}
|
||||
@@ -82,7 +85,7 @@ pub async fn commit_and_log(
|
||||
collection,
|
||||
rkey
|
||||
)
|
||||
.execute(&state.db)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| format!("DB Error (records): {}", e))?;
|
||||
}
|
||||
@@ -126,17 +129,20 @@ pub async fn commit_and_log(
|
||||
&[] as &[String],
|
||||
blocks_cids,
|
||||
)
|
||||
.fetch_one(&state.db)
|
||||
.fetch_one(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| format!("DB Error (repo_seq): {}", e))?;
|
||||
|
||||
sqlx::query(
|
||||
&format!("NOTIFY repo_updates, '{}'", seq_row.seq)
|
||||
)
|
||||
.execute(&state.db)
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| format!("DB Error (notify): {}", e))?;
|
||||
|
||||
tx.commit().await
|
||||
.map_err(|e| format!("Failed to commit transaction: {}", e))?;
|
||||
|
||||
Ok(CommitResult {
|
||||
commit_cid: new_root_cid,
|
||||
rev: rev.to_string(),
|
||||
|
||||
@@ -294,11 +294,20 @@ pub async fn put_record(
|
||||
};
|
||||
|
||||
let new_mst = if existing_cid.is_some() {
|
||||
mst.update(&key, record_cid).await.unwrap()
|
||||
match mst.update(&key, record_cid).await {
|
||||
Ok(m) => m,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to update MST"}))).into_response(),
|
||||
}
|
||||
} else {
|
||||
mst.add(&key, record_cid).await.unwrap()
|
||||
match mst.add(&key, record_cid).await {
|
||||
Ok(m) => m,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to add to MST"}))).into_response(),
|
||||
}
|
||||
};
|
||||
let new_mst_root = match new_mst.persist().await {
|
||||
Ok(c) => c,
|
||||
Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(json!({"error": "InternalError", "message": "Failed to persist MST"}))).into_response(),
|
||||
};
|
||||
let new_mst_root = new_mst.persist().await.unwrap();
|
||||
|
||||
let op = if existing_cid.is_some() {
|
||||
RecordOp::Update { collection: input.collection.clone(), rkey: input.rkey.clone(), cid: record_cid }
|
||||
|
||||
Reference in New Issue
Block a user