use crate::conformance::MemoryConformanceFailure;
use crate::ports::{MemoryReaderPort, MemoryWriteOutcome, MemoryWriterPort};
use crate::value_objects::{
Attributes, MemoryConfidence, MemoryEntryKind, MemoryEvidence, MemoryMoment, MemoryQuestion,
MemoryRelation, MemoryRelationKind, MemoryWrite,
};
mod support;
use support::{entry, expect_unsupported, moment, named, scope, write, Checked};
#[derive(Debug)]
pub struct MemoryConformance;
impl MemoryConformance {
pub async fn run(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Result<Vec<&'static str>, MemoryConformanceFailure> {
let mut passed = Vec::new();
Self::capabilities_are_stable(writer, reader)?;
passed.push("capabilities_are_stable");
Self::an_unwritten_scope_recalls_nothing(reader).await?;
passed.push("an_unwritten_scope_recalls_nothing");
Self::what_is_remembered_can_be_recalled(writer, reader).await?;
passed.push("what_is_remembered_can_be_recalled");
Self::scopes_do_not_bleed_into_each_other(writer, reader).await?;
passed.push("scopes_do_not_bleed_into_each_other");
Self::the_same_write_twice_is_one_memory(writer, reader).await?;
passed.push("the_same_write_twice_is_one_memory");
Self::reasons_survive_the_round_trip(writer, reader).await?;
passed.push("reasons_survive_the_round_trip");
Self::a_chain_of_reasons_can_be_followed(writer, reader).await?;
passed.push("a_chain_of_reasons_can_be_followed");
Self::evidence_survives_the_round_trip(writer, reader).await?;
passed.push("evidence_survives_the_round_trip");
Self::questions_are_answered_or_declined(writer, reader).await?;
passed.push("questions_are_answered_or_declined");
Self::time_travel_is_honoured_or_declined(writer, reader).await?;
passed.push("time_travel_is_honoured_or_declined");
Ok(passed)
}
async fn an_unwritten_scope_recalls_nothing(reader: &dyn MemoryReaderPort) -> Checked {
const PROPERTY: &str = "an_unwritten_scope_recalls_nothing";
let recollection = reader
.recall(&scope("unwritten"))
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !reader.capabilities().recalls() {
return expect_unsupported(PROPERTY, &recollection, "recall");
}
if recollection.entries().is_empty() {
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!(
"a scope nobody wrote to came back with {} entries",
recollection.entries().len()
),
))
}
}
async fn what_is_remembered_can_be_recalled(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "what_is_remembered_can_be_recalled";
let scope = scope("round-trip");
let outcome = writer
.remember(
&scope,
write(vec![entry(
"the rollback was rehearsed in March",
MemoryEntryKind::Observation,
moment(10),
)]),
"conformance:round-trip",
)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !writer.capabilities().remembers() {
return if outcome == MemoryWriteOutcome::NotRemembered {
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!("a backend that does not remember answered {outcome:?}"),
))
};
}
if outcome != MemoryWriteOutcome::Remembered {
return Err(MemoryConformanceFailure::new(
PROPERTY,
format!("a first write answered {outcome:?}"),
));
}
if !reader.capabilities().recalls() {
return Ok(());
}
let recalled = reader
.recall(&scope)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
match recalled.entries() {
[only] if only.summary() == "the rollback was rehearsed in March" => Ok(()),
[] => Err(MemoryConformanceFailure::new(
PROPERTY,
"what was written came back as nothing — a dropped write and an empty scope must not look the same",
)),
others => Err(MemoryConformanceFailure::new(
PROPERTY,
format!("expected exactly what was written, got {} entries", others.len()),
)),
}
}
async fn scopes_do_not_bleed_into_each_other(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "scopes_do_not_bleed_into_each_other";
if !writer.capabilities().remembers() || !reader.capabilities().recalls() {
return Ok(());
}
let mine = scope("mine");
let yours = scope("yours");
writer
.remember(
&mine,
write(vec![entry("mine", MemoryEntryKind::Decision, moment(1))]),
"conformance:mine",
)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
writer
.remember(
&yours,
write(vec![entry("yours", MemoryEntryKind::Decision, moment(2))]),
"conformance:yours",
)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
let recalled = reader
.recall(&mine)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if recalled.entries().iter().any(|e| e.summary() == "yours") {
return Err(MemoryConformanceFailure::new(
PROPERTY,
"one scope's memory surfaced in another's",
));
}
Ok(())
}
async fn the_same_write_twice_is_one_memory(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "the_same_write_twice_is_one_memory";
if !writer.capabilities().remembers() {
return Ok(());
}
let scope = scope("retried");
let entries = || {
write(vec![entry(
"decided once",
MemoryEntryKind::Decision,
moment(5),
)])
};
let first = writer
.remember(&scope, entries(), "conformance:retried")
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
let second = writer
.remember(&scope, entries(), "conformance:retried")
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if first != MemoryWriteOutcome::Remembered {
return Err(MemoryConformanceFailure::new(
PROPERTY,
format!("a first write answered {first:?}"),
));
}
if second != MemoryWriteOutcome::AlreadyRemembered {
return Err(MemoryConformanceFailure::new(
PROPERTY,
format!("the same write repeated answered {second:?}, not AlreadyRemembered"),
));
}
if !reader.capabilities().recalls() {
return Ok(());
}
let recalled = reader
.recall(&scope)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if recalled.entries().len() == 1 {
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!(
"one write made twice left {} entries",
recalled.entries().len()
),
))
}
}
async fn reasons_survive_the_round_trip(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "reasons_survive_the_round_trip";
const WHY: &str = "the queue growth is what made a rollback necessary";
let scope = scope("reasons");
let observation = named(
"conformance:observation",
"the queue was backing up",
MemoryEntryKind::Observation,
moment(40),
);
let decision = named(
"conformance:decision",
"roll back rather than restart",
MemoryEntryKind::Decision,
moment(50),
);
let because = MemoryRelation::new(
decision.id().clone(),
observation.id().clone(),
MemoryRelationKind::ChosenBecause,
WHY,
MemoryConfidence::High,
)
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
let explained = MemoryWrite::new(vec![observation, decision], vec![because])
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
let outcome = writer
.remember(&scope, explained, "conformance:reasons")
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !writer.capabilities().remembers() {
return if outcome == MemoryWriteOutcome::NotRemembered {
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!("a backend that does not remember answered {outcome:?}"),
))
};
}
if !writer.capabilities().keeps_reasons() || !reader.capabilities().recalls() {
return Ok(());
}
let recalled = reader
.recall(&scope)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
match recalled.relations() {
[only] if only.why() == WHY && only.kind() == MemoryRelationKind::ChosenBecause => {
if only.from().as_str() == "conformance:decision"
&& only.to().as_str() == "conformance:observation"
{
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!(
"a reason came back pointing {} -> {}",
only.from(),
only.to()
),
))
}
}
[] => Err(MemoryConformanceFailure::new(
PROPERTY,
"the entries came back and the reason between them did not — \
memory that can be read and not followed",
)),
others => Err(MemoryConformanceFailure::new(
PROPERTY,
format!("expected exactly the reason written, got {}", others.len()),
)),
}
}
async fn a_chain_of_reasons_can_be_followed(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "a_chain_of_reasons_can_be_followed";
let scope = scope("chain");
let observation = named(
"chain:observation",
"the queue was backing up",
MemoryEntryKind::Observation,
moment(60),
);
let decision = named(
"chain:decision",
"roll back rather than restart",
MemoryEntryKind::Decision,
moment(70),
);
let outcome = named(
"chain:outcome",
"the queue drained",
MemoryEntryKind::Outcome,
moment(80),
);
let chain = vec![
MemoryRelation::new(
decision.id().clone(),
observation.id().clone(),
MemoryRelationKind::ChosenBecause,
"the queue growth is what made a rollback necessary",
MemoryConfidence::High,
)
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?,
MemoryRelation::new(
outcome.id().clone(),
decision.id().clone(),
MemoryRelationKind::FollowsFrom,
"the queue drained because the rollback removed the bad revision",
MemoryConfidence::Medium,
)
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?,
];
let from = outcome.id().clone();
let to = observation.id().clone();
if writer.capabilities().remembers() {
let explained = MemoryWrite::new(vec![observation, decision, outcome], chain)
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
writer
.remember(&scope, explained, "conformance:chain")
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
}
let followed = reader
.follow(&scope, &from, &to)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !reader.capabilities().follows_reasons() {
return expect_unsupported(PROPERTY, &followed, "follow");
}
if !writer.capabilities().remembers() || !writer.capabilities().keeps_reasons() {
return Ok(());
}
let reached = followed
.relations()
.iter()
.any(|relation| relation.to() == &to);
if followed.relations().len() >= 2 && reached {
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
format!(
"following an outcome back to what it came from returned {} step(s) \
and {} the far end — a memory that cannot be walked",
followed.relations().len(),
if reached { "reached" } else { "never reached" }
),
))
}
}
async fn evidence_survives_the_round_trip(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "evidence_survives_the_round_trip";
if !writer.capabilities().keeps_evidence() {
return Ok(());
}
let scope = scope("evidence");
let evidenced = entry(
"the queue was empty at 03:20",
MemoryEntryKind::Observation,
moment(20),
)
.with_evidence(vec![MemoryEvidence::new(
"dead-letter count",
Some("dead-letter-queue".to_owned()),
Attributes::empty(),
)
.expect("evidence should be valid")]);
writer
.remember(&scope, write(vec![evidenced]), "conformance:evidence")
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !reader.capabilities().recalls() {
return Ok(());
}
let recalled = reader
.recall(&scope)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
match recalled.entries().first() {
Some(entry) if entry.evidence().len() == 1 => Ok(()),
Some(entry) => Err(MemoryConformanceFailure::new(
PROPERTY,
format!(
"a backend that keeps evidence returned an entry with {} of it",
entry.evidence().len()
),
)),
None => Err(MemoryConformanceFailure::new(
PROPERTY,
"the evidenced entry did not come back at all",
)),
}
}
async fn questions_are_answered_or_declined(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "questions_are_answered_or_declined";
let scope = scope("asked");
if writer.capabilities().remembers() {
writer
.remember(
&scope,
write(vec![entry(
"we restarted the ingester",
MemoryEntryKind::Outcome,
moment(30),
)]),
"conformance:asked",
)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
}
let question =
MemoryQuestion::new("did we restart the ingester?").expect("question should be valid");
let answer = reader
.ask(&scope, &question)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if reader.capabilities().answers_questions() {
Ok(())
} else {
expect_unsupported(PROPERTY, &answer, "ask")
}
}
async fn time_travel_is_honoured_or_declined(
writer: &dyn MemoryWriterPort,
reader: &dyn MemoryReaderPort,
) -> Checked {
const PROPERTY: &str = "time_travel_is_honoured_or_declined";
let scope = scope("as-known-at");
if writer.capabilities().remembers() {
writer
.remember(
&scope,
write(vec![
entry("known early", MemoryEntryKind::Observation, moment(100)),
entry("known later", MemoryEntryKind::Observation, moment(900)),
]),
"conformance:as-known-at",
)
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
}
let recalled = reader
.as_known_at(&scope, MemoryMoment::at(moment(500)))
.await
.map_err(|error| MemoryConformanceFailure::new(PROPERTY, error.to_string()))?;
if !reader.capabilities().travels_in_time() {
return expect_unsupported(PROPERTY, &recalled, "as_known_at");
}
if !writer.capabilities().remembers() {
return Ok(());
}
if recalled
.entries()
.iter()
.any(|e| e.summary() == "known later")
{
return Err(MemoryConformanceFailure::new(
PROPERTY,
"reading memory as of a moment returned something learned after it",
));
}
if recalled
.entries()
.iter()
.any(|e| e.summary() == "known early")
{
Ok(())
} else {
Err(MemoryConformanceFailure::new(
PROPERTY,
"reading memory as of a moment lost what was already known then",
))
}
}
}