Files
leocrm/ARCHITECTURE_PLAN.md
T

2433 lines
99 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# LeoCRM — Architektur-Plan für Phasen BK
> **Erstellt:** 2026-08-20
> **Basis:** Code-Analyse aller vorhandenen Systeme (kommunikation, graph_rag, unified_search, automation, workflows, ai, approval, storage, outbox, wiki)
> **Leitlinie:** Alles aufbauend auf vorhandenem Code. Keine parallelen Systeme.
> **Letzte Migration:** 0128_ai_decision_records.py → nächste: 0129+
---
## Inhaltsverzeichnis
1. [Phase B Lücken (3 Tasks)](#phase-b-lücken)
2. [Phase F Lücken (3 Tasks)](#phase-f-lücken)
3. [Phase G Lücke (1 Task)](#phase-g-lücke)
4. [Phase H Rest (12 Tasks)](#phase-h-rest)
5. [Phase I (25 Tasks)](#phase-i)
6. [Phase J (10 Tasks)](#phase-j)
7. [Phase K (6 Tasks)](#phase-k)
---
## Vorhandene Systeme (Basis für alle Phasen)
### Kommunikation Plugin
- **Pfad:** `app/plugins/builtins/kommunikation/`
- **Models:** `CommConversation` (is_system, is_direct, is_archived, metadata_), `CommParticipant` (participant_type: user/ai/system/gateway), `CommMessage` (sender_type, content, content_format, metadata_, reply_to_id), `CommMessageBlock` (block_type, block_data JSONB, sort_order), `CommMessageAttachment`, `CommMessageReaction`, `CommMessageRead`, `CommConversationMute`, `CommConversationPin`
- **Services:** `send_message()`, `get_conversation()`, `get_messages()`, `parse_mentions()`, `create_plugin_room()`, `post_system_message()`, `get_or_create_system_channel()`
- **Contracts:** `KommunikationContract` registriert in `ContractRegistry` unter `"kommunikation"`
- **MiniAppRegistry:** `register()`, `unregister()`, `unregister_plugin()`, `list_apps()`, `get_app()` — Singleton via `get_miniapp_registry()`
- **ContentTypes:** `BLOCK_TYPES` dict: text, markdown, html, image, audio, video, file, action_card, contact_card, miniapp
- **WebSocket:** `WebSocketManager` für Real-time
- **ParticipantRegistry:** `get_participant_registry()` mit `ParticipantHandler`
### GraphRAG Plugin
- **Pfad:** `app/plugins/builtins/graph_rag/`
- **Models:** `EntityRelationship` (source_type, source_id, target_type, target_id, relationship_type, meta JSONB)
- **Services:** `create_relationship()`, `traverse_graph()` (BFS, max_hops)
- **Provider:** `GraphRAGSearchProvider` registriert in `SearchProviderRegistry`
- **Routes:** `/api/v1/graph`
- **Dependencies:** `unified_search`
### Unified Search Plugin
- **Pfad:** `app/plugins/builtins/unified_search/`
- **BaseSearchProvider:** `search_fts()`, `search_vector()` mit visibility filtering, `supports_fts/vector/rag/graph` flags
- **SearchProviderRegistry:** `register()`, `unregister()`, `get()`, `get_all()` — Singleton via `get_search_registry()`
- **Providers:** contact, company, task, tag, mail, ai_chat, etc. (14+)
- **Embedding:** `llm_embed()` integration, `EMBEDDING_DIMENSIONS = 768`
- **Contracts:** `get_search_registry()` in `unified_search.contracts`
### Automation Plugin
- **Pfad:** `app/plugins/builtins/automation/`
- **Models:** `AgentDefinition`, `AutomationDefinition`
- **Agent Runner:** `app/plugins/builtins/automation/agent_runner.py` — ruft `run_react_loop`, `build_agent_context`, `resolve_agent_permissions`, `get_agent_tools`, `enforce_data_policy`, `mark_as_ai_generated`, `create_decision_record`
- **Agent Loop:** `app/ai/agent_loop.py` — ReAct-Loop mit Tool-Execution, Approval-Pause
- **Agent Comm:** `agent_comm.py``register_agent_comm_tool()` registriert Komm-Tool für Agents
- **Prebuilt:** `prebuilt/email_triage_agent.py`, `prebuilt/follow_up_agent.py`, `prebuilt/contact_enrichment_agent.py`, `prebuilt/report_agent.py`**NICHT in on_activate registriert**
- **on_activate:** Registriert agent_comm_tool, agent_coordinator_tools, MiniApps, cron jobs — **aber KEINE prebuilt agents**
### Workflow Engine
- **Pfad:** `app/workflows/engine.py`
- **Models:** `Workflow` (steps JSONB, trigger_event), `WorkflowInstance` (status, current_step_index, context, resume_at, step_state, lock_owner, retry_count), `WorkflowStepHistory`
- **Step Types:** action, approval, notification, condition, wait, http, mail, calendar, dms, search, agent, crm, event, webhook
- **Decision Guard:** `app/workflows/decision_guard.py` — erstellt ApprovalRequest bei High-Risk
- **Step Handlers:** `app/workflows/step_handlers.py``get_step_handler()`, `StepResult`
### Approval System
- **Pfad:** `app/core/approval.py`
- **Model:** `ApprovalRequest` (entity_type, entity_id, action, requested_by, requested_by_type, status: pending→approved/rejected/expired, request_metadata JSONB)
- **Functions:** `create_approval_request()`, resolve functions
- **Routes:** `app/routes/approvals.py`
### LLM Client
- **Pfad:** `app/ai/llm_client.py`
- **Functions:** `llm_complete()` (chat completion mit retry, cost tracking), `llm_embed()` (text embedding, 768 dims)
- **Provider:** LiteLLM (100+ providers), OpenRouter für embeddings
### Outbox + Event Bus
- **Outbox:** `app/core/outbox.py``enqueue_outbox_event()`, DLQ, replay, monitoring
- **Event Bus:** `app/core/event_bus.py``get_event_bus()`
- **Hooks:** `app/core/hooks.py``do_action()`, `apply_filters()`, `register_action()`, `register_filter()`
### Storage
- **Pfad:** `app/core/storage.py`
- **Classes:** `StorageBackend` (ABC), `LocalStorage`, `S3Storage`
- **Methods:** `save()`, `save_stream()`, `read()`, `delete()`
- **Config:** `STORAGE_BACKEND` env (local/s3), `STORAGE_PATH`, S3 settings
### Wiki Plugin
- **Pfad:** `app/plugins/builtins/wiki/`
- **Models:** WikiArticle, WikiCategory, WikiArticleVersion (versioning with restore)
- **Routes:** `/api/v1/wiki`
- **Plugin:** `WikiPlugin`**kein on_activate** (kein search provider, keine AI tools)
### Notifications (zu deprecieren)
- **Model:** `Notification` (app/models/notification.py) — noch vorhanden
- **Routes:** `app/routes/notifications.py` — noch in main.py aktiv (line 545)
- **Service:** `app/core/notifications.py``post_system_message()` delegiert an kommunikation, `create_notification()` ist deprecated wrapper
- **Frontend:** `NotificationDropdown.tsx`, `NotificationItem.tsx` — noch vorhanden
---
## Phase B Lücken
### B-VEC-IVF: IVFFlat Index-Strategie implementieren
**Basis:** `app/core/storage.py` (nein), `app/plugins/builtins/unified_search/embedding.py` + pgvector
**Was existiert:**
- HNSW ist konfiguriert (`settings.hnsw_ef_search` in `base_provider.py`)
- `test_vector_performance.py` existiert
- Keine IVFFlat-Konfiguration im Code
**Was neu gebaut wird:**
1. **Config-Erweiterung:** `app/config.py`
```python
# Neue Settings
vector_index_type: str = "hnsw" # "hnsw" or "ivfflat"
ivfflat_lists: int = 100 # number of lists for IVFFlat
ivfflat_probes: int = 10 # number of probes at query time
```
2. **Index-Manager:** `app/plugins/builtins/unified_search/index_manager.py` (NEU, ~200 Zeilen)
```python
async def create_ivfflat_index(db, table: str, column: str, lists: int):
"""Create IVFFlat index on vector column."""
await db.execute(text(
f"CREATE INDEX IF NOT EXISTS idx_{table}_{column}_ivfflat "
f"ON {table} USING ivfflat ({column} vector_cosine_ops) WITH (lists = {lists})"
))
async def set_ivfflat_probes(db, probes: int):
await db.execute(text(f"SET LOCAL ivfflat.probes = {probes}"))
async def switch_index_strategy(db, table: str, column: str, strategy: str):
"""Switch between HNSW and IVFFlat."""
# Drop old, create new
```
3. **Migration:** `alembic/versions/0129_ivfflat_index_strategy.py`
- Fügt `vector_index_type` zu system_settings hinzu
- Erstellt IVFFlat-Index alternativ zu HNSW (nicht beide gleichzeitig)
- Downgrade: Drop IVFFlat, restore HNSW
4. **BaseSearchProvider Anpassung:** `app/plugins/builtins/unified_search/base_provider.py`
- In `search_vector()`: Wenn `settings.vector_index_type == "ivfflat"`, set `ivfflat.probes` statt `hnsw.ef_search`
5. **Admin API:** `app/plugins/builtins/unified_search/routes.py`
- `POST /api/v1/search/index/switch` — Switch index strategy (admin only)
- `GET /api/v1/search/index/status` — Current index info
**Verbindungen:**
- `base_provider.py` importiert `index_manager` für probe/ef_search setting
- `config.py` erweitert mit vector_index_type
**Migrationen:**
- `0129_ivfflat_index_strategy.py`
**Tests:**
- `tests/test_ivfflat_index.py` — IVFFlat index creation, query with probes, performance comparison
- `tests/test_vector_performance.py` — erweitert um IVFFlat benchmarks
**Frontend:**
- Settings-Seite: Toggle HNSW ↔ IVFFlat in Admin-Settings
---
### B-STOR-EXT: External Storage Plugin System (WebDAV, Nextcloud)
**Basis:** `app/core/storage.py` (StorageBackend ABC, LocalStorage, S3Storage)
**Was existiert:**
- `StorageBackend` ABC mit `save()`, `save_stream()`, `read()`, `delete()`
- `LocalStorage`, `S3Storage` implementiert
- Kein Plugin-Interface für externe Storage-Provider
**Was neu gebaut wird:**
1. **Storage Provider Registry:** `app/core/storage_registry.py` (NEU, ~150 Zeilen)
```python
class StorageProviderRegistry:
"""Registry for pluggable storage backends."""
def register(self, name: str, backend_class: type[StorageBackend]) -> None: ...
def unregister(self, name: str) -> None: ...
def get(self, name: str) -> type[StorageBackend] | None: ...
def list_providers(self) -> list[str]: ...
_registry: StorageProviderRegistry | None = None
def get_storage_registry() -> StorageProviderRegistry: ...
```
2. **WebDAV Storage Backend:** `app/plugins/builtins/storage_webdav/` (NEU, komplettes Plugin, ~400 Zeilen)
- `plugin.py` — `WebDAVStoragePlugin(BasePlugin)` mit Manifest
- `backend.py` — `WebDAVStorage(StorageBackend)` implementiert save/read/delete via HTTP (PUT/GET/DELETE)
- `routes.py` — `POST /api/v1/storage/webdav/test` — Test connection
- `schemas.py` — WebDAVConfig (url, username, password, base_path)
- `migrations/0001_initial.sql` — storage_provider_configs table
3. **Nextcloud Storage Backend:** `app/plugins/builtins/storage_nextcloud/` (NEU, ~300 Zeilen)
- Baut auf WebDAV auf (Nextcloud hat WebDAV-API)
- `plugin.py` — `NextcloudStoragePlugin(BasePlugin)`
- `backend.py` — `NextcloudStorage(WebDAVStorage)` mit Nextcloud-spezifischen Erweiterungen (sharing, OCS API)
- `routes.py` — `POST /api/v1/storage/nextcloud/test`, `GET /api/v1/storage/nextcloud/shares`
4. **Storage Config Model:** `app/models/storage_config.py` (NEU)
```python
class StorageProviderConfig(Base, TenantMixin):
__tablename__ = "storage_provider_configs"
id: UUID PK
provider_name: str # "local", "s3", "webdav", "nextcloud"
config: dict[str, Any] = JSONB # provider-specific config (encrypted secrets)
is_active: bool = True
priority: int = 0 # for fallback ordering
created_at: TIMESTAMPTZ
```
5. **Storage Factory:** `app/core/storage.py` erweitern
```python
async def get_storage_backend(db, tenant_id) -> StorageBackend:
"""Get configured storage backend for tenant."""
# Query StorageProviderConfig, instantiate via registry
```
6. **Admin API:** `app/routes/storage.py` (NEU)
- `GET /api/v1/storage/providers` — List available providers
- `POST /api/v1/storage/providers` — Configure provider for tenant
- `PUT /api/v1/storage/providers/{id}` — Update config
- `DELETE /api/v1/storage/providers/{id}` — Remove provider
- `POST /api/v1/storage/providers/{id}/test` — Test connection
**Verbindungen:**
- `storage.py` importiert `storage_registry` für dynamische Provider-Auflösung
- DMS/Mail Plugins nutzen `get_storage_backend()` statt direkter LocalStorage/S3Storage
- WebDAV/Nextcloud Plugins registrieren sich in `StorageProviderRegistry` via `on_activate`
**Migrationen:**
- `0130_storage_provider_configs.py` — storage_provider_configs table
- Plugin-Migrations: `storage_webdav/migrations/0001_initial.sql`, `storage_nextcloud/migrations/0001_initial.sql`
**Tests:**
- `tests/test_storage_webdav.py` — WebDAV save/read/delete mit Mock-HTTP-Server
- `tests/test_storage_nextcloud.py` — Nextcloud-spezifische Tests
- `tests/test_storage_registry.py` — Registry register/unregister/get
- `tests/test_storage_provider_config.py` — CRUD API tests
**Frontend:**
- `frontend/src/pages/StorageSettings.tsx` — Provider-Konfiguration
- `frontend/src/components/storage/ProviderConfigForm.tsx` — Config form per provider type
- `frontend/src/components/storage/ConnectionTest.tsx` — Test connection button
---
### B-NOTIF-DEPREC: Notification System zurückbauen
**Basis:** `app/core/notifications.py`, `app/models/notification.py`, `app/routes/notifications.py`
**Was existiert:**
- `Notification` Model in `app/models/notification.py`
- `NotificationType` Model
- `/api/v1/notifications` Routes in `app/routes/notifications.py` (aktiv in main.py:545)
- `NotificationDropdown.tsx`, `NotificationItem.tsx` im Frontend
- `post_system_message()` in `app/core/notifications.py` delegiert bereits an kommunikation
- `create_notification()` ist deprecated wrapper
- `NotificationPreference` Model existiert (wird behalten für preferences)
**Was gemacht wird:**
1. **Routes stilllegen:** `app/routes/notifications.py`
- Alle Endpoints markieren als `@deprecated` in OpenAPI
- `GET /api/v1/notifications` → redirect zu `GET /api/v1/comm/system-channel/messages`
- `POST /api/v1/notifications/read` → redirect zu comm equivalent
- `GET /api/v1/notifications/unread-count` → redirect zu comm equivalent
- Nach 1 Release: Routes aus main.py entfernen (line 545)
2. **Model als deprecated markieren:** `app/models/notification.py`
- `Notification` und `NotificationType` erhalten Docstring `DEPRECATED — use CommConversation is_system=True`
- Keine neuen Writes, nur noch Reads für Migration
3. **Frontend entfernen:**
- `NotificationDropdown.tsx` → ersetzen durch `CommSystemChannel.tsx` (neu, liest aus comm system channel)
- `NotificationItem.tsx` → ersetzen durch `CommMessageItem.tsx` (existiert bereits in kommunikation Frontend)
- Header-Komponente anpassen: Notification-Bell → Comm-Message-Bell
4. **Daten-Migration:** `alembic/versions/0131_migrate_notifications_to_comm.py`
- Liest alle `Notification` records
- Erstellt entsprechende `CommMessage` im system channel
- Markiert Notifications als migrated (neue Spalte `migrated_at`)
- Nach erfolgreichem Migration-Run: Drop `notifications` und `notification_types` tables
5. **Service cleanup:** `app/core/notifications.py`
- `create_notification()` entfernen (deprecated wrapper)
- `post_system_message()` behalten (delegiert an kommunikation)
- Direkte Notification-Queries entfernen
6. **main.py cleanup:**
- `notifications` Router entfernen (line 545)
- `notifications` aus taginfo entfernen (line 432)
- `notifications` aus imports entfernen (line 57)
**Verbindungen:**
- Alle ehemaligen Notification-Consumer nutzen `post_system_message()` aus `app/core/notifications.py`
- Frontend nutzt Comm-System-Channel API
**Migrationen:**
- `0131_migrate_notifications_to_comm.py` — Daten-Migration + Drop tables
**Tests:**
- `tests/test_notification_deprecation.py` — Verify deprecated routes redirect/404
- `tests/test_notification_migration.py` — Verify data migration correctness
- Existierende `tests/test_notifications.py` anpassen (nur noch post_system_message testen)
**Frontend:**
- `frontend/src/components/comm/SystemChannelBell.tsx` (NEU) — ersetzt NotificationDropdown
- Header-Komponente aktualisieren
---
## Phase F Lücken
### F-PREBUILT: Pre-built Agents in plugin.py on_activate registrieren
**Basis:** `app/plugins/builtins/automation/plugin.py` on_activate, `app/plugins/builtins/automation/prebuilt/`
**Was existiert:**
- 4 Pre-built Agent Definitionen: `email_triage_agent.py`, `follow_up_agent.py`, `contact_enrichment_agent.py`, `report_agent.py`
- Jede Datei hat eine `create_*_agent()` Funktion
- `on_activate` in `plugin.py` registriert agent_comm_tool, agent_coordinator_tools, MiniApps, cron jobs — **aber NICHT die prebuilt agents**
**Was neu gebaut wird:**
1. **Prebuilt Agent Registration:** `app/plugins/builtins/automation/plugin.py` on_activate erweitern
```python
async def on_activate(self, db, service_container, event_bus) -> None:
await super().on_activate(db, service_container, event_bus)
# ... existing registrations ...
# Register pre-built agents
try:
from app.plugins.builtins.automation.prebuilt.email_triage_agent import create_email_triage_agent
from app.plugins.builtins.automation.prebuilt.follow_up_agent import create_follow_up_agent
from app.plugins.builtins.automation.prebuilt.contact_enrichment_agent import create_contact_enrichment_agent
from app.plugins.builtins.automation.prebuilt.report_agent import create_report_agent
for create_fn in [create_email_triage_agent, create_follow_up_agent,
create_contact_enrichment_agent, create_report_agent]:
agent = await create_fn(db)
logger.info(f"Pre-built agent registered: {agent.name}")
except Exception:
logger.exception("Failed to register pre-built agents")
```
2. **Idempotency:** Jede `create_*_agent()` Funktion prüft ob Agent bereits existiert (by name + tenant)
- Wenn ja: skip (kein Duplikat)
- Wenn nein: erstelle AgentDefinition
3. **on_deactivate cleanup:** Pre-built agents werden bei Deaktivierung entfernt (oder als system-markiert belassen)
**Verbindungen:**
- `plugin.py on_activate` → `prebuilt/*.create_*_agent()`
- AgentDefinition in DB → verfügbar in Agent Dashboard und Agent Runner
**Migrationen:**
- Keine neue Migration (AgentDefinition Tabelle existiert bereits)
**Tests:**
- `tests/test_prebuilt_agent_registration.py` — Verify 4 agents exist after on_activate
- `tests/test_prebuilt_agent_idempotency.py` — Verify no duplicates on re-activate
**Frontend:**
- Keine Änderung — Agents erscheinen automatisch im Agent Dashboard
---
### F-AGENT-COMM: Agent Run-Results in Kommunikationszentrale posten
**Basis:** `app/plugins/builtins/automation/agent_runner.py`, `app/plugins/builtins/kommunikation/contracts.py`
**Was existiert:**
- `agent_comm.py` registriert ein `send_message` Tool für Agents (postet in comm conversations)
- `agent_runner.py` führt Agent Loop aus und sammelt Results
- KommunikationContract: `send_message()`, `post_system_message()`, `create_plugin_room()`
- **Aber:** Agent Run-Results (Zusammenfassung, Steps, Cost) werden NICHT automatisch in Kommunikationszentrale gepostet
**Was neu gebaut wird:**
1. **Agent Result Poster:** `app/plugins/builtins/automation/agent_result_poster.py` (NEU, ~200 Zeilen)
```python
async def post_agent_run_result(
db: AsyncSession,
tenant_id: uuid.UUID,
agent_run_id: uuid.UUID,
agent_name: str,
result: AgentRunResult,
user_id: uuid.UUID,
) -> None:
"""Post agent run result to kommunikation system channel."""
from app.plugins.builtins.contracts import get_contract
komm = get_contract("kommunikation")
if not komm:
return
# Create rich message with blocks
blocks = [
{"block_type": "markdown", "block_data": {"markdown": f"## Agent: {agent_name}\n\n{result.final_content}"}},
{"block_type": "action_card", "block_data": {
"title": "Agent Run Summary",
"body": f"Steps: {result.steps_taken} | Cost: ${result.total_cost_usd:.4f} | Status: {result.status}",
"actions": [
{"label": "View Details", "action": f"/agents/runs/{agent_run_id}", "type": "link"},
],
}},
]
await komm.send_message(
db=db, tenant_id=tenant_id,
conversation_id=system_channel_id, # get_or_create_system_channel
sender_id=agent_run_id,
sender_type="ai",
content=result.final_content or "Agent run completed",
blocks=blocks,
metadata={"agent_run_id": str(agent_run_id), "agent_name": agent_name},
)
```
2. **Agent Runner Integration:** `app/plugins/builtins/automation/agent_runner.py`
- Nach `run_react_loop()` completion: rufe `post_agent_run_result()` auf
- Bei Error: poste error summary in system channel
- Bei Approval-Pause: poste approval request in system channel
```python
# In agent_runner.py nach run_react_loop:
from app.plugins.builtins.automation.agent_result_poster import post_agent_run_result
await post_agent_run_result(db, tenant_id, agent_run_id, agent_name, result, user_id)
```
3. **CommMessageBlock Types erweitern:** `app/plugins/builtins/kommunikation/content_types.py`
- Neuer block_type: `"agent_result"` mit fields: `agent_run_id`, `agent_name`, `status`, `steps`, `cost`
- Neuer block_type: `"approval_request"` mit fields: `approval_id`, `action`, `entity_type`, `entity_id`
**Verbindungen:**
- `agent_runner.py` → `agent_result_poster.py` → `kommunikation.contracts.send_message()`
- `content_types.py` erweitert um agent_result und approval_request blocks
**Migrationen:**
- Keine (CommMessageBlock nutzt JSONB, schema-flexibel)
**Tests:**
- `tests/test_agent_result_posting.py` — Verify agent run results appear in system channel
- `tests/test_agent_result_blocks.py` — Verify block structure and types
**Frontend:**
- `frontend/src/components/comm/AgentResultBlock.tsx` (NEU) — Rendert agent_result block_type
- `frontend/src/components/comm/ApprovalRequestBlock.tsx` (NEU) — Rendert approval_request block_type
- Block-Renderer in ChatView erweitern
---
### F-WORK: Agent Workstream auf kommunikation Plugin aufbauen
**Basis:** `app/plugins/builtins/kommunikation/` (CommConversation, CommMessage, CommMessageBlock, MiniAppRegistry)
**Was existiert:**
- KommunikationPlugin mit CommConversation (is_system, is_direct, metadata_)
- CommMessageBlock (block_type, block_data JSONB) — unterstützt action_card, contact_card, miniapp
- MiniAppRegistry mit register/unregister
- `create_plugin_room()` für plugin-spezifische Conversations
- Agent Comm Tool (send_message) bereits registriert
**Was neu gebaut wird:**
1. **Agent Workstream Service:** `app/plugins/builtins/automation/agent_workstream.py` (NEU, ~300 Zeilen)
```python
async def create_agent_workstream_room(
db: AsyncSession, tenant_id: uuid.UUID, agent_id: uuid.UUID, user_id: uuid.UUID
) -> dict:
"""Create a dedicated workstream conversation for an agent."""
from app.plugins.builtins.contracts import get_contract
komm = get_contract("kommunikation")
# Create plugin room with agent as participant
room = await komm.create_plugin_room(
db=db, tenant_id=tenant_id,
plugin_name="automation",
room_key=f"agent_{agent_id}",
title=f"Agent Workstream",
metadata={"agent_id": str(agent_id), "type": "agent_workstream"},
)
# Add agent and user as participants
return room
async def post_workstream_update(
db, tenant_id, conversation_id, sender_type, sender_id, content, blocks=None
) -> dict:
"""Post a workstream update with rich blocks."""
from app.plugins.builtins.contracts import get_contract
komm = get_contract("kommunikation")
return await komm.send_message(
db=db, tenant_id=tenant_id,
conversation_id=conversation_id,
sender_id=sender_id, sender_type=sender_type,
content=content, blocks=blocks or [],
)
async def get_agent_workstream(db, tenant_id, agent_id, user_id) -> dict | None:
"""Get or create workstream room for agent."""
# Lookup by metadata.agent_id, create if not exists
```
2. **Workstream Block Types:** `app/plugins/builtins/kommunikation/content_types.py` erweitern
- `"task_card"` — fields: task_id, title, status, assignee, due_date
- `"workflow_card"` — fields: workflow_id, instance_id, status, current_step
- `"knowledge_card"` — fields: entity_type, entity_id, title, source
- `"progress_card"` — fields: current, total, label, percentage
3. **Agent Runner Integration:** `agent_runner.py`
- Bei Agent-Start: `create_agent_workstream_room()` → post "Agent started" message
- Bei jedem Step: `post_workstream_update()` mit progress_card block
- Bei Completion: post result summary (wie F-AGENT-COMM)
- Bei Approval: post approval_request block
4. **Workstream API Routes:** `app/plugins/builtins/automation/agent_routes.py` erweitern
- `GET /api/v1/agents/{id}/workstream` — Get workstream conversation for agent
- `POST /api/v1/agents/{id}/workstream/message` — Post message to workstream
- `GET /api/v1/agents/{id}/workstream/messages` — List workstream messages
**Verbindungen:**
- `agent_runner.py` → `agent_workstream.py` → `kommunikation.contracts.create_plugin_room()` + `send_message()`
- `content_types.py` erweitert um task_card, workflow_card, knowledge_card, progress_card
- `agent_routes.py` erweitert um workstream endpoints
**Migrationen:**
- Keine (nutzt vorhandene comm_conversations, comm_messages, comm_message_blocks)
**Tests:**
- `tests/test_agent_workstream.py` — Workstream room creation, message posting, block rendering
- `tests/test_agent_workstream_integration.py` — Agent run → workstream messages appear
**Frontend:**
- `frontend/src/components/comm/TaskCardBlock.tsx` (NEU)
- `frontend/src/components/comm/WorkflowCardBlock.tsx` (NEU)
- `frontend/src/components/comm/KnowledgeCardBlock.tsx` (NEU)
- `frontend/src/components/comm/ProgressCardBlock.tsx` (NEU)
- `frontend/src/components/comm/WorkstreamBlockRenderer.tsx` (NEU) — Dispatch block_type → component
- AgentChat-Seite erweitert um Workstream-View
---
## Phase G Lücke
### G-WORK: Workflow Workstream auf kommunikation Plugin aufbauen
**Basis:** `app/workflows/engine.py`, `app/plugins/builtins/kommunikation/`
**Was existiert:**
- WorkflowEngine verarbeitet Steps (action, approval, condition, wait, etc.)
- WorkflowInstance hat status, current_step_index, context
- `post_system_message()` in `app/core/notifications.py` delegiert an kommunikation
- engine.py importiert bereits `post_system_message` aus `app.core.notifications`
**Was neu gebaut wird:**
1. **Workflow Workstream Service:** `app/workflows/workstream.py` (NEU, ~250 Zeilen)
```python
async def create_workflow_workstream_room(
db: AsyncSession, tenant_id: uuid.UUID,
workflow_instance_id: uuid.UUID, user_id: uuid.UUID
) -> dict:
"""Create workstream conversation for a workflow instance."""
from app.plugins.builtins.contracts import get_contract
komm = get_contract("kommunikation")
room = await komm.create_plugin_room(
db=db, tenant_id=tenant_id,
plugin_name="workflows",
room_key=f"wf_{workflow_instance_id}",
title=f"Workflow: {workflow_name}",
metadata={"workflow_instance_id": str(workflow_instance_id), "type": "workflow_workstream"},
)
return room
async def post_workflow_step_update(
db, tenant_id, conversation_id, step_index, step_type, status, result=None
) -> dict:
"""Post workflow step progress to workstream."""
blocks = [
{"block_type": "progress_card", "block_data": {
"current": step_index + 1, "total": total_steps,
"label": f"Step {step_index + 1}: {step_type}", "percentage": int((step_index + 1) / total_steps * 100),
}},
{"block_type": "workflow_card", "block_data": {
"workflow_id": str(workflow_id), "instance_id": str(instance_id),
"status": status, "current_step": step_index,
}},
]
await komm.send_message(db, tenant_id, conversation_id, ...)
async def post_workflow_approval_request(
db, tenant_id, conversation_id, approval_id, step_description
) -> dict:
"""Post approval request to workstream."""
blocks = [{"block_type": "approval_request", "block_data": {
"approval_id": str(approval_id), "action": "approve_workflow_step",
"entity_type": "workflow_instance", "entity_id": str(instance_id),
}}]
await komm.send_message(...)
async def post_workflow_completion(db, tenant_id, conversation_id, status, summary):
"""Post workflow completion summary."""
```
2. **Workflow Engine Integration:** `app/workflows/engine.py` erweitern
- Bei Instance-Start: `create_workflow_workstream_room()` → post "Workflow started" message
- Bei jedem Step-Übergang: `post_workflow_step_update()`
- Bei Approval-Step: `post_workflow_approval_request()`
- Bei Completion: `post_workflow_completion()`
```python
# In engine.py process_step():
from app.workflows.workstream import post_workflow_step_update
await post_workflow_step_update(db, tenant_id, conversation_id, step_index, step_type, status, result)
```
3. **Workflow Routes erweitern:** `app/routes/workflows.py`
- `GET /api/v1/workflows/instances/{id}/workstream` — Get workstream conversation
- `POST /api/v1/workflows/instances/{id}/workstream/message` — Post message
**Verbindungen:**
- `engine.py` → `workstream.py` → `kommunikation.contracts.create_plugin_room()` + `send_message()`
- `workstream.py` nutzt workflow_card, progress_card, approval_request block types (aus F-WORK definiert)
**Migrationen:**
- Keine (nutzt vorhandene comm Tabellen)
**Tests:**
- `tests/test_workflow_workstream.py` — Workstream room creation, step updates, approval posts
- `tests/test_workflow_workstream_integration.py` — Full workflow run → workstream messages
**Frontend:**
- `frontend/src/pages/WorkflowDetail.tsx` erweitert um Workstream-Tab
- Nutzt WorkstreamBlockRenderer aus F-WORK
---
## Phase H Rest
### H-SRC: Knowledge Source Adapter — auf unified_search providers aufbauen
**Basis:** `app/plugins/builtins/unified_search/` (BaseSearchProvider, SearchProviderRegistry)
**Was existiert:**
- 14+ Search Providers (contact, company, task, mail, etc.)
- BaseSearchProvider mit search_fts, search_vector, get_embedding_text
- SearchProviderRegistry mit register/unregister/get_all
**Was neu gebaut wird:**
1. **Knowledge Source Adapter:** `app/plugins/builtins/knowledge/source_adapter.py` (NEU, ~300 Zeilen)
```python
class KnowledgeSourceAdapter:
"""Adapts unified_search providers as knowledge sources for extraction."""
async def fetch_source_content(
self, db: AsyncSession, tenant_id: uuid.UUID,
entity_type: str, entity_id: uuid.UUID,
) -> dict[str, Any] | None:
"""Fetch full content for an entity via its search provider."""
from app.plugins.builtins.unified_search.contracts import get_search_registry
registry = get_search_registry()
provider = registry.get(entity_type)
if not provider:
return None
embedding_text = await provider.get_embedding_text(db, entity_id, tenant_id)
return {"entity_type": entity_type, "entity_id": str(entity_id), "content": embedding_text}
async def batch_fetch(
self, db, tenant_id, items: list[dict],
) -> list[dict]:
"""Batch fetch content for multiple entities."""
async def discover_sources(
self, db, tenant_id, since: datetime | None = None,
) -> list[dict]:
"""Discover all entities that could be knowledge sources."""
# Query all providers for recently updated entities
```
2. **Knowledge Source Model:** `app/plugins/builtins/knowledge/models.py` (NEU)
```python
class KnowledgeSource(Base, TenantMixin, OwnedMixin):
__tablename__ = "knowledge_sources"
id: UUID PK
entity_type: str # "contact", "company", "mail", "file", "wiki_article"
entity_id: UUID
content_hash: str # SHA256 of content for change detection
last_extracted_at: TIMESTAMPTZ | None
extraction_status: str # "pending", "extracted", "failed", "stale"
metadata_: dict = JSONB
```
3. **Knowledge Plugin:** `app/plugins/builtins/knowledge/` (NEU, komplettes Plugin)
- `plugin.py` — `KnowledgePlugin(BasePlugin)` mit Manifest, dependencies: ["unified_search", "graph_rag"]
- `routes.py` — `/api/v1/knowledge/sources`, `/api/v1/knowledge/extraction`
- `services.py` — Extraction orchestration
- `schemas.py` — Pydantic schemas
**Verbindungen:**
- `source_adapter.py` → `unified_search.contracts.get_search_registry()` → provider.get_embedding_text()
- Knowledge Plugin dependencies: unified_search, graph_rag
**Migrationen:**
- `0132_knowledge_sources.py` — knowledge_sources table
- Plugin-Migration: `knowledge/migrations/0001_initial.sql`
**Tests:**
- `tests/test_knowledge_source_adapter.py` — Fetch content via providers, batch fetch, discover
- `tests/test_knowledge_sources.py` — CRUD API, change detection
**Frontend:**
- `frontend/src/pages/KnowledgeDashboard.tsx` (NEU) — Source overview, extraction status
---
### H-LLM-REL: LLM Relationship Extraction — auf graph_rag aufbauen
**Basis:** `app/plugins/builtins/graph_rag/services.py` (create_relationship), `app/ai/llm_client.py` (llm_complete)
**Was neu gebaut wird:**
1. **Relationship Extractor:** `app/plugins/builtins/knowledge/relationship_extractor.py` (NEU, ~250 Zeilen)
```python
async def extract_relationships(
db: AsyncSession, tenant_id: uuid.UUID,
entity_type: str, entity_id: uuid.UUID, content: str,
) -> list[dict]:
"""Use LLM to extract relationships from entity content."""
from app.ai.llm_client import llm_complete
prompt = f"""Analyze the following content and extract relationships.
Return JSON array of {{source_type, source_id, target_type, target_id, relationship_type, confidence}}.
Content: {content[:4000]}
Entity: {entity_type} {entity_id}
"""
response = await llm_complete(messages=[{"role": "user", "content": prompt}], ...)
relationships = parse_llm_relationships(response, entity_type, entity_id)
# Create relationships via graph_rag
from app.plugins.builtins.graph_rag.services import create_relationship
for rel in relationships:
await create_relationship(db, tenant_id, **rel)
return relationships
```
2. **LLM Prompt Template:** Strukturiertes Prompt für Relationship-Extraction
- Input: Entity content + context
- Output: JSON array of relationships with confidence scores
- System prompt: Domain-specific relationship types (works_for, has_email, related_to, etc.)
**Verbindungen:**
- `relationship_extractor.py` → `llm_client.llm_complete()` → `graph_rag.services.create_relationship()`
**Migrationen:**
- Keine (EntityRelationship existiert bereits)
**Tests:**
- `tests/test_relationship_extraction.py` — LLM mock, verify relationships created in graph_rag
**Frontend:**
- Knowledge Dashboard: Extracted relationships view
---
### H-ENT: Entity Extraction — auf graph_rag aufbauen
**Basis:** `app/plugins/builtins/graph_rag/models.py`, `app/ai/llm_client.py`
**Was neu gebaut wird:**
1. **Entity Extractor:** `app/plugins/builtins/knowledge/entity_extractor.py` (NEU, ~250 Zeilen)
```python
async def extract_entities(
db: AsyncSession, tenant_id: uuid.UUID,
content: str, source_type: str, source_id: uuid.UUID,
) -> list[dict]:
"""Use LLM to extract named entities from content."""
from app.ai.llm_client import llm_complete
prompt = f"""Extract named entities from the following content.
Return JSON array of {{entity_type, name, attributes, mentions: [{{start, end}}]}}.
Entity types: person, organization, email, phone, date, location, project
Content: {content[:4000]}
"""
response = await llm_complete(...)
entities = parse_llm_entities(response)
# Link entities to existing CRM records or create new ones
for entity in entities:
matched = await match_entity_to_crm(db, tenant_id, entity)
if matched:
# Create relationship: source → matched entity
await create_relationship(db, tenant_id, source_type, source_id, matched.type, matched.id, "mentions")
return entities
```
2. **Entity Matching Service:** `app/plugins/builtins/knowledge/entity_matcher.py` (NEU, ~150 Zeilen)
- Match extracted entities against existing contacts, companies, etc.
- Fuzzy matching by name, email, phone
- Returns matched entity or None
**Verbindungen:**
- `entity_extractor.py` → `llm_client.llm_complete()` → `graph_rag.services.create_relationship()`
- `entity_matcher.py` → Contact/Company models for matching
**Tests:**
- `tests/test_entity_extraction.py` — LLM mock, verify entity extraction and matching
**Frontend:**
- Knowledge Dashboard: Extracted entities view
---
### H-AUTO: Auto-Relationship Creation — auf graph_rag aufbauen
**Basis:** `app/plugins/builtins/graph_rag/services.py`, `app/core/hooks.py`, `app/core/outbox.py`
**Was neu gebaut wird:**
1. **Auto-Relationship Engine:** `app/plugins/builtins/knowledge/auto_relationship.py` (NEU, ~200 Zeilen)
```python
async def auto_create_relationships(
db: AsyncSession, tenant_id: uuid.UUID,
entity_type: str, entity_id: uuid.UUID,
) -> list[dict]:
"""Automatically create relationships based on entity data."""
# Rule-based: contact → company (works_for), mail → contact (sent_by), etc.
rules = get_relationship_rules(entity_type)
relationships = []
for rule in rules:
targets = await rule.find_targets(db, tenant_id, entity_type, entity_id)
for target in targets:
result = await create_relationship(db, tenant_id, entity_type, entity_id, target.type, target.id, rule.relationship_type)
if "error" not in result:
relationships.append(result)
return relationships
```
2. **Hook Integration:** `app/plugins/builtins/knowledge/plugin.py` on_activate
- Register hooks: `contact.after_create`, `contact.after_update`, `mail.received`, `company.after_create`
- On hook fire: `auto_create_relationships()`
3. **Rule Definitions:** `app/plugins/builtins/knowledge/rules.py` (NEU, ~200 Zeilen)
- Contact → Company: if contact.company_id exists, create "works_for" relationship
- Mail → Contact: if mail.from_address matches contact.email, create "sent_by" relationship
- Task → Contact: if task.assigned_to matches contact, create "assigned_to" relationship
- File → Contact: if file.metadata has contact reference, create "belongs_to" relationship
**Verbindungen:**
- `auto_relationship.py` → `graph_rag.services.create_relationship()`
- `plugin.py on_activate` → `hooks.register_action()`
- Hook callbacks → `auto_create_relationships()`
**Tests:**
- `tests/test_auto_relationships.py` — Create contact with company → verify "works_for" relationship
- `tests/test_relationship_rules.py` — Each rule tested independently
**Frontend:**
- Knowledge Dashboard: Auto-created relationships view
---
### H-CONF: Confidence Score + Review Queue — auf graph_rag aufbauen
**Basis:** `app/plugins/builtins/graph_rag/models.py` (EntityRelationship)
**Was neu gebaut wird:**
1. **EntityRelationship erweitern:** Neue Migration fügt confidence und review_status hinzu
```python
# In EntityRelationship (via migration):
confidence: float = 0.0 # 0.0-1.0
review_status: str = "auto" # "auto", "pending_review", "approved", "rejected"
reviewed_by: UUID | None
reviewed_at: TIMESTAMPTZ | None
extraction_method: str = "auto" # "auto", "llm", "manual"
```
2. **Review Queue Service:** `app/plugins/builtins/knowledge/review_queue.py` (NEU, ~200 Zeilen)
```python
async def get_pending_reviews(db, tenant_id, limit=50) -> list[EntityRelationship]: ...
async def approve_relationship(db, tenant_id, rel_id, user_id) -> dict: ...
async def reject_relationship(db, tenant_id, rel_id, user_id, reason) -> dict: ...
async def batch_approve(db, tenant_id, rel_ids, user_id) -> dict: ...
```
3. **Review API:** `app/plugins/builtins/knowledge/routes.py` erweitern
- `GET /api/v1/knowledge/review-queue` — List pending relationships
- `POST /api/v1/knowledge/review/{id}/approve` — Approve
- `POST /api/v1/knowledge/review/{id}/reject` — Reject
- `POST /api/v1/knowledge/review/batch-approve` — Batch approve
**Verbindungen:**
- EntityRelationship erweitert um confidence/review_status
- Review Queue nutzt graph_rag models
**Migrationen:**
- `0133_entity_relationship_confidence.py` — Add confidence, review_status, reviewed_by, reviewed_at, extraction_method
**Tests:**
- `tests/test_review_queue.py` — Pending reviews, approve, reject, batch approve
- `tests/test_confidence_scoring.py` — Confidence threshold filtering
**Frontend:**
- `frontend/src/pages/KnowledgeReviewQueue.tsx` (NEU) — Review queue UI
- `frontend/src/components/knowledge/RelationshipReviewCard.tsx` (NEU)
---
### H-EVT: Event-Driven Extraction — auf Event Bus + graph_rag aufbauen
**Basis:** `app/core/event_bus.py`, `app/core/outbox.py`, `app/core/hooks.py`
**Was neu gebaut wird:**
1. **Event-Driven Extraction Handler:** `app/plugins/builtins/knowledge/event_handler.py` (NEU, ~200 Zeilen)
```python
async def handle_entity_created(event_name: str, payload: dict):
"""Handle entity.created event — trigger extraction."""
entity_type = payload.get("entity_type")
entity_id = uuid.UUID(payload.get("entity_id"))
tenant_id = uuid.UUID(payload.get("tenant_id"))
# Enqueue extraction job
from app.core.jobs import enqueue_job
await enqueue_job("knowledge_extract", {
"entity_type": entity_type, "entity_id": str(entity_id),
"tenant_id": str(tenant_id),
})
async def handle_entity_updated(event_name: str, payload: dict):
"""Handle entity.updated — mark source as stale, re-extract."""
```
2. **Extraction Job:** `app/plugins/builtins/knowledge/jobs.py` (NEU, ~150 Zeilen)
- ARQ job function: `async def knowledge_extract(ctx, entity_type, entity_id, tenant_id)`
- Calls source_adapter → entity_extractor → relationship_extractor → auto_relationship
- Updates KnowledgeSource.extraction_status
3. **Event Registration:** `app/plugins/builtins/knowledge/plugin.py` on_activate
```python
event_bus.subscribe("contact.created", handle_entity_created)
event_bus.subscribe("contact.updated", handle_entity_updated)
event_bus.subscribe("mail.received", handle_entity_created)
event_bus.subscribe("company.created", handle_entity_created)
event_bus.subscribe("company.updated", handle_entity_updated)
```
**Verbindungen:**
- `event_handler.py` → `event_bus.subscribe()` → `jobs.enqueue_job()` → extraction pipeline
- Extraction pipeline: source_adapter → entity_extractor → relationship_extractor → graph_rag
**Tests:**
- `tests/test_event_driven_extraction.py` — Fire event → verify extraction job enqueued and executed
**Frontend:**
- Knowledge Dashboard: Extraction status per source
---
### H-CITE: Evidence/Source References — auf unified_search aufbauen
**Basis:** `app/plugins/builtins/unified_search/` (search results with entity references)
**Was neu gebaut wird:**
1. **Evidence Model:** `app/plugins/builtins/knowledge/models.py` erweitern
```python
class KnowledgeEvidence(Base, TenantMixin):
__tablename__ = "knowledge_evidence"
id: UUID PK
relationship_id: UUID FK → entity_relationships.id
source_type: str # "contact", "mail", "file", "wiki_article"
source_id: UUID
source_snippet: str # Text snippet that supports the relationship
confidence: float
created_at: TIMESTAMPTZ
```
2. **Evidence Service:** `app/plugins/builtins/knowledge/evidence_service.py` (NEU, ~150 Zeilen)
```python
async def add_evidence(db, tenant_id, relationship_id, source_type, source_id, snippet, confidence): ...
async def get_evidence_for_relationship(db, tenant_id, relationship_id) -> list[dict]: ...
async def search_evidence(db, tenant_id, query) -> list[dict]: ...
```
3. **Evidence API:** `app/plugins/builtins/knowledge/routes.py` erweitern
- `GET /api/v1/knowledge/relationships/{id}/evidence` — List evidence for relationship
- `POST /api/v1/knowledge/relationships/{id}/evidence` — Add evidence
**Verbindungen:**
- Evidence linked to EntityRelationship (graph_rag)
- Source references use unified_search entity types
**Migrationen:**
- `0134_knowledge_evidence.py` — knowledge_evidence table
**Tests:**
- `tests/test_knowledge_evidence.py` — Add, list, search evidence
**Frontend:**
- Relationship Detail: Evidence section
---
### H-WIKI-SEARCH: Wiki Search Provider — Wiki als unified_search provider registrieren
**Basis:** `app/plugins/builtins/wiki/` (WikiPlugin, models, routes), `app/plugins/builtins/unified_search/` (BaseSearchProvider)
**Was existiert:**
- WikiPlugin hat kein on_activate (kein search provider registriert)
- WikiArticle Model existiert
- BaseSearchProvider mit search_fts, search_vector
**Was neu gebaut wird:**
1. **Wiki Search Provider:** `app/plugins/builtins/wiki/search_provider.py` (NEU, ~200 Zeilen)
```python
class WikiSearchProvider(BaseSearchProvider):
entity_type = "wiki_article"
supports_fts = True
supports_vector = True
supports_rag = True
async def _search_fts_filtered(self, db, tsquery, tenant_id, limit, visible_ids):
# Search wiki_articles.title and content via FTS
async def _search_vector_filtered(self, db, embedding, tenant_id, limit, visible_ids):
# Search wiki_articles.embedding via vector cosine
async def get_embedding_text(self, db, entity_id, tenant_id) -> str:
# Return article title + content
def to_search_result(self, article) -> dict:
return {"entity_type": "wiki_article", "entity_id": str(article.id), "title": article.title, ...}
```
2. **Wiki Plugin on_activate:** `app/plugins/builtins/wiki/plugin.py` erweitern
```python
async def on_activate(self, db, service_container, event_bus) -> None:
from app.plugins.builtins.unified_search.contracts import get_search_registry
from app.plugins.builtins.wiki.search_provider import WikiSearchProvider
registry = get_search_registry()
registry.register(WikiSearchProvider())
await super().on_activate(db, service_container, event_bus)
async def on_deactivate(self, db, service_container, event_bus) -> None:
from app.plugins.builtins.unified_search.contracts import get_search_registry
get_search_registry().unregister("wiki_article")
await super().on_deactivate(db, service_container, event_bus)
```
3. **Wiki Embedding Index:** `app/plugins/builtins/wiki/models.py` erweitern
- WikiArticle braucht `embedding` vector(768) column (via migration)
- `lifecycle.py` hook: On article save, generate embedding via `llm_embed()`
**Verbindungen:**
- `wiki/plugin.py on_activate` → `unified_search.contracts.get_search_registry().register(WikiSearchProvider())`
- `wiki/search_provider.py` extends `unified_search.base_provider.BaseSearchProvider`
**Migrationen:**
- `0135_wiki_embedding_column.py` — Add embedding vector(768) to wiki_articles
**Tests:**
- `tests/test_wiki_search_provider.py` — FTS and vector search on wiki articles
- `tests/test_wiki_search_integration.py` — Wiki results appear in unified search
**Frontend:**
- Keine Änderung — Wiki results appear in global search
---
### H-WIKI-EMBED: Wiki Embeddings — auf llm_client.llm_embed aufbauen
**Basis:** `app/ai/llm_client.py` (llm_embed), `app/plugins/builtins/wiki/models.py`
**Was neu gebaut wird:**
1. **Wiki Embedding Service:** `app/plugins/builtins/wiki/embedding_service.py` (NEU, ~150 Zeilen)
```python
async def generate_wiki_embedding(db, tenant_id, article_id) -> None:
"""Generate and store embedding for wiki article."""
from app.ai.llm_client import llm_embed
article = await get_article(db, tenant_id, article_id)
text = f"{article.title}\n\n{article.content}"
embedding = await llm_embed(text)
article.embedding = embedding
await db.commit()
async def batch_generate_embeddings(db, tenant_id, batch_size=50) -> int:
"""Generate embeddings for all articles without one."""
```
2. **Wiki Lifecycle Hook:** `app/plugins/builtins/wiki/plugin.py` on_activate
- Register hook: `wiki.article.after_create` → `generate_wiki_embedding()`
- Register hook: `wiki.article.after_update` → `generate_wiki_embedding()` (re-embed)
3. **Batch Job:** `app/plugins/builtins/wiki/jobs.py` (NEU)
- ARQ job: `async def wiki_embed_batch(ctx, tenant_id)`
- Called via cron or manual trigger
**Verbindungen:**
- `embedding_service.py` → `llm_client.llm_embed()` → update WikiArticle.embedding
- `plugin.py on_activate` → `hooks.register_action()`
**Migrationen:**
- Siehe H-WIKI-SEARCH (0135 fügt embedding column hinzu)
**Tests:**
- `tests/test_wiki_embeddings.py` — Generate embedding, verify vector stored, batch job
**Frontend:**
- Wiki Settings: "Re-generate embeddings" button
---
### H-WIKI-LINK: Auto-Linking — Wiki → Entity Links
**Basis:** `app/plugins/builtins/wiki/models.py` (entity links), `app/plugins/builtins/graph_rag/services.py`
**Was neu gebaut wird:**
1. **Wiki Auto-Linker:** `app/plugins/builtins/wiki/auto_linker.py` (NEU, ~200 Zeilen)
```python
async def auto_link_wiki_to_entities(
db: AsyncSession, tenant_id: uuid.UUID, article_id: uuid.UUID,
) -> list[dict]:
"""Analyze wiki article content and auto-link to CRM entities."""
article = await get_article(db, tenant_id, article_id)
# 1. Extract entities from article content
from app.plugins.builtins.knowledge.entity_extractor import extract_entities
entities = await extract_entities(db, tenant_id, article.content, "wiki_article", article_id)
# 2. Create entity links
for entity in entities:
if entity.get("matched_id"):
await create_entity_link(db, tenant_id, "wiki_article", article_id, entity["type"], entity["matched_id"])
# Also create graph_rag relationship
await create_relationship(db, tenant_id, "wiki_article", article_id, entity["type"], entity["matched_id"], "references")
return entities
```
2. **Wiki Hook Integration:** `app/plugins/builtins/wiki/plugin.py` on_activate
- Register hook: `wiki.article.after_create` → `auto_link_wiki_to_entities()`
- Register hook: `wiki.article.after_update` → `auto_link_wiki_to_entities()` (re-link)
**Verbindungen:**
- `auto_linker.py` → `knowledge.entity_extractor.extract_entities()` → `graph_rag.services.create_relationship()`
- `plugin.py on_activate` → `hooks.register_action()`
**Tests:**
- `tests/test_wiki_auto_linking.py` — Create article mentioning contact → verify link created
**Frontend:**
- Wiki Article Detail: "Linked Entities" section (auto-linked + manual)
---
### H-DATA-LIFE: Derived-Data Lifecycle — auf Outbox + Event Bus aufbauen
**Basis:** `app/core/outbox.py` (enqueue_outbox_event), `app/core/event_bus.py`, `app/core/hooks.py`
**Was neu gebaut wird:**
1. **Derived Data Tracker:** `app/plugins/builtins/knowledge/derived_data.py` (NEU, ~250 Zeilen)
```python
class DerivedDataRegistry:
"""Tracks which derived data depends on which source data."""
def register_dependency(self, source_type, source_id, derived_type, derived_id, plugin_name): ...
def get_dependencies(self, source_type, source_id) -> list[dict]: ...
async def invalidate_dependents(self, db, tenant_id, source_type, source_id): ...
```
2. **Derived Data Model:** `app/plugins/builtins/knowledge/models.py` erweitern
```python
class DerivedDataDependency(Base, TenantMixin):
__tablename__ = "derived_data_dependencies"
id: UUID PK
source_type: str
source_id: UUID
derived_type: str # "embedding", "graph_relationship", "search_index", "rag_chunk"
derived_id: UUID
plugin_name: str
is_valid: bool = True
invalidated_at: TIMESTAMPTZ | None
created_at: TIMESTAMPTZ
```
3. **Event Handler:** `app/plugins/builtins/knowledge/lifecycle_handler.py` (NEU, ~150 Zeilen)
```python
async def handle_source_updated(event_name, payload):
"""When source data changes, invalidate derived data."""
source_type = payload["entity_type"]
source_id = uuid.UUID(payload["entity_id"])
tenant_id = uuid.UUID(payload["tenant_id"])
await registry.invalidate_dependents(db, tenant_id, source_type, source_id)
# Enqueue re-computation job
await enqueue_job("recompute_derived_data", {
"source_type": source_type, "source_id": str(source_id),
"tenant_id": str(tenant_id),
})
```
4. **Plugin Registration:** `app/plugins/builtins/knowledge/plugin.py` on_activate
- Subscribe to: `contact.updated`, `company.updated`, `mail.updated`, `file.updated`, `wiki.article.updated`
- On event: `handle_source_updated()`
**Verbindungen:**
- `lifecycle_handler.py` → `outbox.enqueue_outbox_event()` → `event_bus`
- `derived_data.py` → tracks dependencies across plugins
**Migrationen:**
- `0136_derived_data_dependencies.py` — derived_data_dependencies table
**Tests:**
- `tests/test_derived_data_lifecycle.py` — Update contact → verify embeddings invalidated and recomputed
**Frontend:**
- Admin: Derived data status dashboard
---
### H-RET: Knowledge Retention — auf bestehende Retention Patterns aufbauen
**Basis:** `app/core/outbox.py` (DLQ, retention), `app/models/audit.py` (audit log retention)
**Was neu gebaut wird:**
1. **Retention Policy Model:** `app/plugins/builtins/knowledge/models.py` erweitern
```python
class KnowledgeRetentionPolicy(Base, TenantMixin):
__tablename__ = "knowledge_retention_policies"
id: UUID PK
entity_type: str # "knowledge_source", "entity_relationship", "knowledge_evidence"
max_age_days: int | None # None = unlimited
max_count: int | None # None = unlimited
action: str # "archive", "delete", "anonymize"
is_active: bool = True
created_at: TIMESTAMPTZ
```
2. **Retention Service:** `app/plugins/builtins/knowledge/retention.py` (NEU, ~200 Zeilen)
```python
async def apply_retention_policies(db, tenant_id) -> dict:
"""Apply all active retention policies."""
policies = await get_active_policies(db, tenant_id)
results = {}
for policy in policies:
if policy.action == "archive":
count = await archive_old_records(db, tenant_id, policy)
elif policy.action == "delete":
count = await delete_old_records(db, tenant_id, policy)
elif policy.action == "anonymize":
count = await anonymize_old_records(db, tenant_id, policy)
results[policy.entity_type] = count
return results
```
3. **Retention Job:** `app/plugins/builtins/knowledge/jobs.py` erweitern
- ARQ job: `async def apply_retention(ctx, tenant_id)`
- Cron: täglich um 3 Uhr
4. **Retention API:** `app/plugins/builtins/knowledge/routes.py` erweitern
- `GET /api/v1/knowledge/retention/policies` — List policies
- `POST /api/v1/knowledge/retention/policies` — Create policy
- `PUT /api/v1/knowledge/retention/policies/{id}` — Update
- `DELETE /api/v1/knowledge/retention/policies/{id}` — Delete
- `POST /api/v1/knowledge/retention/apply` — Manual apply
**Verbindungen:**
- `retention.py` → KnowledgeSource, EntityRelationship, KnowledgeEvidence models
- Cron job registriert via automation plugin cron_jobs contribution
**Migrationen:**
- `0137_knowledge_retention.py` — knowledge_retention_policies table
**Tests:**
- `tests/test_knowledge_retention.py` — Create old records, apply policy, verify archived/deleted
**Frontend:**
- `frontend/src/pages/KnowledgeRetentionSettings.tsx` (NEU) — Policy management
---
## Phase I
### I-1: Cross-System Integration — Agent→Workflow, Agent→Knowledge, MCP
**Basis:** `app/ai/agent_loop.py`, `app/workflows/engine.py`, `app/plugins/builtins/knowledge/`
**Was neu gebaut wird:**
1. **Agent → Workflow Bridge:** `app/plugins/builtins/automation/agent_workflow_bridge.py` (NEU, ~200 Zeilen)
```python
async def agent_trigger_workflow(
db: AsyncSession, tenant_id: uuid.UUID,
agent_id: uuid.UUID, workflow_id: uuid.UUID, context: dict,
) -> dict:
"""Agent triggers a workflow execution."""
from app.services.workflow_service import create_instance
instance = await create_instance(db, tenant_id, workflow_id, context, initiated_by=agent_id)
# Post to agent workstream
await post_workstream_update(...)
return {"workflow_instance_id": str(instance.id)}
```
- Register as agent tool: `trigger_workflow`
2. **Agent → Knowledge Bridge:** `app/plugins/builtins/automation/agent_knowledge_bridge.py` (NEU, ~200 Zeilen)
```python
async def agent_query_knowledge(
db: AsyncSession, tenant_id: uuid.UUID,
query: str, entity_type: str | None = None,
) -> dict:
"""Agent queries knowledge graph and unified search."""
from app.plugins.builtins.unified_search.contracts import get_search_registry
from app.plugins.builtins.graph_rag.services import traverse_graph
# 1. Unified search for relevant entities
# 2. Graph traversal for related entities
# 3. Combine and return
```
- Register as agent tool: `query_knowledge`
3. **MCP Exposure:** `app/plugins/builtins/mcp/` (NEU, komplettes Plugin, ~400 Zeilen)
- `plugin.py` — `MCPPlugin(BasePlugin)` mit Manifest
- `server.py` — MCP Server mit tool definitions für CRM entities
- `routes.py` — `/api/v1/mcp/tools`, `/api/v1/mcp/call`
- Exposes CRM operations as MCP tools for external AI clients
**Verbindungen:**
- `agent_workflow_bridge.py` → `workflow_service.create_instance()`
- `agent_knowledge_bridge.py` → `unified_search` + `graph_rag`
- MCP Plugin → CRM routes (read-only initially)
**Migrationen:**
- Plugin-Migration: `mcp/migrations/0001_initial.sql`
**Tests:**
- `tests/test_agent_workflow_bridge.py` — Agent triggers workflow
- `tests/test_agent_knowledge_bridge.py` — Agent queries knowledge
- `tests/test_mcp_exposure.py` — MCP tool listing and calling
**Frontend:**
- Agent Editor: Available tools include `trigger_workflow`, `query_knowledge`
- Settings: MCP configuration
---
### I-2 bis I-5: Human-AI Workstream — auf kommunikation Plugin aufbauen
**Basis:** `app/plugins/builtins/kommunikation/` (CommConversation, CommMessageBlock, MiniAppRegistry)
**Was neu gebaut wird:**
1. **Human-AI Workstream Service:** `app/plugins/builtins/kommunikation/human_ai_workstream.py` (NEU, ~300 Zeilen)
```python
async def create_human_ai_session(
db, tenant_id, user_id, agent_id, topic: str,
) -> dict:
"""Create a human-AI collaboration session (CommConversation with type=human_ai)."""
conv = await create_plugin_room(db, tenant_id, "kommunikation", f"hai_{agent_id}_{user_id}",
title=f"Human-AI: {topic}", metadata={"type": "human_ai", "agent_id": str(agent_id)})
# Add user and agent as participants
return conv
async def post_ai_proposal(db, tenant_id, conversation_id, proposal_type, content, actions):
"""Post an AI proposal with action_card block for human review."""
blocks = [{"block_type": "action_card", "block_data": {
"title": f"AI Proposal: {proposal_type}",
"body": content,
"actions": actions, # [{label, action, type: "approve/reject/edit"}]
}}]
await send_message(db, tenant_id, conversation_id, sender_type="ai", ...)
async def handle_human_response(db, tenant_id, conversation_id, message_id, response_type, edits):
"""Process human response to AI proposal."""
```
2. **Workstream Block Types:** `content_types.py` erweitern
- `"ai_proposal"` — fields: proposal_type, content, actions, confidence
- `"human_decision"` — fields: decision, rationale, decided_by
- `"ai_explanation"` — fields: explanation, evidence_ids, confidence
3. **Workstream API:** `app/plugins/builtins/kommunikation/routes.py` erweitern
- `POST /api/v1/comm/workstream/sessions` — Create human-AI session
- `GET /api/v1/comm/workstream/sessions` — List sessions
- `POST /api/v1/comm/workstream/sessions/{id}/proposals` — Post AI proposal
- `POST /api/v1/comm/workstream/sessions/{id}/responses` — Human response
**Verbindungen:**
- `human_ai_workstream.py` → `kommunikation.services.send_message()` + `create_plugin_room()`
- Agent Runner → `post_ai_proposal()` when agent needs human input
**Tests:**
- `tests/test_human_ai_workstream.py` — Session creation, proposal, response flow
**Frontend:**
- `frontend/src/components/comm/AIProposalBlock.tsx` (NEU)
- `frontend/src/components/comm/HumanDecisionBlock.tsx` (NEU)
- `frontend/src/pages/HumanAIWorkstream.tsx` (NEU)
---
### I-6 bis I-8: MiniApp Runtime — auf kommunikation/miniapp_registry aufbauen
**Basis:** `app/plugins/builtins/kommunikation/miniapp_registry.py` (MiniAppRegistry, MiniAppDef)
**Was existiert:**
- MiniAppRegistry mit register/unregister/list_apps/get_app
- MiniAppDef: app_id, name, icon, description, plugin_name, render_schema
- CommMessageBlock block_type="miniapp" bereits definiert
**Was neu gebaut wird:**
1. **MiniApp Runtime Service:** `app/plugins/builtins/kommunikation/miniapp_runtime.py` (NEU, ~300 Zeilen)
```python
class MiniAppRuntime:
"""Executes mini-app actions and manages state."""
async def execute_action(
self, db, tenant_id, app_id, action_name, params, user_id,
) -> dict:
"""Execute a mini-app action."""
app = get_miniapp_registry().get_app(app_id)
if not app:
raise NotFoundError(f"MiniApp {app_id} not found")
# Execute action via plugin's action handler
handler = self._get_action_handler(app.plugin_name, app_id)
result = await handler(db, tenant_id, action_name, params, user_id)
return result
async def get_state(self, db, tenant_id, app_id, context) -> dict:
"""Get current mini-app state for rendering."""
```
2. **MiniApp Action API:** `app/plugins/builtins/kommunikation/routes.py` erweitern
- `GET /api/v1/comm/miniapps` — List all registered mini-apps
- `GET /api/v1/comm/miniapps/{app_id}` — Get mini-app definition
- `POST /api/v1/comm/miniapps/{app_id}/actions` — Execute action
- `GET /api/v1/comm/miniapps/{app_id}/state` — Get state
3. **Plugin Action Handler Convention:** Plugins die MiniApps registrieren, stellen action handlers bereit
- In `plugin.py`: `def get_miniapp_action_handler(self, app_id) -> Callable`
- Automation plugin: Agent control mini-app
- DMS plugin: File preview mini-app
- Tasks plugin: Task board mini-app
**Verbindungen:**
- `miniapp_runtime.py` → `miniapp_registry.get_app()` → plugin action handler
- `routes.py` → `miniapp_runtime.execute_action()`
**Tests:**
- `tests/test_miniapp_runtime.py` — Register app, execute action, get state
**Frontend:**
- `frontend/src/components/comm/MiniAppBlockRenderer.tsx` (NEU) — Renders miniapp blocks
- `frontend/src/components/miniapps/GenericMiniApp.tsx` (NEU) — Schema-based renderer
---
### I-9 bis I-12: Dashboard & Analytics — echte DB-Queries
**Basis:** Alle vorhandenen Models (Contact, Company, Mail, Task, Calendar, DMS, AgentRun, WorkflowInstance, etc.)
**Was neu gebaut wird:**
1. **Dashboard Service:** `app/services/dashboard_service.py` (NEU, ~400 Zeilen)
```python
async def get_dashboard_metrics(db, tenant_id, user_id) -> dict:
"""Get real dashboard metrics from DB."""
return {
"contacts": await count_contacts(db, tenant_id),
"companies": await count_companies(db, tenant_id),
"tasks": {
"open": await count_tasks(db, tenant_id, status="open"),
"overdue": await count_overdue_tasks(db, tenant_id),
"completed_this_week": await count_completed_this_week(db, tenant_id),
},
"mails": {
"unread": await count_unread_mails(db, tenant_id, user_id),
"total_today": await count_mails_today(db, tenant_id),
},
"agents": {
"active_runs": await count_active_agent_runs(db, tenant_id),
"total_runs": await count_total_agent_runs(db, tenant_id),
"success_rate": await calculate_agent_success_rate(db, tenant_id),
},
"workflows": {
"active_instances": await count_active_workflow_instances(db, tenant_id),
"completed_this_month": await count_completed_workflows(db, tenant_id),
},
"knowledge": {
"relationships": await count_relationships(db, tenant_id),
"sources_extracted": await count_extracted_sources(db, tenant_id),
},
}
async def get_activity_feed(db, tenant_id, user_id, limit=50) -> list[dict]:
"""Get recent activity across all systems."""
# Query audit_log, agent_runs, workflow_instances, comm_messages
```
2. **Dashboard API:** `app/routes/dashboard.py` (NEU)
- `GET /api/v1/dashboard/metrics` — Real-time metrics
- `GET /api/v1/dashboard/activity-feed` — Activity feed
- `GET /api/v1/dashboard/trends` — Trend data (7/30/90 days)
**Verbindungen:**
- `dashboard_service.py` → alle Core + Plugin Models (echte SQL queries)
- Keine Mocks, keine hardcoded data
**Tests:**
- `tests/test_dashboard.py` — Verify metrics match real DB counts
**Frontend:**
- `frontend/src/pages/Dashboard.tsx` (NEU/überarbeitet) — Real metrics
- `frontend/src/components/dashboard/MetricCard.tsx` (NEU)
- `frontend/src/components/dashboard/ActivityFeed.tsx` (NEU)
- `frontend/src/components/dashboard/TrendChart.tsx` (NEU)
---
### I-13 bis I-16: DSGVO Export — echte DB-Queries über alle Plugins
**Basis:** Alle Models mit tenant_id + user_id reference
**Was neu gebaut wird:**
1. **DSGVO Export Service:** `app/services/dsgvo_export.py` (NEU, ~400 Zeilen)
```python
async def export_user_data(db, tenant_id, user_id) -> dict:
"""Export all data associated with a user for DSGVO compliance."""
data = {
"user": await get_user_data(db, tenant_id, user_id),
"contacts": await get_user_contacts(db, tenant_id, user_id),
"companies": await get_user_companies(db, tenant_id, user_id),
"tasks": await get_user_tasks(db, tenant_id, user_id),
"calendar_events": await get_user_events(db, tenant_id, user_id),
"mails": await get_user_mails(db, tenant_id, user_id),
"files": await get_user_files(db, tenant_id, user_id),
"audit_logs": await get_user_audit_logs(db, tenant_id, user_id),
"agent_runs": await get_user_agent_runs(db, tenant_id, user_id),
"workflow_instances": await get_user_workflows(db, tenant_id, user_id),
"comm_messages": await get_user_comm_messages(db, tenant_id, user_id),
"knowledge_relationships": await get_user_relationships(db, tenant_id, user_id),
"approvals": await get_user_approvals(db, tenant_id, user_id),
}
return data
async def export_to_zip(data: dict, output_path: str) -> str:
"""Export data as ZIP with JSON files."""
```
2. **DSGVO Export API:** `app/routes/dsgvo.py` (NEU)
- `POST /api/v1/dsgvo/export` — Trigger export (async job)
- `GET /api/v1/dsgvo/export/{job_id}` — Get export status
- `GET /api/v1/dsgvo/export/{job_id}/download` — Download ZIP
- `DELETE /api/v1/dsgvo/export/{job_id}` — Delete export file
3. **DSGVO Delete (Right to be forgotten):**
- `POST /api/v1/dsgvo/delete` — Anonymize user data (soft-delete + PII removal)
- `POST /api/v1/dsgvo/delete/{user_id}` — Admin: delete specific user data
**Verbindungen:**
- `dsgvo_export.py` → alle Plugin Models (echte SQL queries)
- Export job via ARQ
- Sensitive data excluded via `SENSITIVE_FIELDS`
**Tests:**
- `tests/test_dsgvo_export.py` — Export completeness, sensitive data exclusion
- `tests/test_dsgvo_delete.py` — Anonymization correctness
**Frontend:**
- `frontend/src/pages/DSGVOExport.tsx` (NEU) — Export/Download/Delete UI
---
### I-17 bis I-20: Onboarding — Setup Wizard
**Basis:** `app/config.py`, `app/models/system_settings.py`, alle Plugin on_activate
**Was neu gebaut wird:**
1. **Onboarding Service:** `app/services/onboarding.py` (NEU, ~300 Zeilen)
```python
async def get_onboarding_status(db, tenant_id, user_id) -> dict:
"""Check which onboarding steps are completed."""
steps = [
{"id": "profile", "label": "Complete your profile", "done": await is_profile_complete(db, user_id)},
{"id": "import_contacts", "label": "Import contacts", "done": await has_contacts(db, tenant_id)},
{"id": "mail_account", "label": "Connect mail account", "done": await has_mail_account(db, tenant_id, user_id)},
{"id": "first_agent", "label": "Create your first agent", "done": await has_agents(db, tenant_id)},
{"id": "first_workflow", "label": "Create a workflow", "done": await has_workflows(db, tenant_id)},
{"id": "knowledge_extraction", "label": "Run knowledge extraction", "done": await has_knowledge(db, tenant_id)},
]
return {"steps": steps, "completion": sum(s["done"] for s in steps) / len(steps)}
async def complete_step(db, tenant_id, user_id, step_id, data) -> dict:
"""Mark onboarding step as complete."""
```
2. **Onboarding API:** `app/routes/onboarding.py` (NEU)
- `GET /api/v1/onboarding/status` — Get onboarding status
- `POST /api/v1/onboarding/steps/{step_id}/complete` — Complete step
- `POST /api/v1/onboarding/skip` — Skip onboarding
3. **Onboarding Model:** `app/models/onboarding.py` (NEU)
```python
class OnboardingProgress(Base, TenantMixin):
__tablename__ = "onboarding_progress"
id: UUID PK
user_id: UUID FK → users.id
step_id: str
completed_at: TIMESTAMPTZ | None
skipped: bool = False
metadata_: dict = JSONB
```
**Verbindungen:**
- `onboarding.py` → Contact, MailAccount, AgentDefinition, Workflow models
- System settings for onboarding configuration
**Migrationen:**
- `0138_onboarding_progress.py` — onboarding_progress table
**Tests:**
- `tests/test_onboarding.py` — Status check, step completion, skip
**Frontend:**
- `frontend/src/pages/SetupWizard.tsx` (NEU) — Multi-step wizard
- `frontend/src/components/onboarding/StepCard.tsx` (NEU)
- `frontend/src/components/onboarding/ProgressBar.tsx` (NEU)
---
### I-21 bis I-23: Performance Optimization
**Basis:** Alle Routes und Services, `app/core/redis.py`
**Was neu gebaut wird:**
1. **Query Optimization:**
- N+1 query detection und batch loading mit `selectinload()` / `joinedload()`
- Pagination auf alle List-Endpoints (bereits teilweise vorhanden via `app/core/pagination.py`)
- Redis caching für häufige Queries
2. **Cache Service:** `app/core/cache.py` erweitern
```python
async def cached_query(key: str, ttl: int, query_func, *args, **kwargs):
"""Cache query result in Redis."""
redis = await get_redis()
cached = await redis.get(key)
if cached:
return json.loads(cached)
result = await query_func(*args, **kwargs)
await redis.setex(key, ttl, json.dumps(result))
return result
```
3. **Index Optimization:**
- Audit aller DB-Indizes via `scripts/check_indexes.py`
- Fehlende Indizes identifizieren und via Migration hinzufügen
- Unused Indizes entfernen
4. **Frontend Performance:**
- Code splitting: Lazy-load plugin pages
- React Query: staleTime, cacheTime optimization
- Bundle size analysis
**Migrationen:**
- `0139_performance_indexes.py` — Additional indexes based on query analysis
**Tests:**
- `tests/test_performance.py` — Query performance benchmarks
- `tests/test_cache.py` — Cache hit/miss, TTL expiry
**Frontend:**
- `frontend/src/components/common/LazyPage.tsx` (NEU) — Lazy loading wrapper
---
### I-24 bis I-25: Final Polish
**Basis:** Alle Systeme
**Was gemacht wird:**
1. **Error Handling Polish:**
- Alle Routes nutzen `build_error_response()` aus `error_codes.py`
- Frontend Error Boundaries auf allen Plugin-Seiten
- Partial-Failure-Semantik für Batch-Operationen
2. **Documentation Update:**
- `docs/api-documentation.md` — Alle neuen Endpoints
- `docs/plugin-development-guide.md` — Knowledge plugin, MCP plugin
- `README.md` — Updated features list
- `docs/test-strategy.md` — Updated test coverage
3. **PWA Re-activation:**
- Service Worker in `frontend/src/main.tsx` re-aktivieren (aktuell deaktiviert)
- Offline-first für kritische Views
4. **Accessibility Audit:**
- ARIA attributes auf allen interaktiven Elementen
- 44px touch targets
- Keyboard navigation
**Tests:**
- `tests/test_error_handling.py` — Verify structured error responses
- Frontend: `frontend/src/__tests__/accessibility.test.tsx`
**Frontend:**
- Error Boundaries: `frontend/src/components/common/ErrorBoundary.tsx` (NEU)
- PWA: `frontend/src/sw.ts` (NEU/überarbeitet)
---
## Phase J
### J-1: Controlled Self-Improvement — auf echte DB-Tabellen + Services aufbauen
**Basis:** `app/models/audit.py` (AuditLog), `app/plugins/builtins/automation/models.py` (AgentDefinition, AgentRun), `app/models/workflow.py` (WorkflowInstance)
**Was neu gebaut wird:**
1. **Self-Improvement Plugin:** `app/plugins/builtins/self_improvement/` (NEU, komplettes Plugin)
- `plugin.py` — `SelfImprovementPlugin(BasePlugin)` mit Manifest
- `models.py` — ImprovementSignal, ImprovementProposal, ProposalVersion, EvaluationResult
- `services.py` — Signal collection, proposal generation, evaluation
- `routes.py` — `/api/v1/improvements`
- `schemas.py` — Pydantic schemas
- `jobs.py` — ARQ jobs for pattern detection and evaluation
2. **ImprovementSignal Model:**
```python
class ImprovementSignal(Base, TenantMixin):
__tablename__ = "improvement_signals"
id: UUID PK
signal_type: str # "repeated_error", "low_success_rate", "slow_workflow", "manual_feedback"
source_type: str # "agent_run", "workflow_instance", "audit_log", "user_feedback"
source_id: UUID
severity: str # "low", "medium", "high"
description: str
metadata_: dict = JSONB
detected_at: TIMESTAMPTZ
status: str = "new" # "new", "analyzed", "proposal_created", "dismissed"
```
3. **ImprovementProposal Model:**
```python
class ImprovementProposal(Base, TenantMixin, OwnedMixin):
__tablename__ = "improvement_proposals"
id: UUID PK
signal_id: UUID FK → improvement_signals.id
title: str
description: str
proposed_changes: dict = JSONB # What should change
current_version: int = 1
status: str = "draft" # "draft", "in_review", "approved", "rejected", "implemented", "evaluated"
risk_assessment: dict = JSONB
expected_impact: str # "low", "medium", "high"
created_at: TIMESTAMPTZ
updated_at: TIMESTAMPTZ
```
4. **ProposalVersion Model:**
```python
class ProposalVersion(Base, TenantMixin):
__tablename__ = "improvement_proposal_versions"
id: UUID PK
proposal_id: UUID FK → improvement_proposals.id
version: int
content: dict = JSONB # Full proposal content at this version
created_by: UUID
created_at: TIMESTAMPTZ
```
5. **EvaluationResult Model:**
```python
class EvaluationResult(Base, TenantMixin):
__tablename__ = "improvement_evaluation_results"
id: UUID PK
proposal_id: UUID FK → improvement_proposals.id
metric_type: str # "success_rate", "error_count", "execution_time", "user_satisfaction"
before_value: float
after_value: float
improvement_pct: float
evaluated_at: TIMESTAMPTZ
metadata_: dict = JSONB
```
**Verbindungen:**
- Signal collection → AuditLog, AgentRun, WorkflowInstance queries
- Proposal → Approval system (`app/core/approval.py`)
- Evaluation → DB queries before/after implementation
**Migrationen:**
- `0140_self_improvement_tables.py` — improvement_signals, improvement_proposals, improvement_proposal_versions, improvement_evaluation_results
**Tests:**
- `tests/test_self_improvement.py` — Signal creation, proposal lifecycle, evaluation
**Frontend:**
- `frontend/src/pages/ImprovementCenter.tsx` (NEU) — Overview dashboard
- `frontend/src/components/improvement/ProposalCard.tsx` (NEU)
- `frontend/src/components/improvement/SignalList.tsx` (NEU)
---
### J-2: Improvement Signals — auf AuditLog + AgentRun + WorkflowInstance aufbauen
**Basis:** `app/models/audit.py`, `app/plugins/builtins/automation/models.py`, `app/models/workflow.py`
**Was neu gebaut wird:**
1. **Signal Detector:** `app/plugins/builtins/self_improvement/signal_detector.py` (NEU, ~300 Zeilen)
```python
async def detect_repeated_errors(db, tenant_id, time_window_hours=24) -> list[dict]:
"""Detect repeated error patterns in audit logs."""
# Query audit_log for error entries, group by error_type + entity_type
# If count > threshold, create ImprovementSignal
async def detect_low_agent_success_rate(db, tenant_id) -> list[dict]:
"""Detect agents with low success rates."""
# Query agent_runs, calculate success rate per agent
# If below threshold, create signal
async def detect_slow_workflows(db, tenant_id) -> list[dict]:
"""Detect workflows with long execution times."""
# Query workflow_instances, calculate avg duration per workflow
# If above threshold, create signal
async def detect_manual_feedback(db, tenant_id) -> list[dict]:
"""Collect manual user feedback signals."""
# Query feedback entries (if feedback system exists)
```
2. **Signal Collection Job:** `app/plugins/builtins/self_improvement/jobs.py`
- ARQ job: `async def collect_signals(ctx, tenant_id)`
- Cron: stündlich
- Calls all detect_* functions
**Verbindungen:**
- `signal_detector.py` → AuditLog, AgentRun, WorkflowInstance models (echte SQL queries)
- Cron job via automation plugin
**Tests:**
- `tests/test_signal_detection.py` — Create errors, run detection, verify signals
---
### J-3: Pattern Detection — auf echten Signalen aufbauen
**Basis:** `app/plugins/builtins/self_improvement/models.py` (ImprovementSignal)
**Was neu gebaut wird:**
1. **Pattern Detector:** `app/plugins/builtins/self_improvement/pattern_detector.py` (NEU, ~250 Zeilen)
```python
async def detect_patterns(db, tenant_id, signal_ids: list[uuid.UUID]) -> list[dict]:
"""Analyze signals and detect patterns."""
signals = await get_signals(db, tenant_id, signal_ids)
# Group signals by:
# - Entity type (all errors on contacts)
# - Agent (all failures from one agent)
# - Workflow (all slow workflows)
# - Time clustering (errors spike at certain times)
patterns = []
# 1. Frequency analysis
# 2. Correlation analysis
# 3. Time-series analysis
return patterns
```
2. **LLM-assisted Pattern Analysis:**
```python
async def llm_analyze_patterns(patterns: list[dict]) -> list[dict]:
"""Use LLM to suggest root causes and improvements."""
from app.ai.llm_client import llm_complete
prompt = f"Analyze these patterns and suggest root causes: {json.dumps(patterns)}"
response = await llm_complete(messages=[{"role": "user", "content": prompt}])
return parse_llm_analysis(response)
```
**Verbindungen:**
- `pattern_detector.py` → ImprovementSignal queries → `llm_client.llm_complete()`
**Tests:**
- `tests/test_pattern_detection.py` — Create signals, detect patterns
---
### J-4: ImprovementProposal — echte SQLAlchemy Modelle + Migration
**Basis:** `app/plugins/builtins/self_improvement/models.py` (aus J-1)
**Was neu gebaut wird:**
1. **Proposal Generator:** `app/plugins/builtins/self_improvement/proposal_generator.py` (NEU, ~250 Zeilen)
```python
async def generate_proposal(db, tenant_id, signal_id, pattern_analysis) -> dict:
"""Generate improvement proposal from signal and pattern analysis."""
from app.ai.llm_client import llm_complete
prompt = f"""Based on the following signal and pattern analysis,
generate a concrete improvement proposal.
Signal: {signal_description}
Pattern: {pattern_analysis}
Return JSON with: title, description, proposed_changes, risk_assessment, expected_impact
"""
response = await llm_complete(...)
proposal_data = parse_llm_proposal(response)
# Create proposal in DB
proposal = await create_proposal(db, tenant_id, signal_id, proposal_data)
return proposal
```
2. **Proposal API:** `app/plugins/builtins/self_improvement/routes.py`
- `GET /api/v1/improvements/proposals` — List proposals
- `GET /api/v1/improvements/proposals/{id}` — Get proposal
- `POST /api/v1/improvements/proposals` — Create proposal
- `PUT /api/v1/improvements/proposals/{id}` — Update
- `POST /api/v1/improvements/proposals/{id}/submit` — Submit for review
**Tests:**
- `tests/test_proposal_generation.py` — Signal → proposal generation
---
### J-5: Versioned Drafts — auf bestehende Versionierung aufbauen
**Basis:** `app/plugins/builtins/wiki/models.py` (WikiArticleVersion), `app/plugins/builtins/self_improvement/models.py` (ProposalVersion)
**Was neu gebaut wird:**
1. **Version Service:** `app/plugins/builtins/self_improvement/version_service.py` (NEU, ~200 Zeilen)
```python
async def create_version(db, tenant_id, proposal_id, content, user_id) -> dict:
"""Create a new version of a proposal."""
current_version = await get_current_version(db, tenant_id, proposal_id)
new_version = current_version + 1
version = ProposalVersion(proposal_id=proposal_id, version=new_version, content=content, ...)
db.add(version)
await db.commit()
# Update proposal.current_version
return version
async def get_version_history(db, tenant_id, proposal_id) -> list[dict]: ...
async def restore_version(db, tenant_id, proposal_id, version_id) -> dict: ...
async def diff_versions(db, tenant_id, proposal_id, v1, v2) -> dict: ...
```
2. **Version API:** `app/plugins/builtins/self_improvement/routes.py` erweitern
- `GET /api/v1/improvements/proposals/{id}/versions` — Version history
- `GET /api/v1/improvements/proposals/{id}/versions/{v}` — Get specific version
- `POST /api/v1/improvements/proposals/{id}/versions/{v}/restore` — Restore
- `GET /api/v1/improvements/proposals/{id}/diff?v1=1&v2=2` — Diff
**Tests:**
- `tests/test_proposal_versioning.py` — Create, restore, diff versions
---
### J-6: Evaluation/Sandbox — auf Test-DB aufbauen
**Basis:** `tests/conftest.py` (test DB setup), `app/plugins/builtins/self_improvement/models.py`
**Was neu gebaut wird:**
1. **Sandbox Evaluator:** `app/plugins/builtins/self_improvement/sandbox.py` (NEU, ~300 Zeilen)
```python
async def evaluate_proposal_in_sandbox(
proposal_id: uuid.UUID, tenant_id: uuid.UUID,
) -> dict:
"""Evaluate proposal in isolated sandbox environment."""
# 1. Create sandbox DB transaction (savepoint)
# 2. Apply proposed changes
# 3. Run test scenarios
# 4. Measure metrics (success rate, error count, execution time)
# 5. Rollback transaction
# 6. Return before/after comparison
```
2. **Test Scenario Runner:** `app/plugins/builtins/self_improvement/test_runner.py` (NEU, ~200 Zeilen)
```python
async def run_test_scenarios(db, tenant_id, scenarios: list[dict]) -> dict:
"""Run predefined test scenarios and collect metrics."""
results = []
for scenario in scenarios:
result = await execute_scenario(db, tenant_id, scenario)
results.append(result)
return aggregate_results(results)
```
3. **Evaluation API:** `app/plugins/builtins/self_improvement/routes.py` erweitern
- `POST /api/v1/improvements/proposals/{id}/evaluate` — Run sandbox evaluation
- `GET /api/v1/improvements/proposals/{id}/evaluation` — Get evaluation results
**Tests:**
- `tests/test_sandbox_evaluation.py` — Evaluate proposal, verify metrics
---
### J-7: Human Approval — auf approval.py aufbauen
**Basis:** `app/core/approval.py` (ApprovalRequest, create_approval_request)
**Was neu gebaut wird:**
1. **Proposal Approval Integration:** `app/plugins/builtins/self_improvement/approval_handler.py` (NEU, ~200 Zeilen)
```python
async def request_proposal_approval(
db, tenant_id, proposal_id, requested_by, user_id,
) -> dict:
"""Create approval request for improvement proposal."""
from app.core.approval import create_approval_request
approval = await create_approval_request(
db=db, tenant_id=tenant_id,
entity_type="improvement_proposal",
entity_id=proposal_id,
action="implement_proposal",
requested_by=requested_by,
requested_by_type="system",
approver_id=user_id,
metadata={"proposal_id": str(proposal_id), "risk": proposal.risk_assessment},
)
# Post to workstream
await post_workflow_approval_request(...)
return approval
async def handle_approval_resolved(db, tenant_id, approval_id, status, user_id):
"""Handle approval resolution — implement or reject proposal."""
if status == "approved":
await implement_proposal(db, tenant_id, proposal_id)
else:
await reject_proposal(db, tenant_id, proposal_id)
```
2. **Approval Hook:** Register callback for `improvement_proposal` entity type
- When ApprovalRequest resolved → `handle_approval_resolved()`
**Verbindungen:**
- `approval_handler.py` → `app.core.approval.create_approval_request()`
- Approval resolution → `implement_proposal()` or `reject_proposal()`
- Workstream notification via `post_workflow_approval_request()`
**Tests:**
- `tests/test_proposal_approval.py` — Request approval, approve, implement
---
### J-8: Impact Measurement — auf echte DB-Queries aufbauen
**Basis:** `app/plugins/builtins/self_improvement/models.py` (EvaluationResult), `app/models/audit.py`
**Was neu gebaut wird:**
1. **Impact Measurement Service:** `app/plugins/builtins/self_improvement/impact_measurement.py` (NEU, ~250 Zeilen)
```python
async def measure_impact_before(db, tenant_id, proposal_id) -> dict:
"""Measure metrics before proposal implementation."""
return {
"error_count": await count_errors(db, tenant_id, since=proposal.created_at),
"agent_success_rate": await calculate_agent_success_rate(db, tenant_id),
"workflow_avg_duration": await calculate_avg_workflow_duration(db, tenant_id),
"user_satisfaction": await get_satisfaction_score(db, tenant_id),
}
async def measure_impact_after(db, tenant_id, proposal_id, days=7) -> dict:
"""Measure metrics after proposal implementation."""
# Same metrics, but after implementation date
async def calculate_improvement(before: dict, after: dict) -> dict:
"""Calculate improvement percentages."""
improvements = {}
for key in before:
if before[key] != 0:
improvements[key] = ((after[key] - before[key]) / before[key]) * 100
else:
improvements[key] = 0.0
return improvements
```
2. **Impact API:** `app/plugins/builtins/self_improvement/routes.py` erweitern
- `GET /api/v1/improvements/proposals/{id}/impact` — Get impact measurement
- `POST /api/v1/improvements/proposals/{id}/measure-impact` — Trigger measurement
**Verbindungen:**
- `impact_measurement.py` → AuditLog, AgentRun, WorkflowInstance queries (echte SQL)
- Results stored in EvaluationResult model
**Tests:**
- `tests/test_impact_measurement.py` — Before/after metrics, improvement calculation
---
### J-9: Pattern Insight Frontend
**Frontend:**
- `frontend/src/components/improvement/PatternInsight.tsx` (NEU) — Pattern visualization
- `frontend/src/components/improvement/ImpactChart.tsx` (NEU) — Before/after comparison
- `frontend/src/components/improvement/SignalTimeline.tsx` (NEU) — Signal timeline
---
### J-10: Self-Improvement Documentation
- `docs/self-improvement-guide.md` (NEU) — How the self-improvement system works
- `docs/api-documentation.md` — Update with improvement endpoints
- `PROGRESS.md` — Update with J-phase status
---
## Phase K
### K-1: EU Compliance Finalization — AI Use-Case Registry
**Basis:** `app/ai/ai_use_case.py` (existing), `app/ai/data_policy.py`, `app/ai/transparency.py`
**Was neu gebaut wird:**
1. **AI Use-Case Registry erweitern:** `app/ai/ai_use_case.py`
- Vollständige Use-Case-Registration für alle AI-Features
- Pflichtfelder: intended_purpose, owner, agents, models, provider, data_categories, allowed_actions, human_oversight_policy, risk_class
- API: `GET /api/v1/compliance/ai-use-cases`, `POST /api/v1/compliance/ai-use-cases`
2. **Compliance Dashboard Service:** `app/services/compliance_service.py` (NEU, ~300 Zeilen)
```python
async def get_compliance_status(db, tenant_id) -> dict:
return {
"ai_use_cases": await get_all_use_cases(db, tenant_id),
"data_processing_activities": await get_processing_activities(db, tenant_id),
"retention_policies": await get_retention_policies(db, tenant_id),
"data_subject_requests": await get_dsr_status(db, tenant_id),
"risk_assessments": await get_risk_assessments(db, tenant_id),
}
```
**Tests:**
- `tests/test_compliance_ai_use_cases.py` — Use-case registration, validation
**Frontend:**
- `frontend/src/pages/ComplianceDashboard.tsx` (NEU)
---
### K-2: Data Processing Registry
**Basis:** `app/models/audit.py`, `app/core/sensitive_data.py`
**Was neu gebaut wird:**
1. **Processing Activity Model:** `app/models/processing_activity.py` (NEU)
```python
class ProcessingActivity(Base, TenantMixin):
__tablename__ = "processing_activities"
id: UUID PK
name: str
purpose: str
legal_basis: str # DSGVO Art. 6 basis
data_categories: list[str] = JSONB
recipients: list[str] = JSONB
retention_period_days: int | None
dpia_required: bool = False
dpia_completed: bool = False
is_active: bool = True
```
2. **Processing Activity API:** `app/routes/compliance.py` (NEU)
- CRUD endpoints for processing activities
**Migrationen:**
- `0141_processing_activities.py`
**Tests:**
- `tests/test_processing_activities.py`
---
### K-3: Data Subject Request (DSR) Automation
**Basis:** `app/services/dsgvo_export.py` (aus Phase I)
**Was neu gebaut wird:**
1. **DSR Model:** `app/models/data_subject_request.py` (NEU)
```python
class DataSubjectRequest(Base, TenantMixin):
__tablename__ = "data_subject_requests"
id: UUID PK
request_type: str # "access", "rectification", "erasure", "portability", "restriction", "objection"
requested_by: UUID FK → users.id
status: str = "new" # "new", "processing", "completed", "rejected"
requested_at: TIMESTAMPTZ
completed_at: TIMESTAMPTZ | None
result_data: dict = JSONB # Export data or action result
metadata_: dict = JSONB
```
2. **DSR Service:** `app/services/dsr_service.py` (NEU, ~250 Zeilen)
- `create_request()`, `process_request()`, `complete_request()`
- Access: triggers DSGVO export
- Erasure: triggers anonymization
- Rectification: triggers data update workflow
3. **DSR API:** `app/routes/compliance.py` erweitern
- `POST /api/v1/compliance/dsr` — Create request
- `GET /api/v1/compliance/dsr/{id}` — Status
- `GET /api/v1/compliance/dsr` — List requests
**Migrationen:**
- `0142_data_subject_requests.py`
**Tests:**
- `tests/test_dsr.py` — All request types, processing flow
**Frontend:**
- `frontend/src/pages/DataSubjectRequests.tsx` (NEU)
---
### K-4: DPIA (Data Protection Impact Assessment)
**Basis:** `app/models/processing_activity.py` (aus K-2)
**Was neu gebaut wird:**
1. **DPIA Model:** `app/models/dpia.py` (NEU)
```python
class DPIA(Base, TenantMixin):
__tablename__ = "dpias"
id: UUID PK
processing_activity_id: UUID FK → processing_activities.id
risk_level: str # "low", "medium", "high"
assessment: dict = JSONB
mitigation_measures: list[dict] = JSONB
approved_by: UUID | None
approved_at: TIMESTAMPTZ | None
status: str = "draft"
```
2. **DPIA API:** `app/routes/compliance.py` erweitern
- CRUD for DPIAs
**Migrationen:**
- `0143_dpias.py`
**Tests:**
- `tests/test_dpia.py`
---
### K-5: AI Act Compliance
**Basis:** `app/ai/ai_use_case.py`, `app/ai/transparency.py`, `app/ai/oversight.py`
**Was neu gebaut wird:**
1. **AI Act Risk Classification:** `app/ai/ai_act_compliance.py` (NEU, ~200 Zeilen)
```python
class AIActRiskClass:
MINIMAL = "minimal"
LIMITED = "limited"
HIGH = "high"
UNACCEPTABLE = "unacceptable"
async def classify_ai_use_case(use_case) -> str:
"""Classify AI use case according to EU AI Act risk levels."""
# Based on: purpose, data categories, automation level, human oversight
async def get_ai_act_requirements(risk_class: str) -> dict:
"""Get required compliance measures for risk class."""
```
2. **Transparency Requirements:**
- AI-generierte Inhalte müssen gekennzeichnet sein (transparency.py bereits vorhanden)
- AI-Akteure im Workstream als AI gekennzeichnet (kommunikation sender_type="ai")
- Deepfake-Kennzeichnung für AI-generierte Medien
3. **Human Oversight Requirements:**
- High-Risk AI-Use-Cases MÜSSEN Human Approval haben (approval.py bereits vorhanden)
- Oversight-Records für alle High-Risk-Entscheidungen (oversight.py bereits vorhanden)
**Tests:**
- `tests/test_ai_act_compliance.py` — Risk classification, requirements
---
### K-6: Compliance Documentation & Audit Trail
**Basis:** `app/models/audit.py`, `app/services/compliance_service.py`
**Was neu gebaut wird:**
1. **Compliance Audit Trail:**
- Alle Compliance-relevanten Aktionen werden im AuditLog protokolliert
- DSR requests, DPIA approvals, AI use-case changes, data exports
2. **Compliance Report Generator:** `app/services/compliance_report.py` (NEU, ~200 Zeilen)
```python
async def generate_compliance_report(db, tenant_id, period_start, period_end) -> dict:
"""Generate comprehensive compliance report for a period."""
return {
"period": {"start": period_start, "end": period_end},
"ai_use_cases": await get_use_cases_summary(db, tenant_id, period_start, period_end),
"data_exports": await get_exports_summary(db, tenant_id, period_start, period_end),
"dsr_requests": await get_dsr_summary(db, tenant_id, period_start, period_end),
"audit_trail": await get_audit_summary(db, tenant_id, period_start, period_end),
"retention_actions": await get_retention_summary(db, tenant_id, period_start, period_end),
}
```
3. **Compliance Report API:** `app/routes/compliance.py` erweitern
- `GET /api/v1/compliance/report?start=...&end=...` — Generate report
- `POST /api/v1/compliance/report/export` — Export as PDF/JSON
**Tests:**
- `tests/test_compliance_report.py` — Report generation, completeness
**Frontend:**
- `frontend/src/pages/ComplianceReport.tsx` (NEU)
---
## Migrations-Übersicht
| Migration | Phase | Beschreibung |
|-----------|-------|--------------|
| 0129 | B | IVFFlat index strategy config |
| 0130 | B | storage_provider_configs table |
| 0131 | B | Migrate notifications to comm + drop notification tables |
| 0132 | H | knowledge_sources table |
| 0133 | H | entity_relationships: confidence, review_status, reviewed_by |
| 0134 | H | knowledge_evidence table |
| 0135 | H | wiki_articles: embedding vector(768) column |
| 0136 | H | derived_data_dependencies table |
| 0137 | H | knowledge_retention_policies table |
| 0138 | I | onboarding_progress table |
| 0139 | I | Performance indexes |
| 0140 | J | Self-improvement tables (signals, proposals, versions, evaluations) |
| 0141 | K | processing_activities table |
| 0142 | K | data_subject_requests table |
| 0143 | K | dpias table |
## Neue Plugins
| Plugin | Phase | Dependencies |
|--------|-------|-------------|
| storage_webdav | B | [] |
| storage_nextcloud | B | [storage_webdav] |
| knowledge | H | [unified_search, graph_rag] |
| mcp | I | [] |
| self_improvement | J | [automation] |
## Neue Frontend-Seiten
| Seite | Phase | Pfad |
|-------|-------|------|
| StorageSettings | B | /settings/storage |
| KnowledgeDashboard | H | /knowledge |
| KnowledgeReviewQueue | H | /knowledge/review |
| HumanAIWorkstream | I | /workstream |
| Dashboard | I | /dashboard |
| DSGVOExport | I | /settings/dsgvo |
| SetupWizard | I | /onboarding |
| ImprovementCenter | J | /improvements |
| ComplianceDashboard | K | /compliance |
| DataSubjectRequests | K | /compliance/dsr |
| ComplianceReport | K | /compliance/report |
## Neue CommMessageBlock Types
| Block Type | Phase | Verwendung |
|-----------|-------|-----------|
| agent_result | F | Agent run results |
| approval_request | F | Approval requests in workstream |
| task_card | F | Task references in workstream |
| workflow_card | F | Workflow references in workstream |
| knowledge_card | F | Knowledge entity references |
| progress_card | F | Progress indicators |
| ai_proposal | I | AI proposals for human review |
| human_decision | I | Human decision records |
| ai_explanation | I | AI explanations with evidence |
## Verbindungs-Matrix (Wichtigste)
| Von | Nach | Mechanismus |
|-----|------|-----------|
| Agent Runner | Kommunikation | `kommunikation.contracts.send_message()` |
| Agent Runner | Workstream | `agent_workstream.py` → `create_plugin_room()` |
| Workflow Engine | Kommunikation | `workstream.py` → `send_message()` |
| Knowledge Plugin | Unified Search | `get_search_registry().get()` → `get_embedding_text()` |
| Knowledge Plugin | GraphRAG | `graph_rag.services.create_relationship()` |
| Knowledge Plugin | LLM Client | `llm_complete()` für extraction |
| Knowledge Plugin | Event Bus | `event_bus.subscribe()` für event-driven extraction |
| Wiki Plugin | Unified Search | `get_search_registry().register(WikiSearchProvider())` |
| Wiki Plugin | LLM Client | `llm_embed()` für embeddings |
| Self-Improvement | Approval | `create_approval_request()` für proposal approval |
| Self-Improvement | AuditLog | SQL queries für signal detection |
| MCP Plugin | CRM Routes | Read-only CRM operations als MCP tools |
| Storage Plugins | Storage Registry | `get_storage_registry().register()` |
| Automation Plugin | Prebuilt Agents | `create_*_agent()` in `on_activate()` |
---
## Implementierungs-Reihenfolge
1. **Phase B Lücken** (3 Tasks) — Fundament: Storage, Index, Notification cleanup
2. **Phase F Lücken** (3 Tasks) — Agent Integration: Prebuilt registration, Agent→Comm, Workstream
3. **Phase G Lücke** (1 Task) — Workflow Workstream
4. **Phase H Rest** (12 Tasks) — Knowledge System auf graph_rag + unified_search
5. **Phase I** (25 Tasks) — Cross-System Integration, Human-AI, MiniApps, Dashboard, DSGVO, Onboarding
6. **Phase J** (10 Tasks) — Self-Improvement auf echte DB + AuditLog
7. **Phase K** (6 Tasks) — EU Compliance Finalization
**Total: 60 Tasks**
---
> Dieser Plan ist so detailliert, dass direkt mit der Implementierung begonnen werden kann. Jeder Task hat konkrete Dateipfade, Model-Definitionen, Route-Definitionen, und Verbindungs-Punkte zu vorhandenen Systemen.