From 26bd282ac014a44dbde858c6042940bcff138760 Mon Sep 17 00:00:00 2001 From: michaelk1957 Date: Fri, 12 Jun 2026 04:11:57 +0300 Subject: [PATCH] fix(memory): stop oversized/dense chunks from breaking bge-m3 embedding (#3598) Co-authored-by: Claude Opus 4.8 (1M context) --- src/openhuman/memory_queue/handlers/mod.rs | 24 +- .../memory_queue/handlers/mod_tests.rs | 23 + src/openhuman/memory_store/chunks/produce.rs | 441 ++++++++++++++---- src/openhuman/memory_store/chunks/types.rs | 89 ++++ .../memory_tree/score/embed/ollama.rs | 10 + 5 files changed, 507 insertions(+), 80 deletions(-) diff --git a/src/openhuman/memory_queue/handlers/mod.rs b/src/openhuman/memory_queue/handlers/mod.rs index 2c6641aa1..2cf021586 100644 --- a/src/openhuman/memory_queue/handlers/mod.rs +++ b/src/openhuman/memory_queue/handlers/mod.rs @@ -20,7 +20,7 @@ use crate::openhuman::memory_queue::types::{ JobOutcome, NewJob, NodeRef, ReembedBackfillPayload, SealDocumentPayload, SealPayload, }; use crate::openhuman::memory_store::chunks::store as chunk_store; -use crate::openhuman::memory_store::chunks::types::Chunk; +use crate::openhuman::memory_store::chunks::types::{truncate_to_conservative_tokens, Chunk}; use crate::openhuman::memory_store::content::{ self as content_store, read as content_read, tags as content_tags, }; @@ -34,6 +34,21 @@ use crate::openhuman::memory_tree::tree::{LeafRef, TreeFactory}; /// 1 hour means low-volume sources get summaries within a working session. const L0_DEFAULT_FLUSH_AGE_SECS: i64 = 60 * 60; +/// Conservative per-text embed token budget. Each body is truncated to this +/// (via [`cap_embed_text`]) before it joins an `embed_batch` call, so no single +/// input can exceed the embedder's batch / context limit (`EMBED_NUM_CTX` = +/// 8192) and terminally fail the reembed job. The estimate over-counts dense / +/// multilingual text, so this stays safely under the real limit; the chunker +/// keeps normal chunks well below it, so truncation is a last-resort backstop. +const EMBED_SAFE_TOKENS: u32 = 7500; + +/// Truncate one body to [`EMBED_SAFE_TOKENS`] before it is batched for +/// embedding. A backstop for any body that reaches an embed call without having +/// passed through the conservative chunker (`split_by_token_budget`). +fn cap_embed_text(text: &str) -> &str { + truncate_to_conservative_tokens(text, EMBED_SAFE_TOKENS) +} + /// Maximum `extract_chunk` jobs to coalesce in one worker tick. /// /// Scoring/extraction runs per chunk; embedding is deferred to a post-sync @@ -742,7 +757,12 @@ async fn reembed_collect( // Phase B: one batched embed call. Scope `texts` so its borrow on // `readable` ends before we consume `readable` below. let results = { - let texts: Vec<&str> = readable.iter().map(|(_, body)| body.as_str()).collect(); + // Cap each body to the embed budget so no single input overflows the + // embedder's batch/context limit and fails the whole batch. + let texts: Vec<&str> = readable + .iter() + .map(|(_, body)| cap_embed_text(body)) + .collect(); embedder.embed_batch(&texts).await }; if results.len() != readable.len() { diff --git a/src/openhuman/memory_queue/handlers/mod_tests.rs b/src/openhuman/memory_queue/handlers/mod_tests.rs index f346ff920..7984f1468 100644 --- a/src/openhuman/memory_queue/handlers/mod_tests.rs +++ b/src/openhuman/memory_queue/handlers/mod_tests.rs @@ -571,3 +571,26 @@ async fn ensure_reembed_backfill_enqueues_only_when_uncovered() { "re-call must dedupe to a single chain per signature" ); } + +/// `cap_embed_text` must bound an over-budget body to `EMBED_SAFE_TOKENS` before +/// it is batched for embedding, and pass small bodies through verbatim, so no +/// single `embed_batch` input can exceed the embedder's input limit. +#[test] +fn cap_embed_text_bounds_oversized_input() { + use crate::openhuman::memory_store::chunks::types::conservative_token_estimate; + + // Punctuation-dense body (~1 token/char) far over the embed budget. + let big = "x,".repeat(EMBED_SAFE_TOKENS as usize); + assert!(conservative_token_estimate(&big) > EMBED_SAFE_TOKENS); + let capped = cap_embed_text(&big); + assert!(capped.len() < big.len(), "oversized body must be truncated"); + assert!( + conservative_token_estimate(capped) <= EMBED_SAFE_TOKENS, + "capped body must be within the embed budget", + ); + assert!(big.starts_with(capped), "cap must return a leading prefix"); + + // A small body is returned unchanged. + let small = "hello world"; + assert_eq!(cap_embed_text(small), small); +} diff --git a/src/openhuman/memory_store/chunks/produce.rs b/src/openhuman/memory_store/chunks/produce.rs index ac244bc90..f31d960fd 100644 --- a/src/openhuman/memory_store/chunks/produce.rs +++ b/src/openhuman/memory_store/chunks/produce.rs @@ -10,15 +10,22 @@ //! //! - **Chat**: split at `## ` message boundaries. Each message becomes one //! chunk. If a single message exceeds `max_tokens`, fall back to the -//! paragraph/line/char splitter for that unit only and emit each piece with -//! `partial_message = true`. +//! paragraph/sentence/whitespace/char splitter for that unit only and emit +//! each piece with `partial_message = true`. //! - **Email**: split at `---\nFrom:` separators. Each email in the thread //! becomes one chunk. Same oversize fallback as Chat. -//! - **Document**: original paragraph-based greedy packing (unchanged). +//! - **Document**: split by [`split_by_token_budget`] — a conservative +//! token-estimate splitter (paragraph → sentence → whitespace → hard-char) +//! with ~12% overlap between adjacent chunks. +//! +//! Chunk sizes are bounded by [`conservative_token_estimate`], not the GPT +//! `chars/4` heuristic, so dense markdown/hash/code/multilingual content cannot +//! produce an over-budget chunk that overflows a downstream embedder. use crate::openhuman::memory::util::redact::redact; use crate::openhuman::memory_store::chunks::types::{ - approx_token_count, chunk_id, Chunk, Metadata, SourceKind, + approx_token_count, chunk_id, conservative_token_estimate, truncate_to_conservative_tokens, + Chunk, Metadata, SourceKind, }; /// Default upper bound on per-chunk tokens. @@ -68,9 +75,11 @@ pub struct ChunkerInput { /// - **Chat / Email**: split at message/email boundaries, then greedy-pack /// consecutive units into a single chunk until adding the next unit would /// exceed `max_tokens`. Oversize units (a single message > `max_tokens`) -/// fall back to the paragraph/line/char splitter and emit each piece with +/// fall back to [`split_by_token_budget`] and emit each piece with /// `partial_message = true`. -/// - **Document**: original paragraph-based greedy packing (unchanged). +/// - **Document**: split by [`split_by_token_budget`] — sized by the +/// conservative token estimate (paragraph → sentence → whitespace → +/// hard-char) with ~12% overlap between adjacent chunks. pub fn chunk_markdown(input: &ChunkerInput, opts: &ChunkerOptions) -> Vec { let now = chrono::Utc::now(); let max_tokens = opts.max_tokens.max(1); @@ -313,103 +322,243 @@ fn split_email_messages(md: &str) -> Vec { pieces } -/// Split `text` into pieces each ≤ `max_tokens` tokens. +/// Overlap carried between adjacent chunks (≈12%, within the 10–15% band) so a +/// fact straddling a split boundary survives in both neighbours. +const OVERLAP_PERCENT: u32 = 12; +/// Hard cap on overlap (% of budget) so a chunk can never be mostly a duplicate +/// of its predecessor — avoids pathological near-identical chunks. +const OVERLAP_MAX_PERCENT: u32 = 40; + +/// Split `text` into pieces each ≤ `max_tokens` **conservative** tokens +/// (see [`conservative_token_estimate`]), with ~10–15% overlap between adjacent +/// pieces. Sized by the conservative estimate (NOT `chars/4`), so dense +/// markdown/hash/code/multilingual content — which tokenises far above the +/// chars/4 heuristic — never produces an over-budget chunk that overflows the +/// embedder. /// -/// Preference order for split boundaries: -/// 1. Paragraph (`\n\n`) -/// 2. Line (`\n`) -/// 3. Hard character cut (last resort; preserves UTF-8 code points) +/// Boundary preference (a still-oversized piece falls to the next finer level): +/// 1. paragraph (`\n\n`) +/// 2. sentence (`. ` / `! ` / `? ` / line break) +/// 3. whitespace (word) +/// 4. hard character cut (last resort; preserves UTF-8 code points) +/// +/// Ordering is preserved. Overlap repeats the previous piece's trailing whole +/// segments verbatim (snapped to natural boundaries), never the entire chunk. pub(crate) fn split_by_token_budget(text: &str, max_tokens: u32) -> Vec { - let max_tokens = max_tokens.max(1); + let budget = max_tokens.max(1); if text.is_empty() { return vec![String::new()]; } - if approx_token_count(text) <= max_tokens { + if conservative_token_estimate(text) <= budget { return vec![text.to_string()]; } - // Approximate max chars per chunk (4 chars ≈ 1 token). - let max_chars: usize = (max_tokens as usize).saturating_mul(4); + // Reserve headroom so `overlap (≤cap) + any segment (≤seg_budget) ≤ budget`, + // which prevents a flush from ever emitting an overlap-only (duplicate) chunk. + let overlap_budget = (budget * OVERLAP_PERCENT / 100).max(1); + let overlap_cap = (budget * OVERLAP_MAX_PERCENT / 100).max(overlap_budget); + let seg_budget = budget.saturating_sub(overlap_cap).max(1); - // First: try paragraph split. Walk paragraphs, greedy-accumulate into - // chunks ≤ max_chars. - let paragraphs: Vec<&str> = text.split("\n\n").collect(); - if paragraphs.len() > 1 { - if let Some(out) = pack_segments(¶graphs, "\n\n", max_chars) { - return out; - } + // 1. Reduce to in-order atomic segments, each ≤ seg_budget. + let mut segments: Vec<&str> = Vec::new(); + push_atomic_segments(text, seg_budget, &mut segments); + if segments.len() <= 1 { + return vec![text.to_string()]; } - // Fall back to line split. - let lines: Vec<&str> = text.split('\n').collect(); - if lines.len() > 1 { - if let Some(out) = pack_segments(&lines, "\n", max_chars) { - return out; + // 2. Greedy-pack segments into chunks ≤ budget, carrying overlap forward. + let mut chunks: Vec = Vec::new(); + let mut cur: Vec<&str> = Vec::new(); + let mut cur_tokens = 0u32; + for seg in &segments { + let seg_tokens = conservative_token_estimate(seg); + if !cur.is_empty() && cur_tokens.saturating_add(seg_tokens) > budget { + chunks.push(join_segments(&cur)); + let overlap = tail_overlap(&cur, overlap_budget, overlap_cap); + cur_tokens = overlap.iter().map(|s| conservative_token_estimate(s)).sum(); + cur = overlap; } + cur.push(seg); + cur_tokens = cur_tokens.saturating_add(seg_tokens); } - - // Fall back to hard character-count cut preserving UTF-8 boundaries. - hard_split_by_chars(text, max_chars) + if !cur.is_empty() { + chunks.push(join_segments(&cur)); + } + if chunks.is_empty() { + chunks.push(String::new()); + } + chunks } -/// Greedily pack pre-split segments into chunks ≤ max_chars. Returns `None` -/// if any single segment is already too large — caller should try a finer -/// split. -fn pack_segments(segments: &[&str], sep: &str, max_chars: usize) -> Option> { +/// Append in-order atomic segments of `text`, each with +/// `conservative_token_estimate ≤ budget`, using the boundary hierarchy. +fn push_atomic_segments<'a>(text: &'a str, budget: u32, out: &mut Vec<&'a str>) { + if text.is_empty() { + return; + } + if conservative_token_estimate(text) <= budget { + out.push(text); + return; + } + if split_on_separator(text, "\n\n", budget, out) + || split_on_sentences(text, budget, out) + || split_on_whitespace(text, budget, out) + { + return; + } + hard_split(text, budget, out); +} + +/// Split on a literal separator; recurse on each piece. The separator is kept +/// at the END of its preceding piece so the pieces tile `text` exactly (lossless +/// concatenation — see [`join_segments`]). Returns `false` (no progress) when the +/// separator does not occur. +fn split_on_separator<'a>(text: &'a str, sep: &str, budget: u32, out: &mut Vec<&'a str>) -> bool { + if sep.is_empty() || !text.contains(sep) { + return false; + } let sep_len = sep.len(); - let mut out: Vec = Vec::new(); - let mut current = String::new(); - - for seg in segments { - let seg_len = seg.chars().count(); - // A single segment larger than max_chars forces a finer split. - if seg_len > max_chars { - return None; - } - let projected = if current.is_empty() { - seg_len - } else { - current.chars().count() + sep_len + seg_len - }; - if projected > max_chars { - out.push(std::mem::take(&mut current)); - current.push_str(seg); - } else { - if !current.is_empty() { - current.push_str(sep); - } - current.push_str(seg); - } + let mut pieces: Vec<&str> = Vec::new(); + let mut start = 0usize; + let mut search = 0usize; + while let Some(rel) = text[search..].find(sep) { + let end = search + rel + sep_len; // include the separator in this piece + pieces.push(&text[start..end]); + start = end; + search = end; } - if !current.is_empty() { - out.push(current); + if start < text.len() { + pieces.push(&text[start..]); } - if out.is_empty() { - out.push(String::new()); + if pieces.len() <= 1 { + return false; } - Some(out) + for p in pieces { + push_atomic_segments(p, budget, out); + } + true } -/// Hard character-count cut preserving UTF-8 code-point boundaries. -fn hard_split_by_chars(text: &str, max_chars: usize) -> Vec { - let mut out = Vec::new(); - let mut current = String::new(); - let mut count = 0usize; - for ch in text.chars() { - if count + 1 > max_chars { - out.push(std::mem::take(&mut current)); - count = 0; +/// Split on sentence-ish boundaries: after `.`/`!`/`?` followed by a space, and +/// at line breaks. All boundary bytes are ASCII (1 byte), so slicing is always +/// on a char boundary. Returns `false` when no boundary is found. +fn split_on_sentences<'a>(text: &'a str, budget: u32, out: &mut Vec<&'a str>) -> bool { + let bytes = text.as_bytes(); + let mut pieces: Vec<&str> = Vec::new(); + let mut start = 0usize; + for i in 0..bytes.len() { + let c = bytes[i]; + let boundary_end = if c == b'\n' { + Some(i + 1) + } else if (c == b'.' || c == b'!' || c == b'?') + && i + 1 < bytes.len() + && bytes[i + 1] == b' ' + { + Some(i + 1) + } else { + None + }; + if let Some(end) = boundary_end { + pieces.push(&text[start..end]); + start = end; } - current.push(ch); - count += 1; } - if !current.is_empty() { - out.push(current); + if start < text.len() { + pieces.push(&text[start..]); } - if out.is_empty() { - out.push(String::new()); + if pieces.len() <= 1 { + return false; } - out + for p in pieces { + push_atomic_segments(p, budget, out); + } + true +} + +/// Split on ASCII spaces (word boundaries); recurse on each word. Returns +/// `false` when there is no space to split on. +fn split_on_whitespace<'a>(text: &'a str, budget: u32, out: &mut Vec<&'a str>) -> bool { + let bytes = text.as_bytes(); + let mut pieces: Vec<&str> = Vec::new(); + let mut start = 0usize; + for i in 0..bytes.len() { + if bytes[i] == b' ' { + // Keep the space at the end of the word so pieces tile `text` exactly. + pieces.push(&text[start..=i]); + start = i + 1; + } + } + if start < text.len() { + pieces.push(&text[start..]); + } + if pieces.len() <= 1 { + return false; + } + for p in pieces { + push_atomic_segments(p, budget, out); + } + true +} + +/// Last resort: cut a boundary-free run into ≤ `budget`-token pieces on UTF-8 +/// char boundaries. Reuses [`truncate_to_conservative_tokens`] for sizing and +/// always makes progress (≥1 char) so it cannot loop. +fn hard_split<'a>(mut text: &'a str, budget: u32, out: &mut Vec<&'a str>) { + while !text.is_empty() { + let mut piece = truncate_to_conservative_tokens(text, budget); + if piece.is_empty() { + // Budget smaller than one char's weight — take a single char. + let end = text + .char_indices() + .nth(1) + .map(|(i, _)| i) + .unwrap_or(text.len()); + piece = &text[..end]; + } + let plen = piece.len(); + out.push(&text[..plen]); + text = &text[plen..]; + } +} + +/// Join packed segments back into one chunk body. Segments are contiguous, +/// separator-retaining slices of the source, so plain concatenation reproduces +/// the original text exactly (modulo intentional overlap duplication) — no +/// spurious separators are introduced, which keeps each chunk's real token +/// count ≤ the packing budget. +fn join_segments(segs: &[&str]) -> String { + segs.concat() +} + +/// Trailing whole segments of `cur` summing to ~`overlap_budget` tokens (capped +/// at `overlap_cap`), returned in original order. Never returns the entire +/// chunk, so adjacent chunks cannot be duplicates. +fn tail_overlap<'a>(cur: &[&'a str], overlap_budget: u32, overlap_cap: u32) -> Vec<&'a str> { + if cur.len() <= 1 { + return Vec::new(); + } + let mut acc = 0u32; + let mut take = 0usize; + for seg in cur.iter().rev() { + let t = conservative_token_estimate(seg); + // Cap EVERY trailing segment (including the first): if the last segment + // alone exceeds `overlap_cap`, carry no overlap rather than letting a + // near-`seg_budget` segment become the overlap. This keeps + // `overlap <= overlap_cap`, so `overlap + next_segment <= budget` and a + // flush can never emit an overlap-only (duplicate) chunk. + if acc.saturating_add(t) > overlap_cap { + break; + } + acc = acc.saturating_add(t); + take += 1; + if acc >= overlap_budget { + break; + } + } + if take >= cur.len() { + take = cur.len() - 1; // never repeat the whole chunk + } + cur[cur.len() - take..].to_vec() } #[cfg(test)] @@ -813,4 +962,140 @@ mod tests { let rejoined: String = chunks.iter().map(|c| c.content.as_str()).collect(); assert_eq!(rejoined, "abcdef"); } + + // ── Embed-safety: conservative-token splitting + overlap (#oversized-chunk) ── + + /// Every produced piece must be within the conservative token budget — this + /// is the property that keeps requests under bge-m3's batch/context limit. + fn assert_all_within_budget(pieces: &[String], budget: u32) { + for (i, p) in pieces.iter().enumerate() { + assert!( + conservative_token_estimate(p) <= budget, + "piece {i} is {} est-tokens, over budget {budget}", + conservative_token_estimate(p), + ); + } + } + + #[test] + fn dense_hash_content_splits_under_budget() { + // The exact failure mode: hash/path/punctuation-dense ASCII that chars/4 + // badly under-counts. A single 12 KB block of this overflowed bge-m3. + let dense = + "claude-memory:openhuman:MEMORY.md:67d6fe2727d431b16d41630babfdcf1cdf61bda7b9ba99656\n" + .repeat(200); + let pieces = split_by_token_budget(&dense, 3000); + assert!( + pieces.len() > 1, + "dense content must split, got {}", + pieces.len() + ); + assert_all_within_budget(&pieces, 3000); + } + + #[test] + fn hebrew_content_splits_under_budget() { + // Hebrew tokenises ~1 token/char — chars/4 under-counts ~4×. + let he = "איריס דיברה עם העורך דין ועם עידו על הדירה החדשה ".repeat(300); + let pieces = split_by_token_budget(&he, 1000); + assert!(pieces.len() > 1, "hebrew content must split"); + assert_all_within_budget(&pieces, 1000); + } + + #[test] + fn mixed_markdown_code_splits_under_budget() { + let md = "## Section\nSome prose here. And more prose follows.\n\n\ + ```rust\nfn f(x: u32) -> u32 { x * 4 + 3 }\n```\n\n\ + - bullet a1b2c3d4\n- bullet e5f6a7b8\n" + .repeat(120); + let pieces = split_by_token_budget(&md, 2000); + assert!(pieces.len() > 1, "mixed content must split"); + assert_all_within_budget(&pieces, 2000); + } + + #[test] + fn oversize_content_splits_into_multiple_ordered_pieces() { + let text = (0..50) + .map(|i| format!("Paragraph number {i} with several words of filler content here.")) + .collect::>() + .join("\n\n"); + let pieces = split_by_token_budget(&text, 80); + assert!( + pieces.len() > 2, + "expected several pieces, got {}", + pieces.len() + ); + assert_all_within_budget(&pieces, 80); + // Ordering preserved: "number 0" precedes "number 49" in the stream. + let joined = pieces.join("\n"); + let p0 = joined.find("number 0").expect("first paragraph present"); + let p49 = joined.find("number 49").expect("last paragraph present"); + assert!(p0 < p49, "ordering not preserved"); + } + + #[test] + fn adjacent_chunks_overlap_without_duplicating() { + let text = (0..40) + .map(|i| format!("Sentence {i} carries a unique token marker{i} inside it.")) + .collect::>() + .join(" "); + // Budget chosen so each sentence segment is well under overlap_cap + // (40% of budget); overlap is only carried for segments that fit the + // cap, so a tighter budget would (correctly) yield no overlap. + let pieces = split_by_token_budget(&text, 200); + assert!( + pieces.len() >= 3, + "expected several pieces, got {}", + pieces.len() + ); + assert_all_within_budget(&pieces, 200); + // Overlap: at least one marker appears in two adjacent chunks. + let overlap = pieces.windows(2).any(|w| { + (0..40).any(|i| { + let m = format!("marker{i} "); + w[0].contains(&m) && w[1].contains(&m) + }) + }); + assert!(overlap, "expected ~12% overlap between adjacent chunks"); + // No near-duplicates: adjacent chunks are never identical. + for w in pieces.windows(2) { + assert_ne!(w[0], w[1], "adjacent chunks must not be identical"); + } + } + + #[test] + fn normal_small_content_is_single_chunk_unchanged() { + // Behaviour preserved for already-small content: returned verbatim. + let text = "# Title\nA short paragraph that easily fits.\n\nAnother short one."; + let pieces = split_by_token_budget(text, 3000); + assert_eq!(pieces.len(), 1); + assert_eq!(pieces[0], text); + } + + #[test] + fn large_consecutive_segments_stay_within_budget() { + // A small paragraph packs with a large one, then a second large + // paragraph follows. Regression for the tail_overlap cap: previously the + // large tail segment (> overlap_cap) became the overlap, so + // `overlap + next_segment` exceeded the budget and emitted an + // over-budget chunk. Every chunk must stay within budget, and no two + // adjacent chunks may be identical. + let budget = 100; + let text = format!( + "{}\n\n{}\n\n{}", + "a".repeat(40), + "b".repeat(110), + "c".repeat(110), + ); + let pieces = split_by_token_budget(&text, budget); + assert!( + pieces.len() >= 2, + "expected multiple chunks, got {}", + pieces.len() + ); + assert_all_within_budget(&pieces, budget); + for w in pieces.windows(2) { + assert_ne!(w[0], w[1], "adjacent chunks must not be identical"); + } + } } diff --git a/src/openhuman/memory_store/chunks/types.rs b/src/openhuman/memory_store/chunks/types.rs index f96e8f872..9e8040959 100644 --- a/src/openhuman/memory_store/chunks/types.rs +++ b/src/openhuman/memory_store/chunks/types.rs @@ -299,6 +299,58 @@ pub fn approx_token_count(text: &str) -> u32 { chars.saturating_add(3) / 4 } +/// Per-character weight in **quarter-token** units for +/// [`conservative_token_estimate`]. Deliberately pessimistic so the chunker and +/// the embed backstop never under-split: real SentencePiece/WordPiece output for +/// hash-, code-, and markdown-dense text approaches ~1 token/char — far above +/// the `chars/4` GPT heuristic in [`approx_token_count`]. +fn char_token_quarters(ch: char) -> u32 { + if ch.is_ascii_alphanumeric() { + 2 // 0.50 token/char — alphanumeric runs pack ~2-4 chars per token + } else if ch.is_whitespace() { + 1 // 0.25 token/char — whitespace usually merges into adjacent pieces + } else { + 4 // 1.00 token/char — ASCII punctuation/symbols AND all non-ASCII + // (Hebrew/CJK/emoji), which tokenise ~1 piece per char or worse + } +} + +/// Conservative (over-estimating) token count, for embed-safety decisions only. +/// +/// [`approx_token_count`] (`chars/4`) under-counts dense markdown/hash/code by +/// ~5× and let oversized chunks reach bge-m3 (HTTP 500 "input length exceeds the +/// context length"). This weights characters by class so the result is an +/// upper-ish bound on real tokeniser output. It does **not** replace +/// `approx_token_count`, which still drives summariser/seal token budgeting. +pub fn conservative_token_estimate(text: &str) -> u32 { + let quarters: u64 = text + .chars() + .map(|c| u64::from(char_token_quarters(c))) + .sum(); + let tokens = (quarters + 3) / 4; // ceil(quarters / 4) + tokens.min(u64::from(u32::MAX)) as u32 +} + +/// Largest leading slice of `text` whose [`conservative_token_estimate`] is +/// ≤ `budget`, ending on a UTF-8 char boundary. Returns the whole string when +/// already within budget. Used as the embed-path backstop so an over-long body +/// can never be sent to the embedder above its input limit. +pub fn truncate_to_conservative_tokens(text: &str, budget: u32) -> &str { + if conservative_token_estimate(text) <= budget { + return text; + } + let cap = u64::from(budget).saturating_mul(4); // quarter-tokens + let mut acc: u64 = 0; + for (idx, ch) in text.char_indices() { + let q = u64::from(char_token_quarters(ch)); + if acc + q > cap { + return &text[..idx]; + } + acc += q; + } + text +} + mod time_range_serde { use chrono::{DateTime, TimeZone, Utc}; use serde::{Deserialize, Deserializer, Serialize, Serializer}; @@ -348,6 +400,43 @@ mod tests { assert_eq!(a.len(), 32); } + #[test] + fn conservative_estimate_weights_by_char_class() { + assert_eq!(conservative_token_estimate("abcd"), 2); // 4 alnum × 2q / 4 + assert_eq!(conservative_token_estimate(" "), 1); // 4 ws × 1q / 4 + assert_eq!(conservative_token_estimate("....,,,,"), 8); // 8 punct × 4q / 4 + assert_eq!(conservative_token_estimate("שלום"), 4); // 4 non-ascii × 4q / 4 + assert_eq!(conservative_token_estimate(""), 0); + } + + #[test] + fn conservative_estimate_exceeds_approx_for_dense_content() { + // Hash/path/punctuation-dense ASCII — exactly the content that defeated + // the chars/4 heuristic and overflowed bge-m3. + let dense = + "claude-memory:openhuman:MEMORY.md:67d6fe2727d431b16d41630babfdcf1cdf61bda7b9ba\n" + .repeat(40); + assert!( + conservative_token_estimate(&dense) > approx_token_count(&dense), + "conservative estimate must exceed chars/4 on dense content", + ); + } + + #[test] + fn truncate_respects_budget_and_char_boundaries() { + let text = "שלום עולם ".repeat(100); // Hebrew, ~1 token/char + let out = truncate_to_conservative_tokens(&text, 10); + assert!(conservative_token_estimate(out) <= 10); + assert!(text.starts_with(out)); // valid prefix on a char boundary + assert!(out.len() < text.len()); + } + + #[test] + fn truncate_is_noop_within_budget() { + let text = "short and sweet"; + assert_eq!(truncate_to_conservative_tokens(text, 1000), text); + } + #[test] fn chunk_id_varies_with_seq() { let a = chunk_id(SourceKind::Chat, "slack:#eng", 0, "hello"); diff --git a/src/openhuman/memory_tree/score/embed/ollama.rs b/src/openhuman/memory_tree/score/embed/ollama.rs index 60c1e24c1..dd7c19a91 100644 --- a/src/openhuman/memory_tree/score/embed/ollama.rs +++ b/src/openhuman/memory_tree/score/embed/ollama.rs @@ -162,6 +162,14 @@ struct EmbedRequest<'a> { #[derive(Serialize)] struct EmbedOptions { num_ctx: u32, + // Ollama/llama.cpp evaluates an embedding prompt in a SINGLE batch and + // rejects (HTTP 500 "input length exceeds the context length") any prompt + // longer than `n_batch` — whose default is 2048, far below `num_ctx`. + // Raising `num_ctx` alone does NOT lift this; `num_batch` must match, or + // chunks between ~2k and 8k real tokens fail to embed even though they fit + // the context window. Set equal to `num_ctx` so the batch limit never bites + // before the real context limit. + num_batch: u32, } #[derive(Deserialize)] @@ -188,6 +196,7 @@ impl Embedder for OllamaEmbedder { prompt: text, options: EmbedOptions { num_ctx: EMBED_NUM_CTX, + num_batch: EMBED_NUM_CTX, }, }; let resp = self @@ -293,6 +302,7 @@ mod tests { assert_eq!(body["model"], "bge-m3"); assert_eq!(body["prompt"], "hello world"); assert_eq!(body["options"]["num_ctx"], 8192); + assert_eq!(body["options"]["num_batch"], 8192); Json(serde_json::json!({ "embedding": v })) } }),