From cd3a6a227be8c1ad5c651a7960f165641f8c1361 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Tue, 10 Mar 2026 16:24:11 +0100 Subject: [PATCH] fix: strip null bytes from parsed file content before retain (#535) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * doc: split blog index into Hindsight and Hindsight Cloud sections - Tag the document upload post with `hindsight-cloud` - BlogListPage renders two sections, capping Cloud at 3 posts with a "View all →" link - Swizzle BlogTagsPostsPage so /blog/tags/hindsight-cloud uses the custom grid layout * doc: attribute blog posts to Nicolò Boschi with GitHub profile image Replace the generic "Hindsight Team" author with the real author entry (nicoloboschi) across all 15 blog posts. GitHub profile image is loaded from https://github.com/nicoloboschi.png. * doc: add Hindsight Team title to nicoloboschi author * doc: assign blog posts to correct authors based on git blame - Add benfrank241 (Ben Bartholomew) and chrislatimer (Chris Latimer) to authors.yml - Assign 7 posts to Ben, 1 post to Chris, remainder stay with Nicolò * fix: strip null bytes from parsed file content before retain * test: add tests for sanitize_llm_output * fix: retry retain DB transaction on deadlock during parallel document processing --- .../hindsight_api/engine/db_utils.py | 14 +- .../hindsight_api/engine/memory_engine.py | 4 +- .../engine/retain/orchestrator.py | 454 +++++++++--------- hindsight-api/tests/test_llm_wrapper.py | 37 ++ 4 files changed, 282 insertions(+), 227 deletions(-) create mode 100644 hindsight-api/tests/test_llm_wrapper.py diff --git a/hindsight-api/hindsight_api/engine/db_utils.py b/hindsight-api/hindsight_api/engine/db_utils.py index 475262dc..ffb547a4 100644 --- a/hindsight-api/hindsight_api/engine/db_utils.py +++ b/hindsight-api/hindsight_api/engine/db_utils.py @@ -58,10 +58,16 @@ async def retry_with_backoff( last_exception = e if attempt < max_retries: delay = min(base_delay * (2**attempt), max_delay) - logger.warning( - f"Database operation failed (attempt {attempt + 1}/{max_retries + 1}): {e}. " - f"Retrying in {delay:.1f}s..." - ) + if isinstance(e, asyncpg.exceptions.DeadlockDetectedError): + logger.warning( + f"Deadlock detected during parallel document processing — this is expected and will resolve automatically " + f"(attempt {attempt + 1}/{max_retries + 1}, retrying in {delay:.1f}s)" + ) + else: + logger.warning( + f"Database operation failed (attempt {attempt + 1}/{max_retries + 1}): {e}. " + f"Retrying in {delay:.1f}s..." + ) await asyncio.sleep(delay) else: logger.error(f"Database operation failed after {max_retries + 1} attempts: {e}") diff --git a/hindsight-api/hindsight_api/engine/memory_engine.py b/hindsight-api/hindsight_api/engine/memory_engine.py index cbd1414b..ad0de183 100644 --- a/hindsight-api/hindsight_api/engine/memory_engine.py +++ b/hindsight-api/hindsight_api/engine/memory_engine.py @@ -168,7 +168,7 @@ from enum import Enum from ..metrics import get_metrics_collector from ..pg0 import EmbeddedPostgres, parse_pg0_url from .entity_resolver import EntityResolver -from .llm_wrapper import LLMConfig, requires_api_key +from .llm_wrapper import LLMConfig, requires_api_key, sanitize_llm_output from .query_analyzer import QueryAnalyzer from .reflect import run_reflect_agent from .reflect.tools import tool_expand, tool_recall, tool_search_mental_models, tool_search_observations @@ -664,7 +664,7 @@ class MemoryEngine(MemoryEngineInterface): filename=filename, content_type=task_dict.get("content_type"), ) - markdown_content = convert_result.content + markdown_content = sanitize_llm_output(convert_result.content) or "" winning_parser = convert_result.parser_name except Exception as e: # Re-raise with filename context for better error reporting diff --git a/hindsight-api/hindsight_api/engine/retain/orchestrator.py b/hindsight-api/hindsight_api/engine/retain/orchestrator.py index 356f53a9..3bc54c77 100644 --- a/hindsight-api/hindsight_api/engine/retain/orchestrator.py +++ b/hindsight-api/hindsight_api/engine/retain/orchestrator.py @@ -11,7 +11,7 @@ from collections.abc import Awaitable, Callable from datetime import UTC, datetime from typing import Any -from ..db_utils import acquire_with_retry +from ..db_utils import acquire_with_retry, retry_with_backoff from . import bank_utils @@ -269,9 +269,6 @@ async def retain_batch( for extracted_fact, embedding in zip(extracted_facts, embeddings) ] - # Track document IDs for logging - document_ids_added = [] - # Group contents by document_id for document tracking and chunk storage from collections import defaultdict @@ -280,234 +277,249 @@ async def retain_batch( doc_id = content_dict.get("document_id") contents_by_doc[doc_id].append((idx, content_dict)) - # Step 4: Database transaction - async with acquire_with_retry(pool) as conn: - async with conn.transaction(): - # Handle document tracking for all documents - step_start = time.time() - # Map None document_id to generated UUIDs - doc_id_mapping = {} # Maps original doc_id (including None) to actual doc_id used + # Step 4: Database transaction (retried on deadlock) + result_unit_ids: list[list[str]] = [] - if document_id: - # Legacy: single document_id parameter - combined_content = "\n".join([c.get("content", "") for c in contents_dicts]) - retain_params = {} - # Collect tags from all content items and merge with document_tags - all_tags = set(document_tags or []) - for item in contents_dicts: - item_tags = item.get("tags", []) or [] - all_tags.update(item_tags) - merged_tags = list(all_tags) + log_buffer_pre_db = len(log_buffer) - if contents_dicts: - first_item = contents_dicts[0] - if first_item.get("context"): - retain_params["context"] = first_item["context"] - if first_item.get("event_date"): - retain_params["event_date"] = ( - first_item["event_date"].isoformat() - if hasattr(first_item["event_date"], "isoformat") - else str(first_item["event_date"]) - ) - if first_item.get("metadata"): - retain_params["metadata"] = first_item["metadata"] + async def _run_db_work() -> None: + nonlocal result_unit_ids - await fact_storage.handle_document_tracking( - conn, bank_id, document_id, combined_content, is_first_batch, retain_params, merged_tags - ) - document_ids_added.append(document_id) - doc_id_mapping[None] = document_id # For backwards compatibility - else: - # Handle per-item document_ids (create documents if any item has document_id or if chunks exist) - has_any_doc_ids = any(item.get("document_id") for item in contents_dicts) + # Reset per-fact mutations and log buffer so each retry attempt starts clean + del log_buffer[log_buffer_pre_db:] + document_ids_added: list[str] = [] + for pf in processed_facts: + pf.document_id = None + pf.chunk_id = None - if has_any_doc_ids or chunks: - for original_doc_id, doc_contents in contents_by_doc.items(): - actual_doc_id = original_doc_id + async with acquire_with_retry(pool) as conn: + async with conn.transaction(): + # Handle document tracking for all documents + step_start = time.time() + # Map None document_id to generated UUIDs + doc_id_mapping = {} # Maps original doc_id (including None) to actual doc_id used - # Only create document record if: - # 1. Item has explicit document_id, OR - # 2. There are chunks (need document for chunk storage) - should_create_doc = (original_doc_id is not None) or chunks + if document_id: + # Legacy: single document_id parameter + combined_content = "\n".join([c.get("content", "") for c in contents_dicts]) + retain_params = {} + # Collect tags from all content items and merge with document_tags + all_tags = set(document_tags or []) + for item in contents_dicts: + item_tags = item.get("tags", []) or [] + all_tags.update(item_tags) + merged_tags = list(all_tags) - if should_create_doc: - if actual_doc_id is None: - # No document_id but have chunks - generate one - actual_doc_id = str(uuid.uuid4()) - - # Store mapping for later use - doc_id_mapping[original_doc_id] = actual_doc_id - - # Combine content for this document - combined_content = "\n".join([c.get("content", "") for _, c in doc_contents]) - - # Collect tags from all content items for this document and merge with document_tags - all_tags = set(document_tags or []) - for _, item in doc_contents: - item_tags = item.get("tags", []) or [] - all_tags.update(item_tags) - merged_tags = list(all_tags) - - # Extract retain params from first content item - retain_params = {} - if doc_contents: - first_item = doc_contents[0][1] - if first_item.get("context"): - retain_params["context"] = first_item["context"] - if first_item.get("event_date"): - retain_params["event_date"] = ( - first_item["event_date"].isoformat() - if hasattr(first_item["event_date"], "isoformat") - else str(first_item["event_date"]) - ) - if first_item.get("metadata"): - retain_params["metadata"] = first_item["metadata"] - - await fact_storage.handle_document_tracking( - conn, - bank_id, - actual_doc_id, - combined_content, - is_first_batch, - retain_params, - merged_tags, + if contents_dicts: + first_item = contents_dicts[0] + if first_item.get("context"): + retain_params["context"] = first_item["context"] + if first_item.get("event_date"): + retain_params["event_date"] = ( + first_item["event_date"].isoformat() + if hasattr(first_item["event_date"], "isoformat") + else str(first_item["event_date"]) ) - document_ids_added.append(actual_doc_id) + if first_item.get("metadata"): + retain_params["metadata"] = first_item["metadata"] + await fact_storage.handle_document_tracking( + conn, bank_id, document_id, combined_content, is_first_batch, retain_params, merged_tags + ) + document_ids_added.append(document_id) + doc_id_mapping[None] = document_id # For backwards compatibility + else: + # Handle per-item document_ids (create documents if any item has document_id or if chunks exist) + has_any_doc_ids = any(item.get("document_id") for item in contents_dicts) + + if has_any_doc_ids or chunks: + for original_doc_id, doc_contents in contents_by_doc.items(): + actual_doc_id = original_doc_id + + # Only create document record if: + # 1. Item has explicit document_id, OR + # 2. There are chunks (need document for chunk storage) + should_create_doc = (original_doc_id is not None) or chunks + + if should_create_doc: + if actual_doc_id is None: + # No document_id but have chunks - generate one + actual_doc_id = str(uuid.uuid4()) + + # Store mapping for later use + doc_id_mapping[original_doc_id] = actual_doc_id + + # Combine content for this document + combined_content = "\n".join([c.get("content", "") for _, c in doc_contents]) + + # Collect tags from all content items for this document and merge with document_tags + all_tags = set(document_tags or []) + for _, item in doc_contents: + item_tags = item.get("tags", []) or [] + all_tags.update(item_tags) + merged_tags = list(all_tags) + + # Extract retain params from first content item + retain_params = {} + if doc_contents: + first_item = doc_contents[0][1] + if first_item.get("context"): + retain_params["context"] = first_item["context"] + if first_item.get("event_date"): + retain_params["event_date"] = ( + first_item["event_date"].isoformat() + if hasattr(first_item["event_date"], "isoformat") + else str(first_item["event_date"]) + ) + if first_item.get("metadata"): + retain_params["metadata"] = first_item["metadata"] + + await fact_storage.handle_document_tracking( + conn, + bank_id, + actual_doc_id, + combined_content, + is_first_batch, + retain_params, + merged_tags, + ) + document_ids_added.append(actual_doc_id) + + if document_ids_added: + log_buffer.append( + f"[2.5] Document tracking: {len(document_ids_added)} documents in {time.time() - step_start:.3f}s" + ) + + # Store chunks and map to facts for all documents + step_start = time.time() + chunk_id_map_by_doc = {} # Maps (doc_id, chunk_index) -> chunk_id + + if chunks: + # Group chunks by their source document + chunks_by_doc = defaultdict(list) + for chunk in chunks: + # chunk.content_index tells us which content this chunk came from + original_doc_id = contents_dicts[chunk.content_index].get("document_id") + # Map to actual document_id (handles None -> generated UUID mapping) + actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) + if actual_doc_id is None and document_id: + actual_doc_id = document_id + chunks_by_doc[actual_doc_id].append(chunk) + + # Store chunks for each document + for doc_id, doc_chunks in chunks_by_doc.items(): + chunk_id_map = await chunk_storage.store_chunks_batch(conn, bank_id, doc_id, doc_chunks) + # Store mapping with document context + for chunk_idx, chunk_id in chunk_id_map.items(): + chunk_id_map_by_doc[(doc_id, chunk_idx)] = chunk_id + + log_buffer.append( + f"[3] Store chunks: {len(chunks)} chunks for {len(chunks_by_doc)} documents in {time.time() - step_start:.3f}s" + ) + + # Map chunk_ids and document_ids to facts + for fact, processed_fact in zip(extracted_facts, processed_facts): + # Get the original document_id for this fact's source content + original_doc_id = contents_dicts[fact.content_index].get("document_id") + # Map to actual document_id (handles None -> generated UUID mapping) + actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) + if actual_doc_id is None and document_id: + actual_doc_id = document_id + + # Set document_id on the fact + processed_fact.document_id = actual_doc_id + + # Map chunk_id if this fact came from a chunk + if fact.chunk_index is not None: + # Look up chunk_id using (doc_id, chunk_index) + chunk_id = chunk_id_map_by_doc.get((actual_doc_id, fact.chunk_index)) + if chunk_id: + processed_fact.chunk_id = chunk_id + else: + # No chunks - still need to set document_id on facts + for fact, processed_fact in zip(extracted_facts, processed_facts): + original_doc_id = contents_dicts[fact.content_index].get("document_id") + # Map to actual document_id (handles None -> generated UUID mapping) + actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) + if actual_doc_id is None and document_id: + actual_doc_id = document_id + processed_fact.document_id = actual_doc_id + + non_duplicate_facts = processed_facts + + # Insert facts (document_id is now stored per-fact) + step_start = time.time() + unit_ids = await fact_storage.insert_facts_batch(conn, bank_id, non_duplicate_facts) + log_buffer.append(f"[5] Insert facts: {len(unit_ids)} units in {time.time() - step_start:.3f}s") + + # Process entities + step_start = time.time() + # Build map of content_index -> user entities for merging + user_entities_per_content = { + idx: content.entities for idx, content in enumerate(contents) if content.entities + } + entity_links = await entity_processing.process_entities_batch( + entity_resolver, + conn, + bank_id, + unit_ids, + non_duplicate_facts, + log_buffer, + user_entities_per_content=user_entities_per_content, + entity_labels=getattr(config, "entity_labels", None), + ) + log_buffer.append(f"[6] Process entities: {len(entity_links)} links in {time.time() - step_start:.3f}s") + + # Create temporal links + step_start = time.time() + temporal_link_count = await link_creation.create_temporal_links_batch(conn, bank_id, unit_ids) + log_buffer.append(f"[7] Temporal links: {temporal_link_count} links in {time.time() - step_start:.3f}s") + + # Create semantic links + step_start = time.time() + embeddings_for_links = [fact.embedding for fact in non_duplicate_facts] + semantic_link_count = await link_creation.create_semantic_links_batch( + conn, bank_id, unit_ids, embeddings_for_links + ) + log_buffer.append(f"[8] Semantic links: {semantic_link_count} links in {time.time() - step_start:.3f}s") + + # Insert entity links + step_start = time.time() + if entity_links: + await entity_processing.insert_entity_links_batch(conn, entity_links) + log_buffer.append( + f"[9] Entity links: {len(entity_links) if entity_links else 0} links in {time.time() - step_start:.3f}s" + ) + + # Create causal links + step_start = time.time() + causal_link_count = await link_creation.create_causal_links_batch(conn, unit_ids, non_duplicate_facts) + log_buffer.append(f"[10] Causal links: {causal_link_count} links in {time.time() - step_start:.3f}s") + + # Map results back to original content items + result_unit_ids = _map_results_to_contents(contents, extracted_facts, unit_ids) + + # Transactional outbox: queue any side-effect tasks (e.g. webhook deliveries) + # inside the same transaction so they are atomically committed with the retain data. + if outbox_callback: + await outbox_callback(conn) + + # Flush entity stats (mention_count / last_seen) now that the transaction + # has committed. Uses a fresh pool connection — no locks held. + await entity_resolver.flush_pending_stats() + + # Log final summary + total_time = time.time() - start_time + log_buffer.append(f"{'=' * 60}") + log_buffer.append(f"RETAIN_BATCH COMPLETE: {len(unit_ids)} units in {total_time:.3f}s") if document_ids_added: - log_buffer.append( - f"[2.5] Document tracking: {len(document_ids_added)} documents in {time.time() - step_start:.3f}s" - ) + log_buffer.append(f"Documents: {', '.join(document_ids_added)}") + log_buffer.append(f"{'=' * 60}") - # Store chunks and map to facts for all documents - step_start = time.time() - chunk_id_map_by_doc = {} # Maps (doc_id, chunk_index) -> chunk_id + logger.info("\n" + "\n".join(log_buffer) + "\n") - if chunks: - # Group chunks by their source document - chunks_by_doc = defaultdict(list) - for chunk in chunks: - # chunk.content_index tells us which content this chunk came from - original_doc_id = contents_dicts[chunk.content_index].get("document_id") - # Map to actual document_id (handles None -> generated UUID mapping) - actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) - if actual_doc_id is None and document_id: - actual_doc_id = document_id - chunks_by_doc[actual_doc_id].append(chunk) - - # Store chunks for each document - for doc_id, doc_chunks in chunks_by_doc.items(): - chunk_id_map = await chunk_storage.store_chunks_batch(conn, bank_id, doc_id, doc_chunks) - # Store mapping with document context - for chunk_idx, chunk_id in chunk_id_map.items(): - chunk_id_map_by_doc[(doc_id, chunk_idx)] = chunk_id - - log_buffer.append( - f"[3] Store chunks: {len(chunks)} chunks for {len(chunks_by_doc)} documents in {time.time() - step_start:.3f}s" - ) - - # Map chunk_ids and document_ids to facts - for fact, processed_fact in zip(extracted_facts, processed_facts): - # Get the original document_id for this fact's source content - original_doc_id = contents_dicts[fact.content_index].get("document_id") - # Map to actual document_id (handles None -> generated UUID mapping) - actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) - if actual_doc_id is None and document_id: - actual_doc_id = document_id - - # Set document_id on the fact - processed_fact.document_id = actual_doc_id - - # Map chunk_id if this fact came from a chunk - if fact.chunk_index is not None: - # Look up chunk_id using (doc_id, chunk_index) - chunk_id = chunk_id_map_by_doc.get((actual_doc_id, fact.chunk_index)) - if chunk_id: - processed_fact.chunk_id = chunk_id - else: - # No chunks - still need to set document_id on facts - for fact, processed_fact in zip(extracted_facts, processed_facts): - original_doc_id = contents_dicts[fact.content_index].get("document_id") - # Map to actual document_id (handles None -> generated UUID mapping) - actual_doc_id = doc_id_mapping.get(original_doc_id, original_doc_id) - if actual_doc_id is None and document_id: - actual_doc_id = document_id - processed_fact.document_id = actual_doc_id - - non_duplicate_facts = processed_facts - - # Insert facts (document_id is now stored per-fact) - step_start = time.time() - unit_ids = await fact_storage.insert_facts_batch(conn, bank_id, non_duplicate_facts) - log_buffer.append(f"[5] Insert facts: {len(unit_ids)} units in {time.time() - step_start:.3f}s") - - # Process entities - step_start = time.time() - # Build map of content_index -> user entities for merging - user_entities_per_content = { - idx: content.entities for idx, content in enumerate(contents) if content.entities - } - entity_links = await entity_processing.process_entities_batch( - entity_resolver, - conn, - bank_id, - unit_ids, - non_duplicate_facts, - log_buffer, - user_entities_per_content=user_entities_per_content, - entity_labels=getattr(config, "entity_labels", None), - ) - log_buffer.append(f"[6] Process entities: {len(entity_links)} links in {time.time() - step_start:.3f}s") - - # Create temporal links - step_start = time.time() - temporal_link_count = await link_creation.create_temporal_links_batch(conn, bank_id, unit_ids) - log_buffer.append(f"[7] Temporal links: {temporal_link_count} links in {time.time() - step_start:.3f}s") - - # Create semantic links - step_start = time.time() - embeddings_for_links = [fact.embedding for fact in non_duplicate_facts] - semantic_link_count = await link_creation.create_semantic_links_batch( - conn, bank_id, unit_ids, embeddings_for_links - ) - log_buffer.append(f"[8] Semantic links: {semantic_link_count} links in {time.time() - step_start:.3f}s") - - # Insert entity links - step_start = time.time() - if entity_links: - await entity_processing.insert_entity_links_batch(conn, entity_links) - log_buffer.append( - f"[9] Entity links: {len(entity_links) if entity_links else 0} links in {time.time() - step_start:.3f}s" - ) - - # Create causal links - step_start = time.time() - causal_link_count = await link_creation.create_causal_links_batch(conn, unit_ids, non_duplicate_facts) - log_buffer.append(f"[10] Causal links: {causal_link_count} links in {time.time() - step_start:.3f}s") - - # Map results back to original content items - result_unit_ids = _map_results_to_contents(contents, extracted_facts, unit_ids) - - # Transactional outbox: queue any side-effect tasks (e.g. webhook deliveries) - # inside the same transaction so they are atomically committed with the retain data. - if outbox_callback: - await outbox_callback(conn) - - # Flush entity stats (mention_count / last_seen) now that the transaction - # has committed. Uses a fresh pool connection — no locks held. - await entity_resolver.flush_pending_stats() - - # Log final summary - total_time = time.time() - start_time - log_buffer.append(f"{'=' * 60}") - log_buffer.append(f"RETAIN_BATCH COMPLETE: {len(unit_ids)} units in {total_time:.3f}s") - if document_ids_added: - log_buffer.append(f"Documents: {', '.join(document_ids_added)}") - log_buffer.append(f"{'=' * 60}") - - logger.info("\n" + "\n".join(log_buffer) + "\n") - - return result_unit_ids, usage + await retry_with_backoff(_run_db_work) + return result_unit_ids, usage def _map_results_to_contents( diff --git a/hindsight-api/tests/test_llm_wrapper.py b/hindsight-api/tests/test_llm_wrapper.py new file mode 100644 index 00000000..ac5d2c0a --- /dev/null +++ b/hindsight-api/tests/test_llm_wrapper.py @@ -0,0 +1,37 @@ +import pytest + +from hindsight_api.engine.llm_wrapper import sanitize_llm_output + + +@pytest.mark.parametrize( + "input_text, expected", + [ + # Null bytes stripped + ("hello\x00world", "helloworld"), + ("FIRST\u0000PAGE", "FIRSTPAGE"), + # Multiple null bytes + ("\x00\x00text\x00", "text"), + # Other control characters stripped (non-whitespace) + ("text\x01\x02\x03end", "textend"), + ("text\x08end", "textend"), # backspace + ("text\x0cend", "textend"), # form feed + ("text\x0bend", "textend"), # vertical tab + ("text\x1fend", "textend"), # unit separator + ("text\x7fend", "textend"), # DEL + # Whitespace preserved + ("hello\tworld", "hello\tworld"), + ("hello\nworld", "hello\nworld"), + ("hello\r\nworld", "hello\r\nworld"), + # Unicode surrogates stripped + ("text\ud800end", "textend"), + ("text\udfffend", "textend"), + # Clean text unchanged + ("normal text", "normal text"), + ("unicode: café naïve", "unicode: café naïve"), + # Edge cases + ("", ""), + (None, None), + ], +) +def test_sanitize_llm_output(input_text, expected): + assert sanitize_llm_output(input_text) == expected