# How WeKnora Manages Task Queues and Worker Pools: Architecture and Implementation

> Discover how WeKnora manages task queues and worker pools using Redis and PostgreSQL for efficient and atomic document processing. Explore its architecture and implementation.

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

---

**WeKnora uses the asynq library to run Redis-backed task queues with isolated worker pools, while a durable PostgreSQL fallback in `task_pending_ops` guarantees atomic document-level processing for critical workflows.**

WeKnora is a knowledge management platform developed by Tencent that handles complex background processing for document ingestion, wiki updates, and knowledge base operations. Understanding how WeKnora manages its task queue and worker pools reveals a sophisticated dual-layer architecture that combines high-performance Redis queues with durable SQL persistence to ensure reliability across the entire document processing pipeline.

## Queue Topology and Worker Pool Definitions

The single source of truth for WeKnora's task queue layout lives in **[`internal/types/task.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go)**. This file defines the mapping between business-semantic queue names, worker pool assignments, and scheduling weights.

### Business-Semantic Queue Names

WeKnora separates logical queues from physical worker pools using a `queueDefinitions` slice that maps task types to specific processing lanes:

```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,
    }},
    {Name: QueueWiki,       Pool: WorkerPoolWiki,        Weight: 1, TaskTypes: []string{
        TypeWikiIngest, TypeWikiFinalize,
    }},
}

```

- **Worker-pool names** such as `WorkerPoolCore`, `WorkerPoolWiki`, and `WorkerPoolPostProcess` describe isolated processing lanes that prevent resource contention between different workload types.
- The `QueueDefinitions()` function returns a defensive copy to prevent mutation of global configuration, while helper functions like `QueueForTaskType` and `QueueWeightsForPool` expose the mapping to the rest of the codebase.

### Weight-Based Task Scheduling

Each queue definition includes **Weight** and **SharedWeight** fields that control how aggressively worker pools fetch from specific queues. Higher weights indicate preferential allocation of worker capacity to business-critical tasks. The `QueueWeightsForPool` function generates the queue-weight map required by asynq's server configuration, allowing fine-grained control over throughput distribution without redeploying application code.

## Durable Pending-Operations Queue

While asynq provides fast in-memory queueing via Redis, WeKnora maintains a **durable fallback** for critical workflows through the `task_pending_ops` relational table. This design supports the wiki pipeline and a "Lite" mode that operates without Redis availability.

The interface in **[`internal/types/interfaces/task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/interfaces/task_queue.go)** defines the core primitives:

```go
type TaskPendingQueue interface {
    Enqueue(ctx context.Context, op *TaskPendingOp) error
    PeekBatch(ctx context.Context, limit int) ([]*TaskPendingOp, error)
    ClaimBatch(ctx context.Context, taskType TaskType, scope TaskScope, scopeID string, limit int, staleBefore time.Time) ([]*TaskPendingOp, error)
    ReleaseByIDs(ctx context.Context, ids []int64) error
    DeleteByIDs(ctx context.Context, ids []int64) error
    IncrFailCount(ctx context.Context, id int64) error
    PendingCount(ctx context.Context, taskType TaskType) (int64, error)
    DeleteByDedupKey(ctx context.Context, dedupKey string) error
}

```

### Atomic Claiming with Dedup Keys

The implementation in **[`internal/application/repository/task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/internal/application/repository/task_queue.go)** guarantees document-level atomicity through the `ClaimBatch` method. This method groups operations by `dedup_key` to ensure all pending operations for a single document are claimed together, preventing partial processing and race conditions.

### PostgreSQL Row Locking Strategy

The `ClaimBatch` implementation uses a two-step transaction with database-specific optimizations:

```go
// 1. Choose distinct dedup_keys that are either unclaimed or whose claim is stale.
if tx.Dialector.Name() == "postgres" {
    const anchorSQL = `SELECT dedup_key FROM task_pending_ops … FOR UPDATE SKIP LOCKED`
    // fill keys
}

// 2. Resolve exact rows for chosen keys and stamp claimed_at.
var ids []int64
tx.Model(&types.TaskPendingOp{}).
    Where("dedup_key IN ?", keys).
    Where("(claimed_at IS NULL OR claimed_at < ?)", staleBefore).
    Pluck("id", &ids)

tx.Model(&types.TaskPendingOp{}).Where("id IN ?", ids).Update("claimed_at", now)

```

The **FOR UPDATE SKIP LOCKED** clause enables concurrent workers to claim disjoint key sets without blocking, maximizing throughput while maintaining consistency (see lines 88-108 in [`task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/task_queue.go)).

## Worker Pool Concurrency Configuration

Each worker pool's concurrency is driven by environment variables following the pattern `WEKNORA_ASYNQ_*_CONCURRENCY`. Default values are declared in [`internal/types/task.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go):

```go
DefaultCoreWorkerConcurrency        = 8
DefaultPostProcessWorkerConcurrency = 2
DefaultEnrichmentWorkerConcurrency  = 12
DefaultMaintenanceWorkerConcurrency = 4
DefaultSharedWorkerConcurrency      = 6
DefaultWikiWorkerConcurrency        = 8

```

The `ResolveWorkerPoolConcurrency` function (lines 66-85) reads these environment variables, validates positivity, and falls back to defaults. This allows operators to tune throughput for specific workload types without modifying source code.

## Enqueueing Tasks via the TaskEnqueuer Interface

Application code uses the **`TaskEnqueuer`** abstraction defined in [`internal/types/interfaces/task_enqueuer.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/interfaces/task_enqueuer.go) to decouple business logic from the underlying asynq client:

```go
type TaskEnqueuer interface {
    Enqueue(task *asynq.Task, opts ...asynq.Option) (*asynq.TaskInfo, error)
}

```

Typical usage in the wiki ingest pipeline demonstrates dynamic queue selection:

```go
payload := &WikiIngestPayload{
    TracingContext: ctx.Trace(),
    TenantID:       tenantID,
    KnowledgeBaseID: kbID,
}
task := asynq.NewTask(types.TypeWikiIngest, payload)
queueName, _ := types.QueueForTaskType(types.TypeWikiIngest) // → "wiki"
_, err := enqueuer.Enqueue(task, asynq.Queue(queueName), asynq.MaxRetry(5))

```

The `QueueForTaskType` helper ensures tasks route to the correct logical queue based on the type definitions in [`task.go`](https://github.com/Tencent/WeKnora/blob/main/task.go).

## Running Workers and Dead-Letter Handling

Each worker pool instantiates an independent `asynq.Server` during application bootstrap. The server configuration combines resolved concurrency settings with queue-weight maps:

```go
poolCfg := types.ResolveWorkerPoolConcurrency(envReader)
cfg := asynq.Config{
    Concurrency: poolCfg.Core,
    Queues:      types.QueueWeightsForPool(types.WorkerPoolCore),
}
srv := asynq.NewServer(redisConn, cfg)
srv.Run(mux)

```

Workers use the **Langfuse tracing middleware** ([`internal/tracing/langfuse/asynq.go`](https://github.com/Tencent/WeKnora/blob/main/internal/tracing/langfuse/asynq.go)) to propagate trace IDs from HTTP requests into background jobs, maintaining observability across asynchronous boundaries.

When tasks exhaust their retry budgets, the dead-letter middleware persists failures to the `task_dead_letters` table. The repository in [`task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/task_queue.go) (lines 228-254) provides administrative functions for listing, paginating, and manually deleting dead letters through the admin UI.

## Summary

- **Queue definitions** in [`internal/types/task.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go) map business-semantic task types to isolated worker pools with configurable weights.
- **Dual-layer persistence** combines Redis-backed asynq queues for performance with a PostgreSQL `task_pending_ops` table for durability and Lite-mode compatibility.
- **Atomic claiming** via `ClaimBatch` uses `dedup_key` grouping and `FOR UPDATE SKIP LOCKED` to prevent race conditions during batch processing.
- **Environment-driven configuration** via `WEKNORA_ASYNQ_*_CONCURRENCY` variables allows runtime tuning of worker pool sizes without redeployment.
- **Interface abstractions** like `TaskEnqueuer` decouple business logic from queue implementation details, facilitating testing and future migrations.

## Frequently Asked Questions

### What task queue library does WeKnora use?

WeKnora builds its task queue and worker pools on top of **asynq**, a Go library that uses Redis lists for task storage and distribution. The system wraps asynq with custom abstractions in [`internal/types/interfaces/task_enqueuer.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/interfaces/task_enqueuer.go) and augments it with a durable PostgreSQL fallback for critical document processing workflows.

### How does WeKnora ensure atomic processing of related tasks?

The `ClaimBatch` method in [`internal/application/repository/task_queue.go`](https://github.com/Tencent/WeKnora/blob/main/internal/application/repository/task_queue.go) groups pending operations by `dedup_key`, ensuring that all operations for a single document are claimed atomically. On PostgreSQL, the implementation uses `SELECT ... FOR UPDATE SKIP LOCKED` to allow concurrent workers while guaranteeing that no two workers claim operations for the same document simultaneously.

### What happens when a task fails permanently in WeKnora?

When a task exhausts its configured retry budget (set via `asynq.MaxRetry`), the dead-letter middleware captures the failure and inserts a record into the `task_dead_letters` table. Operators can inspect these failed tasks through the admin UI and choose to delete them or manually retry after fixing underlying issues.

### How are worker pool sizes configured in WeKnora?

Worker pool concurrency is controlled through environment variables following the pattern `WEKNORA_ASYNQ_<POOL>_CONCURRENCY`, where `<POOL>` corresponds to the worker pool name (Core, Wiki, PostProcess, etc.). The `ResolveWorkerPoolConcurrency` function in [`internal/types/task.go`](https://github.com/Tencent/WeKnora/blob/main/internal/types/task.go) reads these values, validates them, and falls back to sane defaults ranging from 2 to 12 workers per pool depending on the workload type.