How the WeKnora Runtime Task-Queue Dashboard Schedule Works with Per-Stage Worker-Pool Governance
WeKnora uses Redis-backed Asynq servers with a single source of truth in internal/types/task.go to map queues to dedicated worker pools, enabling the runtime dashboard to expose real-time scheduling weights and per-stage governance while allowing elastic capacity sharing via a shared pool.
The Tencent/WeKnora platform implements a sophisticated task-queue architecture that separates workloads into isolated stages while providing unified observability. The runtime task-queue dashboard schedule visualizes queue depth, latency, and worker allocation by aggregating data from six independent Asynq servers. Each server enforces per-stage worker-pool governance through compile-time definitions that guarantee minimum concurrency for critical paths and elastic borrowing for background tasks.
Queue Topology: The Single Source of Truth
All scheduling logic originates in internal/types/task.go, which declares the worker-pool topology and weight mappings that both the Asynq servers and the dashboard API consume.
Worker-Pool Constants
The system defines six distinct pools to isolate workload categories:
// internal/types/task.go
const (
WorkerPoolCore = "core"
WorkerPoolPostProcess = "postprocess"
WorkerPoolEnrichment = "enrichment"
WorkerPoolMaintenance = "maintenance"
WorkerPoolShared = "shared"
WorkerPoolWiki = "wiki"
)
Queue Definitions and Weight Mapping
The queueDefinitions slice (lines 63-88) pairs each logical queue with its draining pool, a scheduling weight, and an optional SharedWeight for elastic borrowing:
var queueDefinitions = []QueueDefinition{
{Name: QueueDefault, Pool: WorkerPoolCore, Weight: 1, SharedWeight: 3, TaskTypes: []string{TypeDocumentProcess, TypeManualProcess}},
{Name: QueueChatAttachment, Pool: WorkerPoolCore, Weight: 3, SharedWeight: 3, TaskTypes: []string{TypeTemporaryDocumentProcess}},
{Name: QueuePostProcess, Pool: WorkerPoolPostProcess, Weight: 1, TaskTypes: []string{TypeKnowledgePostProcess}},
// ... additional queues
}
Helper functions QueueWeightsForPool(pool) (lines 14-24) and QueueWeightsForSharedPool() (lines 26-38) extract these mappings at runtime. Because both the server construction and the dashboard read from this identical structure, the UI always reflects the actual dequeue priorities configured in the Asynq servers.
Per-Pool Asynq Server Architecture
internal/router/task.go instantiates six independent Asynq servers—one per pool—each with dedicated concurrency limits and queue subscriptions.
Server Instantiation
Each constructor resolves environment-based concurrency and binds the specific weight map for its pool:
// internal/router/task.go
func NewCoreAsynqServer(svc interfaces.SystemSettingService) *asynq.Server {
allocation := resolveWorkerPoolConcurrency(svc)
return newAsynqServer(allocation.Core, types.QueueWeightsForPool(types.WorkerPoolCore))
}
This pattern repeats for PostProcess, Enrichment, Maintenance, Shared, and Wiki. Every server uses the same Redis client options (getAsynqRedisClientOpt) to ensure atomic dequeue operations via BRPOPLPUSH, but each enforces its own concurrency boundary defined by environment variables such as WEKNORA_ASYNQ_CORE_CONCURRENCY (defaulting to DefaultCoreWorkerConcurrency = 8).
Elastic Shared Pool Mechanics
The shared pool operates elastically by subscribing to queues from other pools when they are idle. NewSharedAsynqServer invokes QueueWeightsForSharedPool(), which returns SharedWeight values from the definitions:
// Conceptual flow from internal/router/task.go
sharedWeights := types.QueueWeightsForSharedPool() // includes core and enrichment queues with SharedWeight values
return newAsynqServer(allocation.Shared, sharedWeights)
This design allows the shared pool to absorb spare capacity from the core and enrichment stages without starving their guaranteed minimum workers.
Runtime Dashboard Data Pipeline
The dashboard endpoint (GET /api/v1/system/admin/runtime/queues) is implemented in internal/handler/system.go and relies on internal/application/repository/task_queue.go to aggregate state.
Asynq Inspector Integration
The repository uses the Asynq Inspector to collect per-queue metrics:
- Depth metrics:
Size,Pending,Active,Scheduled,Retry,Archived,Completed - Performance indicators:
LatencyMs(age of oldest pending task),MemoryUsageBytes(Redis memory consumption) - Worker heartbeats:
WorkerServerStatcontaining concurrency limits, active worker counts, and subscribed queues
QueueStat Aggregation
The handler merges live Inspector data with static topology definitions to produce the QueueStat struct defined in internal/types/task.go (lines 91-108):
type QueueStat struct {
Name string `json:"name"`
Pool string `json:"pool"`
Weight int `json:"weight"`
Size int `json:"size"`
Pending int `json:"pending"`
Active int `json:"active"`
Scheduled int `json:"scheduled"`
Retry int `json:"retry"`
Archived int `json:"archived"`
Completed int `json:"completed"`
Processed int `json:"processed"`
Failed int `json:"failed"`
Paused bool `json:"paused"`
LatencyMs int64 `json:"latency_ms"`
MemoryUsageBytes int64 `json:"memory_usage_bytes"`
}
The API injects the static Weight and Pool fields from queueDefinitions into each record, ensuring the UI displays accurate scheduling priorities alongside real-time depth and latency.
Governance and Isolation Guarantees
Per-stage worker-pool governance prevents pipeline stages from monopolizing resources while maximizing hardware utilization.
Stage Isolation
Consider a heavy knowledge-post-process job (type: knowledge:post_process). It enters QueuePostProcess, which exclusively belongs to the postprocess pool. Even if the core pool saturates with document parsing operations, the post-process stage maintains its dedicated worker allocation (allocation.PostProcess), ensuring latency-sensitive fan-out tasks proceed without blocking.
Elastic Capacity Borrowing
When the postprocess pool idles, the shared pool can dequeue from QueuePostProcess because QueueWeightsForSharedPool() includes the SharedWeight defined in the queue's definition. This elastic borrowing prevents capacity waste while respecting the minimum guarantees enforced by the dedicated pool servers.
Failure Observability
Dead-letter handling (newDeadLetterKnowledgeFailer) persists failed tasks to the task_dead_letters table and updates the associated knowledge entry to ParseStatusFailed. The dashboard surfaces these records, enabling operators to inspect, retry, or purge failed tasks directly from the UI without accessing Redis manually.
Practical Code Examples
Enqueuing to the Core Pool
Tasks default to the core pool when enqueued without explicit queue routing:
import (
"github.com/hibiken/asynq"
"github.com/Tencent/WeKnora/internal/types"
)
client := weknora.NewAsyncqClient()
payload := types.DocumentProcessPayload{
RequestId: "req-123",
TenantID: 42,
KnowledgeID: "kn-001",
}
task := asynq.NewTask(types.TypeDocumentProcess, marshal(payload))
_, err := client.Enqueue(task) // Routes to QueueDefault (core pool)
Enqueuing to a Specific Stage
Route tasks to the post-process pool explicitly:
task := asynq.NewTask(types.TypeKnowledgePostProcess, marshal(payload))
_, err := client.Enqueue(task, asynq.Queue(types.QueuePostProcess))
Querying the Runtime API
Retrieve current queue statistics via the admin endpoint:
curl -H "Authorization: Bearer $TOKEN" \
http://localhost:8080/api/v1/system/admin/runtime/queues
The response contains QueueStat objects with current depth, latency, and worker allocation per pool.
Summary
- Single source of truth:
internal/types/task.godefines all queue-to-pool mappings, weights, and shared weights, ensuring the dashboard and Asynq servers remain synchronized. - Isolated worker pools: Six independent Asynq servers in
internal/router/task.goenforce hard concurrency limits per stage, preventing resource starvation between critical and background workloads. - Elastic capacity: The
sharedpool utilizesSharedWeightdefinitions to borrow idle capacity from core and enrichment queues without violating minimum guarantees. - Runtime observability: The dashboard aggregates Asynq Inspector metrics with static topology data to expose real-time queue health, latency, and per-pool governance status.
- Failure handling: Dead-letter recording in
task_dead_lettersintegrates with the dashboard for operational task management.
Frequently Asked Questions
How does WeKnora prevent background tasks from starving interactive workloads?
The core pool maintains a guaranteed minimum concurrency reserved exclusively for interactive queues like QueueDefault and QueueChatAttachment. Because internal/router/task.go constructs separate Asynq servers for each pool, background-heavy pools such as enrichment or postprocess cannot consume threads allocated to the core pool. The shared pool may only borrow capacity when the core pool's workers are idle, ensuring interactive tasks retain priority.
What determines the scheduling weight shown in the runtime dashboard?
The Weight field originates from the queueDefinitions slice in internal/types/task.go. Each queue definition assigns a static integer weight that the Asynq server uses to calculate dequeue probability. The dashboard reads these same definitions via QueueWeightsForPool() and injects them into the QueueStat response, guaranteeing that the UI weight matches the actual runtime priority.
Can operators adjust worker concurrency without restarting WeKnora?
Yes. The resolveWorkerPoolConcurrency() function in internal/router/task.go reads environment variables such as WEKNORA_ASYNQ_CORE_CONCURRENCY at server initialization. While runtime adjustment requires a rolling restart of the specific Asynq server, the shared pool's elastic weights automatically adapt to load changes without configuration changes, absorbing spikes via SharedWeight borrowing.
How does the dashboard calculate queue latency?
The repository layer in internal/application/repository/task_queue.go invokes the Asynq Inspector to retrieve the timestamp of the oldest pending task in each queue. It subtracts this value from the current time to derive LatencyMs, which the dashboard renders as an indicator of queue backpressure. High latency on a dedicated pool queue signals that the stage-specific concurrency limit may require increase.
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 →