Skip to content

Commit a202d32

Browse files
committed
feat: Add twisted instrumentation and unit/integration tests
Signed-off-by: Cagri Yonca <cagri@ibm.com>
1 parent 45ceebb commit a202d32

10 files changed

Lines changed: 1376 additions & 0 deletions

File tree

src/instana/__init__.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -213,6 +213,12 @@ def boot_agent() -> None:
213213
from instana.instrumentation.tornado import (
214214
server as tornado_server, # noqa: F401
215215
)
216+
from instana.instrumentation.twisted import (
217+
client as twisted_client, # noqa: F401
218+
)
219+
from instana.instrumentation.twisted import (
220+
server as twisted_server, # noqa: F401
221+
)
216222

217223

218224
def _start_profiler() -> None:

src/instana/instrumentation/twisted/__init__.py

Whitespace-only changes.
Lines changed: 130 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,130 @@
1+
# (c) Copyright IBM Corp. 2026
2+
3+
try:
4+
from typing import TYPE_CHECKING, Any, Callable, Union
5+
6+
import wrapt
7+
8+
if TYPE_CHECKING:
9+
from twisted.internet.defer import Deferred
10+
from twisted.web.client import Agent
11+
from twisted.web.iweb import IResponse
12+
13+
from instana.span.span import InstanaSpan
14+
15+
from opentelemetry.context import get_current
16+
from opentelemetry.semconv.trace import SpanAttributes
17+
from twisted.python.failure import Failure
18+
from twisted.web.http_headers import Headers as TwistedHeaders
19+
20+
from instana.log import logger
21+
from instana.propagators.format import Format
22+
from instana.singletons import agent, get_tracer
23+
from instana.span.span import get_current_span
24+
from instana.util.secrets import strip_secrets_from_query
25+
from instana.util.traceutils import extract_custom_headers
26+
27+
@wrapt.patch_function_wrapper("twisted.web.client", "Agent.request")
28+
def request_with_instana(
29+
wrapped: Callable[..., Any],
30+
instance: "Agent",
31+
argv: tuple[Any, ...],
32+
kwargs: dict[str, Any],
33+
) -> "Deferred":
34+
try:
35+
parent_span = get_current_span()
36+
37+
# If we're not tracing, just return
38+
if not parent_span.is_recording():
39+
return wrapped(*argv, **kwargs)
40+
41+
# argv: (method, url[, headers[, bodyProducer]])
42+
method = argv[0]
43+
url = argv[1]
44+
headers = argv[2] if len(argv) > 2 else kwargs.get("headers")
45+
46+
method_str = (
47+
method.decode("latin-1") if isinstance(method, bytes) else str(method)
48+
)
49+
url_str = url.decode("latin-1") if isinstance(url, bytes) else str(url)
50+
51+
parent_context = get_current()
52+
tracer = get_tracer()
53+
span = tracer.start_span("twisted-client", context=parent_context)
54+
55+
# Query param scrubbing
56+
parts = url_str.split("?", 1)
57+
span.set_attribute(SpanAttributes.HTTP_URL, parts[0])
58+
if len(parts) > 1 and parts[1]:
59+
cleaned_qp = strip_secrets_from_query(
60+
parts[1],
61+
agent.options.secrets_matcher,
62+
agent.options.secrets_list,
63+
)
64+
span.set_attribute("http.params", cleaned_qp)
65+
66+
span.set_attribute(SpanAttributes.HTTP_METHOD, method_str)
67+
68+
# Build / augment headers with trace correlation
69+
if headers is None or not isinstance(headers, TwistedHeaders):
70+
headers = TwistedHeaders({})
71+
72+
# Capture outgoing request headers
73+
headers_dict = {
74+
k.decode("latin-1"): v[0].decode("utf-8")
75+
for k, v in headers.getAllRawHeaders()
76+
}
77+
extract_custom_headers(span, headers_dict)
78+
79+
# Inject Instana correlation headers
80+
inject_carrier: dict[str, str] = {}
81+
tracer.inject(span.context, Format.HTTP_HEADERS, inject_carrier)
82+
for key, value in inject_carrier.items():
83+
headers.setRawHeaders(key.encode("latin-1"), [value.encode("utf-8")])
84+
85+
# Rebuild argv with the modified headers
86+
new_argv = (argv[0], argv[1], headers) + argv[3:]
87+
88+
deferred = wrapped(*new_argv, **kwargs)
89+
90+
if deferred is not None:
91+
deferred.addBoth(_finish_tracing, span)
92+
93+
return deferred
94+
except Exception:
95+
logger.debug("twisted request_with_instana", exc_info=True)
96+
return wrapped(*argv, **kwargs)
97+
98+
def _finish_tracing(
99+
result: "Union[IResponse, Failure]", span: "InstanaSpan"
100+
) -> "Union[IResponse, Failure]":
101+
"""Callback/errback attached to the Agent.request Deferred."""
102+
try:
103+
if isinstance(result, Failure):
104+
span.record_exception(result.value)
105+
else:
106+
response = result
107+
status_code = response.code
108+
span.set_attribute(SpanAttributes.HTTP_STATUS_CODE, status_code)
109+
110+
if status_code >= 400:
111+
span.mark_as_errored()
112+
113+
# Capture response headers
114+
headers_dict = {
115+
k.decode("latin-1"): v[0].decode("utf-8")
116+
for k, v in response.headers.getAllRawHeaders()
117+
}
118+
extract_custom_headers(span, headers_dict)
119+
120+
except Exception:
121+
logger.debug("twisted _finish_tracing", exc_info=True)
122+
finally:
123+
if span.is_recording():
124+
span.end()
125+
126+
return result
127+
128+
logger.debug("Instrumenting twisted client")
129+
except ImportError:
130+
pass
Lines changed: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
1+
# (c) Copyright IBM Corp. 2026
2+
3+
try:
4+
from typing import TYPE_CHECKING, Any, Callable, Optional
5+
6+
import wrapt
7+
from opentelemetry import context, trace
8+
9+
if TYPE_CHECKING:
10+
from twisted.web.http import Request
11+
from twisted.web.resource import Resource
12+
13+
from opentelemetry.semconv.trace import SpanAttributes
14+
15+
from instana.log import logger
16+
from instana.propagators.format import Format
17+
from instana.singletons import agent, get_tracer
18+
from instana.util.secrets import strip_secrets_from_query
19+
from instana.util.traceutils import extract_custom_headers
20+
21+
@wrapt.patch_function_wrapper("twisted.web.resource", "Resource.render")
22+
def render_with_instana(
23+
wrapped: Callable[..., Any],
24+
instance: "Resource",
25+
argv: tuple[Any, ...],
26+
kwargs: dict[str, Any],
27+
) -> Optional[bytes]:
28+
try:
29+
request: "Request" = argv[0]
30+
tracer = get_tracer()
31+
32+
# Extract parent context from incoming request headers
33+
parent_context = None
34+
if request.requestHeaders:
35+
headers_dict: dict[str, str] = {
36+
k.decode("latin-1"): v[0].decode("utf-8")
37+
for k, v in request.requestHeaders.getAllRawHeaders()
38+
}
39+
parent_context = tracer.extract(Format.HTTP_HEADERS, headers_dict)
40+
41+
span = tracer.start_span("twisted-server", context=parent_context)
42+
43+
# Set span as current so downstream code (e.g. twisted-client) can find it
44+
ctx = trace.set_span_in_context(span)
45+
token = context.attach(ctx)
46+
request._instana_token = token
47+
48+
# Extract the URL components
49+
host = request.getHeader("host") or ""
50+
scheme = "https" if request.isSecure() else "http"
51+
raw_path = request.path
52+
path = (
53+
raw_path.decode("latin-1") if isinstance(raw_path, bytes) else raw_path
54+
)
55+
56+
url = f"{scheme}://{host}{path}"
57+
span.set_attribute(SpanAttributes.HTTP_URL, url)
58+
59+
raw_method = request.method
60+
method = (
61+
raw_method.decode("latin-1")
62+
if isinstance(raw_method, bytes)
63+
else raw_method
64+
)
65+
span.set_attribute(SpanAttributes.HTTP_METHOD, method)
66+
67+
# Query param scrubbing
68+
raw_query = request.uri
69+
query = (
70+
raw_query.decode("latin-1")
71+
if isinstance(raw_query, bytes)
72+
else raw_query
73+
)
74+
if "?" in query:
75+
qs = query.split("?", 1)[1]
76+
if qs:
77+
cleaned_qp = strip_secrets_from_query(
78+
qs,
79+
agent.options.secrets_matcher,
80+
agent.options.secrets_list,
81+
)
82+
span.set_attribute("http.params", cleaned_qp)
83+
84+
# Request header tracking support
85+
extract_custom_headers(span, headers_dict)
86+
87+
# Inject correlation headers into response
88+
response_headers = {}
89+
tracer.inject(span.context, Format.HTTP_HEADERS, response_headers)
90+
for key, value in response_headers.items():
91+
request.setHeader(key.encode("latin-1"), value.encode("utf-8"))
92+
93+
# Store span on request for later retrieval
94+
request._instana = span
95+
96+
# Wrap request.finish to close the span when the response is complete
97+
original_finish = request.finish
98+
99+
def finish_with_instana() -> None:
100+
try:
101+
_span = request._instana
102+
status_code = request.code
103+
if isinstance(status_code, int):
104+
_span.set_attribute(
105+
SpanAttributes.HTTP_STATUS_CODE, status_code
106+
)
107+
if status_code >= 500:
108+
_span.mark_as_errored()
109+
110+
# Capture response headers
111+
response_hdrs = {
112+
k.decode("latin-1"): v[0].decode("utf-8")
113+
for k, v in request.responseHeaders.getAllRawHeaders()
114+
}
115+
extract_custom_headers(_span, response_hdrs)
116+
117+
if _span.is_recording():
118+
_span.end()
119+
except Exception:
120+
logger.debug("twisted finish_with_instana", exc_info=True)
121+
finally:
122+
context.detach(request._instana_token)
123+
original_finish()
124+
125+
request.finish = finish_with_instana
126+
127+
return wrapped(*argv, **kwargs)
128+
except Exception:
129+
logger.debug("twisted render_with_instana", exc_info=True)
130+
return wrapped(*argv, **kwargs)
131+
132+
logger.debug("Instrumenting twisted server")
133+
except ImportError:
134+
pass

src/instana/span/kind.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@
1616
"httpx",
1717
"tornado-client",
1818
"tornado-server",
19+
"twisted-client",
20+
"twisted-server",
1921
"urllib3",
2022
"wsgi",
2123
"asgi",
@@ -31,6 +33,7 @@
3133
"rabbitmq",
3234
"rpc-server",
3335
"tornado-server",
36+
"twisted-server",
3437
"gcps-consumer",
3538
"asgi",
3639
"kafka-consumer",
@@ -57,6 +60,7 @@
5760
"sqlalchemy",
5861
"s3",
5962
"tornado-client",
63+
"twisted-client",
6064
"urllib3",
6165
"pymongo",
6266
"gcs",
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
# (c) Copyright IBM Corp. 2026
2+
3+
import os
4+
import socket
5+
6+
from ...helpers import testenv
7+
from ..utils import launch_background_thread
8+
9+
app_thread = None
10+
11+
12+
def _get_free_port() -> int:
13+
"""Ask the OS for a free port."""
14+
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
15+
s.bind(("127.0.0.1", 0))
16+
return s.getsockname()[1]
17+
18+
19+
if not any((
20+
app_thread,
21+
os.environ.get("GEVENT_TEST"),
22+
os.environ.get("CASSANDRA_TEST"),
23+
)):
24+
testenv["twisted_port"] = _get_free_port()
25+
testenv["twisted_server"] = "http://127.0.0.1:" + str(testenv["twisted_port"])
26+
27+
# Background Twisted application
28+
from .app import run_server
29+
30+
app_thread = launch_background_thread(run_server, "Twisted")

0 commit comments

Comments
 (0)