Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 1 addition & 19 deletions crates/mirror_worker/src/add_entries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1024,7 +1024,7 @@ fn resolve_target_pending(
return Ok(PendingCheckpoint {
size: committed.size,
hash: committed.hash,
signed_note_bytes: committed_checkpoint_note(committed, verifiers)?,
signed_note_bytes: committed.checkpoint_note_bytes.clone(),
witness_published: true,
witness_response_bytes: Vec::new(),
update_request_hash: Hash::default(),
Expand Down Expand Up @@ -1095,24 +1095,6 @@ fn resolve_target_pending(
})
}

fn committed_checkpoint_note(
committed: &CommittedCheckpoint,
verifiers: &signed_note::VerifierList,
) -> std::result::Result<Vec<u8>, &'static str> {
if !committed.checkpoint_note_bytes.is_empty() {
return Ok(committed.checkpoint_note_bytes.clone());
}

let note = Note::from_bytes(&committed.signed_note_bytes)
.map_err(|_| "committed checkpoint is not a valid signed note")?;
let (log_signatures, _) = note
.verify(verifiers)
.map_err(|_| "committed checkpoint has no valid trusted log signature")?;
Note::new(note.text(), &log_signatures)
.map(|note| note.to_bytes())
.map_err(|_| "committed checkpoint note reconstruction failed")
}

/// Verify a single [`EntryPackage`] against the target pending
/// checkpoint at `upload_end`.
///
Expand Down
43 changes: 3 additions & 40 deletions crates/mirror_worker/src/mirror_state_do.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,14 +67,12 @@ pub struct PendingCheckpoint {
#[serde_as(as = "Base64As")]
pub signed_note_bytes: Vec<u8>,
/// Whether the witness checkpoint has been published to R2.
#[serde(default)]
pub witness_published: bool,
/// Serialized witness response for idempotent request retries.
#[serde_as(as = "Base64As")]
#[serde(default)]
pub witness_response_bytes: Vec<u8>,
/// Fingerprint of the accepted update request.
#[serde(default, with = "generic_log_worker::hash_serde::hex")]
#[serde(with = "generic_log_worker::hash_serde::hex")]
pub update_request_hash: Hash,
}

Expand All @@ -96,7 +94,6 @@ pub struct CommittedCheckpoint {
pub hash: Hash,
/// The source log's signed checkpoint note.
#[serde_as(as = "Base64As")]
#[serde(default)]
pub checkpoint_note_bytes: Vec<u8>,
/// The served checkpoint bytes: the log's signed note with the
/// mirror's cosignature line(s) appended, exactly as written to R2.
Expand Down Expand Up @@ -273,19 +270,12 @@ impl MirrorState {
&& let Some(checkpoint) = current.as_mut()
&& !checkpoint.witness_published
{
self.recover_witness_checkpoint(checkpoint).await?;
self.publish_witness_checkpoint(checkpoint).await?;
}
if CONFIG.witness_enabled()
&& let Some(checkpoint) = current.as_ref()
&& !checkpoint.witness_response_bytes.is_empty()
&& (checkpoint.update_request_hash == request_hash
|| (checkpoint.update_request_hash == Hash::default()
&& checkpoint.size == body.new_size
&& checkpoint.hash == body.new_hash
&& checkpoint_text_matches(
&checkpoint.signed_note_bytes,
&body.signed_note_bytes,
)?))
&& checkpoint.update_request_hash == request_hash
{
return Response::from_json(&UpdatePendingResponse {
witness_response_bytes: checkpoint.witness_response_bytes.clone(),
Expand Down Expand Up @@ -369,16 +359,6 @@ impl MirrorState {
Ok((note.to_bytes(), response))
}

async fn recover_witness_checkpoint(&self, checkpoint: &mut PendingCheckpoint) -> Result<()> {
if checkpoint.witness_response_bytes.is_empty() {
let (note, response) = self.cosign_witness_checkpoint(&checkpoint.signed_note_bytes)?;
checkpoint.signed_note_bytes = note;
checkpoint.witness_response_bytes = response;
self.state.storage().put(PENDING_KEY, &*checkpoint).await?;
}
self.publish_witness_checkpoint(checkpoint).await
}

async fn publish_witness_checkpoint(&self, checkpoint: &mut PendingCheckpoint) -> Result<()> {
let origin = self
.state
Expand Down Expand Up @@ -509,14 +489,6 @@ fn update_request_hash(body: &UpdatePendingRequest) -> Hash {
Hash(digest.finalize().into())
}

fn checkpoint_text_matches(left: &[u8], right: &[u8]) -> Result<bool> {
let left = Note::from_bytes(left)
.map_err(|error| Error::from(format!("parse persisted checkpoint note: {error:?}")))?;
let right = Note::from_bytes(right)
.map_err(|error| Error::from(format!("parse submitted checkpoint note: {error:?}")))?;
Ok(left.text() == right.text())
}

trait PublicationStorage {
async fn persist_committed(&self, checkpoint: &CommittedCheckpoint) -> Result<()>;
async fn clear_intent(&self) -> Result<()>;
Expand Down Expand Up @@ -682,15 +654,6 @@ mod tests {
assert_eq!(decoded.signed_note_bytes, b"signed-note-bytes");
}

#[test]
fn committed_checkpoint_accepts_legacy_json() {
let json = r#"{"size":42,"hash":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","signed_note_bytes":"c2VydmVk"}"#;
let decoded: CommittedCheckpoint = serde_json::from_str(json).unwrap();
assert_eq!(decoded.size, 42);
assert!(decoded.checkpoint_note_bytes.is_empty());
assert_eq!(decoded.signed_note_bytes, b"served");
}

#[test]
fn publication_intent_json_format() {
let intent = PublicationIntent(CommittedCheckpoint {
Expand Down
Loading