diff --git a/packages/bub-tapestore-otel/src/bub_tapestore_otel/_compat.py b/packages/bub-tapestore-otel/src/bub_tapestore_otel/_compat.py new file mode 100644 index 0000000..e3be81b --- /dev/null +++ b/packages/bub-tapestore-otel/src/bub_tapestore_otel/_compat.py @@ -0,0 +1,16 @@ +"""Tape types resolved across bub versions. + +bub >= 0.3.10 vendors the tape module as ``bub.tape`` and constructs entries +from it; older bub sources tapes from ``republic``. Import tape types from +this module so the whole package agrees on a single resolution order. +""" + +from __future__ import annotations + +try: + from bub.tape import AsyncTapeStore, TapeEntry, TapeQuery, TapeStore +except ImportError: + from republic import TapeEntry, TapeQuery + from republic.tape import AsyncTapeStore, TapeStore + +__all__ = ["AsyncTapeStore", "TapeEntry", "TapeQuery", "TapeStore"] diff --git a/packages/bub-tapestore-otel/src/bub_tapestore_otel/exporter.py b/packages/bub-tapestore-otel/src/bub_tapestore_otel/exporter.py index 79d44ea..b9f4308 100644 --- a/packages/bub-tapestore-otel/src/bub_tapestore_otel/exporter.py +++ b/packages/bub-tapestore-otel/src/bub_tapestore_otel/exporter.py @@ -8,8 +8,9 @@ from typing import Any from loguru import logger -from pydantic import BaseModel, ConfigDict, Field -from republic import TapeEntry +from pydantic import BaseModel, ConfigDict, Field, SkipValidation + +from bub_tapestore_otel._compat import TapeEntry FORCE_FLUSH_TIMEOUT_MS = 3_000 DEFAULT_AGENT_NAME = "bub" @@ -48,7 +49,9 @@ class TraceMessage(TapeProjectionModel): class TraceProjection(TapeProjectionModel): tape: str - entries: list[TapeEntry] + # entries may be instances of either bub.tape's or republic's TapeEntry; + # validating class identity would silently drop every span on a mismatch. + entries: list[SkipValidation[TapeEntry]] input_messages: list[TraceMessage] output_messages: list[TraceMessage] tool_calls: list[ToolCall] diff --git a/packages/bub-tapestore-otel/src/bub_tapestore_otel/store.py b/packages/bub-tapestore-otel/src/bub_tapestore_otel/store.py index 5c8be4f..7b3de19 100644 --- a/packages/bub-tapestore-otel/src/bub_tapestore_otel/store.py +++ b/packages/bub-tapestore-otel/src/bub_tapestore_otel/store.py @@ -1,12 +1,12 @@ from __future__ import annotations +import inspect from collections.abc import Iterable from typing import Protocol from loguru import logger -from republic import AsyncTapeStore, TapeEntry, TapeQuery -from republic.tape import TapeStore -from republic.tape.store import is_async_tape_store + +from bub_tapestore_otel._compat import AsyncTapeStore, TapeEntry, TapeQuery, TapeStore class TapeExporter(Protocol): @@ -23,17 +23,17 @@ def __init__(self, inner: TapeStore | AsyncTapeStore, exporter: TapeExporter) -> self._exporter = exporter async def list_tapes(self) -> list[str]: - if is_async_tape_store(self._inner): + if _is_async_tape_store(self._inner): return await self._inner.list_tapes() return self._inner.list_tapes() async def fetch_all(self, query: TapeQuery[AsyncTapeStore]) -> Iterable[TapeEntry]: - if is_async_tape_store(self._inner): + if _is_async_tape_store(self._inner): return await self._inner.fetch_all(query) return self._inner.fetch_all(query) async def append(self, tape: str, entry: TapeEntry) -> None: - if is_async_tape_store(self._inner): + if _is_async_tape_store(self._inner): await self._inner.append(tape, entry) else: self._inner.append(tape, entry) @@ -43,7 +43,7 @@ async def append(self, tape: str, entry: TapeEntry) -> None: logger.opt(exception=True).warning("tapestore.otel.export_failed action=append tape={}", tape) async def reset(self, tape: str) -> None: - if is_async_tape_store(self._inner): + if _is_async_tape_store(self._inner): await self._inner.reset(tape) else: self._inner.reset(tape) @@ -51,3 +51,7 @@ async def reset(self, tape: str) -> None: self._exporter.reset(tape) except Exception: logger.opt(exception=True).warning("tapestore.otel.export_failed action=reset tape={}", tape) + + +def _is_async_tape_store(store: TapeStore | AsyncTapeStore) -> bool: + return inspect.iscoroutinefunction(store.append) diff --git a/packages/bub-tapestore-otel/tests/test_exporter.py b/packages/bub-tapestore-otel/tests/test_exporter.py index 6f9b474..d2025ee 100644 --- a/packages/bub-tapestore-otel/tests/test_exporter.py +++ b/packages/bub-tapestore-otel/tests/test_exporter.py @@ -3,8 +3,8 @@ from types import SimpleNamespace import bub_tapestore_otel.exporter as exporter +from bub_tapestore_otel._compat import TapeEntry from bub_tapestore_otel.exporter import OTelTapeExporter, _instrument_trace, _should_flush_batch, build_tape_trace -from republic import TapeEntry def test_build_tape_trace_exports_genai_and_openinference_llm_attributes() -> None: @@ -129,6 +129,25 @@ def test_build_tape_trace_falls_back_to_prompt_when_messages_are_missing() -> No assert trace.llm_attributes["llm.input_messages.0.message.content"] == "plain prompt" +def test_build_tape_trace_accepts_entries_from_a_foreign_tape_entry_class() -> None: + # bub >= 0.3.10 constructs entries from bub.tape while older installs use + # republic's TapeEntry; validating class identity here would silently drop + # every span on a mismatch. + entry = SimpleNamespace( + kind="event", + payload={"name": "loop.step", "data": {"status": "ok", "elapsed_ms": 42}}, + id="entry-1", + date="2026-07-07T00:00:00Z", + ) + + trace = build_tape_trace("drift__1", [entry]) + + assert trace.entries == [entry] + assert trace.status == "ok" + assert trace.duration_ms == 42 + assert exporter._step_span_attributes(trace.steps[0])["bub.tape.entry.first_id"] == "entry-1" + + def test_batch_flushes_on_completed_tape_turn_markers() -> None: assert _should_flush_batch(TapeEntry.event("loop.step", data={"status": "ok"})) assert _should_flush_batch(TapeEntry.event("loop.step", data={"status": "error"})) diff --git a/packages/bub-tapestore-otel/tests/test_store.py b/packages/bub-tapestore-otel/tests/test_store.py index ba43563..7794d25 100644 --- a/packages/bub-tapestore-otel/tests/test_store.py +++ b/packages/bub-tapestore-otel/tests/test_store.py @@ -3,8 +3,8 @@ from collections.abc import Iterable import pytest +from bub_tapestore_otel._compat import TapeEntry, TapeQuery from bub_tapestore_otel.store import OTelTapeStore -from republic import TapeEntry, TapeQuery class MemoryStore: diff --git a/packages/bub-tapestore-sqlite/src/bub_tapestore_sqlite/store.py b/packages/bub-tapestore-sqlite/src/bub_tapestore_sqlite/store.py index a36326e..97ca7f2 100644 --- a/packages/bub-tapestore-sqlite/src/bub_tapestore_sqlite/store.py +++ b/packages/bub-tapestore-sqlite/src/bub_tapestore_sqlite/store.py @@ -7,10 +7,9 @@ from typing import Any import aiosqlite -import bub import sqlite_vec from any_llm import AnyLLM -from republic import RepublicError, TapeContext, TapeEntry, TapeQuery +from republic import RepublicError, TapeEntry, TapeQuery from republic.core.errors import ErrorKind ALLOWED_JOURNAL_MODES = {"DELETE", "TRUNCATE", "PERSIST", "MEMORY", "WAL", "OFF"} @@ -50,15 +49,8 @@ def __init__( synchronous: str = "NORMAL", embedding_model: str | None = None, ) -> None: - from bub.builtin.agent import _build_llm - from bub.builtin.settings import AgentSettings - self._path = Path(path) - self._llm = _build_llm( # type: ignore[arg-type] - bub.ensure_config(AgentSettings), - self, - TapeContext(), - ) + self._embedding_client: AnyLLM | None = None self._path.parent.mkdir(parents=True, exist_ok=True) self._busy_timeout_ms = busy_timeout_ms self._journal_mode = normalize_journal_mode(journal_mode) @@ -200,9 +192,10 @@ async def fetch_all(self, query: TapeQuery) -> Iterable[TapeEntry]: async def _compute_embedding(self, texts: list[str]) -> list[list[float]]: if self._embedding_model is None: raise RuntimeError("No embedding model configured for tape store.") - provider, model = self._embedding_model.split(":", 1) - llm: AnyLLM = self._llm._core.get_client(provider) - response = await llm.aembedding(model, texts) + if self._embedding_client is None: + self._embedding_client = _build_embedding_client(self._embedding_model) + _, model = AnyLLM.split_model_provider(self._embedding_model) + response = await self._embedding_client.aembedding(model, texts) return self._embedding_response_to_vectors(response) async def _fetch_by_semantic_query( @@ -796,3 +789,11 @@ def _raise_missing_for_query(query: TapeQuery) -> None: raise RepublicError( ErrorKind.NOT_FOUND, f"Anchor '{query._after_anchor}' was not found." ) + + +def _build_embedding_client(embedding_model: str) -> AnyLLM: + from bub.builtin.settings import load_settings + + provider, _ = AnyLLM.split_model_provider(embedding_model) + settings = load_settings() + return AnyLLM.create(provider, **settings.model_client_kwargs(provider)) diff --git a/packages/bub-tapestore-sqlite/tests/test_store.py b/packages/bub-tapestore-sqlite/tests/test_store.py index 4912ebf..70b3e0b 100644 --- a/packages/bub-tapestore-sqlite/tests/test_store.py +++ b/packages/bub-tapestore-sqlite/tests/test_store.py @@ -6,7 +6,7 @@ from types import SimpleNamespace import pytest -from republic import RepublicError, TapeContext, TapeEntry, TapeQuery +from republic import RepublicError, TapeEntry, TapeQuery from bub_tapestore_sqlite.store import SQLiteTapeStore @@ -29,23 +29,15 @@ def _embedding_response(vectors: list[list[float]]) -> SimpleNamespace: def _mock_embedding_client(monkeypatch, fake_aembedding): - from bub.builtin import agent as agent_module - - build_calls: list[TapeContext] = [] - - def fake_build_llm(settings, tape_store, tape_context): - del settings, tape_store - assert isinstance(tape_context, TapeContext) - build_calls.append(tape_context) - return SimpleNamespace( - _core=SimpleNamespace( - get_client=lambda provider: SimpleNamespace( - aembedding=fake_aembedding, - ) - ) - ) + import bub_tapestore_sqlite.store as store_module + + build_calls: list[str] = [] + + def fake_build_embedding_client(embedding_model: str): + build_calls.append(embedding_model) + return SimpleNamespace(aembedding=fake_aembedding) - monkeypatch.setattr(agent_module, "_build_llm", fake_build_llm) + monkeypatch.setattr(store_module, "_build_embedding_client", fake_build_embedding_client) return build_calls