# How the WeKnora Runtime Task-Queue Dashboard Schedule Works with Per-Stage Worker-Pool Governance

> Discover how WeKnora's runtime task-queue dashboard leverages Redis Asynq servers and worker-pool governance for real-time scheduling, per-stage control, and elastic capacity sharing.

- Repository: [Tencent/WeKnora](https://github.com/tencent/WeKnora)
- Tags: internals
- Published: 2026-09-12

---

**WeKnora uses Redis-backed Asynq servers with a single source of truth in [`internal/types/task.go`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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:

```go
// 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:

```go
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`](https://github.com/Tencent/WeKnora/blob/main/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:

```go
// 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:

```go
// 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`](https://github.com/Tencent/WeKnora/blob/main/internal/handler/system.go) and relies on [`internal/application/repository/task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/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**: `WorkerServerStat` containing 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`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go) (lines 91-108):

```go
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:

```go
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:

```go
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:

```bash
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.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go) defines 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.go`](https://github.com/Tencent/WeKnora/blob/main/internal/router/task.go) enforce hard concurrency limits per stage, preventing resource starvation between critical and background workloads.
- **Elastic capacity**: The `shared` pool utilizes `SharedWeight` definitions 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_letters` integrates 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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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`](https://github.com/Tencent/WeKnora/blob/main/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.