Bahn: aisupport, Analyse-O2C-C2S, awesome-bahn-mcp-servers, beam-mcp,
Confluence_Bot, db-planet-mcp-server, O2C-Harness, project-audit,
Projekt-KIQ-HP, teamlandkarte-mcp
Dhive: Jury-Voting
Privat: CV, NoteGraph (NOTE: NoteGraph needs complete redo after consolidation)
Shared: AI-Orchestrator, OrgMyLife, power_skills_and_more
Shared/references: symphony (read-only)
Bahn repos remain available as independent remotes - this monorepo
pulls them in via subtree, the originals are untouched.
24 KiB
Design: Capacity Matching MCP Server
Note
: For comprehensive architecture documentation including ADRs, quality requirements, and deployment view, see architecture.md.
Architecture Overview
Update (2026-02): The server uses Azure OpenAI for role/requirements extraction and similarity scoring via embeddings. There are no heuristic fallbacks for these features.
The embedding model is
text-embedding-3-largewith 3072 dimensions and results are cached in a local SQLite embedding cache.
Embedding similarity acceleration (global prefetch + batch embeddings)
Update (2026-02): Embedding similarity is optimized via global bulk prefetch plus true Azure embeddings batch requests. This is primarily implemented in
SimilarityEngineandAzureOpenAIClient.
Motivation
The matching pipeline can be slow if embeddings are requested on-demand from inner
loops. While the system already uses a persistent SQLite embedding cache
(EmbeddingCache), additional optimization is necessary to avoid repeated:
- normalization + cache-key computation
- SQLite reads
- sequential Azure embedding calls for cache misses
Goals
- Bulk prefetch embeddings for a matching run (roles + competences), across all free capacities at once (not person-by-person / incremental embedding).
- Maintain deterministic behavior and strict error semantics.
- Preserve existing similarity contracts (
per_skill/aggregate) and tool outputs.
Data flow
- Collect all texts that will be embedded in the current run:
- required competences
- candidate competences (from all free capacities)
- required roles
- candidate roles (from all free capacities)
- Normalize + deduplicate using existing
_normalize_text()logic:- skip empty/whitespace-only texts
- emit a Python logger warning if all candidate competences normalize to empty
- Resolve embeddings using a layered cache approach:
- in-memory cache (instance-level, attached to
SimilarityEngine; not created/cleared per invocation) - persistent SQLite embedding cache (
EmbeddingCache) - Azure embeddings API (batch only the cache-missing texts; chunked)
- in-memory cache (instance-level, attached to
- Compute cosine similarities locally using the prefetched mapping.
Normalization and cache keys
Cache key behavior remains unchanged to avoid invalidating the on-disk cache:
- normalization: trim + collapse whitespace
- key: SHA256 of
model|dims|normalized(text).lower()
Azure embeddings batch API + chunking
The Azure embeddings API supports embedding multiple inputs per request via
input=[...]. The implementation requirements are:
- preserve stable mapping from input texts → returned vectors
- rely on the
indexfield in the API response if available - otherwise assume stable ordering
- rely on the
- chunk large requests using
azure_openai.embedding_batch_size(default: 128) - strict failure semantics:
- the matching operation fails immediately on error
- log the failing chunk input list via Python logging
- cost tracking: log once per chunk with a count of embeddings requested
Trade-offs
- Batch calls reduce network requests dramatically but increase blast radius: one
failing request affects multiple inputs.
- mitigation: keep retry logic at the chunk level
- Prefetch holds a larger in-memory mapping for the current run.
- mitigation: only store vectors required for that run (after dedup)
Tests
- prefetch/dedup: client is called once per unique normalized text (cache misses only)
- cache layering: in-memory first, then SQLite, then Azure; no duplicate calls
- chunking: multiple Azure requests when
len(missing_texts) > batch_size - empty input handling: warning when all candidate competences normalize to empty
- strict error propagation: batch failures identify/log which chunk failed
(diagram below is historical; “sampling” arrows represent the client’s LLM runtime.
The MCP server itself does not rely on MCP “sampling” for Azure OpenAI calls.)
┌──────────────┐
│ MCP Client │
│ (with LLM) │
└──────┬───────┘
│ MCP protocol (tool calls)
▼
┌───────────────────────────────────────────────────────────────────────┐
│ MCP Server │
│ │
│ Task Management Tools: │
│ - list_open_tasks │
│ - get_task_details (table-first) │
│ - validate_task_requirements (table-first) │
│ - find_capacities_for_task (requires confirm) │
│ - infer_roles │
│ │
│ Requirement Gathering Tools: │
│ - extract_requirements (writes pending) │
│ - update_requirements (writes pending) │
│ - collect_structured_requirement_data │
│ (writes pending) │
│ - start_guided_capture │
│ - guided_set_description │
│ - guided_set_role │
│ - guided_set_time_range (open-ended) │
│ - guided_set_competences (writes pending) │
│ - confirm_requirements │
│ │
│ Search Execution & Refinement: │
│ - find_matching_capacities (hard-gated) │
│ - filter_search_results │
│ - get_results_by_category │
│ │
│ Business Logic: Matcher / Scorer / Caches │
│ │
│ Integrations: │
│ - Trino / Open Data Lake (read-only SQL) │
│ - Azure OpenAI (embeddings only) │
└───────────────────────────────────────────────────────────────────────┘
│
▼
┌────────────────────────────┐
│ Open Data Lake (Trino) │
└────────────────────────────┘
Component Design
1. Configuration Layer (src/teamlandkarte_mcp/config.py)
Purpose: Load and validate configuration from config.toml and environment.
Important:
config.tomlcontains non-secret settings only (endpoints, model/deployment names, cache settings).- credentials come from environment (recommended:
.env) and must not be committed:DATA_LAKE_USERNAMEDATA_LAKE_PASSWORDAZURE_OPENAI_EMBEDDING_API_KEY
config.toml structure (excerpt):
[database]
host = "..."
port = 8446
backend = "trino"
http_scheme = "https"
verify_ssl = true
catalog = "hive"
schema = "tier1_open_lake"
[azure_openai]
endpoint = "https://<resource>.openai.azure.com"
api_version = "2024-02-15-preview"
embedding_deployment = "text-embedding-3-large"
[embedding_cache]
path = ".cache/embeddings.sqlite3"
ttl_days = 30
[matching.similarity]
# `per_skill` = best-match per required competence
# `aggregate` = mean(required) vs mean(candidate)
strategy = "per_skill"
Roundtrip examples (tool-level sequences)
These examples describe the intended end-to-end client/server interaction.
They explicitly include the strict Review → Ask → Confirm step before any
matching execution when matching.require_confirmation = true.
Roundtrip A: Ad-hoc search (underspecified) using guided capture
start_guided_capture()guided_set_description("...")(must capture concrete scope/goal; skill-only is insufficient)guided_set_role("Backend Developer")guided_set_time_range(date_start="2026-04-01", date_end="2026-06-30")guided_set_competences(["Python", "FastAPI", "Docker"])- Review:
show_pending_requirements()(server returns a single review table) - Ask (assistant → user): "Soll ich diese Anforderungen so übernehmen und die Suche starten?" (Ja/Nein)
- Confirm (only on "Ja"):
confirm_requirements(confirm=true) - Execute:
find_matching_capacities(role_name="Backend Developer", competences=[...], date_start="2026-04-01", date_end="2026-06-30")- Server returns deterministic headers (
Using SEARCH_ID=...,SEARCH_ID=...,META=...) plus## Summaryand Top results.
- Server returns deterministic headers (
If the user answers "Nein" in step 7, call confirm_requirements(confirm=false)
and continue capturing/updating requirements.
Roundtrip B: DB task workflow (task_id-based)
list_open_tasks(limit=...)get_task_details(task_id)- Optional helpers:
infer_primary_role(task_id=...)(if role unclear)validate_task_requirements(task_id)(only if the user explicitly requests validation; validation is independent from matching)
- Capture/update requirements as needed (e.g. from task details)
- Review:
show_pending_requirements() - Ask (assistant → user): confirm Yes/No
- Confirm (only on "Yes"):
confirm_requirements(confirm=true) - Execute matching (depending on the client UX):
find_capacities_for_task(task_id)ORfind_matching_capacities(...)
Alternative marker: If the assistant already displayed requirements to the user outside of the server review table, it MAY call
request_requirements_confirmation()instead ofshow_pending_requirements(). The actual confirmation is still performed exclusively viaconfirm_requirements(confirm=true)after user approval.
2. Database Layer (src/teamlandkarte_mcp/database/trino_client.py)
Purpose: Manage Trino/Presto connectivity to the Open Data Lake and execute read-only queries.
Database Schema:
Capacity Tables (existing):
teamlandkarte_v_capacities_latest- Fields:
id,owner_name,role_name,role_level,begin_date,end_date,deletion_reason - Filter:
deletion_reason IS NULL(active capacities only)
- Fields:
teamlandkarte_v_capacity_competences_latest- Fields:
capacity_id(FK),competence_id(FK) - JOIN: Links capacities to competences
- Fields:
teamlandkarte_v_competences_latest- Fields:
id,name - Stores competence/skill names
- Fields:
Task Tables (new for Round 5):
beschaffungstool_kmp_task_latest- Fields:
id,title__c,description__c,startdate__c,enddate__c,createddate,status__c - Filter:
status__c = "Veröffentlicht"(published tasks only) - Note: Role is NOT stored in DB, extracted from description
- Fields:
beschaffungstool_kmp_skill_latest- Fields:
task__c(FK to task.id),skillname__c - JOIN: Links tasks to their required skills
- Multiple rows per task (one per skill)
- Fields:
Note: Trino catalog/schema are configured via
config.toml(catalog/schema). The queries in code use unqualified view names.
Key Classes:
class TrinoClient:
def __init__(self, config: DatabaseConfig):
# Initialize Trino DB-API connection
def get_all_capacities_with_competences(self) -> List[Capacity]:
"""Fetch all active capacities with their competences in a single query"""
query = """
SELECT
cap.id,
cap.owner_name,
cap.role_name,
cap.role_level,
cap.begin_date,
cap.end_date,
comp.name as competence_name
FROM teamlandkarte_v_capacities_latest cap
LEFT JOIN teamlandkarte_v_capacity_competences_latest cap_comp
ON cap.id = cap_comp.capacity_id
LEFT JOIN teamlandkarte_v_competences_latest comp
ON cap_comp.competence_id = comp.id
WHERE cap.deletion_reason IS NULL
ORDER BY cap.id, comp.name
"""
# Execute query with proper parameterization
# Group results by capacity_id in Python
# Return List[Capacity] with competences populated
def get_open_tasks(self, limit: int = 20) -> List[Task]:
"""Fetch latest open tasks with their skills."""
# Apply LIMIT to tasks before joining skills, otherwise one task with
# many skill rows can consume the full LIMIT.
query = """
WITH latest_tasks AS (
SELECT
t.id,
t.title__c,
t.description__c,
t.startdate__c,
t.enddate__c,
t.createddate
FROM beschaffungstool_kmp_task_latest t
WHERE t.status__c = ?
ORDER BY t.createddate DESC
LIMIT ?
)
SELECT
t.id,
t.title__c,
t.description__c,
t.startdate__c,
t.enddate__c,
t.createddate,
s.skillname__c
FROM latest_tasks t
LEFT JOIN beschaffungstool_kmp_skill_latest s
ON t.id = s.task__c
ORDER BY t.createddate DESC
"""
# Execute with limit parameter
# Group results by task_id (multiple rows per task for skills)
# Return List[Task] with skills populated
def get_task_by_id(self, task_id: str) -> Optional[Task]:
"""Fetch a single task with its skills by ID"""
query = """
SELECT
t.id,
t.title__c,
t.description__c,
t.startdate__c,
t.enddate__c,
t.createddate,
s.skillname__c
FROM beschaffungstool_kmp_task_latest t
LEFT JOIN beschaffungstool_kmp_skill_latest s
ON t.id = s.task__c
WHERE t.id = ? AND t.status__c = 'Veröffentlicht'
"""
# Execute with task_id parameter
# Group skill rows into single Task object
# Return Task or None if not found
@dataclass
class Capacity:
id: int
owner_name: str
role_name: str
role_level: str
begin_date: date
end_date: date
competences: List[str] = field(default_factory=list)
@dataclass
class Task:
id: str
title: str
description: str
start_date: Optional[date]
end_date: Optional[date]
created_date: datetime
skills: List[str] = field(default_factory.list)
Safety:
- Connection uses read-only user (enforced by DB permissions)
- Query validation to prevent any INSERT/UPDATE/DELETE
- Parameterized queries to prevent SQL injection
- Single join query reduces round-trips and complexity
Implementation Note: The query returns multiple rows per capacity (one per competence). Python code groups these rows by capacity_id to create the final List[Capacity] structure.
3. Cache Layer (src/cache/)
Purpose: Provide two-tier caching for database queries and search results.
Search session behavior:
- Matching creates a new
search_id(UUID) and stores full results in-memory. - Filtering creates a new
filter_idunder the samesearch_id. - Search-related tools emit a JSON block so clients can reliably capture IDs.
- Tools reject non-UUID
search_idearly to avoid confusing cache-miss flows.
Implementation:
from cachetools import TTLCache
from datetime import timedelta
import json
import uuid
class QueryCache:
"""Database query cache with 12-hour TTL"""
def __init__(self, ttl_hours: int = 12, max_size: int = 100):
self.cache = TTLCache(
maxsize=max_size,
ttl=ttl_hours * 3600
)
self.stats = {"hits": 0, "misses": 0}
def get_or_fetch(self, key: str, fetch_fn: callable) -> Any:
if key in self.cache:
self.stats["hits"] += 1
return self.cache[key]
self.stats["misses"] += 1
result = fetch_fn()
self.cache[key] = result
return result
class SearchCache:
"""Search results cache with 60-minute TTL, keyed by search_id"""
def __init__(self, ttl_minutes: int = 60, max_size: int = 100):
self.cache = TTLCache(
maxsize=max_size,
ttl=ttl_minutes * 60
)
self.stats = {"hits": 0, "misses": 0}
def store_search(self, search_data: dict) -> str:
"""Store search results and return search_id"""
search_id = str(uuid.uuid4())
self.cache[search_id] = json.dumps(search_data)
return search_id
def get_search(self, search_id: str) -> Optional[dict]:
"""Retrieve search results by search_id"""
if search_id in self.cache:
self.stats["hits"] += 1
return json.loads(self.cache[search_id])
self.stats["misses"] += 1
return None
def update_search(self, search_id: str, search_data: dict) -> bool:
"""Update existing search results (e.g., add filters)"""
if search_id in self.cache:
self.cache[search_id] = json.dumps(search_data)
return True
return False
Cache Keys:
QueryCache:
"all_capacities_with_competences"- Complete dataset from database
SearchCache:
<search_id>(UUID) - Stores complete search results and associated filters as JSON:{ "task_id": "T-12345", // NEW: Reference to database task (null for ad-hoc searches) "requirements": { "required_competences": ["Python", "FastAPI"], "preferred_role": "Backend Developer", "candidate_roles": ["Backend Developer", "Software Engineer"], "date_start": "2026-04-01", "date_end": "2026-06-30" }, "results_by_category": { "Top": [...], "Good": [...], "Partial": [...], "Low": [...] }, "filters": { "filter-1": { "criteria": { "role": "Developer", "competences": ["Python", "Docker"], "min_similarity": 0.7 }, "results_by_category": { "Top": [...], "Good": [...], "Partial": [...], "Low": [...] }, "timestamp": "2026-02-11T14:35:00Z" }, "filter-2": { "criteria": { "competences": ["Cloud"], "min_similarity": 0.7 }, "results_by_category": {...}, "timestamp": "2026-02-11T14:40:00Z" } }, "timestamp": "2026-02-11T14:30:00Z" }
Notes:
task_idfield isnullfor ad-hoc searches (viafind_matching_capacities)task_idfield contains database task ID for task-based searches (viafind_capacities_for_task)- This allows tracking the source of each search for audit and debugging purposes
Session State (in-memory, per client session):
"last_requirements"- Stores TaskRequirements object for use withupdate_requirements- Cleared when new search is executed via
find_matching_capacities
4. Matching Engine (src/matching/matcher.py)
Purpose: Analyze tasks and match capacities based on requirements.
Key Components:
Matching Logic
# NOTE (2026-02): This section is historical pseudocode. The current
# implementation uses Azure OpenAI via `AzureOpenAIClient` (chat + embeddings)
# and does not rely on MCP sampling.
class CapacityMatcher:
def __init__(
self,
config: MatchingConfig,
similarity_engine,
):
self.config = config
self.similarity_engine = similarity_engine
async def _match_competences(
self,
capacity_comps: List[str],
required_comps: List[str],
) -> float:
"""Compute similarity via embeddings + cosine similarity."""
details = await self.similarity_engine.compute_competence_similarity(
required=required_comps,
candidate=capacity_comps,
)
if not details:
return 0.0
return sum(d["score"] for d in details.values()) / len(details)
5. Scoring System (src/matching/scorer.py)
Purpose: Categorize match results into Top/Good/Partial/Low.
class ScoreCategory(Enum):
TOP = "Top"
GOOD = "Good"
PARTIAL = "Partial"
LOW = "Low"
class Scorer:
def __init__(self, config: MatchingConfig):
self.thresholds = {
ScoreCategory.TOP: config.top_threshold,
ScoreCategory.GOOD: config.good_threshold,
ScoreCategory.PARTIAL: config.partial_threshold
}
def categorize(self, score: float) -> ScoreCategory:
if score >= self.thresholds[ScoreCategory.TOP]:
return ScoreCategory.TOP
elif score >= self.thresholds[ScoreCategory.GOOD]:
return ScoreCategory.GOOD
elif score >= self.thresholds[ScoreCategory.PARTIAL]:
return ScoreCategory.PARTIAL
else:
return ScoreCategory.LOW
def group_by_category(self, results: List[MatchResult]) -> Dict[ScoreCategory, List[MatchResult]]:
"""Group results by score category"""
grouped = {cat: [] for cat in ScoreCategory}
for result in results:
category = self.categorize(result.overall_score)
result.category = category
grouped[category].append(result)
return grouped
6. MCP Server (src/mcp_server.py)
Purpose: Expose capacity matching as MCP tools.
Workflow Design:
The server provides a multi-step workflow with separate tools for each phase:
-
Task Management Phase (NEW in Round 5):
list_open_tasks()- Browse published tasks from databaseget_task_details()- View task details with extracted rolevalidate_task_requirements()- Compare DB skills vs description-extracted skillsfind_capacities_for_task()- Automated capacity search for database task
-
Requirement Gathering Phase (Ad-hoc):
extract_requirements()- Extract from natural language task descriptioncollect_structured_requirement_data()- Interactive step-by-step collectionupdate_requirements()- Modify previously extracted requirements
-
Search Execution Phase:
find_matching_capacities()- Execute search with confirmed requirements
-
Result Refinement Phase:
filter_search_results()- Apply fuzzy filters to resultsget_results_by_category()- Browse results by score category
Key Design Principles:
- Separation of Concerns: Each tool has a single, clear responsibility
- LLM State Management: The client-side LLM tracks requirement state across tool calls
- Explicit Confirmation: Requirements are returned for user confirmation before search execution
- Filter Chain Tracking: Filters use
filter_idto maintain navigation hierarchy - Database Task Integration: Seamless workflow for published tasks with automatic validation
- Common Extraction Logic: DRY principle with shared
_extract_requirements_from_description()function - Availability as filter: Availability is displayed and can be used as a date-range filter (not part of scoring)
MCP Tools: