How Queue Processor V2 Handles Async Research Requests with Encrypted Databases in Local Deep Research
Queue Processor V2 orchestrates asynchronous research jobs by maintaining a daemon thread that polls encrypted per-user SQLite databases to determine whether requests execute immediately or enter a queue, then manages worker threads and buffers progress updates until database access is available.
The Local Deep Research platform manages sensitive research configurations by storing user data in encrypted SQLite databases. To handle concurrent LLM-driven research workflows without blocking the main application, the system implements Queue Processor V2, a thread-safe orchestration layer that respects per-user encryption keys while efficiently scheduling background tasks. This architecture ensures that the queue processor v2 handles async research requests with encrypted databases through a combination of direct execution shortcuts, persistent queue management, and buffered asynchronous callbacks.
Core Architecture and Components
The QueueProcessorV2 class serves as the central coordinator in src/local_deep_research/web/queue/processor_v2.py. It runs as a long-lived daemon thread that continuously monitors for pending work across all active users.
Key components include:
QueueProcessorV2(lines 39‑53): Maintains the background processing loop, tracks users requiring attention, and manages thread-safe access to shared state.- Encrypted user database: The processor opens per-user SQLite files using
db_manager.open_user_databasewith session passwords stored insession_password_store(lines 14‑22). UserQueueService: Provides CRUD operations forQueuedResearch,UserActiveResearch, andResearchHistorytables (lines 94‑100).SettingsManager: Retrieves per-user configuration such asapp.queue_modeandapp.max_concurrent_researchesfrom the encrypted database (lines 30‑38).- Research service: The
start_research_processfunction insrc/local_deep_research/web/services/research_service.pyexecutes the actual LLM workflow in isolated worker threads. - Notification helpers: Located in
src/local_deep_research/notifications/queue_helpers.py, these send email, Slack, or Discord alerts upon completion or failure.
Request Flow: Direct Execution vs. Queueing
When a user submits a research request via an HTTP endpoint, the system calls queue_processor.notify_research_queued() (lines 13‑81). This method implements a dual-path decision logic:
Direct-mode shortcut: If the user's app.queue_mode setting equals "direct" and their current active research count is below app.max_concurrent_researches, the processor immediately invokes _start_research_directly. This path opens the encrypted database only to verify settings and active task counts before launching the worker thread.
Queue-mode persistence: When limits are reached or queue mode is enforced, the request metadata is saved via UserQueueService.add_task_metadata, and the username is added to the self._users_to_check set through notify_user_activity (lines 98‑100). The request waits in the encrypted database until resources become available.
from local_deep_research.web.queue.processor_v2 import queue_processor
def start_research(username, research_id, session_id, **kwargs):
"""
Called from the API endpoint when a user submits a request.
Triggers either immediate execution or queue persistence.
"""
queue_processor.notify_research_queued(
username,
research_id,
session_id=session_id,
**kwargs
)
The Background Processing Loop
The daemon thread created by start() runs _process_queue_loop (lines 95‑124), which executes every check_interval seconds. During each iteration:
- The loop copies the current
self._users_to_checkset to avoid mutation during iteration. - It invokes
_process_user_queue(username, session_id)for each active user. - It removes users whose queues are empty to reduce unnecessary database connections.
The _process_user_queue method (lines 54‑124 for setup, lines 125‑164 for decision logic) performs the following steps:
- Retrieves the session password from
session_password_store. - Opens the encrypted database via
db_manager.open_user_database. - Reads the
max_concurrent_researchessetting from the user's encrypted configuration. - Queries queue status via
UserQueueService.get_queue_statusto calculateavailable_slots = max_concurrent - active_tasks. - If slots exist, pulls up to that many
QueuedResearchrows ordered bypositionand starts each via_start_queued_researches.
To prevent resource leaks, each loop iteration calls cleanup_current_thread() (lines 28‑38), which closes per-thread SQLAlchemy engines and file descriptors.
def handle_user_request(username, session_id):
# Adds user to check-list and triggers immediate processing
pending = queue_processor.process_user_request(username, session_id)
if pending:
print(f"{pending} queued tasks will be processed")
Starting and Tracking Research Tasks
When the processor decides to start a queued research, _start_queued_researches (lines 138‑183) marks the database row with is_processing=True, updates the task status to "processing", and invokes _start_research.
The _start_research method (lines 190‑237) creates a UserActiveResearch record, updates ResearchHistory status to IN_PROGRESS, extracts the stored settings_snapshot, and launches start_research_process in a separate worker thread. This separation ensures that the main processing loop remains responsive while long-running LLM operations execute in isolation.
Buffered Asynchronous Updates
Because worker threads may generate progress callbacks when the main processing thread lacks immediate database access, Queue Processor V2 implements a buffering mechanism. The queue_progress_update and queue_error_update methods (lines 171‑188 and 194‑213 respectively) store update dictionaries in self.pending_operations, keyed by UUID.
When database access becomes available—typically during the next poll cycle or user-initiated request—the process_pending_operations_for_user method (lines 224‑298) drains the pending operations buffer and writes progress percentages or error messages directly to the ResearchHistory table in the encrypted database. This guarantees eventual consistency without blocking worker threads.
def report_progress(username, research_id, percent):
# Called from within the research worker thread
queue_processor.queue_progress_update(username, research_id, percent)
Completion Handling and Notifications
Upon research completion or failure, the worker thread invokes notify_research_completed or notify_research_failed (lines 106‑124 and 146‑165). These methods:
- Update the task status in the encrypted database.
- Remove the active research record.
- Invoke notification helpers from
src/local_deep_research/notifications/queue_helpers.pyto send alerts via configured channels.
def research_finished(username, research_id, password):
# Signals completion and triggers notifications
queue_processor.notify_research_completed(username, research_id, password)
Summary
- Queue Processor V2 operates as a daemon thread in
processor_v2.py, continuously polling encrypted per-user SQLite databases to manage research workflows. - It supports dual execution paths: immediate direct-mode execution when resources permit, or persistent queueing when limits are reached.
- Per-user encryption is maintained throughout the lifecycle, with database connections opened using session-specific passwords via
db_manager.open_user_database. - Asynchronous updates are buffered in memory and flushed to the encrypted database when access is available, preventing worker thread blocking.
- Resource cleanup occurs automatically via
cleanup_current_thread()to prevent file descriptor leaks across loop iterations. - Notification integration ensures users receive alerts upon research completion or failure through external channels.
Frequently Asked Questions
How does Queue Processor V2 maintain encryption when accessing user databases?
Queue Processor V2 retrieves session passwords from session_password_store and passes them to db_manager.open_user_database within _process_user_queue. This ensures each per-user SQLite database is decrypted only when the processor actively checks that specific user's queue, maintaining isolation between user data stores.
What determines whether a research request runs immediately or enters the queue?
The notify_research_queued method checks the user's app.queue_mode setting and current active research count against app.max_concurrent_researches. If queue mode is disabled and slots are available, the request executes immediately via _start_research_directly; otherwise, it persists to the QueuedResearch table and waits for the background loop.
How are progress updates handled when the database is temporarily unavailable?
Worker threads call queue_progress_update or queue_error_update, which store updates in the self.pending_operations dictionary. The process_pending_operations_for_user method later drains this buffer and writes changes to the encrypted database during the next available cycle, ensuring no updates are lost due to temporary connection constraints.
What happens when a research task fails or completes?
The processor invokes notify_research_completed or notify_research_failed, which updates the ResearchHistory status in the encrypted database and triggers notification helpers in src/local_deep_research/notifications/queue_helpers.py. These helpers send alerts through configured email, Slack, or Discord channels without blocking the main processing thread.
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 →