How WeKnora Manages Task Queues and Worker Pools: Architecture and Implementation
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. 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:
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, andWorkerPoolPostProcessdescribe 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 likeQueueForTaskTypeandQueueWeightsForPoolexpose 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 defines the core primitives:
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 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:
// 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).
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:
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 to decouple business logic from the underlying asynq client:
type TaskEnqueuer interface {
Enqueue(task *asynq.Task, opts ...asynq.Option) (*asynq.TaskInfo, error)
}
Typical usage in the wiki ingest pipeline demonstrates dynamic queue selection:
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.
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:
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) 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 (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.gomap 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_opstable for durability and Lite-mode compatibility. - Atomic claiming via
ClaimBatchusesdedup_keygrouping andFOR UPDATE SKIP LOCKEDto prevent race conditions during batch processing. - Environment-driven configuration via
WEKNORA_ASYNQ_*_CONCURRENCYvariables allows runtime tuning of worker pool sizes without redeployment. - Interface abstractions like
TaskEnqueuerdecouple 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 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 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 reads these values, validates them, and falls back to sane defaults ranging from 2 to 12 workers per pool depending on the workload type.
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 →