chopratejas Claude Opus 4.6 (1M context) commited on
Commit
1e808d1
·
1 Parent(s): 75a9bd3

Reduce compression latency: cache serializations, eager-load all compressors, fix Magika, bump to 0.5.5

Browse files

SmartCrusher: eliminate 5-7x redundant json.dumps by threading cached item_strings
through _crush_array → _create_plan → _plan_* methods, TOIN token counting, and CCR
storage. Move ISO datetime regex to module level. Cache field name hashes in TOIN
semantic detection. Add item_strings param to error detection.

ContentRouter: compile prose detection regex at module level. Extend
eager_load_compressors() to pre-load Magika detector, tree-sitter parsers (8 common
languages), CodeAwareCompressor, and SmartCrusher at startup.

Magika: add as proxy dependency (was never declared in pyproject.toml). Update
detector.py for Magika 1.x API (result.output.label, result.score). Fix batch
detection to use identify_bytes loop (identify_bytes_batch removed in 1.x).

Proxy: simplify startup to use eager_load_compressors() return status dict for
unified component logging.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

headroom/__init__.py CHANGED
@@ -153,7 +153,7 @@ from .transforms import (
153
  TransformPipeline,
154
  )
155
 
156
- __version__ = "0.5.4"
157
 
158
  __all__ = [
159
  # Main client
 
153
  TransformPipeline,
154
  )
155
 
156
+ __version__ = "0.5.5"
157
 
158
  __all__ = [
159
  # Main client
headroom/compression/detector.py CHANGED
@@ -209,8 +209,8 @@ class MagikaDetector:
209
  magika = self._ensure_magika()
210
  result: MagikaResult = magika.identify_bytes(content.encode("utf-8"))
211
 
212
- raw_label = result.output.ct_label
213
- confidence = result.output.score
214
 
215
  # Map to our content type
216
  content_type, language = self._map_label(raw_label)
@@ -233,8 +233,6 @@ class MagikaDetector:
233
  def detect_batch(self, contents: list[str]) -> list[DetectionResult]:
234
  """Detect content types for multiple contents.
235
 
236
- More efficient than calling detect() in a loop.
237
-
238
  Args:
239
  contents: List of content strings to analyze.
240
 
@@ -244,16 +242,9 @@ class MagikaDetector:
244
  if not contents:
245
  return []
246
 
247
- magika = self._ensure_magika()
248
  results = []
249
 
250
- # Convert to bytes for Magika
251
- byte_contents = [c.encode("utf-8") for c in contents]
252
-
253
- # Batch detection
254
- magika_results = magika.identify_bytes_batch(byte_contents)
255
-
256
- for content, magika_result in zip(contents, magika_results):
257
  if not content or not content.strip():
258
  results.append(
259
  DetectionResult(
@@ -264,8 +255,9 @@ class MagikaDetector:
264
  )
265
  continue
266
 
267
- raw_label = magika_result.output.ct_label
268
- confidence = magika_result.output.score
 
269
  content_type, language = self._map_label(raw_label)
270
 
271
  if confidence < self.min_confidence:
 
209
  magika = self._ensure_magika()
210
  result: MagikaResult = magika.identify_bytes(content.encode("utf-8"))
211
 
212
+ raw_label = result.output.label
213
+ confidence = result.score
214
 
215
  # Map to our content type
216
  content_type, language = self._map_label(raw_label)
 
233
  def detect_batch(self, contents: list[str]) -> list[DetectionResult]:
234
  """Detect content types for multiple contents.
235
 
 
 
236
  Args:
237
  contents: List of content strings to analyze.
238
 
 
242
  if not contents:
243
  return []
244
 
 
245
  results = []
246
 
247
+ for content in contents:
 
 
 
 
 
 
248
  if not content or not content.strip():
249
  results.append(
250
  DetectionResult(
 
255
  )
256
  continue
257
 
258
+ magika_result = self._ensure_magika().identify_bytes(content.encode("utf-8"))
259
+ raw_label = magika_result.output.label
260
+ confidence = magika_result.score
261
  content_type, language = self._map_label(raw_label)
262
 
263
  if confidence < self.min_confidence:
headroom/proxy/server.py CHANGED
@@ -1995,36 +1995,32 @@ class HeadroomProxy:
1995
  else:
1996
  logger.info("Smart Routing: DISABLED (legacy sequential mode)")
1997
 
1998
- # Eagerly load ML compressors at startup (avoids download on first request)
1999
- # Kompress requires [ml] extra (torch + transformers). If not installed, skip.
2000
  self._kompress_status = "not installed"
2001
- from headroom.transforms.kompress_compressor import is_kompress_available
2002
 
2003
- if is_kompress_available() and self.config.optimize:
2004
- logger.info("Kompress: Downloading model (first-time only)...")
2005
  for transform in self.anthropic_pipeline.transforms:
2006
  if hasattr(transform, "eager_load_compressors"):
2007
- transform.eager_load_compressors()
2008
- self._kompress_status = "enabled"
2009
  break
2010
- if self._kompress_status == "enabled":
2011
- logger.info("Kompress: ENABLED (ModernBERT token compressor)")
2012
- else:
2013
- if self.config.optimize:
2014
- logger.info(
2015
- "Kompress: not installed (pip install headroom-ai[ml] for ML compression)"
2016
- )
2017
 
2018
- # LLMLingua fallback (only loads if Kompress is not available)
2019
- if self._kompress_status != "enabled" and self.config.llmlingua_enabled:
2020
- for transform in self.anthropic_pipeline.transforms:
2021
- if hasattr(transform, "_get_llmlingua"):
2022
- llmlingua = transform._get_llmlingua()
2023
- if llmlingua:
2024
- self._llmlingua_status = "enabled"
2025
- break
 
 
 
 
 
2026
 
2027
- # LLMLingua status
2028
  if self._llmlingua_status == "enabled":
2029
  logger.info(
2030
  f"LLMLingua: ENABLED (device={self.config.llmlingua_device}, "
@@ -2037,9 +2033,10 @@ class HeadroomProxy:
2037
  elif self._llmlingua_status == "disabled":
2038
  logger.info("LLMLingua: DISABLED")
2039
 
2040
- # Code-aware status
2041
  if self._code_aware_status == "enabled":
2042
  logger.info("Code-Aware: ENABLED (AST-based compression)")
 
 
2043
  elif self._code_aware_status == "lazy":
2044
  logger.info("Code-Aware: LAZY (will load when code content detected)")
2045
  elif self._code_aware_status == "available":
@@ -2049,6 +2046,9 @@ class HeadroomProxy:
2049
  elif self._code_aware_status == "disabled":
2050
  logger.info("Code-Aware: DISABLED")
2051
 
 
 
 
2052
  # CCR status
2053
  ccr_features = []
2054
  if self.config.ccr_inject_tool:
 
1995
  else:
1996
  logger.info("Smart Routing: DISABLED (legacy sequential mode)")
1997
 
1998
+ # Eagerly load ALL compressors, parsers, and detectors at startup
1999
+ # This eliminates cold-start latency spikes on first requests
2000
  self._kompress_status = "not installed"
2001
+ eager_status: dict[str, str] = {}
2002
 
2003
+ if self.config.optimize:
2004
+ logger.info("Pre-loading compressors and parsers...")
2005
  for transform in self.anthropic_pipeline.transforms:
2006
  if hasattr(transform, "eager_load_compressors"):
2007
+ eager_status = transform.eager_load_compressors()
 
2008
  break
 
 
 
 
 
 
 
2009
 
2010
+ # Update internal status from eager loading results
2011
+ if eager_status.get("kompress") == "enabled":
2012
+ self._kompress_status = "enabled"
2013
+ if eager_status.get("llmlingua") == "enabled":
2014
+ self._llmlingua_status = "enabled"
2015
+ if eager_status.get("code_aware") == "enabled":
2016
+ self._code_aware_status = "enabled"
2017
+
2018
+ # Log component status
2019
+ if self._kompress_status == "enabled":
2020
+ logger.info("Kompress: ENABLED (ModernBERT token compressor)")
2021
+ elif self.config.optimize:
2022
+ logger.info("Kompress: not installed (pip install headroom-ai[ml] for ML compression)")
2023
 
 
2024
  if self._llmlingua_status == "enabled":
2025
  logger.info(
2026
  f"LLMLingua: ENABLED (device={self.config.llmlingua_device}, "
 
2033
  elif self._llmlingua_status == "disabled":
2034
  logger.info("LLMLingua: DISABLED")
2035
 
 
2036
  if self._code_aware_status == "enabled":
2037
  logger.info("Code-Aware: ENABLED (AST-based compression)")
2038
+ if "tree_sitter" in eager_status:
2039
+ logger.info(f"Tree-Sitter: {eager_status['tree_sitter']}")
2040
  elif self._code_aware_status == "lazy":
2041
  logger.info("Code-Aware: LAZY (will load when code content detected)")
2042
  elif self._code_aware_status == "available":
 
2046
  elif self._code_aware_status == "disabled":
2047
  logger.info("Code-Aware: DISABLED")
2048
 
2049
+ if eager_status.get("magika") == "enabled":
2050
+ logger.info("Magika: ENABLED (ML content detection)")
2051
+
2052
  # CCR status
2053
  ccr_features = []
2054
  if self.config.ccr_inject_tool:
headroom/transforms/content_router.py CHANGED
@@ -428,6 +428,7 @@ class ContentRouterConfig:
428
  _CODE_FENCE_PATTERN = re.compile(r"^```(\w*)\s*$", re.MULTILINE)
429
  _JSON_BLOCK_START = re.compile(r"^\s*[\[{]", re.MULTILINE)
430
  _SEARCH_RESULT_PATTERN = re.compile(r"^\S+:\d+:", re.MULTILINE)
 
431
 
432
 
433
  def is_mixed_content(content: str) -> bool:
@@ -442,7 +443,7 @@ def is_mixed_content(content: str) -> bool:
442
  indicators = {
443
  "has_code_fences": bool(_CODE_FENCE_PATTERN.search(content)),
444
  "has_json_blocks": bool(_JSON_BLOCK_START.search(content)),
445
- "has_prose": len(re.findall(r"[A-Z][a-z]+\s+\w+\s+\w+", content)) > 5,
446
  "has_search_results": bool(_SEARCH_RESULT_PATTERN.search(content)),
447
  }
448
 
@@ -1168,30 +1169,97 @@ class ContentRouter(Transform):
1168
  logger.debug("HTMLExtractor not available (install trafilatura)")
1169
  return self._html_extractor
1170
 
1171
- def eager_load_compressors(self) -> None:
1172
  """Pre-load compressors at startup to avoid first-request latency.
1173
 
1174
- Call this during proxy startup to load models (~5s)
1175
- before any requests arrive.
 
 
 
1176
  """
1177
- # Prefer Kompress (faster, smaller, better on structured data)
 
 
1178
  if self.config.enable_kompress:
1179
  compressor = self._get_kompress()
1180
  if compressor:
1181
  logger.info("Kompress model pre-loaded at startup")
1182
- return # No need to also load LLMLingua
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1183
 
1184
- if self.config.enable_llmlingua:
1185
- compressor = self._get_llmlingua()
1186
- if compressor:
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1187
  try:
1188
- from .llmlingua_compressor import _get_llmlingua_compressor
1189
-
1190
- device = compressor._resolve_device()
1191
- _get_llmlingua_compressor(compressor.config.model_name, device)
1192
- logger.info("LLMLingua model pre-loaded at startup")
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1193
  except Exception as e:
1194
- logger.warning("Failed to pre-load LLMLingua model: %s", e)
 
 
 
 
 
 
 
 
 
 
1195
 
1196
  def _get_kompress(self) -> Any:
1197
  """Get KompressCompressor (lazy load). Downloads from HuggingFace on first use."""
 
428
  _CODE_FENCE_PATTERN = re.compile(r"^```(\w*)\s*$", re.MULTILINE)
429
  _JSON_BLOCK_START = re.compile(r"^\s*[\[{]", re.MULTILINE)
430
  _SEARCH_RESULT_PATTERN = re.compile(r"^\S+:\d+:", re.MULTILINE)
431
+ _PROSE_PATTERN = re.compile(r"[A-Z][a-z]+\s+\w+\s+\w+")
432
 
433
 
434
  def is_mixed_content(content: str) -> bool:
 
443
  indicators = {
444
  "has_code_fences": bool(_CODE_FENCE_PATTERN.search(content)),
445
  "has_json_blocks": bool(_JSON_BLOCK_START.search(content)),
446
+ "has_prose": len(_PROSE_PATTERN.findall(content)) > 5,
447
  "has_search_results": bool(_SEARCH_RESULT_PATTERN.search(content)),
448
  }
449
 
 
1169
  logger.debug("HTMLExtractor not available (install trafilatura)")
1170
  return self._html_extractor
1171
 
1172
+ def eager_load_compressors(self) -> dict[str, str]:
1173
  """Pre-load compressors at startup to avoid first-request latency.
1174
 
1175
+ Call this during proxy startup to load models and parsers
1176
+ before any requests arrive. Eliminates cold-start latency spikes.
1177
+
1178
+ Returns:
1179
+ Dict of component name -> status string for logging.
1180
  """
1181
+ status: dict[str, str] = {}
1182
+
1183
+ # 1. ML text compressor: Kompress or LLMLingua fallback
1184
  if self.config.enable_kompress:
1185
  compressor = self._get_kompress()
1186
  if compressor:
1187
  logger.info("Kompress model pre-loaded at startup")
1188
+ status["kompress"] = "enabled"
1189
+ else:
1190
+ status["kompress"] = "unavailable"
1191
+ if "kompress" not in status or status["kompress"] != "enabled":
1192
+ if self.config.enable_llmlingua:
1193
+ compressor = self._get_llmlingua()
1194
+ if compressor:
1195
+ try:
1196
+ from .llmlingua_compressor import _get_llmlingua_compressor
1197
+
1198
+ device = compressor._resolve_device()
1199
+ _get_llmlingua_compressor(compressor.config.model_name, device)
1200
+ logger.info("LLMLingua model pre-loaded at startup")
1201
+ status["llmlingua"] = "enabled"
1202
+ except Exception as e:
1203
+ logger.warning("Failed to pre-load LLMLingua model: %s", e)
1204
+ status["llmlingua"] = f"failed: {e}"
1205
+
1206
+ # 2. Magika content detector (avoids 100-200ms on first content detection)
1207
+ try:
1208
+ from ..compression.detector import _get_magika, _magika_available
1209
 
1210
+ if _magika_available():
1211
+ _get_magika() # Initializes the singleton
1212
+ logger.info("Magika content detector pre-loaded at startup")
1213
+ status["magika"] = "enabled"
1214
+ else:
1215
+ status["magika"] = "not installed"
1216
+ except Exception as e:
1217
+ logger.debug("Magika pre-load skipped: %s", e)
1218
+ status["magika"] = "skipped"
1219
+
1220
+ # 3. CodeAware compressor + common tree-sitter parsers
1221
+ if self.config.enable_code_aware:
1222
+ code_compressor = self._get_code_compressor()
1223
+ if code_compressor:
1224
+ status["code_aware"] = "enabled"
1225
+ # Pre-load tree-sitter parsers for common languages
1226
+ # Each parser is ~50ms to load; doing it here avoids 500ms+ on first code hit
1227
  try:
1228
+ from .code_compressor import _check_tree_sitter_available, _get_parser
1229
+
1230
+ if _check_tree_sitter_available():
1231
+ common_languages = [
1232
+ "python",
1233
+ "javascript",
1234
+ "typescript",
1235
+ "go",
1236
+ "rust",
1237
+ "java",
1238
+ "c",
1239
+ "cpp",
1240
+ ]
1241
+ loaded = []
1242
+ for lang in common_languages:
1243
+ try:
1244
+ _get_parser(lang)
1245
+ loaded.append(lang)
1246
+ except (ValueError, ImportError):
1247
+ pass # Language not available, skip
1248
+ if loaded:
1249
+ logger.info("Tree-sitter parsers pre-loaded: %s", ", ".join(loaded))
1250
+ status["tree_sitter"] = f"loaded ({len(loaded)} languages)"
1251
  except Exception as e:
1252
+ logger.debug("Tree-sitter pre-load skipped: %s", e)
1253
+ status["tree_sitter"] = "skipped"
1254
+ else:
1255
+ status["code_aware"] = "not installed"
1256
+
1257
+ # 4. SmartCrusher (lightweight init, but ensures import + TOIN ready)
1258
+ smart_crusher = self._get_smart_crusher()
1259
+ if smart_crusher:
1260
+ status["smart_crusher"] = "ready"
1261
+
1262
+ return status
1263
 
1264
  def _get_kompress(self) -> Any:
1265
  """Get KompressCompressor (lazy load). Downloads from HuggingFace on first use."""
headroom/transforms/smart_crusher.py CHANGED
@@ -92,6 +92,10 @@ _HOSTNAME_PATTERN = re.compile(
92
  _QUOTED_STRING_PATTERN = re.compile(r"['\"]([^'\"]{1,50})['\"]") # Short quoted strings
93
  _EMAIL_PATTERN = re.compile(r"\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Z|a-z]{2,}\b")
94
 
 
 
 
 
95
 
96
  def extract_query_anchors(text: str) -> set[str]:
97
  """Extract query anchors from user text (legacy regex-based method).
@@ -628,7 +632,10 @@ def _detect_rare_status_values(items: list[dict], common_fields: set[str]) -> li
628
  _ERROR_KEYWORDS_FOR_PRESERVATION = ERROR_KEYWORDS
629
 
630
 
631
- def _detect_error_items_for_preservation(items: list[dict]) -> list[int]:
 
 
 
632
  """Detect items containing error keywords for PRESERVATION guarantee.
633
 
634
  This is NOT for crushability analysis - it's for ensuring ALL error items
@@ -636,6 +643,10 @@ def _detect_error_items_for_preservation(items: list[dict]) -> list[int]:
636
  are NEVER dropped, even if errors are common in the dataset.
637
 
638
  Uses keywords because error semantics are well-defined across domains.
 
 
 
 
639
  """
640
  error_indices: list[int] = []
641
 
@@ -643,9 +654,12 @@ def _detect_error_items_for_preservation(items: list[dict]) -> list[int]:
643
  if not isinstance(item, dict):
644
  continue
645
 
646
- # Serialize item to check all content
647
  try:
648
- item_str = json.dumps(item).lower()
 
 
 
649
  except Exception:
650
  continue
651
 
@@ -694,13 +708,18 @@ def _detect_items_by_learned_semantics(
694
  if not confident_semantics:
695
  return []
696
 
 
 
 
697
  for i, item in enumerate(items):
698
  if not isinstance(item, dict):
699
  continue
700
 
701
  for field_name, value in item.items():
702
- # Hash the field name to match TOIN's format
703
- field_hash = hashlib.sha256(field_name.encode()).hexdigest()[:8]
 
 
704
 
705
  if field_hash not in confident_semantics:
706
  continue
@@ -1080,9 +1099,9 @@ class SmartAnalyzer:
1080
 
1081
  Uses STRUCTURAL detection based on value format, not field names.
1082
  """
1083
- # Check string fields for ISO 8601 patterns
1084
- iso_datetime_pattern = re.compile(r"^\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}")
1085
- iso_date_pattern = re.compile(r"^\d{4}-\d{2}-\d{2}$")
1086
 
1087
  for name, stats in field_stats.items():
1088
  if stats.field_type == "string":
@@ -2464,12 +2483,14 @@ class SmartCrusher(Transform):
2464
  # Create compression plan with relevance scoring
2465
  # Pass TOIN preserve_fields so items with those fields get priority
2466
  # Pass effective_max_items for thread-safe compression
 
2467
  plan = self._create_plan(
2468
  analysis,
2469
  items,
2470
  query_context,
2471
  preserve_fields=toin_preserve_fields or None,
2472
  effective_max_items=effective_max_items,
 
2473
  )
2474
 
2475
  # Execute compression
@@ -2483,7 +2504,8 @@ class SmartCrusher(Transform):
2483
  and len(result) < len(items) # Only cache if compression actually happened
2484
  ):
2485
  store = self._get_compression_store()
2486
- original_json = json.dumps(items, default=str)
 
2487
  compressed_json = json.dumps(result, default=str)
2488
 
2489
  ccr_hash = store.store(
@@ -2521,8 +2543,8 @@ class SmartCrusher(Transform):
2521
 
2522
  # TOIN: Record compression event for cross-user learning
2523
  try:
2524
- # Calculate token counts (approximate)
2525
- original_tokens = len(json.dumps(items, default=str)) // 4
2526
  compressed_tokens = len(json.dumps(result, default=str)) // 4
2527
 
2528
  toin.record_compression(
@@ -2573,18 +2595,29 @@ class SmartCrusher(Transform):
2573
  # Universal JSON type handlers (string, number, mixed arrays)
2574
  # =================================================================
2575
 
2576
- def _compute_k_split(self, items: list, bias: float = 1.0) -> tuple[int, int, int, int]:
 
 
 
 
 
2577
  """Compute adaptive K split into first/last/importance slots.
2578
 
2579
  Uses the existing Kneedle-based adaptive_sizer for K_total, then
2580
  splits according to configurable first_fraction / last_fraction.
2581
 
 
 
 
 
 
2582
  Returns:
2583
  (k_total, k_first, k_last, k_importance)
2584
  """
2585
  from .adaptive_sizer import compute_optimal_k
2586
 
2587
- item_strings = [json.dumps(item, default=str) for item in items]
 
2588
  k_total = compute_optimal_k(
2589
  item_strings,
2590
  bias=bias,
@@ -2997,6 +3030,7 @@ class SmartCrusher(Transform):
2997
  query_context: str = "",
2998
  preserve_fields: list[str] | None = None,
2999
  effective_max_items: int | None = None,
 
3000
  ) -> CompressionPlan:
3001
  """Create a detailed compression plan using relevance scoring.
3002
 
@@ -3006,6 +3040,7 @@ class SmartCrusher(Transform):
3006
  query_context: Context string from user messages for relevance scoring.
3007
  preserve_fields: TOIN-learned fields that users commonly retrieve.
3008
  Items with values in these fields get higher priority.
 
3009
  effective_max_items: Thread-safe max items limit (defaults to config value).
3010
  """
3011
  # Use provided effective_max_items or fall back to config
@@ -3027,22 +3062,46 @@ class SmartCrusher(Transform):
3027
 
3028
  if analysis.recommended_strategy == CompressionStrategy.TIME_SERIES:
3029
  plan = self._plan_time_series(
3030
- analysis, items, plan, query_context, preserve_fields, max_items
 
 
 
 
 
 
3031
  )
3032
 
3033
  elif analysis.recommended_strategy == CompressionStrategy.CLUSTER_SAMPLE:
3034
  plan = self._plan_cluster_sample(
3035
- analysis, items, plan, query_context, preserve_fields, max_items
 
 
 
 
 
 
3036
  )
3037
 
3038
  elif analysis.recommended_strategy == CompressionStrategy.TOP_N:
3039
  plan = self._plan_top_n(
3040
- analysis, items, plan, query_context, preserve_fields, max_items
 
 
 
 
 
 
3041
  )
3042
 
3043
  else: # SMART_SAMPLE or NONE
3044
  plan = self._plan_smart_sample(
3045
- analysis, items, plan, query_context, preserve_fields, max_items
 
 
 
 
 
 
3046
  )
3047
 
3048
  return plan
@@ -3055,6 +3114,7 @@ class SmartCrusher(Transform):
3055
  query_context: str = "",
3056
  preserve_fields: list[str] | None = None,
3057
  max_items: int | None = None,
 
3058
  ) -> CompressionPlan:
3059
  """Plan compression for time series data.
3060
 
@@ -3111,7 +3171,12 @@ class SmartCrusher(Transform):
3111
 
3112
  # 5. Items with high relevance to query context (PROBABILISTIC semantic match)
3113
  if query_context:
3114
- item_strs = [json.dumps(item, default=str) for item in items]
 
 
 
 
 
3115
  scores = self._scorer.score_batch(item_strs, query_context)
3116
  for i, score in enumerate(scores):
3117
  if score.score >= self._relevance_threshold:
@@ -3138,6 +3203,7 @@ class SmartCrusher(Transform):
3138
  query_context: str = "",
3139
  preserve_fields: list[str] | None = None,
3140
  max_items: int | None = None,
 
3141
  ) -> CompressionPlan:
3142
  """Plan compression for clusterable data (like logs).
3143
 
@@ -3211,7 +3277,12 @@ class SmartCrusher(Transform):
3211
 
3212
  # 5. Items with high relevance to query context (PROBABILISTIC semantic match)
3213
  if query_context:
3214
- item_strs = [json.dumps(item, default=str) for item in items]
 
 
 
 
 
3215
  scores = self._scorer.score_batch(item_strs, query_context)
3216
  for i, score in enumerate(scores):
3217
  if score.score >= self._relevance_threshold:
@@ -3238,6 +3309,7 @@ class SmartCrusher(Transform):
3238
  query_context: str = "",
3239
  preserve_fields: list[str] | None = None,
3240
  max_items: int | None = None,
 
3241
  ) -> CompressionPlan:
3242
  """Plan compression for scored/ranked data.
3243
 
@@ -3269,7 +3341,13 @@ class SmartCrusher(Transform):
3269
 
3270
  if not score_field:
3271
  return self._plan_smart_sample(
3272
- analysis, items, plan, query_context, preserve_fields, effective_max
 
 
 
 
 
 
3273
  )
3274
 
3275
  plan.sort_field = score_field
@@ -3307,7 +3385,12 @@ class SmartCrusher(Transform):
3307
  # Only add items that are NOT already in top N but match the query strongly
3308
  # Use a higher threshold (0.5) since the score field already captures relevance
3309
  if query_context:
3310
- item_strs = [json.dumps(item, default=str) for item in items]
 
 
 
 
 
3311
  scores = self._scorer.score_batch(item_strs, query_context)
3312
  # Higher threshold and limit count to avoid adding everything
3313
  high_threshold = max(0.5, self._relevance_threshold * 2)
@@ -3340,6 +3423,7 @@ class SmartCrusher(Transform):
3340
  query_context: str = "",
3341
  preserve_fields: list[str] | None = None,
3342
  max_items: int | None = None,
 
3343
  ) -> CompressionPlan:
3344
  """Plan smart statistical sampling using STATISTICAL detection.
3345
 
@@ -3415,7 +3499,12 @@ class SmartCrusher(Transform):
3415
 
3416
  # 6. Items with high relevance to query context (PROBABILISTIC semantic match)
3417
  if query_context:
3418
- item_strs = [json.dumps(item, default=str) for item in items]
 
 
 
 
 
3419
  scores = self._scorer.score_batch(item_strs, query_context)
3420
  for i, score in enumerate(scores):
3421
  if score.score >= self._relevance_threshold:
 
92
  _QUOTED_STRING_PATTERN = re.compile(r"['\"]([^'\"]{1,50})['\"]") # Short quoted strings
93
  _EMAIL_PATTERN = re.compile(r"\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Z|a-z]{2,}\b")
94
 
95
+ # Temporal detection patterns (compiled once, used in SmartAnalyzer._detect_temporal_field)
96
+ _ISO_DATETIME_PATTERN = re.compile(r"^\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}")
97
+ _ISO_DATE_PATTERN = re.compile(r"^\d{4}-\d{2}-\d{2}$")
98
+
99
 
100
  def extract_query_anchors(text: str) -> set[str]:
101
  """Extract query anchors from user text (legacy regex-based method).
 
632
  _ERROR_KEYWORDS_FOR_PRESERVATION = ERROR_KEYWORDS
633
 
634
 
635
+ def _detect_error_items_for_preservation(
636
+ items: list[dict],
637
+ item_strings: list[str] | None = None,
638
+ ) -> list[int]:
639
  """Detect items containing error keywords for PRESERVATION guarantee.
640
 
641
  This is NOT for crushability analysis - it's for ensuring ALL error items
 
643
  are NEVER dropped, even if errors are common in the dataset.
644
 
645
  Uses keywords because error semantics are well-defined across domains.
646
+
647
+ Args:
648
+ items: List of items to check.
649
+ item_strings: Pre-computed JSON serializations to avoid redundant json.dumps.
650
  """
651
  error_indices: list[int] = []
652
 
 
654
  if not isinstance(item, dict):
655
  continue
656
 
657
+ # Reuse cached serialization if available, otherwise serialize
658
  try:
659
+ if item_strings is not None and i < len(item_strings):
660
+ item_str = item_strings[i].lower()
661
+ else:
662
+ item_str = json.dumps(item).lower()
663
  except Exception:
664
  continue
665
 
 
708
  if not confident_semantics:
709
  return []
710
 
711
+ # Pre-compute field name hashes to avoid redundant SHA256 per item
712
+ _field_hash_cache: dict[str, str] = {}
713
+
714
  for i, item in enumerate(items):
715
  if not isinstance(item, dict):
716
  continue
717
 
718
  for field_name, value in item.items():
719
+ # Hash the field name to match TOIN's format (cached per unique field name)
720
+ if field_name not in _field_hash_cache:
721
+ _field_hash_cache[field_name] = _hash_field_name(field_name)
722
+ field_hash = _field_hash_cache[field_name]
723
 
724
  if field_hash not in confident_semantics:
725
  continue
 
1099
 
1100
  Uses STRUCTURAL detection based on value format, not field names.
1101
  """
1102
+ # Check string fields for ISO 8601 patterns (module-level compiled)
1103
+ iso_datetime_pattern = _ISO_DATETIME_PATTERN
1104
+ iso_date_pattern = _ISO_DATE_PATTERN
1105
 
1106
  for name, stats in field_stats.items():
1107
  if stats.field_type == "string":
 
2483
  # Create compression plan with relevance scoring
2484
  # Pass TOIN preserve_fields so items with those fields get priority
2485
  # Pass effective_max_items for thread-safe compression
2486
+ # Pass item_strings to avoid redundant json.dumps across plan methods
2487
  plan = self._create_plan(
2488
  analysis,
2489
  items,
2490
  query_context,
2491
  preserve_fields=toin_preserve_fields or None,
2492
  effective_max_items=effective_max_items,
2493
+ item_strings=item_strings,
2494
  )
2495
 
2496
  # Execute compression
 
2504
  and len(result) < len(items) # Only cache if compression actually happened
2505
  ):
2506
  store = self._get_compression_store()
2507
+ # Reuse cached item_strings to avoid re-serializing
2508
+ original_json = "[" + ", ".join(item_strings) + "]"
2509
  compressed_json = json.dumps(result, default=str)
2510
 
2511
  ccr_hash = store.store(
 
2543
 
2544
  # TOIN: Record compression event for cross-user learning
2545
  try:
2546
+ # Calculate token counts (approximate) - reuse cached item_strings
2547
+ original_tokens = sum(len(s) for s in item_strings) // 4
2548
  compressed_tokens = len(json.dumps(result, default=str)) // 4
2549
 
2550
  toin.record_compression(
 
2595
  # Universal JSON type handlers (string, number, mixed arrays)
2596
  # =================================================================
2597
 
2598
+ def _compute_k_split(
2599
+ self,
2600
+ items: list,
2601
+ bias: float = 1.0,
2602
+ item_strings: list[str] | None = None,
2603
+ ) -> tuple[int, int, int, int]:
2604
  """Compute adaptive K split into first/last/importance slots.
2605
 
2606
  Uses the existing Kneedle-based adaptive_sizer for K_total, then
2607
  splits according to configurable first_fraction / last_fraction.
2608
 
2609
+ Args:
2610
+ items: List of items (used as fallback for serialization).
2611
+ bias: Compression bias multiplier.
2612
+ item_strings: Pre-computed JSON serializations to avoid redundant json.dumps.
2613
+
2614
  Returns:
2615
  (k_total, k_first, k_last, k_importance)
2616
  """
2617
  from .adaptive_sizer import compute_optimal_k
2618
 
2619
+ if item_strings is None:
2620
+ item_strings = [json.dumps(item, default=str) for item in items]
2621
  k_total = compute_optimal_k(
2622
  item_strings,
2623
  bias=bias,
 
3030
  query_context: str = "",
3031
  preserve_fields: list[str] | None = None,
3032
  effective_max_items: int | None = None,
3033
+ item_strings: list[str] | None = None,
3034
  ) -> CompressionPlan:
3035
  """Create a detailed compression plan using relevance scoring.
3036
 
 
3040
  query_context: Context string from user messages for relevance scoring.
3041
  preserve_fields: TOIN-learned fields that users commonly retrieve.
3042
  Items with values in these fields get higher priority.
3043
+ item_strings: Pre-computed JSON serializations to avoid redundant json.dumps.
3044
  effective_max_items: Thread-safe max items limit (defaults to config value).
3045
  """
3046
  # Use provided effective_max_items or fall back to config
 
3062
 
3063
  if analysis.recommended_strategy == CompressionStrategy.TIME_SERIES:
3064
  plan = self._plan_time_series(
3065
+ analysis,
3066
+ items,
3067
+ plan,
3068
+ query_context,
3069
+ preserve_fields,
3070
+ max_items,
3071
+ item_strings=item_strings,
3072
  )
3073
 
3074
  elif analysis.recommended_strategy == CompressionStrategy.CLUSTER_SAMPLE:
3075
  plan = self._plan_cluster_sample(
3076
+ analysis,
3077
+ items,
3078
+ plan,
3079
+ query_context,
3080
+ preserve_fields,
3081
+ max_items,
3082
+ item_strings=item_strings,
3083
  )
3084
 
3085
  elif analysis.recommended_strategy == CompressionStrategy.TOP_N:
3086
  plan = self._plan_top_n(
3087
+ analysis,
3088
+ items,
3089
+ plan,
3090
+ query_context,
3091
+ preserve_fields,
3092
+ max_items,
3093
+ item_strings=item_strings,
3094
  )
3095
 
3096
  else: # SMART_SAMPLE or NONE
3097
  plan = self._plan_smart_sample(
3098
+ analysis,
3099
+ items,
3100
+ plan,
3101
+ query_context,
3102
+ preserve_fields,
3103
+ max_items,
3104
+ item_strings=item_strings,
3105
  )
3106
 
3107
  return plan
 
3114
  query_context: str = "",
3115
  preserve_fields: list[str] | None = None,
3116
  max_items: int | None = None,
3117
+ item_strings: list[str] | None = None,
3118
  ) -> CompressionPlan:
3119
  """Plan compression for time series data.
3120
 
 
3171
 
3172
  # 5. Items with high relevance to query context (PROBABILISTIC semantic match)
3173
  if query_context:
3174
+ # Reuse pre-computed item_strings if available
3175
+ item_strs = (
3176
+ item_strings
3177
+ if item_strings is not None
3178
+ else [json.dumps(item, default=str) for item in items]
3179
+ )
3180
  scores = self._scorer.score_batch(item_strs, query_context)
3181
  for i, score in enumerate(scores):
3182
  if score.score >= self._relevance_threshold:
 
3203
  query_context: str = "",
3204
  preserve_fields: list[str] | None = None,
3205
  max_items: int | None = None,
3206
+ item_strings: list[str] | None = None,
3207
  ) -> CompressionPlan:
3208
  """Plan compression for clusterable data (like logs).
3209
 
 
3277
 
3278
  # 5. Items with high relevance to query context (PROBABILISTIC semantic match)
3279
  if query_context:
3280
+ # Reuse pre-computed item_strings if available
3281
+ item_strs = (
3282
+ item_strings
3283
+ if item_strings is not None
3284
+ else [json.dumps(item, default=str) for item in items]
3285
+ )
3286
  scores = self._scorer.score_batch(item_strs, query_context)
3287
  for i, score in enumerate(scores):
3288
  if score.score >= self._relevance_threshold:
 
3309
  query_context: str = "",
3310
  preserve_fields: list[str] | None = None,
3311
  max_items: int | None = None,
3312
+ item_strings: list[str] | None = None,
3313
  ) -> CompressionPlan:
3314
  """Plan compression for scored/ranked data.
3315
 
 
3341
 
3342
  if not score_field:
3343
  return self._plan_smart_sample(
3344
+ analysis,
3345
+ items,
3346
+ plan,
3347
+ query_context,
3348
+ preserve_fields,
3349
+ effective_max,
3350
+ item_strings=item_strings,
3351
  )
3352
 
3353
  plan.sort_field = score_field
 
3385
  # Only add items that are NOT already in top N but match the query strongly
3386
  # Use a higher threshold (0.5) since the score field already captures relevance
3387
  if query_context:
3388
+ # Reuse pre-computed item_strings if available
3389
+ item_strs = (
3390
+ item_strings
3391
+ if item_strings is not None
3392
+ else [json.dumps(item, default=str) for item in items]
3393
+ )
3394
  scores = self._scorer.score_batch(item_strs, query_context)
3395
  # Higher threshold and limit count to avoid adding everything
3396
  high_threshold = max(0.5, self._relevance_threshold * 2)
 
3423
  query_context: str = "",
3424
  preserve_fields: list[str] | None = None,
3425
  max_items: int | None = None,
3426
+ item_strings: list[str] | None = None,
3427
  ) -> CompressionPlan:
3428
  """Plan smart statistical sampling using STATISTICAL detection.
3429
 
 
3499
 
3500
  # 6. Items with high relevance to query context (PROBABILISTIC semantic match)
3501
  if query_context:
3502
+ # Reuse pre-computed item_strings if available
3503
+ item_strs = (
3504
+ item_strings
3505
+ if item_strings is not None
3506
+ else [json.dumps(item, default=str) for item in items]
3507
+ )
3508
  scores = self._scorer.score_batch(item_strs, query_context)
3509
  for i, score in enumerate(scores):
3510
  if score.score >= self._relevance_threshold:
pyproject.toml CHANGED
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
4
 
5
  [project]
6
  name = "headroom-ai"
7
- version = "0.5.4"
8
  description = "The Context Optimization Layer for LLM Applications - Cut costs by 50-90%"
9
  readme = "README.md"
10
  license = "Apache-2.0"
@@ -60,6 +60,7 @@ proxy = [
60
  "httpx[http2]>=0.24.0",
61
  "openai>=2.14.0", # OpenAI API format support
62
  "mcp>=1.0.0", # MCP server (headroom_compress, retrieve, stats)
 
63
  ]
64
  # AST-based code compression (tree-sitter)
65
  code = [
 
4
 
5
  [project]
6
  name = "headroom-ai"
7
+ version = "0.5.5"
8
  description = "The Context Optimization Layer for LLM Applications - Cut costs by 50-90%"
9
  readme = "README.md"
10
  license = "Apache-2.0"
 
60
  "httpx[http2]>=0.24.0",
61
  "openai>=2.14.0", # OpenAI API format support
62
  "mcp>=1.0.0", # MCP server (headroom_compress, retrieve, stats)
63
+ "magika>=0.6.0", # ML content detection for ContentRouter
64
  ]
65
  # AST-based code compression (tree-sitter)
66
  code = [