deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/runtime/scheduler/element.py
2026-09-07 18:24:55 +08:00

171 lines
6.0 KiB
Python

"""Element / Matrix delivery helpers for scheduled task results."""
from __future__ import annotations
import logging
import os
import uuid
from pathlib import Path
import httpx
logger = logging.getLogger(__name__)
class ElementDeliveryClient:
def __init__(self, *, homeserver_url: str, access_token: str) -> None:
self.homeserver_url = homeserver_url.rstrip("/")
self.access_token = access_token
@classmethod
def from_env(cls) -> ElementDeliveryClient | None:
enabled = os.environ.get("ELEMENT_ENABLED", "true").strip().lower()
if enabled in {"0", "false", "no", "off"}:
return None
homeserver_url = os.environ.get("ELEMENT_HOMESERVER_URL", "").strip()
if not homeserver_url:
return None
access_token = os.environ.get("ELEMENT_ACCESS_TOKEN", "").strip()
if not access_token:
access_token = _load_cached_access_token()
if not access_token:
try:
access_token = _bootstrap_access_token(homeserver_url)
except Exception:
logger.exception("Failed to bootstrap Element bot access token")
return None
if not access_token:
return None
return cls(homeserver_url=homeserver_url, access_token=access_token)
def _headers(self) -> dict[str, str]:
return {"Authorization": f"Bearer {self.access_token}"}
async def send_direct_message(self, element_user_id: str, body: str, *, room_id: str | None = None) -> str:
if not room_id:
room_id = await self._create_dm_room(element_user_id)
await self._send_message(room_id, body)
return room_id
async def _create_dm_room(self, element_user_id: str) -> str:
url = f"{self.homeserver_url}/_matrix/client/v3/createRoom"
payload = {
"preset": "trusted_private_chat",
"is_direct": True,
"invite": [element_user_id],
}
async with httpx.AsyncClient(timeout=20) as client:
response = await client.post(url, headers=self._headers(), json=payload)
response.raise_for_status()
data = response.json()
room_id = data.get("room_id")
if not isinstance(room_id, str) or not room_id:
raise RuntimeError("Matrix createRoom response did not include room_id")
return room_id
async def _send_message(self, room_id: str, body: str) -> None:
txn_id = uuid.uuid4().hex
url = f"{self.homeserver_url}/_matrix/client/v3/rooms/{room_id}/send/m.room.message/{txn_id}"
payload = {
"msgtype": "m.text",
"body": body,
}
async with httpx.AsyncClient(timeout=30) as client:
response = await client.put(url, headers=self._headers(), json=payload)
response.raise_for_status()
def _state_file() -> Path:
path = os.environ.get("ELEMENT_STATE_FILE", "").strip()
if path:
return Path(path)
return Path(".element-bot-token")
def _load_cached_access_token() -> str | None:
path = _state_file()
if not path.exists():
return None
token = path.read_text(encoding="utf-8").strip()
return token or None
def _save_cached_access_token(token: str) -> None:
path = _state_file()
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(token, encoding="utf-8")
def _bootstrap_access_token(homeserver_url: str) -> str | None:
username = os.environ.get("ELEMENT_BOT_USERNAME", "scheduled-bot").strip() or "scheduled-bot"
password = os.environ.get("ELEMENT_BOT_PASSWORD", "").strip()
auto_register = os.environ.get("ELEMENT_AUTO_REGISTER", "false").strip().lower() in {
"1",
"true",
"yes",
"on",
}
if auto_register:
if not password:
password = f"scheduled-bot-{uuid.uuid4().hex}"
token = _register_bot(homeserver_url, username, password)
if token:
_save_cached_access_token(token)
return token
if username and password:
token = _login_bot(homeserver_url, username, password)
if token:
_save_cached_access_token(token)
return token
return None
def _register_bot(homeserver_url: str, username: str, password: str) -> str | None:
url = f"{homeserver_url.rstrip('/')}/_matrix/client/v3/register"
payload = {
"username": username,
"password": password,
"auth": {"type": "m.login.dummy"},
"initial_device_display_name": "Scheduled Task Bot",
}
with httpx.Client(timeout=20) as client:
response = client.post(url, json=payload)
if response.status_code == 401:
session = response.json().get("session")
if session:
payload["auth"]["session"] = session
response = client.post(url, json=payload)
if response.status_code == 400 and response.json().get("errcode") == "M_USER_IN_USE":
password_env = os.environ.get("ELEMENT_BOT_PASSWORD", "").strip()
if password_env:
return _login_bot(homeserver_url, username, password_env)
logger.error(
"Element bot username %s already exists and no ELEMENT_BOT_PASSWORD was provided",
username,
)
return None
response.raise_for_status()
data = response.json()
token = data.get("access_token")
return token if isinstance(token, str) and token else None
def _login_bot(homeserver_url: str, username: str, password: str) -> str | None:
url = f"{homeserver_url.rstrip('/')}/_matrix/client/v3/login"
payload = {
"type": "m.login.password",
"identifier": {"type": "m.id.user", "user": username},
"password": password,
"initial_device_display_name": "Scheduled Task Bot",
}
with httpx.Client(timeout=20) as client:
response = client.post(url, json=payload)
response.raise_for_status()
data = response.json()
token = data.get("access_token")
return token if isinstance(token, str) and token else None