How WeKnora Implements Event Streaming and Tracing: A Deep Dive into the Go Source Code
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 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 theasyncModeflagEmitAndWait– 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 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. 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:
AppendEventandGetEvents– for client-facing streaming dataAppendSteerEvents,GetSteerEvents, andUpdateSteerEventData– for control-plane instructions that influence generation but remain hidden from clientsDeleteSteerEvent– 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) 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) 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. 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:
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:
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:
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:
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
StreamManagerinterface 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.goprovides synchronous and asynchronous event dispatch with middleware support for logging, timing, and recovery. - Tracing is implemented via the
TracingContextstruct andLangfuseTracingCarrierinterface, enabling W3C traceparent propagation from HTTP handlers through Asynq workers to downstream services. - Control-plane steer events use separate methods in the
StreamManagerinterface 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) handles in-process pub/sub communication between components using Go channels and function callbacks, while StreamManager (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, 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.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →