Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
200 changes: 28 additions & 172 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -286,7 +286,7 @@ app = Starlette(
)
```

### FastAPI
#### FastAPI

Docs: https://fastapi.tiangolo.com/tutorial/handling-errors/

Expand Down Expand Up @@ -615,11 +615,7 @@ INFO log [16b61d57f9ff4a85ac80f5cd406e0aa2] root test_get
INFO access [16b61d57f9ff4a85ac80f5cd406e0aa2] uvicorn.access 127.0.0.1:24810 - "GET /test HTTP/1.1" 200
```

# Extensions

In addition to the middleware, we've added a couple of extensions for third-party packages.

## Sentry
# Sentry

If your project has [sentry-sdk](https://pypi.org/project/sentry-sdk/)
installed, correlation IDs will automatically be added to Sentry events as a `transaction_id`.
Expand All @@ -629,183 +625,43 @@ this [blogpost](https://blog.sentry.io/2019/04/04/trace-errors-through-stack-usi
for a little bit of detail. The transaction ID is displayed in the event detail view in Sentry and is just an easy way
to connect logs to a Sentry event.

## Celery
# Integration with Celery

> Note: If you're using the celery integration, install the package with `pip install asgi-correlation-id[celery]`
Celery workers run as separate processes, so correlation IDs are lost when spawning background tasks from requests.
However, you can transfer correlation IDs to workers using Celery signals.

For Celery user's there's one primary issue: workers run as completely separate processes, so correlation IDs are lost
when spawning background tasks from requests.

However, with some Celery signal magic, we can actually transfer correlation IDs to worker processes, like this:
Add the following to your Celery configuration (e.g., `celery.py`):

```python
@before_task_publish.connect()
def transfer_correlation_id(headers) -> None:
# This is called before task.delay() finishes
# Here we're able to transfer the correlation ID via the headers kept in our backend
headers[header_key] = correlation_id.get()


@task_prerun.connect()
def load_correlation_id(task) -> None:
# This is called when the worker picks up the task
# Here we're able to load the correlation ID from the headers
id_value = task.request.get(header_key)
correlation_id.set(id_value)
```

To configure correlation ID transfer, simply import and run the setup function the package provides:

```python
from asgi_correlation_id.extensions.celery import load_correlation_ids

load_correlation_ids()
```

### Taking it one step further - Adding Celery tracing IDs

In addition to transferring request IDs to Celery workers, we've added one more log filter for improving tracing in
celery processes. This is completely separate from correlation ID functionality, but is something we use ourselves, so
keep in the package with the rest of the signals.

The log filter adds an ID, `celery_current_id` for each worker process, and an ID, `celery_parent_id` for the process
that spawned it.

Here's a quick summary of outputs from different scenarios:

| Scenario | Correlation ID | Celery Current ID | Celery Parent ID |
|------------------------------------------ |--------------------|-------------------|------------------|
| Request | ✅ | | |
| Request -> Worker | ✅ | ✅ | |
| Request -> Worker -> Another worker | ✅ | ✅ | ✅ |
| Beat -> Worker | ✅* | ✅ | | |
| Beat -> Worker -> Worker | ✅* | ✅ | ✅ | ✅ |

*When we're in a process spawned separately from an HTTP request, a correlation ID is still spawned for the first
process in the chain, and passed down. You can think of the correlation ID as an origin ID, while the combination of
current and parent-ids as a way of linking the chain.

To add the current and parent IDs, just alter your `celery.py` to this:

```diff
+ from asgi_correlation_id.extensions.celery import load_correlation_ids, load_celery_current_and_parent_ids

load_correlation_ids()
+ load_celery_current_and_parent_ids()
```

If you wish to correlate celery task IDs through the IDs found in your broker (i.e., the celery `task_id`), use the `use_internal_celery_task_id` argument on `load_celery_current_and_parent_ids`
```diff
from asgi_correlation_id.extensions.celery import load_correlation_ids, load_celery_current_and_parent_ids

load_correlation_ids()
+ load_celery_current_and_parent_ids(use_internal_celery_task_id=True)
```
Note: `load_celery_current_and_parent_ids` will ignore the `generator` argument when `use_internal_celery_task_id` is set to `True`

To set up the additional log filters, update your log config like this:

```diff
LOGGING = {
'version': 1,
'disable_existing_loggers': False,
'filters': {
'correlation_id': {
+ '()': 'asgi_correlation_id.CorrelationIdFilter',
+ 'uuid_length': 32,
+ 'default_value': '-',
+ },
+ 'celery_tracing': {
+ '()': 'asgi_correlation_id.CeleryTracingIdsFilter',
+ 'uuid_length': 32,
+ 'default_value': '-',
+ },
},
'formatters': {
'web': {
'class': 'logging.Formatter',
'datefmt': '%H:%M:%S',
'format': '%(levelname)s ... [%(correlation_id)s] %(name)s %(message)s',
},
+ 'celery': {
+ 'class': 'logging.Formatter',
+ 'datefmt': '%H:%M:%S',
+ 'format': '%(levelname)s ... [%(correlation_id)s] [%(celery_parent_id)s-%(celery_current_id)s] %(name)s %(message)s',
+ },
},
'handlers': {
'web': {
'class': 'logging.StreamHandler',
'filters': ['correlation_id'],
'formatter': 'web',
},
+ 'celery': {
+ 'class': 'logging.StreamHandler',
+ 'filters': ['correlation_id', 'celery_tracing'],
+ 'formatter': 'celery',
+ },
},
'loggers': {
'my_project': {
+ 'handlers': ['celery' if any('celery' in i for i in sys.argv) else 'web'],
'level': 'DEBUG',
'propagate': True,
},
},
}
```

With these IDs configured you should be able to:

1. correlate all logs from a single origin, and
2. piece together the order each log was run, and which process spawned which

#### Example

With everything configured, assuming you have a set of tasks like this:
from uuid import uuid4

```python
@celery.task()
def debug_task() -> None:
logger.info('Debug task 1')
second_debug_task.delay()
second_debug_task.delay()
from celery.signals import before_task_publish, task_postrun, task_prerun

from asgi_correlation_id import correlation_id

@celery.task()
def second_debug_task() -> None:
logger.info('Debug task 2')
third_debug_task.delay()
fourth_debug_task.delay()
CORRELATION_ID_HEADER = 'CORRELATION_ID'


@celery.task()
def third_debug_task() -> None:
logger.info('Debug task 3')
fourth_debug_task.delay()
fourth_debug_task.delay()
@before_task_publish.connect(weak=False)
def transfer_correlation_id(headers: dict[str, str], **kwargs) -> None:
"""Transfer correlation ID from the request to the Celery task headers."""
cid = correlation_id.get()
if cid:
headers[CORRELATION_ID_HEADER] = cid


@celery.task()
def fourth_debug_task() -> None:
logger.info('Debug task 4')
```
@task_prerun.connect(weak=False)
def load_correlation_id(task, **kwargs) -> None:
"""Load correlation ID from task headers, or generate a new one."""
id_value = task.request.get(CORRELATION_ID_HEADER)
if id_value:
correlation_id.set(id_value)
else:
correlation_id.set(uuid4().hex)

your logs could look something like this:

```
correlation-id current-id
| parent-id |
| | |
INFO [3b162382e1] [ - ] [93ddf3639c] project.tasks - Debug task 1
INFO [3b162382e1] [93ddf3639c] [24046ab022] project.tasks - Debug task 2
INFO [3b162382e1] [93ddf3639c] [cb5595a417] project.tasks - Debug task 2
INFO [3b162382e1] [24046ab022] [08f5428a66] project.tasks - Debug task 3
INFO [3b162382e1] [24046ab022] [32f40041c6] project.tasks - Debug task 4
INFO [3b162382e1] [cb5595a417] [1c75a4ed2c] project.tasks - Debug task 3
INFO [3b162382e1] [08f5428a66] [578ad2d141] project.tasks - Debug task 4
INFO [3b162382e1] [cb5595a417] [21b2ef77ae] project.tasks - Debug task 4
INFO [3b162382e1] [08f5428a66] [8cad7fc4d7] project.tasks - Debug task 4
INFO [3b162382e1] [1c75a4ed2c] [72a43319f0] project.tasks - Debug task 4
INFO [3b162382e1] [1c75a4ed2c] [ec3cf4113e] project.tasks - Debug task 4
@task_postrun.connect(weak=False)
def cleanup_correlation_id(**kwargs) -> None:
"""Clear the correlation ID after the task completes."""
correlation_id.set(None)
```
7 changes: 2 additions & 5 deletions asgi_correlation_id/__init__.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,9 @@
from asgi_correlation_id.context import celery_current_id, celery_parent_id, correlation_id
from asgi_correlation_id.log_filters import CeleryTracingIdsFilter, CorrelationIdFilter
from asgi_correlation_id.context import correlation_id
from asgi_correlation_id.log_filters import CorrelationIdFilter
from asgi_correlation_id.middleware import CorrelationIdMiddleware

__all__ = (
'CeleryTracingIdsFilter',
'CorrelationIdFilter',
'CorrelationIdMiddleware',
'correlation_id',
'celery_current_id',
'celery_parent_id',
)
5 changes: 0 additions & 5 deletions asgi_correlation_id/context.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,3 @@
from contextvars import ContextVar

# Middleware
correlation_id: ContextVar[str | None] = ContextVar('correlation_id', default=None)

# Celery extension
celery_parent_id: ContextVar[str | None] = ContextVar('celery_parent', default=None)
celery_current_id: ContextVar[str | None] = ContextVar('celery_current', default=None)
111 changes: 0 additions & 111 deletions asgi_correlation_id/extensions/celery.py

This file was deleted.

Loading
Loading