119 KiB
Langflow Projekt-Matching Implementation Plan
For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: Fully automated Langflow workflow that converts incoming freelancermap project e-mails into evaluated EspoCRM opportunities (Must-have-Match > 85 %) with notification/alert mails and Langfuse tracing — per spec docs/superpowers/specs/2026-07-09-langflow-projekt-matching-design.md.
Architecture: A pure-Python package projektmatch (host-testable with pytest, deployed by file copy into the Langflow container's /app/langflow mount) holds ALL logic. Langflow custom components are thin wrappers around package functions. Two flows: Flow 1 "PM Ingest" (Webhook → PMIngest → TextOutput) and Flow 2 "PM Projekt bewerten" (TextInput → PMFetch → PMExtract → PMMatch → PMRules → PMCrm → PMNotify → TextOutput), built and uploaded programmatically via the Langflow API. A systemd user timer POSTs to the webhook every 5 minutes.
Tech Stack: Python 3.11 (stdlib imaplib/smtplib/email + requests + beautifulsoup4 — all present in the Langflow 1.10.0 image), pytest on host venv, Langflow API, EspoCRM REST API, Langfuse public API, vLLM OpenAI-compatible API, systemd user units.
Global Constraints
- Local AI ONLY: vLLM at
http://host.containers.internal:8081/v1(from pod) /http://127.0.0.1:8081/v1(from host), modelAxionML/Qwen3.5-9B-NVFP4. It is a REASONING model: never set a smallmax_tokens; leavemax_tokensunset. - Match numbers, table assembly, Misc rules are deterministic Python — never LLM output.
- Gate:
Must-have-Match > 85(strictly greater) → CRM + notification; else rejected (trace only). - Notification subject:
[Projekt-Match] Freelancermap — <Must> % — <Name>. Alert subject:[Projekt-Match-Fehler] Freelancermap — <n> Projekt(e). Flow 1 skips mails whose subject starts with either prefix. - Terminal-state rule: every project ends
created/rejected/duplicate/failedon FIRST attempt; trigger mail then ALWAYS moves to Trash; failures → ONE aggregated alert mail per trigger mail. No cross-cycle retries (only in-run HTTP retries, 2–3 attempts). - Duplicate key: canonical URL (redirects followed, query+fragment stripped) stored in Opportunity field
cProjektlink, matched withequals. - Secrets (the IMAP password from the project brief, EspoCRM API key from
/home/tlg/mkt/bewerb/.secrets/espocrm-api.md) live in Langflow global variables and chmod-600 host files only — never in git, flow JSON, or chat output. - Mail account:
chancen@destengs.com, IMAPmail.destengs.com:993SSL, SMTPmail.destengs.com:587STARTTLS. Design-phase notify/alert recipient:chancen@destengs.com; after gates 1+2:Thomas.Langer@destengs.com(Task 16). - CRM:
https://crm.creature-go.com(API base<base>/api/v1), headerX-Api-Key. Teams: Projekt→DesTEngS, Arbeitnehmer-Angebot→Arbeitnehmer, ANÜ→ANÜ; team IDs looked up by name, NEVER guessed. Agency →cAccount1Id(noaccountId); direct →accountId(nocAccount1Id). - Description markdown format (exact, from SKILL.md Schritt 5): first line
Must-have-Match: X % · Nice-to-have-Match: Y %(empty category →–), blank line, table header| Nr. | Kat. | ❔ | Anforderung |, rows ordered Must → Nice → Misc, numbered from 1. - IMAP custom keyword
$ProjektCheckedmarks inspected non-project mails (left in INBOX untouched). - All git commits in repo
/home/lwc/bin; project code under/home/lwc/bin/projekt-matching/. - Langflow API auth (v1.10.0): flow CRUD + variables need Bearer JWT from
GET /api/v1/auto_login;POST /api/v1/run/...and webhook needx-api-key. Responses are gzip →curl --compressed. - Langflow data dir: host
~/.local/share/langflow_pod/langflow-data= container/app/langflow. Deployed files need group-read (dir isg+rwxs, uid-1000/gid-0 process).
Deviations discovered during execution
response_format json_schemareplacesguided_json(vLLM 0.22 ignoresguided_json).chat_template_kwargsenable_thinking: falseis required for structured calls.- Component input
ctxwas renamedctx_in(Langflow reservesComponent.ctx). PYTHONPATH=/app/langflowenv var replaces the.pth-file approach.- Langfuse org/project provisioning goes through the session tRPC endpoint (the admin API is EE-gated).
- Flow 1 is a single
PMIngestcomponent (not the originally sketched multi-node graph). PMNotifycallsstage_notifydirectly rather than going through an intermediate layer.
Task-level code snippets further below in this document were NOT retro-edited to reflect these deviations; the projektmatch package source under /home/lwc/bin/projekt-matching/ is authoritative.
File Structure
projekt-matching/
├── .gitignore # .venv/, __pycache__/, *.local.env
├── pytest.ini
├── projektmatch/ # pure-python package (host-tested, container-deployed)
│ ├── __init__.py
│ ├── config.py # Cfg dataclass
│ ├── rules.py # Misc rules, match calc, description markdown (Task 2)
│ ├── mailparse.py # mail detection/splitting, canonical URL (Task 3)
│ ├── espocrm.py # EspoCRM client (Task 4)
│ ├── llm.py # vLLM chat_json + prompts + schemas (Task 5)
│ ├── mailer.py # IMAP MailBox + SMTP send_mail (Task 6)
│ ├── pm_trace.py # custom Langfuse trace via ingestion API (Task 7)
│ ├── stages.py # Flow-2 stage functions + guard (Task 7)
│ └── ingest.py # Flow-1 logic: lock, loop, dispatch, alert, trash (Task 8)
├── components/ # Langflow custom components (thin wrappers, Task 9)
│ ├── pm_fetch.py pm_extract.py pm_match.py pm_rules.py
│ ├── pm_crm.py pm_notify.py pm_ingest.py
├── deploy/
│ ├── deploy_files.sh # copy package+components+vorgaben into langflow-data (Task 10)
│ ├── build_flows.py # create variables, API key, flows via Langflow API (Task 10)
│ ├── setup_langfuse.py # org/project/keys/score-configs (Task 11)
│ ├── secrets.local.env # chmod 600, gitignored (Task 10)
│ ├── trigger_webhook.sh # curl POST webhook (Task 12)
│ ├── projekt-matching.service / .timer (Task 12)
├── tests/ # pytest, run on host venv
│ ├── test_rules.py test_mailparse.py test_espocrm.py test_llm.py
│ ├── test_mailer.py test_stages.py test_ingest.py
│ └── e2e/
│ ├── cleanup_crm.py run_flow2.py send_test_mail.py check_state.py
Interfaces between tasks use one shared shape: ctx dict flowing through Flow 2 stages:
{"canonical": str, "title": str, "status": "ok|failed|rejected|created", "error": str?, "page_text": str?, "extract": {...}?, "mustMatch": int|None?, "niceMatch": int|None?, "description": str?, "decision": "consider|rejected"?, "opportunityId": str?, "crmUrl": str?, "projectName": str?}
Task 1: Scaffold, venv, pytest smoke
Files:
- Create:
projekt-matching/.gitignore,projekt-matching/pytest.ini,projekt-matching/projektmatch/__init__.py,projekt-matching/tests/test_smoke.py
Interfaces: Produces the package skeleton and the test command every later task uses: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest.
- Step 1: Create directories and venv
cd /home/lwc/bin
mkdir -p projekt-matching/{projektmatch,components,deploy,tests/e2e}
cd projekt-matching
python3 -m venv .venv
.venv/bin/pip install --quiet pytest requests beautifulsoup4
Expected: exit 0. (requests/bs4 are already inside the Langflow image; the venv is only for host-side tests.)
- Step 2: Write config files
projekt-matching/.gitignore:
.venv/
__pycache__/
*.pyc
*.local.env
projekt-matching/pytest.ini:
[pytest]
testpaths = tests
projekt-matching/projektmatch/__init__.py:
"""projektmatch — automated freelancermap project matching (Langflow backend)."""
projekt-matching/tests/test_smoke.py:
import projektmatch
def test_package_imports():
assert projektmatch.__doc__
- Step 3: Run tests
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest -q
Expected: 1 passed
- Step 4: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): scaffold package, venv, pytest smoke"
Task 2: rules.py — Misc rules, match calculation, description markdown
Files:
- Create:
projekt-matching/projektmatch/rules.py - Test:
projekt-matching/tests/test_rules.py
Interfaces:
-
Produces:
eval_misc(item: dict, today: datetime.date, geocode_fn=geocode) -> str(returns"yes"|"no"|"unknown");calc_match(rows: list[dict], kat: str) -> int|None;build_description(rows: list[dict]) -> str;order_rows(rows) -> list[dict]; constantsYES/NO/UNKNOWN,SYMBOL,RATING_VALUE. Row shape:{"kat": "Must|Nice|Misc", "rating": "yes|no|unknown", "text": str}. Misc item shape (filled by extraction LLM, Task 5):{"miscType": "start|workload|duration|location|security|other", "startDate": "YYYY-MM-DD"|None, "workloadPercent": int|None, "remotePercent": int|None, "onsiteLocation": str|None}. -
Step 1: Write failing tests
projekt-matching/tests/test_rules.py:
import datetime as dt
from projektmatch import rules
TODAY = dt.date(2026, 7, 9)
def geo_munich(place):
return (48.137, 11.575) # ~40 km from Sauerlach
def geo_hamburg(place):
return (53.551, 9.994) # ~600 km
def geo_fail(place):
raise OSError("network down")
def misc(**kw):
base = {"miscType": "other", "startDate": None, "workloadPercent": None,
"remotePercent": None, "onsiteLocation": None}
base.update(kw)
return base
def test_security_is_always_no():
assert rules.eval_misc(misc(miscType="security"), TODAY) == rules.NO
def test_start_within_8_weeks_yes_after_no_missing_unknown():
assert rules.eval_misc(misc(miscType="start", startDate="2026-08-01"), TODAY) == rules.YES
assert rules.eval_misc(misc(miscType="start", startDate="2026-12-01"), TODAY) == rules.NO
assert rules.eval_misc(misc(miscType="start"), TODAY) == rules.UNKNOWN
assert rules.eval_misc(misc(miscType="start", startDate="asap"), TODAY) == rules.UNKNOWN
def test_workload():
assert rules.eval_misc(misc(miscType="workload", workloadPercent=100), TODAY) == rules.YES
assert rules.eval_misc(misc(miscType="workload", workloadPercent=50), TODAY) == rules.UNKNOWN
assert rules.eval_misc(misc(miscType="workload"), TODAY) == rules.UNKNOWN
def test_duration_always_yes():
assert rules.eval_misc(misc(miscType="duration"), TODAY) == rules.YES
def test_location_rules():
full_remote = misc(miscType="location", remotePercent=100)
assert rules.eval_misc(full_remote, TODAY) == rules.YES
hybrid_near = misc(miscType="location", remotePercent=60, onsiteLocation="München")
assert rules.eval_misc(hybrid_near, TODAY, geocode_fn=geo_munich) == rules.YES
hybrid_far = misc(miscType="location", remotePercent=60, onsiteLocation="Hamburg")
assert rules.eval_misc(hybrid_far, TODAY, geocode_fn=geo_hamburg) == rules.NO
unclear = misc(miscType="location")
assert rules.eval_misc(unclear, TODAY) == rules.NO
geo_broken = misc(miscType="location", remotePercent=40, onsiteLocation="Xyzstadt")
assert rules.eval_misc(geo_broken, TODAY, geocode_fn=geo_fail) == rules.UNKNOWN
def test_calc_match_and_rounding():
rows = [{"kat": "Must", "rating": "yes", "text": "a"},
{"kat": "Must", "rating": "yes", "text": "b"},
{"kat": "Must", "rating": "unknown", "text": "c"},
{"kat": "Nice", "rating": "no", "text": "d"},
{"kat": "Misc", "rating": "yes", "text": "e"}]
assert rules.calc_match(rows, "Must") == 75 # (100+100+25)/3
assert rules.calc_match(rows, "Nice") == 0
assert rules.calc_match([], "Nice") is None
# commercial rounding, half up (not banker's): (100+25)/2 = 62.5 -> 63
two = [{"kat": "Must", "rating": "yes", "text": "a"},
{"kat": "Must", "rating": "unknown", "text": "b"}]
assert rules.calc_match(two, "Must") == 63
def test_build_description_exact_format():
rows = [{"kat": "Misc", "rating": "yes", "text": "Hybrid: 60 % remote"},
{"kat": "Must", "rating": "no", "text": "Mehrjährige LLM-Plattform-Erfahrung"},
{"kat": "Must", "rating": "yes", "text": "RAG-Systeme"},
{"kat": "Nice", "rating": "yes", "text": "vLLM"}]
text = rules.build_description(rules.order_rows(rows))
assert text == (
"Must-have-Match: 50 % · Nice-to-have-Match: 100 %\n"
"\n"
"| Nr. | Kat. | ❔ | Anforderung |\n"
"|---|---|---|---|\n"
"| 1 | Must | ❌ | Mehrjährige LLM-Plattform-Erfahrung |\n"
"| 2 | Must | ✅ | RAG-Systeme |\n"
"| 3 | Nice | ✅ | vLLM |\n"
"| 4 | Misc | ✅ | Hybrid: 60 % remote |")
def test_empty_category_dash():
rows = [{"kat": "Must", "rating": "yes", "text": "a"}]
assert rules.build_description(rows).startswith(
"Must-have-Match: 100 % · Nice-to-have-Match: –")
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_rules.py -q
Expected: FAIL / ERROR with ModuleNotFoundError: No module named 'projektmatch.rules' (import error counts).
- Step 3: Implement
projektmatch/rules.py
"""Deterministic rules: Misc evaluation, match calculation, description markdown.
Replicates SKILL.md Schritt 3b/4/5. No LLM here — the numbers written to the
CRM must never be hallucinated.
"""
from __future__ import annotations
import datetime as dt
import json
import urllib.parse
import urllib.request
from math import asin, cos, radians, sin, sqrt
SAUERLACH = (47.9721357, 11.6528398)
NOMINATIM_UA = "projekt-matching/1.0 (Thomas.Langer@destengs.com)"
YES, NO, UNKNOWN = "yes", "no", "unknown"
SYMBOL = {YES: "✅", NO: "❌", UNKNOWN: "❔"}
RATING_VALUE = {YES: 100, UNKNOWN: 25, NO: 0}
KAT_ORDER = {"Must": 0, "Nice": 1, "Misc": 2}
def geocode(place: str):
"""Nominatim lookup -> (lat, lon) or None. Policy: max 1 request/s."""
url = "https://nominatim.openstreetmap.org/search?" + urllib.parse.urlencode(
{"q": place, "format": "json", "countrycodes": "de", "limit": 1})
req = urllib.request.Request(url, headers={"User-Agent": NOMINATIM_UA})
data = json.load(urllib.request.urlopen(req, timeout=20))
return (float(data[0]["lat"]), float(data[0]["lon"])) if data else None
def distance_to_sauerlach_km(place: str, geocode_fn=geocode):
"""Haversine distance; None on geocoding failure/ambiguity."""
try:
pos = geocode_fn(place)
except Exception:
return None
if pos is None:
return None
dlat = radians(pos[0] - SAUERLACH[0])
dlon = radians(pos[1] - SAUERLACH[1])
h = (sin(dlat / 2) ** 2
+ cos(radians(SAUERLACH[0])) * cos(radians(pos[0])) * sin(dlon / 2) ** 2)
return 2 * 6371 * asin(sqrt(h))
def eval_misc(item: dict, today: dt.date, geocode_fn=geocode) -> str:
kind = item.get("miscType") or "other"
if kind == "security":
return NO # K.-o.-Kriterium (SÜ dauert >= 3 Monate)
if kind == "start":
iso = item.get("startDate")
if not iso:
return UNKNOWN
try:
start = dt.date.fromisoformat(iso)
except ValueError:
return UNKNOWN
return YES if start <= today + dt.timedelta(days=56) else NO
if kind == "workload":
pct = item.get("workloadPercent")
if pct is None:
return UNKNOWN
return YES if 75 <= pct <= 100 else UNKNOWN
if kind == "duration":
return YES
if kind == "location":
if item.get("remotePercent") == 100:
return YES
place = item.get("onsiteLocation")
if not place:
return NO # Einsatzort/Remote-Anteil unklar
km = distance_to_sauerlach_km(place, geocode_fn)
if km is None:
return UNKNOWN
if km <= 50:
return YES
if km <= 60:
return UNKNOWN
return NO
return UNKNOWN
def calc_match(rows: list, kat: str):
vals = [RATING_VALUE[r["rating"]] for r in rows if r["kat"] == kat]
if not vals:
return None
return int(sum(vals) / len(vals) + 0.5) # half-up, not banker's rounding
def fmt_match(value):
return "–" if value is None else f"{value} %"
def order_rows(rows: list) -> list:
return sorted(rows, key=lambda r: KAT_ORDER[r["kat"]]) # stable sort
def build_description(rows: list) -> str:
"""rows already ordered Must -> Nice -> Misc (use order_rows)."""
must, nice = calc_match(rows, "Must"), calc_match(rows, "Nice")
lines = [f"Must-have-Match: {fmt_match(must)} · Nice-to-have-Match: {fmt_match(nice)}",
"", "| Nr. | Kat. | ❔ | Anforderung |", "|---|---|---|---|"]
for i, r in enumerate(rows, 1):
lines.append(f"| {i} | {r['kat']} | {SYMBOL[r['rating']]} | {r['text']} |")
return "\n".join(lines)
- Step 4: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_rules.py -q
Expected: 8 passed
- Step 5: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): deterministic misc rules, match calc, description markdown"
Task 3: mailparse.py — mail detection, project splitting, canonical URL
Files:
- Create:
projekt-matching/projektmatch/mailparse.py - Test:
projekt-matching/tests/test_mailparse.py
Interfaces:
-
Produces:
is_own_mail(subject: str) -> bool;bodies(msg: email.message.Message) -> tuple[str, str](html, text);split_projects(html: str, text: str) -> list[dict]with items{"title": str, "url": str}deduped per project path;canonical_url(url: str, session=None) -> str(follows redirects, strips query+fragment; raisesrequests.RequestExceptionon failure); constantBROWSER_UA. -
Step 1: Write failing tests
projekt-matching/tests/test_mailparse.py:
import email
from email.policy import default as default_policy
from unittest import mock
from projektmatch import mailparse
HTML = """
<html><body>
<h2><a href="https://www.freelancermap.de/nproj/3020338.html?utm_source=x&t=1">
Python Entwickler für eine KI-Anwendung (m/w/d)</a></h2>
<p>Beschreibung ...</p>
<a href="https://www.freelancermap.de/nproj/3020338.html?utm_source=x&t=1">Zum Projekt</a>
<h2><a href="https://www.freelancermap.de/nproj/3020995.html?agent=2">
Test Engineer (m/w/d)</a></h2>
<a href="https://www.freelancermap.de/other/page.html">Impressum</a>
</body></html>
"""
TEXT = """Neue Projekte fuer Sie:
Python Entwickler für eine KI-Anwendung (m/w/d)
https://www.freelancermap.de/nproj/3020338.html?utm_source=x
Test Engineer (m/w/d)
https://www.freelancermap.de/nproj/3020995.html?agent=2
"""
def test_is_own_mail():
assert mailparse.is_own_mail("[Projekt-Match] Freelancermap — 86 % — X")
assert mailparse.is_own_mail("[Projekt-Match-Fehler] Freelancermap — 2 Projekt(e)")
assert not mailparse.is_own_mail("Example 1 for Langflow process")
assert not mailparse.is_own_mail(None)
def test_split_projects_html_dedupes_and_prefers_long_title():
items = mailparse.split_projects(HTML, "")
assert [i["title"] for i in items] == [
"Python Entwickler für eine KI-Anwendung (m/w/d)", "Test Engineer (m/w/d)"]
assert "3020338" in items[0]["url"]
def test_split_projects_text_fallback():
items = mailparse.split_projects("", TEXT)
assert len(items) == 2
assert items[0]["title"] == "Python Entwickler für eine KI-Anwendung (m/w/d)"
def test_split_projects_none_for_normal_mail():
assert mailparse.split_projects("<p>Hallo</p>", "Hallo") == []
def test_bodies_multipart():
msg = email.message.EmailMessage(policy=default_policy)
msg["Subject"] = "x"
msg.set_content("plain body")
msg.add_alternative("<p>html body</p>", subtype="html")
html, text = mailparse.bodies(msg)
assert "html body" in html and "plain body" in text
def test_canonical_url_strips_query():
session = mock.Mock()
session.get.return_value = mock.Mock(
url="https://www.freelancermap.de/projekt/python-entwickler?ref=1",
raise_for_status=lambda: None)
out = mailparse.canonical_url("https://www.freelancermap.de/nproj/1.html?x=1",
session=session)
assert out == "https://www.freelancermap.de/projekt/python-entwickler"
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_mailparse.py -q
Expected: ModuleNotFoundError: No module named 'projektmatch.mailparse'
- Step 3: Implement
projektmatch/mailparse.py
"""E-mail parsing: own-mail detection, project splitting, canonical URLs."""
from __future__ import annotations
import re
from urllib.parse import urlsplit, urlunsplit
import requests
from bs4 import BeautifulSoup
OWN_PREFIXES = ("[Projekt-Match]", "[Projekt-Match-Fehler]")
PROJECT_URL_RE = re.compile(
r"https?://(?:www\.)?freelancermap\.de/(?:nproj/|projekt/)[^\s\"'<>)\]]+")
BROWSER_UA = ("Mozilla/5.0 (X11; Linux x86_64; rv:128.0) "
"Gecko/20100101 Firefox/128.0")
def is_own_mail(subject) -> bool:
return bool(subject) and subject.strip().startswith(OWN_PREFIXES)
def bodies(msg):
"""Return (html, text) of the first text/html and text/plain parts."""
html = text = ""
for part in msg.walk():
if part.get_content_maintype() == "multipart":
continue
try:
payload = part.get_content()
except Exception:
continue
ctype = part.get_content_type()
if ctype == "text/html" and not html:
html = payload
elif ctype == "text/plain" and not text:
text = payload
return html, text
def _url_key(url: str) -> str:
return urlsplit(url).path
def split_projects(html: str, text: str) -> list:
"""Extract project items {title, url} from a mail body. HTML preferred;
per project (URL path) the LONGEST anchor text wins (skips 'Zum Projekt'
buttons). Plain-text fallback: URL line + nearest preceding non-empty line."""
best = {} # key -> {"title", "url"}
order = []
if html:
soup = BeautifulSoup(html, "html.parser")
for a in soup.find_all("a", href=PROJECT_URL_RE):
url = a["href"]
key = _url_key(url)
title = " ".join(a.get_text(" ", strip=True).split())
if key not in best:
best[key] = {"title": title, "url": url}
order.append(key)
elif len(title) > len(best[key]["title"]):
best[key]["title"] = title
if not best and text:
last_line = ""
for line in text.splitlines():
match = PROJECT_URL_RE.search(line)
stripped = line.strip()
if match:
key = _url_key(match.group(0))
if key not in best:
best[key] = {"title": last_line or match.group(0),
"url": match.group(0)}
order.append(key)
elif stripped:
last_line = stripped
return [best[k] for k in order if best[k]["title"]]
def canonical_url(url: str, session=None, timeout=20) -> str:
"""Follow redirects, strip query string and fragment."""
sess = session or requests.Session()
resp = sess.get(url, timeout=timeout, allow_redirects=True,
headers={"User-Agent": BROWSER_UA})
resp.raise_for_status()
parts = urlsplit(resp.url)
return urlunsplit((parts.scheme, parts.netloc, parts.path, "", ""))
- Step 4: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_mailparse.py -q
Expected: 6 passed
- Step 5: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): mail parsing, project splitting, canonical URLs"
Task 4: espocrm.py — EspoCRM client
Files:
- Create:
projekt-matching/projektmatch/espocrm.py - Test:
projekt-matching/tests/test_espocrm.py
Interfaces:
-
Produces:
class EspoClient(base: str, api_key: str, retries=2, backoff=2.0)with methodssearch(entity, wtype, attr, value, select="name", max_size=50) -> list,find_opportunity_by_link(url) -> dict|None,team_id(name) -> str(raisesEspoErrorif missing),ensure_account(name, acc_type) -> str,ensure_contact(first, last, account_id) -> str,unique_opportunity_name(name) -> str,create_opportunity(payload) -> dict,get_opportunity(oid) -> dict,delete(entity, oid); helperscore_token(name) -> str,split_person(full) -> (first, last); exceptionEspoError. -
Consumes: nothing from earlier tasks.
-
Step 1: Write failing tests
projekt-matching/tests/test_espocrm.py:
from unittest import mock
import pytest
from projektmatch import espocrm
def make_client(responses):
"""Client whose session returns queued mock responses."""
client = espocrm.EspoClient("https://crm.example/api/v1", "k", retries=1,
backoff=0)
client.session = mock.Mock()
client.session.request.side_effect = responses
return client
def resp(status=200, body=None, reason=""):
r = mock.Mock()
r.status_code = status
r.json.return_value = body if body is not None else {}
r.text = "x" if body is not None else ""
r.headers = {"X-Status-Reason": reason}
return r
def test_find_opportunity_by_link():
client = make_client([resp(body={"list": [{"id": "1", "name": "P"}]})])
assert client.find_opportunity_by_link("https://x")["id"] == "1"
params = client.session.request.call_args.kwargs["params"]
assert params["where[0][type]"] == "equals"
assert params["where[0][attribute]"] == "cProjektlink"
def test_team_id_missing_raises():
client = make_client([resp(body={"list": []})])
with pytest.raises(espocrm.EspoError, match="DesTEngS"):
client.team_id("DesTEngS")
def test_ensure_account_exact_match_and_create():
hits = {"list": [{"id": "a1", "name": "Aristo Group", "type": "Reseller"}]}
client = make_client([resp(body=hits)])
assert client.ensure_account("Aristo Group", "Reseller") == "a1"
client = make_client([resp(body={"list": []}), resp(body={"id": "a2"})])
assert client.ensure_account("Neue GmbH", "Customer") == "a2"
payload = client.session.request.call_args.kwargs["json"]
assert payload == {"name": "Neue GmbH", "type": "Customer"}
def test_unique_opportunity_name_suffix():
hits = {"list": [{"name": "Projekt X"}, {"name": "Projekt X (2)"}]}
client = make_client([resp(body=hits)])
assert client.unique_opportunity_name("Projekt X") == "Projekt X (3)"
def test_retry_on_500_then_success():
client = make_client([resp(status=500), resp(body={"list": []})])
assert client.search("Team", "equals", "name", "X") == []
def test_400_raises_with_reason():
client = make_client([resp(status=400, body={}, reason="bad field")])
with pytest.raises(espocrm.EspoError, match="bad field"):
client.search("Team", "equals", "name", "X")
def test_helpers():
assert espocrm.core_token("Aristo Group GmbH") == "Aristo"
assert espocrm.split_person("Max Muster Mann") == ("Max Muster", "Mann")
assert espocrm.split_person("Mann") == ("", "Mann")
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_espocrm.py -q
Expected: ModuleNotFoundError
- Step 3: Implement
projektmatch/espocrm.py
"""Minimal EspoCRM REST client (X-Api-Key auth, in-run retries)."""
from __future__ import annotations
import re
import time
import requests
class EspoError(RuntimeError):
pass
class EspoClient:
def __init__(self, base: str, api_key: str, retries: int = 2,
backoff: float = 2.0):
self.base = base.rstrip("/")
self.session = requests.Session()
self.session.headers["X-Api-Key"] = api_key
self.retries = retries
self.backoff = backoff
def _req(self, method: str, path: str, **kw):
last = None
for attempt in range(self.retries + 1):
try:
r = self.session.request(method, f"{self.base}/{path}",
timeout=30, **kw)
if r.status_code < 500:
if r.status_code >= 400:
raise EspoError(
f"{method} {path} -> {r.status_code} "
f"{r.headers.get('X-Status-Reason', '')}")
return r.json() if r.text else {}
last = EspoError(f"{method} {path} -> {r.status_code}")
except (requests.ConnectionError, requests.Timeout) as exc:
last = exc
time.sleep(self.backoff * (attempt + 1))
raise EspoError(str(last))
def search(self, entity, wtype, attr, value, select="name", max_size=50):
params = {"where[0][type]": wtype, "where[0][attribute]": attr,
"where[0][value]": value, "select": select,
"maxSize": max_size}
return self._req("GET", entity, params=params).get("list", [])
def find_opportunity_by_link(self, url):
hits = self.search("Opportunity", "equals", "cProjektlink", url)
return hits[0] if hits else None
def team_id(self, name):
hits = self.search("Team", "equals", "name", name)
if not hits:
raise EspoError(
f"Team '{name}' nicht gefunden — cowork-api braucht Lese-"
f"Zugriff auf Team, und das Team muss existieren")
return hits[0]["id"]
def ensure_account(self, name, acc_type):
for hit in self.search("Account", "contains", "name",
core_token(name), select="name,type"):
if hit["name"].strip().lower() == name.strip().lower():
return hit["id"]
return self._req("POST", "Account",
json={"name": name, "type": acc_type})["id"]
def ensure_contact(self, first, last, account_id):
full = f"{first} {last}".strip().lower()
for hit in self.search("Contact", "contains", "name", last,
select="name,accountName"):
if hit["name"].strip().lower() == full:
return hit["id"]
return self._req("POST", "Contact",
json={"firstName": first, "lastName": last,
"accountId": account_id})["id"]
def unique_opportunity_name(self, name):
existing = {h["name"] for h in self.search(
"Opportunity", "startsWith", "name", name, max_size=100)}
if name not in existing:
return name
n = 2
while f"{name} ({n})" in existing:
n += 1
return f"{name} ({n})"
def create_opportunity(self, payload):
return self._req("POST", "Opportunity", json=payload)
def get_opportunity(self, oid):
return self._req("GET", f"Opportunity/{oid}")
def delete(self, entity, oid):
return self._req("DELETE", f"{entity}/{oid}")
def core_token(name: str) -> str:
for tok in re.split(r"[^\wÄÖÜäöüß]+", name or ""):
if len(tok) > 2:
return tok
return name
def split_person(full: str):
parts = (full or "").split()
if not parts:
return ("", "")
return (" ".join(parts[:-1]), parts[-1])
- Step 4: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_espocrm.py -q
Expected: 7 passed
- Step 5: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): EspoCRM client with dedup helpers and retries"
Task 5: llm.py — vLLM client, prompts, schemas
Files:
- Create:
projekt-matching/projektmatch/llm.py - Test:
projekt-matching/tests/test_llm.py
Interfaces:
-
Produces:
chat_json(base: str, model: str, messages: list, schema: dict) -> dict(guided_json with response_format fallback,<think>-stripping, NO max_tokens);extract_project(base, model, page_text: str) -> dict(validated againstEXTRACT_SCHEMA);match_cv(base, model, cv_text: str, requirements: list[dict]) -> list[dict](input items{"nr": int, "text": str}, returns per nr{"nr", "rating": "yes|no|unknown", "reason"}, one retry on incomplete coverage); exceptionLlmError. -
Consumes: nothing from earlier tasks.
-
Step 1: Write failing tests
projekt-matching/tests/test_llm.py:
import json
from unittest import mock
import pytest
from projektmatch import llm
def fake_post(payloads):
"""Return a mock for requests.post yielding chat completions."""
responses = []
for p in payloads:
r = mock.Mock()
r.status_code = 200
r.json.return_value = {"choices": [{"message": {"content": p}}]}
responses.append(r)
return mock.Mock(side_effect=responses)
def test_chat_json_strips_think_and_parses():
content = "<think>lange Kette</think>{\"a\": 1}"
with mock.patch("projektmatch.llm.requests.post", fake_post([content])) as p:
out = llm.chat_json("http://v/v1", "m", [{"role": "user", "content": "x"}],
{"type": "object"})
assert out == {"a": 1}
body = p.call_args.kwargs["json"]
assert "max_tokens" not in body
assert body["guided_json"] == {"type": "object"}
def test_chat_json_fallback_to_response_format():
bad = mock.Mock(status_code=400,
text="Unknown parameter: 'guided_json'")
good = mock.Mock(status_code=200)
good.json.return_value = {"choices": [{"message": {"content": "{}"}}]}
with mock.patch("projektmatch.llm.requests.post",
mock.Mock(side_effect=[bad, good])) as p:
assert llm.chat_json("http://v/v1", "m", [], {"type": "object"}) == {}
assert "response_format" in p.call_args.kwargs["json"]
def test_match_cv_retries_on_missing_nr_then_raises():
reqs = [{"nr": 1, "text": "Python"}, {"nr": 2, "text": "K8s"}]
partial = json.dumps({"ratings": [
{"nr": 1, "rating": "yes", "reason": "ok"}]})
with mock.patch("projektmatch.llm.requests.post",
fake_post([partial, partial])):
with pytest.raises(llm.LlmError, match="unvollständig"):
llm.match_cv("http://v/v1", "m", "CV", reqs)
def test_match_cv_ok():
reqs = [{"nr": 1, "text": "Python"}]
full = json.dumps({"ratings": [
{"nr": 1, "rating": "unknown", "reason": "Teilevidenz"}]})
with mock.patch("projektmatch.llm.requests.post", fake_post([full])):
out = llm.match_cv("http://v/v1", "m", "CV", reqs)
assert out[0]["rating"] == "unknown"
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_llm.py -q
Expected: ModuleNotFoundError
- Step 3: Implement
projektmatch/llm.py
"""vLLM (OpenAI-compatible) calls with guided_json + German prompts.
The model is a REASONING model: max_tokens stays UNSET so long thinking
chains cannot truncate the final answer (65k context window).
"""
from __future__ import annotations
import json
import re
import requests
THINK_RE = re.compile(r"<think>.*?</think>", re.S)
PAGE_TEXT_LIMIT = 24000
class LlmError(RuntimeError):
pass
EXTRACT_SCHEMA = {
"type": "object",
"properties": {
"projectName": {"type": "string"},
"offerType": {"enum": ["Projekt", "Arbeitnehmer-Angebot", "ANÜ"]},
"buyerType": {"enum": ["agency", "direct"]},
"companyName": {"type": ["string", "null"]},
"contactPerson": {"type": ["string", "null"]},
"requirements": {
"type": "array",
"items": {
"type": "object",
"properties": {
"text": {"type": "string"},
"kat": {"enum": ["Must", "Nice", "Misc"]},
"miscType": {"enum": ["start", "workload", "duration",
"location", "security", "other"]},
"startDate": {"type": ["string", "null"]},
"workloadPercent": {"type": ["integer", "null"]},
"remotePercent": {"type": ["integer", "null"]},
"onsiteLocation": {"type": ["string", "null"]},
},
"required": ["text", "kat", "miscType", "startDate",
"workloadPercent", "remotePercent",
"onsiteLocation"],
},
},
},
"required": ["projectName", "offerType", "buyerType", "companyName",
"contactPerson", "requirements"],
}
MATCH_SCHEMA = {
"type": "object",
"properties": {
"ratings": {
"type": "array",
"items": {
"type": "object",
"properties": {
"nr": {"type": "integer"},
"rating": {"enum": ["yes", "no", "unknown"]},
"reason": {"type": "string"},
},
"required": ["nr", "rating", "reason"],
},
},
},
"required": ["ratings"],
}
EXTRACT_SYSTEM = """Du extrahierst Anforderungen aus deutschen \
Freiberufler-Projektausschreibungen. Antworte NUR mit JSON nach Schema.
Regeln:
- projectName: Titel der Ausschreibung OHNE Portal-Zusatz (z. B. ohne "auf \
www.freelancermap.de").
- requirements: jede Anforderung einzeln, im Originalwortlaut (behutsames \
Kürzen erlaubt, Bedeutung nie verändern). Rahmenbedingungen (Start, \
Einsatzort/Remote-Anteil, Auslastung, Laufzeit) sind Anforderungen der \
Kategorie Misc.
- kat: Must = zwingend (Abschnitt "Must-haves"/"Anforderungen"; "zwingend", \
"erforderlich", "vorausgesetzt", "sehr gute Kenntnisse"). Nice = optional \
(Abschnitt "Nice-to-haves"; "von Vorteil", "wünschenswert", "idealerweise", \
"plus"). Misc = Rahmenbedingungen und alles, was weder Muss noch \
Wunsch-Qualifikation ist. Explizite Abschnittsüberschriften haben Vorrang \
vor Signalwörtern; "idealerweise" INNERHALB einer Must-Zeile lässt die \
Zeile Must bleiben.
- miscType nur für Misc-Zeilen relevant (sonst "other"): start = \
Projektstart/Verfügbarkeit (startDate als ISO-Datum YYYY-MM-DD, wenn ein \
konkretes Datum genannt ist, sonst null); workload = Auslastung \
(workloadPercent 0-100 oder null); duration = Laufzeit; location = \
Einsatzort/Remote (remotePercent 0-100 oder null; onsiteLocation = \
Ortsname oder null); security = Sicherheitsüberprüfung (SÜ, SÜ1/SÜ2/SÜ3, \
Ü2, Geheimschutz).
- offerType: "Projekt" = Freiberufler-/Werkauftrag (auch über Agentur). \
"Arbeitnehmer-Angebot" bei Festanstellung ("Festanstellung", "unbefristet", \
"Gehalt", "Arbeitsvertrag"). "ANÜ" bei Arbeitnehmerüberlassung ("ANÜ", \
"AÜG", "Überlassung", "Zeitarbeit"). Im Zweifel "Projekt".
- buyerType: "agency" bei Personaldienstleistern/Vermittlern (Hays, \
GULP/Randstad, SThree, Computer Futures, Aristo, freelancermap-Vermittler, \
"im Auftrag unseres Kunden", "für unseren Kunden"), sonst "direct". Im \
Zweifel "agency".
- companyName: Name der Agentur bzw. des Endkunden, sonst null. \
contactPerson: vollständiger Name der Ansprechperson, sonst null."""
MATCH_SYSTEM = """Du bewertest nüchtern und streng, ob ein Lebenslauf \
einzelne Projekt-Anforderungen abdeckt. Antworte NUR mit JSON nach Schema: \
für JEDE übergebene Nr. genau ein Eintrag in ratings.
Bewertung:
- "yes" NUR bei klarer Evidenz im Lebenslauf.
- "no" wenn der Lebenslauf nichts Belastbares hergibt. Streng bleiben — \
eine geschönte Bewertung macht die Match-Zahlen wertlos. Kalibrierung: \
Proof-of-Concept-Erfahrung deckt "produktiven Betrieb" NICHT ab; \
"mehrjährig" wörtlich nehmen; ein Produktname (z. B. "Azure DevOps \
Server") belegt KEINE Cloud-Plattform-Erfahrung.
- "unknown" NUR bei echter Teilevidenz, wenn die Entscheidung von Wissen \
abhängt, das nur der Kandidat selbst hat.
- "wie z. B."-Aufzählungen: gleichwertige Alternativen zählen als Abdeckung \
(Beispiel: Ollama/llama.cpp/Transformers decken "LLM-Inference-Stacks wie \
z. B. vLLM, TGI, Triton" ab).
- reason: EIN kurzer deutscher Satz mit der Begründung."""
def _post(base, model, messages, body_extra):
body = {"model": model, "messages": messages, "temperature": 0.1}
body.update(body_extra)
return requests.post(f"{base.rstrip('/')}/chat/completions", json=body,
timeout=1500)
def chat_json(base, model, messages, schema):
resp = _post(base, model, messages, {"guided_json": schema})
if resp.status_code == 400 and "guided_json" in getattr(resp, "text", ""):
resp = _post(base, model, messages, {"response_format": {
"type": "json_schema",
"json_schema": {"name": "out", "schema": schema}}})
if resp.status_code != 200:
raise LlmError(f"vLLM HTTP {resp.status_code}: "
f"{getattr(resp, 'text', '')[:300]}")
content = resp.json()["choices"][0]["message"]["content"] or ""
content = THINK_RE.sub("", content).strip()
try:
return json.loads(content)
except json.JSONDecodeError as exc:
raise LlmError(f"LLM lieferte kein JSON: {exc}: {content[:200]}")
def extract_project(base, model, page_text):
messages = [{"role": "system", "content": EXTRACT_SYSTEM},
{"role": "user", "content":
"Ausschreibungstext:\n\n" + page_text[:PAGE_TEXT_LIMIT]}]
out = chat_json(base, model, messages, EXTRACT_SCHEMA)
if not out.get("requirements"):
raise LlmError("Extraktion ohne Anforderungen")
return out
def match_cv(base, model, cv_text, requirements):
listing = "\n".join(f"{r['nr']}. {r['text']}" for r in requirements)
messages = [{"role": "system", "content": MATCH_SYSTEM},
{"role": "user", "content":
f"Lebenslauf:\n\n{cv_text}\n\nAnforderungen:\n{listing}"}]
wanted = {r["nr"] for r in requirements}
for attempt in range(2):
out = chat_json(base, model, messages, MATCH_SCHEMA)
got = {r["nr"]: r for r in out.get("ratings", []) if r["nr"] in wanted}
if set(got) == wanted:
return [got[n] for n in sorted(got)]
missing = sorted(wanted - set(got))
messages = messages + [
{"role": "assistant", "content": json.dumps(out)},
{"role": "user", "content":
f"Es fehlen Bewertungen für Nr. {missing}. Antworte erneut "
f"mit ratings für ALLE Nummern."}]
raise LlmError(f"Matching unvollständig, fehlend: {missing}")
- Step 4: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_llm.py -q
Expected: 4 passed
- Step 5: Live smoke test against local vLLM (host)
cd /home/lwc/bin/projekt-matching && .venv/bin/python - <<'PY'
from projektmatch import llm
out = llm.extract_project(
"http://127.0.0.1:8081/v1", "AxionML/Qwen3.5-9B-NVFP4",
"Projekttitel: Python Entwickler KI\n\nMust-haves:\n- Python\n- LLM-Erfahrung\n"
"Nice-to-haves:\n- Kubernetes\n\nStart: 01.09.2026, 100% remote, 6 Monate, Auslastung 100%")
print(out["projectName"], out["offerType"], len(out["requirements"]))
for r in out["requirements"]:
print(r["kat"], r["miscType"], r["text"][:60])
PY
Expected: prints a plausible project name, Projekt, and ≥ 6 requirements with Must/Nice/Misc categories (start/workload/duration/location misc rows). If HTTP 400 mentions guided_json, the fallback path must have handled it — investigate only if the call raises.
- Step 6: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): vLLM guided-json client with extraction and matching prompts"
Task 6: config.py + mailer.py — configuration, IMAP, SMTP
Files:
- Create:
projekt-matching/projektmatch/config.py,projekt-matching/projektmatch/mailer.py - Test:
projekt-matching/tests/test_mailer.py
Interfaces:
-
Produces
config.Cfgdataclass (all fields with defaults):imap_user="chancen@destengs.com",imap_password="",espo_base="",espo_api_key="",notify_to="",alert_to="",threshold=85,vllm_base="",vllm_model="",flow2_id="",langflow_api_key="",langflow_base="http://127.0.0.1:7860",data_dir="/app/langflow",crm_web_base="https://crm.creature-go.com". -
Produces
mailer.MailBox(user, password, conn=None)withunchecked_uids() -> list[str],fetch(uid) -> email.message.Message,flag_checked(uid),move_to_trash(uid),close();mailer.send_mail(user, password, to, subject, body); constantsKEYWORD="$ProjektChecked",IMAP_HOST/PORT,SMTP_HOST/PORT. -
Step 1: Write
projektmatch/config.py
"""Runtime configuration passed from Langflow component inputs."""
from dataclasses import dataclass
@dataclass
class Cfg:
imap_user: str = "chancen@destengs.com"
imap_password: str = ""
espo_base: str = ""
espo_api_key: str = ""
notify_to: str = ""
alert_to: str = ""
threshold: int = 85
vllm_base: str = ""
vllm_model: str = ""
flow2_id: str = ""
langflow_api_key: str = ""
langflow_base: str = "http://127.0.0.1:7860"
data_dir: str = "/app/langflow"
crm_web_base: str = "https://crm.creature-go.com"
- Step 2: Write failing tests
projekt-matching/tests/test_mailer.py:
import email
from unittest import mock
from projektmatch import mailer
def fake_conn():
conn = mock.Mock()
conn.uid.return_value = ("OK", [b"3 7"])
conn.list.return_value = ("OK", [
b'(\\HasNoChildren) "." "INBOX"',
b'(\\HasNoChildren \\Trash) "." "Trash"'])
return conn
def test_unchecked_uids_search():
box = mailer.MailBox("u", "p", conn=fake_conn())
assert box.unchecked_uids() == ["3", "7"]
box.conn.uid.assert_called_with("SEARCH", None,
"UNKEYWORD $ProjektChecked")
def test_flag_checked():
box = mailer.MailBox("u", "p", conn=fake_conn())
box.flag_checked("3")
box.conn.uid.assert_called_with("STORE", "3", "+FLAGS",
"($ProjektChecked)")
def test_trash_folder_by_special_use():
box = mailer.MailBox("u", "p", conn=fake_conn())
assert box.trash_folder() == "Trash"
def test_move_uses_move_then_fallback():
conn = fake_conn()
box = mailer.MailBox("u", "p", conn=conn)
conn.uid.return_value = ("OK", [b""])
box.move_to_trash("3")
assert conn.uid.call_args_list[-1].args[:2] == ("MOVE", "3")
conn.uid.side_effect = [("NO", [b""]), ("OK", [b""]), ("OK", [b""])]
box.move_to_trash("4")
assert conn.expunge.called
def test_send_mail_starttls():
with mock.patch("projektmatch.mailer.smtplib.SMTP") as smtp:
server = smtp.return_value.__enter__.return_value
mailer.send_mail("u@x", "pw", "to@x", "Subj", "Body")
server.starttls.assert_called_once()
server.login.assert_called_once_with("u@x", "pw")
msg = server.send_message.call_args.args[0]
assert msg["Subject"] == "Subj" and msg["To"] == "to@x"
- Step 3: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_mailer.py -q
Expected: ModuleNotFoundError: No module named 'projektmatch.mailer'
- Step 4: Implement
projektmatch/mailer.py
"""IMAP mailbox handling and SMTP sending for chancen@destengs.com."""
from __future__ import annotations
import email
import imaplib
import smtplib
import ssl
from email.message import EmailMessage
from email.policy import default as default_policy
IMAP_HOST, IMAP_PORT = "mail.destengs.com", 993
SMTP_HOST, SMTP_PORT = "mail.destengs.com", 587
KEYWORD = "$ProjektChecked"
class MailBox:
def __init__(self, user, password, conn=None):
self.conn = conn or imaplib.IMAP4_SSL(IMAP_HOST, IMAP_PORT)
if conn is None:
self.conn.login(user, password)
self.conn.select("INBOX")
self._trash = None
def unchecked_uids(self):
typ, data = self.conn.uid("SEARCH", None, f"UNKEYWORD {KEYWORD}")
if typ != "OK" or not data or not data[0]:
return []
return [u.decode() for u in data[0].split()]
def fetch(self, uid):
typ, data = self.conn.uid("FETCH", uid, "(BODY.PEEK[])")
return email.message_from_bytes(data[0][1], policy=default_policy)
def flag_checked(self, uid):
self.conn.uid("STORE", uid, "+FLAGS", f"({KEYWORD})")
def trash_folder(self):
if self._trash:
return self._trash
typ, data = self.conn.list()
candidates = []
for raw in data or []:
line = raw.decode() if isinstance(raw, bytes) else str(raw)
name = line.rsplit(" ", 1)[-1].strip('"')
if "\\Trash" in line.split(")")[0]:
self._trash = name
return name
candidates.append(name)
for cand in ("Trash", "INBOX.Trash"):
if cand in candidates:
self._trash = cand
return cand
self._trash = "Trash"
return self._trash
def move_to_trash(self, uid):
folder = self.trash_folder()
typ, _ = self.conn.uid("MOVE", uid, folder)
if typ != "OK":
self.conn.uid("COPY", uid, folder)
self.conn.uid("STORE", uid, "+FLAGS", "(\\Deleted)")
self.conn.expunge()
def close(self):
try:
self.conn.logout()
except Exception:
pass
def send_mail(user, password, to, subject, body):
msg = EmailMessage()
msg["From"], msg["To"], msg["Subject"] = user, to, subject
msg.set_content(body)
with smtplib.SMTP(SMTP_HOST, SMTP_PORT, timeout=30) as server:
server.starttls(context=ssl.create_default_context())
server.login(user, password)
server.send_message(msg)
- Step 5: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_mailer.py -q
Expected: 5 passed
- Step 6: Live IMAP smoke test — verify the server accepts custom keywords
cd /home/lwc/bin/projekt-matching && PM_PW='<IMAP password from the project brief>' .venv/bin/python - <<'PY'
import os
import imaplib
conn = imaplib.IMAP4_SSL("mail.destengs.com", 993)
conn.login("chancen@destengs.com", os.environ["PM_PW"])
typ, data = conn.select("INBOX")
print("select:", typ)
typ, data = conn.response("PERMANENTFLAGS")
print("permanentflags:", data) # must contain \* (custom keywords allowed)
typ, data = conn.list()
print([l.decode() for l in data if b"Trash" in l])
conn.logout()
PY
Expected: select: OK; PERMANENTFLAGS contains \*; a Trash folder line appears. If \* is missing the server refuses custom keywords — then change KEYWORD handling to use the standard \Flagged flag instead (unchecked_uids → UNFLAGGED, flag_checked → +FLAGS (\Flagged)) and adjust the two mailer tests accordingly.
- Step 7: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): Cfg dataclass, IMAP mailbox with keyword/trash handling, SMTP send"
Task 7: pm_trace.py + stages.py — Flow-2 stage functions
Files:
- Create:
projekt-matching/projektmatch/pm_trace.py,projekt-matching/projektmatch/stages.py - Test:
projekt-matching/tests/test_stages.py
Interfaces:
-
Consumes:
rules,llm,espocrm,mailer,config.Cfgfrom Tasks 2–6. -
Produces:
run_stage(name: str, fn, ctx: dict, cfg: Cfg) -> dict(skips whenctx["status"] == "failed"; catches exceptions intostatus="failed", error=...); stage functionsstage_fetch,stage_extract,stage_match,stage_rules,stage_crm,stage_notify— each(ctx, cfg) -> ctxoperating on the ctx dict shape from the File Structure section;TEAM_BY_OFFERmapping;pm_trace.post_trace(ctx)(readsLANGFUSE_PUBLIC_KEY/SECRET_KEY/HOSTenv vars set on the Langflow container; never raises). -
Step 1: Write failing tests
projekt-matching/tests/test_stages.py:
import datetime as dt
from unittest import mock
from projektmatch import stages
from projektmatch.config import Cfg
CFG = Cfg(imap_password="pw", espo_base="https://crm/api/v1",
espo_api_key="k", notify_to="n@x", alert_to="a@x",
vllm_base="http://v/v1", vllm_model="m", data_dir="/tmp")
EXTRACT = {
"projectName": "Python Entwickler KI",
"offerType": "Projekt", "buyerType": "agency",
"companyName": "Aristo Group", "contactPerson": "Max Mann",
"requirements": [
{"text": "Python", "kat": "Must", "miscType": "other",
"startDate": None, "workloadPercent": None, "remotePercent": None,
"onsiteLocation": None},
{"text": "K8s", "kat": "Nice", "miscType": "other",
"startDate": None, "workloadPercent": None, "remotePercent": None,
"onsiteLocation": None},
{"text": "100 % remote", "kat": "Misc", "miscType": "location",
"startDate": None, "workloadPercent": None, "remotePercent": 100,
"onsiteLocation": None},
],
}
def ctx_after_match(ratings):
ctx = {"canonical": "https://x/projekt/p", "title": "t", "status": "ok",
"extract": dict(EXTRACT)}
ctx["ratings"] = ratings
return ctx
def test_run_stage_catches_and_skips():
def boom(ctx, cfg):
raise ValueError("kaputt")
ctx = stages.run_stage("fetch", boom, {"status": "ok"}, CFG)
assert ctx["status"] == "failed" and "fetch: kaputt" in ctx["error"]
untouched = stages.run_stage("extract", boom, dict(ctx), CFG)
assert untouched["status"] == "failed" # boom not called again
def test_stage_rules_consider_and_reject():
yes = [{"nr": 1, "rating": "yes", "reason": "ok"},
{"nr": 2, "rating": "yes", "reason": "ok"}]
ctx = stages.stage_rules(ctx_after_match(yes), CFG)
assert ctx["decision"] == "consider" and ctx["mustMatch"] == 100
assert ctx["status"] == "ok"
assert "| 1 | Must | ✅ | Python |" in ctx["description"]
no = [{"nr": 1, "rating": "no", "reason": "fehlt"},
{"nr": 2, "rating": "yes", "reason": "ok"}]
ctx2 = stages.stage_rules(ctx_after_match(no), CFG)
assert ctx2["status"] == "rejected" and ctx2["decision"] == "rejected"
assert ctx2["mustMatch"] == 0
def test_stage_crm_agency_linking_and_verify():
ctx = ctx_after_match([{"nr": 1, "rating": "yes", "reason": "ok"},
{"nr": 2, "rating": "yes", "reason": "ok"}])
ctx = stages.stage_rules(ctx, CFG)
espo = mock.Mock()
espo.team_id.return_value = "T1"
espo.ensure_account.return_value = "A1"
espo.ensure_contact.return_value = "C1"
espo.unique_opportunity_name.return_value = "Python Entwickler KI"
espo.create_opportunity.return_value = {"id": "O1"}
espo.get_opportunity.return_value = {
"id": "O1", "name": "Python Entwickler KI",
"cProjektlink": "https://x/projekt/p",
"description": ctx["description"], "cAccount1Id": "A1",
"accountId": None, "teamsIds": ["T1"]}
out = stages.stage_crm(ctx, CFG, espo=espo)
assert out["status"] == "created" and out["opportunityId"] == "O1"
assert out["crmUrl"].endswith("#Opportunity/view/O1")
payload = espo.create_opportunity.call_args.args[0]
assert payload["cAccount1Id"] == "A1" and "accountId" not in payload
assert payload["teamsIds"] == ["T1"]
espo.ensure_account.assert_called_with("Aristo Group", "Reseller")
def test_stage_crm_skips_when_rejected():
ctx = {"status": "rejected", "decision": "rejected"}
espo = mock.Mock()
assert stages.stage_crm(ctx, CFG, espo=espo)["status"] == "rejected"
espo.team_id.assert_not_called()
def test_stage_notify_only_on_created():
sent = []
with mock.patch("projektmatch.stages.mailer.send_mail",
lambda *a: sent.append(a)), \
mock.patch("projektmatch.stages.pm_trace.post_trace", lambda c: None):
ctx = {"status": "created", "projectName": "P", "mustMatch": 90,
"niceMatch": None, "crmUrl": "https://crm/#Opportunity/view/O1"}
stages.stage_notify(dict(ctx), CFG)
assert sent[0][2] == "n@x"
assert sent[0][3] == "[Projekt-Match] Freelancermap — 90 % — P"
sent.clear()
stages.stage_notify({"status": "rejected"}, CFG)
assert not sent
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_stages.py -q
Expected: ModuleNotFoundError
- Step 3: Implement
projektmatch/pm_trace.py
"""Post one consolidated Langfuse trace per project (ingestion API).
Reads the LANGFUSE_* env vars that create_pod_langflow.sh already sets on
the Langflow container. Never raises — tracing must not fail a run."""
from __future__ import annotations
import datetime as dt
import os
import uuid
import requests
def post_trace(ctx: dict) -> None:
host = os.environ.get("LANGFUSE_HOST")
pk = os.environ.get("LANGFUSE_PUBLIC_KEY")
sk = os.environ.get("LANGFUSE_SECRET_KEY")
if not (host and pk and sk):
return
now = dt.datetime.now(dt.timezone.utc).isoformat()
body = {
"id": str(uuid.uuid4()),
"timestamp": now,
"name": "projekt-match",
"input": {"canonical": ctx.get("canonical"), "title": ctx.get("title")},
"output": {"decision": ctx.get("decision"),
"status": ctx.get("status"),
"mustMatch": ctx.get("mustMatch"),
"niceMatch": ctx.get("niceMatch"),
"description": ctx.get("description")},
"metadata": {"opportunityId": ctx.get("opportunityId"),
"crmUrl": ctx.get("crmUrl"),
"error": ctx.get("error"),
"projectName": ctx.get("projectName")},
"tags": [ctx.get("status") or "unknown"],
}
event = {"batch": [{"id": str(uuid.uuid4()), "type": "trace-create",
"timestamp": now, "body": body}]}
try:
requests.post(f"{host.rstrip('/')}/api/public/ingestion", json=event,
auth=(pk, sk), timeout=15)
except Exception:
pass
- Step 4: Implement
projektmatch/stages.py
"""Flow 2 stage functions. Each takes (ctx, cfg) and returns ctx.
ctx["status"]: "ok" -> pipeline continues; "failed" -> all later stages
skip; "rejected" -> only CRM/notify skip; "created" -> set by stage_crm.
"""
from __future__ import annotations
import datetime as dt
import json
import time
from pathlib import Path
import requests
from bs4 import BeautifulSoup
from . import espocrm, llm, mailer, pm_trace, rules
from .mailparse import BROWSER_UA
TEAM_BY_OFFER = {"Projekt": "DesTEngS",
"Arbeitnehmer-Angebot": "Arbeitnehmer",
"ANÜ": "ANÜ"}
MIN_PAGE_CHARS = 200
def run_stage(name, fn, ctx, cfg, **kw):
if ctx.get("status") == "failed":
return ctx
try:
return fn(ctx, cfg, **kw)
except Exception as exc: # noqa: BLE001
ctx.update(status="failed", error=f"{name}: {exc}")
return ctx
def stage_fetch(ctx, cfg):
last = None
for attempt in range(3):
try:
resp = requests.get(ctx["canonical"], timeout=30,
headers={"User-Agent": BROWSER_UA})
resp.raise_for_status()
soup = BeautifulSoup(resp.text, "html.parser")
for tag in soup(["script", "style", "noscript"]):
tag.decompose()
text = "\n".join(line for line in
soup.get_text("\n", strip=True).splitlines()
if line.strip())
if len(text) < MIN_PAGE_CHARS:
raise ValueError(
f"Seitentext zu kurz ({len(text)} Zeichen) — "
f"Login-Wall oder leere Seite?")
ctx.update(page_text=text[:llm.PAGE_TEXT_LIMIT], status="ok")
return ctx
except Exception as exc: # noqa: BLE001
last = exc
time.sleep(2 * (attempt + 1))
raise RuntimeError(f"Abruf fehlgeschlagen: {last}")
def stage_extract(ctx, cfg):
ctx["extract"] = llm.extract_project(cfg.vllm_base, cfg.vllm_model,
ctx["page_text"])
ctx["projectName"] = ctx["extract"]["projectName"]
return ctx
def stage_match(ctx, cfg):
reqs = [{"nr": i + 1, "text": r["text"]}
for i, r in enumerate(ctx["extract"]["requirements"])
if r["kat"] in ("Must", "Nice")]
if not reqs:
ctx["ratings"] = []
return ctx
cv_path = Path(cfg.data_dir) / "vorgaben" / \
"Lebenslauf_Dr-Ing_Thomas_Langer.md"
cv_text = cv_path.read_text(encoding="utf-8")
ctx["ratings"] = llm.match_cv(cfg.vllm_base, cfg.vllm_model, cv_text, reqs)
return ctx
def stage_rules(ctx, cfg):
today = dt.date.today()
rated = {r["nr"]: r["rating"] for r in ctx.get("ratings", [])}
rows, nr = [], 0
for req in ctx["extract"]["requirements"]:
if req["kat"] in ("Must", "Nice"):
nr += 1
rating = rated.get(nr, rules.UNKNOWN)
else:
rating = rules.eval_misc(req, today)
rows.append({"kat": req["kat"], "rating": rating, "text": req["text"]})
rows = rules.order_rows(rows)
ctx["mustMatch"] = rules.calc_match(rows, "Must")
ctx["niceMatch"] = rules.calc_match(rows, "Nice")
ctx["description"] = rules.build_description(rows)
consider = ctx["mustMatch"] is not None and \
ctx["mustMatch"] > int(cfg.threshold)
ctx["decision"] = "consider" if consider else "rejected"
if not consider:
ctx["status"] = "rejected"
return ctx
def stage_crm(ctx, cfg, espo=None):
if ctx.get("status") != "ok" or ctx.get("decision") != "consider":
return ctx
espo = espo or espocrm.EspoClient(cfg.espo_base, cfg.espo_api_key)
ex = ctx["extract"]
team = espo.team_id(TEAM_BY_OFFER[ex["offerType"]])
agency = ex["buyerType"] == "agency"
account_id = None
if ex.get("companyName"):
account_id = espo.ensure_account(
ex["companyName"], "Reseller" if agency else "Customer")
contact_id = None
if ex.get("contactPerson") and account_id:
first, last = espocrm.split_person(ex["contactPerson"])
contact_id = espo.ensure_contact(first, last, account_id)
name = espo.unique_opportunity_name(ex["projectName"])
payload = {"name": name, "description": ctx["description"],
"cProjektlink": ctx["canonical"], "teamsIds": [team]}
if account_id:
payload["cAccount1Id" if agency else "accountId"] = account_id
if contact_id:
payload["contactsIds"] = [contact_id]
created = espo.create_opportunity(payload)
oid = created["id"]
back = espo.get_opportunity(oid)
problems = []
if back.get("name") != name:
problems.append("name")
if back.get("cProjektlink") != ctx["canonical"]:
problems.append("cProjektlink")
if back.get("description") != ctx["description"]:
problems.append("description")
if account_id and agency and not back.get("cAccount1Id"):
problems.append("cAccount1Id")
if account_id and not agency and not back.get("accountId"):
problems.append("accountId")
if problems:
raise RuntimeError(f"CRM-Verifikation fehlgeschlagen: {problems}")
ctx.update(status="created", opportunityId=oid, projectName=name,
crmUrl=f"{cfg.crm_web_base}/#Opportunity/view/{oid}")
return ctx
def stage_notify(ctx, cfg):
if ctx.get("status") == "created":
subject = (f"[Projekt-Match] Freelancermap — {ctx['mustMatch']} % — "
f"{ctx['projectName']}")
body = (f"Neues passendes Projekt gefunden.\n\n"
f"Projekt: {ctx['projectName']}\n"
f"Must-have-Match: {rules.fmt_match(ctx['mustMatch'])}\n"
f"Nice-to-have-Match: {rules.fmt_match(ctx.get('niceMatch'))}\n\n"
f"CRM-Verkaufschance: {ctx['crmUrl']}\n")
mailer.send_mail(cfg.imap_user, cfg.imap_password, cfg.notify_to,
subject, body)
pm_trace.post_trace(ctx)
return ctx
def summary(ctx):
"""Compact JSON result string for TextOutput / Flow 1."""
keys = ("status", "error", "decision", "mustMatch", "niceMatch",
"opportunityId", "crmUrl", "projectName", "canonical")
return json.dumps({k: ctx.get(k) for k in keys}, ensure_ascii=False)
- Step 5: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_stages.py -q
Expected: 5 passed
- Step 6: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): flow-2 stages with gate, CRM verify, notification, langfuse trace"
Task 8: ingest.py — Flow-1 logic
Files:
- Create:
projekt-matching/projektmatch/ingest.py - Test:
projekt-matching/tests/test_ingest.py
Interfaces:
-
Consumes:
mailer.MailBox,mailer.send_mail,mailparse,espocrm.EspoClient,config.Cfg. -
Produces:
run_ingest(cfg, mailbox=None, dispatch=None, espo=None, resolve=None) -> dict(summary counters);dispatch_flow2(cfg, item: dict) -> dict(POST run API, returns parsed status dict);parse_run_result(payload) -> dict. -
Step 1: Write failing tests
projekt-matching/tests/test_ingest.py:
import json
from unittest import mock
from projektmatch import ingest
from projektmatch.config import Cfg
PROJECT_HTML = ('<a href="https://www.freelancermap.de/nproj/1.html?x=1">'
'Projekt Eins</a>')
def cfg(tmp_path):
return Cfg(imap_password="pw", alert_to="a@x", flow2_id="F2",
langflow_api_key="key", data_dir=str(tmp_path))
def run(tmp_path, mails, dup=None, dispatch_result=None):
box = mock.Mock()
box.unchecked_uids.return_value = list(mails)
box.fetch.side_effect = lambda uid: mails[uid]
espo = mock.Mock()
espo.find_opportunity_by_link.return_value = dup
sent = []
with mock.patch("projektmatch.ingest.mailparse.bodies",
side_effect=lambda m: (m["html"], "")), \
mock.patch("projektmatch.ingest.mailer.send_mail",
lambda *a: sent.append(a)):
summary = ingest.run_ingest(
cfg(tmp_path), mailbox=box, espo=espo,
resolve=lambda url: "https://www.freelancermap.de/projekt/eins",
dispatch=lambda item: dispatch_result or {"status": "created"})
return summary, box, sent
def mail(subject, html=""):
return {"Subject": subject, "html": html}
def test_non_project_mail_flagged_not_trashed(tmp_path):
mails = {"1": mail("Newsletter", "<p>Hallo</p>")}
summary, box, sent = run(tmp_path, mails)
box.flag_checked.assert_called_once_with("1")
box.move_to_trash.assert_not_called()
assert summary["mails"] == 0 and not sent
def test_own_mail_skipped(tmp_path):
mails = {"1": mail("[Projekt-Match] Freelancermap — 90 % — X",
PROJECT_HTML)}
summary, box, sent = run(tmp_path, mails)
box.flag_checked.assert_called_once_with("1")
box.move_to_trash.assert_not_called()
def test_project_mail_created_and_trashed(tmp_path):
mails = {"1": mail("Example 1 for Langflow process", PROJECT_HTML)}
summary, box, sent = run(tmp_path, mails)
assert summary == {"mails": 1, "projects": 1, "created": 1,
"rejected": 0, "duplicate": 0, "failed": 0}
box.move_to_trash.assert_called_once_with("1")
assert not sent
def test_duplicate_skips_dispatch(tmp_path):
mails = {"1": mail("Mail", PROJECT_HTML)}
summary, box, sent = run(tmp_path, mails, dup={"id": "X"})
assert summary["duplicate"] == 1 and summary["created"] == 0
box.move_to_trash.assert_called_once()
def test_failed_dispatch_alerts_and_trashes(tmp_path):
mails = {"1": mail("Mail", PROJECT_HTML)}
summary, box, sent = run(tmp_path, mails,
dispatch_result={"status": "failed",
"error": "vLLM down"})
assert summary["failed"] == 1
assert len(sent) == 1
_, _, to, subject, body = sent[0]
assert to == "a@x"
assert subject == "[Projekt-Match-Fehler] Freelancermap — 1 Projekt(e)"
assert "vLLM down" in body
box.move_to_trash.assert_called_once()
def test_lock_prevents_second_run(tmp_path):
lock = ingest.acquire_lock(str(tmp_path))
assert lock is not None
summary = ingest.run_ingest(cfg(tmp_path), mailbox=mock.Mock())
assert summary == {"skipped": "locked"}
ingest.release_lock(lock)
def test_parse_run_result_finds_status_json():
inner = json.dumps({"status": "created", "mustMatch": 90})
payload = {"outputs": [{"outputs": [{"results": {"text": {
"data": {"text": inner}}}}]}]}
assert ingest.parse_run_result(payload)["status"] == "created"
- Step 2: Run tests to verify they fail
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest tests/test_ingest.py -q
Expected: ModuleNotFoundError
- Step 3: Implement
projektmatch/ingest.py
"""Flow 1: poll inbox, split projects, dedup, dispatch Flow 2, alert, trash."""
from __future__ import annotations
import email.header
import fcntl
import json
import os
import requests
from . import espocrm, mailer, mailparse
SUMMARY_KEYS = ("created", "rejected", "duplicate", "failed")
LOCK_NAME = "projekt-matching.lock"
def acquire_lock(data_dir):
path = os.path.join(data_dir, LOCK_NAME)
handle = open(path, "w") # noqa: SIM115
try:
fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
return handle
except OSError:
handle.close()
return None
def release_lock(handle):
fcntl.flock(handle, fcntl.LOCK_UN)
handle.close()
def decode_subject(msg):
raw = msg["Subject"] or ""
try:
parts = email.header.decode_header(raw)
return "".join(p.decode(enc or "utf-8", "replace")
if isinstance(p, bytes) else p for p, enc in parts)
except Exception:
return raw
def dispatch_flow2(cfg, item):
resp = requests.post(
f"{cfg.langflow_base}/api/v1/run/{cfg.flow2_id}?stream=false",
headers={"x-api-key": cfg.langflow_api_key},
json={"input_value": json.dumps(
{"canonical": item["canonical"], "title": item["title"]},
ensure_ascii=False),
"input_type": "text", "output_type": "text"},
timeout=1800)
resp.raise_for_status()
return parse_run_result(resp.json())
def parse_run_result(payload):
"""Depth-first search for the stage summary JSON in the run response."""
stack = [payload]
while stack:
node = stack.pop()
if isinstance(node, dict):
stack.extend(node.values())
elif isinstance(node, list):
stack.extend(node)
elif isinstance(node, str) and node.lstrip().startswith("{"):
try:
data = json.loads(node)
except json.JSONDecodeError:
continue
if isinstance(data, dict) and "status" in data:
return data
raise ValueError("Kein Status-JSON in der Run-Antwort gefunden")
def run_ingest(cfg, mailbox=None, dispatch=None, espo=None, resolve=None):
lock = acquire_lock(cfg.data_dir)
if lock is None:
return {"skipped": "locked"}
try:
box = mailbox or mailer.MailBox(cfg.imap_user, cfg.imap_password)
crm = espo or espocrm.EspoClient(cfg.espo_base, cfg.espo_api_key)
send = dispatch or (lambda item: dispatch_flow2(cfg, item))
resolve = resolve or mailparse.canonical_url
summary = {"mails": 0, "projects": 0,
**{k: 0 for k in SUMMARY_KEYS}}
for uid in box.unchecked_uids():
msg = box.fetch(uid)
subject = decode_subject(msg)
html, text = mailparse.bodies(msg)
items = [] if mailparse.is_own_mail(subject) else \
mailparse.split_projects(html, text)
if not items:
box.flag_checked(uid)
continue
summary["mails"] += 1
results = []
for item in items:
summary["projects"] += 1
result = dict(item)
try:
result["canonical"] = resolve(item["url"])
dup = crm.find_opportunity_by_link(result["canonical"])
if dup:
result["status"] = "duplicate"
else:
result.update(send(result))
except Exception as exc: # noqa: BLE001
result.update(status="failed", error=str(exc))
results.append(result)
summary[result["status"]] += 1
failed = [r for r in results if r["status"] == "failed"]
if failed:
lines = "\n\n".join(
f"- {f.get('title', '?')}\n"
f" {f.get('canonical', f.get('url', '?'))}\n"
f" Fehler: {f.get('error', '?')}" for f in failed)
mailer.send_mail(
cfg.imap_user, cfg.imap_password, cfg.alert_to,
f"[Projekt-Match-Fehler] Freelancermap — "
f"{len(failed)} Projekt(e)",
"Folgende Projekte konnten nicht verarbeitet "
"werden:\n\n" + lines + "\n")
box.move_to_trash(uid)
box.close()
return summary
finally:
release_lock(lock)
- Step 4: Run tests to verify they pass
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/pytest -q
Expected: all tests pass (Tasks 1–8 suites, ~36 tests).
- Step 5: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): flow-1 ingest with lock, dedup, dispatch, alert aggregation"
Task 9: Langflow custom components (thin wrappers)
Files:
- Create:
projekt-matching/components/pm_fetch.py,pm_extract.py,pm_match.py,pm_rules.py,pm_crm.py,pm_notify.py,pm_ingest.py
Interfaces:
-
Consumes:
projektmatch.stages,projektmatch.ingest,projektmatch.config.Cfg(deployed at/app/langflow/projektmatchinside the container, Task 10). -
Produces: seven Langflow
Componentclasses. Class names / output names are consumed bybuild_flows.py(Task 10): every component has exactly one output namedout, methodbuild_out. Global-variable names referenced viaload_from_db:PM_IMAP_PASSWORD,PM_ESPO_API_KEY,PM_LANGFLOW_API_KEY(Credential, SecretStrInput);PM_ESPO_BASE,PM_NOTIFY_TO,PM_ALERT_TO,PM_THRESHOLD,PM_VLLM_BASE,PM_VLLM_MODEL,PM_FLOW2_ID(Generic, MessageTextInput). -
No host-side unit tests (wrappers only; logic is tested in Tasks 2–8). Verified live in Task 10 Step 6 and Task 13.
-
Step 1: Write
components/pm_fetch.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
import json
from langflow.custom import Component
from langflow.io import MessageTextInput, Output
from langflow.schema import Data
from projektmatch import stages
class PMFetch(Component):
display_name = "PM 1 Fetch"
description = "Projektseite abrufen (canonical URL -> Seitentext)"
inputs = [MessageTextInput(name="payload", display_name="Payload JSON")]
outputs = [Output(name="out", display_name="Context", method="build_out")]
def build_out(self) -> Data:
ctx = json.loads(self.payload)
ctx.setdefault("status", "ok")
ctx = stages.run_stage("fetch", stages.stage_fetch, ctx, None)
self.status = ctx.get("status", "")
return Data(data=ctx)
(stage_fetch ignores cfg, so None is fine there.)
- Step 2: Write
components/pm_extract.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output
from langflow.schema import Data
from projektmatch import stages
from projektmatch.config import Cfg
class PMExtract(Component):
display_name = "PM 2 Extract"
description = "LLM 1: Anforderungen strukturiert extrahieren"
inputs = [
DataInput(name="ctx", display_name="Context"),
MessageTextInput(name="vllm_base", display_name="vLLM Base",
value="PM_VLLM_BASE", load_from_db=True,
advanced=True),
MessageTextInput(name="vllm_model", display_name="vLLM Model",
value="PM_VLLM_MODEL", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Context", method="build_out")]
def build_out(self) -> Data:
cfg = Cfg(vllm_base=self.vllm_base, vllm_model=self.vllm_model)
ctx = stages.run_stage("extract", stages.stage_extract,
dict(self.ctx.data), cfg)
self.status = ctx.get("status", "")
return Data(data=ctx)
- Step 3: Write
components/pm_match.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output
from langflow.schema import Data
from projektmatch import stages
from projektmatch.config import Cfg
class PMMatch(Component):
display_name = "PM 3 Match CV"
description = "LLM 2: Anforderungen gegen Lebenslauf bewerten"
inputs = [
DataInput(name="ctx", display_name="Context"),
MessageTextInput(name="vllm_base", display_name="vLLM Base",
value="PM_VLLM_BASE", load_from_db=True,
advanced=True),
MessageTextInput(name="vllm_model", display_name="vLLM Model",
value="PM_VLLM_MODEL", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Context", method="build_out")]
def build_out(self) -> Data:
cfg = Cfg(vllm_base=self.vllm_base, vllm_model=self.vllm_model)
ctx = stages.run_stage("match", stages.stage_match,
dict(self.ctx.data), cfg)
self.status = ctx.get("status", "")
return Data(data=ctx)
- Step 4: Write
components/pm_rules.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output
from langflow.schema import Data
from projektmatch import stages
from projektmatch.config import Cfg
class PMRules(Component):
display_name = "PM 4 Rules+Gate"
description = ("Deterministisch: Misc-Regeln, Match-Berechnung, "
"Beschreibungs-Markdown, Gate > Schwellwert")
inputs = [
DataInput(name="ctx", display_name="Context"),
MessageTextInput(name="threshold", display_name="Threshold",
value="PM_THRESHOLD", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Context", method="build_out")]
def build_out(self) -> Data:
cfg = Cfg(threshold=int(self.threshold))
ctx = stages.run_stage("rules", stages.stage_rules,
dict(self.ctx.data), cfg)
self.status = f"{ctx.get('decision')} ({ctx.get('mustMatch')} %)"
return Data(data=ctx)
- Step 5: Write
components/pm_crm.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output, SecretStrInput
from langflow.schema import Data
from projektmatch import stages
from projektmatch.config import Cfg
class PMCrm(Component):
display_name = "PM 5 CRM"
description = "EspoCRM: Team, Firma, Kontakt, Verkaufschance + Verify"
inputs = [
DataInput(name="ctx", display_name="Context"),
MessageTextInput(name="espo_base", display_name="Espo Base",
value="PM_ESPO_BASE", load_from_db=True,
advanced=True),
SecretStrInput(name="espo_api_key", display_name="Espo API Key",
value="PM_ESPO_API_KEY", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Context", method="build_out")]
def build_out(self) -> Data:
cfg = Cfg(espo_base=self.espo_base, espo_api_key=self.espo_api_key)
ctx = stages.run_stage("crm", stages.stage_crm,
dict(self.ctx.data), cfg)
self.status = ctx.get("status", "")
return Data(data=ctx)
- Step 6: Write
components/pm_notify.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output, SecretStrInput
from langflow.schema.message import Message
from projektmatch import stages
from projektmatch.config import Cfg
class PMNotify(Component):
display_name = "PM 6 Notify"
description = "Benachrichtigungs-Mail bei created + Langfuse-Trace"
inputs = [
DataInput(name="ctx", display_name="Context"),
SecretStrInput(name="imap_password", display_name="Mail Password",
value="PM_IMAP_PASSWORD", load_from_db=True,
advanced=True),
MessageTextInput(name="notify_to", display_name="Notify To",
value="PM_NOTIFY_TO", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Result", method="build_out")]
def build_out(self) -> Message:
cfg = Cfg(imap_password=self.imap_password, notify_to=self.notify_to)
# stage_notify is called DIRECTLY (not via run_stage): it must run for
# every terminal status so failed runs still post their Langfuse trace.
ctx = stages.stage_notify(dict(self.ctx.data), cfg)
self.status = ctx.get("status", "")
return Message(text=stages.summary(ctx))
- Step 7: Write
components/pm_ingest.py
import sys
if "/app/langflow" not in sys.path:
sys.path.insert(0, "/app/langflow")
import json
from langflow.custom import Component
from langflow.io import DataInput, MessageTextInput, Output, SecretStrInput
from langflow.schema.message import Message
from projektmatch import ingest
from projektmatch.config import Cfg
class PMIngest(Component):
display_name = "PM Ingest"
description = ("IMAP-Postfach abrufen, Projekte splitten, CRM-Dedup, "
"Flow 2 je Projekt, Alert-Mail, Trigger-Mail -> Trash")
inputs = [
DataInput(name="trigger", display_name="Webhook Trigger"),
SecretStrInput(name="imap_password", display_name="Mail Password",
value="PM_IMAP_PASSWORD", load_from_db=True,
advanced=True),
MessageTextInput(name="espo_base", display_name="Espo Base",
value="PM_ESPO_BASE", load_from_db=True,
advanced=True),
SecretStrInput(name="espo_api_key", display_name="Espo API Key",
value="PM_ESPO_API_KEY", load_from_db=True,
advanced=True),
MessageTextInput(name="alert_to", display_name="Alert To",
value="PM_ALERT_TO", load_from_db=True,
advanced=True),
MessageTextInput(name="flow2_id", display_name="Flow 2 ID",
value="PM_FLOW2_ID", load_from_db=True,
advanced=True),
SecretStrInput(name="langflow_api_key", display_name="Langflow Key",
value="PM_LANGFLOW_API_KEY", load_from_db=True,
advanced=True),
]
outputs = [Output(name="out", display_name="Summary", method="build_out")]
def build_out(self) -> Message:
cfg = Cfg(imap_password=self.imap_password,
espo_base=self.espo_base,
espo_api_key=self.espo_api_key,
alert_to=self.alert_to,
flow2_id=self.flow2_id,
langflow_api_key=self.langflow_api_key)
summary = ingest.run_ingest(cfg)
self.status = json.dumps(summary)
return Message(text=json.dumps(summary))
- Step 8: Syntax check all components (host)
Run: cd /home/lwc/bin/projekt-matching && for f in components/*.py; do python3 -m py_compile "$f" && echo "OK $f"; done
Expected: OK for all seven files (imports are NOT executed by py_compile, so missing langflow on the host is fine).
- Step 9: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): seven Langflow wrapper components"
Task 10: Deploy files + build flows via Langflow API
Files:
- Create:
projekt-matching/deploy/deploy_files.sh,projekt-matching/deploy/secrets.local.env(chmod 600, gitignored),projekt-matching/deploy/build_flows.py
Interfaces:
-
Consumes: component files (Task 9), package (Tasks 2–8).
-
Produces: package + vorgaben deployed under
~/.local/share/langflow_pod/langflow-data/; Langflow global variablesPM_*(list in Task 9); flowsPM Projekt bewerten(endpointpm-flow2) andPM Ingest(endpointpm-ingest); Langflow API key stored in~/.config/projekt-matching/envasLANGFLOW_API_KEY=...plusFLOW1_WEBHOOK=http://127.0.0.1:8090/api/v1/webhook/pm-ingest; variablePM_FLOW2_IDset to the real Flow-2 UUID. Task 12/13 consume~/.config/projekt-matching/env. -
Step 1: Write
deploy/deploy_files.sh
#!/bin/bash
# Copy the projektmatch package and vorgaben files into the Langflow data dir
# (mounted at /app/langflow in the container). Idempotent.
set -e
SRC="$(cd "$(dirname "$0")/.." && pwd)"
DEST="$HOME/.local/share/langflow_pod/langflow-data"
rsync -a --delete --exclude __pycache__ "$SRC/projektmatch/" "$DEST/projektmatch/"
mkdir -p "$DEST/vorgaben"
cp /home/tlg/mkt/bewerb/vorgaben/Lebenslauf_Dr-Ing_Thomas_Langer.md \
/home/tlg/mkt/bewerb/vorgaben/rahmenbedingungen.md "$DEST/vorgaben/"
# Langflow runs as uid 1000 gid 0 -> needs group read
chmod -R g+rX "$DEST/projektmatch" "$DEST/vorgaben"
echo "Deployed to $DEST"
Run: chmod +x projekt-matching/deploy/deploy_files.sh && projekt-matching/deploy/deploy_files.sh
Expected: Deployed to /home/lwc/.local/share/langflow_pod/langflow-data, then verify import inside the container:
podman exec langflow_ctr python -c "import sys; sys.path.insert(0,'/app/langflow'); import projektmatch.stages; print('import ok')"
Expected: import ok
- Step 2: Write
deploy/secrets.local.env(gitignored; chmod 600)
PM_IMAP_PASSWORD='<the IMAP password from the project brief>'
PM_ESPO_API_KEY="$(grep -oE '[0-9a-f]{32}' /home/tlg/mkt/bewerb/.secrets/espocrm-api.md | head -1)"
PM_ESPO_BASE="$(grep -oE 'https://[^ `]+/api/v1' /home/tlg/mkt/bewerb/.secrets/espocrm-api.md | head -1)"
Run: chmod 600 projekt-matching/deploy/secrets.local.env && bash -c 'source projekt-matching/deploy/secrets.local.env && [ -n "$PM_ESPO_API_KEY" ] && [ -n "$PM_ESPO_BASE" ] && echo secrets-ok'
Expected: secrets-ok
- Step 3: Inspect the existing smoke-test flow to confirm the edge JSON format
Langflow's edge sourceHandle/targetHandle are JSON strings with " replaced by œ (U+0153). Confirm against the existing flow before building:
TOKEN=$(curl -s --compressed http://127.0.0.1:8090/api/v1/auto_login | python3 -c 'import json,sys; print(json.load(sys.stdin)["access_token"])')
FLOW=$(curl -s --compressed -H "Authorization: Bearer $TOKEN" http://127.0.0.1:8090/api/v1/flows/ | python3 -c 'import json,sys; flows=json.load(sys.stdin); print([f["id"] for f in flows if "Smoke" in f["name"]][0])')
curl -s --compressed -H "Authorization: Bearer $TOKEN" "http://127.0.0.1:8090/api/v1/flows/$FLOW" \
| python3 -c 'import json,sys; e=json.load(sys.stdin)["data"]["edges"][0]; print(json.dumps(e, indent=1)[:900])'
Expected: an edge object with source, target, sourceHandle containing œdataTypeœ etc. If the format differs (plain JSON handles, different keys), adapt handle_str()/make_edge() in Step 4 to mirror EXACTLY what this dump shows.
- Step 4: Write
deploy/build_flows.py
#!/usr/bin/env python3
"""Create Langflow global variables, API key, and both PM flows via API.
Idempotent: variables are upserted; existing flows with the same name are
deleted and recreated. Run AFTER deploy_files.sh. Needs secrets.local.env
sourced into the environment (PM_IMAP_PASSWORD, PM_ESPO_API_KEY,
PM_ESPO_BASE)."""
import json
import os
import sys
from pathlib import Path
import requests
BASE = "http://127.0.0.1:8090/api/v1"
COMPONENTS = Path(__file__).resolve().parent.parent / "components"
ENV_FILE = Path.home() / ".config" / "projekt-matching" / "env"
GENERIC_VARS = {
"PM_ESPO_BASE": os.environ.get("PM_ESPO_BASE", ""),
"PM_NOTIFY_TO": "chancen@destengs.com",
"PM_ALERT_TO": "chancen@destengs.com",
"PM_THRESHOLD": "85",
"PM_VLLM_BASE": "http://host.containers.internal:8081/v1",
"PM_VLLM_MODEL": "AxionML/Qwen3.5-9B-NVFP4",
"PM_FLOW2_ID": "placeholder",
}
SECRET_VARS = ("PM_IMAP_PASSWORD", "PM_ESPO_API_KEY")
def login():
r = requests.get(f"{BASE}/auto_login", timeout=15)
r.raise_for_status()
token = r.json()["access_token"]
return {"Authorization": f"Bearer {token}"}
def upsert_variable(hdr, name, value, vtype):
existing = {v["name"]: v["id"] for v in
requests.get(f"{BASE}/variables/", headers=hdr,
timeout=15).json()}
body = {"name": name, "value": value, "type": vtype,
"default_fields": []}
if name in existing:
requests.patch(f"{BASE}/variables/{existing[name]}", headers=hdr,
json={"id": existing[name], "name": name,
"value": value}, timeout=15).raise_for_status()
else:
requests.post(f"{BASE}/variables/", headers=hdr, json=body,
timeout=15).raise_for_status()
def get_api_key(hdr):
r = requests.post(f"{BASE}/api_key/", headers=hdr,
json={"name": "projekt-matching"}, timeout=15)
r.raise_for_status()
return r.json()["api_key"]
def component_template(hdr, code):
r = requests.post(f"{BASE}/custom_component", headers=hdr,
json={"code": code}, timeout=60)
r.raise_for_status()
out = r.json()
return out["data"], out["type"]
def builtin_template(hdr, category, name):
r = requests.get(f"{BASE}/all", headers=hdr, timeout=60)
r.raise_for_status()
return r.json()[category][name]
def handle_str(d):
return json.dumps(d, separators=(",", ":")).replace('"', "œ")
def make_node(node_id, comp_type, template, x):
return {"id": node_id, "type": "genericNode",
"position": {"x": x, "y": 0},
"data": {"type": comp_type, "id": node_id, "node": template,
"showNode": True}}
def make_edge(src_id, src_type, out_name, out_types, tgt_id, field,
input_types, field_type):
sh = {"dataType": src_type, "id": src_id, "name": out_name,
"output_types": out_types}
th = {"fieldName": field, "id": tgt_id, "inputTypes": input_types,
"type": field_type}
return {"id": f"reactflow__edge-{src_id}-{tgt_id}",
"source": src_id, "target": tgt_id,
"sourceHandle": handle_str(sh), "targetHandle": handle_str(th),
"data": {"sourceHandle": sh, "targetHandle": th},
"className": ""}
def create_flow(hdr, name, endpoint, nodes, edges):
flows = requests.get(f"{BASE}/flows/", headers=hdr, timeout=30).json()
for f in flows:
if f["name"] == name:
requests.delete(f"{BASE}/flows/{f['id']}", headers=hdr,
timeout=15)
body = {"name": name, "endpoint_name": endpoint, "is_component": False,
"data": {"nodes": nodes, "edges": edges,
"viewport": {"x": 0, "y": 0, "zoom": 0.5}}}
r = requests.post(f"{BASE}/flows/", headers=hdr, json=body, timeout=30)
r.raise_for_status()
return r.json()["id"]
def main():
hdr = login()
for name in SECRET_VARS:
upsert_variable(hdr, name, os.environ[name], "Credential")
for name, value in GENERIC_VARS.items():
if name != "PM_FLOW2_ID":
upsert_variable(hdr, name, value, "Generic")
api_key = get_api_key(hdr)
# PMIngest reads the run-API key via load_from_db -> must exist as a
# Credential global variable, not only in the host env file.
upsert_variable(hdr, "PM_LANGFLOW_API_KEY", api_key, "Credential")
def custom(fname):
code = (COMPONENTS / fname).read_text()
return component_template(hdr, code)
# ---- Flow 2: TextInput -> PMFetch -> ... -> PMNotify -> TextOutput
text_in = builtin_template(hdr, "input_output", "TextInput")
text_out = builtin_template(hdr, "input_output", "TextOutput")
chain = [("TextInput-pm2i", "TextInput", text_in)]
for fname, nid in [("pm_fetch.py", "PMFetch-pm2a"),
("pm_extract.py", "PMExtract-pm2b"),
("pm_match.py", "PMMatch-pm2c"),
("pm_rules.py", "PMRules-pm2d"),
("pm_crm.py", "PMCrm-pm2e"),
("pm_notify.py", "PMNotify-pm2f")]:
tpl, ctype = custom(fname)
chain.append((nid, ctype, tpl))
chain.append(("TextOutput-pm2o", "TextOutput", text_out))
nodes2 = [make_node(nid, ctype, tpl, x=i * 380)
for i, (nid, ctype, tpl) in enumerate(chain)]
edges2 = [
make_edge("TextInput-pm2i", "TextInput", "text", ["Message"],
"PMFetch-pm2a", "payload", ["Message"], "str"),
make_edge("PMFetch-pm2a", "PMFetch", "out", ["Data"],
"PMExtract-pm2b", "ctx", ["Data"], "other"),
make_edge("PMExtract-pm2b", "PMExtract", "out", ["Data"],
"PMMatch-pm2c", "ctx", ["Data"], "other"),
make_edge("PMMatch-pm2c", "PMMatch", "out", ["Data"],
"PMRules-pm2d", "ctx", ["Data"], "other"),
make_edge("PMRules-pm2d", "PMRules", "out", ["Data"],
"PMCrm-pm2e", "ctx", ["Data"], "other"),
make_edge("PMCrm-pm2e", "PMCrm", "out", ["Data"],
"PMNotify-pm2f", "ctx", ["Data"], "other"),
make_edge("PMNotify-pm2f", "PMNotify", "out", ["Message"],
"TextOutput-pm2o", "input_value", ["Message"], "str"),
]
flow2_id = create_flow(hdr, "PM Projekt bewerten", "pm-flow2",
nodes2, edges2)
upsert_variable(hdr, "PM_FLOW2_ID", flow2_id, "Generic")
# ---- Flow 1: Webhook -> PMIngest -> TextOutput
webhook = builtin_template(hdr, "data", "Webhook")
tpl_ing, ctype_ing = custom("pm_ingest.py")
nodes1 = [make_node("Webhook-pm1w", "Webhook", webhook, 0),
make_node("PMIngest-pm1a", ctype_ing, tpl_ing, 380),
make_node("TextOutput-pm1o", "TextOutput", text_out, 760)]
edges1 = [
make_edge("Webhook-pm1w", "Webhook", "output_data", ["Data"],
"PMIngest-pm1a", "trigger", ["Data"], "other"),
make_edge("PMIngest-pm1a", "PMIngest", "out", ["Message"],
"TextOutput-pm1o", "input_value", ["Message"], "str"),
]
flow1_id = create_flow(hdr, "PM Ingest", "pm-ingest", nodes1, edges1)
ENV_FILE.parent.mkdir(parents=True, exist_ok=True)
ENV_FILE.write_text(
f"LANGFLOW_API_KEY={api_key}\n"
f"FLOW1_ID={flow1_id}\nFLOW2_ID={flow2_id}\n"
f"FLOW1_WEBHOOK=http://127.0.0.1:8090/api/v1/webhook/pm-ingest\n")
ENV_FILE.chmod(0o600)
print(f"flow1={flow1_id} flow2={flow2_id}\nenv -> {ENV_FILE}")
if __name__ == "__main__":
sys.exit(main())
Note for the implementer: the exact builtin_template category keys ("input_output", "data") and output names ("text", "output_data") must be verified against GET /api/v1/all and the Step-3 edge dump — Langflow 1.10 sometimes nests templates as {"template": ...}. Adjust the two helper functions to match reality; the node/edge SHAPE built here mirrors what Step 3 shows.
- Step 5: Run deployment
cd /home/lwc/bin/projekt-matching
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python deploy/build_flows.py
Expected: flow1=<uuid> flow2=<uuid> and env -> /home/lwc/.config/projekt-matching/env. Open http://127.0.0.1:8090 manually later if debugging is needed — both flows must render without red error nodes.
- Step 6: Verify Flow 2 runs end-to-end against a REAL freelancermap page (no CRM yet: use a rejected-style URL)
source ~/.config/projekt-matching/env
curl -s --compressed -X POST \
"http://127.0.0.1:8090/api/v1/run/$FLOW2_ID?stream=false" \
-H "x-api-key: $LANGFLOW_API_KEY" -H 'Content-Type: application/json' \
-d '{"input_value": "{\"canonical\": \"https://www.freelancermap.de/projekt/test-engineer-m-w-d-3020995\", \"title\": \"Test Engineer\"}", "input_type": "text", "output_type": "text"}' \
| python3 -m json.tool | grep -E '"(status|mustMatch|decision)"' | head
Expected: run completes (~1–5 min, reasoning model); the summary JSON shows "status": "rejected" (Must-have-Match < 50 % for the Test-Engineer project) — proving fetch → extract → match → rules → gate work in-container. If it errors, read podman logs --tail 100 langflow_ctr.
- Step 7: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): deployment scripts and programmatic flow builder"
Task 11: Langfuse — dedicated project, keys, score configs
Files:
- Modify:
/home/lwc/bin/create_pod_langflow.sh(lines 61–75 area and langfuse-web env) - Create:
projekt-matching/deploy/setup_langfuse.py
Interfaces:
-
Produces: Langfuse org
projekt-matching-org+ projectprojekt-matchingwith API key pair; Langflow container envLANGFUSE_PUBLIC_KEY/SECRET_KEYswitched to the new keys (auto-tracing of all flow runs goes to the new project);pm_trace.post_trace(Task 7) inherits the same envs; score configsextraction-correctandmatching-correct. -
Step 1: Add an admin API key and parameterized PM keys to the pod script
In create_pod_langflow.sh, after line 66 (LANGFUSE_INIT_PASSWORD=...) add:
# Admin key for the Organization Management API (used once by
# projekt-matching/deploy/setup_langfuse.py to create the projekt-matching
# project); PM_* keys are the Langfuse project keys Langflow traces to.
LANGFUSE_ADMIN_API_KEY="admin-pm-3f61c2a89d4e"
PM_LANGFUSE_PUBLIC_KEY="$LANGFUSE_INIT_PUBLIC_KEY"
PM_LANGFUSE_SECRET_KEY="$LANGFUSE_INIT_SECRET_KEY"
Change the two Langflow-container lines
-e LANGFUSE_PUBLIC_KEY="$LANGFUSE_INIT_PUBLIC_KEY" → -e LANGFUSE_PUBLIC_KEY="$PM_LANGFUSE_PUBLIC_KEY" and
-e LANGFUSE_SECRET_KEY="$LANGFUSE_INIT_SECRET_KEY" → -e LANGFUSE_SECRET_KEY="$PM_LANGFUSE_SECRET_KEY".
In the langfuse-web container run command add: -e ADMIN_API_KEY="$LANGFUSE_ADMIN_API_KEY".
Run: bash /home/lwc/bin/create_pod_langflow.sh (recreates the pod, ~2 min; data persists in Postgres).
Expected: ends with Langflow Web UI is reachable....
- Step 2: Write
deploy/setup_langfuse.py
#!/usr/bin/env python3
"""Create Langfuse org/project/keys + score configs for projekt-matching.
Uses the Organization Management API (ADMIN_API_KEY) — Langfuse v3.195.
Prints the project keys; Step 3 wires them into the pod script."""
import sys
import requests
BASE = "http://127.0.0.1:8091"
ADMIN = {"Authorization": "Bearer admin-pm-3f61c2a89d4e"}
def post(path, json, auth=None, headers=None):
r = requests.post(f"{BASE}{path}", json=json, auth=auth,
headers=headers, timeout=30)
if r.status_code == 404:
sys.exit(f"ENDPOINT MISSING: {path} — create org/project/keys "
f"manually in the UI (login admin@example.com), then rerun "
f"score-config part with the new keys.")
r.raise_for_status()
return r.json()
def main():
org = post("/api/admin/organizations",
{"name": "projekt-matching-org"}, headers=ADMIN)
org_key = post(f"/api/admin/organizations/{org['id']}/apiKeys", {},
headers=ADMIN)
org_auth = (org_key["publicKey"], org_key["secretKey"])
project = post("/api/public/projects",
{"name": "projekt-matching", "retention": 0},
auth=org_auth)
keys = post(f"/api/public/projects/{project['id']}/apiKeys", {},
auth=org_auth)
proj_auth = (keys["publicKey"], keys["secretKey"])
for name, desc in (
("extraction-correct",
"Must/Nice/Misc korrekt aus der Ausschreibung abgeleitet?"),
("matching-correct",
"✅/❌/❔-Bewertungen gegen den Lebenslauf korrekt?")):
post("/api/public/score-configs",
{"name": name, "dataType": "CATEGORICAL", "description": desc,
"categories": [{"label": "correct", "value": 1},
{"label": "partially-correct", "value": 0.5},
{"label": "wrong", "value": 0}]},
auth=proj_auth)
print(f"PM_LANGFUSE_PUBLIC_KEY={keys['publicKey']}")
print(f"PM_LANGFUSE_SECRET_KEY={keys['secretKey']}")
if __name__ == "__main__":
main()
Run: cd /home/lwc/bin/projekt-matching && .venv/bin/python deploy/setup_langfuse.py
Expected: two PM_LANGFUSE_*=pk-lf-.../sk-lf-... lines. Fallback if an admin endpoint is missing in this Langfuse version: log into http://127.0.0.1:8091 as admin@example.com / password from the pod script, create org projekt-matching-org + project projekt-matching + API keys in the UI (~1 min), create the two score configs under Settings → Score Configs with the categories above, and use those keys in Step 3.
- Step 3: Wire the new keys into the pod script and restart
Edit create_pod_langflow.sh: set PM_LANGFUSE_PUBLIC_KEY="pk-lf-..." and PM_LANGFUSE_SECRET_KEY="sk-lf-..." to the printed values. Run bash /home/lwc/bin/create_pod_langflow.sh again.
Verify: podman exec langflow_ctr env | grep LANGFUSE shows the NEW keys, and after Task 13's first Flow-2 run a trace appears in project projekt-matching at http://127.0.0.1:8091.
- Step 4: Annotation queue (UI, one minute — API for queues is not stable across versions)
In the Langfuse UI, project projekt-matching → Annotation → Queues → New queue projekt-matching-review, attach both score configs. This is Thomas's manual-evaluation workflow (spec §6): open queue, score extraction-correct and matching-correct per trace.
- Step 5: Commit
cd /home/lwc/bin && git add create_pod_langflow.sh projekt-matching && git commit -m "feat(projekt-matching): langfuse project, keys, score configs; admin API in pod script"
Task 12: systemd timer + webhook trigger
Files:
- Create:
projekt-matching/deploy/trigger_webhook.sh,projekt-matching/deploy/projekt-matching.service,projekt-matching/deploy/projekt-matching.timer,projekt-matching/tests/e2e/flag_inbox.py
Interfaces:
-
Consumes:
~/.config/projekt-matching/env(Task 10). -
Produces: user units
projekt-matching.timer/.servicefiring every 5 min. -
Step 1: Write
deploy/trigger_webhook.sh
#!/bin/bash
# Fire the PM Ingest webhook. Carries no logic; failures are harmless
# (state lives in the IMAP inbox, next tick retries).
set -u
source "$HOME/.config/projekt-matching/env"
curl -sS -m 15 -X POST "$FLOW1_WEBHOOK" \
-H "x-api-key: $LANGFLOW_API_KEY" -H 'Content-Type: application/json' \
-d '{"source": "systemd-timer"}' || true
echo
deploy/projekt-matching.service:
[Unit]
Description=Trigger PM Ingest Langflow webhook
[Service]
Type=oneshot
ExecStart=%h/bin/projekt-matching/deploy/trigger_webhook.sh
deploy/projekt-matching.timer:
[Unit]
Description=Poll chancen@ inbox via PM Ingest every 5 minutes
[Timer]
OnCalendar=*:00/5
Persistent=false
[Install]
WantedBy=timers.target
- Step 2: Write
tests/e2e/flag_inbox.py— mark ALL current INBOX mails$ProjektCheckedso the system starts fresh (pre-existing mails, including the original example mails, are never auto-processed; they stay in INBOX for reference):
#!/usr/bin/env python3
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from projektmatch import mailer # noqa: E402
box = mailer.MailBox("chancen@destengs.com", os.environ["PM_IMAP_PASSWORD"])
uids = box.unchecked_uids()
for uid in uids:
box.flag_checked(uid)
print(f"flagged {len(uids)} mails as checked")
box.close()
Run:
cd /home/lwc/bin/projekt-matching
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python tests/e2e/flag_inbox.py
Expected: flagged <n> mails as checked.
- Step 3: Install and enable the timer; one manual trigger first
chmod +x projekt-matching/deploy/trigger_webhook.sh
projekt-matching/deploy/trigger_webhook.sh # manual run: expect {"message":...} JSON, no error
podman logs --tail 30 langflow_ctr # PM Ingest ran; summary {"mails": 0, ...}
cp projekt-matching/deploy/projekt-matching.{service,timer} ~/.config/systemd/user/
systemctl --user daemon-reload
systemctl --user enable --now projekt-matching.timer
systemctl --user list-timers | grep projekt-matching
Expected: timer listed with NEXT within 5 minutes. After the next tick: journalctl --user -u projekt-matching.service -n 5 shows the curl output.
- Step 4: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "feat(projekt-matching): systemd timer, webhook trigger, inbox baseline script"
Task 13: E2E helpers + component-level gate tests (Flow 2 direct)
Files:
- Create:
projekt-matching/tests/e2e/cleanup_crm.py,projekt-matching/tests/e2e/run_flow2.py,projekt-matching/tests/e2e/check_state.py
Interfaces:
- Consumes:
~/.config/projekt-matching/env,deploy/secrets.local.env, package modules. - Produces: reusable scripts for Tasks 14–16.
Known references: example 1 project 4 canonical https://www.freelancermap.de/projekt/python-entwickler-fuer-eine-ki-anwendung-m-w-d (skill scored 86 % Must); example 2 canonical https://www.freelancermap.de/projekt/test-engineer-m-w-d-3020995 (25 % Must).
- Step 1: Write
tests/e2e/cleanup_crm.py
#!/usr/bin/env python3
"""Delete test opportunities whose cProjektlink contains a token."""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from projektmatch.espocrm import EspoClient # noqa: E402
TOKENS = sys.argv[1:] or [
"python-entwickler-fuer-eine-ki-anwendung",
"test-engineer-m-w-d-3020995",
"nproj/3020338", "nproj/3020995"]
espo = EspoClient(os.environ["PM_ESPO_BASE"], os.environ["PM_ESPO_API_KEY"])
for token in TOKENS:
for hit in espo.search("Opportunity", "contains", "cProjektlink", token,
select="name,cProjektlink", max_size=100):
espo.delete("Opportunity", hit["id"])
print(f"deleted {hit['id']} {hit['name']}")
print("cleanup done")
- Step 2: Write
tests/e2e/run_flow2.py
#!/usr/bin/env python3
"""Run Flow 2 directly via the Langflow run API from the host."""
import json
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from projektmatch.config import Cfg # noqa: E402
from projektmatch.ingest import dispatch_flow2 # noqa: E402
env = dict(line.split("=", 1) for line in open(
os.path.expanduser("~/.config/projekt-matching/env"))
if "=" in line.strip())
cfg = Cfg(flow2_id=env["FLOW2_ID"].strip(),
langflow_api_key=env["LANGFLOW_API_KEY"].strip(),
langflow_base="http://127.0.0.1:8090")
result = dispatch_flow2(cfg, {"canonical": sys.argv[1],
"title": sys.argv[2] if len(sys.argv) > 2
else "Test"})
print(json.dumps(result, indent=2, ensure_ascii=False))
- Step 3: Write
tests/e2e/check_state.py
#!/usr/bin/env python3
"""Report current system state: CRM hits for a token, INBOX/Trash mails."""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from projektmatch import mailer # noqa: E402
from projektmatch.espocrm import EspoClient # noqa: E402
token = sys.argv[1] if len(sys.argv) > 1 else ""
espo = EspoClient(os.environ["PM_ESPO_BASE"], os.environ["PM_ESPO_API_KEY"])
if token:
hits = espo.search("Opportunity", "contains", "cProjektlink", token,
select="name,cProjektlink")
print(f"CRM hits for '{token}': {len(hits)}")
for h in hits:
print(f" {h['id']} {h['name']}")
box = mailer.MailBox("chancen@destengs.com", os.environ["PM_IMAP_PASSWORD"])
for folder in ("INBOX", box.trash_folder()):
box.conn.select(folder)
typ, data = box.conn.uid("SEARCH", None, "ALL")
uids = data[0].split() if data and data[0] else []
print(f"{folder}: {len(uids)} mails")
for uid in uids[-8:]:
typ, d = box.conn.uid("FETCH", uid.decode(),
"(BODY.PEEK[HEADER.FIELDS (SUBJECT)])")
print(" ", d[0][1].decode(errors="replace").strip())
box.close()
- Step 4: Component gate test — example 1 project 4 (consider path)
cd /home/lwc/bin/projekt-matching
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python tests/e2e/cleanup_crm.py
.venv/bin/python tests/e2e/run_flow2.py \
"https://www.freelancermap.de/projekt/python-entwickler-fuer-eine-ki-anwendung-m-w-d" \
"Python Entwickler für eine KI-Anwendung (m/w/d)"
Expected JSON: "status": "created", "mustMatch" > 85, "crmUrl": "https://crm.creature-go.com/#Opportunity/view/<id>". Then verify:
-
.venv/bin/python tests/e2e/check_state.py python-entwickler-fuer-eine-ki-anwendung→CRM hits: 1; open the crmUrl description — exact table format (match line,| Nr. | Kat. | ❔ | Anforderung |, Must→Nice→Misc). -
A
[Projekt-Match] Freelancermap — <Must> % — ...mail is in INBOX (design-phase recipient chancen@) — visible in check_state output. -
A
projekt-matchtrace exists:curl -s -u "pk-lf-...:sk-lf-..." "http://127.0.0.1:8091/api/public/traces?name=projekt-match&limit=3"(use the Task-11 keys) → the run appears with output.decision "consider". -
If
mustMatch≤ 85: inspect the trace's table — usually a Must item was miscategorized or over-strictly rated. TuneEXTRACT_SYSTEM/MATCH_SYSTEMprompts (Task 5), redeploy (deploy/deploy_files.sh+ restart langflow containerpodman restart langflow_ctr), rerun. Iterate until stable > 85 across 2 consecutive runs. -
Step 5: Component gate test — example 2 (reject path)
.venv/bin/python tests/e2e/cleanup_crm.py
.venv/bin/python tests/e2e/run_flow2.py \
"https://www.freelancermap.de/projekt/test-engineer-m-w-d-3020995" "Test Engineer (m/w/d)"
Expected: "status": "rejected", "mustMatch" < 50, no crmUrl; check_state.py test-engineer → CRM hits: 0; NO new [Projekt-Match] mail; trace tagged rejected.
- Step 6: Clean up and commit
.venv/bin/python tests/e2e/cleanup_crm.py
cd /home/lwc/bin && git add projekt-matching && git commit -m "test(projekt-matching): e2e helper scripts; component gates verified"
Task 14: Gate 1 — end-to-end consider path via e-mail
Files:
-
Create:
projekt-matching/tests/e2e/send_test_mail.py -
Step 1: Write
tests/e2e/send_test_mail.py
#!/usr/bin/env python3
"""Re-send an original example mail (found by subject) as a fresh test mail
from chancen@ to chancen@. Fallback: synthetic single-project mail."""
import os
import smtplib
import ssl
import sys
import time
from email.message import EmailMessage
from email.policy import default as default_policy
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
from projektmatch import mailer, mailparse # noqa: E402
USER = "chancen@destengs.com"
PW = os.environ["PM_IMAP_PASSWORD"]
SUBJECT = sys.argv[1] # e.g. "Example 1 for Langflow process"
FALLBACK_URL = sys.argv[2] if len(sys.argv) > 2 else ""
box = mailer.MailBox(USER, PW)
html = ""
for folder in ("INBOX", box.trash_folder()):
box.conn.select(folder)
typ, data = box.conn.uid("SEARCH", None,
f'SUBJECT "{SUBJECT}"')
uids = data[0].split() if data and data[0] else []
if uids:
msg = box.fetch(uids[-1].decode())
html, _ = mailparse.bodies(msg)
break
box.close()
out = EmailMessage(policy=default_policy)
out["From"], out["To"] = USER, USER
out["Subject"] = f"Test {SUBJECT} {int(time.time())}"
if html:
out.set_content("siehe HTML")
out.add_alternative(html, subtype="html")
print(f"re-sending original '{SUBJECT}' body")
else:
if not FALLBACK_URL:
sys.exit(f"original mail '{SUBJECT}' not found and no fallback URL")
out.set_content(f"Neues Projekt:\n\nTestprojekt\n{FALLBACK_URL}\n")
print("original not found -> synthetic fallback mail")
with smtplib.SMTP(mailer.SMTP_HOST, mailer.SMTP_PORT, timeout=30) as server:
server.starttls(context=ssl.create_default_context())
server.login(USER, PW)
server.send_message(out)
print(f"sent: {out['Subject']}")
- Step 2: Run Gate 1
cd /home/lwc/bin/projekt-matching
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python tests/e2e/cleanup_crm.py
.venv/bin/python tests/e2e/send_test_mail.py "Example 1 for Langflow process" \
"https://www.freelancermap.de/nproj/3020338.html?utm_source=systemmail&utm_medium=email&utm_campaign=projektagent&agent=233342&t=1783432300&html=1"
Wait ≤ 10 min (timer + processing; 5 projects × 2 LLM calls). Then:
.venv/bin/python tests/e2e/check_state.py python-entwickler-fuer-eine-ki-anwendung
Gate 1 passes when ALL of: CRM hits for the python-entwickler link = 1; a [Projekt-Match] Freelancermap — <Must> % — Python Entwickler... mail (Must > 85) is in INBOX; the Test Example 1... trigger mail is in Trash (not INBOX); Langfuse shows one projekt-match trace per project in the mail (5 with original body, 1 with fallback), python-entwickler tagged created, others rejected/created per their real match. If the other four projects create unwanted CRM entries, that is CORRECT behavior (they passed the gate); note them for Thomas.
- Step 3: Commit
cd /home/lwc/bin && git add projekt-matching && git commit -m "test(projekt-matching): gate 1 e2e passed (consider path via e-mail)"
Task 15: Gate 2 (reject path) + duplicate suppression
- Step 1: Gate 2 — reject path
cd /home/lwc/bin/projekt-matching
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python tests/e2e/cleanup_crm.py test-engineer-m-w-d-3020995 nproj/3020995
.venv/bin/python tests/e2e/send_test_mail.py "Example 2 for Langflow process" \
"https://www.freelancermap.de/nproj/3020995.html?utm_source=systemmail&utm_medium=email&utm_campaign=projektagent&agent=233344&t=1783520448&html=1"
Wait ≤ 10 min, then .venv/bin/python tests/e2e/check_state.py test-engineer.
Gate 2 passes when ALL of: CRM hits for test-engineer = 0; NO new [Projekt-Match] mail since Gate 1; trigger mail in Trash; Langfuse trace shows the full table + matches with Must < 50 and tag rejected.
- Step 2: Duplicate suppression
Re-send the Gate-1 mail WITHOUT cleanup:
.venv/bin/python tests/e2e/send_test_mail.py "Example 1 for Langflow process" \
"https://www.freelancermap.de/nproj/3020338.html?utm_source=systemmail&utm_medium=email&utm_campaign=projektagent&agent=233342&t=1783432300&html=1"
Wait ≤ 10 min, then check_state.py python-entwickler-fuer-eine-ki-anwendung.
Expected: CRM hits STILL 1 (no second opportunity, no (2) suffix entry), no additional notification mail, trigger mail in Trash.
- Step 3: Commit
cd /home/lwc/bin && git add -A projekt-matching && git commit -m "test(projekt-matching): gate 2 and dedup e2e passed"
Task 16: Recipient flip + Gate 3 (production recipient)
- Step 1: Flip notify + alert recipients to Thomas.Langer@destengs.com
cd /home/lwc/bin/projekt-matching && .venv/bin/python - <<'PY'
import requests
BASE = "http://127.0.0.1:8090/api/v1"
hdr = {"Authorization": "Bearer " + requests.get(
f"{BASE}/auto_login", timeout=15).json()["access_token"]}
for v in requests.get(f"{BASE}/variables/", headers=hdr, timeout=15).json():
if v["name"] in ("PM_NOTIFY_TO", "PM_ALERT_TO"):
requests.patch(f"{BASE}/variables/{v['id']}", headers=hdr,
json={"id": v["id"], "name": v["name"],
"value": "Thomas.Langer@destengs.com"},
timeout=15).raise_for_status()
print(f"{v['name']} -> Thomas.Langer@destengs.com")
PY
Expected: both variables reported flipped. Design phase ends here.
- Step 2: Gate 3 — final consider run with production recipient
set -a && source deploy/secrets.local.env && set +a
.venv/bin/python tests/e2e/cleanup_crm.py
.venv/bin/python tests/e2e/send_test_mail.py "Example 1 for Langflow process" \
"https://www.freelancermap.de/nproj/3020338.html?utm_source=systemmail&utm_medium=email&utm_campaign=projektagent&agent=233342&t=1783432300&html=1"
Wait ≤ 10 min, then check_state.py python-entwickler-fuer-eine-ki-anwendung.
Gate 3 passes when: CRM hit = 1 with Must > 85; trigger mail in Trash; NO [Projekt-Match] mail in chancen@ INBOX (recipient is now Thomas); the Langfuse trace/Langflow log shows stage_notify completed without error (SMTP accepted). Ask Thomas to confirm receipt at Thomas.Langer@destengs.com — the automation cannot read that inbox.
-
Step 3: Final housekeeping
-
Leave the Gate-3 opportunity in the CRM (it is a real, correctly-evaluated project) unless Thomas wants it removed.
-
Update memory file
langflow-pod-gotchas.mdif new gotchas surfaced (e.g. edge-format deviations, IMAP keyword support). -
cd /home/lwc/bin && git add -A && git commit -m "feat(projekt-matching): gate 3 passed — production recipient live"
Self-Review Notes (already applied)
- Spec coverage: §3 architecture → Tasks 9–12; §4 Flow 1 → Task 8; §5 Flow 2 → Tasks 2–7; §6 Langfuse → Tasks 7 (pm_trace) + 11; §7 config/secrets → Tasks 6 + 10; §8 error handling → Tasks 5 (LLM retries), 7 (run_stage/fetch retries), 8 (terminal states, alert aggregation, lock); §9 test plan → Tasks 13–16; §10 out-of-scope respected (no snippet fallback, freelancermap-only detection).
- Flow 1 is one
PMIngestcomponent (not three) because IMAP connections and per-mail state cannot cross Langflow component boundaries cleanly; the spec's Flow-1 steps all live iningest.run_ingest, which is fully unit-tested. Flow 2 keeps the visible 6-stage pipeline. - Types checked: ctx dict keys consistent across stages/components/ingest (
status/error/decision/mustMatch/niceMatch/description/opportunityId/crmUrl/projectName/canonical/title);Cfgfields match component wiring;rules.YES/NO/UNKNOWNstrings matchllm.MATCH_SCHEMAenum. - Known verification points intentionally left to execution (marked in steps): Langflow 1.10 edge/handle exact encoding (Task 10 Step 3), Langfuse admin endpoints (Task 11 Step 2 fallback), IMAP custom-keyword support (Task 6 Step 6 fallback).