mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-08-04 04:28:49 -04:00
Compare commits
21 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fb8c391a88 | |||
| 0de76c4056 | |||
| 25c9e735ef | |||
| 28c333e647 | |||
| 84709a00d9 | |||
| 578312200a | |||
| f23221420f | |||
| 6a84398e75 | |||
| 3250a4ce68 | |||
| cb0f6af002 | |||
| 9297bed5b9 | |||
| 2e631ad816 | |||
| d183fe545b | |||
| 9914651cc9 | |||
| 46905ab9b0 | |||
| 61c138d9e7 | |||
| 25a4d134b1 | |||
| 98e4d8451b | |||
| 5104a9a967 | |||
| 01790c2f08 | |||
| d96c7af3df |
@@ -189,6 +189,7 @@ SEARXNG_INSTANCE=http://localhost:8080
|
|||||||
# ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=26214400 # email compose attachment (25 MB)
|
# ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=26214400 # email compose attachment (25 MB)
|
||||||
# ODYSSEUS_STT_MAX_AUDIO_BYTES=26214400 # speech-to-text audio (25 MB)
|
# ODYSSEUS_STT_MAX_AUDIO_BYTES=26214400 # speech-to-text audio (25 MB)
|
||||||
# ODYSSEUS_ICS_MAX_BYTES=10485760 # calendar .ics import (10 MB)
|
# ODYSSEUS_ICS_MAX_BYTES=10485760 # calendar .ics import (10 MB)
|
||||||
|
# ODYSSEUS_TTS_CACHE_MAX_BYTES=524288000 # TTS cache (500 MB)
|
||||||
|
|
||||||
# ============================================================
|
# ============================================================
|
||||||
# Host Docker access (explicit opt-in)
|
# Host Docker access (explicit opt-in)
|
||||||
|
|||||||
@@ -153,6 +153,16 @@ module.exports = async ({ github, context, core }) => {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const LABEL_BAD = 'needs more info';
|
||||||
|
const LABEL_GOOD = 'ready for review';
|
||||||
|
|
||||||
|
// Closed issues are no longer awaiting review.
|
||||||
|
// This also prevents later edits to closed issues from restoring the label.
|
||||||
|
if (issue.state === 'closed') {
|
||||||
|
await dropLabel(LABEL_GOOD);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
// ── Find existing bot comment to update in-place ──────────────────────────
|
// ── Find existing bot comment to update in-place ──────────────────────────
|
||||||
const MARKER = '<!-- issue-description-check -->';
|
const MARKER = '<!-- issue-description-check -->';
|
||||||
const { data: comments } = await github.rest.issues.listComments({
|
const { data: comments } = await github.rest.issues.listComments({
|
||||||
@@ -160,9 +170,6 @@ module.exports = async ({ github, context, core }) => {
|
|||||||
});
|
});
|
||||||
const existing = comments.find(c => c.user.type === 'Bot' && c.body.includes(MARKER));
|
const existing = comments.find(c => c.user.type === 'Bot' && c.body.includes(MARKER));
|
||||||
|
|
||||||
const LABEL_BAD = 'needs more info';
|
|
||||||
const LABEL_GOOD = 'ready for review';
|
|
||||||
|
|
||||||
if (failures.length === 0) {
|
if (failures.length === 0) {
|
||||||
if (existing) {
|
if (existing) {
|
||||||
await github.rest.issues.deleteComment({ owner, repo, comment_id: existing.id });
|
await github.rest.issues.deleteComment({ owner, repo, comment_id: existing.id });
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ name: ci / issue description check
|
|||||||
|
|
||||||
on:
|
on:
|
||||||
issues:
|
issues:
|
||||||
types: [opened, edited, reopened]
|
types: [opened, edited, reopened, closed]
|
||||||
|
|
||||||
permissions:
|
permissions:
|
||||||
issues: write
|
issues: write
|
||||||
|
|||||||
@@ -692,7 +692,7 @@ from routes.history.history_routes import setup_history_routes
|
|||||||
app.include_router(setup_history_routes(session_manager, upload_handler=upload_handler))
|
app.include_router(setup_history_routes(session_manager, upload_handler=upload_handler))
|
||||||
|
|
||||||
# Search
|
# Search
|
||||||
from routes.search_routes import setup_search_routes
|
from routes.search.search_routes import setup_search_routes
|
||||||
app.include_router(setup_search_routes(config))
|
app.include_router(setup_search_routes(config))
|
||||||
|
|
||||||
# Presets
|
# Presets
|
||||||
@@ -820,7 +820,7 @@ set_ai_rag_manager(rag_manager, personal_docs_mgr)
|
|||||||
logger.info("AI interaction tools initialized (session, memory, RAG, UI control)")
|
logger.info("AI interaction tools initialized (session, memory, RAG, UI control)")
|
||||||
|
|
||||||
# Webhooks
|
# Webhooks
|
||||||
from routes.webhook_routes import setup_webhook_routes
|
from routes.webhook.webhook_routes import setup_webhook_routes
|
||||||
app.include_router(setup_webhook_routes(webhook_manager, auth_manager, session_manager, api_key_manager))
|
app.include_router(setup_webhook_routes(webhook_manager, auth_manager, session_manager, api_key_manager))
|
||||||
|
|
||||||
# API Tokens
|
# API Tokens
|
||||||
@@ -852,7 +852,7 @@ app.include_router(setup_codex_routes(
|
|||||||
))
|
))
|
||||||
app.include_router(setup_claude_routes())
|
app.include_router(setup_claude_routes())
|
||||||
|
|
||||||
from routes.vault_routes import setup_vault_routes
|
from routes.vault.vault_routes import setup_vault_routes
|
||||||
app.include_router(setup_vault_routes())
|
app.include_router(setup_vault_routes())
|
||||||
|
|
||||||
# Contacts (CardDAV)
|
# Contacts (CardDAV)
|
||||||
|
|||||||
@@ -67,6 +67,7 @@ services:
|
|||||||
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
||||||
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
||||||
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
||||||
|
- ODYSSEUS_TTS_CACHE_MAX_BYTES=${ODYSSEUS_TTS_CACHE_MAX_BYTES}
|
||||||
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
||||||
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
||||||
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
||||||
|
|||||||
@@ -66,6 +66,7 @@ services:
|
|||||||
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
||||||
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
||||||
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
||||||
|
- ODYSSEUS_TTS_CACHE_MAX_BYTES=${ODYSSEUS_TTS_CACHE_MAX_BYTES}
|
||||||
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
||||||
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
||||||
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
||||||
|
|||||||
@@ -55,6 +55,7 @@ services:
|
|||||||
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
- ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES=${ODYSSEUS_EMAIL_COMPOSE_UPLOAD_MAX_BYTES:-26214400}
|
||||||
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
- ODYSSEUS_STT_MAX_AUDIO_BYTES=${ODYSSEUS_STT_MAX_AUDIO_BYTES:-26214400}
|
||||||
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
- ODYSSEUS_ICS_MAX_BYTES=${ODYSSEUS_ICS_MAX_BYTES:-10485760}
|
||||||
|
- ODYSSEUS_TTS_CACHE_MAX_BYTES=${ODYSSEUS_TTS_CACHE_MAX_BYTES}
|
||||||
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
- DATA_BRAVE_API_KEY=${DATA_BRAVE_API_KEY:-}
|
||||||
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
- GOOGLE_API_KEY=${GOOGLE_API_KEY:-}
|
||||||
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
- GOOGLE_PSE_CX=${GOOGLE_PSE_CX:-}
|
||||||
|
|||||||
+4
-1
@@ -38,7 +38,10 @@ python-dateutil
|
|||||||
caldav
|
caldav
|
||||||
cryptography
|
cryptography
|
||||||
bcrypt
|
bcrypt
|
||||||
mcp
|
# Built-in servers use the v1 low-level Server decorator API. MCP SDK v2 is a
|
||||||
|
# breaking rewrite, so keep fresh installs on the maintained v1 line until the
|
||||||
|
# servers are migrated together.
|
||||||
|
mcp<2
|
||||||
pyotp
|
pyotp
|
||||||
qrcode[pil]
|
qrcode[pil]
|
||||||
croniter
|
croniter
|
||||||
|
|||||||
@@ -0,0 +1,5 @@
|
|||||||
|
"""Search route domain package (slice 2j, #4082/#4071).
|
||||||
|
|
||||||
|
Contains search_routes.py, migrated from the flat routes/ directory.
|
||||||
|
Backward-compat shim at routes/search_routes.py re-exports from here.
|
||||||
|
"""
|
||||||
@@ -0,0 +1,111 @@
|
|||||||
|
"""Search routes — /api/search/config GET, /api/search POST."""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
from typing import Dict, Any
|
||||||
|
|
||||||
|
from fastapi import APIRouter, Request
|
||||||
|
|
||||||
|
import time
|
||||||
|
|
||||||
|
from services.search import get_search_config, comprehensive_web_search, PROVIDER_INFO
|
||||||
|
from services.search.core import _call_provider
|
||||||
|
from services.search.providers import _get_provider_key, _get_search_instance
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
async def _request_values(request: Request) -> Dict[str, Any]:
|
||||||
|
"""Accept JSON, form data, or query params for search endpoints.
|
||||||
|
|
||||||
|
The browser UI posts FormData, while the agent's generic app_api tool
|
||||||
|
posts JSON. FastAPI Form(...) rejects JSON with a 422 before our handler
|
||||||
|
runs, which made the model think SearXNG was broken.
|
||||||
|
"""
|
||||||
|
values: Dict[str, Any] = dict(request.query_params)
|
||||||
|
content_type = (request.headers.get("content-type") or "").lower()
|
||||||
|
try:
|
||||||
|
if "application/json" in content_type:
|
||||||
|
body = await request.json()
|
||||||
|
if isinstance(body, dict):
|
||||||
|
values.update(body)
|
||||||
|
else:
|
||||||
|
form = await request.form()
|
||||||
|
values.update(dict(form))
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return values
|
||||||
|
|
||||||
|
|
||||||
|
def setup_search_routes(config) -> APIRouter:
|
||||||
|
router = APIRouter(tags=["search"])
|
||||||
|
|
||||||
|
@router.get("/api/search/config")
|
||||||
|
async def get_search_settings() -> Dict[str, Any]:
|
||||||
|
return get_search_config()
|
||||||
|
|
||||||
|
@router.post("/api/search")
|
||||||
|
async def do_web_search(request: Request) -> Dict[str, Any]:
|
||||||
|
"""Standalone web search — returns context string + source list.
|
||||||
|
|
||||||
|
Used by Compare mode to pre-search once and share results across panes.
|
||||||
|
"""
|
||||||
|
values = await _request_values(request)
|
||||||
|
query = str(values.get("query") or values.get("q") or "").strip()
|
||||||
|
if not query:
|
||||||
|
return {"context": "", "sources": [], "error": "query is required"}
|
||||||
|
time_filter = values.get("time_filter") or values.get("freshness")
|
||||||
|
if time_filter is not None:
|
||||||
|
time_filter = str(time_filter).strip() or None
|
||||||
|
try:
|
||||||
|
context, sources = comprehensive_web_search(
|
||||||
|
query, return_sources=True, time_filter=time_filter,
|
||||||
|
)
|
||||||
|
return {"context": context, "sources": sources}
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Standalone web search failed: {e}")
|
||||||
|
return {"context": "", "sources": [], "error": str(e)}
|
||||||
|
|
||||||
|
@router.get("/api/search/providers")
|
||||||
|
async def list_search_providers():
|
||||||
|
"""Return available search providers with config status."""
|
||||||
|
providers = []
|
||||||
|
for pid, (label, needs_key, needs_url) in PROVIDER_INFO.items():
|
||||||
|
if pid == "disabled":
|
||||||
|
continue
|
||||||
|
available = True
|
||||||
|
if needs_key and not _get_provider_key(pid):
|
||||||
|
available = False
|
||||||
|
if needs_url and pid == "searxng" and not _get_search_instance():
|
||||||
|
available = False
|
||||||
|
providers.append({
|
||||||
|
"id": pid,
|
||||||
|
"label": label,
|
||||||
|
"available": available,
|
||||||
|
})
|
||||||
|
return providers
|
||||||
|
|
||||||
|
@router.post("/api/search/query")
|
||||||
|
async def search_with_provider(request: Request) -> Dict[str, Any]:
|
||||||
|
"""Search using a specific provider. Used by compare search mode."""
|
||||||
|
values = await _request_values(request)
|
||||||
|
query = str(values.get("query") or values.get("q") or "").strip()
|
||||||
|
provider = str(values.get("provider") or "").strip()
|
||||||
|
try:
|
||||||
|
count = int(values.get("count") or values.get("limit") or 10)
|
||||||
|
except Exception:
|
||||||
|
count = 10
|
||||||
|
if not query:
|
||||||
|
return {"results": [], "provider": provider, "error": "query is required"}
|
||||||
|
if provider not in PROVIDER_INFO or provider == "disabled":
|
||||||
|
return {"results": [], "provider": provider, "error": "Unknown provider"}
|
||||||
|
t0 = time.time()
|
||||||
|
try:
|
||||||
|
results = _call_provider(provider, query, min(count, 20))
|
||||||
|
elapsed = round(time.time() - t0, 2)
|
||||||
|
return {"results": results, "provider": provider, "time": elapsed}
|
||||||
|
except Exception as e:
|
||||||
|
elapsed = round(time.time() - t0, 2)
|
||||||
|
logger.error(f"Search provider {provider} failed: {e}")
|
||||||
|
return {"results": [], "provider": provider, "time": elapsed, "error": str(e)}
|
||||||
|
|
||||||
|
return router
|
||||||
+9
-107
@@ -1,111 +1,13 @@
|
|||||||
"""Search routes — /api/search/config GET, /api/search POST."""
|
"""Backward-compat shim — canonical location is routes/search/search_routes.py.
|
||||||
|
|
||||||
import logging
|
This module is replaced in ``sys.modules`` by the canonical module object so
|
||||||
from typing import Dict, Any
|
that ``import routes.search_routes`` and ``from routes.search_routes import X``
|
||||||
|
keep resolving to the canonical module. Keeps existing import paths working
|
||||||
|
after slice 2j (#4082/#4071).
|
||||||
|
"""
|
||||||
|
|
||||||
from fastapi import APIRouter, Request
|
import sys as _sys
|
||||||
|
|
||||||
import time
|
from routes.search import search_routes as _canonical # noqa: F401
|
||||||
|
|
||||||
from services.search import get_search_config, comprehensive_web_search, PROVIDER_INFO
|
_sys.modules[__name__] = _canonical
|
||||||
from services.search.core import _call_provider
|
|
||||||
from services.search.providers import _get_provider_key, _get_search_instance
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
async def _request_values(request: Request) -> Dict[str, Any]:
|
|
||||||
"""Accept JSON, form data, or query params for search endpoints.
|
|
||||||
|
|
||||||
The browser UI posts FormData, while the agent's generic app_api tool
|
|
||||||
posts JSON. FastAPI Form(...) rejects JSON with a 422 before our handler
|
|
||||||
runs, which made the model think SearXNG was broken.
|
|
||||||
"""
|
|
||||||
values: Dict[str, Any] = dict(request.query_params)
|
|
||||||
content_type = (request.headers.get("content-type") or "").lower()
|
|
||||||
try:
|
|
||||||
if "application/json" in content_type:
|
|
||||||
body = await request.json()
|
|
||||||
if isinstance(body, dict):
|
|
||||||
values.update(body)
|
|
||||||
else:
|
|
||||||
form = await request.form()
|
|
||||||
values.update(dict(form))
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
return values
|
|
||||||
|
|
||||||
|
|
||||||
def setup_search_routes(config) -> APIRouter:
|
|
||||||
router = APIRouter(tags=["search"])
|
|
||||||
|
|
||||||
@router.get("/api/search/config")
|
|
||||||
async def get_search_settings() -> Dict[str, Any]:
|
|
||||||
return get_search_config()
|
|
||||||
|
|
||||||
@router.post("/api/search")
|
|
||||||
async def do_web_search(request: Request) -> Dict[str, Any]:
|
|
||||||
"""Standalone web search — returns context string + source list.
|
|
||||||
|
|
||||||
Used by Compare mode to pre-search once and share results across panes.
|
|
||||||
"""
|
|
||||||
values = await _request_values(request)
|
|
||||||
query = str(values.get("query") or values.get("q") or "").strip()
|
|
||||||
if not query:
|
|
||||||
return {"context": "", "sources": [], "error": "query is required"}
|
|
||||||
time_filter = values.get("time_filter") or values.get("freshness")
|
|
||||||
if time_filter is not None:
|
|
||||||
time_filter = str(time_filter).strip() or None
|
|
||||||
try:
|
|
||||||
context, sources = comprehensive_web_search(
|
|
||||||
query, return_sources=True, time_filter=time_filter,
|
|
||||||
)
|
|
||||||
return {"context": context, "sources": sources}
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Standalone web search failed: {e}")
|
|
||||||
return {"context": "", "sources": [], "error": str(e)}
|
|
||||||
|
|
||||||
@router.get("/api/search/providers")
|
|
||||||
async def list_search_providers():
|
|
||||||
"""Return available search providers with config status."""
|
|
||||||
providers = []
|
|
||||||
for pid, (label, needs_key, needs_url) in PROVIDER_INFO.items():
|
|
||||||
if pid == "disabled":
|
|
||||||
continue
|
|
||||||
available = True
|
|
||||||
if needs_key and not _get_provider_key(pid):
|
|
||||||
available = False
|
|
||||||
if needs_url and pid == "searxng" and not _get_search_instance():
|
|
||||||
available = False
|
|
||||||
providers.append({
|
|
||||||
"id": pid,
|
|
||||||
"label": label,
|
|
||||||
"available": available,
|
|
||||||
})
|
|
||||||
return providers
|
|
||||||
|
|
||||||
@router.post("/api/search/query")
|
|
||||||
async def search_with_provider(request: Request) -> Dict[str, Any]:
|
|
||||||
"""Search using a specific provider. Used by compare search mode."""
|
|
||||||
values = await _request_values(request)
|
|
||||||
query = str(values.get("query") or values.get("q") or "").strip()
|
|
||||||
provider = str(values.get("provider") or "").strip()
|
|
||||||
try:
|
|
||||||
count = int(values.get("count") or values.get("limit") or 10)
|
|
||||||
except Exception:
|
|
||||||
count = 10
|
|
||||||
if not query:
|
|
||||||
return {"results": [], "provider": provider, "error": "query is required"}
|
|
||||||
if provider not in PROVIDER_INFO or provider == "disabled":
|
|
||||||
return {"results": [], "provider": provider, "error": "Unknown provider"}
|
|
||||||
t0 = time.time()
|
|
||||||
try:
|
|
||||||
results = _call_provider(provider, query, min(count, 20))
|
|
||||||
elapsed = round(time.time() - t0, 2)
|
|
||||||
return {"results": results, "provider": provider, "time": elapsed}
|
|
||||||
except Exception as e:
|
|
||||||
elapsed = round(time.time() - t0, 2)
|
|
||||||
logger.error(f"Search provider {provider} failed: {e}")
|
|
||||||
return {"results": [], "provider": provider, "time": elapsed, "error": str(e)}
|
|
||||||
|
|
||||||
return router
|
|
||||||
|
|||||||
@@ -1409,7 +1409,7 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter:
|
|||||||
|
|
||||||
# Prefer the configured DEFAULT (→ Utility) model — not the current chat
|
# Prefer the configured DEFAULT (→ Utility) model — not the current chat
|
||||||
# session's model. Fall back to the caller's session model only if unset.
|
# session's model. Fall back to the caller's session model only if unset.
|
||||||
url, model, headers = resolve_endpoint("default", owner=user)
|
url, model, headers = resolve_endpoint("utility", owner=user)
|
||||||
if not url or not model:
|
if not url or not model:
|
||||||
url = url or ((body.get("endpoint_url") or "").strip() or None)
|
url = url or ((body.get("endpoint_url") or "").strip() or None)
|
||||||
model = model or ((body.get("model") or "").strip() or None)
|
model = model or ((body.get("model") or "").strip() or None)
|
||||||
|
|||||||
@@ -0,0 +1,5 @@
|
|||||||
|
"""Vault route domain package (slice 2k, #4082/#4071).
|
||||||
|
|
||||||
|
Contains vault_routes.py, migrated from the flat routes/ directory.
|
||||||
|
Backward-compat shim at routes/vault_routes.py re-exports from here.
|
||||||
|
"""
|
||||||
@@ -0,0 +1,242 @@
|
|||||||
|
"""
|
||||||
|
vault_routes.py
|
||||||
|
|
||||||
|
Vaultwarden / Bitwarden CLI integration — config and unlock endpoints.
|
||||||
|
Stores the BW_SESSION key in data/vault.json with restrictive permissions.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
import asyncio
|
||||||
|
from pathlib import Path
|
||||||
|
from datetime import datetime
|
||||||
|
from fastapi import APIRouter, Request
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
from core.middleware import require_admin
|
||||||
|
from core.platform_compat import IS_WINDOWS, safe_chmod, which_tool
|
||||||
|
from src.constants import VAULT_FILE as _VAULT_FILE
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
VAULT_FILE = Path(_VAULT_FILE)
|
||||||
|
|
||||||
|
|
||||||
|
def _find_bw() -> str:
|
||||||
|
"""Locate the bw binary, checking PATH and common npm-global locations.
|
||||||
|
|
||||||
|
On Windows the Bitwarden CLI shim is `bw.cmd`/`bw.exe`, resolved by
|
||||||
|
which_tool via PATHEXT.
|
||||||
|
"""
|
||||||
|
p = which_tool("bw")
|
||||||
|
if p:
|
||||||
|
return p
|
||||||
|
if IS_WINDOWS:
|
||||||
|
appdata = os.environ.get("APPDATA", os.path.expanduser("~"))
|
||||||
|
for candidate in (
|
||||||
|
os.path.join(appdata, "npm", "bw.cmd"),
|
||||||
|
os.path.join(appdata, "npm", "bw.exe"),
|
||||||
|
):
|
||||||
|
if os.path.isfile(candidate):
|
||||||
|
return candidate
|
||||||
|
return "bw"
|
||||||
|
home = os.path.expanduser("~")
|
||||||
|
for candidate in (
|
||||||
|
f"{home}/.npm-global/bin/bw",
|
||||||
|
f"{home}/.nvm/versions/node/*/bin/bw",
|
||||||
|
"/usr/local/bin/bw",
|
||||||
|
"/opt/homebrew/bin/bw",
|
||||||
|
):
|
||||||
|
if "*" in candidate:
|
||||||
|
import glob
|
||||||
|
for m in glob.glob(candidate):
|
||||||
|
if os.path.isfile(m) and os.access(m, os.X_OK):
|
||||||
|
return m
|
||||||
|
elif os.path.isfile(candidate) and os.access(candidate, os.X_OK):
|
||||||
|
return candidate
|
||||||
|
return "bw" # fall back to PATH lookup (will FileNotFoundError, handled below)
|
||||||
|
|
||||||
|
|
||||||
|
def _load_config() -> dict:
|
||||||
|
if VAULT_FILE.exists():
|
||||||
|
try:
|
||||||
|
data = json.loads(VAULT_FILE.read_text(encoding="utf-8"))
|
||||||
|
return data if isinstance(data, dict) else {}
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
return {}
|
||||||
|
|
||||||
|
|
||||||
|
def _save_config(cfg: dict):
|
||||||
|
VAULT_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
VAULT_FILE.write_text(json.dumps(cfg, indent=2), encoding="utf-8")
|
||||||
|
# POSIX: restrict the BW_SESSION store to 0o600. Windows: no-op (profile dir
|
||||||
|
# is ACL-restricted already).
|
||||||
|
safe_chmod(str(VAULT_FILE), 0o600)
|
||||||
|
|
||||||
|
|
||||||
|
async def _run_bw(args: list, session: str = None, input_text: str = None,
|
||||||
|
bw_password: str = None) -> tuple:
|
||||||
|
env = {}
|
||||||
|
env.update(os.environ)
|
||||||
|
if session:
|
||||||
|
env["BW_SESSION"] = session
|
||||||
|
# Secrets must never be passed as argv — process arguments are world-readable
|
||||||
|
# via `ps` / `/proc/<pid>/cmdline` to any local user. Keep --passwordenv
|
||||||
|
# support for bw commands that need it; unlock/login callers should prefer
|
||||||
|
# stdin so the master password is not left in the child environment either.
|
||||||
|
if bw_password is not None:
|
||||||
|
env["BW_PASSWORD"] = bw_password
|
||||||
|
bw_path = _find_bw()
|
||||||
|
try:
|
||||||
|
proc = await asyncio.create_subprocess_exec(
|
||||||
|
bw_path, *args,
|
||||||
|
stdin=asyncio.subprocess.PIPE if input_text else None,
|
||||||
|
stdout=asyncio.subprocess.PIPE,
|
||||||
|
stderr=asyncio.subprocess.PIPE,
|
||||||
|
env=env,
|
||||||
|
)
|
||||||
|
except FileNotFoundError:
|
||||||
|
return "", "bw CLI not installed (install `nodejs-bitwarden-cli` or `bitwarden-cli`)", 127
|
||||||
|
except Exception as e:
|
||||||
|
return "", f"Failed to launch bw: {e}", 1
|
||||||
|
try:
|
||||||
|
stdout, stderr = await proc.communicate(input=input_text.encode() if input_text else None)
|
||||||
|
except Exception as e:
|
||||||
|
return "", f"bw subprocess error: {e}", 1
|
||||||
|
return stdout.decode(errors="replace").strip(), stderr.decode(errors="replace").strip(), proc.returncode
|
||||||
|
|
||||||
|
|
||||||
|
class VaultConfig(BaseModel):
|
||||||
|
server_url: str = ""
|
||||||
|
email: str = ""
|
||||||
|
|
||||||
|
|
||||||
|
class VaultUnlockRequest(BaseModel):
|
||||||
|
master_password: str
|
||||||
|
|
||||||
|
|
||||||
|
class VaultLoginRequest(BaseModel):
|
||||||
|
email: str
|
||||||
|
master_password: str
|
||||||
|
|
||||||
|
|
||||||
|
def setup_vault_routes():
|
||||||
|
router = APIRouter(prefix="/api/vault", tags=["vault"])
|
||||||
|
|
||||||
|
@router.get("/config")
|
||||||
|
async def get_config(request: Request):
|
||||||
|
"""Return vault config (no sensitive fields)."""
|
||||||
|
require_admin(request)
|
||||||
|
cfg = _load_config()
|
||||||
|
return {
|
||||||
|
"server_url": cfg.get("server_url", ""),
|
||||||
|
"email": cfg.get("email", ""),
|
||||||
|
"unlocked": bool(cfg.get("session")),
|
||||||
|
"unlocked_at": cfg.get("unlocked_at", ""),
|
||||||
|
"bw_installed": await _check_bw_installed(),
|
||||||
|
}
|
||||||
|
|
||||||
|
@router.post("/config")
|
||||||
|
async def save_config(req: VaultConfig, request: Request):
|
||||||
|
"""Save vault URL + email. Runs 'bw config server' to point at Vaultwarden."""
|
||||||
|
require_admin(request)
|
||||||
|
cfg = _load_config()
|
||||||
|
cfg["server_url"] = req.server_url.strip().rstrip("/")
|
||||||
|
cfg["email"] = req.email.strip()
|
||||||
|
|
||||||
|
if cfg["server_url"]:
|
||||||
|
_, stderr, rc = await _run_bw(["config", "server", cfg["server_url"]])
|
||||||
|
if rc != 0:
|
||||||
|
return {"ok": False, "error": f"bw config failed: {stderr[:300]}"}
|
||||||
|
|
||||||
|
_save_config(cfg)
|
||||||
|
return {"ok": True}
|
||||||
|
|
||||||
|
@router.post("/login")
|
||||||
|
async def login(req: VaultLoginRequest, request: Request):
|
||||||
|
"""Log in to Vaultwarden (required once per account)."""
|
||||||
|
require_admin(request)
|
||||||
|
cfg = _load_config()
|
||||||
|
# Update email
|
||||||
|
cfg["email"] = req.email
|
||||||
|
_save_config(cfg)
|
||||||
|
|
||||||
|
stdout, stderr, rc = await _run_bw(
|
||||||
|
["login", req.email, "--raw"],
|
||||||
|
input_text=req.master_password + "\n",
|
||||||
|
)
|
||||||
|
if rc != 0:
|
||||||
|
# Already logged in is OK
|
||||||
|
if "already logged in" in stderr.lower():
|
||||||
|
return {"ok": True, "already": True}
|
||||||
|
return {"ok": False, "error": f"Login failed: {stderr[:300]}"}
|
||||||
|
# bw login --raw prints session key on success (when 2FA disabled)
|
||||||
|
if stdout:
|
||||||
|
cfg["session"] = stdout
|
||||||
|
cfg["unlocked_at"] = datetime.utcnow().isoformat()
|
||||||
|
_save_config(cfg)
|
||||||
|
return {"ok": True}
|
||||||
|
|
||||||
|
@router.post("/unlock")
|
||||||
|
async def unlock(req: VaultUnlockRequest, request: Request):
|
||||||
|
"""Unlock the vault and save the session key."""
|
||||||
|
require_admin(request)
|
||||||
|
# Pass the master password on stdin, not argv. argv is visible through
|
||||||
|
# `ps` / /proc/<pid>/cmdline; stdin also avoids leaving the secret in
|
||||||
|
# the child process environment.
|
||||||
|
stdout, stderr, rc = await _run_bw(
|
||||||
|
["unlock", "--raw"],
|
||||||
|
input_text=req.master_password + "\n",
|
||||||
|
)
|
||||||
|
if rc != 0:
|
||||||
|
return {"ok": False, "error": f"Unlock failed: {stderr[:300]}"}
|
||||||
|
session = stdout.strip()
|
||||||
|
if not session:
|
||||||
|
return {"ok": False, "error": "bw returned empty session"}
|
||||||
|
cfg = _load_config()
|
||||||
|
cfg["session"] = session
|
||||||
|
cfg["unlocked_at"] = datetime.utcnow().isoformat()
|
||||||
|
_save_config(cfg)
|
||||||
|
return {"ok": True, "message": "Vault unlocked"}
|
||||||
|
|
||||||
|
@router.post("/lock")
|
||||||
|
async def lock(request: Request):
|
||||||
|
"""Lock the vault (clear session from config)."""
|
||||||
|
require_admin(request)
|
||||||
|
cfg = _load_config()
|
||||||
|
cfg.pop("session", None)
|
||||||
|
cfg.pop("unlocked_at", None)
|
||||||
|
_save_config(cfg)
|
||||||
|
# Also tell bw to lock
|
||||||
|
await _run_bw(["lock"])
|
||||||
|
return {"ok": True, "message": "Vault locked"}
|
||||||
|
|
||||||
|
@router.post("/logout")
|
||||||
|
async def logout(request: Request):
|
||||||
|
"""Log out of the Bitwarden CLI completely."""
|
||||||
|
require_admin(request)
|
||||||
|
await _run_bw(["logout"])
|
||||||
|
cfg = _load_config()
|
||||||
|
cfg.pop("session", None)
|
||||||
|
cfg.pop("email", None)
|
||||||
|
cfg.pop("unlocked_at", None)
|
||||||
|
_save_config(cfg)
|
||||||
|
return {"ok": True}
|
||||||
|
|
||||||
|
return router
|
||||||
|
|
||||||
|
|
||||||
|
async def _check_bw_installed() -> bool:
|
||||||
|
try:
|
||||||
|
proc = await asyncio.create_subprocess_exec(
|
||||||
|
_find_bw(), "--version",
|
||||||
|
stdout=asyncio.subprocess.PIPE,
|
||||||
|
stderr=asyncio.subprocess.PIPE,
|
||||||
|
)
|
||||||
|
await proc.communicate()
|
||||||
|
return proc.returncode == 0
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
+9
-237
@@ -1,242 +1,14 @@
|
|||||||
"""
|
"""Backward-compat shim — canonical location is routes/vault/vault_routes.py.
|
||||||
vault_routes.py
|
|
||||||
|
|
||||||
Vaultwarden / Bitwarden CLI integration — config and unlock endpoints.
|
This module is replaced in ``sys.modules`` by the canonical module object so
|
||||||
Stores the BW_SESSION key in data/vault.json with restrictive permissions.
|
that ``import routes.vault_routes``, ``from routes.vault_routes import X``,
|
||||||
|
and the ``import ... as vr`` + ``monkeypatch.setattr(vr, ...)`` pattern used
|
||||||
|
by test_vault_password_not_in_argv.py all operate on the *same* object.
|
||||||
|
Keeps existing import paths working after slice 2k (#4082/#4071).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import json
|
import sys as _sys
|
||||||
import logging
|
|
||||||
import os
|
|
||||||
import shutil
|
|
||||||
import asyncio
|
|
||||||
from pathlib import Path
|
|
||||||
from datetime import datetime
|
|
||||||
from fastapi import APIRouter, Request
|
|
||||||
from pydantic import BaseModel
|
|
||||||
|
|
||||||
from core.middleware import require_admin
|
from routes.vault import vault_routes as _canonical # noqa: F401
|
||||||
from core.platform_compat import IS_WINDOWS, safe_chmod, which_tool
|
|
||||||
from src.constants import VAULT_FILE as _VAULT_FILE
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
_sys.modules[__name__] = _canonical
|
||||||
|
|
||||||
VAULT_FILE = Path(_VAULT_FILE)
|
|
||||||
|
|
||||||
|
|
||||||
def _find_bw() -> str:
|
|
||||||
"""Locate the bw binary, checking PATH and common npm-global locations.
|
|
||||||
|
|
||||||
On Windows the Bitwarden CLI shim is `bw.cmd`/`bw.exe`, resolved by
|
|
||||||
which_tool via PATHEXT.
|
|
||||||
"""
|
|
||||||
p = which_tool("bw")
|
|
||||||
if p:
|
|
||||||
return p
|
|
||||||
if IS_WINDOWS:
|
|
||||||
appdata = os.environ.get("APPDATA", os.path.expanduser("~"))
|
|
||||||
for candidate in (
|
|
||||||
os.path.join(appdata, "npm", "bw.cmd"),
|
|
||||||
os.path.join(appdata, "npm", "bw.exe"),
|
|
||||||
):
|
|
||||||
if os.path.isfile(candidate):
|
|
||||||
return candidate
|
|
||||||
return "bw"
|
|
||||||
home = os.path.expanduser("~")
|
|
||||||
for candidate in (
|
|
||||||
f"{home}/.npm-global/bin/bw",
|
|
||||||
f"{home}/.nvm/versions/node/*/bin/bw",
|
|
||||||
"/usr/local/bin/bw",
|
|
||||||
"/opt/homebrew/bin/bw",
|
|
||||||
):
|
|
||||||
if "*" in candidate:
|
|
||||||
import glob
|
|
||||||
for m in glob.glob(candidate):
|
|
||||||
if os.path.isfile(m) and os.access(m, os.X_OK):
|
|
||||||
return m
|
|
||||||
elif os.path.isfile(candidate) and os.access(candidate, os.X_OK):
|
|
||||||
return candidate
|
|
||||||
return "bw" # fall back to PATH lookup (will FileNotFoundError, handled below)
|
|
||||||
|
|
||||||
|
|
||||||
def _load_config() -> dict:
|
|
||||||
if VAULT_FILE.exists():
|
|
||||||
try:
|
|
||||||
data = json.loads(VAULT_FILE.read_text(encoding="utf-8"))
|
|
||||||
return data if isinstance(data, dict) else {}
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
return {}
|
|
||||||
|
|
||||||
|
|
||||||
def _save_config(cfg: dict):
|
|
||||||
VAULT_FILE.parent.mkdir(parents=True, exist_ok=True)
|
|
||||||
VAULT_FILE.write_text(json.dumps(cfg, indent=2), encoding="utf-8")
|
|
||||||
# POSIX: restrict the BW_SESSION store to 0o600. Windows: no-op (profile dir
|
|
||||||
# is ACL-restricted already).
|
|
||||||
safe_chmod(str(VAULT_FILE), 0o600)
|
|
||||||
|
|
||||||
|
|
||||||
async def _run_bw(args: list, session: str = None, input_text: str = None,
|
|
||||||
bw_password: str = None) -> tuple:
|
|
||||||
env = {}
|
|
||||||
env.update(os.environ)
|
|
||||||
if session:
|
|
||||||
env["BW_SESSION"] = session
|
|
||||||
# Secrets must never be passed as argv — process arguments are world-readable
|
|
||||||
# via `ps` / `/proc/<pid>/cmdline` to any local user. Keep --passwordenv
|
|
||||||
# support for bw commands that need it; unlock/login callers should prefer
|
|
||||||
# stdin so the master password is not left in the child environment either.
|
|
||||||
if bw_password is not None:
|
|
||||||
env["BW_PASSWORD"] = bw_password
|
|
||||||
bw_path = _find_bw()
|
|
||||||
try:
|
|
||||||
proc = await asyncio.create_subprocess_exec(
|
|
||||||
bw_path, *args,
|
|
||||||
stdin=asyncio.subprocess.PIPE if input_text else None,
|
|
||||||
stdout=asyncio.subprocess.PIPE,
|
|
||||||
stderr=asyncio.subprocess.PIPE,
|
|
||||||
env=env,
|
|
||||||
)
|
|
||||||
except FileNotFoundError:
|
|
||||||
return "", "bw CLI not installed (install `nodejs-bitwarden-cli` or `bitwarden-cli`)", 127
|
|
||||||
except Exception as e:
|
|
||||||
return "", f"Failed to launch bw: {e}", 1
|
|
||||||
try:
|
|
||||||
stdout, stderr = await proc.communicate(input=input_text.encode() if input_text else None)
|
|
||||||
except Exception as e:
|
|
||||||
return "", f"bw subprocess error: {e}", 1
|
|
||||||
return stdout.decode(errors="replace").strip(), stderr.decode(errors="replace").strip(), proc.returncode
|
|
||||||
|
|
||||||
|
|
||||||
class VaultConfig(BaseModel):
|
|
||||||
server_url: str = ""
|
|
||||||
email: str = ""
|
|
||||||
|
|
||||||
|
|
||||||
class VaultUnlockRequest(BaseModel):
|
|
||||||
master_password: str
|
|
||||||
|
|
||||||
|
|
||||||
class VaultLoginRequest(BaseModel):
|
|
||||||
email: str
|
|
||||||
master_password: str
|
|
||||||
|
|
||||||
|
|
||||||
def setup_vault_routes():
|
|
||||||
router = APIRouter(prefix="/api/vault", tags=["vault"])
|
|
||||||
|
|
||||||
@router.get("/config")
|
|
||||||
async def get_config(request: Request):
|
|
||||||
"""Return vault config (no sensitive fields)."""
|
|
||||||
require_admin(request)
|
|
||||||
cfg = _load_config()
|
|
||||||
return {
|
|
||||||
"server_url": cfg.get("server_url", ""),
|
|
||||||
"email": cfg.get("email", ""),
|
|
||||||
"unlocked": bool(cfg.get("session")),
|
|
||||||
"unlocked_at": cfg.get("unlocked_at", ""),
|
|
||||||
"bw_installed": await _check_bw_installed(),
|
|
||||||
}
|
|
||||||
|
|
||||||
@router.post("/config")
|
|
||||||
async def save_config(req: VaultConfig, request: Request):
|
|
||||||
"""Save vault URL + email. Runs 'bw config server' to point at Vaultwarden."""
|
|
||||||
require_admin(request)
|
|
||||||
cfg = _load_config()
|
|
||||||
cfg["server_url"] = req.server_url.strip().rstrip("/")
|
|
||||||
cfg["email"] = req.email.strip()
|
|
||||||
|
|
||||||
if cfg["server_url"]:
|
|
||||||
_, stderr, rc = await _run_bw(["config", "server", cfg["server_url"]])
|
|
||||||
if rc != 0:
|
|
||||||
return {"ok": False, "error": f"bw config failed: {stderr[:300]}"}
|
|
||||||
|
|
||||||
_save_config(cfg)
|
|
||||||
return {"ok": True}
|
|
||||||
|
|
||||||
@router.post("/login")
|
|
||||||
async def login(req: VaultLoginRequest, request: Request):
|
|
||||||
"""Log in to Vaultwarden (required once per account)."""
|
|
||||||
require_admin(request)
|
|
||||||
cfg = _load_config()
|
|
||||||
# Update email
|
|
||||||
cfg["email"] = req.email
|
|
||||||
_save_config(cfg)
|
|
||||||
|
|
||||||
stdout, stderr, rc = await _run_bw(
|
|
||||||
["login", req.email, "--raw"],
|
|
||||||
input_text=req.master_password + "\n",
|
|
||||||
)
|
|
||||||
if rc != 0:
|
|
||||||
# Already logged in is OK
|
|
||||||
if "already logged in" in stderr.lower():
|
|
||||||
return {"ok": True, "already": True}
|
|
||||||
return {"ok": False, "error": f"Login failed: {stderr[:300]}"}
|
|
||||||
# bw login --raw prints session key on success (when 2FA disabled)
|
|
||||||
if stdout:
|
|
||||||
cfg["session"] = stdout
|
|
||||||
cfg["unlocked_at"] = datetime.utcnow().isoformat()
|
|
||||||
_save_config(cfg)
|
|
||||||
return {"ok": True}
|
|
||||||
|
|
||||||
@router.post("/unlock")
|
|
||||||
async def unlock(req: VaultUnlockRequest, request: Request):
|
|
||||||
"""Unlock the vault and save the session key."""
|
|
||||||
require_admin(request)
|
|
||||||
# Pass the master password on stdin, not argv. argv is visible through
|
|
||||||
# `ps` / /proc/<pid>/cmdline; stdin also avoids leaving the secret in
|
|
||||||
# the child process environment.
|
|
||||||
stdout, stderr, rc = await _run_bw(
|
|
||||||
["unlock", "--raw"],
|
|
||||||
input_text=req.master_password + "\n",
|
|
||||||
)
|
|
||||||
if rc != 0:
|
|
||||||
return {"ok": False, "error": f"Unlock failed: {stderr[:300]}"}
|
|
||||||
session = stdout.strip()
|
|
||||||
if not session:
|
|
||||||
return {"ok": False, "error": "bw returned empty session"}
|
|
||||||
cfg = _load_config()
|
|
||||||
cfg["session"] = session
|
|
||||||
cfg["unlocked_at"] = datetime.utcnow().isoformat()
|
|
||||||
_save_config(cfg)
|
|
||||||
return {"ok": True, "message": "Vault unlocked"}
|
|
||||||
|
|
||||||
@router.post("/lock")
|
|
||||||
async def lock(request: Request):
|
|
||||||
"""Lock the vault (clear session from config)."""
|
|
||||||
require_admin(request)
|
|
||||||
cfg = _load_config()
|
|
||||||
cfg.pop("session", None)
|
|
||||||
cfg.pop("unlocked_at", None)
|
|
||||||
_save_config(cfg)
|
|
||||||
# Also tell bw to lock
|
|
||||||
await _run_bw(["lock"])
|
|
||||||
return {"ok": True, "message": "Vault locked"}
|
|
||||||
|
|
||||||
@router.post("/logout")
|
|
||||||
async def logout(request: Request):
|
|
||||||
"""Log out of the Bitwarden CLI completely."""
|
|
||||||
require_admin(request)
|
|
||||||
await _run_bw(["logout"])
|
|
||||||
cfg = _load_config()
|
|
||||||
cfg.pop("session", None)
|
|
||||||
cfg.pop("email", None)
|
|
||||||
cfg.pop("unlocked_at", None)
|
|
||||||
_save_config(cfg)
|
|
||||||
return {"ok": True}
|
|
||||||
|
|
||||||
return router
|
|
||||||
|
|
||||||
|
|
||||||
async def _check_bw_installed() -> bool:
|
|
||||||
try:
|
|
||||||
proc = await asyncio.create_subprocess_exec(
|
|
||||||
_find_bw(), "--version",
|
|
||||||
stdout=asyncio.subprocess.PIPE,
|
|
||||||
stderr=asyncio.subprocess.PIPE,
|
|
||||||
)
|
|
||||||
await proc.communicate()
|
|
||||||
return proc.returncode == 0
|
|
||||||
except Exception:
|
|
||||||
return False
|
|
||||||
|
|||||||
@@ -0,0 +1,5 @@
|
|||||||
|
"""Webhook route domain package (slice 2l, #4082/#4071).
|
||||||
|
|
||||||
|
Contains webhook_routes.py, migrated from the flat routes/ directory.
|
||||||
|
Backward-compat shim at routes/webhook_routes.py re-exports from here.
|
||||||
|
"""
|
||||||
@@ -0,0 +1,395 @@
|
|||||||
|
"""Webhook, API Token, and sync chat routes."""
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
import logging
|
||||||
|
from typing import Optional
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
from fastapi import APIRouter, HTTPException, Request, Form
|
||||||
|
from pydantic import BaseModel, Field
|
||||||
|
|
||||||
|
from core.database import SessionLocal, Webhook, ModelEndpoint
|
||||||
|
from src.auth_helpers import owner_filter
|
||||||
|
from src.url_security import validate_public_http_url
|
||||||
|
from src.webhook_manager import WebhookManager, validate_webhook_url, validate_events
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
router = APIRouter(prefix="/api", tags=["webhooks"])
|
||||||
|
|
||||||
|
# Input limits
|
||||||
|
MAX_NAME_LEN = 100
|
||||||
|
MAX_URL_LEN = 2048
|
||||||
|
MAX_SECRET_LEN = 256
|
||||||
|
MAX_MESSAGE_LEN = 32_000
|
||||||
|
|
||||||
|
|
||||||
|
from core.middleware import require_admin as _require_admin
|
||||||
|
|
||||||
|
|
||||||
|
def _select_api_chat_fallback_endpoint(db, token_owner: Optional[str]):
|
||||||
|
"""First enabled ModelEndpoint visible to token_owner — their own rows plus
|
||||||
|
legacy null-owner ("shared") rows. Owner-scoped: an unscoped .first() would
|
||||||
|
let a chat-scoped token fall back onto another user's private endpoint and
|
||||||
|
silently spend that owner's API key/quota. Prefer owner rows before shared
|
||||||
|
rows. Fails closed to null-owner rows only when token_owner is absent.
|
||||||
|
Does not validate base_url — admin-configured local/LAN endpoints remain allowed.
|
||||||
|
"""
|
||||||
|
query = db.query(ModelEndpoint).filter(ModelEndpoint.is_enabled == True) # noqa: E712
|
||||||
|
if token_owner:
|
||||||
|
query = owner_filter(query, ModelEndpoint, token_owner)
|
||||||
|
return query.order_by(ModelEndpoint.owner.desc(), ModelEndpoint.created_at).first()
|
||||||
|
return query.filter(ModelEndpoint.owner == None).order_by(ModelEndpoint.created_at).first() # noqa: E711
|
||||||
|
|
||||||
|
|
||||||
|
def _caller_owns_session(sess_owner, caller) -> bool:
|
||||||
|
"""Strict session-ownership gate for the token-authenticated sync-chat
|
||||||
|
endpoint (`POST /api/v1/chat`).
|
||||||
|
|
||||||
|
Mirrors ``_verify_session_owner`` in session_routes.py and the null-owner
|
||||||
|
gates in notes/calendar/gallery: a caller may resume a session ONLY when
|
||||||
|
its owner matches them exactly. A null/empty session owner (legacy or
|
||||||
|
migrated rows) is deliberately NOT resumable by an arbitrary token — the
|
||||||
|
old ``sess_owner and sess_owner != caller`` form skipped the check whenever
|
||||||
|
``sess_owner`` was falsy, so any chat-scoped token (e.g. a paired mobile
|
||||||
|
device) could resume such a session, inject a message, and read back its
|
||||||
|
history and reuse the owner's endpoint credentials. Fail closed: an
|
||||||
|
unresolvable caller also returns False.
|
||||||
|
"""
|
||||||
|
if not caller:
|
||||||
|
return False
|
||||||
|
return sess_owner == caller
|
||||||
|
|
||||||
|
|
||||||
|
def setup_webhook_routes(
|
||||||
|
webhook_manager: WebhookManager,
|
||||||
|
auth_manager,
|
||||||
|
session_manager=None,
|
||||||
|
api_key_manager=None,
|
||||||
|
) -> APIRouter:
|
||||||
|
|
||||||
|
@router.get("/webhooks")
|
||||||
|
def list_webhooks(request: Request):
|
||||||
|
_require_admin(request)
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
hooks = db.query(Webhook).all()
|
||||||
|
return [
|
||||||
|
{
|
||||||
|
"id": w.id,
|
||||||
|
"name": w.name,
|
||||||
|
"url": w.url,
|
||||||
|
"has_secret": bool(w.secret),
|
||||||
|
"events": w.events.split(",") if w.events else [],
|
||||||
|
"is_active": w.is_active,
|
||||||
|
"last_triggered_at": w.last_triggered_at.isoformat() if w.last_triggered_at else None,
|
||||||
|
"last_status_code": w.last_status_code,
|
||||||
|
"last_error": w.last_error,
|
||||||
|
"created_at": w.created_at.isoformat() if w.created_at else None,
|
||||||
|
}
|
||||||
|
for w in hooks
|
||||||
|
]
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
@router.post("/webhooks")
|
||||||
|
def create_webhook(
|
||||||
|
request: Request,
|
||||||
|
name: str = Form(""),
|
||||||
|
url: str = Form(""),
|
||||||
|
secret: str = Form(""),
|
||||||
|
events: str = Form(""),
|
||||||
|
):
|
||||||
|
_require_admin(request)
|
||||||
|
name = name.strip()[:MAX_NAME_LEN]
|
||||||
|
if not name:
|
||||||
|
raise HTTPException(400, "Webhook name is required")
|
||||||
|
try:
|
||||||
|
url = validate_webhook_url(url)
|
||||||
|
except ValueError as e:
|
||||||
|
raise HTTPException(400, str(e))
|
||||||
|
try:
|
||||||
|
events = validate_events(events)
|
||||||
|
except ValueError as e:
|
||||||
|
raise HTTPException(400, str(e))
|
||||||
|
|
||||||
|
secret_val = secret.strip()[:MAX_SECRET_LEN] or None
|
||||||
|
# Encrypt the secret at rest using the same Fernet key as API keys
|
||||||
|
encrypted_secret = None
|
||||||
|
if secret_val and api_key_manager:
|
||||||
|
encrypted_secret = api_key_manager.encrypt_api_key(secret_val)
|
||||||
|
elif secret_val:
|
||||||
|
encrypted_secret = secret_val # Fallback if no encryption available
|
||||||
|
|
||||||
|
webhook_id = str(uuid.uuid4())[:8]
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
db.add(Webhook(
|
||||||
|
id=webhook_id,
|
||||||
|
name=name,
|
||||||
|
url=url,
|
||||||
|
secret=encrypted_secret,
|
||||||
|
events=events,
|
||||||
|
is_active=True,
|
||||||
|
))
|
||||||
|
db.commit()
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
return {"id": webhook_id, "name": name}
|
||||||
|
|
||||||
|
@router.post("/webhooks/{webhook_id}/test")
|
||||||
|
async def test_webhook(request: Request, webhook_id: str):
|
||||||
|
_require_admin(request)
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
wh = db.query(Webhook).filter(Webhook.id == webhook_id).first()
|
||||||
|
if not wh:
|
||||||
|
raise HTTPException(404, "Webhook not found")
|
||||||
|
url, secret = wh.url, wh.secret
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
await webhook_manager.deliver_test(webhook_id, url, secret)
|
||||||
|
return {"status": "sent"}
|
||||||
|
|
||||||
|
@router.patch("/webhooks/{webhook_id}")
|
||||||
|
def toggle_webhook(request: Request, webhook_id: str):
|
||||||
|
_require_admin(request)
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
wh = db.query(Webhook).filter(Webhook.id == webhook_id).first()
|
||||||
|
if not wh:
|
||||||
|
raise HTTPException(404, "Webhook not found")
|
||||||
|
wh.is_active = not wh.is_active
|
||||||
|
db.commit()
|
||||||
|
return {"id": webhook_id, "is_active": wh.is_active}
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
@router.delete("/webhooks/{webhook_id}")
|
||||||
|
def delete_webhook(request: Request, webhook_id: str):
|
||||||
|
_require_admin(request)
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
deleted = db.query(Webhook).filter(Webhook.id == webhook_id).delete()
|
||||||
|
db.commit()
|
||||||
|
if not deleted:
|
||||||
|
raise HTTPException(404, "Webhook not found")
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
return {"status": "deleted"}
|
||||||
|
|
||||||
|
# ================================================================
|
||||||
|
# Sync Chat Endpoint (for n8n / Make / Activepieces)
|
||||||
|
# ================================================================
|
||||||
|
|
||||||
|
# Known provider base URLs — auto-resolved from api_key prefix or model name
|
||||||
|
KNOWN_PROVIDERS = {
|
||||||
|
"deepseek": "https://api.deepseek.com/v1",
|
||||||
|
"openai": "https://api.openai.com/v1",
|
||||||
|
"mistral": "https://api.mistral.ai/v1",
|
||||||
|
"groq": "https://api.groq.com/openai/v1",
|
||||||
|
"together": "https://api.together.xyz/v1",
|
||||||
|
"openrouter": "https://openrouter.ai/api/v1",
|
||||||
|
"ollama": "https://ollama.com/api",
|
||||||
|
"opencode-zen": "https://opencode.ai/zen/v1",
|
||||||
|
"opencode-go": "https://opencode.ai/zen/go/v1",
|
||||||
|
"fireworks": "https://api.fireworks.ai/inference/v1",
|
||||||
|
"venice": "https://api.venice.ai/api/v1",
|
||||||
|
"kimi-code": "https://api.kimi.com/coding/v1",
|
||||||
|
"kimicode": "https://api.kimi.com/coding/v1",
|
||||||
|
}
|
||||||
|
|
||||||
|
# Model prefix → provider mapping for auto-detection
|
||||||
|
MODEL_PROVIDER_MAP = {
|
||||||
|
"deepseek": "deepseek",
|
||||||
|
"gpt-": "openai",
|
||||||
|
"o1": "openai",
|
||||||
|
"o3": "openai",
|
||||||
|
"o4": "openai",
|
||||||
|
"mistral": "mistral",
|
||||||
|
"llama": "groq",
|
||||||
|
"mixtral": "groq",
|
||||||
|
"kimi-for-coding": "kimi-code",
|
||||||
|
"kimi": "kimi-code",
|
||||||
|
}
|
||||||
|
|
||||||
|
def _resolve_base_url(model: Optional[str], provider: Optional[str]) -> Optional[str]:
|
||||||
|
"""Try to auto-resolve a base URL from provider name or model prefix."""
|
||||||
|
if provider and provider.lower() in KNOWN_PROVIDERS:
|
||||||
|
return KNOWN_PROVIDERS[provider.lower()]
|
||||||
|
if model:
|
||||||
|
model_lower = model.lower()
|
||||||
|
for prefix, prov in MODEL_PROVIDER_MAP.items():
|
||||||
|
if model_lower.startswith(prefix):
|
||||||
|
return KNOWN_PROVIDERS[prov]
|
||||||
|
return None
|
||||||
|
|
||||||
|
class SyncChatRequest(BaseModel):
|
||||||
|
message: str = Field(..., max_length=MAX_MESSAGE_LEN)
|
||||||
|
model: Optional[str] = Field(None, max_length=200)
|
||||||
|
session: Optional[str] = Field(None, max_length=100)
|
||||||
|
api_key: Optional[str] = Field(None, max_length=256)
|
||||||
|
base_url: Optional[str] = Field(None, max_length=MAX_URL_LEN)
|
||||||
|
provider: Optional[str] = Field(None, max_length=50)
|
||||||
|
|
||||||
|
@router.post("/v1/chat")
|
||||||
|
async def sync_chat(request: Request, body: SyncChatRequest):
|
||||||
|
if not getattr(request.state, "api_token", False):
|
||||||
|
raise HTTPException(403, "This endpoint requires an API token")
|
||||||
|
scopes = set(getattr(request.state, "api_token_scopes", []) or [])
|
||||||
|
if "chat" not in scopes:
|
||||||
|
raise HTTPException(403, "API token is not scoped for chat")
|
||||||
|
token_owner = getattr(request.state, "api_token_owner", None)
|
||||||
|
|
||||||
|
from core.models import ChatMessage
|
||||||
|
from src.llm_core import llm_call_async
|
||||||
|
from src.endpoint_resolver import build_chat_url, build_headers, build_models_url, normalize_base
|
||||||
|
|
||||||
|
message = body.message.strip()
|
||||||
|
if not message:
|
||||||
|
raise HTTPException(400, "Message is required")
|
||||||
|
|
||||||
|
session_id = body.session
|
||||||
|
sess = None
|
||||||
|
|
||||||
|
# --- Case 1: Resume an existing session ---
|
||||||
|
if session_id and session_manager:
|
||||||
|
try:
|
||||||
|
sess = session_manager.get_session(session_id)
|
||||||
|
except (KeyError, Exception):
|
||||||
|
raise HTTPException(404, "Session not found")
|
||||||
|
# SECURITY: verify the API-token's user owns this session — without
|
||||||
|
# this any token holder could resume any user's chat by passing its
|
||||||
|
# ID. The token's user is on request.state.user (set by API-token
|
||||||
|
# middleware); fall back to require_user if not present.
|
||||||
|
try:
|
||||||
|
from src.auth_helpers import get_current_user as _gcu
|
||||||
|
_tok_user = token_owner or getattr(request.state, "user", None) or _gcu(request)
|
||||||
|
except Exception:
|
||||||
|
_tok_user = None
|
||||||
|
# Strict ownership (see _caller_owns_session): fail closed so a
|
||||||
|
# null-owner / cross-owner session can't be resumed by an arbitrary
|
||||||
|
# chat-scoped token.
|
||||||
|
_sess_owner = getattr(sess, "owner", None)
|
||||||
|
if not _caller_owns_session(_sess_owner, _tok_user):
|
||||||
|
raise HTTPException(404, "Session not found")
|
||||||
|
|
||||||
|
# --- Case 2: Direct API key + model (no pre-configured endpoint needed) ---
|
||||||
|
if not sess and body.api_key:
|
||||||
|
api_key = body.api_key.strip()
|
||||||
|
model = body.model or "deepseek-chat"
|
||||||
|
|
||||||
|
# Validate only token-supplied direct base_url; auto-resolved known-provider
|
||||||
|
# URLs are not subject to extra local/LAN blocking beyond existing provider logic.
|
||||||
|
direct_base_url = body.base_url.strip().rstrip("/") if body.base_url else None
|
||||||
|
if direct_base_url:
|
||||||
|
try:
|
||||||
|
base_url = validate_public_http_url(direct_base_url)
|
||||||
|
except ValueError as e:
|
||||||
|
detail = str(e).replace("URL", "base_url", 1)
|
||||||
|
raise HTTPException(400, detail)
|
||||||
|
else:
|
||||||
|
base_url = _resolve_base_url(model, body.provider)
|
||||||
|
if not base_url:
|
||||||
|
raise HTTPException(400,
|
||||||
|
"Could not auto-detect provider. Pass base_url (e.g. 'https://api.deepseek.com/v1') "
|
||||||
|
"or provider ('deepseek', 'openai', 'groq', etc.)")
|
||||||
|
base_url = normalize_base(base_url)
|
||||||
|
endpoint_url = build_chat_url(base_url)
|
||||||
|
|
||||||
|
if not session_manager:
|
||||||
|
raise HTTPException(500, "Session manager not available")
|
||||||
|
|
||||||
|
sid = str(uuid.uuid4())
|
||||||
|
sess = session_manager.create_session(
|
||||||
|
session_id=sid, name="API Chat", endpoint_url=endpoint_url,
|
||||||
|
model=model, owner=token_owner,
|
||||||
|
)
|
||||||
|
sess.headers = build_headers(api_key, base_url)
|
||||||
|
session_manager.save_sessions()
|
||||||
|
session_id = sid
|
||||||
|
|
||||||
|
# --- Case 3: Fall back to first configured ModelEndpoint ---
|
||||||
|
if not sess:
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
ep = _select_api_chat_fallback_endpoint(db, token_owner)
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
if not ep:
|
||||||
|
raise HTTPException(400,
|
||||||
|
"No session, api_key, or configured endpoints. "
|
||||||
|
"Pass api_key + model, or configure an endpoint in Admin.")
|
||||||
|
|
||||||
|
base_url = normalize_base(ep.base_url)
|
||||||
|
endpoint_url = build_chat_url(base_url)
|
||||||
|
model = body.model or "auto"
|
||||||
|
api_key = ep.api_key
|
||||||
|
if getattr(ep, "provider_auth_id", None):
|
||||||
|
try:
|
||||||
|
from src.endpoint_resolver import resolve_endpoint_runtime
|
||||||
|
base_url, api_key = resolve_endpoint_runtime(ep, owner=token_owner)
|
||||||
|
endpoint_url = build_chat_url(base_url)
|
||||||
|
except Exception:
|
||||||
|
raise HTTPException(500, "Could not resolve endpoint credentials")
|
||||||
|
|
||||||
|
if model == "auto":
|
||||||
|
try:
|
||||||
|
async with httpx.AsyncClient(timeout=5) as client:
|
||||||
|
models_url = build_models_url(base_url)
|
||||||
|
hdrs = build_headers(api_key, base_url)
|
||||||
|
if models_url:
|
||||||
|
resp = await client.get(models_url, headers=hdrs)
|
||||||
|
resp.raise_for_status()
|
||||||
|
data = resp.json()
|
||||||
|
items = data if isinstance(data, list) else (data.get("data") or [])
|
||||||
|
ids = [m.get("id") for m in items if isinstance(m, dict) and m.get("id")]
|
||||||
|
if not ids and isinstance(data, dict):
|
||||||
|
ids = [
|
||||||
|
m.get("name") or m.get("model")
|
||||||
|
for m in (data.get("models") or [])
|
||||||
|
if m.get("name") or m.get("model")
|
||||||
|
]
|
||||||
|
else:
|
||||||
|
import json as _json
|
||||||
|
ids = _json.loads(ep.cached_models or "[]")
|
||||||
|
model = ids[0] if ids else "auto"
|
||||||
|
except Exception:
|
||||||
|
raise HTTPException(500, "Could not discover models from endpoint")
|
||||||
|
|
||||||
|
if not session_manager:
|
||||||
|
raise HTTPException(500, "Session manager not available")
|
||||||
|
|
||||||
|
sid = str(uuid.uuid4())
|
||||||
|
sess = session_manager.create_session(
|
||||||
|
session_id=sid, name="API Chat", endpoint_url=endpoint_url,
|
||||||
|
model=model, owner=token_owner,
|
||||||
|
)
|
||||||
|
if api_key:
|
||||||
|
sess.headers = build_headers(api_key, base_url)
|
||||||
|
session_manager.save_sessions()
|
||||||
|
session_id = sid
|
||||||
|
|
||||||
|
# --- Send message and get response ---
|
||||||
|
sess.add_message(ChatMessage("user", message))
|
||||||
|
|
||||||
|
messages = [{"role": m.role, "content": m.content} for m in sess.history]
|
||||||
|
|
||||||
|
reply = await llm_call_async(
|
||||||
|
sess.endpoint_url, sess.model, messages,
|
||||||
|
headers=sess.headers, timeout=120,
|
||||||
|
)
|
||||||
|
sess.add_message(ChatMessage("assistant", reply))
|
||||||
|
session_manager.save_sessions()
|
||||||
|
|
||||||
|
webhook_manager.fire_and_forget("chat.completed", {
|
||||||
|
"session_id": session_id, "model": sess.model,
|
||||||
|
"user_message": message[:2000], "response": reply[:2000],
|
||||||
|
})
|
||||||
|
|
||||||
|
return {"response": reply, "session_id": session_id, "model": sess.model}
|
||||||
|
|
||||||
|
return router
|
||||||
+12
-391
@@ -1,395 +1,16 @@
|
|||||||
"""Webhook, API Token, and sync chat routes."""
|
"""Backward-compat shim — canonical location is routes/webhook/webhook_routes.py.
|
||||||
|
|
||||||
import uuid
|
This module is replaced in ``sys.modules`` by the canonical module object so
|
||||||
import logging
|
that ``import routes.webhook_routes``, ``from routes.webhook_routes import X``,
|
||||||
from typing import Optional
|
``importlib.import_module("routes.webhook_routes")``, and the
|
||||||
|
``__import__("routes.webhook_routes", fromlist=[...])`` + ``setattr(wh_mod,
|
||||||
|
...)`` pattern used by test_null_owner_gates.py all operate on the *same*
|
||||||
|
object. Keeps existing import paths working after slice 2l (#4082/#4071).
|
||||||
|
Source-introspection tests read the canonical file by path.
|
||||||
|
"""
|
||||||
|
|
||||||
import httpx
|
import sys as _sys
|
||||||
from fastapi import APIRouter, HTTPException, Request, Form
|
|
||||||
from pydantic import BaseModel, Field
|
|
||||||
|
|
||||||
from core.database import SessionLocal, Webhook, ModelEndpoint
|
from routes.webhook import webhook_routes as _canonical # noqa: F401
|
||||||
from src.auth_helpers import owner_filter
|
|
||||||
from src.url_security import validate_public_http_url
|
|
||||||
from src.webhook_manager import WebhookManager, validate_webhook_url, validate_events
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
_sys.modules[__name__] = _canonical
|
||||||
|
|
||||||
router = APIRouter(prefix="/api", tags=["webhooks"])
|
|
||||||
|
|
||||||
# Input limits
|
|
||||||
MAX_NAME_LEN = 100
|
|
||||||
MAX_URL_LEN = 2048
|
|
||||||
MAX_SECRET_LEN = 256
|
|
||||||
MAX_MESSAGE_LEN = 32_000
|
|
||||||
|
|
||||||
|
|
||||||
from core.middleware import require_admin as _require_admin
|
|
||||||
|
|
||||||
|
|
||||||
def _select_api_chat_fallback_endpoint(db, token_owner: Optional[str]):
|
|
||||||
"""First enabled ModelEndpoint visible to token_owner — their own rows plus
|
|
||||||
legacy null-owner ("shared") rows. Owner-scoped: an unscoped .first() would
|
|
||||||
let a chat-scoped token fall back onto another user's private endpoint and
|
|
||||||
silently spend that owner's API key/quota. Prefer owner rows before shared
|
|
||||||
rows. Fails closed to null-owner rows only when token_owner is absent.
|
|
||||||
Does not validate base_url — admin-configured local/LAN endpoints remain allowed.
|
|
||||||
"""
|
|
||||||
query = db.query(ModelEndpoint).filter(ModelEndpoint.is_enabled == True) # noqa: E712
|
|
||||||
if token_owner:
|
|
||||||
query = owner_filter(query, ModelEndpoint, token_owner)
|
|
||||||
return query.order_by(ModelEndpoint.owner.desc(), ModelEndpoint.created_at).first()
|
|
||||||
return query.filter(ModelEndpoint.owner == None).order_by(ModelEndpoint.created_at).first() # noqa: E711
|
|
||||||
|
|
||||||
|
|
||||||
def _caller_owns_session(sess_owner, caller) -> bool:
|
|
||||||
"""Strict session-ownership gate for the token-authenticated sync-chat
|
|
||||||
endpoint (`POST /api/v1/chat`).
|
|
||||||
|
|
||||||
Mirrors ``_verify_session_owner`` in session_routes.py and the null-owner
|
|
||||||
gates in notes/calendar/gallery: a caller may resume a session ONLY when
|
|
||||||
its owner matches them exactly. A null/empty session owner (legacy or
|
|
||||||
migrated rows) is deliberately NOT resumable by an arbitrary token — the
|
|
||||||
old ``sess_owner and sess_owner != caller`` form skipped the check whenever
|
|
||||||
``sess_owner`` was falsy, so any chat-scoped token (e.g. a paired mobile
|
|
||||||
device) could resume such a session, inject a message, and read back its
|
|
||||||
history and reuse the owner's endpoint credentials. Fail closed: an
|
|
||||||
unresolvable caller also returns False.
|
|
||||||
"""
|
|
||||||
if not caller:
|
|
||||||
return False
|
|
||||||
return sess_owner == caller
|
|
||||||
|
|
||||||
|
|
||||||
def setup_webhook_routes(
|
|
||||||
webhook_manager: WebhookManager,
|
|
||||||
auth_manager,
|
|
||||||
session_manager=None,
|
|
||||||
api_key_manager=None,
|
|
||||||
) -> APIRouter:
|
|
||||||
|
|
||||||
@router.get("/webhooks")
|
|
||||||
def list_webhooks(request: Request):
|
|
||||||
_require_admin(request)
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
hooks = db.query(Webhook).all()
|
|
||||||
return [
|
|
||||||
{
|
|
||||||
"id": w.id,
|
|
||||||
"name": w.name,
|
|
||||||
"url": w.url,
|
|
||||||
"has_secret": bool(w.secret),
|
|
||||||
"events": w.events.split(",") if w.events else [],
|
|
||||||
"is_active": w.is_active,
|
|
||||||
"last_triggered_at": w.last_triggered_at.isoformat() if w.last_triggered_at else None,
|
|
||||||
"last_status_code": w.last_status_code,
|
|
||||||
"last_error": w.last_error,
|
|
||||||
"created_at": w.created_at.isoformat() if w.created_at else None,
|
|
||||||
}
|
|
||||||
for w in hooks
|
|
||||||
]
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
|
|
||||||
@router.post("/webhooks")
|
|
||||||
def create_webhook(
|
|
||||||
request: Request,
|
|
||||||
name: str = Form(""),
|
|
||||||
url: str = Form(""),
|
|
||||||
secret: str = Form(""),
|
|
||||||
events: str = Form(""),
|
|
||||||
):
|
|
||||||
_require_admin(request)
|
|
||||||
name = name.strip()[:MAX_NAME_LEN]
|
|
||||||
if not name:
|
|
||||||
raise HTTPException(400, "Webhook name is required")
|
|
||||||
try:
|
|
||||||
url = validate_webhook_url(url)
|
|
||||||
except ValueError as e:
|
|
||||||
raise HTTPException(400, str(e))
|
|
||||||
try:
|
|
||||||
events = validate_events(events)
|
|
||||||
except ValueError as e:
|
|
||||||
raise HTTPException(400, str(e))
|
|
||||||
|
|
||||||
secret_val = secret.strip()[:MAX_SECRET_LEN] or None
|
|
||||||
# Encrypt the secret at rest using the same Fernet key as API keys
|
|
||||||
encrypted_secret = None
|
|
||||||
if secret_val and api_key_manager:
|
|
||||||
encrypted_secret = api_key_manager.encrypt_api_key(secret_val)
|
|
||||||
elif secret_val:
|
|
||||||
encrypted_secret = secret_val # Fallback if no encryption available
|
|
||||||
|
|
||||||
webhook_id = str(uuid.uuid4())[:8]
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
db.add(Webhook(
|
|
||||||
id=webhook_id,
|
|
||||||
name=name,
|
|
||||||
url=url,
|
|
||||||
secret=encrypted_secret,
|
|
||||||
events=events,
|
|
||||||
is_active=True,
|
|
||||||
))
|
|
||||||
db.commit()
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
|
|
||||||
return {"id": webhook_id, "name": name}
|
|
||||||
|
|
||||||
@router.post("/webhooks/{webhook_id}/test")
|
|
||||||
async def test_webhook(request: Request, webhook_id: str):
|
|
||||||
_require_admin(request)
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
wh = db.query(Webhook).filter(Webhook.id == webhook_id).first()
|
|
||||||
if not wh:
|
|
||||||
raise HTTPException(404, "Webhook not found")
|
|
||||||
url, secret = wh.url, wh.secret
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
|
|
||||||
await webhook_manager.deliver_test(webhook_id, url, secret)
|
|
||||||
return {"status": "sent"}
|
|
||||||
|
|
||||||
@router.patch("/webhooks/{webhook_id}")
|
|
||||||
def toggle_webhook(request: Request, webhook_id: str):
|
|
||||||
_require_admin(request)
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
wh = db.query(Webhook).filter(Webhook.id == webhook_id).first()
|
|
||||||
if not wh:
|
|
||||||
raise HTTPException(404, "Webhook not found")
|
|
||||||
wh.is_active = not wh.is_active
|
|
||||||
db.commit()
|
|
||||||
return {"id": webhook_id, "is_active": wh.is_active}
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
|
|
||||||
@router.delete("/webhooks/{webhook_id}")
|
|
||||||
def delete_webhook(request: Request, webhook_id: str):
|
|
||||||
_require_admin(request)
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
deleted = db.query(Webhook).filter(Webhook.id == webhook_id).delete()
|
|
||||||
db.commit()
|
|
||||||
if not deleted:
|
|
||||||
raise HTTPException(404, "Webhook not found")
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
return {"status": "deleted"}
|
|
||||||
|
|
||||||
# ================================================================
|
|
||||||
# Sync Chat Endpoint (for n8n / Make / Activepieces)
|
|
||||||
# ================================================================
|
|
||||||
|
|
||||||
# Known provider base URLs — auto-resolved from api_key prefix or model name
|
|
||||||
KNOWN_PROVIDERS = {
|
|
||||||
"deepseek": "https://api.deepseek.com/v1",
|
|
||||||
"openai": "https://api.openai.com/v1",
|
|
||||||
"mistral": "https://api.mistral.ai/v1",
|
|
||||||
"groq": "https://api.groq.com/openai/v1",
|
|
||||||
"together": "https://api.together.xyz/v1",
|
|
||||||
"openrouter": "https://openrouter.ai/api/v1",
|
|
||||||
"ollama": "https://ollama.com/api",
|
|
||||||
"opencode-zen": "https://opencode.ai/zen/v1",
|
|
||||||
"opencode-go": "https://opencode.ai/zen/go/v1",
|
|
||||||
"fireworks": "https://api.fireworks.ai/inference/v1",
|
|
||||||
"venice": "https://api.venice.ai/api/v1",
|
|
||||||
"kimi-code": "https://api.kimi.com/coding/v1",
|
|
||||||
"kimicode": "https://api.kimi.com/coding/v1",
|
|
||||||
}
|
|
||||||
|
|
||||||
# Model prefix → provider mapping for auto-detection
|
|
||||||
MODEL_PROVIDER_MAP = {
|
|
||||||
"deepseek": "deepseek",
|
|
||||||
"gpt-": "openai",
|
|
||||||
"o1": "openai",
|
|
||||||
"o3": "openai",
|
|
||||||
"o4": "openai",
|
|
||||||
"mistral": "mistral",
|
|
||||||
"llama": "groq",
|
|
||||||
"mixtral": "groq",
|
|
||||||
"kimi-for-coding": "kimi-code",
|
|
||||||
"kimi": "kimi-code",
|
|
||||||
}
|
|
||||||
|
|
||||||
def _resolve_base_url(model: Optional[str], provider: Optional[str]) -> Optional[str]:
|
|
||||||
"""Try to auto-resolve a base URL from provider name or model prefix."""
|
|
||||||
if provider and provider.lower() in KNOWN_PROVIDERS:
|
|
||||||
return KNOWN_PROVIDERS[provider.lower()]
|
|
||||||
if model:
|
|
||||||
model_lower = model.lower()
|
|
||||||
for prefix, prov in MODEL_PROVIDER_MAP.items():
|
|
||||||
if model_lower.startswith(prefix):
|
|
||||||
return KNOWN_PROVIDERS[prov]
|
|
||||||
return None
|
|
||||||
|
|
||||||
class SyncChatRequest(BaseModel):
|
|
||||||
message: str = Field(..., max_length=MAX_MESSAGE_LEN)
|
|
||||||
model: Optional[str] = Field(None, max_length=200)
|
|
||||||
session: Optional[str] = Field(None, max_length=100)
|
|
||||||
api_key: Optional[str] = Field(None, max_length=256)
|
|
||||||
base_url: Optional[str] = Field(None, max_length=MAX_URL_LEN)
|
|
||||||
provider: Optional[str] = Field(None, max_length=50)
|
|
||||||
|
|
||||||
@router.post("/v1/chat")
|
|
||||||
async def sync_chat(request: Request, body: SyncChatRequest):
|
|
||||||
if not getattr(request.state, "api_token", False):
|
|
||||||
raise HTTPException(403, "This endpoint requires an API token")
|
|
||||||
scopes = set(getattr(request.state, "api_token_scopes", []) or [])
|
|
||||||
if "chat" not in scopes:
|
|
||||||
raise HTTPException(403, "API token is not scoped for chat")
|
|
||||||
token_owner = getattr(request.state, "api_token_owner", None)
|
|
||||||
|
|
||||||
from core.models import ChatMessage
|
|
||||||
from src.llm_core import llm_call_async
|
|
||||||
from src.endpoint_resolver import build_chat_url, build_headers, build_models_url, normalize_base
|
|
||||||
|
|
||||||
message = body.message.strip()
|
|
||||||
if not message:
|
|
||||||
raise HTTPException(400, "Message is required")
|
|
||||||
|
|
||||||
session_id = body.session
|
|
||||||
sess = None
|
|
||||||
|
|
||||||
# --- Case 1: Resume an existing session ---
|
|
||||||
if session_id and session_manager:
|
|
||||||
try:
|
|
||||||
sess = session_manager.get_session(session_id)
|
|
||||||
except (KeyError, Exception):
|
|
||||||
raise HTTPException(404, "Session not found")
|
|
||||||
# SECURITY: verify the API-token's user owns this session — without
|
|
||||||
# this any token holder could resume any user's chat by passing its
|
|
||||||
# ID. The token's user is on request.state.user (set by API-token
|
|
||||||
# middleware); fall back to require_user if not present.
|
|
||||||
try:
|
|
||||||
from src.auth_helpers import get_current_user as _gcu
|
|
||||||
_tok_user = token_owner or getattr(request.state, "user", None) or _gcu(request)
|
|
||||||
except Exception:
|
|
||||||
_tok_user = None
|
|
||||||
# Strict ownership (see _caller_owns_session): fail closed so a
|
|
||||||
# null-owner / cross-owner session can't be resumed by an arbitrary
|
|
||||||
# chat-scoped token.
|
|
||||||
_sess_owner = getattr(sess, "owner", None)
|
|
||||||
if not _caller_owns_session(_sess_owner, _tok_user):
|
|
||||||
raise HTTPException(404, "Session not found")
|
|
||||||
|
|
||||||
# --- Case 2: Direct API key + model (no pre-configured endpoint needed) ---
|
|
||||||
if not sess and body.api_key:
|
|
||||||
api_key = body.api_key.strip()
|
|
||||||
model = body.model or "deepseek-chat"
|
|
||||||
|
|
||||||
# Validate only token-supplied direct base_url; auto-resolved known-provider
|
|
||||||
# URLs are not subject to extra local/LAN blocking beyond existing provider logic.
|
|
||||||
direct_base_url = body.base_url.strip().rstrip("/") if body.base_url else None
|
|
||||||
if direct_base_url:
|
|
||||||
try:
|
|
||||||
base_url = validate_public_http_url(direct_base_url)
|
|
||||||
except ValueError as e:
|
|
||||||
detail = str(e).replace("URL", "base_url", 1)
|
|
||||||
raise HTTPException(400, detail)
|
|
||||||
else:
|
|
||||||
base_url = _resolve_base_url(model, body.provider)
|
|
||||||
if not base_url:
|
|
||||||
raise HTTPException(400,
|
|
||||||
"Could not auto-detect provider. Pass base_url (e.g. 'https://api.deepseek.com/v1') "
|
|
||||||
"or provider ('deepseek', 'openai', 'groq', etc.)")
|
|
||||||
base_url = normalize_base(base_url)
|
|
||||||
endpoint_url = build_chat_url(base_url)
|
|
||||||
|
|
||||||
if not session_manager:
|
|
||||||
raise HTTPException(500, "Session manager not available")
|
|
||||||
|
|
||||||
sid = str(uuid.uuid4())
|
|
||||||
sess = session_manager.create_session(
|
|
||||||
session_id=sid, name="API Chat", endpoint_url=endpoint_url,
|
|
||||||
model=model, owner=token_owner,
|
|
||||||
)
|
|
||||||
sess.headers = build_headers(api_key, base_url)
|
|
||||||
session_manager.save_sessions()
|
|
||||||
session_id = sid
|
|
||||||
|
|
||||||
# --- Case 3: Fall back to first configured ModelEndpoint ---
|
|
||||||
if not sess:
|
|
||||||
db = SessionLocal()
|
|
||||||
try:
|
|
||||||
ep = _select_api_chat_fallback_endpoint(db, token_owner)
|
|
||||||
finally:
|
|
||||||
db.close()
|
|
||||||
|
|
||||||
if not ep:
|
|
||||||
raise HTTPException(400,
|
|
||||||
"No session, api_key, or configured endpoints. "
|
|
||||||
"Pass api_key + model, or configure an endpoint in Admin.")
|
|
||||||
|
|
||||||
base_url = normalize_base(ep.base_url)
|
|
||||||
endpoint_url = build_chat_url(base_url)
|
|
||||||
model = body.model or "auto"
|
|
||||||
api_key = ep.api_key
|
|
||||||
if getattr(ep, "provider_auth_id", None):
|
|
||||||
try:
|
|
||||||
from src.endpoint_resolver import resolve_endpoint_runtime
|
|
||||||
base_url, api_key = resolve_endpoint_runtime(ep, owner=token_owner)
|
|
||||||
endpoint_url = build_chat_url(base_url)
|
|
||||||
except Exception:
|
|
||||||
raise HTTPException(500, "Could not resolve endpoint credentials")
|
|
||||||
|
|
||||||
if model == "auto":
|
|
||||||
try:
|
|
||||||
async with httpx.AsyncClient(timeout=5) as client:
|
|
||||||
models_url = build_models_url(base_url)
|
|
||||||
hdrs = build_headers(api_key, base_url)
|
|
||||||
if models_url:
|
|
||||||
resp = await client.get(models_url, headers=hdrs)
|
|
||||||
resp.raise_for_status()
|
|
||||||
data = resp.json()
|
|
||||||
items = data if isinstance(data, list) else (data.get("data") or [])
|
|
||||||
ids = [m.get("id") for m in items if isinstance(m, dict) and m.get("id")]
|
|
||||||
if not ids and isinstance(data, dict):
|
|
||||||
ids = [
|
|
||||||
m.get("name") or m.get("model")
|
|
||||||
for m in (data.get("models") or [])
|
|
||||||
if m.get("name") or m.get("model")
|
|
||||||
]
|
|
||||||
else:
|
|
||||||
import json as _json
|
|
||||||
ids = _json.loads(ep.cached_models or "[]")
|
|
||||||
model = ids[0] if ids else "auto"
|
|
||||||
except Exception:
|
|
||||||
raise HTTPException(500, "Could not discover models from endpoint")
|
|
||||||
|
|
||||||
if not session_manager:
|
|
||||||
raise HTTPException(500, "Session manager not available")
|
|
||||||
|
|
||||||
sid = str(uuid.uuid4())
|
|
||||||
sess = session_manager.create_session(
|
|
||||||
session_id=sid, name="API Chat", endpoint_url=endpoint_url,
|
|
||||||
model=model, owner=token_owner,
|
|
||||||
)
|
|
||||||
if api_key:
|
|
||||||
sess.headers = build_headers(api_key, base_url)
|
|
||||||
session_manager.save_sessions()
|
|
||||||
session_id = sid
|
|
||||||
|
|
||||||
# --- Send message and get response ---
|
|
||||||
sess.add_message(ChatMessage("user", message))
|
|
||||||
|
|
||||||
messages = [{"role": m.role, "content": m.content} for m in sess.history]
|
|
||||||
|
|
||||||
reply = await llm_call_async(
|
|
||||||
sess.endpoint_url, sess.model, messages,
|
|
||||||
headers=sess.headers, timeout=120,
|
|
||||||
)
|
|
||||||
sess.add_message(ChatMessage("assistant", reply))
|
|
||||||
session_manager.save_sessions()
|
|
||||||
|
|
||||||
webhook_manager.fire_and_forget("chat.completed", {
|
|
||||||
"session_id": session_id, "model": sess.model,
|
|
||||||
"user_message": message[:2000], "response": reply[:2000],
|
|
||||||
})
|
|
||||||
|
|
||||||
return {"response": reply, "session_id": session_id, "model": sess.model}
|
|
||||||
|
|
||||||
return router
|
|
||||||
|
|||||||
@@ -50,7 +50,7 @@ import json
|
|||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from datetime import datetime
|
from datetime import datetime, timezone
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -441,4 +441,4 @@ class Skill:
|
|||||||
|
|
||||||
|
|
||||||
def _now_iso() -> str:
|
def _now_iso() -> str:
|
||||||
return datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ")
|
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
"""Multi-provider TTS service — dispatches to local Kokoro, OpenAI-compatible API, or browser."""
|
"""Multi-provider TTS service — dispatches to local Kokoro, OpenAI-compatible API, or browser."""
|
||||||
|
|
||||||
import io
|
import io
|
||||||
|
import os
|
||||||
import wave
|
import wave
|
||||||
import logging
|
import logging
|
||||||
import hashlib
|
import hashlib
|
||||||
@@ -41,6 +42,11 @@ class TTSService:
|
|||||||
self.cache_dir = Path(cache_dir)
|
self.cache_dir = Path(cache_dir)
|
||||||
self.cache_dir.mkdir(parents=True, exist_ok=True)
|
self.cache_dir.mkdir(parents=True, exist_ok=True)
|
||||||
self._kokoro = None # lazy-init
|
self._kokoro = None # lazy-init
|
||||||
|
|
||||||
|
try:
|
||||||
|
self.max_cache_bytes = int(os.getenv("ODYSSEUS_TTS_CACHE_MAX_BYTES", 500 * 1024 * 1024))
|
||||||
|
except ValueError:
|
||||||
|
self.max_cache_bytes = 500 * 1024 * 1024
|
||||||
|
|
||||||
# ── Settings ──
|
# ── Settings ──
|
||||||
|
|
||||||
@@ -89,6 +95,53 @@ class TTSService:
|
|||||||
ext = ".mp3" if (len(data) >= 3 and (data[:3] == b'ID3' or (data[0] == 0xff and (data[1] & 0xe0) == 0xe0))) else ".wav"
|
ext = ".mp3" if (len(data) >= 3 and (data[:3] == b'ID3' or (data[0] == 0xff and (data[1] & 0xe0) == 0xe0))) else ".wav"
|
||||||
(self.cache_dir / f"{key}{ext}").write_bytes(data)
|
(self.cache_dir / f"{key}{ext}").write_bytes(data)
|
||||||
|
|
||||||
|
self._enforce_cache_limit()
|
||||||
|
|
||||||
|
def _enforce_cache_limit(self):
|
||||||
|
"""Evicts oldest files if the cache exceeds the configured byte limit."""
|
||||||
|
if self.max_cache_bytes <= 0:
|
||||||
|
return
|
||||||
|
|
||||||
|
try:
|
||||||
|
files = []
|
||||||
|
total_size = 0
|
||||||
|
|
||||||
|
# Safely scan files and sum sizes, ignoring files deleted mid-scan
|
||||||
|
for f in self.cache_dir.iterdir():
|
||||||
|
try:
|
||||||
|
if f.is_file() and f.suffix.lower() in (".mp3", ".wav"):
|
||||||
|
files.append(f)
|
||||||
|
total_size += f.stat().st_size
|
||||||
|
except OSError:
|
||||||
|
continue
|
||||||
|
|
||||||
|
if total_size > self.max_cache_bytes:
|
||||||
|
logger.info(
|
||||||
|
f"TTS cache ({total_size} bytes) exceeded limit ({self.max_cache_bytes} bytes). Evicting oldest files."
|
||||||
|
)
|
||||||
|
|
||||||
|
# Sort files by modification time (oldest first)
|
||||||
|
try:
|
||||||
|
files.sort(key=lambda f: f.stat().st_mtime)
|
||||||
|
except OSError as e:
|
||||||
|
logger.warning(f"Failed to sort cache files by mtime: {e}")
|
||||||
|
|
||||||
|
# Trim down to 80% of max capacity
|
||||||
|
target_size = self.max_cache_bytes * 0.8
|
||||||
|
|
||||||
|
while files and total_size > target_size:
|
||||||
|
f = files.pop(0)
|
||||||
|
try:
|
||||||
|
size = f.stat().st_size
|
||||||
|
f.unlink()
|
||||||
|
total_size -= size
|
||||||
|
except OSError as e:
|
||||||
|
logger.warning(f"Failed to evict cache file {f}: {e}")
|
||||||
|
continue
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning(f"Error enforcing TTS cache limit: {e}", exc_info=True)
|
||||||
|
|
||||||
def clear_cache(self):
|
def clear_cache(self):
|
||||||
count = 0
|
count = 0
|
||||||
for f in self.cache_dir.glob("*.*"):
|
for f in self.cache_dir.glob("*.*"):
|
||||||
|
|||||||
+1
-1
@@ -12,7 +12,7 @@ import json
|
|||||||
import re
|
import re
|
||||||
import time
|
import time
|
||||||
import logging
|
import logging
|
||||||
from typing import AsyncGenerator, List, Dict, Optional, Set
|
from typing import Any, AsyncGenerator, List, Dict, Optional, Set
|
||||||
from urllib.parse import urlparse
|
from urllib.parse import urlparse
|
||||||
|
|
||||||
from src.llm_core import (
|
from src.llm_core import (
|
||||||
|
|||||||
+19
-7
@@ -1237,15 +1237,27 @@ def _anthropic_rejects_temperature(model: str) -> bool:
|
|||||||
return False
|
return False
|
||||||
# `(?<![a-z])` anchors "opus" to a word boundary so a substring match like
|
# `(?<![a-z])` anchors "opus" to a word boundary so a substring match like
|
||||||
# `oct-opus`/`octopus-4-8` can't be read as Opus (it would otherwise strip
|
# `oct-opus`/`octopus-4-8` can't be read as Opus (it would otherwise strip
|
||||||
# temperature). Cap the minor at 1-2 digits and forbid a trailing digit so a
|
# temperature). Both version components are capped at 1-2 digits and forbid a
|
||||||
# dated id like `claude-opus-4-20250514` (Opus 4.0) parses as major-only (no
|
# trailing digit, so an 8-digit date can never be read as a version number:
|
||||||
# minor match, kept) instead of reading the date `20250514` as a giant minor
|
# `claude-opus-4-20250514` (Opus 4.0) parses as major-only rather than reading
|
||||||
# that would falsely test >= 4.7. Dated 4.7+ snapshots (`claude-opus-4-7-
|
# `20250514` as a giant minor, and `claude-3-opus-20240229` (legacy Claude 3
|
||||||
# 20260201`) keep their explicit minor and are still matched.
|
# Opus, date directly after "opus-") fails to match at all rather than reading
|
||||||
match = re.search(r"(?<![a-z])opus[-_]?(\d+)[-_.](\d{1,2})(?!\d)", model.lower())
|
# the date as a giant major. Dated 4.7+ snapshots (`claude-opus-4-7-20260201`)
|
||||||
|
# keep their explicit minor and are still matched.
|
||||||
|
#
|
||||||
|
# The minor is optional and a missing minor reads as `.0`, so major-only ids
|
||||||
|
# like `claude-opus-5` are correctly treated as >= 4.7 (issue #5753). Without
|
||||||
|
# this, every Opus 5 call kept `temperature` and failed with HTTP 400 — visible
|
||||||
|
# only on paths that pass a temperature, e.g. scheduled tasks inheriting
|
||||||
|
# `stream_agent_loop`'s 0.3 default, which returned empty responses.
|
||||||
|
match = re.search(
|
||||||
|
r"(?<![a-z])opus[-_]?(\d{1,2})(?!\d)(?:[-_.](\d{1,2})(?!\d))?", model.lower()
|
||||||
|
)
|
||||||
if not match:
|
if not match:
|
||||||
return False
|
return False
|
||||||
return (int(match.group(1)), int(match.group(2))) >= (4, 7)
|
major = int(match.group(1))
|
||||||
|
minor = int(match.group(2)) if match.group(2) else 0
|
||||||
|
return (major, minor) >= (4, 7)
|
||||||
|
|
||||||
# Reasoning effort level sent to Mistral thinking-capable models. Mistral's
|
# Reasoning effort level sent to Mistral thinking-capable models. Mistral's
|
||||||
# API accepts "high", "medium", "low", "none" — see
|
# API accepts "high", "medium", "low", "none" — see
|
||||||
|
|||||||
+13
-7
@@ -758,30 +758,36 @@ export function mdToHtml(src, opts) {
|
|||||||
// Remove empty paragraphs
|
// Remove empty paragraphs
|
||||||
s = s.replace(/<p><\/p>/g, '');
|
s = s.replace(/<p><\/p>/g, '');
|
||||||
|
|
||||||
|
// Every restore below passes a function replacer rather than the block string
|
||||||
|
// itself. With a string replacement, `String.replace` reads `$&`, `` $` ``,
|
||||||
|
// `$'` and `$$` in the *replacement* as substitution patterns, so a restored
|
||||||
|
// block containing them is corrupted: `$&` re-inserts the placeholder, `` $` ``
|
||||||
|
// and `$'` splice in the surrounding document, and `$$` collapses to `$`. Those
|
||||||
|
// sequences are ordinary content in fenced code (`perl -pe 's/x/$& y/'`,
|
||||||
|
// `echo "$$USD"`). A function replacer inserts its return value verbatim.
|
||||||
|
|
||||||
// CRITICAL: Restore allowed HTML blocks first
|
// CRITICAL: Restore allowed HTML blocks first
|
||||||
allowedHtmlBlocks.forEach((block, index) => {
|
allowedHtmlBlocks.forEach((block, index) => {
|
||||||
s = s.replace(`___ALLOWED_HTML_${index}___`, block);
|
s = s.replace(`___ALLOWED_HTML_${index}___`, () => block);
|
||||||
});
|
});
|
||||||
|
|
||||||
// Restore math blocks
|
// Restore math blocks
|
||||||
mathBlocks.forEach((block, index) => {
|
mathBlocks.forEach((block, index) => {
|
||||||
s = s.replace(`___MATH_BLOCK_${index}___`, block);
|
s = s.replace(`___MATH_BLOCK_${index}___`, () => block);
|
||||||
});
|
});
|
||||||
|
|
||||||
// Restore mermaid diagram blocks
|
// Restore mermaid diagram blocks
|
||||||
mermaidBlocks.forEach((block, index) => {
|
mermaidBlocks.forEach((block, index) => {
|
||||||
s = s.replace(`___MERMAID_BLOCK_${index}___`, block);
|
s = s.replace(`___MERMAID_BLOCK_${index}___`, () => block);
|
||||||
});
|
});
|
||||||
|
|
||||||
// CRITICAL: Restore code blocks at the end
|
// CRITICAL: Restore code blocks at the end
|
||||||
codeBlocks.forEach((block, index) => {
|
codeBlocks.forEach((block, index) => {
|
||||||
s = s.replace(`___CODE_BLOCK_${index}___`, block);
|
s = s.replace(`___CODE_BLOCK_${index}___`, () => block);
|
||||||
});
|
});
|
||||||
|
|
||||||
// Restore inline code spans last, so placeholders carried inside restored
|
// Restore inline code spans last, so placeholders carried inside restored
|
||||||
// <a>/allowed-HTML blocks are resolved too. The function replacer keeps the
|
// <a>/allowed-HTML blocks are resolved too.
|
||||||
// escaped code literal — e.g. a shell snippet like `echo $1` is not treated
|
|
||||||
// as a regex back-reference.
|
|
||||||
inlineCodeBlocks.forEach((block, index) => {
|
inlineCodeBlocks.forEach((block, index) => {
|
||||||
s = s.replace(`___INLINE_CODE_${index}___`, () => block);
|
s = s.replace(`___INLINE_CODE_${index}___`, () => block);
|
||||||
});
|
});
|
||||||
|
|||||||
+25
-22
@@ -3031,12 +3031,14 @@ async function initEmailAccountsSettings() {
|
|||||||
const body = {
|
const body = {
|
||||||
name: el('eaf-name').value.trim() || el('eaf-from').value.trim(),
|
name: el('eaf-name').value.trim() || el('eaf-from').value.trim(),
|
||||||
from_address: el('eaf-from').value.trim(),
|
from_address: el('eaf-from').value.trim(),
|
||||||
|
display_name: el('eaf-display-name').value.trim(),
|
||||||
imap_host: el('eaf-imap-host').value.trim(),
|
imap_host: el('eaf-imap-host').value.trim(),
|
||||||
imap_port: parseInt(el('eaf-imap-port').value) || 993,
|
imap_port: parseInt(el('eaf-imap-port').value) || 993,
|
||||||
imap_user: el('eaf-imap-user').value.trim(),
|
imap_user: el('eaf-imap-user').value.trim(),
|
||||||
imap_starttls: el('eaf-imap-starttls').checked,
|
imap_starttls: el('eaf-imap-starttls').checked,
|
||||||
smtp_host: el('eaf-smtp-host').value.trim(),
|
smtp_host: el('eaf-smtp-host').value.trim(),
|
||||||
smtp_port: parseInt(el('eaf-smtp-port').value) || 587,
|
smtp_port: parseInt(el('eaf-smtp-port').value) || 587,
|
||||||
|
smtp_security: el('eaf-smtp-security').value,
|
||||||
smtp_user: el('eaf-imap-user').value.trim(),
|
smtp_user: el('eaf-imap-user').value.trim(),
|
||||||
};
|
};
|
||||||
if (!body.name) { el('eaf-msg').textContent = 'Enter a Name or Email first'; el('eaf-msg').style.color = 'var(--red)'; return; }
|
if (!body.name) { el('eaf-msg').textContent = 'Enter a Name or Email first'; el('eaf-msg').style.color = 'var(--red)'; return; }
|
||||||
@@ -5788,29 +5790,30 @@ export function close() {
|
|||||||
window.history.replaceState(null, '', clean);
|
window.history.replaceState(null, '', clean);
|
||||||
const success = sp.has('email_oauth_success');
|
const success = sp.has('email_oauth_success');
|
||||||
const errMsg = sp.get('email_oauth_error') || '';
|
const errMsg = sp.get('email_oauth_error') || '';
|
||||||
// Open settings → integrations after the app has initialised.
|
// Open settings → integrations once the document is ready. This module owns
|
||||||
function _tryOpen() {
|
// the open() API, so it does not need to wait for a window-level alias.
|
||||||
if (window.settingsModule && typeof window.settingsModule.open === 'function') {
|
function _showResult() {
|
||||||
window.settingsModule.open('integrations');
|
open('integrations');
|
||||||
// Brief toast-style banner.
|
// Brief toast-style banner.
|
||||||
const banner = document.createElement('div');
|
const banner = document.createElement('div');
|
||||||
banner.textContent = success
|
banner.textContent = success
|
||||||
? '✓ Google account connected — email is ready'
|
? 'Google account connected — email is ready'
|
||||||
: `Google OAuth failed: ${errMsg || 'unknown error'}`;
|
: `Google OAuth failed: ${errMsg || 'unknown error'}`;
|
||||||
Object.assign(banner.style, {
|
Object.assign(banner.style, {
|
||||||
position: 'fixed', bottom: '24px', left: '50%', transform: 'translateX(-50%)',
|
position: 'fixed', bottom: '24px', left: '50%', transform: 'translateX(-50%)',
|
||||||
background: success ? 'var(--accent, #50fa7b)' : 'var(--red, #ff5555)',
|
background: success ? 'var(--accent, #50fa7b)' : 'var(--red, #ff5555)',
|
||||||
color: '#000', padding: '8px 18px', borderRadius: '6px', fontSize: '12px',
|
color: '#000', padding: '8px 18px', borderRadius: '6px', fontSize: '12px',
|
||||||
fontWeight: '600', zIndex: '99999', pointerEvents: 'none',
|
fontWeight: '600', zIndex: '99999', pointerEvents: 'none',
|
||||||
boxShadow: '0 2px 12px rgba(0,0,0,0.3)',
|
boxShadow: '0 2px 12px rgba(0,0,0,0.3)',
|
||||||
});
|
});
|
||||||
document.body.appendChild(banner);
|
document.body.appendChild(banner);
|
||||||
setTimeout(() => banner.remove(), 4000);
|
setTimeout(() => banner.remove(), 4000);
|
||||||
} else {
|
}
|
||||||
setTimeout(_tryOpen, 100);
|
if (document.readyState === 'loading') {
|
||||||
}
|
document.addEventListener('DOMContentLoaded', _showResult, { once: true });
|
||||||
|
} else {
|
||||||
|
_showResult();
|
||||||
}
|
}
|
||||||
_tryOpen();
|
|
||||||
})();
|
})();
|
||||||
|
|
||||||
const settingsModule = { open, close, initIntegrations, initUnifiedIntegrations, syncAdminVisibility, refreshAiModelEndpoints };
|
const settingsModule = { open, close, initIntegrations, initUnifiedIntegrations, syncAdminVisibility, refreshAiModelEndpoints };
|
||||||
|
|||||||
@@ -76,7 +76,7 @@ def _load_webhook_routes_for_test(monkeypatch):
|
|||||||
module_name = "routes.webhook_routes_under_test"
|
module_name = "routes.webhook_routes_under_test"
|
||||||
spec = importlib.util.spec_from_file_location(
|
spec = importlib.util.spec_from_file_location(
|
||||||
module_name,
|
module_name,
|
||||||
Path(__file__).resolve().parent.parent / "routes" / "webhook_routes.py",
|
Path(__file__).resolve().parent.parent / "routes" / "webhook" / "webhook_routes.py",
|
||||||
)
|
)
|
||||||
module = importlib.util.module_from_spec(spec)
|
module = importlib.util.module_from_spec(spec)
|
||||||
spec.loader.exec_module(module)
|
spec.loader.exec_module(module)
|
||||||
|
|||||||
@@ -0,0 +1,15 @@
|
|||||||
|
"""Regression coverage for SMTP security saved before Google OAuth."""
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
_REPO = Path(__file__).resolve().parents[1]
|
||||||
|
|
||||||
|
|
||||||
|
def test_email_tab_oauth_connect_persists_selected_smtp_security():
|
||||||
|
source = (_REPO / "static" / "js" / "settings.js").read_text(encoding="utf-8")
|
||||||
|
start = source.index("el('eaf-oauth-btn').addEventListener")
|
||||||
|
handler_body = source[start:source.index("if (!body.name)", start)]
|
||||||
|
|
||||||
|
assert "smtp_security: el('eaf-smtp-security').value" in handler_body
|
||||||
|
assert "display_name: el('eaf-display-name').value.trim()" in handler_body
|
||||||
@@ -0,0 +1,19 @@
|
|||||||
|
"""Regression coverage for the settings UI after Google OAuth redirects."""
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
_REPO = Path(__file__).resolve().parents[1]
|
||||||
|
|
||||||
|
|
||||||
|
def test_oauth_redirect_uses_the_module_local_settings_api():
|
||||||
|
source = (_REPO / "static" / "js" / "settings.js").read_text(encoding="utf-8")
|
||||||
|
handler = source[
|
||||||
|
source.index("(function _handleOauthRedirect"):
|
||||||
|
source.index("const settingsModule =")
|
||||||
|
]
|
||||||
|
|
||||||
|
assert "open('integrations');" in handler
|
||||||
|
assert "window.settingsModule" not in handler
|
||||||
|
assert "window.__odysseusAppStarted" not in handler
|
||||||
|
assert "document.addEventListener('DOMContentLoaded', _showResult, { once: true })" in handler
|
||||||
@@ -0,0 +1,86 @@
|
|||||||
|
"""Regression coverage for issue-description label lifecycle events."""
|
||||||
|
|
||||||
|
import json
|
||||||
|
import shutil
|
||||||
|
import subprocess
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
|
||||||
|
_REPO = Path(__file__).resolve().parent.parent
|
||||||
|
_CHECKER = _REPO / ".github" / "scripts" / "check-issue-description.js"
|
||||||
|
_WORKFLOW = _REPO / ".github" / "workflows" / "issue-description-check.yml"
|
||||||
|
pytestmark = pytest.mark.skipif(not shutil.which("node"), reason="node not on PATH")
|
||||||
|
|
||||||
|
|
||||||
|
def _run_closed_issue(action):
|
||||||
|
harness = r"""
|
||||||
|
const checkIssueDescription = require(process.argv[1]);
|
||||||
|
const action = process.argv[2];
|
||||||
|
const calls = [];
|
||||||
|
const unexpected = (name) => async () => {
|
||||||
|
throw new Error(`${name} should not be called for a closed issue`);
|
||||||
|
};
|
||||||
|
|
||||||
|
const github = {
|
||||||
|
rest: {
|
||||||
|
issues: {
|
||||||
|
removeLabel: async (params) => calls.push({ method: 'removeLabel', params }),
|
||||||
|
getLabel: unexpected('getLabel'),
|
||||||
|
addLabels: unexpected('addLabels'),
|
||||||
|
listComments: unexpected('listComments'),
|
||||||
|
createComment: unexpected('createComment'),
|
||||||
|
updateComment: unexpected('updateComment'),
|
||||||
|
deleteComment: unexpected('deleteComment'),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
};
|
||||||
|
const context = {
|
||||||
|
payload: {
|
||||||
|
action,
|
||||||
|
issue: { number: 42, state: 'closed', body: '', labels: [] },
|
||||||
|
},
|
||||||
|
repo: { owner: 'odysseus-dev', repo: 'odysseus' },
|
||||||
|
};
|
||||||
|
const core = {
|
||||||
|
warning: unexpected('core.warning'),
|
||||||
|
setFailed: unexpected('core.setFailed'),
|
||||||
|
};
|
||||||
|
|
||||||
|
checkIssueDescription({ github, context, core })
|
||||||
|
.then(() => process.stdout.write(JSON.stringify(calls)))
|
||||||
|
.catch((error) => {
|
||||||
|
console.error(error);
|
||||||
|
process.exitCode = 1;
|
||||||
|
});
|
||||||
|
"""
|
||||||
|
proc = subprocess.run(
|
||||||
|
["node", "-e", harness, str(_CHECKER), action],
|
||||||
|
capture_output=True,
|
||||||
|
text=True,
|
||||||
|
cwd=str(_REPO),
|
||||||
|
timeout=30,
|
||||||
|
)
|
||||||
|
assert proc.returncode == 0, proc.stderr
|
||||||
|
return json.loads(proc.stdout)
|
||||||
|
|
||||||
|
|
||||||
|
def test_workflow_handles_issue_closures():
|
||||||
|
workflow = _WORKFLOW.read_text()
|
||||||
|
assert "types: [opened, edited, reopened, closed]" in workflow
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("action", ["closed", "edited"])
|
||||||
|
def test_closed_issue_only_drops_ready_for_review(action):
|
||||||
|
assert _run_closed_issue(action) == [
|
||||||
|
{
|
||||||
|
"method": "removeLabel",
|
||||||
|
"params": {
|
||||||
|
"owner": "odysseus-dev",
|
||||||
|
"repo": "odysseus",
|
||||||
|
"issue_number": 42,
|
||||||
|
"name": "ready for review",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
]
|
||||||
@@ -29,6 +29,13 @@ from src.llm_core import _anthropic_rejects_temperature, _build_anthropic_payloa
|
|||||||
"anthropic/claude-opus-4-7", # tolerate a provider-prefixed id
|
"anthropic/claude-opus-4-7", # tolerate a provider-prefixed id
|
||||||
"claude-opus-4-10", # future minor still >= 4.7
|
"claude-opus-4-10", # future minor still >= 4.7
|
||||||
"claude-opus-5-0", # future major
|
"claude-opus-5-0", # future major
|
||||||
|
# Major-only ids: a missing minor reads as `.0`, so these are >= 4.7 too
|
||||||
|
# (issue #5753). Before the fix the version pattern required a minor, so
|
||||||
|
# these fell through to "accepts temperature" and every call 400'd.
|
||||||
|
"claude-opus-5",
|
||||||
|
"claude-opus-5-20260101", # major-only + dated snapshot
|
||||||
|
"anthropic/claude-opus-5", # major-only behind a provider prefix
|
||||||
|
"claude-opus-6", # future major-only
|
||||||
],
|
],
|
||||||
)
|
)
|
||||||
def test_opus_47_plus_rejects_temperature(model):
|
def test_opus_47_plus_rejects_temperature(model):
|
||||||
@@ -48,7 +55,10 @@ def test_opus_47_plus_rejects_temperature(model):
|
|||||||
"claude-opus-4-6-20251201", # dated 4.6 snapshot — older, still keeps temperature
|
"claude-opus-4-6-20251201", # dated 4.6 snapshot — older, still keeps temperature
|
||||||
"claude-sonnet-4-6",
|
"claude-sonnet-4-6",
|
||||||
"claude-3-5-sonnet",
|
"claude-3-5-sonnet",
|
||||||
"claude-3-opus-20240229", # legacy Claude 3 Opus — no opus-N-M pattern, kept
|
"claude-3-opus-20240229", # legacy Claude 3 Opus — date directly after
|
||||||
|
# "opus-", so the major must not swallow it as version 20240229 (that is
|
||||||
|
# what makes capping the major at 1-2 digits necessary once the minor
|
||||||
|
# became optional in #5753).
|
||||||
"claude-haiku-4-5",
|
"claude-haiku-4-5",
|
||||||
"claude-x",
|
"claude-x",
|
||||||
"octopus-4-8", # "opus" only as a substring of another word — must not match
|
"octopus-4-8", # "opus" only as a substring of another word — must not match
|
||||||
@@ -87,6 +97,20 @@ def test_payload_keeps_temperature_for_older_models():
|
|||||||
assert _payload("claude-3-5-sonnet", 1.2)["temperature"] == 1.0
|
assert _payload("claude-3-5-sonnet", 1.2)["temperature"] == 1.0
|
||||||
|
|
||||||
|
|
||||||
|
def test_payload_omits_temperature_for_major_only_opus_5():
|
||||||
|
# Issue #5753: the scheduled-task path calls stream_agent_loop() without a
|
||||||
|
# temperature and inherits its 0.3 default, so `claude-opus-5` 400'd on every
|
||||||
|
# run and surfaced as "the model returned an empty response". Interactive chat
|
||||||
|
# leaves temperature None and never hit it.
|
||||||
|
assert "temperature" not in _payload("claude-opus-5", 0.3)
|
||||||
|
|
||||||
|
|
||||||
|
def test_payload_keeps_temperature_for_legacy_claude_3_opus():
|
||||||
|
# Guards the major-digit cap: `opus-20240229` must not parse as version
|
||||||
|
# 20240229, or Claude 3 Opus would silently lose the caller's temperature.
|
||||||
|
assert _payload("claude-3-opus-20240229", 0.5)["temperature"] == 0.5
|
||||||
|
|
||||||
|
|
||||||
def test_payload_keeps_temperature_for_dated_opus_4_0():
|
def test_payload_keeps_temperature_for_dated_opus_4_0():
|
||||||
# Anthropic's dated id for Opus 4.0 (claude-opus-4-20250514) is in this repo's
|
# Anthropic's dated id for Opus 4.0 (claude-opus-4-20250514) is in this repo's
|
||||||
# ANTHROPIC_MODELS list. The date must not be misread as a >= 4.7 minor, or the
|
# ANTHROPIC_MODELS list. The date must not be misread as a >= 4.7 minor, or the
|
||||||
|
|||||||
@@ -214,6 +214,50 @@ def test_inline_code_content_is_html_escaped(node_available):
|
|||||||
assert "<b>" not in html
|
assert "<b>" not in html
|
||||||
|
|
||||||
|
|
||||||
|
def test_fenced_code_keeps_dollar_ampersand(node_available):
|
||||||
|
# Issue #5663: the block-restore pass used a string replacement, so `$&` in a
|
||||||
|
# restored block was read as "the matched text" and re-inserted the
|
||||||
|
# placeholder. `perl -pe 's/world/$& again/'` rendered as
|
||||||
|
# "s/world/___CODE_BLOCK_0___amp; again/" — the trailing "amp;" is the orphan
|
||||||
|
# left behind after `$&` consumed the `$&` of the escaped `$&`.
|
||||||
|
html = _run_markdown_case(
|
||||||
|
"```sh\necho \"hello world\" | perl -pe 's/world/$& again/'\n```"
|
||||||
|
)
|
||||||
|
|
||||||
|
assert "___CODE_BLOCK_" not in html
|
||||||
|
assert "s/world/$& again/" in html
|
||||||
|
assert "amp; again" not in html.replace("$& again", "")
|
||||||
|
|
||||||
|
|
||||||
|
def test_fenced_code_keeps_dollar_backtick_and_quote(node_available):
|
||||||
|
# `` $` `` and `$'` splice the text before/after the placeholder into the
|
||||||
|
# block. Unlike `$&` these leave no placeholder behind — the characters just
|
||||||
|
# vanish — so assert the content survives verbatim.
|
||||||
|
html = _run_markdown_case("```sh\nsed \"s/$`/x/\" && sed \"s/$'/y/\"\n```")
|
||||||
|
|
||||||
|
assert "___CODE_BLOCK_" not in html
|
||||||
|
assert "s/$`/x/" in html
|
||||||
|
assert "s/$'/y/" in html
|
||||||
|
|
||||||
|
|
||||||
|
def test_fenced_code_keeps_double_dollar(node_available):
|
||||||
|
# `$$` collapsed to a single `$` in the restored block.
|
||||||
|
html = _run_markdown_case('```sh\necho "$$USD and $$"\n```')
|
||||||
|
|
||||||
|
assert "$$USD and $$" in html
|
||||||
|
|
||||||
|
|
||||||
|
def test_mermaid_block_keeps_dollar_ampersand(node_available):
|
||||||
|
# The mermaid restore site had the same hazard: a node label containing `$&`
|
||||||
|
# re-inserted the ___MERMAID_BLOCK_n___ placeholder into the diagram source,
|
||||||
|
# which then fails to parse. The math and allowed-HTML sites are fixed the
|
||||||
|
# same way; they need KaTeX/sanitizer conditions this harness doesn't set up.
|
||||||
|
html = _run_markdown_case('```mermaid\ngraph TD; A["$&"] --> B;\n```')
|
||||||
|
|
||||||
|
assert "___MERMAID_BLOCK_" not in html
|
||||||
|
assert "$&" in html
|
||||||
|
|
||||||
|
|
||||||
def test_currency_dollar_amounts_are_not_rendered_as_math(node_available):
|
def test_currency_dollar_amounts_are_not_rendered_as_math(node_available):
|
||||||
# "$5 to $10" used to pair the two dollar signs as inline-math delimiters
|
# "$5 to $10" used to pair the two dollar signs as inline-math delimiters
|
||||||
# and render "5 to" through KaTeX. Pandoc-style rules now reject it: the
|
# and render "5 to" through KaTeX. Pandoc-style rules now reject it: the
|
||||||
|
|||||||
@@ -0,0 +1,15 @@
|
|||||||
|
"""Regression coverage for the built-in MCP servers' SDK compatibility line."""
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
REQUIREMENTS = Path(__file__).resolve().parents[1] / "requirements.txt"
|
||||||
|
|
||||||
|
|
||||||
|
def test_mcp_requirement_excludes_breaking_v2_sdk():
|
||||||
|
requirements = [
|
||||||
|
line.split("#", 1)[0].strip().replace(" ", "")
|
||||||
|
for line in REQUIREMENTS.read_text(encoding="utf-8").splitlines()
|
||||||
|
]
|
||||||
|
|
||||||
|
assert "mcp<2" in requirements
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
"""Regression test for the search route shim (slice 2j, #4082/#4071)."""
|
||||||
|
|
||||||
|
import importlib
|
||||||
|
|
||||||
|
import routes.search_routes as _shim_search # noqa: F401
|
||||||
|
|
||||||
|
|
||||||
|
def test_legacy_and_canonical_search_module_are_same_object():
|
||||||
|
legacy = importlib.import_module("routes.search_routes")
|
||||||
|
canonical = importlib.import_module("routes.search.search_routes")
|
||||||
|
assert legacy is canonical
|
||||||
@@ -0,0 +1,58 @@
|
|||||||
|
"""Regression for issue #5697 — skill timestamps must not use ``datetime.utcnow()``.
|
||||||
|
|
||||||
|
``_now_iso()`` builds the ``created`` value in skill frontmatter. ``utcnow()``
|
||||||
|
returns a *naive* datetime and has been deprecated since Python 3.12, scheduled
|
||||||
|
for removal. The replacement must stay timezone-aware while keeping the
|
||||||
|
serialized ``YYYY-MM-DDTHH:MM:SSZ`` shape, so skill files written by older
|
||||||
|
versions keep parsing.
|
||||||
|
|
||||||
|
The UTC check matters on its own: a bare ``datetime.now()`` also produces the
|
||||||
|
right shape, but emits local wall time, which would silently backdate or
|
||||||
|
postdate skills for every user outside UTC.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
import time
|
||||||
|
import warnings
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from services.memory.skill_format import _now_iso
|
||||||
|
|
||||||
|
_ISO_Z = re.compile(r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z$")
|
||||||
|
|
||||||
|
|
||||||
|
def test_now_iso_keeps_serialized_shape():
|
||||||
|
assert _ISO_Z.match(_now_iso())
|
||||||
|
|
||||||
|
|
||||||
|
def test_now_iso_emits_no_deprecation_warning():
|
||||||
|
with warnings.catch_warnings(record=True) as caught:
|
||||||
|
warnings.simplefilter("always")
|
||||||
|
_now_iso()
|
||||||
|
assert not [w for w in caught if issubclass(w.category, DeprecationWarning)]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.skipif(
|
||||||
|
not hasattr(time, "tzset"),
|
||||||
|
reason="time.tzset is unavailable on this platform",
|
||||||
|
)
|
||||||
|
def test_now_iso_is_utc_not_local_time():
|
||||||
|
"""Pin UTC under a non-UTC local timezone, where the two visibly diverge."""
|
||||||
|
original_tz = os.environ.get("TZ")
|
||||||
|
os.environ["TZ"] = "Asia/Amman" # UTC+3, never UTC
|
||||||
|
time.tzset()
|
||||||
|
try:
|
||||||
|
emitted = datetime.strptime(_now_iso(), "%Y-%m-%dT%H:%M:%SZ").replace(
|
||||||
|
tzinfo=timezone.utc
|
||||||
|
)
|
||||||
|
drift = abs((emitted - datetime.now(timezone.utc)).total_seconds())
|
||||||
|
assert drift < 60, f"timestamp is {drift}s off UTC — local time leaked in"
|
||||||
|
finally:
|
||||||
|
if original_tz is None:
|
||||||
|
os.environ.pop("TZ", None)
|
||||||
|
else:
|
||||||
|
os.environ["TZ"] = original_tz
|
||||||
|
time.tzset()
|
||||||
@@ -0,0 +1,97 @@
|
|||||||
|
import os
|
||||||
|
import time
|
||||||
|
from pathlib import Path
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
# Adjust the import path if your file is directly in ./services instead of ./services/tts
|
||||||
|
from services.tts.tts_service import TTSService
|
||||||
|
|
||||||
|
def test_cache_under_limit(tmp_path, monkeypatch):
|
||||||
|
"""Test that writing a file under the size limit does not trigger eviction."""
|
||||||
|
# Set a tiny limit: 100 bytes
|
||||||
|
monkeypatch.setenv("ODYSSEUS_TTS_CACHE_MAX_BYTES", "100")
|
||||||
|
|
||||||
|
# Initialize service with pytest's temporary directory
|
||||||
|
service = TTSService(cache_dir=str(tmp_path))
|
||||||
|
|
||||||
|
# Write a 40-byte file (under the 100-byte limit)
|
||||||
|
service._put_cache("test_key", b"x" * 40)
|
||||||
|
|
||||||
|
# Verify the file was written and nothing was deleted
|
||||||
|
files = list(tmp_path.glob("*.*"))
|
||||||
|
assert len(files) == 1
|
||||||
|
assert sum(f.stat().st_size for f in files) == 40
|
||||||
|
|
||||||
|
def test_cache_exceeds_limit_triggers_eviction(tmp_path, monkeypatch):
|
||||||
|
"""Test that exceeding the limit evicts the oldest files down to 80% capacity."""
|
||||||
|
# Set limit to 100 bytes. 80% target capacity will be 80 bytes.
|
||||||
|
monkeypatch.setenv("ODYSSEUS_TTS_CACHE_MAX_BYTES", "100")
|
||||||
|
service = TTSService(cache_dir=str(tmp_path))
|
||||||
|
|
||||||
|
# 1. Setup: Manually create two older files (40 bytes each)
|
||||||
|
file1 = tmp_path / "oldest.wav"
|
||||||
|
file2 = tmp_path / "middle.wav"
|
||||||
|
|
||||||
|
file1.write_bytes(b"a" * 40)
|
||||||
|
file2.write_bytes(b"b" * 40)
|
||||||
|
|
||||||
|
# Spoof timestamps so file1 is explicitly older than file2
|
||||||
|
now = time.time()
|
||||||
|
os.utime(file1, (now - 100, now - 100)) # 100 seconds ago
|
||||||
|
os.utime(file2, (now - 50, now - 50)) # 50 seconds ago
|
||||||
|
|
||||||
|
# 2. Action: Write a 3rd file using the service method (40 bytes)
|
||||||
|
# Total cache is now 120 bytes, which exceeds 100.
|
||||||
|
# It should delete oldest (file1) to drop to 80 bytes (which matches the 80% target).
|
||||||
|
service._put_cache("newest", b"c" * 40)
|
||||||
|
|
||||||
|
# 3. Assertions
|
||||||
|
# The newest file should exist (saved as .wav because it lacks MP3 magic bytes)
|
||||||
|
newest_file = tmp_path / "newest.wav"
|
||||||
|
|
||||||
|
assert not file1.exists(), "The oldest file should have been evicted."
|
||||||
|
assert file2.exists(), "The middle file should still exist."
|
||||||
|
assert newest_file.exists(), "The newest file should have been saved."
|
||||||
|
|
||||||
|
# Verify the final directory size is <= 80 bytes
|
||||||
|
total_size = sum(f.stat().st_size for f in tmp_path.glob("*.*"))
|
||||||
|
assert total_size <= 80
|
||||||
|
|
||||||
|
def test_cache_limit_disabled(tmp_path, monkeypatch):
|
||||||
|
"""Test that setting max bytes to 0 disables eviction."""
|
||||||
|
monkeypatch.setenv("ODYSSEUS_TTS_CACHE_MAX_BYTES", "0")
|
||||||
|
service = TTSService(cache_dir=str(tmp_path))
|
||||||
|
|
||||||
|
# Write 3 large files that would normally trigger eviction
|
||||||
|
service._put_cache("file1", b"x" * 1000)
|
||||||
|
service._put_cache("file2", b"x" * 1000)
|
||||||
|
service._put_cache("file3", b"x" * 1000)
|
||||||
|
|
||||||
|
# Ensure nothing was deleted
|
||||||
|
files = list(tmp_path.glob("*.*"))
|
||||||
|
assert len(files) == 3
|
||||||
|
assert sum(f.stat().st_size for f in files) == 3000
|
||||||
|
|
||||||
|
def test_cache_eviction_handles_unlink_error_gracefully(tmp_path, monkeypatch):
|
||||||
|
"""Test that if unlinking a file fails, _put_cache still succeeds without raising."""
|
||||||
|
service = TTSService(cache_dir=str(tmp_path))
|
||||||
|
service.max_cache_bytes = 50
|
||||||
|
|
||||||
|
# Create a file to evict
|
||||||
|
old_file = tmp_path / "old.wav"
|
||||||
|
old_file.write_bytes(b"x" * 40)
|
||||||
|
|
||||||
|
# Monkeypatch unlink on Path objects to simulate a PermissionError / file-lock failure
|
||||||
|
def mock_unlink(self_path):
|
||||||
|
raise OSError("Permission denied / file locked")
|
||||||
|
|
||||||
|
monkeypatch.setattr(Path, "unlink", mock_unlink)
|
||||||
|
|
||||||
|
# Writing a new file triggers eviction which encounters the mocked unlink error
|
||||||
|
try:
|
||||||
|
service._put_cache("new_key", b"y" * 40)
|
||||||
|
except Exception as e:
|
||||||
|
pytest.fail(f"_put_cache raised an exception during failed eviction: {e}")
|
||||||
|
|
||||||
|
# The new file should still be written successfully
|
||||||
|
assert (tmp_path / "new_key.wav").exists()
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
"""Regression test for the vault route shim (slice 2k, #4082/#4071)."""
|
||||||
|
|
||||||
|
import importlib
|
||||||
|
|
||||||
|
import routes.vault_routes as _shim_vault # noqa: F401
|
||||||
|
|
||||||
|
|
||||||
|
def test_legacy_and_canonical_vault_module_are_same_object():
|
||||||
|
legacy = importlib.import_module("routes.vault_routes")
|
||||||
|
canonical = importlib.import_module("routes.vault.vault_routes")
|
||||||
|
assert legacy is canonical
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
"""Regression test for the webhook route shim (slice 2l, #4082/#4071)."""
|
||||||
|
|
||||||
|
import importlib
|
||||||
|
|
||||||
|
import routes.webhook_routes as _shim_webhook # noqa: F401
|
||||||
|
|
||||||
|
|
||||||
|
def test_legacy_and_canonical_webhook_module_are_same_object():
|
||||||
|
legacy = importlib.import_module("routes.webhook_routes")
|
||||||
|
canonical = importlib.import_module("routes.webhook.webhook_routes")
|
||||||
|
assert legacy is canonical
|
||||||
Reference in New Issue
Block a user