use rabs_protocol::authority::{ClusterId, CoordinatorAuthority, CoordinatorIncarnationId};
use rabs_protocol::generation::{
ActionGeneration, ActionGenerationId, AttemptAuthority, AttemptId, ExecutionLeaseId,
LeaseRenewalSeq, WorkerBootGeneration, WorkerIncarnationId,
};
use rabs_protocol::raw_bytes::RawBytes;
use rabs_protocol::result_identity::{
AttemptEvidenceBundle, CanonicalActionResultManifest, DigestAlgorithm, LogicalOutput, ObjectId,
OutputRole, ResultKind, TypedDigest,
};
use rabs_protocol::wire_time::PeerId;
use rabs_protocol::worker_fence::WorkerSessionOffer;
use crate::metadata_store::{ActionEntryRow, AuthorityRow, RabsMetadataStore};
use crate::publication::{
OBSERVABLE_PROJECTION_DOMAIN, OfferPreparedActionResult, ProvisionalAncestorRef,
SEMANTIC_PROJECTION_DOMAIN, authority_digest,
};
const ACTION_KEY_DOMAIN: &str = "rabs.action-key.sha256.v1";
#[must_use]
pub fn tagged_digest(domain: &'static str, tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: [tag; 32],
}
}
#[must_use]
pub fn tagged_object(tag: u8) -> ObjectId {
ObjectId(tagged_digest("rabs.object.sha256.v1", tag))
}
#[must_use]
pub fn sample_action_key() -> TypedDigest {
tagged_digest(ACTION_KEY_DOMAIN, 7)
}
#[must_use]
pub fn sample_coordinator_authority() -> CoordinatorAuthority {
CoordinatorAuthority {
cluster_id: ClusterId("cluster-a".to_owned()),
credential_generation: 1,
term: 3,
incarnation_id: CoordinatorIncarnationId(77),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct FixtureAttemptIds {
pub generation: u128,
pub attempt: u128,
pub lease: u128,
}
pub const SAMPLE_ATTEMPT_IDS: FixtureAttemptIds = FixtureAttemptIds {
generation: 11,
attempt: 20,
lease: 30,
};
#[must_use]
pub fn sample_attempt_authority() -> AttemptAuthority {
attempt_authority_for(&sample_coordinator_authority())
}
#[must_use]
pub fn attempt_authority_for(coordinator: &CoordinatorAuthority) -> AttemptAuthority {
attempt_authority_with_ids(coordinator, SAMPLE_ATTEMPT_IDS)
}
#[must_use]
fn attempt_authority_with_ids(
coordinator: &CoordinatorAuthority,
ids: FixtureAttemptIds,
) -> AttemptAuthority {
let coordinator = coordinator.clone();
let created_under = authority_digest(&coordinator);
AttemptAuthority {
coordinator,
action_key: sample_action_key(),
action_generation: ActionGeneration {
generation_id: ActionGenerationId(ids.generation),
per_key_ordinal: 1,
created_under_authority_digest: created_under,
},
attempt_id: AttemptId(ids.attempt),
execution_lease_id: ExecutionLeaseId(ids.lease),
lease_renewal_seq: LeaseRenewalSeq(1),
worker_peer_id: PeerId("worker-a".to_owned()),
worker_boot_generation: WorkerBootGeneration(1),
worker_incarnation_id: WorkerIncarnationId(5),
}
}
#[must_use]
pub fn sample_manifest() -> CanonicalActionResultManifest {
manifest_with_output(41)
}
#[must_use]
pub fn manifest_with_output(output_tag: u8) -> CanonicalActionResultManifest {
manifest_with_output_object(&tagged_object(output_tag))
}
#[must_use]
pub fn manifest_with_output_object(object: &ObjectId) -> CanonicalActionResultManifest {
let object = object.clone();
CanonicalActionResultManifest {
action_key: sample_action_key(),
canonical_descriptor_digest: tagged_digest("rabs.descriptor.sha256.v1", 8),
key_epoch: 1,
projection_epoch: 1,
result_kind: ResultKind::Success,
artifact_bundle_root: Some(tagged_object(40)),
logical_outputs: vec![LogicalOutput {
role: OutputRole::Materializable,
virtual_path: RawBytes::new(b"out/lib.rlib".to_vec()),
object,
}],
semantic_result_digest: tagged_digest(SEMANTIC_PROJECTION_DOMAIN, 0),
observable_result_digest: tagged_digest(OBSERVABLE_PROJECTION_DOMAIN, 0),
}
}
#[must_use]
pub fn sample_evidence(manifest_id: &ObjectId) -> AttemptEvidenceBundle {
AttemptEvidenceBundle {
action_key: sample_action_key(),
canonical_result_manifest_id: manifest_id.clone(),
execution_snapshot_root: tagged_object(60),
observed_input_report: tagged_object(61),
raw_process_and_event_evidence: tagged_object(62),
provenance_receipt: tagged_object(63),
incremental_snapshot: None,
}
}
#[must_use]
pub fn sample_declared() -> Vec<(OutputRole, RawBytes)> {
vec![(
OutputRole::Materializable,
RawBytes::new(b"out/lib.rlib".to_vec()),
)]
}
#[must_use]
pub fn sample_expected_descriptor() -> TypedDigest {
tagged_digest("rabs.descriptor.sha256.v1", 8)
}
#[must_use]
pub fn sample_offer() -> OfferPreparedActionResult {
sample_offer_with_ancestors(Vec::new())
}
#[must_use]
pub fn sample_offer_with_ancestors(
ancestors: Vec<ProvisionalAncestorRef>,
) -> OfferPreparedActionResult {
offer_for(&sample_coordinator_authority(), 41, 50, 51, ancestors)
}
#[must_use]
pub fn offer_under(coordinator: &CoordinatorAuthority) -> OfferPreparedActionResult {
offer_for(coordinator, 41, 50, 51, Vec::new())
}
#[must_use]
pub fn divergent_offer_under(coordinator: &CoordinatorAuthority) -> OfferPreparedActionResult {
offer_for(coordinator, 42, 52, 53, Vec::new())
}
#[must_use]
pub fn offer_with_manifest_bytes(
coordinator: &CoordinatorAuthority,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes(coordinator, &tagged_object(41), 51)
}
#[must_use]
pub fn offer_serving_object(
coordinator: &CoordinatorAuthority,
output: &ObjectId,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes(coordinator, output, 51)
}
#[must_use]
pub fn divergent_offer_with_manifest_bytes(
coordinator: &CoordinatorAuthority,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes(coordinator, &tagged_object(42), 53)
}
fn offer_with_bytes(
coordinator: &CoordinatorAuthority,
output: &ObjectId,
evidence_tag: u8,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes_with_ids(coordinator, output, evidence_tag, SAMPLE_ATTEMPT_IDS)
}
fn offer_with_bytes_with_ids(
coordinator: &CoordinatorAuthority,
output: &ObjectId,
evidence_tag: u8,
ids: FixtureAttemptIds,
) -> (OfferPreparedActionResult, Vec<u8>) {
let stamped = offer_for_object(coordinator, output, 50, evidence_tag, Vec::new(), ids);
let bytes = crate::manifest_codec::encode_manifest_v1(&stamped.manifest);
let manifest_id = ObjectId(
crate::digest_set::digest_set(&bytes, crate::digest_set::DigestRequest::default(), None)
.expect("digest the manifest bytes")
.atp_content_id,
);
let offer = OfferPreparedActionResult::build(
attempt_authority_with_ids(coordinator, ids),
stamped.manifest,
manifest_id.clone(),
sample_evidence(&manifest_id),
tagged_object(evidence_tag),
tagged_digest("rabs.observation-stream.sha256.v1", 9),
&sample_declared(),
Vec::new(),
)
.expect("sample offer fixtures are internally consistent");
assert_eq!(
crate::manifest_codec::encode_manifest_v1(&offer.manifest),
bytes,
"re-stamping must not change the manifest bytes its id was taken from"
);
(offer, bytes)
}
fn offer_for(
coordinator: &CoordinatorAuthority,
output_tag: u8,
manifest_tag: u8,
evidence_tag: u8,
ancestors: Vec<ProvisionalAncestorRef>,
) -> OfferPreparedActionResult {
offer_for_object(
coordinator,
&tagged_object(output_tag),
manifest_tag,
evidence_tag,
ancestors,
SAMPLE_ATTEMPT_IDS,
)
}
#[allow(clippy::too_many_arguments)]
fn offer_for_object(
coordinator: &CoordinatorAuthority,
output: &ObjectId,
manifest_tag: u8,
evidence_tag: u8,
ancestors: Vec<ProvisionalAncestorRef>,
ids: FixtureAttemptIds,
) -> OfferPreparedActionResult {
let manifest_id = tagged_object(manifest_tag);
OfferPreparedActionResult::build(
attempt_authority_with_ids(coordinator, ids),
manifest_with_output_object(output),
manifest_id.clone(),
sample_evidence(&manifest_id),
tagged_object(evidence_tag),
tagged_digest("rabs.observation-stream.sha256.v1", 9),
&sample_declared(),
ancestors,
)
.expect("sample offer fixtures are internally consistent")
}
pub fn install_ready_store(store: &mut dyn RabsMetadataStore) {
let coordinator = sample_coordinator_authority();
store
.acquire_authority(&AuthorityRow {
digest: authority_digest(&coordinator),
cluster_id: "cluster-a".to_owned(),
incarnation: 77,
term: 3,
acquired_seq: 1,
})
.expect("acquire authority");
install_admission_world(store, &coordinator);
install_offer_closure(store, &offer_under(&coordinator));
}
pub fn install_admission_world(
store: &mut dyn RabsMetadataStore,
coordinator: &CoordinatorAuthority,
) {
install_admission_world_with_ids(store, coordinator, SAMPLE_ATTEMPT_IDS);
}
pub fn install_admission_world_with_ids(
store: &mut dyn RabsMetadataStore,
coordinator: &CoordinatorAuthority,
ids: FixtureAttemptIds,
) {
let auth = authority_digest(coordinator);
let attempt_authority = attempt_authority_with_ids(coordinator, ids);
store
.upsert_action_entry(&ActionEntryRow {
action_key: sample_action_key(),
key_epoch: 1,
projection_epoch: 1,
})
.expect("upsert action entry");
store
.create_bound_generation(
&auth,
&attempt_authority.action_generation,
&sample_action_key(),
)
.expect("create generation");
store
.admit_worker_session(
&auth,
&WorkerSessionOffer {
worker_peer_id: attempt_authority.worker_peer_id.clone(),
boot_generation: attempt_authority.worker_boot_generation,
incarnation: attempt_authority.worker_incarnation_id,
reenrollment_proof: None,
},
1,
)
.expect("worker session");
store
.admit_attempt_lease(&attempt_authority, 1, 60_000)
.expect("attempt lease");
}
#[must_use]
pub fn offer_with_manifest_bytes_with_ids(
coordinator: &CoordinatorAuthority,
ids: FixtureAttemptIds,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes_with_ids(coordinator, &tagged_object(41), 51, ids)
}
#[must_use]
pub fn divergent_offer_with_manifest_bytes_with_ids(
coordinator: &CoordinatorAuthority,
ids: FixtureAttemptIds,
) -> (OfferPreparedActionResult, Vec<u8>) {
offer_with_bytes_with_ids(coordinator, &tagged_object(42), 53, ids)
}
pub fn install_offer_closure(store: &mut dyn RabsMetadataStore, offer: &OfferPreparedActionResult) {
let mut closure = vec![
offer.manifest_id.0.clone(),
offer.evidence_id.0.clone(),
offer.evidence.execution_snapshot_root.0.clone(),
offer.evidence.observed_input_report.0.clone(),
offer.evidence.raw_process_and_event_evidence.0.clone(),
offer.evidence.provenance_receipt.0.clone(),
];
if let Some(root) = &offer.manifest.artifact_bundle_root {
closure.push(root.0.clone());
}
for output in &offer.manifest.logical_outputs {
closure.push(output.object.0.clone());
}
if let Some(snapshot) = &offer.evidence.incremental_snapshot {
closure.push(snapshot.0.clone());
}
for id in &closure {
if store
.object_durably_located(id)
.expect("check object location")
{
continue;
}
let key = crate::metadata_store::digest_key(id);
store.record_object(id, 64).expect("record object");
store
.add_location(id, &format!("/cas/{key}"), Some(1), "raw", true)
.expect("locate object");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{RusqliteEngine, SqlMetadataStore};
use crate::publication::{CommitDurabilityProfile, PublicationOutcome, process_offer};
#[test]
fn sample_offer_commits_into_a_ready_store() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
install_ready_store(&mut store);
let outcome = process_offer(
&mut store,
&sample_offer(),
&sample_expected_descriptor(),
|_| None, 900, 1, CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.expect("offer must be accepted");
assert!(
matches!(outcome, PublicationOutcome::Committed(_)),
"expected a real commit, got {outcome:?}"
);
}
}