Skip to content

Commit 1d608f1

Browse files
committed
add set_value_transformer
1 parent 1eb1d46 commit 1d608f1

3 files changed

Lines changed: 326 additions & 11 deletions

File tree

src/git/versioned_store/mod.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -225,7 +225,9 @@ pub use namespaced::{
225225
};
226226

227227
#[cfg(feature = "proximity")]
228-
pub use namespaced::{ProximityNamespaceHandle, TextIndexAudit, TextNamespaceHandle};
228+
pub use namespaced::{
229+
ProximityNamespaceHandle, TextIndexAudit, TextNamespaceHandle, ValueTransformer,
230+
};
229231

230232
#[cfg(feature = "rocksdb_storage")]
231233
pub use namespaced::RocksDBNamespacedKvStore;

src/git/versioned_store/namespaced.rs

Lines changed: 93 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,12 @@ pub struct MigrationReport {
9595
/// The default namespace name used for backward-compatible flat API calls.
9696
pub const DEFAULT_NAMESPACE: &str = "default";
9797

98+
/// Type alias for the per-target value transformer closure used by cascade.
99+
/// Takes the raw primary-tree value bytes and returns `Some(text)` to embed,
100+
/// or `None` to opt this id out of cascade for this index.
101+
#[cfg(feature = "proximity")]
102+
pub type ValueTransformer = Arc<dyn Fn(&[u8]) -> Option<String> + Send + Sync>;
103+
98104
// ---------------------------------------------------------------------------
99105
// NamespacedKvStore
100106
// ---------------------------------------------------------------------------
@@ -150,6 +156,14 @@ pub struct NamespacedKvStore<
150156
/// are silently skipped — drift can be detected via `audit_text_index`.
151157
#[cfg(feature = "proximity")]
152158
pub(crate) cascade_lists: HashMap<String, Vec<String>>,
159+
/// Per-target value transformers consulted during cascade. When
160+
/// `(ns, idx_name)` is keyed here, the closure runs on the raw value
161+
/// bytes to produce the text that will be embedded. Returning `None`
162+
/// opts the id out of cascade for this index — useful for "this row
163+
/// isn't text-indexable" decisions on structured payloads. Absence of
164+
/// a transformer falls back to UTF-8 interpretation.
165+
#[cfg(feature = "proximity")]
166+
pub(crate) text_transformers: HashMap<(String, String), ValueTransformer>,
153167
/// Threshold (in bytes) above which inserted values are externalised as
154168
/// separate blobs via [`crate::storage::NodeStorage::insert_blob`].
155169
/// `None` (default) disables externalisation. Runtime config only —
@@ -293,26 +307,43 @@ impl<'a, const N: usize, S: NodeStorage<N>, M: MetadataBackend> NamespaceHandle<
293307
}
294308

295309
/// Cascade an insert to every configured text sub-index for this
296-
/// namespace. Silent no-op for namespaces without a cascade list, for
297-
/// non-UTF-8 values, and for targets whose embedder/proximity-index
298-
/// isn't currently loaded.
310+
/// namespace. Per target the text to embed comes from:
311+
///
312+
/// 1. A registered `ValueTransformer` (PR 4d follow-up) if one exists —
313+
/// its `Option<String>` return value determines whether this id
314+
/// participates (`None` opts out for this index).
315+
/// 2. Otherwise UTF-8 interpretation of the bytes (silent skip on
316+
/// non-UTF-8).
317+
///
318+
/// Silent no-op cases otherwise: no cascade list configured, target
319+
/// whose embedder/proximity-index isn't currently loaded, or the
320+
/// embedder itself returns an error.
299321
#[cfg(feature = "proximity")]
300322
fn cascade_insert(&mut self, key: &[u8], value: &[u8]) {
301323
let Some(cascade_list) = self.store.cascade_lists.get(&self.ns_name).cloned() else {
302324
return;
303325
};
304-
let Ok(text) = std::str::from_utf8(value) else {
305-
return;
306-
};
307326
for idx_name in cascade_list {
308-
let embedder_key = (self.ns_name.clone(), idx_name.clone());
327+
let target_key = (self.ns_name.clone(), idx_name.clone());
309328
let prox_key = (self.ns_name.clone(), text_inner_proximity_name(&idx_name));
310329

311-
// Resolve the embedder first (immutable borrow ends with the clone).
312-
let Some(embedder) = self.store.text_embedders.get(&embedder_key).cloned() else {
330+
// 1. Determine the text to embed.
331+
let text: String = match self.store.text_transformers.get(&target_key).cloned() {
332+
Some(transformer) => match transformer(value) {
333+
Some(t) => t,
334+
None => continue, // transformer explicitly opted out
335+
},
336+
None => match std::str::from_utf8(value) {
337+
Ok(s) => s.to_string(),
338+
Err(_) => continue, // non-UTF-8 + no transformer
339+
},
340+
};
341+
342+
// 2. Resolve the embedder (immutable borrow ends with the clone).
343+
let Some(embedder) = self.store.text_embedders.get(&target_key).cloned() else {
313344
continue;
314345
};
315-
let Ok(vec) = embedder.embed(text) else {
346+
let Ok(vec) = embedder.embed(&text) else {
316347
continue;
317348
};
318349
if let Some(idx) = self.store.proximity_indexes.get_mut(&prox_key) {
@@ -934,6 +965,8 @@ impl<const N: usize> NamespacedKvStore<N, GitNodeStorage<N>, GitMetadataBackend>
934965
text_embedders: HashMap::new(),
935966
#[cfg(feature = "proximity")]
936967
cascade_lists: HashMap::new(),
968+
#[cfg(feature = "proximity")]
969+
text_transformers: HashMap::new(),
937970
externalize_threshold: None,
938971
};
939972

@@ -975,6 +1008,8 @@ impl<const N: usize> NamespacedKvStore<N, GitNodeStorage<N>, GitMetadataBackend>
9751008
text_embedders: HashMap::new(),
9761009
#[cfg(feature = "proximity")]
9771010
cascade_lists: HashMap::new(),
1011+
#[cfg(feature = "proximity")]
1012+
text_transformers: HashMap::new(),
9781013
externalize_threshold: None,
9791014
};
9801015

@@ -1862,6 +1897,8 @@ impl<const N: usize> NamespacedKvStore<N, InMemoryNodeStorage<N>, GitMetadataBac
18621897
text_embedders: HashMap::new(),
18631898
#[cfg(feature = "proximity")]
18641899
cascade_lists: HashMap::new(),
1900+
#[cfg(feature = "proximity")]
1901+
text_transformers: HashMap::new(),
18651902
externalize_threshold: None,
18661903
};
18671904

@@ -1900,6 +1937,8 @@ impl<const N: usize> NamespacedKvStore<N, InMemoryNodeStorage<N>, GitMetadataBac
19001937
text_embedders: HashMap::new(),
19011938
#[cfg(feature = "proximity")]
19021939
cascade_lists: HashMap::new(),
1940+
#[cfg(feature = "proximity")]
1941+
text_transformers: HashMap::new(),
19031942
externalize_threshold: None,
19041943
};
19051944

@@ -1950,6 +1989,8 @@ impl<const N: usize> NamespacedKvStore<N, FileNodeStorage<N>, GitMetadataBackend
19501989
text_embedders: HashMap::new(),
19511990
#[cfg(feature = "proximity")]
19521991
cascade_lists: HashMap::new(),
1992+
#[cfg(feature = "proximity")]
1993+
text_transformers: HashMap::new(),
19531994
externalize_threshold: None,
19541995
};
19551996

@@ -1988,6 +2029,8 @@ impl<const N: usize> NamespacedKvStore<N, FileNodeStorage<N>, GitMetadataBackend
19882029
text_embedders: HashMap::new(),
19892030
#[cfg(feature = "proximity")]
19902031
cascade_lists: HashMap::new(),
2032+
#[cfg(feature = "proximity")]
2033+
text_transformers: HashMap::new(),
19912034
externalize_threshold: None,
19922035
};
19932036

@@ -2661,6 +2704,7 @@ where
26612704
let embedder_key = (self.ns_name.clone(), idx_name.to_string());
26622705
self.store.dirty_proximity_indexes.remove(&inner_idx_key);
26632706
self.store.text_embedders.remove(&embedder_key);
2707+
self.store.text_transformers.remove(&embedder_key);
26642708
self.store
26652709
.proximity_indexes
26662710
.remove(&inner_idx_key)
@@ -2707,6 +2751,45 @@ where
27072751
pub fn cascade_for_namespace(&self, ns_name: &str) -> Option<&[String]> {
27082752
self.cascade_lists.get(ns_name).map(|v| v.as_slice())
27092753
}
2754+
2755+
/// Register a value transformer for a `(namespace, text_index)` target.
2756+
///
2757+
/// During cascade, the transformer runs on the raw primary-tree value
2758+
/// bytes for each id. It returns either:
2759+
///
2760+
/// - `Some(text)` — the text to embed and upsert into this index.
2761+
/// - `None` — opt this id out of cascade for this index. The primary
2762+
/// tree still gets the bytes; only the named text index is skipped.
2763+
///
2764+
/// Without a registered transformer, cascade falls back to interpreting
2765+
/// the bytes as UTF-8 (and silently skips non-UTF-8 values).
2766+
///
2767+
/// Runtime configuration only — not persisted; users register their
2768+
/// transformers each process. Typical use: extract a field from a JSON
2769+
/// payload, or stringify a Protobuf message.
2770+
pub fn set_value_transformer<F>(&mut self, ns_name: &str, idx_name: &str, transformer: F)
2771+
where
2772+
F: Fn(&[u8]) -> Option<String> + Send + Sync + 'static,
2773+
{
2774+
self.text_transformers.insert(
2775+
(ns_name.to_string(), idx_name.to_string()),
2776+
Arc::new(transformer),
2777+
);
2778+
}
2779+
2780+
/// Remove the value transformer for a `(namespace, text_index)` target.
2781+
/// Returns whether one was registered.
2782+
pub fn clear_value_transformer(&mut self, ns_name: &str, idx_name: &str) -> bool {
2783+
self.text_transformers
2784+
.remove(&(ns_name.to_string(), idx_name.to_string()))
2785+
.is_some()
2786+
}
2787+
2788+
/// True when a value transformer is registered for the given target.
2789+
pub fn has_value_transformer(&self, ns_name: &str, idx_name: &str) -> bool {
2790+
self.text_transformers
2791+
.contains_key(&(ns_name.to_string(), idx_name.to_string()))
2792+
}
27102793
}
27112794

27122795
// ---------------------------------------------------------------------------

0 commit comments

Comments
 (0)