# How WeKnora Implements Event Streaming and Tracing: A Deep Dive into the Go Source Code

> Explore how WeKnora implements event streaming and tracing with its dual-subsystem architecture. Learn about the StreamManager interface and TracingContext for enhanced observability.

- Repository: [Tencent/WeKnora](https://github.com/tencent/WeKnora)
- Tags: deep-dive
- Published: 2026-09-13

---

**WeKnora implements event streaming and tracing through a dual-subsystem architecture that uses a `StreamManager` interface with Redis and in-memory backends for real-time event propagation, and a `TracingContext` struct with Langfuse integration for cross-process observability.**

Tencent's WeKnora provides real-time observability for AI agent workflows by combining an append-only event streaming pipeline with distributed tracing capabilities. The implementation leverages Go interfaces to abstract storage backends while maintaining strict correlation between HTTP requests, background tasks, and LLM generation steps. This article examines the source code structure in `internal/event`, `internal/stream`, and `internal/types` to reveal how the system handles high-throughput streaming and end-to-end tracing.

## Event Streaming Architecture

### Event Bus and Type Definitions

All domain events in WeKnora are enumerated in [`internal/event/event.go`](https://github.com/Tencent/WeKnora/blob/main/internal/event/event.go) as the `EventType` string constants. These include `query.received`, `chat.stream`, `agent.thought`, and `error`, covering the full lifecycle of AI interactions.

The **EventBus** struct maintains a registry mapping each `EventType` to a slice of `EventHandler` functions. It provides two primary dispatch mechanisms:

- **`Emit`** – supports synchronous or asynchronous delivery based on the `asyncMode` flag
- **`EmitAndWait`** – always fires handlers concurrently and aggregates errors before returning

Middleware support is implemented via `ApplyMiddleware`, allowing cross-cutting concerns like logging, timing, and panic recovery to wrap any handler. The middleware chain is defined in [`internal/event/middleware.go`](https://github.com/Tencent/WeKnora/blob/main/internal/event/middleware.go) and can be composed functionally around the core event handlers.

### Stream Manager Interface

The abstraction for persistent event storage is defined in [`internal/types/interfaces/stream_manager.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/interfaces/stream_manager.go). The **StreamManager** interface treats each event stream as an append-only log indexed by a composite key of `<sessionID, messageID>`.

The interface declares six key methods:

- `AppendEvent` and `GetEvents` – for client-facing streaming data
- `AppendSteerEvents`, `GetSteerEvents`, and `UpdateSteerEventData` – for control-plane instructions that influence generation but remain hidden from clients
- `DeleteSteerEvent` – for cleanup of control signals

The `StreamEvent` struct serves as the universal payload, carrying fields for `ID`, `Content`, `Done` status, `Timestamp`, token usage metadata, and arbitrary `Data` fields.

### Redis and In-Memory Implementations

WeKnora provides two concrete implementations of the `StreamManager` interface.

**RedisStreamManager** (in [`internal/stream/redis_manager.go`](https://github.com/Tencent/WeKnora/blob/main/internal/stream/redis_manager.go)) persists events using Redis **List** structures. It utilizes `LPUSH` to append new chunks and `LRANGE` for paginated retrieval. The implementation maintains a "live-run" key with TTL-based expiration for automatic cleanup of abandoned sessions, alongside a dedicated "steer" list per session for control-plane events.

**MemoryStreamManager** (in [`internal/stream/memory_manager.go`](https://github.com/Tencent/WeKnora/blob/main/internal/stream/memory_manager.go)) stores events in an in-process nested map (`map[string]map[string]*memoryStreamData`). This implementation is optimized for unit tests and single-process deployments where external Redis infrastructure is unavailable.

## Tracing Architecture

### TracingContext and Carrier Interface

Distributed tracing logic resides in [`internal/types/tracing.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/tracing.go). The **TracingContext** struct carries Langfuse trace identifiers including `lf_trace_id`, `lf_traceparent`, and user/session context fields (`LangfuseUserID`, `LangfuseSessionID`).

To enable generic propagation without reflection, WeKnora defines the **LangfuseTracingCarrier** interface with `SetLangfuseTracing` and `GetLangfuseTracing` methods. Any struct embedding `TracingContext` automatically satisfies this interface, allowing the tracing subsystem to inject or extract trace data from arbitrary payloads.

The struct maintains backward compatibility by retaining legacy fields (`lf_trace_id`, `lf_parent_obs_id`), ensuring older serialized payloads remain byte-compatible during rolling deployments.

### Cross-Process Trace Propagation

WeKnora traces flow across process boundaries through explicit carrier injection. When an HTTP request arrives, the W3C `traceparent` header is extracted and stored in a `TracingContext`. This context is then **injected** into Asynq task payloads via the carrier interface when enqueueing background jobs.

Inside the async worker, the payload's `TracingContext` is read via `GetLangfuseTracing`, and the original `traceparent` is re-attached to outgoing HTTP calls or further async jobs. This creates a single contiguous trace tree in Langfuse that links the initial HTTP request, background processing, and downstream LLM API calls.

## Practical Code Examples

### Publishing Events to Redis Streams

The following handler snippet demonstrates appending LLM output chunks to a persistent stream:

```go
func streamAnswerChunk(ctx context.Context, sm interfaces.StreamManager,
    sessionID, messageID, chunk string) error {

    ev := interfaces.StreamEvent{
        ID:        uuid.New().String(),
        Content:   chunk,
        Done:      false,
        Timestamp: time.Now(),
    }
    return sm.AppendEvent(ctx, sessionID, messageID, ev)
}

```

### Consuming Events with Offset Tracking

Clients retrieve new events using an offset-based pagination pattern:

```go
func fetchEvents(ctx context.Context, sm interfaces.StreamManager,
    sessionID, messageID string, from int) ([]interfaces.StreamEvent, int, error) {

    events, next, err := sm.GetEvents(ctx, sessionID, messageID, from)
    if err != nil {
        return nil, from, err
    }
    // `events` holds all new chunks; `next` is the offset for the next poll
    return events, next, nil
}

```

### Injecting Tracing into Async Tasks

Embed `TracingContext` into task payloads to preserve trace continuity across the Asynq boundary:

```go
type MyTaskPayload struct {
    types.TracingContext // embeds tracing data
    Data string
}

func enqueueTask(ctx context.Context, client *asynq.Client, payload MyTaskPayload) error {
    payload.SetLangfuseTracing(types.TracingContext{
        LangfuseTraceparent: ctx.Value("lf_traceparent").(string),
        LangfuseUserID:      ctx.Value("user_id").(string),
        LangfuseSessionID:   ctx.Value("session_id").(string),
    })
    task := asynq.NewTask("my_task", marshal(payload))
    _, err := client.Enqueue(task)
    return err
}

```

### Extracting Context in Background Workers

Workers recover the original trace context to maintain observability across async boundaries:

```go
func handleMyTask(ctx context.Context, t *asynq.Task) error {
    var p MyTaskPayload
    if err := json.Unmarshal(t.Payload(), &p); err != nil {
        return err
    }
    trace := p.GetLangfuseTracing()
    // Re-attach to outgoing calls
    ctx = context.WithValue(ctx, "lf_traceparent", trace.LangfuseTraceparent)
    return process(ctx, p.Data)
}

```

## Summary

- **Event streaming** in WeKnora relies on the `StreamManager` interface with dual Redis and in-memory backends, using Redis Lists for persistence and offset-based pagination for client delivery.
- The **EventBus** in [`internal/event/event.go`](https://github.com/Tencent/WeKnora/blob/main/internal/event/event.go) provides synchronous and asynchronous event dispatch with middleware support for logging, timing, and recovery.
- **Tracing** is implemented via the `TracingContext` struct and `LangfuseTracingCarrier` interface, enabling W3C traceparent propagation from HTTP handlers through Asynq workers to downstream services.
- Control-plane **steer events** use separate methods in the `StreamManager` interface to influence agent behavior without exposing control signals to the client event stream.
- All subsystems maintain backward compatibility and support both production Redis deployments and lightweight in-memory configurations for testing.

## Frequently Asked Questions

### What is the difference between EventBus and StreamManager in WeKnora?

**EventBus** ([`internal/event/event.go`](https://github.com/Tencent/WeKnora/blob/main/internal/event/event.go)) handles in-process pub/sub communication between components using Go channels and function callbacks, while **StreamManager** ([`internal/types/interfaces/stream_manager.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/interfaces/stream_manager.go)) provides persistent, session-scoped storage for streaming data that survives process restarts. The EventBus is ephemeral and synchronous/asynchronous, whereas StreamManager implementations (Redis or memory) maintain append-only logs indexed by session and message IDs.

### How does WeKnora handle trace continuity across asynchronous boundaries?

WeKnora uses the **LangfuseTracingCarrier** interface to serialize `TracingContext` data into Asynq task payloads. When a background worker dequeues the task, it extracts the original Langfuse `traceparent` and user context, then re-injects these values into the new context for outgoing HTTP calls. This carrier pattern ensures a single trace tree spans HTTP requests, async processing, and external API calls.

### Why does WeKnora use Redis Lists instead of Redis Streams for event storage?

According to the source code in [`internal/stream/redis_manager.go`](https://github.com/Tencent/WeKnora/blob/main/internal/stream/redis_manager.go), WeKnora uses standard Redis **List** operations (`LPUSH`, `LRANGE`) rather than Redis Streams (XADD/XREAD) because the offset-based pagination model aligns with the `StreamManager` interface's `GetEvents` method signature. Lists provide simple integer indexing for client polling with `from` offsets, while the "steer" sub-system uses separate list keys for control-plane isolation. This design prioritizes simplicity and compatibility over the consumer-group semantics of Redis Streams.

### How does the steer event system work in WeKnora's streaming architecture?

**Steer events** are control-plane signals (e.g., user interruptions or parameter adjustments) stored via `AppendSteerEvents` and retrieved via `GetSteerEvents`. Unlike standard stream events, steer events are never returned to the client through `GetEvents`; instead, the AI agent logic polls them separately to modify ongoing generation. This separation ensures that control instructions remain internal while streaming output flows to users, enabling real-time steering without exposing implementation details.