Files
OpenViking/tests/test_telemetry_runtime.py
runyunzhouandqin-ctx d370ca6a7d feat(config): isolate account runtime resources (#5323)
* feat(config): isolate account runtime resources

* fix(runtime): harden account resource lifecycle

* test: streamline account runtime coverage

* refactor(config): remove redundant account runtime paths

* fix(config): make runtime publication loop-safe

* fix(config): order runtime change notifications

* test: align vector fixtures after rebase

* test: add account runtime isolation e2e coverage

* test: consolidate account runtime coverage into existing contracts

---------

Co-authored-by: qin-ctx <qinhaojie.exe@bytedance.com>
2026-09-24 15:08:56 +08:00

911 lines
30 KiB
Python

# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
# SPDX-License-Identifier: AGPL-3.0
from __future__ import annotations
import asyncio
import json
from concurrent.futures import ThreadPoolExecutor
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
from openviking.models.embedder.base import DenseEmbedderBase, EmbedResult, embed_compat
from openviking.models.vlm.base import VLMBase
from openviking.observability.context import (
bind_operation_observability_context,
bind_root_observability_context,
get_operation_observability_context,
get_root_observability_context,
reset_operation_observability_context,
reset_root_observability_context,
)
from openviking.storage.collection_schemas import TextEmbeddingHandler
from openviking.storage.queuefs.process_result import ProcessOutcome
from openviking.storage.queuefs.semantic_executor import SemanticTreeStats
from openviking.storage.queuefs.semantic_msg import SemanticMsg
from openviking.storage.queuefs.semantic_processor import SemanticProcessor
from openviking.telemetry import (
get_current_telemetry,
register_telemetry,
tracer_module,
unregister_telemetry,
)
from openviking.telemetry.backends.memory import MemoryOperationTelemetry
from openviking.telemetry.context import bind_telemetry, bind_telemetry_stage
from openviking.telemetry.snapshot import TelemetrySnapshot
from openviking.telemetry.span_models import OperationSpanAttributes, RootSpanAttributes
from openviking_cli.utils import logger as logger_module
def test_root_observability_context_bind_and_reset():
root = RootSpanAttributes(http_method="GET", http_route="/demo", request_id="req-1")
token = bind_root_observability_context(root)
try:
assert get_root_observability_context() is root
finally:
reset_root_observability_context(token)
assert get_root_observability_context() is None
def test_operation_observability_context_bind_and_reset():
operation = OperationSpanAttributes(operation="search.find", telemetry_id="tm-demo")
token = bind_operation_observability_context(operation)
try:
assert get_operation_observability_context() is operation
finally:
reset_operation_observability_context(token)
assert get_operation_observability_context() is None
def test_telemetry_snapshot_to_dict_supports_summary_only():
snapshot = TelemetrySnapshot(
telemetry_id="tm_demo",
summary={"duration_ms": 1.2, "tokens": {"total": 3}},
)
payload = snapshot.to_dict(include_summary=True)
assert payload == {
"id": "tm_demo",
"summary": {"duration_ms": 1.2, "tokens": {"total": 3}},
}
def test_telemetry_summary_breaks_down_llm_and_embedding_token_usage():
telemetry = MemoryOperationTelemetry(operation="resources.add_resource", enabled=True)
telemetry.record_token_usage("llm", 11, 7)
telemetry.record_token_usage("embedding", 13, 0)
summary = telemetry.finish().summary
assert telemetry.telemetry_id
assert telemetry.telemetry_id.startswith("tm_")
assert summary["tokens"]["total"] == 31
assert summary["duration_ms"] >= 0
assert summary["tokens"]["llm"] == {
"input": 11,
"output": 7,
"total": 18,
}
assert summary["tokens"]["embedding"] == {"total": 13}
assert "queue" not in summary
assert "vector" not in summary
assert "semantic_nodes" not in summary
assert "memory" not in summary
assert "errors" not in summary
def test_telemetry_summary_breaks_down_stage_token_usage():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
telemetry.record_token_usage("embedding", 11, 0, stage="embed_query")
telemetry.record_token_usage("rerank", 7, 0, stage="rerank")
telemetry.record_token_usage("llm", 5, 3, stage="vlm")
summary = telemetry.finish().summary
assert summary["tokens"]["total"] == 26
assert summary["tokens"]["rerank"] == {"total": 7}
assert summary["tokens"]["stages"]["embed_query"]["embedding"] == {"total": 11}
assert summary["tokens"]["stages"]["rerank"]["rerank"] == {"total": 7}
assert summary["tokens"]["stages"]["vlm"]["llm"] == {
"input": 5,
"output": 3,
"total": 8,
}
@pytest.mark.asyncio
async def test_bind_telemetry_stage_propagates_across_async_tasks():
telemetry = MemoryOperationTelemetry(operation="resource.process", enabled=True)
async def _worker() -> None:
await asyncio.sleep(0)
get_current_telemetry().add_token_usage(6, 4)
with bind_telemetry(telemetry):
with bind_telemetry_stage("resource_summarize"):
await asyncio.create_task(_worker())
summary = telemetry.finish().summary
assert summary["tokens"]["stages"]["resource_summarize"]["llm"] == {
"input": 6,
"output": 4,
"total": 10,
}
@pytest.mark.asyncio
async def test_embed_compat_binds_query_stage_for_embedding_tokens():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
query = " ".join(f"token-{idx}" for idx in range(200))
class _TelemetryAwareAsyncEmbedder(DenseEmbedderBase):
def __init__(self):
super().__init__("telemetry-test", config={"max_input_tokens": 20})
def embed(self, text: str, is_query: bool = False) -> EmbedResult:
raise AssertionError("embed_async should be used")
async def embed_async(self, text: str, is_query: bool = False) -> EmbedResult:
assert is_query is True
assert text.endswith("...(truncated for embedding)")
assert "token-199" not in text
get_current_telemetry().record_token_usage("embedding", 9, 0)
return EmbedResult(dense_vector=[0.1, 0.2])
def get_dimension(self) -> int:
return 2
with bind_telemetry(telemetry):
await embed_compat(_TelemetryAwareAsyncEmbedder(), query, is_query=True)
summary = telemetry.finish().summary
assert summary["tokens"]["stages"]["embed_query"]["embedding"] == {"total": 9}
def test_vlm_base_defaults_operation_tokens_to_vlm_stage():
class _DummyVLM(VLMBase):
def get_completion(self, *args, **kwargs):
raise NotImplementedError()
async def get_completion_async(self, *args, **kwargs):
raise NotImplementedError()
def get_vision_completion(self, *args, **kwargs):
raise NotImplementedError()
async def get_vision_completion_async(self, *args, **kwargs):
raise NotImplementedError()
telemetry = MemoryOperationTelemetry(operation="session.commit", enabled=True)
with bind_telemetry(telemetry):
_DummyVLM({"provider": "openai", "model": "gpt-4o-mini"}).update_token_usage(
model_name="gpt-4o-mini",
provider="openai",
prompt_tokens=7,
completion_tokens=5,
)
summary = telemetry.finish().summary
assert summary["tokens"]["stages"]["vlm"]["llm"] == {
"input": 7,
"output": 5,
"total": 12,
}
def test_disabled_telemetry_still_has_request_id():
telemetry = MemoryOperationTelemetry(operation="resources.add_resource", enabled=False)
assert telemetry.telemetry_id
assert telemetry.telemetry_id.startswith("tm_")
def test_telemetry_summary_uses_simplified_internal_metric_keys():
summary = MemoryOperationTelemetry(
operation="search.find",
enabled=True,
)
summary.count("vector.searches", 2)
summary.count("vector.scored", 5)
summary.count("vector.passed", 3)
summary.set("vector.returned", 2)
summary.count("vector.scanned", 5)
summary.set("vector.scan_reason", "")
summary.set("semantic_nodes.total", 4)
summary.set("semantic_nodes.done", 3)
summary.set("semantic_nodes.pending", 1)
summary.set("semantic_nodes.running", 0)
summary.set("memory.extracted", 6)
result = summary.finish().summary
assert result["vector"] == {
"searches": 2,
"scored": 5,
"passed": 3,
"returned": 2,
"scanned": 5,
"scan_reason": "",
}
assert result["semantic_nodes"] == {
"total": 4,
"done": 3,
"pending": 1,
}
assert result["memory"] == {"extracted": 6}
def test_telemetry_summary_includes_cuvs_route_and_stage_timings():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
telemetry.record_cuvs_search(
{
"algorithm": "brute_force",
"dtype": "float16",
"max_concurrent_gpu_searches": 2,
"auto_mode": True,
"route_reason": "native_filter_threshold",
"filter_kind": "path",
"filter_cache_hit": True,
"native_filter_reused": True,
"build_performed": False,
"eligible_count": 12,
"records_generation": 3,
"index_size": 1000,
"memory_estimated_peak_bytes": 4096,
"memory_free_bytes": 8192,
"memory_usable_bytes": 6144,
"total_ms": 1.25,
"preflight_ms": 0.4,
"native_search_ms": 0.7,
}
)
cuvs = telemetry.finish().summary["vector"]["cuvs"]
assert cuvs == {
"searches": 1,
"algorithms": {"brute_force": 1},
"dtypes": {"float16": 1},
"max_concurrent_gpu_searches": 2,
"auto_mode_searches": 1,
"routes": {"native_filter_threshold": 1},
"filter_kinds": {"path": 1},
"filter_cache_hits": 1,
"native_filter_reuses": 1,
"eligible_count_max": 12,
"records_generation_max": 3,
"index_size_max": 1000,
"memory": {
"estimated_peak_bytes_max": 4096,
"free_bytes_min": 8192,
"usable_bytes_min": 6144,
},
"timings_ms": {
"total": {"sum": 1.25, "max": 1.25},
"preflight": {"sum": 0.4, "max": 0.4},
"native_search": {"sum": 0.7, "max": 0.7},
},
}
def test_telemetry_summary_includes_cuvs_micro_batching_fields():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
for batch_size, batch_wait_ms in ((4, 0.8), (4, 0.7), (1, 1.0)):
telemetry.record_cuvs_search(
{
"algorithm": "brute_force",
"dtype": "float32",
"route_reason": "cuvs",
"filter_kind": "none",
"micro_batching_enabled": True,
"micro_batching_warm_fast_path": batch_size == 4,
"batch_size": batch_size,
"batch_wait_ms": batch_wait_ms,
}
)
cuvs = telemetry.finish().summary["vector"]["cuvs"]
assert cuvs["micro_batching_searches"] == 3
assert cuvs["micro_batched_searches"] == 2
assert cuvs["micro_batching_warm_fast_path_searches"] == 2
assert cuvs["batch_size_max"] == 4
assert cuvs["searches_by_batch_size"] == {"1": 1, "4": 2}
assert cuvs["timings_ms"]["batch_wait"] == {"sum": 2.5, "max": 1.0}
def test_cuvs_telemetry_aggregation_is_completion_order_independent():
samples = [
{
"algorithm": "brute_force",
"dtype": "float32",
"max_concurrent_gpu_searches": 1,
"auto_mode": False,
"route_reason": "cuvs",
"filter_kind": "none",
"filter_cache_eviction_fallback": True,
"filter_words_packed": True,
"build_performed": True,
"records_generation": 2,
"index_size": 100,
"memory_estimated_peak_bytes": 4000,
"memory_free_bytes": 8000,
"memory_usable_bytes": 7000,
"total_ms": 12,
"queue_ms": 2,
"gpu_gate_queue_ms": 1.5,
"build_ms": 5,
"gpu_search_ms": 5,
},
{
"algorithm": "brute_force",
"dtype": "float32",
"max_concurrent_gpu_searches": 2,
"auto_mode": True,
"route_reason": "native_filter_threshold",
"filter_kind": "path",
"filter_cache_hit": True,
"native_filter_reused": True,
"eligible_count": 3,
"records_generation": 2,
"index_size": 100,
"memory_free_bytes": 6000,
"memory_usable_bytes": 5000,
"total_ms": 3,
"preflight_ms": 1,
"native_search_ms": 2,
},
]
def aggregate(order):
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
for index in order:
telemetry.record_cuvs_search(samples[index])
return telemetry.finish().summary["vector"]["cuvs"]
forward = aggregate([0, 1])
reverse = aggregate([1, 0])
assert forward == reverse
assert forward["routes"] == {"cuvs": 1, "native_filter_threshold": 1}
assert forward["builds"] == 1
assert forward["filter_cache_eviction_fallbacks"] == 1
assert forward["packed_filter_queries"] == 1
assert forward["memory"]["free_bytes_min"] == 6000
assert forward["timings_ms"]["total"] == {"sum": 15.0, "max": 12.0}
assert forward["timings_ms"]["gpu_gate_queue"] == {"sum": 1.5, "max": 1.5}
def test_cuvs_telemetry_timing_sum_is_strictly_order_independent():
# A large first value makes ordinary floating-point associativity loss observable.
durations_ms = [10_000_000_000_000.0, 0.001, 0.001]
def aggregate(order):
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
for index in order:
telemetry.record_cuvs_search(
{
"algorithm": "brute_force",
"dtype": "float32",
"route_reason": "cuvs",
"filter_kind": "none",
"total_ms": durations_ms[index],
}
)
return telemetry.finish().summary["vector"]["cuvs"]["timings_ms"]["total"]
forward = aggregate([0, 1, 2])
reverse = aggregate([2, 1, 0])
assert (
forward
== reverse
== {
"sum": 10_000_000_000_000.002,
"max": 10_000_000_000_000.0,
}
)
def test_cuvs_telemetry_aggregation_is_thread_safe():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
samples = [
{
"algorithm": "brute_force",
"dtype": "float32",
"route_reason": "cuvs",
"filter_kind": "none",
"memory_free_bytes": 8000,
"total_ms": 0.001,
"gpu_search_ms": 0.001,
},
{
"algorithm": "cagra",
"dtype": "float16",
"route_reason": "native_filter_threshold",
"filter_kind": "path",
"memory_free_bytes": 6000,
"total_ms": 0.002,
"native_search_ms": 0.002,
},
] * 100
with ThreadPoolExecutor(max_workers=8) as executor:
list(executor.map(telemetry.record_cuvs_search, samples))
cuvs = telemetry.finish().summary["vector"]["cuvs"]
assert cuvs["searches"] == 200
assert cuvs["algorithms"] == {"brute_force": 100, "cagra": 100}
assert cuvs["dtypes"] == {"float16": 100, "float32": 100}
assert cuvs["routes"] == {"cuvs": 100, "native_filter_threshold": 100}
assert cuvs["filter_kinds"] == {"none": 100, "path": 100}
assert cuvs["memory"]["free_bytes_min"] == 6000
assert cuvs["timings_ms"]["total"] == {"sum": 0.3, "max": 0.002}
def test_cuvs_telemetry_bounds_dimensions_and_prunes_unobserved_values():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
telemetry.record_cuvs_search(
{
"algorithm": "tenant-specific-algorithm",
"dtype": "tenant-specific-dtype",
"route_reason": "tenant/1234",
"filter_kind": "tenant-specific-filter",
"auto_mode": "false",
"filter_cache_hit": "false",
"native_filter_reused": "false",
"build_performed": "false",
"eligible_count": None,
"memory_estimated_peak_bytes": float("inf"),
"total_ms": float("nan"),
"gpu_search_ms": -1,
}
)
cuvs = telemetry.finish().summary["vector"]["cuvs"]
assert cuvs == {
"searches": 1,
"algorithms": {"other": 1},
"dtypes": {"other": 1},
"max_concurrent_gpu_searches": 1,
"routes": {"other": 1},
"filter_kinds": {"other": 1},
}
def test_cuvs_telemetry_distinguishes_zero_memory_from_unobserved_memory():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
telemetry.record_cuvs_search(
{
"algorithm": "brute_force",
"dtype": "float32",
"route_reason": "native_memory_budget",
"filter_kind": "none",
"memory_free_bytes": 0,
"memory_usable_bytes": 0,
}
)
cuvs = telemetry.finish().summary["vector"]["cuvs"]
assert cuvs["memory"] == {"free_bytes_min": 0, "usable_bytes_min": 0}
assert "estimated_peak_bytes_max" not in cuvs["memory"]
def test_init_tracer_forwards_headers_to_grpc_exporter(monkeypatch):
captured = {}
class FakeExporter:
def __init__(self, **kwargs):
captured.update(kwargs)
monkeypatch.setattr(tracer_module, "OTLPGrpcSpanExporter", FakeExporter)
monkeypatch.setattr(tracer_module, "BatchSpanProcessor", lambda exporter, **kwargs: exporter)
class FakeTracerProvider:
def __init__(self, resource=None):
self.resource = resource
def add_span_processor(self, _processor):
return None
monkeypatch.setattr(tracer_module, "TracerProvider", FakeTracerProvider)
monkeypatch.setattr(
tracer_module,
"Resource",
SimpleNamespace(create=lambda attrs: attrs),
)
monkeypatch.setattr(
tracer_module,
"otel_trace",
SimpleNamespace(
set_tracer_provider=lambda _provider: None,
get_tracer=lambda service_name: f"tracer:{service_name}",
),
)
monkeypatch.setattr(
tracer_module,
"TraceContextTextMapPropagator",
lambda: "propagator",
)
monkeypatch.setattr(tracer_module, "_setup_logging", lambda: None)
monkeypatch.setattr(tracer_module, "_init_asyncio_instrumentation", lambda: None)
tracer_module.init_tracer(
endpoint="apmplus-cn-beijing.ivolces.com:4317",
service_name="memorydb",
protocol="grpc",
insecure=True,
headers={"x-byteapm-appkey": "trace-appkey"},
enabled=True,
)
assert captured["endpoint"] == "apmplus-cn-beijing.ivolces.com:4317"
assert captured["insecure"] is True
assert captured["headers"] == {"x-byteapm-appkey": "trace-appkey"}
def test_init_tracer_forwards_headers_to_http_exporter(monkeypatch):
captured = {}
class FakeExporter:
def __init__(self, **kwargs):
captured.update(kwargs)
monkeypatch.setattr(tracer_module, "OTLPHttpSpanExporter", FakeExporter)
monkeypatch.setattr(tracer_module, "BatchSpanProcessor", lambda exporter, **kwargs: exporter)
class FakeTracerProvider:
def __init__(self, resource=None):
self.resource = resource
def add_span_processor(self, _processor):
return None
monkeypatch.setattr(tracer_module, "TracerProvider", FakeTracerProvider)
monkeypatch.setattr(
tracer_module,
"Resource",
SimpleNamespace(create=lambda attrs: attrs),
)
monkeypatch.setattr(
tracer_module,
"otel_trace",
SimpleNamespace(
set_tracer_provider=lambda _provider: None,
get_tracer=lambda service_name: f"tracer:{service_name}",
),
)
monkeypatch.setattr(
tracer_module,
"TraceContextTextMapPropagator",
lambda: "propagator",
)
monkeypatch.setattr(tracer_module, "_setup_logging", lambda: None)
monkeypatch.setattr(tracer_module, "_init_asyncio_instrumentation", lambda: None)
tracer_module.init_tracer(
endpoint="https://apmplus-cn-beijing.ivolces.com/api/otlp/v1/traces",
service_name="memorydb",
protocol="http",
headers={"X-ByteAPM-AppKey": "trace-appkey"},
enabled=True,
)
assert captured["endpoint"] == "https://apmplus-cn-beijing.ivolces.com/api/otlp/v1/traces"
assert captured["headers"] == {"X-ByteAPM-AppKey": "trace-appkey"}
def test_init_otel_log_handler_forwards_headers_to_grpc_exporter(monkeypatch):
captured = {}
class FakeExporter:
def __init__(self, **kwargs):
captured.update(kwargs)
class FakeLoggerProvider:
def __init__(self, resource=None):
self.resource = resource
def add_log_record_processor(self, _processor):
return None
monkeypatch.setattr(logger_module, "_otel_log_handler_initialized", False)
monkeypatch.setattr(logger_module, "_otel_log_handler", None)
monkeypatch.setattr(logger_module, "OTLPGrpcLogExporter", FakeExporter)
monkeypatch.setattr(
logger_module,
"BatchLogRecordProcessor",
lambda exporter: exporter,
)
monkeypatch.setattr(logger_module, "LoggerProvider", FakeLoggerProvider)
monkeypatch.setattr(
logger_module,
"LoggingHandler",
lambda **kwargs: SimpleNamespace(**kwargs),
)
monkeypatch.setattr(
logger_module,
"Resource",
SimpleNamespace(create=lambda attrs: attrs),
)
monkeypatch.setattr(logger_module, "set_logger_provider", lambda _provider: None)
monkeypatch.setattr(
logger_module,
"get_logger",
lambda _name: SimpleNamespace(
info=lambda *args, **kwargs: None, warning=lambda *args, **kwargs: None
),
)
handler = logger_module.init_otel_log_handler(
protocol="grpc",
endpoint="apmplus-cn-beijing.ivolces.com:4317",
service_name="memorydb",
insecure=True,
headers={"x-byteapm-appkey": "log-appkey"},
enabled=True,
)
assert handler is not None
assert captured["endpoint"] == "apmplus-cn-beijing.ivolces.com:4317"
assert captured["insecure"] is True
assert captured["headers"] == {"x-byteapm-appkey": "log-appkey"}
def test_init_otel_log_handler_forwards_headers_to_http_exporter(monkeypatch):
captured = {}
class FakeExporter:
def __init__(self, **kwargs):
captured.update(kwargs)
class FakeLoggerProvider:
def __init__(self, resource=None):
self.resource = resource
def add_log_record_processor(self, _processor):
return None
monkeypatch.setattr(logger_module, "_otel_log_handler_initialized", False)
monkeypatch.setattr(logger_module, "_otel_log_handler", None)
monkeypatch.setattr(logger_module, "OTLPHttpLogExporter", FakeExporter)
monkeypatch.setattr(
logger_module,
"BatchLogRecordProcessor",
lambda exporter: exporter,
)
monkeypatch.setattr(logger_module, "LoggerProvider", FakeLoggerProvider)
monkeypatch.setattr(
logger_module,
"LoggingHandler",
lambda **kwargs: SimpleNamespace(**kwargs),
)
monkeypatch.setattr(
logger_module,
"Resource",
SimpleNamespace(create=lambda attrs: attrs),
)
monkeypatch.setattr(logger_module, "set_logger_provider", lambda _provider: None)
monkeypatch.setattr(
logger_module,
"get_logger",
lambda _name: SimpleNamespace(
info=lambda *args, **kwargs: None, warning=lambda *args, **kwargs: None
),
)
handler = logger_module.init_otel_log_handler(
protocol="http",
endpoint="https://apmplus-cn-beijing.ivolces.com/api/otlp/v1/logs",
service_name="memorydb",
headers={"X-ByteAPM-AppKey": "log-appkey"},
enabled=True,
)
assert handler is not None
assert captured["endpoint"] == "https://apmplus-cn-beijing.ivolces.com/api/otlp/v1/logs"
assert captured["headers"] == {"X-ByteAPM-AppKey": "log-appkey"}
def test_telemetry_summary_detects_groups_by_prefix_without_static_key_lists():
telemetry = MemoryOperationTelemetry(operation="search.find", enabled=True)
telemetry.set("vector.debug_probe", 1)
telemetry.set("queue.semantic.processed", 2)
telemetry.set("memory.extracted", 1)
result = telemetry.finish().summary
assert "vector" in result
assert "queue" in result
assert "memory" in result
@pytest.mark.asyncio
async def test_semantic_processor_binds_registered_operation_telemetry(monkeypatch):
telemetry = MemoryOperationTelemetry(operation="resources.add_resource", enabled=True)
register_telemetry(telemetry)
resolver = SimpleNamespace(get_vlm=AsyncMock())
processor = SemanticProcessor(vlm_resolver=resolver)
class FakeVikingFS:
async def exists(self, uri, ctx=None):
return True
async def ls(self, uri, ctx=None):
return []
class _FakeTreeExecutor:
stale = False
def __init__(self, **kwargs):
assert kwargs["processor"]._vlm_resolver is resolver
assert kwargs["ctx"].account_id == "default"
async def run(self, root_uri):
assert get_current_telemetry() is telemetry
get_current_telemetry().record_token_usage("llm", 11, 7)
def get_stats(self):
return SemanticTreeStats()
monkeypatch.setattr(
"openviking.storage.queuefs.semantic_processor.get_viking_fs",
lambda: FakeVikingFS(),
)
monkeypatch.setattr(
"openviking.storage.queuefs.semantic_processor.SemanticTreeExecutor",
lambda **kwargs: _FakeTreeExecutor(**kwargs),
)
try:
outcome = await processor.on_dequeue(
SemanticMsg(
uri="viking://resources/demo",
context_type="resource",
recursive=False,
propagate_to_parent=False,
telemetry_id=telemetry.telemetry_id,
).to_dict()
)
finally:
unregister_telemetry(telemetry.telemetry_id)
assert outcome.outcome is ProcessOutcome.SUCCESS
result = telemetry.finish()
summary = result.summary
assert summary["tokens"]["total"] == 18
assert summary["tokens"]["llm"]["total"] == 18
assert "embedding" not in summary["tokens"]
@pytest.mark.asyncio
async def test_semantic_processor_binds_metric_account_context(monkeypatch):
resolver = SimpleNamespace(get_vlm=AsyncMock())
processor = SemanticProcessor(vlm_resolver=resolver)
ran = {"value": False}
class FakeVikingFS:
async def exists(self, uri, ctx=None):
return True
async def ls(self, uri, ctx=None):
return []
class _FakeTreeExecutor:
stale = False
def __init__(self, **kwargs):
assert kwargs["processor"]._vlm_resolver is resolver
assert kwargs["ctx"].account_id == "acct-semantic"
async def run(self, root_uri):
ran["value"] = True
root_context = get_root_observability_context()
assert root_context is not None
assert root_context.account_id == "acct-semantic"
def get_stats(self):
return SemanticTreeStats()
monkeypatch.setattr(
"openviking.storage.queuefs.semantic_processor.get_viking_fs",
lambda: FakeVikingFS(),
)
monkeypatch.setattr(
"openviking.storage.queuefs.semantic_processor.SemanticTreeExecutor",
lambda **kwargs: _FakeTreeExecutor(**kwargs),
)
outcome = await processor.on_dequeue(
SemanticMsg(
uri="viking://resources/demo",
context_type="resource",
recursive=False,
propagate_to_parent=False,
account_id="acct-semantic",
).to_dict()
)
assert outcome.outcome is ProcessOutcome.SUCCESS
assert ran["value"] is True
@pytest.mark.asyncio
async def test_embedding_handler_binds_registered_operation_telemetry(monkeypatch):
telemetry = MemoryOperationTelemetry(operation="resources.add_resource", enabled=True)
register_telemetry(telemetry)
class _TelemetryAwareEmbedder:
def prepare_embedding_input(self, content):
return content
def embed(self, text: str, is_query: bool = False) -> EmbedResult:
assert text == "hello"
assert is_query is False
get_current_telemetry().record_token_usage("embedding", 9, 0)
return EmbedResult(dense_vector=[0.1, 0.2])
async def embed_async(self, text: str, is_query: bool = False) -> EmbedResult:
return self.embed(text, is_query=is_query)
class _DummyVikingDB:
is_closing = False
uses_content_field = False
async def account_uses_content_field(self, account_id):
assert account_id == "default"
return False
async def upsert(self, _data, *, ctx=None, options=None):
return "rec-1"
provider = SimpleNamespace(bind=Mock(return_value=_TelemetryAwareEmbedder()))
handler = TextEmbeddingHandler(
_DummyVikingDB(), embedding_provider=provider,
)
payload = {
"data": json.dumps(
{
"id": "msg-1",
"message": "hello",
"telemetry_id": telemetry.telemetry_id,
"context_data": {
"id": "id-1",
"uri": "viking://resources/sample",
"account_id": "default",
"abstract": "sample",
},
}
)
}
try:
outcome = await handler.on_dequeue(payload)
finally:
unregister_telemetry(telemetry.telemetry_id)
assert outcome.outcome is ProcessOutcome.SUCCESS
provider.bind.assert_called_once_with("default")
result = telemetry.finish()
summary = result.summary
assert summary["tokens"]["embedding"] == {"total": 9}
def test_telemetry_summary_includes_only_memory_group_when_memory_metrics_exist():
telemetry = MemoryOperationTelemetry(operation="session.commit", enabled=True)
telemetry.record_token_usage("llm", 5, 3)
telemetry.set("memory.extracted", 4)
summary = telemetry.finish().summary
assert summary["memory"] == {"extracted": 4}
assert "queue" not in summary
assert "vector" not in summary
assert "semantic_nodes" not in summary
assert "errors" not in summary