How to Debug Semantic Processing Pipeline Issues in OpenViking
To debug semantic processing pipeline issues in OpenViking, enable DEBUG logging to trace message flow, verify QueueManager worker thread health, inspect DagStats for stuck nodes, and instrument the _file_summary_task and _overview_task methods in semantic_dag.py to isolate failures.
The semantic processing pipeline in the volcengine/OpenViking repository automatically generates .abstract.md and .overview.md files for directories and pushes them to the embedding queue. Because the system relies on asynchronous workers and a lazy DAG executor, failures often manifest as stalled queues, missing summaries, or silently skipped directories. This guide walks you through systematic debugging using the actual source code paths and diagnostic hooks exposed by the SemanticQueue, SemanticProcessor, and SemanticDagExecutor components.
Understanding the Pipeline Architecture
Before debugging, map the three core components that handle every semantic message:
- SemanticQueue (
openviking/storage/queuefs/semantic_queue.py): WrapsNamedQueueto enqueueSemanticMsgobjects containing target URIs and recursion flags. - SemanticProcessor (
openviking/storage/queuefs/semantic_processor.py): The dequeue handler that receives messages and initiates the DAG executor. - SemanticDagExecutor (
openviking/storage/queuefs/semantic_dag.py): Executes a lazy DAG where directory overviews wait for child file summaries to complete.
The flow is orchestrated by QueueManager (openviking/storage/queuefs/queue_manager.py), which creates the queues and manages worker threads. When a message enters SemanticQueue.enqueue(), a worker eventually triggers SemanticProcessor.on_dequeue(), launching the recursive walk via SemanticDagExecutor.run().
Enable Verbose Logging
OpenViking uses the standard Python logging module centralized in openviking_cli/utils/logger.py. Switch from INFO to DEBUG to see state transitions in SemanticProcessor and SemanticDagExecutor.
from openviking_cli.utils.logger import get_logger
logger = get_logger("openviking")
logger.setLevel("DEBUG")
With debug logging active, you will see entries such as Processing semantic generation for: viking://user/1234/project (recursive=True) in semantic_processor.py, plus Dispatching directory … and File summary completed for … entries from the DAG executor. If these logs stop appearing mid-stream, the worker thread likely encountered an unhandled exception.
Verify Queue Health and Worker Status
Inspect Queue Depth
If enqueued messages never process, check whether the queue is actually draining:
import asyncio
from openviking.storage.queuefs import get_queue_manager
qm = get_queue_manager()
semantic_q = qm.get_queue(qm.SEMANTIC)
print("Queue size:", asyncio.run(semantic_q.size()))
A static size after enqueueing indicates the worker is stalled or dead.
Check Worker Thread Vitality
Each queue type has a dedicated worker thread managed by QueueManager._queue_worker_loop. Verify liveness directly:
import threading
for name, thread in qm._queue_threads.items():
print(f"{name} alive: {thread.is_alive()}")
If the SEMANTIC thread is not alive, inspect the logs for stack traces originating from semantic_processor.py or semantic_dag.py.
Examine DAG Statistics for Stuck Nodes
SemanticDagExecutor exposes runtime statistics via get_stats(), returning a DagStats dataclass defined in semantic_dag.py. Access this during a debugging session to pinpoint bottlenecks:
# Assuming you have a reference to the running processor
executor = processor._dag_executor # Set during on_dequeue
stats = executor.get_stats()
print(f"Total: {stats.total_nodes}, Pending: {stats.pending_nodes}, "
f"In-progress: {stats.in_progress_nodes}, Done: {stats.done_nodes}")
Interpret the metrics as follows:
pending_nodes> 0 indefinitely: A child directory never reported completion. Check the_overview_taskinsemantic_dag.py(lines 40–73) for the parent URI.in_progress_nodesstuck: A file summary task raised an exception without calling_on_file_done, leaving the DAG waiting forever.
Validate Individual Task Execution
Instrument File Summary Tasks
The _file_summary_task method (lines 44–58 in semantic_dag.py) generates LLM summaries for individual files. Wrap it to capture exceptions:
async def _file_summary_task_debug(self, parent_uri, file_path):
try:
await self._file_summary_task(parent_uri, file_path)
except Exception as e:
from openviking_cli.utils.logger import get_logger
logger = get_logger("openviking")
logger.error(f"DEBUG: file summary failed for {file_path}: {e}", exc_info=True)
raise
# Patch at runtime for debugging
processor._file_summary_task = _file_summary_task_debug.__get__(processor)
Monitor Overview Generation
Similarly, instrument _overview_task (lines 40–73 in semantic_dag.py) to log the exact file lists and child URIs it receives before building the directory overview. This reveals when the LLM receives empty or malformed context.
Check LLM Concurrency and Configuration
The pipeline throttles LLM calls via a semaphore initialized from max_concurrent_llm in openviking_cli/utils/config/vlm_config.py. If this value is too low for large directories, tasks block indefinitely waiting for the semaphore (self._llm_sem in the executor).
Verify the current limit:
from openviking_cli.utils.config import get_openviking_config
cfg = get_openviking_config()
print(f"Max concurrent LLM calls: {cfg.vlm.max_concurrent_llm}")
Rate-limit errors or network timeouts surface as warnings inside _generate_single_file_summary in semantic_processor.py (lines 46–55). Ensure the LLM client in openviking/utils/llm.py is correctly configured and responsive.
Isolate with a Minimal Test Case
Create a tiny directory with two or three text files and trigger the pipeline manually to rule out data-scale issues:
import asyncio
from openviking.storage.queuefs.semantic_msg import SemanticMsg
from openviking.storage.queuefs.semantic_queue import SemanticQueue
from openviking_cli.utils import VikingURI
async def minimal_test():
uri = VikingURI("viking://user/1234/tmp_test").uri
msg = SemanticMsg(
uri=uri,
recursive=True,
account_id="1234",
user_id="1234",
agent_id="",
role="ROOT",
context_type="default",
)
q = SemanticQueue(agfs=None) # In production, QueueManager injects agfs
await q.enqueue(msg)
asyncio.run(minimal_test())
Observe the logs and DagStats—if this succeeds while your production batch fails, the issue is likely resource exhaustion or a specific malformed file in the larger dataset.
Key Source Files and Debug Entry Points
Keep these files open when tracing failures:
| File | Key Function / Class | Debug Purpose |
|---|---|---|
openviking/storage/queuefs/semantic_queue.py |
SemanticQueue.enqueue() |
Verify message enqueue logic |
openviking/storage/queuefs/semantic_processor.py |
on_dequeue() (lines 53–84) |
Entry point for DAG creation |
openviking/storage/queuefs/semantic_dag.py |
run(), _dispatch_dir() (lines 71–98) |
Recursive directory dispatch |
openviking/storage/queuefs/semantic_dag.py |
_file_summary_task() (lines 44–58) |
File-level LLM summary |
openviking/storage/queuefs/semantic_dag.py |
_overview_task() (lines 40–73) |
Directory overview generation |
openviking/storage/queuefs/semantic_dag.py |
DagStats, get_stats() |
Health metrics |
openviking/storage/queuefs/queue_manager.py |
_queue_worker_loop() (lines 64–84) |
Worker thread management |
openviking_cli/utils/logger.py |
get_logger() |
Logging configuration |
openviking_cli/utils/config/vlm_config.py |
max_concurrent_llm |
Concurrency limits |
Summary
- Enable DEBUG logging via
openviking_cli.utils.loggerto trace the enqueue-to-completion lifecycle insemantic_processor.py. - Verify worker health by checking
QueueManager._queue_threadsliveness andSemanticQueue.size()trends. - Inspect DagStats from
SemanticDagExecutor.get_stats()to identify stuck pending or in-progress nodes. - Instrument tasks by wrapping
_file_summary_taskand_overview_taskinsemantic_dag.pyto catch silent exceptions. - Validate LLM configuration in
vlm_config.pyto ensure the semaphore does not block execution. - Test minimally with a single small directory to isolate data-specific corruption from systemic failures.
Frequently Asked Questions
Why does my semantic queue size never decrease after enqueueing?
The worker thread responsible for the SEMANTIC queue may be dead or blocked. Check qm._queue_threads for thread liveness and review logs for unhandled exceptions in SemanticProcessor.on_dequeue(). If the thread is alive but idle, the DAG executor may be waiting on a semaphore or stuck LLM call.
How do I identify which file is causing the pipeline to hang?
Access the running SemanticDagExecutor instance via processor._dag_executor and call get_stats(). If in_progress_nodes is non-zero, the stuck node corresponds to a file currently inside _file_summary_task. Temporarily patch that method to log every file_path entry before awaiting the LLM call.
What causes empty .overview.md files in large directories?
Empty overviews usually indicate that _overview_task ran before child file summaries completed, or the LLM returned an empty response. Check DagStats.pending_nodes—if the count dropped to zero prematurely, a child task may have crashed without signaling completion. Also verify max_concurrent_llm is sufficient for your directory size to prevent timeouts.
Where is the LLM client configured for the semantic pipeline?
The LLM client is configured in openviking/utils/llm.py and consumed by semantic_processor.py. Concurrency is controlled by max_concurrent_llm defined in openviking_cli/utils/config/vlm_config.py. Reduce this value if you hit rate limits, or increase it if the pipeline stalls on IO-bound LLM calls.
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 →