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
12 changes: 12 additions & 0 deletions changelog.d/blob-cell-put-clear.added.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
- An embedded `Session` can replace or clear one Blob value of an existing
node or edge by exact id with `put_blob_at_as` and `clear_blob_at_as`. Each
write is one graph commit under the `change` action, never reads the old
value, and carries the row's other cells as an `update` does. An optional
precondition, `BlobPrecondition::Tags` (the current ETag is one of the given
tags) or `BlobPrecondition::AnyExisting` (the cell is not null), fails with
`BlobWritePreconditionFailed`, which names the current ETag, and changes
nothing. A put returns the value's length, its ETag and the commit; the ETag
is the one a read at that commit reports. One put holds at most the
session's `write_max_bytes` (32 MiB by default), inclusive; a larger one
fails before any table is opened with resource `Blob write payload bytes`.
The CLI and HTTP server do not expose these writes yet.
14 changes: 14 additions & 0 deletions crates/omnigraph-core/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,16 @@ pub enum OmniError {
expected: String,
actual: Option<String>,
},
/// A Blob write's `If-Match` precondition did not hold at the write's
/// pinned base, so the write had no effect. Distinct from the graph-head
/// `PreconditionFailed`: it compares the cell's strong validator, and a
/// null or external cell has none. `current_etag` is the managed cell's
/// validator at that base, or `None` when the cell has none.
#[error(
"Blob write precondition failed; the cell's current validator is {}",
current_etag.as_deref().unwrap_or("<none>")
)]
BlobWritePreconditionFailed { current_etag: Option<String> },
/// Engine-layer policy enforcement (MR-722). Wraps either a policy
/// denial ("you can't do that") or a policy-evaluation failure
/// ("the policy engine itself blew up"). The HTTP layer maps
Expand Down Expand Up @@ -863,6 +873,10 @@ impl OmniError {
}
}

pub fn blob_write_precondition_failed(current_etag: Option<String>) -> Self {
Self::BlobWritePreconditionFailed { current_etag }
}

pub fn recovery_required(operation_id: impl Into<String>, reason: impl Into<String>) -> Self {
Self::RecoveryRequired {
operation_id: operation_id.into(),
Expand Down
3 changes: 3 additions & 0 deletions crates/omnigraph-server/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1472,6 +1472,9 @@ impl ApiError {
),
api::PreconditionFailureOutput { expected, actual },
),
err @ OmniError::BlobWritePreconditionFailed { .. } => {
Self::blob_precondition_failed(err.to_string())
}
err @ OmniError::ExternalBlobPolicy { .. } => Self::bad_request(err.to_string()),
err @ OmniError::StoredExternalBlobDenied { .. } => Self::bad_request(err.to_string()),
OmniError::ExternalBlobSource { uri, reason } => {
Expand Down
1 change: 1 addition & 0 deletions crates/omnigraph-server/src/redacted_cause.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ impl RedactedCause {
OmniError::BlobRangeNotSatisfiable { .. } => "BlobRangeNotSatisfiable",
OmniError::RecoveryRequired { .. } => "RecoveryRequired",
OmniError::PreconditionFailed { .. } => "PreconditionFailed",
OmniError::BlobWritePreconditionFailed { .. } => "BlobWritePreconditionFailed",
OmniError::Policy(_) => "Policy",
OmniError::AlreadyInitialized { .. } => "AlreadyInitialized",
OmniError::InitializationCommitted { .. } => "InitializationCommitted",
Expand Down
234 changes: 152 additions & 82 deletions crates/omnigraph/src/blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,13 @@ pub struct BlobCell {
pub struct BlobEtag(String);

impl BlobEtag {
/// A caller-supplied strong entity tag, for a [`BlobPrecondition::Tags`]
/// comparison. It is compared exactly; a tag the engine never issued
/// matches nothing.
pub fn from_tag(tag: impl Into<String>) -> Self {
Self(tag.into())
}

pub fn as_str(&self) -> &str {
&self.0
}
Expand All @@ -74,6 +81,34 @@ impl fmt::Display for BlobEtag {
}
}

/// The condition a Blob write requires of the cell at its pinned write base.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BlobPrecondition {
/// A non-null value exists, managed or external. A null cell fails it.
AnyExisting,
/// The managed value's validator equals one of these. A null or external
/// cell has no validator and fails it.
Tags(Vec<BlobEtag>),
}

/// The result of a Blob write: the cell's new state and the exact commit of
/// the publication that made it visible.
#[derive(Debug, Clone)]
pub enum BlobWriteOutcome {
/// The cell now holds managed bytes. `etag` equals what a read at
/// `commit` returns.
Managed {
length: u64,
etag: BlobEtag,
commit: crate::db::GraphCommit,
},
/// The cell is now null. `commit` is `None` when it already was, and the
/// clear published nothing.
Null {
commit: Option<crate::db::GraphCommit>,
},
}

/// Descriptor-only external Blob reference. Reading the cell never probes or
/// dereferences this URI.
#[derive(Debug, Clone, PartialEq, Eq)]
Expand Down Expand Up @@ -1210,11 +1245,109 @@ impl<'a> BlobDescriptorDecoder<'a> {
}
}

struct ResolvedBlobCell {
table_key: String,
pub(crate) struct ResolvedBlobCell {
pub(crate) table_key: String,
pub(crate) stable_table_id: u64,
pub(crate) table_incarnation_id: u64,
pub(crate) stable_property_id: u64,
}

/// Find one entity's row in an opened table version and classify its Blob
/// cell from the descriptor alone, without opening any payload: the stable
/// row id and the descriptor, or `None` when no row has the id. A duplicate id
/// or a malformed descriptor is a Blob integrity failure.
pub(crate) async fn locate_blob_cell(
dataset: &lance::Dataset,
id_column: &str,
cell: &BlobCell,
) -> Result<Option<(u64, BlobDescriptor)>> {
let mut scanner = dataset.scan();
scanner
.project(&[cell.property.as_str()])
.map_err(OmniError::storage)?;
scanner.filter_expr(col(id_column).eq(lit(cell.id.clone())));
scanner.blob_handling(BlobHandling::BlobsDescriptions);
scanner.with_row_id();
scanner.limit(Some(2), None).map_err(OmniError::storage)?;
let batches = scanner
.try_into_stream()
.await
.map_err(OmniError::storage)?
.try_collect::<Vec<_>>()
.await
.map_err(OmniError::storage)?;

let mut selected = None;
for batch in &batches {
let descriptions = batch
.column_by_name(&cell.property)
.and_then(|column| column.as_any().downcast_ref::<StructArray>())
.ok_or_else(|| {
OmniError::blob_integrity(format!(
"Blob property '{}.{}' did not scan as a Blob-v2 descriptor",
cell.type_name, cell.property
))
})?;
let row_ids = batch
.column_by_name("_rowid")
.and_then(|column| column.as_any().downcast_ref::<UInt64Array>())
.ok_or_else(|| {
OmniError::blob_integrity(format!(
"Blob lookup for '{}.{}' omitted its stable row id",
cell.type_name, cell.property
))
})?;
if descriptions.len() != row_ids.len() {
return Err(OmniError::blob_integrity(format!(
"Blob descriptor and stable-row-id cardinalities differ for '{}.{}'",
cell.type_name, cell.property
)));
}
let decoder = BlobDescriptorDecoder::try_new(descriptions)?;
for row in 0..batch.num_rows() {
if selected.is_some() {
return Err(OmniError::blob_integrity(format!(
"entity id '{}' appears more than once in {} type '{}'",
cell.id,
entity_label(cell.entity),
cell.type_name
)));
}
if row_ids.is_null(row) {
return Err(OmniError::blob_integrity(format!(
"Blob lookup for '{}.{}' returned a null stable row id",
cell.type_name, cell.property
)));
}
selected = Some((row_ids.value(row), decoder.classify(row)?));
}
}
Ok(selected)
}

/// The strong validator of a managed cell in this exact opened table version:
/// table and property identity, the opened version, the row's stable id and
/// the transaction file named by the version's immutable manifest. A write
/// computes it from its own detached commit, a read from the version it opened,
/// and the two agree for the same version and row.
pub(crate) fn managed_blob_etag(
dataset: &lance::Dataset,
stable_table_id: u64,
table_incarnation_id: u64,
stable_property_id: u64,
stable_row_id: u64,
cell: &BlobCell,
) -> Result<BlobEtag> {
let transaction_file =
immutable_transaction_witness(dataset.manifest().transaction_file.as_deref(), cell)?;
Ok(blob_etag(
stable_table_id,
table_incarnation_id,
stable_property_id,
dataset.version().version,
stable_row_id,
transaction_file,
))
}

decide_seam! {
Expand Down Expand Up @@ -1452,77 +1585,17 @@ impl Omnigraph {
None => {}
}

let mut scanner = dataset.scan();
scanner
.project(&[cell.property.as_str()])
.map_err(OmniError::storage)?;
scanner.filter_expr(col(catalog.system_columns.id).eq(lit(cell.id.clone())));
scanner.blob_handling(BlobHandling::BlobsDescriptions);
scanner.with_row_id();
scanner.limit(Some(2), None).map_err(OmniError::storage)?;
let stream = scanner
.try_into_stream()
.await
.map_err(OmniError::storage)?;
let batches = stream
.try_collect::<Vec<_>>()
.await
.map_err(OmniError::storage)?;

let mut selected = None;
for batch in &batches {
let descriptions = batch
.column_by_name(&cell.property)
.and_then(|column| column.as_any().downcast_ref::<StructArray>())
.ok_or_else(|| {
OmniError::blob_integrity(format!(
"Blob property '{}.{}' did not scan as a Blob-v2 descriptor",
cell.type_name, cell.property
))
})?;
let row_ids = batch
.column_by_name("_rowid")
.and_then(|column| column.as_any().downcast_ref::<UInt64Array>())
let (stable_row_id, descriptor) =
locate_blob_cell(&dataset, catalog.system_columns.id, &cell)
.await?
.ok_or_else(|| {
OmniError::blob_integrity(format!(
"Blob lookup for '{}.{}' omitted its stable row id",
cell.type_name, cell.property
OmniError::manifest_not_found(format!(
"no {} '{}' with id '{}' found",
entity_label(cell.entity),
cell.type_name,
cell.id
))
})?;
if descriptions.len() != row_ids.len() {
return Err(OmniError::blob_integrity(format!(
"Blob descriptor and stable-row-id cardinalities differ for '{}.{}'",
cell.type_name, cell.property
)));
}
let decoder = BlobDescriptorDecoder::try_new(descriptions)?;
for row in 0..batch.num_rows() {
if selected.is_some() {
return Err(OmniError::blob_integrity(format!(
"entity id '{}' appears more than once in {} type '{}'",
cell.id,
entity_label(cell.entity),
cell.type_name
)));
}
if row_ids.is_null(row) {
return Err(OmniError::blob_integrity(format!(
"Blob lookup for '{}.{}' returned a null stable row id",
cell.type_name, cell.property
)));
}
selected = Some((row_ids.value(row), decoder.classify(row)?));
}
}

let (stable_row_id, descriptor) = selected.ok_or_else(|| {
OmniError::manifest_not_found(format!(
"no {} '{}' with id '{}' found",
entity_label(cell.entity),
cell.type_name,
cell.id
))
})?;

let content = match descriptor {
BlobDescriptor::Null => {
Expand All @@ -1546,8 +1619,12 @@ impl Omnigraph {
// after branch delete/recreate. The transaction-file identity
// is stored in this exact immutable manifest and carries the
// UUID of the transaction that produced this table version.
let transaction_file = immutable_transaction_witness(
dataset.manifest().transaction_file.as_deref(),
let etag = managed_blob_etag(
&dataset,
stable_table_id,
table_incarnation_id,
resolved_cell.stable_property_id,
stable_row_id,
&cell,
)?;
let mut files = dataset
Expand Down Expand Up @@ -1579,14 +1656,7 @@ impl Omnigraph {
}
BlobContent::Managed {
length,
etag: blob_etag(
stable_table_id,
table_incarnation_id,
resolved_cell.stable_property_id,
actual_table_version,
stable_row_id,
transaction_file,
),
etag,
reader: BlobReader {
_dataset: Arc::clone(&dataset),
file: Arc::new(file),
Expand All @@ -1604,7 +1674,7 @@ impl Omnigraph {
}
}

fn resolve_blob_cell(
pub(crate) fn resolve_blob_cell(
catalog: &omnigraph_compiler::catalog::Catalog,
cell: &BlobCell,
) -> Result<ResolvedBlobCell> {
Expand Down Expand Up @@ -1671,7 +1741,7 @@ fn resolve_blob_cell(
})
}

fn entity_label(kind: EntityKind) -> &'static str {
pub(crate) fn entity_label(kind: EntityKind) -> &'static str {
match kind {
EntityKind::Node => "node",
EntityKind::Edge => "edge",
Expand Down
Loading
Loading