forked from forthfate/openorbit
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathobservability.py
More file actions
68 lines (58 loc) · 3.06 KB
/
Copy pathobservability.py
File metadata and controls
68 lines (58 loc) · 3.06 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
"""OpenTelemetry setup for every workflow action and process lifecycle event."""
from __future__ import annotations
import json
from datetime import UTC, datetime
from pathlib import Path
from threading import Lock
from opentelemetry import trace
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter, SpanExportResult
class JsonlSpanExporter(SpanExporter):
"""Local OTEL exporter; each record is an exported span, never an ad-hoc log."""
def __init__(self, path: Path) -> None:
self.path = path
self.lock = Lock()
self.path.parent.mkdir(parents=True, exist_ok=True)
def export(self, spans) -> SpanExportResult: # type: ignore[no-untyped-def]
with self.lock, self.path.open("a", encoding="utf-8") as stream:
for span in spans:
record = {
"resourceSpans": [
{
"resource": {key: str(value) for key, value in span.resource.attributes.items()},
"scopeSpans": [
{
"spans": [
{
"name": span.name,
"traceId": f"{span.context.trace_id:032x}",
"spanId": f"{span.context.span_id:016x}",
"parentSpanId": f"{span.parent.span_id:016x}"
if span.parent
else None,
"startTime": span.start_time,
"endTime": span.end_time,
"attributes": dict(span.attributes),
"status": span.status.status_code.name,
"events": [
{"name": event.name, "attributes": dict(event.attributes)}
for event in span.events
],
}
]
}
],
}
],
"exportedAt": datetime.now(UTC).isoformat(),
}
stream.write(json.dumps(record, default=str, ensure_ascii=False) + "\n")
return SpanExportResult.SUCCESS
def shutdown(self) -> None:
return None
def configure_telemetry(path: Path): # type: ignore[no-untyped-def]
provider = TracerProvider(resource=Resource.create({"service.name": "orbit-agent-console"}))
provider.add_span_processor(BatchSpanProcessor(JsonlSpanExporter(path), schedule_delay_millis=200))
trace.set_tracer_provider(provider)
return trace.get_tracer("orbit.console")