Compare commits

..
8 Commits
Author SHA1 Message Date
jpmschweitzerandClaude Opus 4.5 9e7d8394f3 release: v1.4.8 - Paperless orphan cleanup
Build and Push / build (release) Successful in 29s
🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 17:04:24 +01:00
jpmschweitzerandClaude Opus 4.5 983a934b85 feat: add Paperless orphan cleanup endpoint
- POST /maintenance/cleanup/paperless - detect and clean orphaned Paperless documents
- Checks indexed documents against Paperless API
- Removes vectors and graph nodes for deleted documents
- Supports dry_run mode for preview

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 17:00:56 +01:00
jpmschweitzerandClaude Opus 4.5 6d5760c297 release: v1.4.7 - Paperless custom field fix
Build and Push / build (release) Successful in 29s
🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 16:40:28 +01:00
jpmschweitzerandClaude Opus 4.5 867de65354 fix: use field ID for Paperless custom field updates
Paperless API requires field ID (integer) not field name (string)
when updating custom fields. Now looks up field ID by name before
updating library_indexed custom field.

Also includes webhook debugging endpoint for development.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 16:37:26 +01:00
jpmschweitzerandClaude Opus 4.5 f2b8c7d111 fix: update Paperless webhook payload to match include_document format
Build and Push / build (release) Successful in 30s
- Change model field from document_id to id (Paperless sends id)
- Add content, created, modified, added, original_file_name, owner fields
- Add extra="ignore" config to handle additional Paperless fields
- Update sync service to use content from webhook payload
- Skip Paperless API call when content already provided

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 14:38:40 +01:00
jpmschweitzer f4352841a2 Merge feature/document-storage: Paperless-ngx integration
Build and Push / build (release) Successful in 30s
2025-12-25 14:17:11 +01:00
jpmschweitzerandClaude Opus 4.5 4ff3fc4c7a feat: add Paperless-ngx document storage integration
- Add /documents router with webhook, upload, search, health endpoints
- Create DocumentSyncService for indexing documents to vectors/graph
- Add PaperlessClient for REST API integration
- Configure dependency injection for Paperless client
- Add document models for webhook payloads and responses
- Event-driven architecture via Paperless workflow webhooks

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-25 14:16:40 +01:00
jpmschweitzerandClaude Opus 4.5 e6e65d6d78 feat: add test data cleanup endpoint
Build and Push / build (release) Successful in 28s
Add POST /maintenance/cleanup/test-data endpoint to purge LLM test data
from wiki, graph, and vectors. Security-restricted to test user namespace
only (users/llm-tester/*, users/llm_tester/*).

- Supports dry_run=true (default) to preview before deleting
- Cleans vectors, graph nodes, and wiki pages
- Scheduler task configured for weekly cleanup (Sunday 3:00 AM)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-24 21:12:16 +01:00
15 changed files with 2501 additions and 95 deletions
+3 -1
View File
@@ -7,6 +7,7 @@ QDRANT_PORT=6333
OLLAMA_URL=http://192.168.86.149:11434
SEARXNG_URL=http://192.168.86.149:8080
REDIS_HOST=192.168.86.149
PAPERLESS_URL=http://192.168.86.149:8091
OLLAMA_MODEL=mistral-nemo-large:latest
OLLAMA_EMBEDDING_MODEL=nomic-embed-text
@@ -20,4 +21,5 @@ WIKI_GRAPHQL_API=your_jwt_token_here
LIBRARY_API_KEY=key_here
NEO4J_PASSWORD=key_here
WIKIJS_DB_PASSWORD=key_here
SCHEDULER_API_KEY=key_here
SCHEDULER_API_KEY=key_here
PAPERLESS_TOKEN=key_here
+70
View File
@@ -5,6 +5,76 @@ All notable changes to Library Desk will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
## [1.4.8] - 2025-12-25
### Added
- **Paperless Orphan Cleanup** - `POST /maintenance/cleanup/paperless` endpoint
- Detects documents deleted from Paperless but still indexed in Library Desk
- Removes orphaned vectors and graph nodes
- Supports `dry_run=true` for preview mode
## [1.4.7] - 2025-12-25
### Fixed
- **Paperless Custom Field Update** - Fixed 400 error when marking documents as indexed
- Paperless API requires field ID (integer) not field name (string)
- Now looks up `library_indexed` field ID before updating
- Webhook params format: `doc_url` and `title` from Jinja templates
### Added
- **Webhook Debug Endpoint** - `POST /documents/webhook-capture` for development testing
## [1.4.6] - 2025-12-25
### Fixed
- **Paperless Webhook Payload Format** - Updated model to match Paperless `include_document=true` format
- Paperless sends `id` instead of `document_id`
- Paperless sends full document data including `content`, `title`, `tags`, etc.
- Webhook now uses content from payload, skipping extra Paperless API call
- Added `extra = "ignore"` to handle additional Paperless fields
## [1.4.5] - 2025-12-25
### Added
- **Document Storage Integration** - Paperless-ngx integration for PDFs, images, and documents
- Event-driven architecture via Paperless webhooks
- `POST /documents/webhook` - Receive document events from Paperless workflows
- `POST /documents/upload` - Upload files directly to Paperless
- `POST /documents/upload-url` - Download and upload documents from URL
- `POST /documents/search` - Semantic search across indexed documents
- `GET /documents/health` - Paperless connectivity health check
- **DocumentSyncService** - Indexes Paperless documents into vectors and graph
- Fetches document content via Paperless API
- Chunks text and generates embeddings for Qdrant
- Creates Document nodes in Neo4j knowledge graph
- Supports multi-tenancy via user parameter in webhook URL
- **PaperlessClient** - REST API client for Paperless-ngx
- Document retrieval, upload, and update operations
- Health check support
- **Paperless Workflow Configuration**
- Production workflow: Document Added (NOT tagged llm-test) → webhook to Library Desk
- Test workflow: Document Added (tagged llm-test) → webhook with test user
### Changed
- Updated `src/config.py` with Paperless configuration settings
- Added `PaperlessDep` dependency injection for document endpoints
## [1.4.4] - 2025-12-24
### Added
- **Test Data Cleanup Endpoint** - `POST /maintenance/cleanup/test-data`
- Purges LLM test data from wiki, graph, and vectors
- Security-restricted to test user namespace only (`users/llm-tester/*`, `users/llm_tester/*`)
- Supports `dry_run=true` (default) to preview before deleting
- Scheduler task configured for weekly cleanup (Sunday 3:00 AM)
## [1.4.3] - 2025-12-24
### Changed
+509
View File
@@ -0,0 +1,509 @@
# Phase 3: Document Storage System - Implementation Plan
## Overview
Document storage tier for Library Desk - storing and indexing PDFs, images, videos, and git documentation mirrors.
**User Decisions:**
- Paperless-ngx container for OCR
- Ebooks deferred to future phase
- Video.js player deferred to after core implementation
| Phase | Status | Version |
|-------|--------|---------|
| Phase 1: Cleanup System | Complete | v1.4.0 |
| Phase 2: Volatile Memory | Complete | v1.4.3 |
| Phase 3: Document Storage | Planning | - |
| Phase 4: Test Data Cleanup | Complete | v1.4.4 |
---
## Architecture
**Paperless-ngx as primary document store** (no SeaweedFS needed):
```
┌─────────────────────────────────────────────────────────────────┐
│ External Sources │
│ ┌─────────┐ ┌────────────┐ ┌──────────────┐ │
│ │ GitHub │ │ Direct │ │ Email/Folder │ │
│ │ Docs │ │ Upload │ │ Ingestion │ │
│ └────┬────┘ └─────┬──────┘ └──────┬───────┘ │
└───────┼─────────────┼────────────────┼──────────────────────────┘
│ │ │
▼ ▼ ▼
┌─────────────────────────────────────────────────────────────────┐
│ Paperless-ngx │
│ ┌───────────────────────────────────────────────────────────┐ │
│ │ - Document storage (PDFs, images, videos) │ │
│ │ - OCR via Tesseract (PDFs, images) │ │
│ │ - Web UI for browsing/tagging │ │
│ │ - REST API for integration │ │
│ └─────────────────────────┬─────────────────────────────────┘ │
└────────────────────────────┼────────────────────────────────────┘
│ REST API (sync)
┌─────────────────────────────────────────────────────────────────┐
│ Library Desk │
│ ┌───────────────────────────────────────────────────────────┐ │
│ │ DocumentSyncService │ │
│ │ - Polls Paperless for new/updated docs │ │
│ │ - Extracts text + metadata via API │ │
│ │ - Sends to vector/graph pipelines │ │
│ └─────────────────────────┬─────────────────────────────────┘ │
│ │ │
│ ┌────────────────┼────────────────┐ │
│ ▼ ▼ ▼ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Qdrant │ │ Neo4j │ │ Wiki.js │ │
│ │ (vectors)│ │ (graph) │ │ (catalog)│ │
│ └──────────┘ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────────┘
```
**File handling by type:**
| File Type | Paperless | Library Desk |
|-----------|-----------|--------------|
| PDFs | OCR → text | Index text → vectors/graph |
| Images | OCR → text | Index text → vectors/graph |
| Videos | Storage only | Index metadata → vectors/graph |
---
## Technology Stack
| Component | Purpose | Rationale |
|-----------|---------|-----------|
| **Paperless-ngx** | Document storage + OCR | All-in-one: storage, OCR, web UI, REST API |
| **ClamAV** | Virus scanning | Host OS install, pyclamd integration, better isolation |
| **PDF.js** | PDF viewer | Embeddable in Wiki.js (deferred) |
**Why Paperless-ngx as primary store:**
- Eliminates need for separate blob storage (SeaweedFS/MinIO)
- Built-in web UI for browsing and tagging
- Tesseract OCR with 100+ language support
- REST API for Library Desk integration
- Handles videos as raw files (no OCR, but stored)
- Email and folder watching for automatic ingestion
- Active community, well-maintained
---
## Paperless-ngx API Deep Dive
### Authentication
```
POST /api/token/
Body: {"username": "...", "password": "..."}
Response: {"token": "..."}
Header: Authorization: Token <token>
```
### Document Upload (for HybridRAG → Paperless)
```
POST /api/documents/post_document/
Content-Type: multipart/form-data
Fields:
- document (file, required)
- title (string)
- created (datetime)
- correspondent (ID)
- document_type (ID)
- storage_path (ID)
- tags (repeatable IDs)
- custom_fields (JSON array)
Response: {"task_id": "uuid"}
```
Track consumption: `GET /api/tasks/?task_id={uuid}` → returns document ID when complete
### Document Search
```
GET /api/documents/?query=search+terms # Full-text search
GET /api/documents/?more_like_id=123 # Similarity search
Response includes __search_hit__:
{
"score": 0.95,
"highlights": "<span>matched</span> text",
"rank": 0
}
```
### Custom Field Filtering
```
GET /api/documents/?custom_field_query=field_name__operation=value
Operations:
- exact, in, isnull, exists (all types)
- icontains, istartswith, iendswith (text)
- gt, gte, lt, lte, range (numeric/date)
- contains (document links)
```
### Bulk Operations
```
POST /api/documents/bulk_edit/
{
"documents": [1, 2, 3],
"method": "add_tag|remove_tag|set_correspondent|set_document_type|merge|split|...",
"parameters": {...}
}
```
### Webhooks (Push to Library Desk!)
Paperless workflows can trigger webhooks on document events:
| Trigger | When | Available Data |
|---------|------|----------------|
| Consumption Started | Before OCR | file_path, source, filename |
| Document Added | After OCR | content, tags, doc_type, correspondent, `{doc_url}` |
| Document Updated | On change | Same as Added |
| Scheduled | Time-based | Date offsets from document dates |
**Webhook Action**: POST to Library Desk endpoint with document data
### Organization Features
| Feature | Purpose | API Endpoint |
|---------|---------|--------------|
| Tags | Nested labels (5 levels deep) | `/api/tags/` |
| Correspondents | Source/destination | `/api/correspondents/` |
| Document Types | Classification | `/api/document_types/` |
| Storage Paths | File organization | `/api/storage_paths/` |
| Custom Fields | Extensible metadata | `/api/custom_fields/` |
### Custom Fields We Should Create
| Field Name | Type | Purpose |
|------------|------|---------|
| `source_url` | URL | Original download URL (for HybridRAG uploads) |
| `library_indexed` | Boolean | Sync status with Library Desk |
| `library_doc_id` | Text | Library Desk document reference |
| `collection` | Text | Logical grouping (e.g., "fastapi-docs") |
### External LLM Add-ons (Optional)
Community tools exist for Ollama integration:
- **[paperless-ai](https://github.com/clusterzx/paperless-ai)** - Auto-tagging, RAG chat
- **[paperless-gpt](https://github.com/icereed/paperless-gpt)** - LLM-enhanced OCR, auto-titling
**Recommendation:** Skip these - Library Desk already has Ollama integration for:
- Embedding (nomic-embed-text)
- LLM analysis (mistral-nemo)
- Entity extraction
- HybridRAG
We'll do our own classification/tagging via Library Desk after sync.
---
## Virus Scanning Integration
**ClamAV daemon + pyclamd** (no third-party REST wrappers):
```
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ File Upload │────►│ Library Desk │────►│ ClamAV Daemon │
│ (URL or file) │ │ (pyclamd) │ │ (clamd:3310) │
└─────────────────┘ └────────┬────────┘ └─────────────────┘
┌────────────┴────────────┐
▼ ▼
┌──────────┐ ┌──────────┐
│ Clean │ │ Infected │
│ ✓ │ │ ✗ │
└────┬─────┘ └────┬─────┘
│ │
▼ ▼
Upload to Paperless Reject + Log
```
### ClamAV Deployment (Host OS)
ClamAV runs on the host OS (not containerized) for better security isolation:
```bash
# Installed via apt on Ubuntu/Debian
# Config: /etc/clamav/clamd.conf
# TCPSocket 3310
# TCPAddr 0.0.0.0
```
Benefits: scans outside container isolation, single virus DB, survives container restarts.
### Library Desk Integration
```python
# src/clients/clamav_client.py
import pyclamd
class ClamAVClient:
def __init__(self, host: str, port: int = 3310):
self.cd = pyclamd.ClamdNetworkSocket(host, port)
async def scan_bytes(self, data: bytes) -> ScanResult:
"""Scan file bytes, return clean/infected status."""
result = self.cd.scan_stream(data)
if result is None:
return ScanResult(clean=True)
return ScanResult(clean=False, virus_name=result['stream'][1])
def ping(self) -> bool:
"""Health check."""
return self.cd.ping()
```
### Scan Points
| Location | When | Action on Infected |
|----------|------|-------------------|
| `/documents/upload` | Before Paperless upload | Reject with 400, log threat |
| HybridRAG web fetch | Before saving PDF | Skip file, log threat |
| `/documents/webhook` | Optional re-scan | Quarantine in Paperless |
### Config Settings
```python
# src/config.py
CLAMAV_HOST: str = "192.168.86.149" # Host OS IP (not container)
CLAMAV_PORT: int = 3310
CLAMAV_ENABLED: bool = True # Bypass for testing
CLAMAV_TIMEOUT: int = 30 # seconds
```
---
## Integration Strategy
### Option A: Webhook Push (Preferred)
```
Paperless Workflow → POST webhook → Library Desk /documents/webhook
```
- Real-time indexing when documents added/updated
- Configure in Paperless: Workflow → Document Added → Webhook Action
- Library Desk receives document ID, fetches content via API
### Option B: Polling Pull (Fallback)
```
Scheduler → POST /documents/sync → Library Desk polls Paperless
```
- Periodic sync for missed webhooks or initial bulk import
- Track `library_indexed` custom field to skip already-processed docs
### Option C: HybridRAG Upload (New!)
```
HybridRAG web search → finds PDF → POST to Paperless → webhook → indexed
```
- When HybridRAG finds a relevant PDF/document in web results
- Download and upload to Paperless with `source_url` custom field
- Paperless OCRs it, triggers webhook, Library Desk indexes
---
## Library Desk API Design
### Documents Router (`/documents`)
| Endpoint | Method | Purpose |
|----------|--------|---------|
| `/documents/webhook` | POST | Receive Paperless webhook (Document Added/Updated) |
| `/documents/sync` | POST | Pull new/updated docs from Paperless → index |
| `/documents/upload` | POST | Upload file to Paperless (for HybridRAG) |
| `/documents/sync-from-git` | POST | Pull docs from Gitea → upload to Paperless → index |
| `/documents/{document_id}` | GET | Get document metadata |
| `/documents/{document_id}/text` | GET | Get extracted text |
| `/documents/search` | POST | Semantic search across documents |
| `/documents/collection/{name}` | GET | List documents in collection |
| `/documents/collection/{name}/catalog` | POST | Generate wiki catalog page |
**Upload flow (HybridRAG → Paperless):**
1. HybridRAG finds PDF in web results
2. POST `/documents/upload` with URL or file
3. Library Desk downloads, uploads to Paperless with metadata
4. Returns task_id for async tracking
5. Paperless webhook triggers indexing when OCR complete
### Viewers Router (`/viewers`) - Deferred
| Endpoint | Method | Purpose |
|----------|--------|---------|
| `/viewers/pdf/{document_id}` | GET | Serve PDF.js viewer |
| `/viewers/image/{document_id}` | GET | Serve image lightbox |
| `/viewers/video/{document_id}` | GET | Serve Video.js player |
---
## Data Flow: Document Processing Pipeline
```
1. INTAKE (Paperless-ngx handles this)
└─ Upload via Paperless UI, email, or folder watch
└─ Paperless assigns document ID and stores file
2. OCR EXTRACTION (Paperless-ngx handles this)
├─ PDFs → Tesseract → Plain text
├─ Images → Tesseract → Plain text
└─ Videos → Metadata only (no OCR)
3. SYNC TO LIBRARY DESK (scheduled or manual)
└─ Poll Paperless API for new/updated documents
└─ Fetch text content + metadata
4. TEXT CHUNKING
└─ VectorService._chunk_text() (existing)
5. EMBEDDING
└─ OllamaClient.embed() (existing)
6. VECTOR STORAGE (Qdrant)
└─ Payload: {doc_type: "document", paperless_id, ...}
7. GRAPH STORAGE (Neo4j)
└─ Document node + MENTIONS relationships
8. WIKI CATALOG (optional)
└─ Auto-generate catalog page via ConsolidationService
```
---
## Git Docs Integration
Extends existing `scheduler/src/executors/doc_sync_executor.py`:
1. **Scheduler** syncs docs from GitHub → Gitea (existing)
2. **Post-sync hook** calls `POST /documents/sync-from-git`
3. **Library Desk** indexes docs into vectors/graph
4. **Auto-generate** wiki catalog page for collection
---
## Wiki.js Viewer Integration
Since Wiki.js v2 requires disabled HTML sanitization for iframes:
```markdown
<!-- In wiki catalog page -->
## Document Preview
<iframe
src="http://library-desk:8089/viewers/pdf/abc123"
width="100%" height="600px">
</iframe>
```
**Wiki.js Settings Required:**
- `Administration > Security > Allowed HTML Elements: iframe`
- `Content Security Policy: frame-src http://library-desk:8089`
---
## Implementation Phases
### Phase 3.1: Infrastructure Setup
- [ ] Deploy Paperless-ngx container (Docker Compose)
- [x] ClamAV installed on host OS (port 3310)
- [ ] Configure Paperless: storage path, OCR settings, API token
- [ ] Create custom fields in Paperless: `source_url`, `library_indexed`, `library_doc_id`, `collection`
- [ ] Create `src/clients/paperless_client.py`
- [ ] Create `src/clients/clamav_client.py` (pyclamd wrapper)
- [ ] Create `src/models/document.py`
- [ ] Add config settings to `src/config.py` (PAPERLESS_*, CLAMAV_*)
### Phase 3.2: Webhook Integration (Push)
- [ ] Create `src/routers/documents.py`
- [ ] Implement `/documents/webhook` endpoint (receives Paperless events)
- [ ] Configure Paperless Workflow: Document Added → Webhook → Library Desk
- [ ] Create `src/services/document_sync_service.py`
- [ ] Implement document indexing pipeline (fetch text → chunk → embed → graph)
### Phase 3.3: Polling Sync (Pull Fallback)
- [ ] Implement `/documents/sync` endpoint
- [ ] Poll Paperless for docs where `library_indexed=false`
- [ ] Track sync state (last_sync timestamp in Redis)
- [ ] Update `library_indexed` after successful indexing
### Phase 3.4: HybridRAG Upload Integration
- [ ] Implement `/documents/upload` endpoint
- [ ] Download file from URL
- [ ] **Virus scan before upload** (reject if infected, log threat)
- [ ] Upload clean files to Paperless with metadata
- [ ] Set `source_url` custom field
- [ ] Extend HybridRAG service to detect and upload relevant PDFs
- [ ] Add `save_to_documents` option to HybridRAG config
### Phase 3.5: Indexing Pipeline
- [ ] Extend VectorService for `doc_type: "document"`
- [ ] Extend GraphService for Document nodes (link to Paperless ID)
- [ ] Implement `/documents/search` endpoint
- [ ] Add dependency injection
### Phase 3.6: Git Docs Integration
- [ ] Create `src/clients/gitea_client.py`
- [ ] Implement `/documents/sync-from-git` → bulk upload to Paperless
- [ ] Create collection auto-cataloging (wiki pages)
- [ ] Add scheduler task for periodic git sync
### Phase 3.7: Viewers (Deferred)
*After core implementation is working*
- [ ] Create `static/pdf-viewer.html` (PDF.js)
- [ ] Create `static/image-viewer.html`
- [ ] Create `static/video-player.html` (Video.js)
- [ ] Create `src/routers/viewers.py`
### Phase 3.8: Maintenance & Testing
- [ ] Extend cleanup for document orphans
- [ ] Add document orphan detection (Paperless deleted but still in Qdrant/Neo4j)
- [ ] Create `tests/test_document_sync.py`
- [ ] Create `tests/test_paperless_client.py`
---
## Files to Create
| Path | Purpose |
|------|---------|
| `src/clients/paperless_client.py` | Paperless-ngx REST API client |
| `src/clients/clamav_client.py` | ClamAV scanner (pyclamd wrapper) |
| `src/clients/gitea_client.py` | Gitea repo access |
| `src/models/document.py` | Document/Collection/ScanResult models |
| `src/services/document_sync_service.py` | Sync orchestrator |
| `src/routers/documents.py` | Document endpoints (webhook, sync, upload, search) |
| `tests/test_document_sync.py` | Sync service tests |
| `tests/test_paperless_client.py` | API client tests |
| `tests/test_clamav_client.py` | Virus scanner tests |
| `docker/docker-compose.documents.yml` | Paperless + ClamAV deployment |
**Deferred files (Phase 3.7):**
| Path | Purpose |
|------|---------|
| `src/routers/viewers.py` | Viewer endpoints |
| `static/pdf-viewer.html` | PDF.js viewer |
| `static/image-viewer.html` | Image lightbox |
| `static/video-player.html` | Video.js player |
## Files to Modify
| Path | Changes |
|------|---------|
| `src/config.py` | `PAPERLESS_*`, `CLAMAV_*` settings |
| `src/core/dependencies.py` | DocumentSyncService, PaperlessClient, ClamAVClient DI |
| `src/main.py` | Register documents router |
| `src/services/vector_service.py` | `doc_type: "document"` handling |
| `src/services/graph_service.py` | Document node with Paperless ID |
| `src/services/hybrid_rag_service.py` | Add `save_to_documents` option + virus scan |
| `src/models/hybrid_rag.py` | Add `save_to_documents` config |
| `src/routers/maintenance.py` | Document orphan cleanup, ClamAV health check |
| `requirements.txt` | Add `pyclamd` |
## Paperless Custom Fields Setup
Create these in Paperless UI (Administration → Custom Fields):
| Field | Type | Purpose |
|-------|------|---------|
| `source_url` | URL | Original download URL |
| `library_indexed` | Boolean | Sync status |
| `library_doc_id` | Text | Library Desk reference |
| `collection` | Text | Logical grouping |
+139 -90
View File
@@ -6,15 +6,24 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
| Tier | Storage | Purpose | TTL |
|------|---------|---------|-----|
| **Volatile** | Redis | Weather, news, financial, ephemeral context | 5min - 2hr |
| **Documents** | TBD (research) | Git mirrors, PDFs, video, images | Permanent |
| **Volatile** | Qdrant (vectors) | Weather, news, financial, ephemeral context | 5min - 2hr |
| **Documents** | Paperless-ngx + ClamAV (host) | Git mirrors, PDFs, video, images | Permanent |
| **Knowledge** | Wiki + Neo4j | Personal dossiers, research, summaries | Permanent |
**Implementation Priority**: Cleanup → Volatile → Documents
**Implementation Priority**: Cleanup → Volatile → Documents → Test Data Cleanup
### Phase Status
| Phase | Status | Version |
|-------|--------|---------|
| Phase 1: Cleanup System | ✅ Complete | v1.4.0 |
| Phase 2: Volatile Memory | ✅ Complete | v1.4.3 |
| Phase 3: Document Storage | ✅ Planned | See [DOCUMENT_STORAGE_PLAN.md](DOCUMENT_STORAGE_PLAN.md) |
| Phase 4: Test Data Cleanup | ✅ Complete | v1.4.4 |
---
## Phase 1: Cleanup System Completion
## Phase 1: Cleanup System Completion
### Current State
- **COMPLETE** - All Phase 1 tasks implemented
@@ -57,9 +66,9 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
---
## Phase 2: Volatile Memory System
## Phase 2: Volatile Memory System
### Architecture
### Architecture (Final Implementation)
```
┌─────────────────┐ ┌──────────────┐ ┌─────────────────┐
@@ -71,88 +80,46 @@ A three-tier memory architecture for Library Desk with intelligent orchestration
┌─────────────────┐
Redis
(DB 4, TTL)
Qdrant
(volatile_{user})
└─────────────────┘
```
### Data Model
**Key design decisions:**
- Vector storage in Qdrant (not Redis) for semantic search
- Collection per user: `volatile_{user}`
- TTL via `ttl_expiry` timestamp in payload
- Natural language conversion for embedding structured data
- Integrated into HybridRAG with priority boost
```python
class VolatileRecord(BaseModel):
key: str # e.g., "weather:rotterdam"
namespace: str # e.g., "weather", "news", "financial"
data: dict # Actual content
source: Optional[str] # Origin API/service
created_at: datetime
updated_at: datetime
ttl: int # Seconds until expiration
refresh_schedule: Optional[str] # Cron expression, if repeating
user: str # Multi-tenant isolation
```
**Key pattern**: `{user}:volatile:{namespace}:{key_hash}`
### Implementation Order: Integration-First
1. **Start with Consolidation Hook** - Understand data flow through existing system
2. **Build Service Layer** - VolatileCacheService with Redis operations
3. **Add API Endpoints** - REST interface for volatile data
4. **Biographer Integration** - Query user preferences for relevance
### Tasks
#### 2.1 Integrate with Consolidation (FIRST)
**New file**: `src/services/volatile_service.py`
```python
class VolatileCacheService:
async def get(user, namespace, key) -> Optional[VolatileRecord]
async def set(user, namespace, key, data, ttl, refresh_schedule=None)
async def delete(user, namespace, key)
async def list_namespace(user, namespace) -> List[str]
async def get_scheduled(user) -> List[VolatileRecord] # For scheduler
```
#### 2.2 Create Volatile API Router
**New file**: `src/routers/volatile.py`
### Endpoints (Implemented)
| Endpoint | Method | Purpose |
|----------|--------|---------|
| `/volatile/{namespace}/{key}` | GET | Retrieve record |
| `/volatile/{namespace}/{key}` | POST | Store/update record |
| `/volatile/search?q=...` | GET | Semantic search across volatile data |
| `/volatile/store?namespace=...&key=...` | POST | Store/update record |
| `/volatile/{namespace}/{key}` | GET | Retrieve specific record |
| `/volatile/{namespace}/{key}` | DELETE | Remove record |
| `/volatile/{namespace}` | GET | List keys in namespace |
| `/volatile/scheduled` | GET | List records needing refresh |
| `/volatile/stats` | GET | Cache statistics |
| `/volatile/scheduled` | GET | Records needing refresh |
| `/volatile/namespaces` | GET | List available namespaces |
| `/maintenance/cleanup/volatile` | POST | Purge expired records |
#### 2.3 Integrate with Consolidation
**File**: `src/services/consolidation_service.py`
### Namespaces
Add relevance trigger detection:
1. During consolidation, analyze search results for location/interest patterns
2. Query tatlock's Biographer collection for user preferences
3. If match found, create/update volatile refresh schedule
#### 2.4 Biographer Integration
**File**: `src/core/dependencies.py`
```python
def get_biographer_qdrant() -> QdrantClientWrapper:
"""Direct access to tatlock's Biographer collection."""
# Configure to connect to tatlock's Qdrant
```
#### 2.5 Scheduler-Side Configuration
Document required scheduler tasks:
```json
{
"task_name": "volatile_refresh",
"schedule": "*/15 * * * *",
"endpoint": "GET /volatile/scheduled",
"follow_up": "For each record, call refresh endpoint with record.refresh_schedule"
}
```
| Namespace | Default TTL | Use Case |
|-----------|-------------|----------|
| weather | 30 min | Current conditions, forecasts |
| news | 1 hour | Headlines, breaking news |
| financial | 5 min | Stock prices, exchange rates |
| transit | 5 min | Train/bus schedules, delays |
| traffic | 10 min | Commute times, road conditions |
| air_quality | 1 hour | Pollution, pollen counts |
| sports | 1 min | Live scores, matches |
| social | 10 min | Social notifications |
| system | 1 min | Service health status |
| context | 1 hour | Session state |
| custom | 1 hour | User-defined data |
---
@@ -219,34 +186,116 @@ Add LLM-powered category descriptor generation:
---
## Files to Modify/Create
## Phase 4: LLM Tester Data Cleanup ✅
### Phase 1 (Cleanup)
- `src/routers/maintenance.py` - Add timestamp tracking
### Problem
LLM testing creates accumulated cruft across the system:
- Wiki.js pages under `llm-tester/` and `llm_tester/` paths
- Graph nodes (Document, Entity) linked to test pages
- Vector chunks in Qdrant for test content
This data accumulates over time and clutters Wiki.js visually (no separate tenant scope for tests).
### Solution
Add a maintenance endpoint to purge all LLM tester artifacts across wiki, graph, and vectors.
### Tasks
#### 4.1 Identify Test Data Patterns ✅
**Patterns matched** (security-restricted to test user namespace):
- `users/llm-tester/*`
- `users/llm_tester/*`
#### 4.2 Add Cleanup Endpoint ✅
**File**: `src/routers/maintenance.py`
```python
@router.post("/cleanup/test-data")
async def cleanup_test_data(
dry_run: bool = Query(default=True),
wiki: WikiJSDep = None,
vector_service: VectorServiceDep = None,
graph_service: GraphServiceDep = None,
api_key: str = Depends(verify_api_key)
):
"""
Purge LLM tester data from wiki, graph, and vectors.
**Security**: Only deletes pages in the test user namespace:
- users/llm-tester/*
- users/llm_tester/*
Use dry_run=true to preview what would be deleted.
"""
```
#### 4.3 Implementation Steps ✅
1. **Wiki cleanup**: Delete pages via GraphQL mutation
2. **Graph cleanup**: Delete Document nodes using `delete_page()` method
3. **Vector cleanup**: Delete chunks using `delete_page_chunks()` method
#### 4.4 Scheduler Integration ✅
**Recommended schedule**: Weekly (Sunday 3:00 AM)
```json
{
"task_name": "test_data_cleanup",
"schedule": "0 3 * * 0",
"endpoint": "POST /maintenance/cleanup/test-data?dry_run=false",
"description": "Weekly cleanup of LLM test data"
}
```
### Files to Modify
- `src/routers/maintenance.py` - Add cleanup endpoint
- `src/services/wiki_service.py` - Add bulk delete by path pattern (if needed)
- `src/services/graph_service.py` - May need pattern-based node deletion
- `src/services/vector_service.py` - Add pattern-based chunk deletion
---
## Files Modified/Created
### Phase 1 (Cleanup) ✅
- `src/routers/maintenance.py` - Timestamp tracking, cleanup endpoints
- `src/services/graph_service.py` - Bidirectional validation
- `src/services/vector_service.py` - Cross-reference checks
- `LIBRARIAN_INTEGRATION.md` - Scheduler config docs
### Phase 2 (Volatile)
- `src/services/volatile_service.py` - **NEW**
- `src/routers/volatile.py` - **NEW**
- `src/models/volatile.py` - **NEW**
- `src/core/dependencies.py` - Add Biographer client
- `src/services/consolidation_service.py` - Relevance triggers
- `tests/test_volatile.py` - **NEW**
### Phase 2 (Volatile)
- `src/services/volatile_service.py` - Qdrant-based volatile cache
- `src/routers/volatile.py` - Simplified endpoints
- `src/models/volatile.py` - Namespaces and models
- `src/models/hybrid_rag.py` - Volatile config options
- `src/services/hybrid_rag_service.py` - Volatile integration
- `src/clients/qdrant_client.py` - Expiry filter methods
- `tests/test_volatile.py` - 37 tests
### Phase 3 (Documents)
- `docs/DOCUMENT_STORAGE_RESEARCH.md` - **NEW**
- `src/services/document_store_service.py` - **NEW** (post-research)
- `src/routers/documents.py` - **NEW** (post-research)
### Phase 4 (Test Data Cleanup)
- `src/routers/maintenance.py` - Add cleanup endpoint
- `src/services/wiki_service.py` - Bulk delete by path pattern
- `src/services/graph_service.py` - Pattern-based node deletion
- `src/services/vector_service.py` - Pattern-based chunk deletion
---
## Resolved Design Decisions
1. **Biographer Qdrant**: Same Qdrant instance, different collection. Library-Desk queries directly.
2. **Scheduler API**: Has REST API for task registration. Library-Desk can programmatically create refresh schedules.
3. **External API calls**: Library-Desk routes through SearXNG for web search. Consider dedicated API integrations for high-value volatiles (weather, financial) for consistent quality.
1. **Volatile Storage**: Qdrant vectors (not Redis) for semantic search capability
2. **Collection Naming**: `volatile_{user}` for per-user isolation
3. **TTL Mechanism**: `ttl_expiry` timestamp in payload, background cleanup job
4. **HybridRAG Integration**: Volatile as third source with RRF priority boost
5. **Biographer Qdrant**: Same Qdrant instance, different collection
6. **Scheduler API**: Has REST API for task registration
---
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "library-desk"
version = "1.4.3"
version = "1.4.8"
description = "Coordination service for The Library system - HybridRAG queries, document ingestion, entity extraction, and knowledge consolidation"
readme = "README.md"
requires-python = ">=3.12"
+488
View File
@@ -0,0 +1,488 @@
"""
Paperless-ngx API client for Library Desk.
Provides async document management via Paperless-ngx:
- Document upload and retrieval
- Search and filtering
- Custom field management
- Task status tracking
"""
import httpx
from typing import Optional, List, Dict, Any
from dataclasses import dataclass
import logging
logger = logging.getLogger(__name__)
@dataclass
class PaperlessDocument:
"""Represents a document from Paperless-ngx."""
id: int
title: str
content: str
created: Optional[str] = None
modified: Optional[str] = None
added: Optional[str] = None
correspondent: Optional[int] = None
document_type: Optional[int] = None
storage_path: Optional[int] = None
tags: List[int] = None
archive_serial_number: Optional[int] = None
original_file_name: Optional[str] = None
archived_file_name: Optional[str] = None
custom_fields: List[Dict[str, Any]] = None
def __post_init__(self):
if self.tags is None:
self.tags = []
if self.custom_fields is None:
self.custom_fields = []
@dataclass
class SearchHit:
"""Search result with relevance info."""
document: PaperlessDocument
score: float
rank: int
highlights: Optional[str] = None
class PaperlessClient:
"""
Paperless-ngx REST API client.
Documentation: https://docs.paperless-ngx.com/api/
"""
def __init__(self, base_url: str, token: str, timeout: int = 30):
"""
Initialize Paperless-ngx client.
Args:
base_url: Paperless-ngx base URL (e.g., "http://paperless:8000")
token: API token for authentication
timeout: Request timeout in seconds
"""
self.base_url = base_url.rstrip("/")
self.api_url = f"{self.base_url}/api"
self.headers = {
"Authorization": f"Token {token}",
"Accept": "application/json",
}
self.client = httpx.AsyncClient(timeout=float(timeout), headers=self.headers)
logger.info(f"Initialized Paperless client: {base_url}")
async def close(self):
"""Close HTTP client."""
await self.client.aclose()
# =========================================================================
# Document Operations
# =========================================================================
async def get_document(self, document_id: int) -> Optional[PaperlessDocument]:
"""
Get a document by ID.
Args:
document_id: Paperless document ID
Returns:
PaperlessDocument or None if not found
"""
try:
response = await self.client.get(f"{self.api_url}/documents/{document_id}/")
response.raise_for_status()
data = response.json()
return self._parse_document(data)
except httpx.HTTPStatusError as e:
if e.response.status_code == 404:
return None
logger.error(f"Failed to get document {document_id}: {e}")
raise
except Exception as e:
logger.error(f"Failed to get document {document_id}: {e}")
raise
async def get_document_content(self, document_id: int) -> Optional[str]:
"""
Get extracted text content of a document.
Args:
document_id: Paperless document ID
Returns:
Text content or None if not found
"""
doc = await self.get_document(document_id)
return doc.content if doc else None
async def list_documents(
self,
page: int = 1,
page_size: int = 25,
ordering: str = "-added",
correspondent: Optional[int] = None,
document_type: Optional[int] = None,
tags: Optional[List[int]] = None,
) -> Dict[str, Any]:
"""
List documents with pagination and filtering.
Args:
page: Page number (starts at 1)
page_size: Results per page
ordering: Sort order (prefix with - for descending)
correspondent: Filter by correspondent ID
document_type: Filter by document type ID
tags: Filter by tag IDs
Returns:
Paginated response with count, next, previous, results
"""
params = {
"page": page,
"page_size": page_size,
"ordering": ordering,
}
if correspondent:
params["correspondent__id"] = correspondent
if document_type:
params["document_type__id"] = document_type
if tags:
params["tags__id__in"] = ",".join(str(t) for t in tags)
try:
response = await self.client.get(f"{self.api_url}/documents/", params=params)
response.raise_for_status()
data = response.json()
return {
"count": data.get("count", 0),
"next": data.get("next"),
"previous": data.get("previous"),
"results": [self._parse_document(d) for d in data.get("results", [])],
}
except Exception as e:
logger.error(f"Failed to list documents: {e}")
raise
async def search_documents(
self,
query: str,
page: int = 1,
page_size: int = 25,
) -> List[SearchHit]:
"""
Full-text search documents.
Args:
query: Search query string
page: Page number
page_size: Results per page
Returns:
List of SearchHit with document and relevance info
"""
params = {
"query": query,
"page": page,
"page_size": page_size,
}
try:
response = await self.client.get(f"{self.api_url}/documents/", params=params)
response.raise_for_status()
data = response.json()
results = []
for item in data.get("results", []):
doc = self._parse_document(item)
hit_info = item.get("__search_hit__", {})
results.append(SearchHit(
document=doc,
score=hit_info.get("score", 0.0),
rank=hit_info.get("rank", 0),
highlights=hit_info.get("highlights"),
))
return results
except Exception as e:
logger.error(f"Search failed for '{query}': {e}")
raise
async def upload_document(
self,
file_content: bytes,
filename: str,
title: Optional[str] = None,
correspondent: Optional[int] = None,
document_type: Optional[int] = None,
tags: Optional[List[int]] = None,
custom_fields: Optional[List[Dict[str, Any]]] = None,
) -> str:
"""
Upload a document to Paperless-ngx.
Args:
file_content: File bytes
filename: Original filename
title: Document title (optional, derived from filename if not set)
correspondent: Correspondent ID
document_type: Document type ID
tags: List of tag IDs
custom_fields: List of custom field values
Returns:
Task UUID for tracking consumption status
"""
files = {"document": (filename, file_content)}
data = {}
if title:
data["title"] = title
if correspondent:
data["correspondent"] = correspondent
if document_type:
data["document_type"] = document_type
if tags:
# Tags need to be sent multiple times for multiple values
data["tags"] = tags
if custom_fields:
data["custom_fields"] = custom_fields
try:
response = await self.client.post(
f"{self.api_url}/documents/post_document/",
files=files,
data=data,
)
response.raise_for_status()
result = response.json()
task_id = result.get("task_id", "")
logger.info(f"Uploaded document '{filename}', task_id: {task_id}")
return task_id
except Exception as e:
logger.error(f"Failed to upload document '{filename}': {e}")
raise
async def get_task_status(self, task_id: str) -> Dict[str, Any]:
"""
Get status of a consumption task.
Args:
task_id: Task UUID from upload
Returns:
Task status with state, result, etc.
"""
try:
response = await self.client.get(
f"{self.api_url}/tasks/",
params={"task_id": task_id},
)
response.raise_for_status()
data = response.json()
results = data.get("results", [])
if results:
return results[0]
return {"status": "NOT_FOUND"}
except Exception as e:
logger.error(f"Failed to get task status {task_id}: {e}")
raise
async def update_document(
self,
document_id: int,
title: Optional[str] = None,
correspondent: Optional[int] = None,
document_type: Optional[int] = None,
tags: Optional[List[int]] = None,
custom_fields: Optional[List[Dict[str, Any]]] = None,
) -> PaperlessDocument:
"""
Update a document's metadata.
Args:
document_id: Document ID to update
title: New title
correspondent: New correspondent ID
document_type: New document type ID
tags: New tag IDs (replaces existing)
custom_fields: New custom field values
Returns:
Updated document
"""
data = {}
if title is not None:
data["title"] = title
if correspondent is not None:
data["correspondent"] = correspondent
if document_type is not None:
data["document_type"] = document_type
if tags is not None:
data["tags"] = tags
if custom_fields is not None:
data["custom_fields"] = custom_fields
try:
response = await self.client.patch(
f"{self.api_url}/documents/{document_id}/",
json=data,
)
response.raise_for_status()
return self._parse_document(response.json())
except Exception as e:
logger.error(f"Failed to update document {document_id}: {e}")
raise
# =========================================================================
# Custom Fields
# =========================================================================
async def list_custom_fields(self) -> List[Dict[str, Any]]:
"""
List all custom fields.
Returns:
List of custom field definitions
"""
try:
response = await self.client.get(f"{self.api_url}/custom_fields/")
response.raise_for_status()
return response.json().get("results", [])
except Exception as e:
logger.error(f"Failed to list custom fields: {e}")
raise
async def get_custom_field_by_name(self, name: str) -> Optional[Dict[str, Any]]:
"""
Get a custom field by name.
Args:
name: Custom field name
Returns:
Custom field definition or None
"""
fields = await self.list_custom_fields()
for field in fields:
if field.get("name") == name:
return field
return None
# =========================================================================
# Tags, Correspondents, Document Types
# =========================================================================
async def list_tags(self) -> List[Dict[str, Any]]:
"""List all tags."""
try:
response = await self.client.get(f"{self.api_url}/tags/")
response.raise_for_status()
return response.json().get("results", [])
except Exception as e:
logger.error(f"Failed to list tags: {e}")
raise
async def list_correspondents(self) -> List[Dict[str, Any]]:
"""List all correspondents."""
try:
response = await self.client.get(f"{self.api_url}/correspondents/")
response.raise_for_status()
return response.json().get("results", [])
except Exception as e:
logger.error(f"Failed to list correspondents: {e}")
raise
async def list_document_types(self) -> List[Dict[str, Any]]:
"""List all document types."""
try:
response = await self.client.get(f"{self.api_url}/document_types/")
response.raise_for_status()
return response.json().get("results", [])
except Exception as e:
logger.error(f"Failed to list document types: {e}")
raise
# =========================================================================
# Bulk Operations
# =========================================================================
async def bulk_edit(
self,
document_ids: List[int],
method: str,
parameters: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
"""
Bulk edit documents.
Args:
document_ids: List of document IDs
method: Operation (add_tag, remove_tag, set_correspondent, etc.)
parameters: Operation parameters
Returns:
Operation result
"""
data = {
"documents": document_ids,
"method": method,
}
if parameters:
data["parameters"] = parameters
try:
response = await self.client.post(
f"{self.api_url}/documents/bulk_edit/",
json=data,
)
response.raise_for_status()
return response.json()
except Exception as e:
logger.error(f"Bulk edit failed: {e}")
raise
# =========================================================================
# Health Check
# =========================================================================
async def health_check(self) -> bool:
"""
Check if Paperless-ngx is responding.
Returns:
True if service is healthy
"""
try:
response = await self.client.get(f"{self.api_url}/", timeout=5.0)
return response.status_code < 400
except Exception as e:
logger.error(f"Paperless health check failed: {e}")
return False
# =========================================================================
# Helpers
# =========================================================================
def _parse_document(self, data: Dict[str, Any]) -> PaperlessDocument:
"""Parse API response into PaperlessDocument."""
return PaperlessDocument(
id=data.get("id", 0),
title=data.get("title", ""),
content=data.get("content", ""),
created=data.get("created"),
modified=data.get("modified"),
added=data.get("added"),
correspondent=data.get("correspondent"),
document_type=data.get("document_type"),
storage_path=data.get("storage_path"),
tags=data.get("tags", []),
archive_serial_number=data.get("archive_serial_number"),
original_file_name=data.get("original_file_name"),
archived_file_name=data.get("archived_file_name"),
custom_fields=data.get("custom_fields", []),
)
+5
View File
@@ -101,6 +101,11 @@ class Settings(BaseSettings):
content_extraction_timeout: int = Field(default=5, ge=1, le=30, description="Trafilatura per-URL timeout in seconds")
content_max_length: int = Field(default=2000, ge=500, le=10000, description="Max extracted content length per result")
# Paperless-ngx Configuration
paperless_url: str = Field(default="http://paperless:8000", description="Paperless-ngx URL")
paperless_token: str = Field(default="", description="Paperless-ngx API token")
paperless_timeout: int = Field(default=30, ge=5, le=120, description="Paperless API timeout in seconds")
# Document Store Configuration
document_store_enabled: bool = Field(default=True, description="Enable document store feature")
document_catalog_path_prefix: str = Field(default="docs", description="Wiki path prefix for catalog pages")
+53 -1
View File
@@ -22,6 +22,7 @@ from src.clients.wikijs_client import WikiJSClient
from src.clients.searxng_client import SearXNGClient
from src.clients.ollama_client import OllamaClient
from src.clients.content_extractor import ContentExtractor
from src.clients.paperless_client import PaperlessClient
logger = logging.getLogger(__name__)
@@ -155,6 +156,28 @@ def get_content_extractor() -> ContentExtractor:
return extractor
@lru_cache
def get_paperless_client() -> PaperlessClient:
"""
Get Paperless-ngx client singleton.
Returns:
Initialized Paperless-ngx REST API client
Note: Returns None-like client if paperless_token is not configured
"""
settings = get_settings()
if not settings.paperless_token:
logger.warning("Paperless token not configured - document storage disabled")
client = PaperlessClient(
base_url=settings.paperless_url,
token=settings.paperless_token,
timeout=settings.paperless_timeout
)
logger.debug(f"Created Paperless client: {settings.paperless_url}")
return client
# Type aliases for FastAPI endpoint dependencies
# Usage: def my_endpoint(neo4j: Neo4jDep):
Neo4jDep = Annotated[Neo4jClient, Depends(get_neo4j_client)]
@@ -164,6 +187,7 @@ SearXNGDep = Annotated[SearXNGClient, Depends(get_searxng_client)]
OllamaDep = Annotated[OllamaClient, Depends(get_ollama_client)]
RedisDep = Annotated[aioredis.Redis, Depends(get_redis_client)]
ContentExtractorDep = Annotated[ContentExtractor, Depends(get_content_extractor)]
PaperlessDep = Annotated[PaperlessClient, Depends(get_paperless_client)]
# Lifecycle management functions
@@ -202,6 +226,21 @@ async def startup_clients():
logger.error(f"✗ Ollama health check failed: {e}")
pass
# Check Paperless availability
settings = get_settings()
if settings.paperless_token:
try:
paperless = get_paperless_client()
is_healthy = await paperless.health_check()
if is_healthy:
logger.info(f"✓ Paperless-ngx ready: {settings.paperless_url}")
else:
logger.warning("✗ Paperless-ngx not responding")
except Exception as e:
logger.error(f"✗ Paperless health check failed: {e}")
else:
logger.info("○ Paperless-ngx not configured (document storage disabled)")
# Qdrant, Wiki.js, SearXNG are lazy-initialized
logger.info("Service clients startup complete")
@@ -230,7 +269,8 @@ async def shutdown_clients():
clients_to_close = [
("Wiki.js", get_wikijs_client()),
("SearXNG", get_searxng_client()),
("Ollama", get_ollama_client())
("Ollama", get_ollama_client()),
("Paperless", get_paperless_client()),
]
for name, client in clients_to_close:
@@ -312,6 +352,18 @@ async def check_service_health() -> dict:
logger.error(f"Ollama health check failed: {e}")
health["ollama"] = False
# Paperless-ngx
settings = get_settings()
if settings.paperless_token:
try:
paperless = get_paperless_client()
health["paperless"] = await paperless.health_check()
except Exception as e:
logger.error(f"Paperless health check failed: {e}")
health["paperless"] = False
else:
health["paperless"] = None # Not configured
return health
+2 -1
View File
@@ -51,7 +51,7 @@ app.add_middleware(
from src.routers import (
wiki, tools, graph, vector, hybrid_rag, consolidation,
ingestion, entity_linking, webhooks, rag_search, content,
maintenance, volatile
maintenance, volatile, documents
)
app.include_router(wiki.router)
@@ -67,6 +67,7 @@ app.include_router(rag_search.router)
app.include_router(content.router)
app.include_router(maintenance.router)
app.include_router(volatile.router)
app.include_router(documents.router)
# Mount static files directory for Wiki.js integration scripts
static_dir = Path(__file__).parent.parent / "static"
+207
View File
@@ -0,0 +1,207 @@
"""
Document storage models for Library Desk.
Models for Paperless-ngx document management, virus scanning,
and document sync operations.
"""
from pydantic import BaseModel, Field
from typing import Dict, Any, Optional, List
from datetime import datetime
from enum import Enum
class DocumentType(str, Enum):
"""Types of documents supported in the document store."""
PDF = "pdf"
IMAGE = "image"
VIDEO = "video"
TEXT = "text"
ARCHIVE = "archive"
OTHER = "other"
class SyncStatus(str, Enum):
"""Status of document sync with Library Desk."""
PENDING = "pending"
INDEXED = "indexed"
FAILED = "failed"
SKIPPED = "skipped"
# =============================================================================
# Document Models
# =============================================================================
class DocumentMetadata(BaseModel):
"""Metadata for a document in Paperless-ngx."""
paperless_id: int = Field(..., description="Paperless-ngx document ID")
title: str = Field(..., description="Document title")
filename: Optional[str] = Field(None, description="Original filename")
content: Optional[str] = Field(None, description="Extracted text content")
created: Optional[datetime] = Field(None, description="Document creation date")
modified: Optional[datetime] = Field(None, description="Last modification date")
added: Optional[datetime] = Field(None, description="Date added to Paperless")
correspondent: Optional[str] = Field(None, description="Correspondent name")
document_type: Optional[str] = Field(None, description="Document type name")
tags: List[str] = Field(default_factory=list, description="Tag names")
custom_fields: Dict[str, Any] = Field(default_factory=dict, description="Custom field values")
class DocumentRecord(BaseModel):
"""A document record with sync status."""
metadata: DocumentMetadata = Field(..., description="Document metadata from Paperless")
sync_status: SyncStatus = Field(default=SyncStatus.PENDING, description="Library Desk sync status")
indexed_at: Optional[datetime] = Field(None, description="When indexed in Library Desk")
collection: Optional[str] = Field(None, description="Collection name (e.g., 'fastapi-docs')")
source_url: Optional[str] = Field(None, description="Original source URL if uploaded via HybridRAG")
# =============================================================================
# Upload Request/Response Models
# =============================================================================
class DocumentUploadRequest(BaseModel):
"""Request to upload a document to Paperless-ngx."""
url: Optional[str] = Field(None, description="URL to download document from")
title: Optional[str] = Field(None, description="Document title (derived from filename if not set)")
collection: Optional[str] = Field(None, description="Collection to add document to")
tags: List[str] = Field(default_factory=list, description="Tags to apply")
correspondent: Optional[str] = Field(None, description="Correspondent name")
document_type: Optional[str] = Field(None, description="Document type name")
class DocumentUploadResponse(BaseModel):
"""Response from document upload."""
task_id: str = Field(..., description="Paperless task ID for tracking")
filename: str = Field(..., description="Uploaded filename")
message: str = Field(..., description="Status message")
# =============================================================================
# Webhook Models
# =============================================================================
class PaperlessWebhookPayload(BaseModel):
"""
Payload from Paperless-ngx webhook.
Supports Jinja template format:
- doc_url: Contains document ID in URL path (e.g., http://paperless:8000/documents/123/)
- title: Document title from {{ doc_title }}
"""
doc_url: str = Field(..., description="Paperless document URL containing ID")
title: Optional[str] = Field(None, description="Document title")
class Config:
extra = "ignore" # Ignore extra fields
@property
def document_id(self) -> int:
"""Extract document ID from doc_url."""
import re
match = re.search(r'/documents/(\d+)/?', self.doc_url)
if match:
return int(match.group(1))
raise ValueError(f"Cannot extract document ID from URL: {self.doc_url}")
class WebhookResponse(BaseModel):
"""Response to webhook processing."""
document_id: int = Field(..., description="Processed document ID")
status: str = Field(..., description="Processing status")
indexed: bool = Field(..., description="Whether document was indexed")
message: Optional[str] = Field(None, description="Additional details")
# =============================================================================
# Sync Models
# =============================================================================
class SyncRequest(BaseModel):
"""Request to sync documents from Paperless-ngx."""
since: Optional[datetime] = Field(None, description="Only sync documents modified after this time")
collection: Optional[str] = Field(None, description="Only sync documents in this collection")
limit: int = Field(default=100, ge=1, le=1000, description="Maximum documents to sync")
force_reindex: bool = Field(default=False, description="Re-index already indexed documents")
class SyncResult(BaseModel):
"""Result of a sync operation."""
documents_found: int = Field(..., description="Total documents matching criteria")
documents_indexed: int = Field(..., description="Successfully indexed")
documents_skipped: int = Field(..., description="Skipped (already indexed)")
documents_failed: int = Field(..., description="Failed to index")
errors: List[str] = Field(default_factory=list, description="Error messages")
duration_seconds: float = Field(..., description="Sync duration")
# =============================================================================
# Collection Models
# =============================================================================
class Collection(BaseModel):
"""A logical grouping of documents."""
name: str = Field(..., description="Collection name (e.g., 'fastapi-docs')")
description: Optional[str] = Field(None, description="Collection description")
document_count: int = Field(default=0, description="Number of documents")
source: Optional[str] = Field(None, description="Source (e.g., 'github.com/tiangolo/fastapi')")
last_sync: Optional[datetime] = Field(None, description="Last sync timestamp")
wiki_page: Optional[str] = Field(None, description="Wiki catalog page path")
class CollectionListResponse(BaseModel):
"""Response listing all collections."""
collections: List[Collection] = Field(..., description="List of collections")
total_documents: int = Field(..., description="Total documents across all collections")
# =============================================================================
# Search Models
# =============================================================================
class DocumentSearchRequest(BaseModel):
"""Request to search documents."""
query: str = Field(..., min_length=1, description="Search query")
collection: Optional[str] = Field(None, description="Limit to collection")
document_type: Optional[DocumentType] = Field(None, description="Filter by type")
limit: int = Field(default=10, ge=1, le=50, description="Maximum results")
include_content: bool = Field(default=False, description="Include full text content")
class DocumentSearchHit(BaseModel):
"""A document search result."""
paperless_id: int = Field(..., description="Paperless document ID")
title: str = Field(..., description="Document title")
score: float = Field(..., description="Relevance score")
highlights: Optional[str] = Field(None, description="Highlighted matching text")
collection: Optional[str] = Field(None, description="Collection name")
document_type: Optional[str] = Field(None, description="Document type")
content_preview: Optional[str] = Field(None, description="Content preview if requested")
class DocumentSearchResponse(BaseModel):
"""Response from document search."""
query: str = Field(..., description="Original query")
hits: List[DocumentSearchHit] = Field(..., description="Search results")
total: int = Field(..., description="Total matching documents")
duration_ms: int = Field(..., description="Search duration in milliseconds")
# =============================================================================
# Health Check Models
# =============================================================================
class DocumentStoreHealth(BaseModel):
"""Health status of document storage components."""
paperless_healthy: bool = Field(..., description="Paperless-ngx responding")
paperless_version: Optional[str] = Field(None, description="Paperless version")
total_documents: Optional[int] = Field(None, description="Total documents in Paperless")
indexed_documents: Optional[int] = Field(None, description="Documents indexed in Library Desk")
+424
View File
@@ -0,0 +1,424 @@
"""
Document storage router for Library Desk API.
Event-driven integration with Paperless-ngx:
- Webhook receiver triggers indexing after Paperless virus scan passes
- Upload endpoint sends files to Paperless for processing
- Search across indexed documents
"""
from fastapi import APIRouter, HTTPException, Depends, Query, UploadFile, File, Request
from typing import Optional
import logging
import time
from src.models.document import (
PaperlessWebhookPayload,
WebhookResponse,
DocumentUploadRequest,
DocumentUploadResponse,
DocumentSearchRequest,
DocumentSearchResponse,
DocumentStoreHealth,
)
from src.core.dependencies import (
verify_api_key,
PaperlessDep,
QdrantDep,
OllamaDep,
Neo4jDep,
WikiJSDep,
)
from src.core.multi_tenancy import DEFAULT_USER
from src.config import get_settings
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/documents", tags=["Documents"])
# =============================================================================
# Webhook Endpoint (primary integration - event-driven)
# =============================================================================
@router.post("/webhook", response_model=WebhookResponse)
async def receive_webhook(
payload: PaperlessWebhookPayload,
paperless: PaperlessDep,
qdrant: QdrantDep,
ollama: OllamaDep,
neo4j: Neo4jDep,
wiki: WikiJSDep,
user: str = Query(default=DEFAULT_USER, description="User identifier"),
):
"""
Receive webhook events from Paperless-ngx.
This is the primary integration point. Configure Paperless workflow:
1. Trigger: Document Added (after consumption completes)
2. Condition: Document passed virus scan (ClamAV in Paperless)
3. Action: Webhook POST to this endpoint
Library Desk indexes the document into vectors and graph.
"""
from src.services.document_sync_service import DocumentSyncService
doc_id = payload.document_id
logger.info(f"Webhook received: document_id={doc_id}, title={payload.title}")
settings = get_settings()
if not settings.document_store_enabled:
return WebhookResponse(
document_id=doc_id,
status="skipped",
indexed=False,
message="Document store is disabled"
)
try:
sync_service = DocumentSyncService(
paperless_client=paperless,
qdrant_client=qdrant,
ollama_client=ollama,
neo4j_client=neo4j,
wiki_client=wiki,
settings=settings
)
# Fetch content from Paperless (template only provides doc_url and title)
result = await sync_service.index_document(
document_id=doc_id,
user=user,
)
return WebhookResponse(
document_id=doc_id,
status="indexed" if result.success else "failed",
indexed=result.success,
message=result.error if not result.success else f"Indexed: {result.title}"
)
except Exception as e:
logger.error(f"Webhook processing failed for document {doc_id}: {e}", exc_info=True)
return WebhookResponse(
document_id=doc_id,
status="error",
indexed=False,
message=str(e)
)
# =============================================================================
# Debug Capture Endpoint
# =============================================================================
@router.post("/webhook-capture")
async def capture_webhook(request: Request):
"""Capture raw webhook payload for debugging."""
import json
from pathlib import Path
from datetime import datetime
# Get raw body
body = await request.body()
headers = dict(request.headers)
query_params = dict(request.query_params)
# Build capture data
capture = {
"timestamp": datetime.now().isoformat(),
"method": request.method,
"url": str(request.url),
"query_params": query_params,
"headers": headers,
"content_type": headers.get("content-type", "unknown"),
"body_raw": body.decode("utf-8", errors="replace"),
}
# Try to parse as JSON
try:
capture["body_json"] = json.loads(body)
except:
capture["body_json"] = None
# Write to file
capture_file = Path("logs/webhook_capture.json")
capture_file.parent.mkdir(exist_ok=True)
with open(capture_file, "w") as f:
json.dump(capture, f, indent=2, default=str)
logger.info(f"Captured webhook: {capture['body_raw'][:200]}")
return {"status": "captured", "file": str(capture_file)}
# =============================================================================
# Simple Webhook (URL parameters only)
# =============================================================================
@router.post("/webhook-simple", response_model=WebhookResponse)
async def receive_webhook_simple(
doc_url: str = Query(..., description="Paperless document URL containing ID"),
title: str = Query(default="", description="Document title"),
user: str = Query(default=DEFAULT_USER, description="User identifier"),
paperless: PaperlessDep = None,
qdrant: QdrantDep = None,
ollama: OllamaDep = None,
neo4j: Neo4jDep = None,
wiki: WikiJSDep = None,
):
"""
Simple webhook endpoint accepting URL parameters.
Used when Paperless Jinja templates don't work with JSON body.
URL format: /webhook-simple?doc_url=http://...&title=...&user=...
"""
from src.services.document_sync_service import DocumentSyncService
import re
# Extract document ID from URL
match = re.search(r'/documents/(\d+)/?', doc_url)
if not match:
return WebhookResponse(
document_id=0,
status="error",
indexed=False,
message=f"Cannot extract document ID from URL: {doc_url}"
)
doc_id = int(match.group(1))
logger.info(f"Webhook-simple received: document_id={doc_id}, title={title}")
settings = get_settings()
if not settings.document_store_enabled:
return WebhookResponse(
document_id=doc_id,
status="skipped",
indexed=False,
message="Document store is disabled"
)
try:
sync_service = DocumentSyncService(
paperless_client=paperless,
qdrant_client=qdrant,
ollama_client=ollama,
neo4j_client=neo4j,
wiki_client=wiki,
settings=settings
)
result = await sync_service.index_document(
document_id=doc_id,
user=user,
)
return WebhookResponse(
document_id=doc_id,
status="indexed" if result.success else "failed",
indexed=result.success,
message=result.error if not result.success else f"Indexed: {result.title}"
)
except Exception as e:
logger.error(f"Webhook-simple failed for document {doc_id}: {e}", exc_info=True)
return WebhookResponse(
document_id=doc_id,
status="error",
indexed=False,
message=str(e)
)
# =============================================================================
# Upload Endpoints
# =============================================================================
@router.post("/upload", response_model=DocumentUploadResponse)
async def upload_document(
file: UploadFile = File(...),
title: Optional[str] = Query(None, description="Document title"),
collection: Optional[str] = Query(None, description="Collection name"),
paperless: PaperlessDep = None,
api_key: str = Depends(verify_api_key),
):
"""
Upload a document to Paperless-ngx.
Paperless handles virus scanning. If clean, Paperless webhook
triggers indexing back to Library Desk.
"""
settings = get_settings()
if not settings.document_store_enabled:
raise HTTPException(status_code=503, detail="Document store is disabled")
content = await file.read()
filename = file.filename or "document"
custom_fields = []
if collection:
custom_fields.append({"field": "collection", "value": collection})
try:
task_id = await paperless.upload_document(
file_content=content,
filename=filename,
title=title,
custom_fields=custom_fields if custom_fields else None,
)
return DocumentUploadResponse(
task_id=task_id,
filename=filename,
message=f"Uploaded to Paperless, task {task_id}. Indexing via webhook after scan."
)
except Exception as e:
logger.error(f"Upload failed for '{filename}': {e}")
raise HTTPException(status_code=500, detail=f"Upload failed: {e}")
@router.post("/upload-url", response_model=DocumentUploadResponse)
async def upload_from_url(
request: DocumentUploadRequest,
paperless: PaperlessDep = None,
api_key: str = Depends(verify_api_key),
):
"""
Download document from URL and upload to Paperless-ngx.
Used by HybridRAG to save discovered PDFs. Paperless scans and
webhooks back for indexing.
"""
import httpx
settings = get_settings()
if not settings.document_store_enabled:
raise HTTPException(status_code=503, detail="Document store is disabled")
if not request.url:
raise HTTPException(status_code=400, detail="URL is required")
try:
async with httpx.AsyncClient(timeout=60.0) as client:
response = await client.get(request.url, follow_redirects=True)
response.raise_for_status()
content = response.content
filename = request.url.split("/")[-1].split("?")[0] or "document"
except Exception as e:
logger.error(f"Download failed from {request.url}: {e}")
raise HTTPException(status_code=400, detail=f"Download failed: {e}")
try:
custom_fields = [{"field": "source_url", "value": request.url}]
if request.collection:
custom_fields.append({"field": "collection", "value": request.collection})
task_id = await paperless.upload_document(
file_content=content,
filename=filename,
title=request.title,
custom_fields=custom_fields,
)
return DocumentUploadResponse(
task_id=task_id,
filename=filename,
message=f"Uploaded from URL, task {task_id}. Indexing via webhook after scan."
)
except Exception as e:
logger.error(f"Upload failed for URL '{request.url}': {e}")
raise HTTPException(status_code=500, detail=f"Upload failed: {e}")
# =============================================================================
# Search
# =============================================================================
@router.post("/search", response_model=DocumentSearchResponse)
async def search_documents(
request: DocumentSearchRequest,
qdrant: QdrantDep,
ollama: OllamaDep,
user: str = Query(default=DEFAULT_USER, description="User identifier"),
api_key: str = Depends(verify_api_key),
):
"""
Semantic search across indexed documents.
"""
from src.services.vector_service import VectorService
from src.core.dependencies import get_wikijs_client
from src.models.document import DocumentSearchHit
settings = get_settings()
if not settings.document_store_enabled:
raise HTTPException(status_code=503, detail="Document store is disabled")
start_time = time.time()
try:
wiki = get_wikijs_client()
vector_service = VectorService(qdrant, wiki, ollama)
results = await vector_service.search(
query=request.query,
user=user,
limit=request.limit,
score_threshold=0.5,
doc_type="document"
)
hits = []
for result in results.get("results", []):
hits.append(DocumentSearchHit(
paperless_id=result.get("metadata", {}).get("paperless_id", 0),
title=result.get("title", ""),
score=result.get("score", 0.0),
highlights=result.get("chunk_text", "")[:200] if request.include_content else None,
collection=result.get("metadata", {}).get("collection"),
document_type=result.get("metadata", {}).get("document_type"),
content_preview=result.get("chunk_text", "")[:500] if request.include_content else None,
))
return DocumentSearchResponse(
query=request.query,
hits=hits,
total=len(hits),
duration_ms=int((time.time() - start_time) * 1000)
)
except Exception as e:
logger.error(f"Document search failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
# =============================================================================
# Health
# =============================================================================
@router.get("/health", response_model=DocumentStoreHealth)
async def document_store_health(paperless: PaperlessDep):
"""Check Paperless-ngx connectivity."""
settings = get_settings()
paperless_healthy = False
if settings.paperless_token:
try:
paperless_healthy = await paperless.health_check()
except Exception as e:
logger.error(f"Paperless health check failed: {e}")
return DocumentStoreHealth(
paperless_healthy=paperless_healthy,
paperless_version="connected" if paperless_healthy else None,
total_documents=None,
indexed_documents=None
)
+232 -1
View File
@@ -19,9 +19,10 @@ from src.services.graph_service import GraphService
from src.services.volatile_service import VolatileCacheService
from src.core.dependencies import (
VectorServiceDep, GraphServiceDep, WikiJSDep, RedisDep,
QdrantDep, OllamaDep, verify_api_key
QdrantDep, OllamaDep, PaperlessDep, verify_api_key
)
from src.config import get_settings
from src.core.multi_tenancy import DEFAULT_USER
from datetime import datetime, timezone
logger = logging.getLogger(__name__)
@@ -220,6 +221,47 @@ class VolatileCleanupResponse(BaseModel):
duration_ms: float
class TestDataCleanupResponse(BaseModel):
"""Response from test data cleanup operation."""
success: bool
dry_run: bool
wiki_pages_deleted: int
graph_nodes_deleted: int
vector_chunks_deleted: int
pages_found: List[Dict[str, Any]] = Field(default_factory=list)
duration_ms: float
class PaperlessCleanupResponse(BaseModel):
"""Response from Paperless orphan cleanup operation."""
success: bool
dry_run: bool
paperless_ids_checked: int = Field(description="Total Paperless IDs found in indexes")
orphans_found: int = Field(description="Documents deleted from Paperless but still indexed")
orphan_ids: List[int] = Field(default_factory=list, description="Paperless IDs that are orphans")
vector_chunks_deleted: int = Field(description="Vector chunks removed")
graph_nodes_deleted: int = Field(description="Graph Document nodes removed")
duration_ms: float
# Test data path patterns - restricted to test user namespace only
# These are the only paths that can be cleaned up for safety
TEST_USER_PATH_PREFIXES = [
"users/llm-tester/",
"users/llm_tester/",
]
def _matches_test_user_path(path: str) -> bool:
"""Check if a path is in the test user namespace.
Only matches paths that START with test user prefixes for safety.
This prevents accidental deletion of non-test data.
"""
path_lower = path.lower()
return any(path_lower.startswith(prefix) for prefix in TEST_USER_PATH_PREFIXES)
# ========== Endpoints ==========
@router.post("/cleanup/vectors", response_model=VectorCleanupResponse)
@@ -556,6 +598,195 @@ async def cleanup_volatile(
raise HTTPException(status_code=500, detail=str(e))
@router.post("/cleanup/test-data", response_model=TestDataCleanupResponse)
async def cleanup_test_data(
dry_run: bool = Query(default=True, description="Preview only, don't delete"),
wiki: WikiJSDep = None,
vector_service: VectorServiceDep = None,
graph_service: GraphServiceDep = None,
api_key: str = Depends(verify_api_key)
):
"""
Purge LLM tester data from wiki, graph, and vectors.
**Security**: Only deletes pages in the test user namespace:
- users/llm-tester/*
- users/llm_tester/*
This endpoint cannot delete data outside these paths.
**Use dry_run=true (default) to preview what would be deleted.**
**Scheduler Integration:**
```json
{
"task_name": "test_data_cleanup",
"schedule": "0 3 * * 0",
"endpoint": "POST /maintenance/cleanup/test-data?dry_run=false",
"description": "Weekly cleanup of LLM test data"
}
```
"""
start_time = time.time()
try:
# List all wiki pages
all_pages = await wiki.list_all_pages(batch_size=500)
# Filter for test user paths only (security: restricted to test namespace)
test_pages = [
{"id": p["id"], "path": p["path"], "title": p.get("title", "")}
for p in all_pages
if _matches_test_user_path(p.get("path", ""))
]
logger.info(f"Found {len(test_pages)} test pages matching patterns: {TEST_USER_PATH_PREFIXES}")
wiki_deleted = 0
graph_deleted = 0
vector_deleted = 0
if not dry_run and test_pages:
for page in test_pages:
page_id = page["id"]
page_path = page["path"]
try:
# Delete vector chunks for this page (using DEFAULT_USER collection)
chunks_removed = await vector_service.delete_page_chunks(page_id, DEFAULT_USER)
vector_deleted += chunks_removed
# Delete graph node for this page (returns count, may be 0 if no node)
graph_removed = await graph_service.delete_page(page_id, DEFAULT_USER)
graph_deleted += graph_removed
# Delete wiki page (raises exception on failure, returns None on success)
await wiki.delete_page(page_id)
wiki_deleted += 1
logger.info(f"Deleted test page: {page_path} (id={page_id})")
except Exception as e:
logger.error(f"Failed to delete page {page_path}: {e}")
continue
duration_ms = (time.time() - start_time) * 1000
return TestDataCleanupResponse(
success=True,
dry_run=dry_run,
wiki_pages_deleted=wiki_deleted,
graph_nodes_deleted=graph_deleted,
vector_chunks_deleted=vector_deleted,
pages_found=test_pages,
duration_ms=duration_ms
)
except Exception as e:
logger.error(f"Test data cleanup failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@router.post("/cleanup/paperless", response_model=PaperlessCleanupResponse)
async def cleanup_paperless_orphans(
user: str = Query(..., description="User identifier"),
dry_run: bool = Query(default=True, description="Preview only, don't delete"),
vector_service: VectorServiceDep = None,
graph_service: GraphServiceDep = None,
paperless: PaperlessDep = None,
api_key: str = Depends(verify_api_key)
):
"""
Find and clean up Paperless document orphans.
Detects documents that were indexed in Library Desk but have since been
deleted from Paperless-ngx. Removes orphaned vectors and graph nodes.
**Use dry_run=true (default) to preview what would be deleted.**
**Scheduler Integration:**
```json
{
"task_name": "paperless_orphan_cleanup",
"schedule": "0 5 * * *",
"endpoint": "POST /maintenance/cleanup/paperless?user=jpmschweitzer&dry_run=false",
"description": "Daily cleanup of orphaned Paperless documents"
}
```
"""
start_time = time.time()
try:
settings = get_settings()
if not settings.paperless_token:
raise HTTPException(status_code=503, detail="Paperless not configured")
# Get all document chunks from vectors with doc_type="document"
chunk_refs = await vector_service.get_all_chunk_references(user)
doc_chunks = [ref for ref in chunk_refs if ref.get("doc_type") == "document"]
# Extract unique paperless_ids
paperless_ids = list(set(
ref.get("paperless_id") for ref in doc_chunks
if ref.get("paperless_id")
))
logger.info(f"Found {len(paperless_ids)} unique Paperless IDs in indexes")
# Check each against Paperless API
orphan_ids = []
for pid in paperless_ids:
try:
doc = await paperless.get_document(pid)
if doc is None:
orphan_ids.append(pid)
except Exception as e:
# Document not found or API error - treat as orphan
logger.debug(f"Paperless document {pid} not found: {e}")
orphan_ids.append(pid)
logger.info(f"Found {len(orphan_ids)} orphaned Paperless documents")
# Delete orphans if not dry run
vectors_deleted = 0
graph_deleted = 0
if not dry_run and orphan_ids:
for pid in orphan_ids:
try:
# Delete vector chunks for this paperless_id
chunks_removed = await vector_service.delete_paperless_document_chunks(pid, user)
vectors_deleted += chunks_removed
# Delete graph node for this paperless_id
graph_removed = await graph_service.delete_paperless_document(pid, user)
graph_deleted += graph_removed
logger.info(f"Cleaned up orphaned Paperless document {pid}: {chunks_removed} chunks, {graph_removed} nodes")
except Exception as e:
logger.error(f"Failed to cleanup Paperless document {pid}: {e}")
duration_ms = (time.time() - start_time) * 1000
return PaperlessCleanupResponse(
success=True,
dry_run=dry_run,
paperless_ids_checked=len(paperless_ids),
orphans_found=len(orphan_ids),
orphan_ids=orphan_ids,
vector_chunks_deleted=vectors_deleted,
graph_nodes_deleted=graph_deleted,
duration_ms=duration_ms
)
except HTTPException:
raise
except Exception as e:
logger.error(f"Paperless orphan cleanup failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@router.get("/health", response_model=HealthCheckResponse)
async def maintenance_health(
user: str = Query(..., description="User identifier"),
+293
View File
@@ -0,0 +1,293 @@
"""
Document sync service for Library Desk.
Handles indexing of Paperless-ngx documents into vectors and graph.
Called by webhook when Paperless completes document processing.
"""
import logging
import re
import hashlib
import uuid
from typing import Optional, List
from dataclasses import dataclass
from src.clients.paperless_client import PaperlessClient
from src.clients.qdrant_client import QdrantClientWrapper
from src.clients.ollama_client import OllamaClient
from src.clients.neo4j_client import Neo4jClient
from src.clients.wikijs_client import WikiJSClient
from src.core.multi_tenancy import get_qdrant_collection_name
from src.config import Settings
logger = logging.getLogger(__name__)
@dataclass
class IndexResult:
"""Result of indexing a single document."""
success: bool
document_id: int
title: str = ""
chunks_created: int = 0
error: Optional[str] = None
class DocumentSyncService:
"""
Service for syncing Paperless documents to Library Desk indexes.
Handles:
- Fetching document content from Paperless API
- Chunking and embedding into Qdrant
- Creating graph nodes in Neo4j
"""
def __init__(
self,
paperless_client: PaperlessClient,
qdrant_client: QdrantClientWrapper,
ollama_client: OllamaClient,
neo4j_client: Neo4jClient,
wiki_client: WikiJSClient,
settings: Settings,
chunk_size: int = 500,
chunk_overlap: int = 50
):
self.paperless = paperless_client
self.qdrant = qdrant_client
self.ollama = ollama_client
self.neo4j = neo4j_client
self.wiki = wiki_client
self.settings = settings
self.chunk_size = chunk_size
self.chunk_overlap = chunk_overlap
def _chunk_text(self, text: str) -> List[str]:
"""Chunk text into overlapping segments."""
text = re.sub(r'\s+', ' ', text).strip()
words = text.split()
if len(words) <= self.chunk_size:
return [text] if text else []
chunks = []
start = 0
while start < len(words):
end = start + self.chunk_size
chunk_words = words[start:end]
chunks.append(' '.join(chunk_words))
start = end - self.chunk_overlap
return chunks
async def index_document(
self,
document_id: int,
user: str,
content: Optional[str] = None,
title: Optional[str] = None,
) -> IndexResult:
"""
Index a single document from Paperless into vectors and graph.
Args:
document_id: Paperless document ID
user: User identifier for multi-tenancy
content: Optional document content (if provided, skip Paperless API call)
title: Optional document title (if provided, skip Paperless API call)
Returns:
IndexResult with success status and details
"""
logger.info(f"Indexing document {document_id} for user {user}")
try:
# If content and title provided (from webhook), skip API call
if content is not None and title is not None:
doc_title = title
doc_content = content
original_filename = None
correspondent = None
document_type = None
tags = []
else:
# Fetch document from Paperless
doc = await self.paperless.get_document(document_id)
if not doc:
return IndexResult(
success=False,
document_id=document_id,
error="Document not found in Paperless"
)
doc_title = doc.title
doc_content = doc.content or ""
original_filename = doc.original_file_name
correspondent = doc.correspondent
document_type = doc.document_type
tags = doc.tags
if not doc_content.strip():
logger.warning(f"Document {document_id} has no text content")
return IndexResult(
success=True,
document_id=document_id,
title=doc_title,
chunks_created=0,
error="No text content (possibly image/video only)"
)
# Index vectors
chunks_created = await self._index_vectors(
document_id=document_id,
title=doc_title,
content=doc_content,
user=user,
metadata={
"paperless_id": document_id,
"original_filename": original_filename,
"correspondent": correspondent,
"document_type": document_type,
"tags": tags,
}
)
# Index graph node
await self._index_graph(
document_id=document_id,
title=doc_title,
content=doc_content,
user=user,
)
# Mark as indexed in Paperless (optional - if custom field exists)
try:
await self._mark_indexed(document_id)
except Exception as e:
logger.debug(f"Could not mark document as indexed: {e}")
logger.info(f"Successfully indexed document {document_id}: {chunks_created} chunks")
return IndexResult(
success=True,
document_id=document_id,
title=doc_title,
chunks_created=chunks_created
)
except Exception as e:
logger.error(f"Failed to index document {document_id}: {e}", exc_info=True)
return IndexResult(
success=False,
document_id=document_id,
error=str(e)
)
async def _index_vectors(
self,
document_id: int,
title: str,
content: str,
user: str,
metadata: dict,
) -> int:
"""Create vector embeddings for document content."""
collection = get_qdrant_collection_name(user)
self.qdrant.ensure_collection(collection)
# Delete existing chunks for this document
try:
self.qdrant.client.delete(
collection_name=collection,
points_selector={
"filter": {
"must": [
{"key": "doc_type", "match": {"value": "document"}},
{"key": "paperless_id", "match": {"value": document_id}},
]
}
}
)
except Exception as e:
logger.debug(f"No existing chunks to delete: {e}")
# Chunk content
chunks = self._chunk_text(content)
if not chunks:
return 0
# Generate embeddings
embeddings = await self.ollama.embed_batch(chunks)
# Build points
points = []
for i, (chunk, embedding) in enumerate(zip(chunks, embeddings)):
point_id = str(uuid.uuid4())
content_hash = hashlib.md5(chunk.encode()).hexdigest()
points.append({
"id": point_id,
"vector": embedding,
"payload": {
"doc_type": "document",
"paperless_id": document_id,
"title": title,
"chunk_text": chunk,
"chunk_index": i,
"content_hash": content_hash,
**metadata
}
})
# Upsert to Qdrant
if points:
self.qdrant.client.upsert(
collection_name=collection,
points=points
)
return len(points)
async def _index_graph(
self,
document_id: int,
title: str,
content: str,
user: str,
):
"""Create graph node for document."""
# Create Document node in Neo4j
query = """
MERGE (d:Document {paperless_id: $paperless_id, user: $user})
SET d.title = $title,
d.doc_type = 'document',
d.updated_at = datetime()
RETURN d
"""
await self.neo4j.execute_query(
query,
{
"paperless_id": document_id,
"user": user,
"title": title,
}
)
# TODO: Extract entities from content and create relationships
# This could use the same entity extraction as wiki pages
async def _mark_indexed(self, document_id: int):
"""Mark document as indexed in Paperless custom field."""
# Try to update library_indexed custom field if it exists
try:
# Look up field ID by name (Paperless requires ID, not name)
field = await self.paperless.get_custom_field_by_name("library_indexed")
if field:
await self.paperless.update_document(
document_id=document_id,
custom_fields=[{"field": field["id"], "value": True}]
)
except Exception:
# Field might not exist, that's OK
pass
+42
View File
@@ -1306,6 +1306,48 @@ Feel free to expand it with more details!
logger.error(f"Failed to delete document {document_id} from graph: {e}", exc_info=True)
return 0
async def delete_paperless_document(
self,
paperless_id: int,
user: str
) -> int:
"""
Delete a Paperless document node and all its relationships.
Args:
paperless_id: Paperless-ngx document ID
user: User identifier
Returns:
Number of nodes deleted (1 if successful, 0 if not found)
"""
user_doc_label = get_neo4j_user_label(user)
delete_query = f"""
MATCH (d:{user_doc_label}:Document {{paperless_id: $paperless_id}})
DETACH DELETE d
RETURN count(d) as deleted_count
"""
try:
result = await self.neo4j.execute_query(
delete_query,
{"paperless_id": paperless_id}
)
deleted_count = result[0]["deleted_count"] if result else 0
if deleted_count > 0:
logger.info(f"Deleted Document node for Paperless document {paperless_id}")
else:
logger.debug(f"No Document node found for Paperless document {paperless_id}")
return deleted_count
except Exception as e:
logger.error(f"Failed to delete Paperless document {paperless_id} from graph: {e}", exc_info=True)
return 0
async def delete_collection_node(
self,
collection_id: str,
+33
View File
@@ -389,6 +389,39 @@ class VectorService:
logger.error(f"Failed to delete chunks for document {document_id}: {e}", exc_info=True)
return 0
async def delete_paperless_document_chunks(
self,
paperless_id: int,
user: str
) -> int:
"""
Delete all chunks for a Paperless document.
Args:
paperless_id: Paperless-ngx document ID
user: User identifier
Returns:
Number of chunks deleted
"""
collection_name = get_qdrant_collection_name(user)
try:
deleted_count = await self.qdrant.delete_by_filter(
collection_name=collection_name,
filter_conditions={
"doc_type": "document",
"paperless_id": paperless_id
}
)
logger.info(f"Deleted chunks for Paperless document {paperless_id}")
return deleted_count
except Exception as e:
logger.error(f"Failed to delete chunks for Paperless document {paperless_id}: {e}", exc_info=True)
return 0
async def delete_collection_chunks(
self,
collection_id: str,