How the Capture Queue Manages Worker Claims in Claude-Obsidian

The Claude-Obsidian capture queue ensures exclusive worker ownership of ingestion tasks by atomically updating a durable JSON queue file with cryptographically random claim tokens, process metadata, and host identifiers, while detecting concurrent modifications through strict inode and timestamp validation.

The capture queue in the AgriciDaniel/claude-obsidian repository orchestrates file ingestion and URL downloading by allowing multiple workers to safely claim discrete work items without race conditions. Understanding how this system manages worker claims is essential for building reliable automation around the vault's capture functionality. The implementation relies on atomic filesystem operations and durable state tracking stored in the vault's metadata directory.

Queue Storage Architecture and Entry States

The Durable JSON Queue Location

Persistent queue state lives at .vault-meta/capture/queue.json within the vault root. Each entry maintains an id, state field (one of queued, claimed, completed, or failed), and metadata. When a worker acquires a task, the entry gains a claim sub-object containing ownership verification data.

The Worker Claim Lifecycle

Locating Entries via Atomic Reads

The CaptureQueue.claim(item_id, worker) method in claude_obsidian/capture.py (around line 1500) begins by invoking _read_runtime_regular to load the queue file atomically. This ensures the worker sees a consistent snapshot before attempting modification.

Cryptographic Token Generation

Before recording ownership, the system generates a 32-character lowercase hexadecimal token using os.urandom(16).hex() (lines 1528-1534). This cryptographically random value acts as a bearer token that workers must present to release or update the claim.

Recording Claim Metadata

The claim object is constructed with fields validated by _validate_queue_claim (lines 1467-1486):

  • token: The generated hex string
  • pid: Current process ID from os.getpid()
  • host: Hostname via socket.gethostname()
  • worker: Optional CLI identifier from --worker argument
  • claimed_at: ISO-8601 timestamp from _utc_now()
  • claimed_epoch: Monotonic time from time.time()

Atomic Updates and Conflict Detection

The critical section uses _atomic_runtime_write (lines 354-380) to write queue changes:

  1. Serialize modified JSON to a temporary file
  2. Verify no concurrent modifications by comparing st_dev, st_ino, st_size, st_mtime_ns, and st_mode against initial stats
  3. Execute os.replace(..., src_dir_fd=runtime_fd, dst_dir_fd=runtime_fd) for atomic rename

If stat fields differ, the system raises CaptureConflict (line 393), preventing two workers from claiming the same item simultaneously. Upon success, claim() returns the token (line 1552).

Verifying and Releasing Claims

Token Validation in Completion Methods

The CaptureQueue.complete(item_id, token, result) and CaptureQueue.fail(item_id, token, error) methods verify ownership before state transitions. The validation logic (lines 1608-1625) confirms:

  1. The supplied token matches the stored claim.token
  2. The owning process remains alive via _process_alive(pid)

Recovering Stale Claims

If a worker crashes, the PID check fails, marking the claim as stale. Other workers can then reclaim the item using queue.resume(item_id, stale_after=3600), which bypasses the token requirement for orphaned entries while respecting the staleness threshold.

Implementation Example: Claiming Work from the Queue

from claude_obsidian.capture import CaptureQueue

# Initialize queue pointing to vault root

queue = CaptureQueue("/path/to/vault")

# Claim a specific work item

item_id = "00a1b2c3d4e5f6a7b8c9d0e1"
claim_token = queue.claim(item_id=item_id, worker="worker-1")
print(f"Acquired claim token: {claim_token}")  # 32-char hex string

# Process the capture work...

result = {"status": "ok", "ingested_files": ["note.md"]}

# Release claim and mark complete

queue.complete(item_id=item_id, token=claim_token, result=result)

# Handle worker crashes: reclaim stale work after 1 hour

queue.resume(item_id=item_id, stale_after=3600)

Summary

  • The capture queue resides at .vault-meta/capture/queue.json and tracks four states: queued, claimed, completed, and failed.
  • Atomic file operations via os.replace with directory file descriptors ensure no interleaved writes corrupt the queue.
  • Each claim contains a 32-character cryptographic token, process ID, hostname, and dual timestamps for robust ownership verification.
  • Conflict detection compares inode and modification time metadata to prevent race conditions during concurrent claims.
  • Stale claims from crashed workers are identified via PID liveness checks, allowing safe reclamation by other processes.

Frequently Asked Questions

What prevents two workers from claiming the same queue item simultaneously?

The system employs optimistic concurrency control: before committing a claim, _atomic_runtime_write verifies that the queue file's st_ino, st_size, and st_mtime_ns remain unchanged since reading. If another process modified the file, CaptureConflict aborts the operation.

How does the capture queue detect crashed workers?

During completion or failure operations, the queue validates the original claim's PID using _process_alive(). If the process no longer exists, the claim is considered stale and can be reclaimed via resume() without presenting the original token.

Why does the claim token use os.urandom instead of UUID?

The implementation specifically uses os.urandom(16).hex() to generate 256 bits of entropy without external dependencies, creating an unpredictable bearer token that workers must possess to modify queue state, preventing unauthorized state transitions even if item IDs are guessable.

Where is the claim validation logic located?

Claim structure validation occurs in _validate_queue_claim at lines 1467-1486 of claude_obsidian/capture.py, while runtime verification (token matching and PID checks) is implemented in the complete and fail methods around lines 1608-1625.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →