Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-observability/src/lexigram/ai/observability/metrics/core.py: 17%

60 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1""" 

2Metrics collection for Lexigram Intelligence operations. 

3 

4This module provides comprehensive metrics collection for: 

5- LLM operations (requests, tokens, duration, costs) 

6- Vector store operations (add, search, delete) 

7- Cache operations (hits, misses) 

8- RAG pipeline operations (queries, retrievals, latency) 

9 

10All metrics use lexigram-monitor's MetricsCollectorProtocol for consistent observability. 

11""" 

12 

13from __future__ import annotations 

14 

15from typing import TYPE_CHECKING, Annotated 

16 

17from lexigram.di.decorators import inject 

18from lexigram.di.markers import Inject 

19 

20if TYPE_CHECKING: 

21 from lexigram.contracts.observability.metrics import ( 

22 MetricsCollectorProtocol as MetricsCollectorProtocol, 

23 ) 

24 

25 

26@inject 

27class AIMetrics: 

28 """Centralized metrics collection for intelligence operations. 

29 

30 Provides counters, gauges, and histograms for tracking: 

31 - LLM API calls, tokens, costs, and latency 

32 - Vector store operations and performance 

33 - Embedding cache hit rates 

34 - RAG pipeline end-to-end performance 

35 

36 Example: 

37 >>> metrics = AIMetrics() 

38 >>> # Track LLM request 

39 >>> metrics.llm_requests_total.increment( 

40 ... labels={"provider": "openai", "model": "gpt-4", "status": "success"} 

41 ... ) 

42 >>> # Track tokens 

43 >>> metrics.llm_tokens_total.increment( 

44 ... amount=1500, 

45 ... labels={"provider": "openai", "model": "gpt-4", "type": "completion"} 

46 ... ) 

47 >>> # Track duration 

48 >>> metrics.llm_duration_seconds.observe( 

49 ... value=0.523, 

50 ... labels={"provider": "openai", "model": "gpt-4"} 

51 ... ) 

52 """ 

53 

54 def __init__( 

55 self, 

56 collector: Annotated[MetricsCollectorProtocol, Inject] | None = None, 

57 ) -> None: 

58 """Initialize intelligence metrics. 

59 

60 Args: 

61 collector: Metrics collector to use (DI-injected). 

62 """ 

63 if collector is None: 

64 raise ValueError( 

65 "No metrics collector provided and lexigram-monitor not available", 

66 ) 

67 

68 self._collector = collector 

69 

70 # LLM Metrics 

71 self.llm_requests_total = self._collector.create_counter( 

72 name="intelligence_llm_requests_total", 

73 description="Total number of LLM API requests", 

74 ) 

75 

76 self.llm_tokens_total = self._collector.create_counter( 

77 name="intelligence_llm_tokens_total", 

78 description="Total number of tokens processed (prompt + completion)", 

79 ) 

80 

81 self.llm_duration_seconds = self._collector.create_histogram( 

82 name="intelligence_llm_duration_seconds", 

83 description="Duration of LLM API calls in seconds", 

84 buckets=[0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0], 

85 ) 

86 

87 self.llm_cost_dollars = self._collector.create_counter( 

88 name="intelligence_llm_cost_dollars", 

89 description="Total cost of LLM API calls in dollars", 

90 ) 

91 

92 self.llm_active_requests = self._collector.create_gauge( 

93 name="intelligence_llm_active_requests", 

94 description="Number of currently active LLM requests", 

95 ) 

96 

97 # Vector Store Metrics 

98 self.vector_operations_total = self._collector.create_counter( 

99 name="intelligence_vector_operations_total", 

100 description="Total number of vector store operations", 

101 ) 

102 

103 self.vector_duration_seconds = self._collector.create_histogram( 

104 name="intelligence_vector_duration_seconds", 

105 description="Duration of vector store operations in seconds", 

106 buckets=[0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0], 

107 ) 

108 

109 self.vector_documents_total = self._collector.create_counter( 

110 name="intelligence_vector_documents_total", 

111 description="Total number of documents in vector operations", 

112 ) 

113 

114 self.vector_collection_size = self._collector.create_gauge( 

115 name="intelligence_vector_collection_size", 

116 description="Number of documents in each vector collection", 

117 ) 

118 

119 # Embedding Cache Metrics 

120 self.embedding_cache_hits = self._collector.create_counter( 

121 name="intelligence_embedding_cache_hits", 

122 description="Number of embedding cache hits", 

123 ) 

124 

125 self.embedding_cache_misses = self._collector.create_counter( 

126 name="intelligence_embedding_cache_misses", 

127 description="Number of embedding cache misses", 

128 ) 

129 

130 self.embedding_cache_size = self._collector.create_gauge( 

131 name="intelligence_embedding_cache_size", 

132 description="Number of embeddings cached", 

133 ) 

134 

135 # RAG Pipeline Metrics 

136 self.rag_queries_total = self._collector.create_counter( 

137 name="intelligence_rag_queries_total", 

138 description="Total number of RAG queries processed", 

139 ) 

140 

141 self.rag_duration_seconds = self._collector.create_histogram( 

142 name="intelligence_rag_duration_seconds", 

143 description="Duration of RAG pipeline stages in seconds", 

144 buckets=[0.1, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0], 

145 ) 

146 

147 self.rag_documents_retrieved = self._collector.create_histogram( 

148 name="intelligence_rag_documents_retrieved", 

149 description="Number of documents retrieved per RAG query", 

150 buckets=[1, 3, 5, 10, 20, 50, 100], 

151 ) 

152 

153 self.rag_active_queries = self._collector.create_gauge( 

154 name="intelligence_rag_active_queries", 

155 description="Number of currently active RAG queries", 

156 ) 

157 

158 # Embedding Metrics 

159 self.embedding_operations_total = self._collector.create_counter( 

160 name="intelligence_embedding_operations_total", 

161 description="Total number of embedding operations", 

162 ) 

163 

164 self.embedding_duration_seconds = self._collector.create_histogram( 

165 name="intelligence_embedding_duration_seconds", 

166 description="Duration of embedding operations in seconds", 

167 buckets=[0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5], 

168 ) 

169 

170 self.embedding_batch_size = self._collector.create_histogram( 

171 name="intelligence_embedding_batch_size", 

172 description="Number of texts embedded in a single batch", 

173 buckets=[1, 5, 10, 25, 50, 100, 250, 500], 

174 ) 

175 

176 # Document ingestion metrics 

177 self.document_ingestion_jobs_submitted = self._collector.create_counter( 

178 "intelligence_document_ingestion_jobs_submitted_total", 

179 "Total number of document ingestion jobs submitted", 

180 ) 

181 self.document_ingestion_jobs_completed = self._collector.create_counter( 

182 "intelligence_document_ingestion_jobs_completed_total", 

183 "Total number of document ingestion jobs completed", 

184 ) 

185 self.document_ingestion_jobs_failed = self._collector.create_counter( 

186 "intelligence_document_ingestion_jobs_failed_total", 

187 "Total number of document ingestion jobs failed", 

188 ) 

189 self.document_ingestion_duration_seconds = self._collector.create_histogram( 

190 "intelligence_document_ingestion_duration_seconds", 

191 "Document ingestion duration in seconds", 

192 ) 

193 self.document_chunks_created_total = self._collector.create_counter( 

194 "intelligence_document_chunks_created_total", 

195 "Total number of document chunks created", 

196 ) 

197 self.document_ingestion_workers_active = self._collector.create_gauge( 

198 "intelligence_document_ingestion_workers_active", 

199 "Number of active document ingestion workers", 

200 ) 

201 

202 # Batch embedding metrics 

203 self.batch_embedding_jobs_submitted = self._collector.create_counter( 

204 "intelligence_batch_embedding_jobs_submitted_total", 

205 "Total number of batch embedding jobs submitted", 

206 ) 

207 self.batch_embedding_jobs_completed = self._collector.create_counter( 

208 "intelligence_batch_embedding_jobs_completed_total", 

209 "Total number of batch embedding jobs completed", 

210 ) 

211 self.batch_embedding_jobs_failed = self._collector.create_counter( 

212 "intelligence_batch_embedding_jobs_failed_total", 

213 "Total number of batch embedding jobs failed", 

214 ) 

215 self.batch_embedding_duration_seconds = self._collector.create_histogram( 

216 "intelligence_batch_embedding_duration_seconds", 

217 "Batch embedding duration in seconds", 

218 ) 

219 self.batch_embedding_texts_processed = self._collector.create_counter( 

220 "intelligence_batch_embedding_texts_processed_total", 

221 "Total number of texts processed for embedding", 

222 ) 

223 self.batch_embedding_workers_active = self._collector.create_gauge( 

224 "intelligence_batch_embedding_workers_active", 

225 "Number of active batch embedding workers", 

226 ) 

227 

228 # Maintenance metrics 

229 self.maintenance_workers_active = self._collector.create_gauge( 

230 "intelligence_maintenance_workers_active", 

231 "Number of active maintenance workers", 

232 ) 

233 self.maintenance_tasks_completed = self._collector.create_counter( 

234 "intelligence_maintenance_tasks_completed_total", 

235 "Total number of maintenance tasks completed", 

236 ) 

237 self.maintenance_tasks_failed = self._collector.create_counter( 

238 "intelligence_maintenance_tasks_failed_total", 

239 "Total number of maintenance tasks failed", 

240 ) 

241 self.maintenance_task_duration_seconds = self._collector.create_histogram( 

242 "intelligence_maintenance_task_duration_seconds", 

243 "Maintenance task duration in seconds", 

244 ) 

245 

246 # Dead Letter Queue metrics 

247 self.dlq_workers_active = self._collector.create_gauge( 

248 "intelligence_dlq_workers_active", 

249 "Number of active DLQ workers", 

250 ) 

251 self.dlq_items_total = self._collector.create_gauge( 

252 "intelligence_dlq_items_total", 

253 "Total number of items in DLQ", 

254 ) 

255 self.dlq_items_added = self._collector.create_counter( 

256 "intelligence_dlq_items_added_total", 

257 "Total number of items added to DLQ", 

258 ) 

259 self.dlq_items_retried = self._collector.create_counter( 

260 "intelligence_dlq_items_retried_total", 

261 "Total number of DLQ items retried", 

262 ) 

263 self.dlq_items_archived = self._collector.create_counter( 

264 "intelligence_dlq_items_archived_total", 

265 "Total number of DLQ items archived", 

266 ) 

267 self.dlq_items_deleted = self._collector.create_counter( 

268 "intelligence_dlq_items_deleted_total", 

269 "Total number of DLQ items deleted", 

270 ) 

271 self.dlq_notifications_sent = self._collector.create_counter( 

272 "intelligence_dlq_notifications_sent_total", 

273 "Total number of DLQ notifications sent", 

274 ) 

275 

276 def get_collector(self) -> MetricsCollectorProtocol: 

277 """Get the underlying metrics collector. 

278 

279 Returns: 

280 The MetricsCollectorProtocol instance for advanced usage. 

281 """ 

282 return self._collector 

283 

284 def record_completion( 

285 self, 

286 provider: str, 

287 model: str, 

288 tokens: int, 

289 cost: float, 

290 ) -> None: 

291 """Record a successful LLM completion. 

292 

293 Args: 

294 provider: Provider name. 

295 model: Model identifier. 

296 tokens: Total tokens consumed. 

297 cost: Estimated dollar cost. 

298 """ 

299 self.llm_requests_total.increment( 

300 labels={"provider": provider, "model": model, "status": "success"}, 

301 ) 

302 self.llm_tokens_total.increment( 

303 amount=tokens, 

304 labels={"provider": provider, "model": model, "type": "completion"}, 

305 ) 

306 self.llm_cost_dollars.increment( 

307 amount=cost, 

308 labels={"provider": provider, "model": model}, 

309 ) 

310 

311 def record_error(self, provider: str, error_type: str) -> None: 

312 """Record an LLM or vector store error. 

313 

314 Args: 

315 provider: Provider name. 

316 error_type: Short error category string. 

317 """ 

318 self.llm_requests_total.increment( 

319 labels={ 

320 "provider": provider, 

321 "model": "unknown", 

322 "status": "error", 

323 "error_type": error_type, 

324 }, 

325 )