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: 20 additions & 0 deletions crates/tracedecay-graph-db/src/generation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ pub use identity::{
#[cfg(test)]
pub(crate) use recovered::recovered_generation_digest_chunked;
pub(crate) use recovered::recovered_generation_digest_from_database;
pub(crate) use recovered::recovered_relation_lanes;
pub(crate) use replay::InlineOnlyGraphGenerationManifestProvider;
use replay::validate_sealed_replay;
pub use replay::{
Expand Down Expand Up @@ -157,6 +158,10 @@ pub struct GraphGenerationManifestIdentity {
/// it. Re-validated against `dependencies` on every read and invisible to
/// equality and clones.
digest_memo: DependencyClosureDigestMemo,
/// The namespace this identity's rows are stored under when another
/// projection sealed them: a sibling scope's base, read and proven
/// under this projection. `None` for rows stored under their own.
stored_namespace: Option<GraphNamespace>,
}

impl GraphGenerationManifestIdentity {
Expand All @@ -177,9 +182,20 @@ impl GraphGenerationManifestIdentity {
watermark,
dependencies,
digest_memo: DependencyClosureDigestMemo::default(),
stored_namespace: None,
}
}

/// This identity over rows stored under `physical`, the namespace its
/// generation was sealed under, possibly by another projection.
pub(crate) fn stored_under(mut self, physical: GraphNamespace) -> Result<Self, GraphDbError> {
self.stored_namespace = None;
if physical != self.physical_namespace()? {
self.stored_namespace = Some(physical);
}
Ok(self)
}

pub fn dependency_closure_digest(
&self,
check: &dyn Fn() -> Result<(), GraphDbError>,
Expand All @@ -188,6 +204,9 @@ impl GraphGenerationManifestIdentity {
}

pub(crate) fn physical_namespace(&self) -> Result<GraphNamespace, GraphDbError> {
if let Some(stored) = &self.stored_namespace {
return Ok(stored.clone());
}
physical_namespace(
&self.projection.namespace,
&self.projection.projection,
Expand Down Expand Up @@ -592,6 +611,7 @@ impl GraphGenerationManifest {
.digest_memo
.dependency_closure
.propagated(&self.dependencies),
stored_namespace: None,
}
}

Expand Down
46 changes: 43 additions & 3 deletions crates/tracedecay-graph-db/src/generation/recovered.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,9 +20,9 @@ use crate::{GraphDbError, GraphNamespace};

use super::{
CheckedDigestWriter, CheckedVecWriter, GraphEntityRef, GraphGenerationManifestIdentity,
GraphGenerationRelation, GraphProjectionIdentity, frame_length_headers,
physical_namespace_projection_map, recovered_entity_ref, write_canonical_row_frame,
write_generation_identity_frames,
GraphGenerationRelation, GraphProjectionIdentity, RowLanes, frame_length_headers,
physical_namespace_projection_map, recovered_entity_ref, row_frame_lanes,
write_canonical_row_frame, write_generation_identity_frames,
};

/// Rows per encode chunk. Sized so one chunk is a few milliseconds of decode
Expand Down Expand Up @@ -144,6 +144,46 @@ pub(crate) fn recovered_generation_digest_chunked(
Ok((encode_lowercase_hex(&digest.finalize()), canonical_bytes))
}

/// Streams each stored relation of `identity`'s generation to `emit` with
/// the row-sum lanes of its frame as `identity` recovers it, in identity
/// order. A relation's frame names its endpoints' projection, so the same
/// stored rows hash differently under each projection that reads them.
pub(crate) fn recovered_relation_lanes(
database: &GrafeoDB,
identity: &GraphGenerationManifestIdentity,
check: &dyn Fn() -> Result<(), GraphDbError>,
emit: &mut dyn FnMut(&str, RowLanes) -> Result<(), GraphDbError>,
) -> Result<(), GraphDbError> {
let store = database.graph_store();
let relations = projection_relation_nodes_sorted_checked(
database,
&identity.physical_namespace()?,
&identity.projection.projection,
check,
)?;
let namespace_projection = physical_namespace_projection_map(identity)?;
let mut canonical = CheckedVecWriter::new(check, MAX_GRAPH_REPLAY_SOURCE_BYTES_V1)?;
let mut endpoints = EndpointIdentityCache::default();
let mut endpoint_refs = HashMap::new();
for (sorted_identity, locator) in &relations {
check()?;
let relation = decode_sorted_relation(
store.as_ref(),
sorted_identity,
*locator,
&namespace_projection,
&mut endpoints,
&mut endpoint_refs,
)?;
let bytes = canonical.encode(&relation, "recovered generation relation")?;
emit(
sorted_identity.as_str(),
row_frame_lanes("relation", bytes)?,
)?;
}
Ok(())
}

/// The single-pass stream for generations at or below one chunk: one decoded
/// row resident at a time, every frame hashed as it is encoded.
#[tracing::instrument(
Expand Down
2 changes: 1 addition & 1 deletion crates/tracedecay-graph-db/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ pub use runtime::{GraphDb, GraphDbRuntimeState, GraphServingEnginePin, GraphSnap
pub use schema::graph_stable_identity;
pub use sealed_layer::{
GraphLayeredRowSpill, GraphLayeredRowsV1, GraphSealedBaseAbsenceV1, GraphSealedBaseV1,
LayeredGraphGeneration,
GraphSiblingSealedBaseV1, LayeredGraphGeneration,
};
pub use sealed_store::{SealedStoreCensusV1, census_sealed_store};

Expand Down
26 changes: 26 additions & 0 deletions crates/tracedecay-graph-db/src/registry/publication.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1306,6 +1306,32 @@ impl GraphDbRegistry {
operation.database().layered_row_spill(projection, base)
}

/// The sealed cold bases other scopes of this store serve for the same
/// projector, which a scope with no parent graph may layer over.
pub fn sibling_sealed_bases(
&self,
registration: GraphDbRegistration,
projection: &GraphProjectionIdentity,
check: &dyn Fn() -> Result<(), GraphDbError>,
) -> Result<Vec<crate::GraphSiblingSealedBaseV1>, GraphDbError> {
let operation = self.registered_operation(registration)?;
operation.database().sibling_sealed_bases(projection, check)
}

/// A row spill for a delta over a sibling scope's base, which it pins
/// until it seals.
pub fn sibling_layered_row_spill(
&self,
registration: GraphDbRegistration,
projection: GraphProjectionIdentity,
sibling: crate::GraphSiblingSealedBaseV1,
) -> Result<crate::GraphLayeredRowSpill, GraphDbError> {
let operation = self.registered_operation(registration)?;
operation
.database()
.sibling_layered_row_spill(projection, sibling)
}

/// Publishes through an already-issued, registry-validated graph lease.
///
/// The caller retains the exact operation lease through the publication;
Expand Down
59 changes: 58 additions & 1 deletion crates/tracedecay-graph-db/src/row_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -292,13 +292,26 @@ impl RowIndex {
count: u64,
key: RowKey,
) -> Result<Option<(RowLanes, u32, u32)>, GraphDbError> {
Ok(self
.position(first, count, key)?
.map(|(_, lanes, first_word, second_word)| (lanes, first_word, second_word)))
}

fn position(
&self,
first: u64,
count: u64,
key: RowKey,
) -> Result<Option<(u64, RowLanes, u32, u32)>, GraphDbError> {
let (mut low, mut high) = (0_u64, count);
while low < high {
let middle = low + (high - low) / 2;
let (found, lanes, first_word, second_word) =
record_parts(&self.record(first + middle)?);
match found.cmp(&key) {
std::cmp::Ordering::Equal => return Ok(Some((lanes, first_word, second_word))),
std::cmp::Ordering::Equal => {
return Ok(Some((first + middle, lanes, first_word, second_word)));
}
std::cmp::Ordering::Less => low = middle + 1,
std::cmp::Ordering::Greater => high = middle,
}
Expand All @@ -322,6 +335,50 @@ impl RowIndex {
.map(|(lanes, from, to)| IndexedRelation { lanes, from, to }))
}

/// Writes this index to `path` with every relation's lanes replaced by
/// the lanes `relane` emits for it, entity records and endpoint ordinals
/// unchanged: the index of the same rows read under another projection.
/// `relane` streams each recorded relation exactly once, in strictly
/// increasing identity order, so only one record is resident at a time.
pub(crate) fn write_relaned(
&self,
path: &Path,
relane: impl FnOnce(
&mut dyn FnMut(&str, RowLanes) -> Result<(), GraphDbError>,
) -> Result<(), GraphDbError>,
) -> Result<(), GraphDbError> {
std::fs::copy(&self.path, path).map_err(|error| index_io("copy", error))?;
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
let mut target = std::fs::OpenOptions::new()
.write(true)
.open(path)
.map_err(|error| index_io("open", error))?;
let mut relaned = 0_u64;
let mut previous: Option<String> = None;
relane(&mut |identity, lanes| {
if previous.as_deref().is_some_and(|last| last >= identity) {
return Err(corrupt("relanes relations out of identity order"));
}
previous = Some(identity.to_owned());
let Some((position, ..)) =
self.position(self.entities, self.relations, row_key("relation", identity))?
else {
return Err(corrupt("relanes a relation it does not record"));
};
target
.seek(SeekFrom::Start(HEADER_BYTES + position * RECORD_BYTES + 16))
.map_err(|error| index_io("seek", error))?;
target
.write_all(&lanes_bytes(lanes))
.map_err(|error| index_io("write", error))?;
relaned += 1;
Ok(())
})?;
if relaned != self.relations {
return Err(corrupt("relanes a different relation set than it records"));
}
target.sync_all().map_err(|error| index_io("sync", error))
}

/// The sum of every recorded row: equal to the generation's row sum
/// exactly when the index records the generation's rows.
pub(crate) fn row_sum(
Expand Down
Loading
Loading