diff --git a/.slop-allow b/.slop-allow new file mode 100644 index 000000000..4e161fbc8 --- /dev/null +++ b/.slop-allow @@ -0,0 +1 @@ +named diff --git a/knot2/crates/knot-atproto/src/jwt.rs b/knot2/crates/knot-atproto/src/jwt.rs index 99f7e8ca2..20d6eb063 100644 --- a/knot2/crates/knot-atproto/src/jwt.rs +++ b/knot2/crates/knot-atproto/src/jwt.rs @@ -72,6 +72,8 @@ pub enum JwtError { NotAJwt, #[error("issuer {value:?} isn't valid account DID")] MalformedIssuer { value: String }, + #[error("Audience {value:?} isn't a repo DID.")] + MalformedAudience { value: String }, #[error("token has no jti, refusing write without replay protection")] MissingNonce, #[error("{len}-byte jti exceeds {MAX_NONCE_BYTES}-byte nonce limit")] @@ -159,6 +161,15 @@ pub fn issuer(parsed: &ParsedJwt) -> Result { }) } +fn repo_audience(parsed: &ParsedJwt) -> Result { + let aud = parsed.claims().aud.audience().as_str().to_owned(); + RepoDid::new(&aud).map_err(|_| JwtError::MalformedAudience { value: aud }) +} + +pub fn token_audience(token: &ServiceJwt) -> Result { + repo_audience(&parse(token)?) +} + fn verifying_key(key: &CryptoKey<'_>) -> Result { match key.codec { KeyCodec::Secp256k1 => VerifyKey::from_k256_bytes(&key.bytes) diff --git a/knot2/crates/knot-atproto/src/lib.rs b/knot2/crates/knot-atproto/src/lib.rs index 9e9354bc1..62b2c2ae6 100644 --- a/knot2/crates/knot-atproto/src/lib.rs +++ b/knot2/crates/knot-atproto/src/lib.rs @@ -12,7 +12,7 @@ pub use auth::{OauthAuthorizer, PointerAuth, PointerAuthorizer, ServiceAuth}; pub use identity::{ IdentityError, MintNonce, PreparedRepoDid, knot_did_document, prepare_repo_did, }; -pub use jwt::{JwtError, JwtNonce, ServiceJwt}; +pub use jwt::{JwtError, JwtNonce, ServiceJwt, token_audience}; pub use plc::{PDS_SERVICE_TYPE, RepoTarget, update_operation}; pub use pointer::PointerReceipt; pub use pubkeys::{KeyParseError, parse_authorized_key}; diff --git a/knot2/crates/knot-xrpc/src/atproto.rs b/knot2/crates/knot-xrpc/src/atproto.rs index fe04d6a3d..5f85a0862 100644 --- a/knot2/crates/knot-xrpc/src/atproto.rs +++ b/knot2/crates/knot-xrpc/src/atproto.rs @@ -62,7 +62,7 @@ fn json_of(ipld: &Ipld) -> Result { .map(Value::Number) .ok_or_else(|| "a non-finite float in the record doesn't fit JSON".to_string()), Ipld::String(string) => Ok(Value::String(string.clone())), - Ipld::Link(cid) => Ok(Value::String(cid.to_string())), + Ipld::Link(cid) => Ok(serde_json::json!({ "$link": cid.to_string() })), Ipld::Bytes(bytes) => { use base64::Engine as _; Ok(serde_json::json!({ @@ -105,9 +105,9 @@ fn storage(error: RepoError) -> XrpcError { /// One atproto repository this knot can serve: /// the knot's own meta-repo under its `did:web`, /// or a hosted repository under its repo DID. -struct AtprotoRepo { - did: RepoDid, - repo: Repo, +pub(crate) struct AtprotoRepo { + pub(crate) did: RepoDid, + pub(crate) repo: Repo, } impl AtprotoRepo { @@ -153,7 +153,7 @@ impl<'de> Deserialize<'de> for RepoSubject { } } -fn resolve( +pub(crate) fn resolve( state: &XrpcState, subject: &RepoSubject, ) -> Result { @@ -236,44 +236,117 @@ where } #[derive(Deserialize)] -pub(crate) struct DidParams { +pub(crate) struct RepoExportParams { pub(crate) did: RepoSubject, + #[serde(default)] + pub(crate) since: Option, } +const EXPORT_FRAME_WINDOW: usize = 5_000; + pub(crate) async fn sync_get_repo( State(state): State>>, - axum::extract::Query(params): axum::extract::Query, + ValidatedQuery(params): ValidatedQuery, ) -> Result { let found = resolve(&state, ¶ms.did)?; - serve_car(found, state.byte_limits.response.get(), |found, tip| { - let tree = chain::load_tree(&found.repo, tip); - let store = found.blocks(tip); - let root = tip.data; - futures::executor::block_on(async move { - let mut nodes = tree.collect_node_cids().await?; - nodes.push(root); - let leaves: Vec = tree - .leaves() - .await? - .into_iter() - .map(|(_, cid)| cid) - .collect(); - futures::stream::iter(nodes) - .map(Ok) - .and_then(|cid| block_or_corrupt(&store, cid)) - .chain( - futures::stream::iter(leaves) - .map(Ok) - .try_filter_map(|cid| block_if_present(&store, cid)), - ) - .try_collect::>() - .await - }) - .map_err(|error| XrpcError::internal(error.to_string())) - }) + let response_limit = state.byte_limits.response.get(); + serve_car( + found, + state.byte_limits.response.get(), + move |found, tip| match ¶ms.since { + Some(since) => frames_since_blocks(found, tip, since.tid().clone(), response_limit), + None => whole_tree_blocks(found, tip), + }, + ) .await } +fn whole_tree_blocks( + found: &AtprotoRepo, + tip: &ChainTip, +) -> Result, XrpcError> { + let tree = chain::load_tree(&found.repo, tip); + let store = found.blocks(tip); + let root = tip.data; + futures::executor::block_on(async move { + let mut nodes = tree.collect_node_cids().await?; + nodes.push(root); + let leaves: Vec = tree + .leaves() + .await? + .into_iter() + .map(|(_, cid)| cid) + .collect(); + futures::stream::iter(nodes) + .map(Ok) + .and_then(|cid| block_or_corrupt(&store, cid)) + .chain( + futures::stream::iter(leaves) + .map(Ok) + .try_filter_map(|cid| block_if_present(&store, cid)), + ) + .try_collect::>() + .await + }) + .map_err(|error: RepoError| XrpcError::internal(error.to_string())) +} + +fn frames_since_blocks( + found: &AtprotoRepo, + tip: &ChainTip, + since: jacquard_common::types::tid::Tid, + limit: usize, +) -> Result, XrpcError> { + let traversed = chain::frames_after( + &found.repo, + tip.git, + knot_types::UnixMicros::new(since.timestamp()), + EXPORT_FRAME_WINDOW, + ); + match traversed.end { + chain::TraversalEnd::Limited => { + return Err(XrpcError::invalid_request( + "`since` is older than this knot's export window. Please fetch the whole repo \ + without `since` instead, thanks.", + )); + } + chain::TraversalEnd::Broken => { + return Err(XrpcError::internal( + "The chain wouldn't traverse back to `since`, so the export diff is \ + unavailable. Please fetch the whole repo instead.", + )); + } + chain::TraversalEnd::Exhausted => (), + } + let mut blocks = BTreeMap::new(); + traversed.found.iter().try_for_each(|(rev, frame)| { + let (_header, message) = crate::firehose::commit_message(frame).ok_or_else(|| { + XrpcError::internal(format!( + "A stored frame at {rev} wouldn't parse as a commit." + )) + })?; + let car = knot_record::car_blocks(&message.blocks).map_err(|error| { + XrpcError::internal(format!( + "A stored frame wouldn't read back as a CAR: {error}" + )) + })?; + blocks.extend(car.blocks); + if knot_record::chain::blocks_bytes(&blocks) > limit { + return Err(XrpcError::request_too_large( + "This response exceeds the configured maximum size. Please narrow the request, \ + thanks.", + )); + } + Ok::<_, XrpcError>(()) + })?; + Ok(blocks) +} + +#[derive(Deserialize)] +pub(crate) struct DidParams { + pub(crate) did: RepoSubject, +} + pub(crate) async fn sync_get_latest_commit( State(state): State>>, axum::extract::Query(params): axum::extract::Query, @@ -293,9 +366,9 @@ pub(crate) async fn sync_get_latest_commit( #[derive(Deserialize)] pub(crate) struct SyncRecordParams { - did: RepoSubject, - collection: RecordCollection, - rkey: UriRkey, + pub(crate) did: RepoSubject, + pub(crate) collection: RecordCollection, + pub(crate) rkey: UriRkey, } pub(crate) async fn sync_get_record( diff --git a/knot2/crates/knot-xrpc/src/blobs.rs b/knot2/crates/knot-xrpc/src/blobs.rs new file mode 100644 index 000000000..cb298ef98 --- /dev/null +++ b/knot2/crates/knot-xrpc/src/blobs.rs @@ -0,0 +1,259 @@ +use std::collections::{BTreeMap, HashMap, HashSet}; +use std::sync::{Arc, Mutex}; + +use axum::body::Bytes; +use axum::extract::State; +use axum::http::{HeaderName, HeaderValue, StatusCode, header}; +use axum::response::{IntoResponse, Response}; +use http::HeaderMap; +use ipld_core::ipld::Ipld; +use jacquard_common::types::tid::Tid; +use knot_git::{Layout, Repo}; +use knot_lfs::{ClaimedSize, LfsError, LfsOid, LfsStore, Pool, RootsDenied, UploadAdmission}; +use knot_record::attachment::{ + Attachment, AttachmentMime, AttachmentSize, AttachmentTooBig, OCTET_STREAM, SNIFF_BYTES, +}; +use knot_record::chain; +use knot_runtime::{Clock, HttpTransport}; +use knot_types::{BlobCid, Decision, RecordAddress, RepoDid, Sha256Digest}; +use serde::Deserialize; +use serde_json::json; + +use crate::XrpcState; +use crate::error::XrpcError; +use crate::query::{Limit, Offset, Since, ValidatedQuery}; +use crate::run_blocking; + +pub(crate) const UPLOAD_BLOB_ROUTE: &str = "/xrpc/com.atproto.repo.uploadBlob"; +pub(crate) const GET_BLOB_ROUTE: &str = "/xrpc/com.atproto.sync.getBlob"; +pub(crate) const LIST_BLOBS_ROUTE: &str = "/xrpc/com.atproto.sync.listBlobs"; + +const BLOB_CSP: &str = "default-src 'none'; sandbox"; +const INDEX_FRAME_WINDOW: usize = 50_000; +const LIST_BLOBS_DEFAULT: usize = 500; +const LIST_BLOBS_MAX: usize = 1000; + +pub(crate) fn store_absent() -> XrpcError { + XrpcError::named( + StatusCode::SERVICE_UNAVAILABLE, + "BlobStoreUnavailable", + "This knot has no blob store; please set `lfs.store_path` to enable one, thanks.", + ) +} + +fn blob_not_found() -> XrpcError { + XrpcError::named( + StatusCode::NOT_FOUND, + "BlobNotFound", + "No such blob for this repo.", + ) +} + +struct UploadedAttachment { + digest: Sha256Digest, + mime: AttachmentMime, + size: AttachmentSize, +} + +impl UploadedAttachment { + fn parse(body: &Bytes, hint: &str) -> Result { + if body.is_empty() { + return Err(XrpcError::invalid_request( + "This upload's body is empty; please send the attachment bytes as the body.", + )); + } + let size = AttachmentSize::new(body.len() as u64).map_err( + |AttachmentTooBig { sent, max_bytes }| { + XrpcError::request_too_large(format!( + "An attachment is at most {max_bytes} bytes, and this upload sends {sent}." + )) + }, + )?; + Ok(Self { + digest: Sha256Digest::hash(body), + mime: AttachmentMime::detected(body, hint), + size, + }) + } + + fn attachment(&self) -> Attachment { + Attachment::new( + BlobCid::from_digest(self.digest), + self.mime.clone(), + self.size, + ) + } + + fn oid(&self) -> LfsOid { + LfsOid::from_digest(self.digest) + } + + fn claimed(&self) -> ClaimedSize { + ClaimedSize::new(self.size.get()) + } +} + +fn sniffed_kind(store: &knot_lfs::DiskStore, repo: &RepoDid, oid: &LfsOid) -> Option<&'static str> { + let mut body = store.read(Pool::Attachments, repo, oid).ok()?; + let mut prefix = Vec::new(); + let mut limited = std::io::Read::take(&mut body, SNIFF_BYTES as u64); + std::io::Read::read_to_end(&mut limited, &mut prefix).ok()?; + knot_record::attachment::sniffed(&prefix) +} + +pub(crate) async fn upload_blob( + State(state): State>>, + headers: HeaderMap, + method: crate::Method, + body: Bytes, +) -> Result { + let repo = state.audience_repo(&headers).await?; + let session = crate::social::session(&state, &headers, &method, repo).await?; + let contributor = session.require_contribution("uploads")?; + let lfs = state.lfs.as_ref().ok_or_else(store_absent)?; + let hint = headers + .get(header::CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + .unwrap_or(""); + let uploaded = UploadedAttachment::parse(&body, hint)?; + let claimed = uploaded.claimed(); + let permit = lfs + .handle + .admission + .admit(claimed) + .map_err(|error| match error { + LfsError::SizeLimitExceeded { declared, limit } => { + XrpcError::request_too_large(format!( + "This knot's LFS object limit is {limit} bytes, and this attachment is \ + {declared}. Please try a smaller upload, thanks." + )) + } + LfsError::FreeSpaceDenied { free, floor } => XrpcError::named( + StatusCode::SERVICE_UNAVAILABLE, + "BlobStoreFull", + format!( + "The blob store's disk is down to {free} bytes, under the {floor}-byte \ + floor. Please try again once there's room, thanks." + ), + ), + other => XrpcError::internal(other.to_string()), + })?; + let oid = uploaded.oid(); + let stored = Arc::clone(&lfs.handle.store); + let upload_repo = contributor.session.repo.clone(); + run_blocking(move || { + let outcome = stored.put( + Pool::Attachments, + &upload_repo, + &oid, + claimed, + &mut &body[..], + ); + drop(permit); + outcome.map_err(|error| XrpcError::internal(error.to_string())) + }) + .await?; + let response = json!({ + "blob": serde_json::to_value(uploaded.attachment()).expect("an attachment serializes") + }); + crate::reads::json(response, state.byte_limits.response.get()) +} + +#[derive(Deserialize)] +pub(crate) struct GetBlobParams { + did: crate::atproto::RepoSubject, + cid: BlobCid, +} + +pub(crate) async fn get_blob( + State(state): State>>, + ValidatedQuery(params): ValidatedQuery, +) -> Result { + let found = crate::atproto::resolve(&state, ¶ms.did)?; + let lfs = state.lfs.as_ref().ok_or_else(store_absent)?; + let oid = LfsOid::from_digest(params.cid.digest()); + let stored = Arc::clone(&lfs.handle.store); + let repo = found.did.clone(); + let bytes = run_blocking(move || { + let read = stored + .read(Pool::Attachments, &repo, &oid) + .map_err(|error| match error { + LfsError::NotFound { .. } => blob_not_found(), + other => XrpcError::internal(other.to_string()), + })?; + let mut body = read; + let mut bytes = Vec::new(); + std::io::Read::read_to_end(&mut body, &mut bytes) + .map_err(|error| XrpcError::internal(error.to_string()))?; + Ok::<_, XrpcError>(bytes) + }) + .await?; + let mime = knot_record::attachment::sniffed(&bytes).unwrap_or(OCTET_STREAM); + Ok(( + StatusCode::OK, + [ + ( + header::CONTENT_TYPE, + HeaderValue::from_str(mime).expect("a sniffed mime is a header value"), + ), + (header::CONTENT_LENGTH, HeaderValue::from(bytes.len())), + ( + HeaderName::from_static("x-content-type-options"), + HeaderValue::from_static("nosniff"), + ), + ( + HeaderName::from_static("content-security-policy"), + HeaderValue::from_static(BLOB_CSP), + ), + ], + bytes, + ) + .into_response()) +} + +#[derive(Deserialize)] +pub(crate) struct ListBlobsParams { + did: crate::atproto::RepoSubject, + #[serde(default)] + since: Option, + #[serde(default)] + limit: Limit, + #[serde(default)] + cursor: Offset, +} + +pub(crate) async fn list_blobs( + State(state): State>>, + ValidatedQuery(params): ValidatedQuery, +) -> Result { + let found = crate::atproto::resolve(&state, ¶ms.did)?; + let index = index_of(&state, &found.did).await?; + match index.coverage { + Coverage::Whole => (), + Coverage::Short(shortness) => return Err(index_refused(shortness)), + } + let eligible: Vec<&(Tid, LfsOid)> = index + .listed + .iter() + .filter(|(rev, _)| { + params + .since + .as_ref() + .is_none_or(|since| rev.newer_than(since.tid())) + }) + .collect(); + let total = eligible.len(); + let taken = params.cursor.get().saturating_add(params.limit.get()); + let page: Vec = eligible + .into_iter() + .skip(params.cursor.get()) + .take(params.limit.get()) + .map(|(_, oid)| BlobCid::from_digest(oid.digest()).to_string()) + .collect(); + let response = json!({ + "cids": page, + "cursor": (taken < total).then(|| taken.to_string()), + }); + crate::reads::json(response, state.byte_limits.response.get()) +} + diff --git a/knot2/crates/knot-xrpc/src/firehose.rs b/knot2/crates/knot-xrpc/src/firehose.rs index a70bde48c..21c50e803 100644 --- a/knot2/crates/knot-xrpc/src/firehose.rs +++ b/knot2/crates/knot-xrpc/src/firehose.rs @@ -297,7 +297,7 @@ impl InfoCode { } } #[derive(Serialize, Deserialize)] -struct Header { +pub(crate) struct Header { op: i64, t: FrameKind, } @@ -320,7 +320,7 @@ struct FrameInfo<'a> { } #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -enum RepoAction { +pub(crate) enum RepoAction { #[serde(rename = "create")] Create, #[serde(rename = "update")] @@ -330,26 +330,26 @@ enum RepoAction { } #[derive(Serialize, Deserialize)] -struct RepoOp { - action: RepoAction, - path: String, - cid: Option, +pub(crate) struct RepoOp { + pub(crate) action: RepoAction, + pub(crate) path: String, + pub(crate) cid: Option, #[serde(skip_serializing_if = "Option::is_none")] prev: Option, } #[derive(Serialize, Deserialize)] -struct CommitMessage { - seq: UnixMicros, +pub(crate) struct CommitMessage { + pub(crate) seq: UnixMicros, rebase: bool, #[serde(rename = "tooBig")] too_big: bool, - repo: RepoDid, - commit: Cid, - rev: Tid, + pub(crate) repo: RepoDid, + pub(crate) commit: Cid, + pub(crate) rev: Tid, since: Option, #[serde(with = "serde_bytes")] - blocks: Vec, - ops: Vec, + pub(crate) blocks: Vec, + pub(crate) ops: Vec, blobs: Vec, #[serde(rename = "prevData", skip_serializing_if = "Option::is_none")] prev_data: Option, @@ -444,8 +444,7 @@ fn glued(header: &H, payload: &P) -> B fn frame(t: FrameKind, payload: &impl Serialize) -> Bytes { glued(&Header { op: 1, t }, payload) } -#[cfg(test)] -fn commit_message(frame: &[u8]) -> Option<(Header, CommitMessage)> { +pub(crate) fn commit_message(frame: &[u8]) -> Option<(Header, CommitMessage)> { let mut read = frame; let header: Header = serde_ipld_dagcbor::de::from_reader_once(&mut read).ok()?; if header.t != FrameKind::Commit { @@ -671,7 +670,7 @@ fn index_repo(feed: &mut ReplayFeed, slot: usize, repo: &Repo, cursor: UnixMicro if UnixMicros::new(tip.rev.timestamp()) > cursor { let budget = cap.saturating_sub(feed.index.len()); let walked = chain::framed_revs(repo, tip.git, cursor, budget); - feed.capped |= walked.hit_limit; + feed.capped |= walked.end == chain::TraversalEnd::Limited; feed.index .extend(walked.found.into_iter().map(|(micros, oid)| { IndexEntry::Chain(FrameAt { @@ -1565,7 +1564,10 @@ mod tests { }; assert_eq!(sig.len(), 64, "secp256k1 signatures are 64 compact bytes"); assert!( - matches!(payload.remove("seq"), Some(ipld_core::ipld::Ipld::Integer(_))), + matches!( + payload.remove("seq"), + Some(ipld_core::ipld::Ipld::Integer(_)) + ), "a push event still carries seq on the wire, just outside the signature" ); assert!( @@ -1630,4 +1632,3 @@ mod tests { ); } } - diff --git a/knot2/crates/knot-xrpc/src/query.rs b/knot2/crates/knot-xrpc/src/query.rs index da2453fee..d2bc12bc0 100644 --- a/knot2/crates/knot-xrpc/src/query.rs +++ b/knot2/crates/knot-xrpc/src/query.rs @@ -74,6 +74,23 @@ impl<'de> Deserialize<'de> for Offset { } } +pub(crate) struct Since(jacquard_common::types::tid::Tid); + +impl Since { + pub(crate) fn tid(&self) -> &jacquard_common::types::tid::Tid { + &self.0 + } +} + +impl<'de> Deserialize<'de> for Since { + fn deserialize>(deserializer: D) -> Result { + let raw = String::deserialize(deserializer)?; + jacquard_common::types::tid::Tid::new(&raw) + .map(Self) + .map_err(|_| de::Error::custom("`since` must be a TID, the rev of an export.")) + } +} + knot_types::scalar_newtype! { pub(crate) struct Total(usize); } @@ -182,7 +199,9 @@ impl<'de> Deserialize<'de> for Revspec { let raw = String::deserialize(deserializer)?; match raw.len() <= MAX_REVSPEC_BYTES && !raw.chars().any(char::is_control) { true => Ok(Self(raw)), - false => Err(de::Error::custom("invalid revision")), + false => Err(de::Error::custom( + "`rev` must be a short plain ref. Please send a TID or a branch name, thanks.", + )), } } } diff --git a/lexicons/issue/issue.json b/lexicons/issue/issue.json index e18006499..dff640701 100644 --- a/lexicons/issue/issue.json +++ b/lexicons/issue/issue.json @@ -43,8 +43,7 @@ "type": "array", "items": { "type": "blob", - "accept": ["image/*"], - "maxSize": 1000000 + "maxSize": 200000000 } } } diff --git a/lexicons/markup/markdown.json b/lexicons/markup/markdown.json index fcf86963e..e09fce027 100644 --- a/lexicons/markup/markdown.json +++ b/lexicons/markup/markdown.json @@ -19,8 +19,7 @@ "type": "array", "items": { "type": "blob", - "accept": ["image/*"], - "maxSize": 1000000 + "maxSize": 200000000 }, "description": "list of blobs referenced in markdown" } diff --git a/lexicons/pulls/pull.json b/lexicons/pulls/pull.json index 9ec13a178..4f1b69aac 100644 --- a/lexicons/pulls/pull.json +++ b/lexicons/pulls/pull.json @@ -72,8 +72,7 @@ "type": "array", "items": { "type": "blob", - "accept": ["image/*"], - "maxSize": 1000000 + "maxSize": 200000000 } } } diff --git a/localinfra/Corefile b/localinfra/Corefile index 9ec135305..d594f4c4f 100644 --- a/localinfra/Corefile +++ b/localinfra/Corefile @@ -10,5 +10,6 @@ tngl.boltless.dev:53 { .:53 { errors + rewrite stop name regex ^([a-z0-9-]+)$ {1}.dns.podman forward . /etc/resolv.conf } diff --git a/localinfra/knot2.Dockerfile b/localinfra/knot2.Dockerfile index 4f3136838..f67d6c66a 100644 --- a/localinfra/knot2.Dockerfile +++ b/localinfra/knot2.Dockerfile @@ -43,6 +43,7 @@ fi # the knot creates its host key and sealed store (parent dirs and all) on # demand, but scan_path has to already exist before it will start mkdir -p "${KNOT_SCAN_PATH:-/home/git/repositories}" +[ -n "${KNOT_LFS_STORE_PATH:-}" ] && mkdir -p "${KNOT_LFS_STORE_PATH}" exec /usr/local/bin/knot-server EOF