Consolidate Odysseus agent harness and tool contracts

This commit is contained in:
pewdiepie-archdaemon
2026-09-17 10:07:40 +00:00
parent 84aa9a91de
commit 218d762427
229 changed files with 28899 additions and 1551 deletions
+58 -5
View File
@@ -124,7 +124,9 @@ class CompletionDecision:
_ARTIFACT_PATH = r"(?:/|\./|\.\./)?[A-Za-z0-9_.-]+(?:/[A-Za-z0-9_.-]+)*\.[A-Za-z0-9]{1,12}"
_ARTIFACT_REQUEST_RE = re.compile(
rf"\b(?:write|create|make|save|produce|generate|export|edit|modify|update|fix|put|place)\b"
rf"\b(?:writ(?:e|ten)|creat(?:e|ed)|make|made|sav(?:e|ed)|produc(?:e|ed)|"
rf"generat(?:e|ed)|export(?:ed)?|edit(?:ed)?|modif(?:y|ied)|updat(?:e|ed)|"
rf"fix(?:ed)?|put|plac(?:e|ed))\b"
rf"[^\n]{{0,80}}?(?P<path>{_ARTIFACT_PATH})",
re.IGNORECASE,
)
@@ -133,7 +135,7 @@ _OUTPUT_PATH_RE = re.compile(
re.IGNORECASE,
)
_EXPLICIT_OUTPUT_FILE_RE = re.compile(
rf"\b(?:to|at|as)\s+(?:the\s+)?(?:file|path)\s+(?P<path>{_ARTIFACT_PATH})",
rf"\b(?:to|at|as|into)\s+(?:the\s+|a\s+)?(?:single\s+)?(?:file|path)\s+(?P<path>{_ARTIFACT_PATH})",
re.IGNORECASE,
)
_NAMED_OUTPUT_FILE_RE = re.compile(
@@ -141,8 +143,9 @@ _NAMED_OUTPUT_FILE_RE = re.compile(
re.IGNORECASE,
)
_EXPLICIT_OUTPUT_DIRECTORY_RE = re.compile(
r"\b(?:save|write|create|make|produce|generate|export|put|place)\b"
r"[^\n]{0,100}?\b(?:into|to|under|inside)\s+"
r"\b(?:sav(?:e|ed)|writ(?:e|ten)|creat(?:e|ed)|make|made|produc(?:e|ed)|"
r"generat(?:e|ed)|export(?:ed)?|put|plac(?:e|ed))\b"
r"[^\n]{0,100}?\b(?:in|into|to|under|inside)\s+"
r"[`'\"]?(?P<path>/(?:[A-Za-z0-9_.-]+/)*[A-Za-z0-9_.-]+/?)"
r"(?=[`'\"\s.,;:]|$)",
re.IGNORECASE,
@@ -208,6 +211,34 @@ def _is_prose_abbreviation(value: str) -> bool:
return _clean_path(value).lower() in {"e.g", "i.e"}
def _artifact_match_is_negated(instruction: str, match: re.Match[str]) -> bool:
"""Reject paths attached to an explicitly negated mutation verb."""
prefix = instruction[max(0, match.start() - 32):match.start()]
return bool(re.search(r"(?:do\s+not|don't|must\s+not|never)\s+$", prefix, re.IGNORECASE))
def _artifact_match_is_callable(instruction: str, match: re.Match[str], path: str) -> bool:
"""Reject dotted callable names such as ``json.dumps(...)`` as artifacts."""
if "/" in path or "\\" in path:
return False
if instruction[match.end("path"):].startswith("("):
return True
# Procedural prompts often name existence helpers without parentheses,
# e.g. "verify with os.path.exists or ls". They are code references, not
# output filenames, even though the generic path regex sees an extension.
return bool(re.fullmatch(r"(?:os\.path|pathlib\.Path|Path)\.[A-Za-z_]\w*", path))
def _artifact_match_is_email_host(instruction: str, match: re.Match[str]) -> bool:
"""Reject the domain portion of an email address as an output path."""
start = match.start("path")
prefix = instruction[max(0, start - 80):start]
return bool(re.search(r"[A-Za-z0-9_.+-]+@$", prefix))
def infer_completion_requirements(
instruction: str,
*,
@@ -216,6 +247,7 @@ def infer_completion_requirements(
) -> CompletionRequirements:
"""Infer only explicitly requested output/edit paths from an instruction."""
text = str(instruction or "")
paths: list[str] = []
for pattern in (
_ARTIFACT_REQUEST_RE,
@@ -226,8 +258,14 @@ def infer_completion_requirements(
_EXPLICIT_OUTPUT_DIRECTORY_RE,
_LOCALIZED_OUTPUT_DIRECTORY_RE,
):
for match in pattern.finditer(str(instruction or "")):
for match in pattern.finditer(text):
path = _clean_path(match.group("path"))
if _artifact_match_is_negated(text, match):
continue
if _artifact_match_is_callable(text, match, path):
continue
if _artifact_match_is_email_host(text, match):
continue
if path and not _is_prose_abbreviation(path) and path not in paths:
paths.append(path)
paths = [path.rstrip("/") if path != "/" else path for path in paths]
@@ -249,6 +287,13 @@ def infer_completion_requirements(
if path in explicit_directories
or any(path.startswith(directory.rstrip("/") + "/") for directory in explicit_directories)
]
explicit_files = [path for path in paths if Path(path).suffix]
if explicit_files:
paths = [
path for path in paths
if path not in explicit_directories
or not any(file.startswith(path.rstrip("/") + "/") for file in explicit_files)
]
cleaned_verifier_commands = tuple(dict.fromkeys(
str(command or "").strip()
for command in verifier_commands
@@ -335,6 +380,14 @@ def _artifact_path_matches_required(artifact_path: str, required_path: str) -> b
def _explicit_tool_paths(tool: str, command: str) -> list[str]:
if tool == "write_file":
try:
args = json.loads(command or "{}")
except (TypeError, json.JSONDecodeError):
args = None
if isinstance(args, Mapping):
path = _clean_path(str(args.get("path") or ""))
return [path] if path else []
# Keep compatibility with the legacy ``path\ncontent`` transport.
path = _clean_path(str(command or "").splitlines()[0] if command else "")
return [path] if path else []
if tool == "edit_file":
+1251 -124
View File
File diff suppressed because it is too large Load Diff
+2 -1
View File
@@ -560,7 +560,8 @@ async def do_manage_settings(content: str, owner: Optional[str] = None) -> Dict:
"hard max": "agent_input_token_hard_max",
"token budget cap": "agent_input_token_hard_max",
"input budget cap": "agent_input_token_hard_max",
"writing style": "email_writing_style", "email writing style": "email_writing_style",
"writing style": "document_writing_style", "document writing style": "document_writing_style",
"email writing style": "email_writing_style",
"reply writing style": "email_writing_style", "email reply writing style": "email_writing_style",
}
def _resolve(k):
+56 -2
View File
@@ -11,6 +11,38 @@ from src.upload_handler import reserve_upload_references
logger = logging.getLogger(__name__)
_DOCUMENT_SEARCH_STOPWORDS = frozenset({
'a', 'an', 'and', 'any', 'about', 'document', 'documents', 'for', 'in',
'my', 'of', 'on', 'or', 'plans', 'the', 'to',
})
def _document_search_tokens(value: str) -> list[str]:
return [
token for token in re.findall(r'[a-z0-9]+', str(value or '').lower())
if token not in _DOCUMENT_SEARCH_STOPWORDS
]
def _rank_document_search(docs, search_text: str):
"""Prefer phrase/all-term matches, then broaden to any meaningful term."""
query = str(search_text or '').strip().lower()
terms = _document_search_tokens(query)
scored = []
for position, doc in enumerate(docs):
haystack = ' '.join((
str(getattr(doc, 'title', '') or ''),
str(getattr(doc, 'current_content', '') or ''),
)).lower()
haystack_terms = set(_document_search_tokens(haystack))
matched = sum(term in haystack_terms for term in terms)
strict = bool(query and query in haystack) or bool(terms and matched == len(terms))
scored.append((doc, strict, matched, position))
strict_matches = [row for row in scored if row[1]]
candidates = strict_matches or [row for row in scored if row[2] > 0]
return [row[0] for row in sorted(candidates, key=lambda row: (-row[2], row[3]))]
def _missing_document_upload(owner: Optional[str], content: Any) -> Optional[str]:
"""Reserve explicit upload URLs before an agent persists document text."""
return reserve_upload_references(get_upload_handler(), owner, content)
@@ -629,6 +661,12 @@ class UpdateDocumentTool:
if is_email_doc:
doc.language = "email"
if new_content == (doc.current_content or ""):
return {
"error": "No update applied — replacement content is unchanged",
"exit_code": 1,
}
missing_id = _missing_document_upload(owner, new_content)
if missing_id:
return {
@@ -761,6 +799,10 @@ class EditDocumentTool:
skipped = 0
for edit in edits:
_find = edit["find"]
if _find == edit["replace"]:
logger.warning("edit_document: skipping no-op FIND/REPLACE block")
skipped += 1
continue
if _find in updated_content:
updated_content = updated_content.replace(_find, edit["replace"], 1)
applied += 1
@@ -941,10 +983,22 @@ class ManageDocumentTool:
search_text = re.sub(
r"\s+(?:instead|please)\s*$", "", search_text, flags=re.IGNORECASE
).strip()
q = q.filter(Document.title.ilike(f"%{search_text}%"))
if args.get("language"):
q = q.filter(Document.language == args["language"])
docs = q.order_by(Document.updated_at.desc()).limit(args.get("limit", 50)).all()
requested_limit = args.get("limit", 50)
try:
requested_limit = max(1, min(int(requested_limit), 200))
except (TypeError, ValueError):
requested_limit = 50
q = q.order_by(Document.updated_at.desc())
# A plain listing must not load the entire document library
# (including every document body) before applying its limit.
if not search_text:
q = q.limit(requested_limit)
docs = q.all()
if search_text:
docs = _rank_document_search(docs, search_text)
docs = docs[:requested_limit]
if not docs:
msg = "No documents found" + (f" matching '{search_text}'" if search_text else "") + "."
return {"response": msg, "documents": [], "exit_code": 0}
+18
View File
@@ -20,6 +20,10 @@ _CODENAV_MAX_LINE = 400
_STRUCTURED_DOCUMENT_SUFFIXES = frozenset({
".doc", ".docx", ".epub", ".pdf", ".pptx", ".xls", ".xlsx",
})
_BINARY_ARTIFACT_SUFFIXES = _STRUCTURED_DOCUMENT_SUFFIXES | frozenset({
".bmp", ".gif", ".ico", ".jpeg", ".jpg", ".mp3", ".mp4", ".ogg",
".png", ".wav", ".webm", ".webp", ".zip",
})
def _glob_to_regex(pat: str) -> "re.Pattern":
@@ -275,6 +279,20 @@ class WriteFileTool:
),
"exit_code": 1,
}
# write_file is a UTF-8 text writer. Refuse to silently destroy an
# existing PDF, image, archive, or media artifact produced by a
# format-aware tool, especially after the agent has verified it.
suffix = os.path.splitext(path)[1].casefold()
if suffix in _BINARY_ARTIFACT_SUFFIXES:
target_existed = os.path.isfile(path)
return {
"error": (
f"write_file: refusing UTF-8 text for binary artifact path {path}. "
"Use Python or a format-specific creation tool, then inspect the result."
),
"exit_code": 1,
"binary_artifact_preserved": target_existed,
}
try:
def _write():
old = ""
+72 -5
View File
@@ -433,6 +433,12 @@ class ExtractTextTool:
return {"error": "extract_text unknown argument(s): " + ", ".join(unknown), "exit_code": 1}
try:
raw_path = str(args.get("path") or '')
# Some native-schema models serialize a workspace path using the
# same URI shape as uploads. This alias grants no extra access:
# convert it back to /workspace and let the normal confinement
# resolver enforce the active root.
if raw_path.startswith('odysseus://workspace/'):
raw_path = '/workspace/' + raw_path[len('odysseus://workspace/'):]
if raw_path.startswith('odysseus://'):
# Upload access is independent of a filesystem workspace and
# must never inherit an administrator's cross-owner override.
@@ -450,8 +456,9 @@ class ExtractTextTool:
path = _resolve_media_path(raw_path, tool_name="extract_text")
except ValueError as exc:
return {"error": str(exc), "exit_code": 1}
if path.suffix.casefold() not in _IMAGE_SUFFIXES:
return {"error": "extract_text currently supports local image files", "exit_code": 1}
suffix = path.suffix.casefold()
if suffix not in _IMAGE_SUFFIXES | _PDF_SUFFIXES:
return {"error": "extract_text supports local image and PDF files", "exit_code": 1}
mode = str(args.get("mode") or "all").strip().casefold()
try:
minimum, maximum = float(args.get("min_confidence", .5)), int(args.get("max_results", 512))
@@ -461,7 +468,56 @@ class ExtractTextTool:
return {"error": "invalid extract_text mode or bounds", "exit_code": 1}
try:
from .ocr_engine import extract_image_text
evidence = await asyncio.to_thread(extract_image_text, path, include_layout=bool(args.get("include_layout", False)), numeric_only=mode == "numbers", min_confidence=minimum, max_results=maximum)
if suffix in _PDF_SUFFIXES:
def _extract_pdf_pages():
try:
import pypdfium2 as pdfium
except ImportError as exc:
raise RuntimeError(
"PDF OCR requires the optional pypdfium2 package"
) from exc
document = pdfium.PdfDocument(str(path))
page_count = len(document)
lines, accepted = [], 0
# Keep one OCR call bounded while covering ordinary
# documents completely. Larger PDFs can be inspected in
# page ranges with inspect_media.
rendered_count = min(page_count, 12)
with tempfile.TemporaryDirectory(prefix="odysseus-pdf-ocr-") as temp_dir:
for index in range(rendered_count):
rendered = document[index].render(scale=2.0).to_pil().convert("RGB")
image_path = Path(temp_dir) / f"page-{index + 1}.png"
rendered.save(image_path, "PNG")
remaining = max(1, maximum - len(lines))
page_evidence = extract_image_text(
image_path,
include_layout=bool(args.get("include_layout", False)),
numeric_only=mode == "numbers",
min_confidence=minimum,
max_results=remaining,
)
accepted += int(page_evidence.get("count") or 0)
for line in page_evidence.get("lines") or []:
if len(lines) >= maximum:
break
lines.append({"page": index + 1, **line})
return {
"legend": {
"page": "one-based PDF page",
"t": "text",
"p": "confidence",
"xy": "pixel center",
},
"page_count": page_count,
"pages_processed": rendered_count,
"count": accepted,
"returned": len(lines),
"truncated": accepted > len(lines) or page_count > rendered_count,
"lines": lines,
}
evidence = await asyncio.to_thread(_extract_pdf_pages)
else:
evidence = await asyncio.to_thread(extract_image_text, path, include_layout=bool(args.get("include_layout", False)), numeric_only=mode == "numbers", min_confidence=minimum, max_results=maximum)
except Exception as exc:
return {"error": f"extract_text failed: {exc}", "exit_code": 1}
return {"output": json.dumps(evidence, ensure_ascii=False, separators=(",", ":")), "exit_code": 0, "ocr": evidence}
@@ -715,9 +771,16 @@ class InspectMediaTool:
suffix = path.suffix.lower()
if suffix in _SVG_SUFFIXES:
renderer = shutil.which("rsvg-convert")
renderer_kind = "rsvg"
if not renderer:
renderer = shutil.which("convert")
renderer_kind = "imagemagick"
if not renderer:
return {
"error": "inspect_media SVG rendering requires rsvg-convert",
"error": (
"inspect_media SVG rendering requires rsvg-convert "
"or ImageMagick convert"
),
"exit_code": 1,
}
raw_output = str(args.get("output_path") or "").strip()
@@ -740,7 +803,11 @@ class InspectMediaTool:
output = Path(temporary.name)
rendered = await asyncio.to_thread(
_run,
[renderer, "--output", str(output), str(path)],
(
[renderer, "--output", str(output), str(path)]
if renderer_kind == "rsvg"
else [renderer, str(path), str(output)]
),
60,
)
if rendered.returncode != 0 or not output.is_file() or output.stat().st_size == 0:
@@ -133,6 +133,62 @@ async def list_models(content: str, session_id: Optional[str] = None, owner: Opt
keyword = content.strip().lower() if content.strip() else None
# ``list_models`` historically treated every filter as a literal model-ID
# substring. For recommendation terms that produced an empty catalog even
# though Odysseus already has a hardware detector and fit ranker. Preserve
# the catalog behavior for real model/provider filters, but give these
# semantic filters their expected read-only meaning.
if keyword in {
"recommended", "recommendation", "recommendations",
"compatible", "hardware", "hardware fit", "best fit",
}:
from src.tools.system import do_app_api
fit_result = await do_app_api(json.dumps({
"action": "call",
"method": "GET",
"path": "/api/hwfit/models",
"query": {"fit_only": "true", "limit": 5, "sort": "fit"},
}), owner=owner)
payload = fit_result.get("json") if isinstance(fit_result, dict) else None
system = payload.get("system") if isinstance(payload, dict) else None
models = payload.get("models") if isinstance(payload, dict) else None
if isinstance(system, dict) and isinstance(models, list):
gpu = system.get("gpu_name") or "No GPU detected"
vram = system.get("gpu_vram_gb")
count = system.get("gpu_count")
backend = system.get("backend") or "unknown"
lines = [
"Detected hardware:",
f"- GPU: {gpu}; count={count}; total VRAM={vram} GB; backend={backend}",
f"- CPU: {system.get('cpu_name') or 'unknown'}; RAM={system.get('total_ram_gb')} GB",
"Ranked compatible models:",
]
compact_models = []
for model_row in models[:5]:
if not isinstance(model_row, dict):
continue
compact = {
key: model_row.get(key)
for key in (
"name", "parameter_count", "quant", "required_gb",
"fit_level", "run_mode", "speed_tps", "score", "context",
)
}
compact_models.append(compact)
lines.append(
"- {name}: params={parameter_count}, quant={quant}, required={required_gb} GB, "
"fit={fit_level}, mode={run_mode}, speed={speed_tps} tok/s, score={score}, context={context}".format(
**compact
)
)
return {
"output": "\n".join(lines),
"system": system,
"models": compact_models,
"exit_code": 0,
}
return fit_result
db = SessionLocal()
try:
query = db.query(ModelEndpoint).filter(ModelEndpoint.is_enabled == True)
+17 -3
View File
@@ -874,6 +874,15 @@ def _python_with_visible_final_expression(content: str) -> str:
return ast.unparse(tree)
def _python_with_configured_import_paths(content: str, env: dict | None) -> str:
"""Expose only explicitly configured package roots under Python ``-I``."""
raw = str((env or {}).get("ODYSSEUS_PYTHON_TOOL_SITE_PACKAGES", ""))
paths = [item for item in raw.split(os.pathsep) if item and os.path.isabs(item)]
if not paths:
return content
return f"import site\n[site.addsitedir(path) for path in {paths!r}]\nexec(compile({content!r}, '<odysseus-python-tool>', 'exec'))"
class PythonTool:
async def execute(self, content: str, ctx: dict) -> dict:
from src.tool_execution import agent_cwd, _truncate
@@ -917,7 +926,9 @@ class PythonTool:
# process-global `/workspace` symlink would break concurrent tasks.
# Give Python the same per-task namespace Bash receives so both inline
# code and loaded scripts see the stable virtual workspace root.
namespaced_content = _python_with_visible_final_expression(content)
namespaced_content = _python_with_configured_import_paths(
_python_with_visible_final_expression(content), _subproc_env
)
python_command = shlex.join((sys.executable or "python", "-I", "-c", namespaced_content))
# Code that explicitly uses the public /workspace path runs inside a
# namespace whose stable cwd is that same bind. Host workspaces under
@@ -944,8 +955,11 @@ class PythonTool:
else:
# Platforms without a usable namespace still receive the same
# alias contract through a conservative source rewrite.
content = _python_with_visible_final_expression(
_replace_workspace_alias(content, agent_cwd())
content = _python_with_configured_import_paths(
_python_with_visible_final_expression(
_replace_workspace_alias(content, agent_cwd())
),
_subproc_env,
)
proc = await asyncio.create_subprocess_exec(
(sys.executable or "python"), "-I", "-c", content,
+49
View File
@@ -270,6 +270,32 @@ class WebSearchTool:
timeout=30,
)
except asyncio.TimeoutError:
# Comprehensive search also downloads several result pages. A
# slow or hostile publisher must not erase the ranked search
# evidence that was already available. Fall back to the metadata
# path so the agent can choose a source and continue with
# web_fetch/private_browser. Keep this bounded independently: the
# abandoned executor thread may still be winding down.
try:
results = await asyncio.wait_for(
loop.run_in_executor(
None,
lambda: searxng_search_results(query, max_pages),
),
timeout=12,
)
text, sources = _format_search_metadata(query, results)
if sources:
output = text[:MAX_OUTPUT_CHARS] if len(text) > MAX_OUTPUT_CHARS else text
output += "\n\n<!-- SOURCES:" + json.dumps(sources) + " -->"
return {
"output": output,
"exit_code": 0,
"evidence_status": "available",
"degraded_mode": "metadata_after_content_timeout",
}
except Exception:
pass
return {
"error": f"web_search timed out after 30s: {query[:200]}",
"exit_code": 1,
@@ -2190,8 +2216,13 @@ class PrivateBrowserTool:
"batch",
}
_AUTO_SCREENSHOT_ACTIONS = {
"open",
"snapshot",
"batch",
"click",
"fill",
"press",
"scroll",
}
@staticmethod
@@ -2972,6 +3003,15 @@ class PrivateBrowserTool:
for command in commands:
if isinstance(command, list) and command:
action = str(command[0]).strip().lower()
if action == "wait":
# Compact/OpenAI schemas sometimes preserve an omitted
# selector as null and put the timeout in the next slot:
# ["wait", null, 2500]. agent-browser accepts only arrays
# of strings, so recover the intended timeout instead of
# rejecting the whole browser batch.
wait_args = [value for value in command[1:] if value is not None]
normalized.append(["wait", *[str(value) for value in wait_args]])
continue
if action in {"open", "read"} and len(command) >= 2:
candidate_url = str(command[1] or "").strip()
if (
@@ -2985,6 +3025,15 @@ class PrivateBrowserTool:
*command[2:],
])
continue
if action == "read" and not re.match(
r"^(?:https?|file)://", candidate_url, re.IGNORECASE
):
# The top-level read action treats target/selector as
# DOM text extraction. Keep batch semantics identical;
# agent-browser's bare `read h1` instead interprets h1
# as a URL/path and fails before the model can answer.
normalized.append(["get", "text", candidate_url])
continue
if action == "evaluate":
normalized.append(["eval", *command[1:]])
continue
+22
View File
@@ -368,6 +368,28 @@ def decode_native_trace(
call_id = selected["call_id"]
if round_no is None:
round_no = selected["round"]
elif event.get("execution_attempted") is False:
# Preview guards return a protocol-level tool result for a
# model-proposed call that was rejected before dispatch (for
# example, an exact duplicate). It is still a real attempted
# model action and must have a correlated call in the trace;
# treating it as an orphan falsely invalidates otherwise
# complete runs. The explicit marker keeps genuinely
# unpaired legacy outputs fail-closed below.
call_id = explicit_call_id or f"native-rejected-{len(builder.events)}"
builder.add(
TraceKind.TOOL_CALL,
{
"tool_name": tool,
"arguments": command,
"command": command,
"execution_attempted": False,
"rejected_before_execution": True,
},
timestamp_s=timestamp,
round=round_no,
correlation_id=call_id,
)
else:
call_id = explicit_call_id or f"native-orphan-{len(builder.events)}"
builder.gap("tool_result_call_unmatched", f"{tool}:{call_id}")
+70 -5
View File
@@ -697,7 +697,8 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
switch_model <model> — Change the model for the current session
set_theme <preset> — Apply a built-in theme preset (dark, light, midnight, paper, cyberpunk, retrowave, forest, ocean, ume, copper, terminal, organs, lavender, gpt, claude, cute)
create_theme <name> <bg> <fg> <panel> <border> <accent> [key=val ...] — Create custom theme. Optional key=val: advanced color overrides AND background effects: bgPattern=<none|dots|synapse|rain|constellations|perlin-flow|petals|sparkles|embers>, bgEffectColor=#RRGGBB, bgEffectIntensity=<num>, bgEffectSize=<num>, frosted=true|false
open_panel <name> — Open a panel (documents, gallery, calendar, email, sessions, notes, memories, skills, settings, theme, cookbook)
get_theme — Return the last server-synchronized theme for this user
open_panel <name> [view] — Open a panel; Cookbook views are download/models, launch/serve, active/running, dependencies, settings
open_email_reply <uid> [folder] [reply|reply-all|ai-reply] [body text] — Open a reply draft document for an email; does not send. ALWAYS append the body text when the user told you what to say (one-shot draft); only omit body when the user just asked to "open a reply" without content.
get_toggles — Return current toggle states (server-side knowledge)
"""
@@ -803,14 +804,27 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
]
custom_themes = {}
try:
from routes.prefs_routes import _load as _load_prefs
custom_themes = _load_prefs().get("custom-themes", {}) or {}
from routes.prefs_routes import _load_for_user
custom_themes = _load_for_user(owner).get("custom-themes", {}) or {}
except Exception:
pass
all_known = set(known_presets) | set(custom_themes.keys())
if theme_name not in all_known:
custom_label = f" | Custom: {', '.join(sorted(custom_themes.keys()))}" if custom_themes else ""
return {"error": f"Unknown theme '{theme_name}'. Available: {', '.join(sorted(known_presets))}{custom_label}"}
try:
from routes.prefs_routes import _load_for_user, _save_for_user
prefs = _load_for_user(owner)
previous = prefs.get("theme") if isinstance(prefs.get("theme"), dict) else {}
stored = {"name": theme_name}
if previous.get("name") == theme_name and isinstance(previous.get("colors"), dict):
stored["colors"] = previous["colors"]
elif isinstance(custom_themes.get(theme_name), dict):
stored["colors"] = custom_themes[theme_name]
prefs["theme"] = stored
_save_for_user(owner, prefs)
except Exception:
pass
return {
"ui_event": "set_theme",
"theme_name": theme_name,
@@ -868,6 +882,17 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
bg["frosted"] = av.lower() in ("true", "1", "yes", "on")
if advanced:
colors["advanced"] = advanced
try:
from routes.prefs_routes import _load_for_user, _save_for_user
prefs = _load_for_user(owner)
custom_themes = prefs.get("custom-themes")
custom_themes = dict(custom_themes) if isinstance(custom_themes, dict) else {}
custom_themes[name] = dict(colors)
prefs["custom-themes"] = custom_themes
prefs["theme"] = {"name": name, "colors": dict(colors)}
_save_for_user(owner, prefs)
except Exception:
pass
return {
"ui_event": "create_theme",
"theme_name": name,
@@ -901,6 +926,7 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
# calendar, email, sessions, notes, memories, skills, settings, theme, cookbook.
panel = parts[1].lower() if len(parts) > 1 else ""
view = ""
view_label = ""
target_date = ""
_panel_aliases = {
"documents": "documents",
@@ -943,6 +969,23 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
target = _panel_aliases.get(panel)
if not target:
return {"error": f"Unknown panel '{panel}'. Valid: documents, gallery, calendar, email, sessions, notes, memories, skills, settings, theme, cookbook."}
if target == "cookbook":
cookbook_views = {
"models": ("Search", "models"), "model": ("Search", "models"),
"download": ("Search", "models"), "search": ("Search", "models"),
"serve": ("Serve", "launch"), "serving": ("Serve", "launch"),
"launch": ("Serve", "launch"),
"active": ("Running", "running"), "running": ("Running", "running"),
"dependencies": ("Dependencies", "dependencies"),
"dependency": ("Dependencies", "dependencies"),
"settings": ("Settings", "settings"),
}
requested_view = parts[2].strip().lower() if len(parts) > 2 else ""
# A panel alias can carry the subview intent by itself. Previously
# `models` and `serve` were silently collapsed to bare Cookbook.
resolved_view = cookbook_views.get(requested_view) or cookbook_views.get(panel)
if resolved_view:
view, view_label = resolved_view
if target == "calendar":
view_words = {"day", "week", "month", "year", "agenda"}
tail_text = ""
@@ -964,9 +1007,13 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
"panel": target,
"results": f"Opening {target} panel",
}
if panel != target:
payload["requested_panel"] = panel
if view:
payload["view"] = view
payload["results"] = f"Opening {target} panel in {view} view"
if view_label:
payload["view_label"] = view_label
payload["results"] = f"Opening {target} panel in {view_label or view} view"
if target_date:
payload["target_date"] = target_date
return payload
@@ -1021,6 +1068,24 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
result["body"] = body
return result
elif action == "get_theme":
try:
from routes.prefs_routes import _load_for_user
saved = _load_for_user(owner).get("theme")
except Exception:
saved = None
name = str(saved.get("name") or "").strip() if isinstance(saved, dict) else ""
if not name:
return {
"results": "The current client theme has not been synchronized to the server.",
"theme_known": False,
}
return {
"results": f"Current theme: {name}",
"current_theme": name,
"theme_known": True,
}
elif action == "get_toggles":
return {
"results": (
@@ -1031,7 +1096,7 @@ async def do_ui_control(content: str, session_id: Optional[str] = None, owner: O
}
else:
return {"error": f"Unknown action '{action}'. Use: toggle, set_mode, switch_model, set_theme, highlight, clear_highlight, get_toggles"}
return {"error": f"Unknown action '{action}'. Use: toggle, set_mode, switch_model, set_theme, create_theme, get_theme, highlight, clear_highlight, get_toggles"}
# ---------------------------------------------------------------------------
+83 -19
View File
@@ -690,7 +690,11 @@ async def action_consolidate_memory(owner: str, **kwargs) -> Tuple[str, bool]:
return False
from src.task_endpoint import resolve_task_candidates
candidates = resolve_task_candidates(owner=group_owner or None)
candidates = resolve_task_candidates(
owner=group_owner or None,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
if not candidates:
return False
@@ -1143,6 +1147,8 @@ async def action_summarize_emails(owner: str, **kwargs) -> Tuple[str, bool]:
do_summary=True,
do_reply=False,
account_id=_email_task_account_id(kwargs),
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
if _result_is_config_error(result):
return result, False
@@ -1164,6 +1170,8 @@ async def action_draft_email_replies(owner: str, **kwargs) -> Tuple[str, bool]:
account_id=_email_task_account_id(kwargs),
days_back=7,
progress_cb=kwargs.get("progress_cb"),
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
if _result_is_config_error(result):
return result, False
@@ -1297,20 +1305,37 @@ async def action_email_auto_translate(owner: str, **kwargs) -> Tuple[str, bool]:
},
],
owner=owner,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
temperature=0.2,
max_tokens=8192,
timeout=180,
)
content = (content or "").strip()
content = _extract_reply(content)
if "<<<SAME_LANGUAGE>>>" in content:
return "", True
marker = _re.search(r"<<<TRANSLATION>>>\s*(.*?)\s*<<<END>>>", content, _re.S | _re.I)
if marker:
content = marker.group(1).strip()
# Translation markers are distinct from the reply/summary markers
# handled by _extract_reply. Some reasoning-capable models repeat
# the opening marker or omit END, so anchor on the first opening
# marker and tolerate either response shape.
marker_open = _re.search(r"<<<\s*TRANSLATION\s*>>>", content, _re.I)
if marker_open:
translated_body = content[marker_open.end():]
marker_close = _re.search(r"<<<\s*END\s*>>>", translated_body, _re.I)
content = translated_body[:marker_close.start()] if marker_close else translated_body
else:
content = _re.sub(r"^\s*<<<TRANSLATION>>>\s*", "", content, flags=_re.I).strip()
content = _re.sub(r"\s*<<<END>>>\s*$", "", content, flags=_re.I).strip()
content = _extract_reply(content)
content = _re.sub(r"<<<\s*(?:TRANSLATION|END)\s*>>>", "", content, flags=_re.I).strip()
# Avoid caching duplicated output when a model emits the same
# translation twice while repairing its requested format.
paragraphs = [p.strip() for p in _re.split(r"\n\s*\n", content) if p.strip()]
if len(paragraphs) >= 2 and paragraphs[-1] == paragraphs[-2]:
paragraphs.pop()
content = "\n\n".join(paragraphs)
elif len(content) > 1 and len(content) % 2 == 0:
midpoint = len(content) // 2
if content[:midpoint].strip() == content[midpoint:].strip():
content = content[:midpoint].strip()
return content, False
since = (_dt.utcnow() - _td(days=days_back)).strftime("%d-%b-%Y")
@@ -1507,7 +1532,11 @@ async def action_classify_events(owner: str, **kwargs) -> Tuple[str, bool]:
return "No upcoming events to classify", True
from src.task_endpoint import resolve_task_candidates
llm_candidates = resolve_task_candidates(owner=owner)
llm_candidates = resolve_task_candidates(
owner=owner,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
llm_available = bool(llm_candidates)
# Pull user memories so the LLM has personal context (relationships,
@@ -1594,12 +1623,18 @@ async def action_classify_events(owner: str, **kwargs) -> Tuple[str, bool]:
from src.text_helpers import strip_think as _st
raw = _st(raw or "", prose=False, prompt_echo=False)
raw = _re.sub(r"^```(?:json)?\s*|\s*```$", "", raw, flags=_re.MULTILINE).strip()
m = _re.search(r"\[.*\]", raw, _re.DOTALL)
if not m:
# Native Qwen/Heretic responses can append a short
# explanation after an otherwise valid JSON array. Decode
# the first complete array instead of using a greedy regex
# that turns the suffix into `json.loads` Extra data.
start = raw.find("[")
if start < 0:
logger.warning(f"[classify-llm] no JSON array in response: {raw[:300]!r}")
failed += len(batch)
continue
arr = _json.loads(m.group())
arr, _end = _json.JSONDecoder().raw_decode(raw[start:])
if not isinstance(arr, list):
raise ValueError("calendar classifier returned a non-array JSON value")
by_idx = {x.get("i"): x for x in arr if isinstance(x, dict)}
for idx, ev in enumerate(batch):
x = by_idx.get(idx)
@@ -1671,6 +1706,8 @@ async def action_extract_email_events(owner: str, **kwargs) -> Tuple[str, bool]:
days_back=days_back,
account_id=account_id,
max_process=max_process,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
),
timeout=timeout,
)
@@ -1802,7 +1839,11 @@ async def action_learn_sender_signatures(owner: str, **kwargs) -> Tuple[str, boo
return "All sender sigs already cached (or no eligible senders)", True
from src.task_endpoint import resolve_task_candidates
candidates = resolve_task_candidates(owner=owner)
candidates = resolve_task_candidates(
owner=owner,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
if not candidates:
return "No LLM endpoint available", False
model = candidates[0][1]
@@ -2063,7 +2104,11 @@ async def action_test_skills(owner: str, **kwargs) -> Tuple[str, bool]:
raise TaskNoop("no skills to test")
from src.task_endpoint import resolve_task_candidates
candidates = resolve_task_candidates(owner=owner)
candidates = resolve_task_candidates(
owner=owner,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
if not candidates:
return "No Default/Utility model configured — set one in Settings.", False
@@ -2194,7 +2239,17 @@ async def action_audit_skills(owner: str, **kwargs) -> Tuple[str, bool]:
if not names:
raise TaskNoop("no unaudited skills")
url, model, headers, teacher = _resolve_audit_models(owner=owner)
try:
url, model, headers, teacher = _resolve_audit_models(
owner=owner,
model_spec=kwargs.get("model"),
endpoint_url=kwargs.get("endpoint_url"),
)
except ValueError as e:
# A missing Utility/Default model is a temporary configuration
# problem, not a completed audit. Let the scheduler retry without
# consuming the daily run or advancing the normal schedule.
raise TaskDeferred(str(e), delay_seconds=20 * 60) from e
try:
from src.llm_core import seconds_since_model_activity
recent = seconds_since_model_activity(url, model)
@@ -2432,7 +2487,11 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]:
# gate until after authoritative account cleanup. State retirement must
# still run when no model is configured.
from src.task_endpoint import resolve_task_candidates
candidates = resolve_task_candidates(owner=owner)
candidates = resolve_task_candidates(
owner=owner,
override_url=kwargs.get("endpoint_url"),
override_model=kwargs.get("model"),
)
target_account_id = _email_task_account_id(kwargs)
# ── 1. Enumerate enabled accounts. Match this task's owner AND fall
@@ -2755,10 +2814,15 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]:
triage_version=TRIAGE_VERSION,
category_tags=CATEGORY_TAGS,
)
cache.setdefault("uids", {})[item["uid"]] = verdict
per_uid_scores[key] = verdict
saved_classifications += 1
continue
# Keep deterministic handling for clearly categorized mail,
# but let ambiguous messages reach the configured task model.
# The unconditional continue here previously made the LLM
# classifier below unreachable for every email.
if verdict.get("tags") or verdict.get("reason") != "categorized by email metadata":
cache.setdefault("uids", {})[item["uid"]] = verdict
per_uid_scores[key] = verdict
saved_classifications += 1
continue
# ── LLM-classify. JSON-only response; bullet-proof parse.
llm_attempts += 1
prompt = (
+11
View File
@@ -495,6 +495,17 @@ class ChatProcessor:
f"Content from {url}:\n\n{content}",
provenance_origin="external",
))
# Automatic exact-URL reads are real network evidence even
# though they happen before the agent loop. Publish the
# source through the same provenance channel as web search
# so the UI and persisted message do not make a grounded
# answer look like an unsupported no-tool response.
if not any(source.get("url") == url for source in web_sources):
web_sources.append({
"url": url,
"title": str(result.get("title") or url),
"acquisition": "automatic_url_fetch",
})
else:
# A failed automatic URL fetch is context too. Never pass
# exception text or response-controlled diagnostics back to
+2928 -156
View File
File diff suppressed because it is too large Load Diff
+2 -2
View File
@@ -158,7 +158,7 @@ def internal_api_base() -> str:
running server over HTTP. Resolution order:
1. ODYSSEUS_INTERNAL_BASE - explicit override (e.g. behind a TLS proxy).
2. APP_PORT - http://127.0.0.1:$APP_PORT (docker-compose).
3. Fallback http://127.0.0.1:7000 - legacy default.
3. Fallback http://127.0.0.1:7011 - matches app.py's bind default.
127.0.0.1 (not "localhost") avoids IPv6/DNS ambiguity for a strictly-local
call. Without this, loopback tools fail with "All connection attempts
@@ -167,4 +167,4 @@ def internal_api_base() -> str:
override = os.environ.get("ODYSSEUS_INTERNAL_BASE")
if override:
return override.rstrip("/")
return f"http://127.0.0.1:{os.environ.get('APP_PORT', '7000')}"
return f"http://127.0.0.1:{os.environ.get('APP_PORT', '7011')}"
+261 -16
View File
@@ -83,6 +83,18 @@ Return ONLY a JSON array of query strings, nothing else.
Example: ["query one", "query two", "query three"]
"""
SMALL_MODEL_QUERY_GEN_PROMPT = """\
You choose web searches for a research task.
Today: {today}
Question: {question}
Round: {round_num}
Return ONLY a JSON array containing {num_queries} short search-query strings.
Use the question's exact topic. Do not explain your answer.
Example: ["topic latest news", "topic official sources"]
"""
RESEARCH_ACTION_PROMPT = """\
You are controlling a bounded research navigator. Choose the next actions that will best answer the user's question.
@@ -354,6 +366,7 @@ class DeepResearcher:
):
self.llm_endpoint = llm_endpoint
self.llm_model = llm_model
self.simple_research_mode = self._looks_like_small_local_model(llm_model)
self.llm_headers = llm_headers
self.search_provider_override = search_provider
self.category = category
@@ -396,6 +409,29 @@ class DeepResearcher:
"""Request cooperative cancellation of the research loop."""
self._cancelled = True
@staticmethod
def _looks_like_small_local_model(model: str) -> bool:
"""Recognize model names that commonly need a lower-complexity loop."""
name = str(model or "").lower()
for match in re.finditer(r"(?<![\w.])(\d+(?:\.\d+)?)\s*b(?!\w)", name):
try:
if 0 < float(match.group(1)) <= 10:
return True
except ValueError:
continue
return bool(
any(marker in name for marker in ("odysseus", "heretic", "trial55"))
)
@staticmethod
def _looks_like_simple_fact_question(question: str) -> bool:
"""Recognize questions that do not need iterative report writing."""
text = re.sub(r"\s+", " ", str(question or "").strip().lower())
return bool(re.match(
r"^(?:where is|what is|who is|when was|when is|how many|how far is)\b",
text,
))
# ------------------------------------------------------------------
# Public API
# ------------------------------------------------------------------
@@ -415,20 +451,44 @@ class DeepResearcher:
prior_urls: URLs already visited (won't be re-fetched).
"""
self._start_time = time.time()
self.fast_fact_mode = (
self.simple_research_mode and self._looks_like_simple_fact_question(question)
)
if self.fast_fact_mode:
# A small local model spends most of its time on synthesis rather
# than retrieval for simple factual questions. One search round
# with a compact deterministic report is both faster and safer.
self.max_rounds = min(self.max_rounds, 1)
self.min_rounds = 1
self.extraction_concurrency = min(self.extraction_concurrency, 2)
logger.info("Using fast factual research path for small model %s", self.llm_model)
findings: List[Dict] = list(prior_findings) if prior_findings else []
report = prior_report or ""
# PLAN: Analyze the question and create a research strategy
if not prior_report:
self._emit(phase="planning")
self.research_plan = await self._create_plan(question)
if self.simple_research_mode:
self.research_plan = (
"Use direct web searches for the user's question and gather "
"current, source-backed evidence."
)
logger.info("Using simplified research loop for model %s", self.llm_model)
else:
self.research_plan = await self._create_plan(question)
logger.info(f"Research plan: {self.research_plan[:200]}")
else:
# Continuation — plan around the follow-up
self._emit(phase="planning")
self.research_plan = await self._create_plan(question)
if self.simple_research_mode:
self.research_plan = (
"Use direct web searches for the user's question and gather "
"current, source-backed evidence."
)
else:
self.research_plan = await self._create_plan(question)
logger.info(f"Continuation plan: {self.research_plan[:200]}")
if not self.category and not prior_report:
if not self.category and not prior_report and not self.simple_research_mode:
self.category = await self._classify_category(question, self.research_plan)
if self.category:
logger.info(f"Auto-detected category: {self.category}")
@@ -501,6 +561,10 @@ class DeepResearcher:
# SYNTHESIZE
if findings:
if self.fast_fact_mode:
report = self._compact_fact_report(question, findings)
self.evolving_report = report
break
self._emit(phase="analyzing", round=round_num,
total_sources=len(self.urls_fetched),
total_findings=len(findings),
@@ -541,6 +605,13 @@ class DeepResearcher:
return "No information could be gathered for this question."
self.evolving_report = report # preserve pre-synthesis report
if self.fast_fact_mode:
# The compact factual path is already the final report. Sending it
# through _final_report would add another slow generation pass on
# small local models and can make a successful lookup appear to
# hang or fail.
logger.info("Research complete via fast factual report")
return report
final = await self._final_report(question, report)
elapsed = time.time() - self._start_time
logger.info(
@@ -663,20 +734,28 @@ class DeepResearcher:
"that the report doesn't yet cover well."
)
prompt = current_date_context() + QUERY_GEN_PROMPT.format(
question=question,
research_plan=self.research_plan or "(No plan — search broadly.)",
report=report or "(No findings yet.)",
round_num=round_num,
num_queries=num_queries,
round_instruction=round_instruction,
)
if getattr(self, "simple_research_mode", False):
prompt = SMALL_MODEL_QUERY_GEN_PROMPT.format(
today=datetime.now().astimezone().strftime("%Y-%m-%d"),
question=question,
round_num=round_num,
num_queries=num_queries,
)
else:
prompt = current_date_context() + QUERY_GEN_PROMPT.format(
question=question,
research_plan=self.research_plan or "(No plan — search broadly.)",
report=report or "(No findings yet.)",
round_num=round_num,
num_queries=num_queries,
round_instruction=round_instruction,
)
try:
response = await self._llm(
[{"role": "user", "content": prompt}],
temperature=0.5,
max_tokens=4096,
max_tokens=512 if getattr(self, "simple_research_mode", False) else 4096,
timeout=getattr(self, "query_timeout", 120),
)
queries = self._parse_json_array(response)
@@ -685,6 +764,33 @@ class DeepResearcher:
q for q in queries
if q not in self.queries_used and not _is_meta_search_query(q)
]
# A weak/local model can return an empty response or malformed
# JSON even when the question is perfectly searchable. Never let
# that silently terminate research with zero sources: the user's
# question is a valid broad discovery query and gives the next
# stage a chance to recover.
if not new_queries:
fallback = self._deterministic_search_topic(question)
fallback_queries = [
fallback,
f"{fallback} fact check",
f"{fallback} reliable sources",
]
new_queries = [
query for query in fallback_queries
if query and not _is_meta_search_query(query)
and query not in self.queries_used
][:num_queries]
if new_queries:
logger.warning(
"Round %s query planner returned no usable queries; "
"using deterministic fallback searches: %s",
round_num, new_queries,
)
self._emit(
phase="warning",
message="Search planning returned no usable queries; trying fallback searches.",
)
self.queries_used.update(new_queries)
logger.info(f"Round {round_num} queries: {new_queries}")
return new_queries
@@ -696,6 +802,11 @@ class DeepResearcher:
async def _plan_research_actions(self, question: str, report: str,
round_num: int) -> List[ResearchAction]:
"""Let the model choose bounded search/fetch/browser actions."""
if getattr(self, "simple_research_mode", False):
# Small local models are much more reliable at producing a short
# query list than a nested tool/action protocol. The caller will
# use _generate_queries instead.
return []
try:
from src.settings import get_setting
@@ -861,7 +972,74 @@ class DeepResearcher:
parsed = urllib.parse.urlparse(str(url or ""))
return (parsed.netloc or parsed.path.split("/", 1)[0]).lower().removeprefix("www.")
def _prioritize_search_results(self, results: List[Dict], *, limit: int) -> List[Dict]:
@staticmethod
def _topic_terms(question: str) -> Set[str]:
"""Return meaningful topic anchors from a research question.
Search engines frequently return pages that match only a generic word
such as ``best`` or ``Boston``. Those pages are especially dangerous
for small models: the extractor can turn an unrelated page into a
plausible-looking answer. Keep this deliberately conservative and
use the same anchors for search-result and fetched-page gates.
"""
stopwords = {
"a", "about", "an", "and", "are", "be", "can", "does", "for",
"from", "how", "in", "is", "it", "latest", "of", "on", "or",
"prone", "should", "the", "this", "to", "was", "were", "what",
"when", "where", "which", "why", "with", "would",
}
return {
token for token in re.findall(r"[^\W_]+", str(question or "").casefold())
if len(token) >= 2 and token not in stopwords
}
@classmethod
def _topic_overlap(cls, question: str, text: str) -> int:
"""Count distinct question anchors present in text."""
terms = cls._topic_terms(question)
haystack = str(text or "").lower()
overlap = 0
for term in terms:
variants = [term]
if term.endswith("s") and len(term) > 3:
variants.append(term[:-1])
if any(re.search(rf"(?<![a-z0-9]){re.escape(variant)}(?![a-z0-9])", haystack)
for variant in variants):
overlap += 1
return overlap
@classmethod
def _topic_relevant(cls, question: str, text: str) -> bool:
"""Require enough topical overlap to let a page reach the model."""
# This English lexical heuristic cannot decide cross-language
# relevance or segment unspaced scripts. Defer those to extraction.
if not str(question or "").isascii() or not str(text or "").isascii():
return True
terms = cls._topic_terms(question)
if not terms:
return True
overlap = cls._topic_overlap(question, text)
# A one-word topic such as "Sweden" is sufficient on its own. For
# multi-anchor questions, one shared word is not evidence of relevance
# ("Boston safety" must not qualify for Boston Terrier neurology).
return overlap >= (1 if len(terms) <= 1 else 2)
@staticmethod
def _deterministic_search_topic(question: str) -> str:
"""Turn a failed planner question into a clean search topic."""
topic = re.sub(r"\s+", " ", str(question or "").strip())
topic = re.sub(
r"^(?:please\s+)?(?:what is|what are|where is|where are|who is|"
r"when was|when is|how does|how do|can you explain)\s+",
"",
topic,
flags=re.IGNORECASE,
)
topic = re.sub(r"[?!.,;:]+$", "", topic).strip()
return topic or re.sub(r"[?!.,;:]+$", "", str(question or "").strip())
def _prioritize_search_results(self, results: List[Dict], *, limit: int,
question: str = "") -> List[Dict]:
"""Prefer stronger and more diverse search hits before extraction.
Search providers often rank broad SEO pages above primary sources. This
@@ -872,6 +1050,17 @@ class DeepResearcher:
return []
candidates = []
stopwords = {
"about", "after", "also", "best", "between", "could", "does",
"from", "have", "into", "most", "only", "people", "should",
"still", "that", "their", "there", "these", "this", "what",
"when", "where", "which", "with", "would", "your", "common",
}
definition_question = bool(re.search(
r"\b(?:define|definition|meaning|mean|what is)\b",
str(question or "").lower(),
))
question_terms = self._topic_terms(question)
seen_urls = set()
for idx, result in enumerate(results or []):
if not isinstance(result, dict):
@@ -879,18 +1068,38 @@ class DeepResearcher:
url = str(result.get("url") or "").strip()
if not url or url in seen_urls or url in self.urls_fetched:
continue
host = self._result_host(url)
if not definition_question and any(token in host for token in (
"dictionary", "wiktionary", "merriam-webster", "collinsdictionary",
)):
continue
seen_urls.add(url)
title = str(result.get("title") or "")
summary = str(result.get("content") or result.get("snippet") or "")
searchable_text = " ".join((title, summary, url)).lower()
result_terms = set(re.findall(r"[a-z0-9]+", searchable_text))
relevance = len(question_terms & result_terms)
assessment = assess_source(url, title=title, summary=summary)
candidates.append({
"idx": idx,
"host": self._result_host(url),
"host": host,
"assessment": assessment,
"relevance": relevance,
"result": result,
})
candidates.sort(key=lambda c: (-c["assessment"].score, c["host"], c["idx"]))
# If the provider returned at least one topic-relevant hit, do not
# spend extraction slots on generic dictionary/listicle results that
# only matched a word such as "best". If every hit lacks metadata or
# overlap, retain the old quality-based behavior rather than returning
# nothing.
relevant = [candidate for candidate in candidates if self._topic_relevant(
question,
" ".join((candidate["result"].get("title") or "", candidate["result"].get("content") or candidate["result"].get("snippet") or "", candidate["result"].get("url") or "")),
)]
if relevant:
candidates = relevant
candidates.sort(key=lambda c: (-c["relevance"], -c["assessment"].score, c["host"], c["idx"]))
picked = []
picked_ids = set()
used_hosts = set()
@@ -970,7 +1179,9 @@ class DeepResearcher:
raw_search_hits.append(r)
search_limit = self.max_urls_per_round * max(1, len(queries))
for r in self._prioritize_search_results(raw_search_hits, limit=search_limit):
for r in self._prioritize_search_results(
raw_search_hits, limit=search_limit, question=question
):
url = str(r.get("url") or "").strip()
if not url or url in self.urls_fetched:
continue
@@ -1093,6 +1304,29 @@ class DeepResearcher:
else:
return None
# Do this before asking the LLM to extract anything. A weak local
# model may confidently answer the goal from an unrelated page even
# when the page itself says it contains no relevant information.
page_topic_text = " ".join((page.title or title or "", page.content or "", url))
if (getattr(self, "simple_research_mode", False)
and not self._topic_relevant(question, page_topic_text)
and page.retrieval != "browser"):
browser_page = await self._browser_fallback(url, title, page)
if browser_page and browser_page.success and browser_page.content:
page = browser_page
page_topic_text = " ".join((page.title or title or "", page.content or "", url))
if (getattr(self, "simple_research_mode", False)
and not self._topic_relevant(question, page_topic_text)):
logger.info("Skipping topically unrelated research page %s", url)
self._record_navigation(
"browser_read" if page.retrieval == "browser" else requested_tool,
url=url,
title=title or page.title,
status="topic_mismatch",
retrieval=page.retrieval,
)
return None
tried_browser_after_weak_extract = False
while True:
content = page.content
@@ -1715,6 +1949,17 @@ class DeepResearcher:
f"{self._format_findings(findings)}"
)
def _compact_fact_report(self, question: str, findings: List[Dict]) -> str:
"""Build a useful answer without a second slow local-model pass."""
rows = []
for finding in findings[:4]:
title = finding.get("title") or finding.get("url") or "Source"
summary = finding.get("summary") or finding.get("evidence") or ""
url = finding.get("url") or ""
if summary:
rows.append(f"- **{title}**: {summary.strip()} [{url}]({url})")
return f"## {question.strip()}\n\n" + "\n\n".join(rows)
def get_stats(self) -> Dict:
"""Return research statistics."""
elapsed = time.time() - self._start_time if self._start_time else 0
+12 -1
View File
@@ -243,11 +243,22 @@ def _process_office_document(
if session_id:
try:
from src.office_doc import create_office_document
is_docx = str(path).lower().endswith(".docx")
stored_body = markdown
if is_docx:
# Keep the original upload addressable so the document
# pane can render a Word-style preview instead of only
# exposing the extracted Markdown.
stored_body = (
f'<!-- docx_source upload_id="{os.path.basename(path)}" -->\n'
f'{markdown}'
)
doc_id = create_office_document(
session_id=session_id,
upload_id=os.path.basename(path),
title=title,
body_text=markdown,
body_text=stored_body,
language="docx" if is_docx else "markdown",
)
if doc_id and auto_opened_docs is not None:
from src.database import SessionLocal, Document
+191
View File
@@ -0,0 +1,191 @@
"""Apply email invitation revisions without treating cancellations as creates."""
import asyncio
import errno
import hashlib
import json
import os
import uuid
from contextlib import asynccontextmanager
from datetime import datetime, timezone
from email.utils import parseaddr
from pathlib import Path
@asynccontextmanager
async def _invitation_lock(owner, sender, source_uid):
"""Serialize a series across pollers/workers, including detached instances.
File locks survive awaits without blocking the loop, release on process
exit, and don't require holding a database transaction across tool calls.
Fixed stripes bound disk usage. Never unlink lock files: another process
may already be waiting on the same inode.
"""
from src.constants import DATA_DIR
identity = json.dumps([str(owner or ""), parseaddr(sender)[1].strip().casefold(), str(source_uid).strip()])
stripe = int(hashlib.sha256(identity.encode()).hexdigest(), 16) % 64
directory = Path(DATA_DIR) / ".calendar-import-locks"
directory.mkdir(mode=0o700, parents=True, exist_ok=True)
fd = os.open(directory / f"{stripe:02x}.lock", os.O_RDWR | os.O_CREAT | getattr(os, "O_NOFOLLOW", 0), 0o600)
try:
if os.name == "nt":
import msvcrt
if os.fstat(fd).st_size == 0:
os.write(fd, b"0")
os.lseek(fd, 0, os.SEEK_SET)
acquire = lambda: msvcrt.locking(fd, msvcrt.LK_NBLCK, 1)
else:
import fcntl
acquire = lambda: fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
while True:
try:
acquire()
break
except OSError as exc:
if exc.errno not in {errno.EACCES, errno.EAGAIN, errno.EDEADLK}:
raise
await asyncio.sleep(0.025)
yield
finally:
os.close(fd)
async def apply_invitation(component, method, *, owner, sender, args):
async with _invitation_lock(owner, sender, component.get("uid") or ""):
return await _apply_invitation(component, method, owner=owner, sender=sender, args=args)
async def _apply_invitation(component, method, *, owner, sender, args):
from core.database import SessionLocal, CalendarCal, CalendarEvent, EmailCalendarInvitation
from src.tool_implementations import do_manage_calendar
from routes.calendar_routes import (
_delete_calendar_reminders_for_event, _push_caldav_event_after_commit,
_ics_naive_dtstart, _recurrence_exdates,
)
source_uid = str(component.get("uid") or "").strip()
if not source_uid:
raise ValueError("Calendar invitation is missing its UID")
sender = parseaddr(sender)[1].strip().casefold()
if not sender:
raise ValueError("Calendar invitation is missing its sender")
owner = str(owner or "")
# Untrusted ICS UIDs must never address arbitrary database event IDs.
identity = hashlib.sha256(json.dumps([owner, sender, source_uid]).encode()).hexdigest()
master_identity = identity
recurrence = component.get("recurrence-id")
recurrence_id = ""
if recurrence is not None:
if str(recurrence.params.get("RANGE", "")).upper() == "THISANDFUTURE":
raise ValueError("THISANDFUTURE invitation updates require a replacement series")
original = _ics_naive_dtstart(recurrence.dt)
recurrence_id = original.isoformat()[:16] if isinstance(recurrence.dt, datetime) else original.date().isoformat()
identity = hashlib.sha256(json.dumps([owner, sender, source_uid, recurrence_id]).encode()).hexdigest()
sequence = int(component.get("sequence", 0))
stamp_value = component.get("dtstamp")
stamp = getattr(stamp_value, "dt", None)
if isinstance(stamp, datetime):
stamp = stamp.replace(tzinfo=timezone.utc) if stamp.tzinfo is None else stamp
stamp = stamp.astimezone(timezone.utc).isoformat()
else:
stamp = ""
cancelled = str(method).upper() == "CANCEL" or str(component.get("status", "")).upper() == "CANCELLED"
# Replies describe an attendee's response, not a replacement event.
if str(method).upper() not in {"", "PUBLISH", "REQUEST", "CANCEL"}:
return {"exit_code": 0, "duplicate": True}
db = SessionLocal()
try:
master = db.get(EmailCalendarInvitation, master_identity) if recurrence_id else None
if master and master.cancelled and (sequence, stamp) <= (master.sequence, master.stamp):
return {"exit_code": 0, "duplicate": True}
state = db.get(EmailCalendarInvitation, identity)
if state and (sequence, stamp) < (state.sequence, state.stamp):
return {"exit_code": 0, "duplicate": True, "uid": state.event_uid or ""}
if state and (sequence, stamp) == (state.sequence, state.stamp):
# A cancellation wins ties; a replay must never resurrect it.
if state.cancelled or not cancelled:
return {"exit_code": 0, "duplicate": True, "uid": state.event_uid or ""}
event = None
if state and state.event_uid:
event = db.query(CalendarEvent).join(CalendarCal).filter(
CalendarEvent.uid == state.event_uid, CalendarCal.owner == owner,
).first()
if state is None:
state = EmailCalendarInvitation(id=identity, owner=owner, sender=sender, source_uid=source_uid, recurrence_id=recurrence_id)
db.add(state)
push_uids = []
def exclude_occurrence():
if master and master.event_uid:
parent = db.query(CalendarEvent).join(CalendarCal).filter(
CalendarEvent.uid == master.event_uid, CalendarCal.owner == owner,
).first()
if parent:
parent.recurrence_exdates = json.dumps(sorted(set(_recurrence_exdates(parent)) | {recurrence_id}))
push_uids.append(parent.uid)
if cancelled:
exclude_occurrence()
if event:
event.status = "cancelled"
_delete_calendar_reminders_for_event(db, owner, event)
# Retain a tombstone even if cancellation arrived before invite.
state.sequence, state.stamp, state.cancelled = sequence, stamp, True
if not recurrence_id:
# Cancelling a series also hides its detached replacements.
children = db.query(EmailCalendarInvitation).filter_by(owner=owner, sender=sender, source_uid=source_uid).all()
for child in children:
if not child.recurrence_id or (child.sequence, child.stamp) > (sequence, stamp):
continue
child.cancelled, child.sequence, child.stamp = True, sequence, stamp
child_event = db.query(CalendarEvent).join(CalendarCal).filter(
CalendarEvent.uid == child.event_uid, CalendarCal.owner == owner,
).first()
if child_event:
child_event.status = "cancelled"
_delete_calendar_reminders_for_event(db, owner, child_event)
push_uids.append(child_event.uid)
db.commit()
if event:
await _push_caldav_event_after_commit(owner, event.uid, "update")
for push_uid in push_uids:
await _push_caldav_event_after_commit(owner, push_uid, "update")
return {"exit_code": 0, "duplicate": True, "uid": state.event_uid or ""}
if not args.get("dtstart"):
raise ValueError("Calendar invitation is missing DTSTART")
action_args = dict(args)
if recurrence_id:
action_args["rrule"] = ""
if event:
action_args.update(action="update_event", uid=event.uid)
result = await do_manage_calendar(
json.dumps(action_args), owner=owner,
import_event_uid=str(uuid.uuid5(uuid.NAMESPACE_URL, "email-invitation:" + identity)),
)
if result.get("exit_code", 0) != 0:
raise RuntimeError(result.get("error") or "Calendar invitation write failed")
uid = str(result.get("uid") or (event.uid if event else ""))
if not uid:
raise RuntimeError("Calendar invitation write returned no event UID")
state.event_uid = uid
state.sequence, state.stamp, state.cancelled = sequence, stamp, False
exclude_occurrence()
if event:
event.status = "confirmed"
if not recurrence_id:
children = db.query(EmailCalendarInvitation).filter_by(owner=owner, sender=sender, source_uid=source_uid).all()
parent = db.get(CalendarEvent, uid)
if parent:
parent.recurrence_exdates = json.dumps(sorted(set(_recurrence_exdates(parent)) | {
child.recurrence_id for child in children if child.recurrence_id
}))
push_uids.append(uid)
db.commit()
if event:
await _push_caldav_event_after_commit(owner, uid, "update")
for push_uid in set(push_uids):
await _push_caldav_event_after_commit(owner, push_uid, "update")
return {**result, "uid": uid, "duplicate": bool(event) or result.get("duplicate", False)}
except Exception:
db.rollback()
raise
finally:
db.close()
+18
View File
@@ -263,6 +263,24 @@ def normalize_base(url: str) -> str:
return url
def same_endpoint_base(left, right) -> bool:
"""Allow credential reuse only for the exact API origin and base path."""
def identity(value):
parsed = urlparse(normalize_base(value))
if (parsed.scheme not in {"http", "https"} or not parsed.hostname
or parsed.username is not None or parsed.password is not None
or parsed.query or parsed.fragment or parsed.params):
return None
return (parsed.scheme, parsed.hostname.lower(),
parsed.port or (443 if parsed.scheme == "https" else 80),
parsed.path.rstrip("/"))
try:
expected = identity(right)
return expected is not None and identity(left) == expected
except ValueError:
return False
def _validated_endpoint_base(url: str) -> str:
"""Return a base URL that is safe for endpoint path appends."""
base = (url or "").strip().rstrip("/")
+13 -1
View File
@@ -19,6 +19,11 @@ logger = logging.getLogger(__name__)
_task_scheduler = None
def _event_automation_enabled_for_owner(owner: Optional[str]) -> bool:
"""Synthetic fixture activity must not auto-fire durable user tasks."""
return not str(owner or "").strip().casefold().startswith("sft_")
def set_task_scheduler(scheduler):
"""Wire up the scheduler reference (called from app.py on startup)."""
global _task_scheduler
@@ -37,7 +42,12 @@ def fire_event(event_name: str, owner: Optional[str] = None):
"""
try:
loop = asyncio.get_running_loop()
loop.create_task(_handle_event(event_name, owner))
# Let the request that emitted the event finish before automation can
# start model work on the same event loop. Otherwise a document create
# can appear to hang while an event-triggered task is running.
# Keep the handoff outside the response flush window. Event-triggered
# tasks may still perform synchronous work before their first await.
loop.call_later(1.0, lambda: loop.create_task(_handle_event(event_name, owner)))
except RuntimeError:
# No running loop — run in a new one (shouldn't happen in FastAPI)
asyncio.run(_handle_event(event_name, owner))
@@ -74,6 +84,8 @@ async def _handle_event(event_name: str, owner: Optional[str] = None):
from core.database import SessionLocal, ScheduledTask
resolved_owner = _resolve_event_owner(owner)
if not _event_automation_enabled_for_owner(resolved_owner):
return
db = SessionLocal()
try:
filters = [
+8
View File
@@ -0,0 +1,8 @@
"""Validation shared by direct native generation paths."""
import math
def validate_temperature(value):
if type(value) not in (int, float) or not math.isfinite(value) or value < 0:
raise ValueError('temperature must be a finite nonnegative number')
return float(value)
+64 -20
View File
@@ -15,6 +15,7 @@ from contextlib import asynccontextmanager
from fastapi import HTTPException
from typing import Optional, Dict, List, Tuple
from src.model_context import get_context_length, DEFAULT_CONTEXT, is_local_endpoint
from src.model_profiles import is_odysseus_merged_tools_model
from urllib.parse import urlparse
logger = logging.getLogger(__name__)
@@ -38,9 +39,9 @@ def _is_managed_stream_endpoint(url: str) -> bool:
except ValueError:
return False
_LOCAL_MODEL_LOCK = asyncio.Lock()
_LOCAL_MODEL_WAITING_FOREGROUND = 0
_LOCAL_MODEL_CURRENT: Dict[str, object] = {}
_LOCAL_MODEL_LOCKS: Dict[str, asyncio.Lock] = {}
_LOCAL_MODEL_WAITING_FOREGROUND: Dict[str, int] = {}
_LOCAL_MODEL_CURRENT: Dict[str, Dict[str, object]] = {}
def _normalize_usage_counts(input_value=0, output_value=0):
@@ -94,6 +95,14 @@ def _local_model_gate_enabled() -> bool:
return os.getenv("ODYSSEUS_LOCAL_MODEL_GATE", "true").lower() not in {"0", "false", "no", "off"}
def _local_model_gate_key(target_url: str) -> str:
"""Identify one independently schedulable local inference endpoint."""
parsed = urlparse(str(target_url or ""))
host = (parsed.hostname or "").lower()
port = parsed.port or (443 if parsed.scheme == "https" else 80)
return f"{parsed.scheme.lower()}://{host}:{port}"
def _gate_workload(workload: Optional[str]) -> str:
return "background" if str(workload or "").lower() == "background" else "foreground"
@@ -111,12 +120,15 @@ async def _local_model_slot(target_url: str, model: str, workload: Optional[str]
yield
return
global _LOCAL_MODEL_WAITING_FOREGROUND
gate_key = _local_model_gate_key(target_url)
gate_lock = _LOCAL_MODEL_LOCKS.setdefault(gate_key, asyncio.Lock())
kind = _gate_workload(workload)
current_task = asyncio.current_task()
if kind == "foreground":
_LOCAL_MODEL_WAITING_FOREGROUND += 1
current = dict(_LOCAL_MODEL_CURRENT)
_LOCAL_MODEL_WAITING_FOREGROUND[gate_key] = (
_LOCAL_MODEL_WAITING_FOREGROUND.get(gate_key, 0) + 1
)
current = dict(_LOCAL_MODEL_CURRENT.get(gate_key, {}))
if current.get("workload") == "background":
task = current.get("task")
if isinstance(task, asyncio.Task) and not task.done():
@@ -132,32 +144,38 @@ async def _local_model_slot(target_url: str, model: str, workload: Optional[str]
from src.interactive_gate import has_foreground_activity
except Exception:
has_foreground_activity = lambda: False # type: ignore
while _LOCAL_MODEL_WAITING_FOREGROUND > 0 or has_foreground_activity():
while (
_LOCAL_MODEL_WAITING_FOREGROUND.get(gate_key, 0) > 0
or has_foreground_activity()
):
await asyncio.sleep(0.25)
acquired = False
try:
await _LOCAL_MODEL_LOCK.acquire()
await gate_lock.acquire()
acquired = True
if kind == "foreground":
_LOCAL_MODEL_WAITING_FOREGROUND = max(0, _LOCAL_MODEL_WAITING_FOREGROUND - 1)
_LOCAL_MODEL_CURRENT.clear()
_LOCAL_MODEL_CURRENT.update({
_LOCAL_MODEL_WAITING_FOREGROUND[gate_key] = max(
0, _LOCAL_MODEL_WAITING_FOREGROUND.get(gate_key, 0) - 1
)
_LOCAL_MODEL_CURRENT[gate_key] = {
"task": current_task,
"workload": kind,
"url": target_url,
"model": model,
"started": time.time(),
})
}
yield
finally:
if kind == "foreground":
_LOCAL_MODEL_WAITING_FOREGROUND = max(0, _LOCAL_MODEL_WAITING_FOREGROUND - 1)
if acquired and _LOCAL_MODEL_LOCK.locked():
owner = _LOCAL_MODEL_CURRENT.get("task")
if kind == "foreground" and not acquired:
_LOCAL_MODEL_WAITING_FOREGROUND[gate_key] = max(
0, _LOCAL_MODEL_WAITING_FOREGROUND.get(gate_key, 0) - 1
)
if acquired and gate_lock.locked():
owner = _LOCAL_MODEL_CURRENT.get(gate_key, {}).get("task")
if owner is current_task:
_LOCAL_MODEL_CURRENT.clear()
_LOCAL_MODEL_LOCK.release()
_LOCAL_MODEL_CURRENT.pop(gate_key, None)
gate_lock.release()
class LLMConfig:
"""Configuration constants for LLM operations."""
@@ -1166,7 +1184,7 @@ def _is_odysseus_qwen_tool_router_model(model: str) -> bool:
or "qwen35-9b-tool-router" in value
or "qwen3.5-9b-tool-router" in value
or "odysseus-qwen3.5-9b" in value
or value.startswith("odysseus-qwen3.5-tools-")
or is_odysseus_merged_tools_model(value)
or "qwen35-email" in value
or "qwen3.5-email" in value
or "qwen35-calendar" in value
@@ -2519,6 +2537,24 @@ async def llm_call_async(
else:
messages_copy = non_sys
# Non-streaming background callers historically inherited the 32k global
# default even when the selected local endpoint exposed a smaller context
# window. Streaming requests already apply this bound; enforce the same
# invariant here before cache-key construction and payload creation.
if max_tokens and max_tokens > 0:
try:
from src.generation_budget import fit_output_token_budget
max_tokens = fit_output_token_budget(
max_tokens,
get_context_length(url, model),
messages_copy,
)
except Exception:
# Context discovery is best-effort. Preserve the established call
# path when endpoint metadata is unavailable.
pass
cache_key = _get_cache_key(
url, model, messages_copy, temperature, max_tokens, headers=headers,
thinking_mode=thinking_mode,
@@ -3554,7 +3590,15 @@ async def _stream_llm_inner(url: str, model: str, messages: List[Dict], temperat
if thinking_part:
reasoning = (reasoning + thinking_part) if reasoning else thinking_part
content = text_part
if reasoning and _normalize_thinking_mode(thinking_mode) != "off":
# DeepSeek may return reasoning_content even when the
# caller requests thinking=off, and its API requires that
# exact field on subsequent tool rounds. Preserve it in the
# reasoning channel for protocol continuity; consumers keep
# reasoning out of the visible final answer.
if reasoning and (
_normalize_thinking_mode(thinking_mode) != "off"
or "deepseek" in str(model or "").lower()
):
_degenerate = degenerate_guard.check(reasoning)
if _degenerate:
yield _degenerate
+4
View File
@@ -148,6 +148,10 @@ KNOWN_CONTEXT_WINDOWS = {
'deepseek-v3': 64000,
'deepseek-v2': 64000,
'deepseek-v4': 64000,
# Provider aliases used by configured Odysseus endpoints may omit the
# generation name. Keep them out of the unknown/small-model fallback,
# which otherwise trims multi-turn tool history to ~1K tokens.
'deepseek-flash': 64000,
# --- Google ---
'gemini-2.5-pro': 1048576,
+61
View File
@@ -0,0 +1,61 @@
"""Stable runtime profiles for models with Odysseus-specific contracts."""
from pathlib import PurePosixPath
import re
AJAX_C375_MODEL_ID = "ajax_c375"
TRIAL55_BASE_MODEL_ID = "odysseus-qwen3.5-heretic-trial55-base"
GENERIC_TOOL_SCHEMA_PROFILE = "generic"
ODYSSEUS_COMPACT_TOOL_SCHEMA_PROFILE = "odysseus_compact"
_ODYSSEUS_TOOL_PROFILE_TOKEN = re.compile(
r"(?:^|[^a-z0-9])(?:odysseus|ajax)(?:[^a-z0-9]|$)",
re.IGNORECASE,
)
def model_id_leaf(value: object) -> str:
"""Normalize a model id while preserving provider/path aliases."""
normalized = str(value or "").strip().lower().rstrip("/")
return PurePosixPath(normalized).name
def is_odysseus_tool_profile_model(value: object) -> bool:
"""Return whether a model name opts into the Odysseus tool runtime."""
return bool(_ODYSSEUS_TOOL_PROFILE_TOKEN.search(model_id_leaf(value)))
def tool_schema_profile(value: object) -> str:
"""Select the sole schema contract for a model before turn routing."""
if is_odysseus_tool_profile_model(value):
return ODYSSEUS_COMPACT_TOOL_SCHEMA_PROFILE
return GENERIC_TOOL_SCHEMA_PROFILE
def is_odysseus_merged_tools_model(value: object) -> bool:
"""Compatibility alias for the Odysseus tool runtime profile."""
return is_odysseus_tool_profile_model(value)
def uses_odysseus_progressive_thinking(value: object) -> bool:
"""Models whose native Qwen thinking is selected from the turn surface."""
return is_odysseus_tool_profile_model(value)
def supports_user_thinking_toggle(value: object) -> bool:
"""Whether the chat UI may expose an explicit thinking on/off switch."""
leaf = model_id_leaf(value)
if not leaf or uses_odysseus_progressive_thinking(leaf):
return False
if leaf.startswith(("gpt", "o1", "o3", "o4")):
return False
return any(pattern in leaf for pattern in (
"qwen3", "qwq", "deepseek-r1", "deepseek-reasoner",
"minimax", "m2-reap", "gemma", "stepfun", "step-3", "step3",
"magistral", "mistral-small", "mistral-medium",
))
+8 -3
View File
@@ -18,8 +18,11 @@ def create_office_document(
upload_id: str,
title: str,
body_text: Optional[str] = None,
language: str = "markdown",
*,
owner: Optional[str] = None,
) -> Optional[str]:
"""Create a markdown Document for an Office attachment and set it active.
"""Create a Document for an Office attachment and set it active.
Returns the new doc_id, or None on failure / empty body. The full
extracted body lives in `current_content`, so the agent can fetch
@@ -42,15 +45,17 @@ def create_office_document(
doc_id = str(uuid.uuid4())
ver_id = str(uuid.uuid4())
sess = db.query(DbSession).filter(DbSession.id == session_id).first()
if owner and sess and sess.owner != owner:
raise ValueError("Office document session belongs to a different owner")
doc = Document(
id=doc_id,
session_id=session_id,
title=title,
language="markdown",
language=language or "markdown",
current_content=body_text,
version_count=1,
is_active=True,
owner=sess.owner if sess else None,
owner=owner or (sess.owner if sess else None),
)
ver = DocumentVersion(
id=ver_id,
+9
View File
@@ -49,6 +49,13 @@ LOW_QUALITY_MARKERS = [
"copyright notice",
"copyright footer",
"all rights reserved",
# Common small-model extraction leakage: these are process narration, not
# evidence from the fetched page.
"the user wants me to extract",
"provided source data",
"i need to create",
"i will create",
"generic request",
]
@@ -58,6 +65,8 @@ def is_low_quality(summary: str) -> bool:
if not isinstance(summary, str) or not summary:
return True
low = summary.lower()
if low.strip() in {"(no content)", "no content", "(no relevant content)"}:
return True
return any(marker in low for marker in LOW_QUALITY_MARKERS)
except Exception:
return False # fail open
+40 -4
View File
@@ -3,6 +3,7 @@
from src.endpoint_resolver import (
resolve_endpoint,
resolve_utility_fallback_candidates,
same_endpoint_base as _same_endpoint_base,
)
from src.llm_core import llm_call_async_with_fallback
from src.interactive_gate import wait_for_interactive_quiet
@@ -22,6 +23,9 @@ def resolve_task_candidates(
fallback_url=None,
fallback_model=None,
fallback_headers=None,
override_url=None,
override_model=None,
override_headers=None,
owner=None,
):
"""Return ordered background-task LLM candidates.
@@ -42,6 +46,26 @@ def resolve_task_candidates(
return
candidates.append((url, model, headers or {}))
if override_url and override_model:
headers = override_headers or {}
try:
from src.database import ModelEndpoint, SessionLocal
from src.endpoint_resolver import normalize_base, resolve_endpoint_runtime, build_headers
db = SessionLocal()
try:
from src.auth_helpers import owner_filter
query = db.query(ModelEndpoint).filter(ModelEndpoint.is_enabled == True)
for ep in owner_filter(query, ModelEndpoint, owner).all():
base = normalize_base(getattr(ep, "base_url", "") or "")
if _same_endpoint_base(override_url, base):
runtime_base, api_key = resolve_endpoint_runtime(ep, owner=owner)
headers = build_headers(api_key, runtime_base or base)
break
finally:
db.close()
except Exception:
pass
_append(override_url, override_model, headers)
_append(*resolve_task_endpoint(fallback_url, fallback_model, fallback_headers, owner=owner))
_append(*resolve_endpoint("utility", owner=owner))
_append(*resolve_endpoint("default", owner=owner))
@@ -56,15 +80,27 @@ async def task_llm_call_async(
fallback_url=None,
fallback_model=None,
fallback_headers=None,
override_url=None,
override_model=None,
override_headers=None,
owner=None,
**kwargs,
):
"""Call the shared background-task LLM candidate chain."""
resolver_kwargs = {
"fallback_url": fallback_url,
"fallback_model": fallback_model,
"fallback_headers": fallback_headers,
"owner": owner,
}
if override_url is not None:
resolver_kwargs["override_url"] = override_url
if override_model is not None:
resolver_kwargs["override_model"] = override_model
if override_headers is not None:
resolver_kwargs["override_headers"] = override_headers
candidates = resolve_task_candidates(
fallback_url=fallback_url,
fallback_model=fallback_model,
fallback_headers=fallback_headers,
owner=owner,
**resolver_kwargs,
)
if not candidates:
raise RuntimeError("No LLM endpoint available for background task")
+29 -6
View File
@@ -21,6 +21,17 @@ from src.task_action_policy import (
logger = logging.getLogger(__name__)
def _is_sft_fixture_owner(owner: str | None) -> bool:
"""Synthetic SFT accounts may manage tasks but must never auto-fire them."""
return str(owner or "").strip().lower().startswith("sft_")
def _background_owner_filter(column):
"""SQL predicate matching real/ownerless accounts, excluding SFT fixtures."""
from sqlalchemy import or_
return or_(column.is_(None), ~column.like("sft\\_%", escape="\\"))
def _utcnow() -> datetime:
"""Return naive UTC for task DB fields without using deprecated APIs."""
return datetime.now(timezone.utc).replace(tzinfo=None)
@@ -552,6 +563,7 @@ class TaskScheduler:
_ST.status == "active",
_ST.next_run.isnot(None),
_ST.next_run < now,
_background_owner_filter(_ST.owner),
).all()
if overdue:
for t in overdue:
@@ -628,6 +640,7 @@ class TaskScheduler:
ScheduledTask.status == "active",
ScheduledTask.trigger_type == "schedule",
ScheduledTask.next_run.isnot(None),
_background_owner_filter(ScheduledTask.owner),
).all()
buckets: Dict[str, list] = {}
for r in rows:
@@ -713,7 +726,7 @@ class TaskScheduler:
try:
owners = set()
for r in db.query(ScheduledTask.owner).distinct().all():
if r[0]:
if r[0] and not _is_sft_fixture_owner(r[0]):
owners.add(r[0])
note_q = db.query(Note.owner).filter(
Note.due_date.isnot(None),
@@ -721,7 +734,7 @@ class TaskScheduler:
Note.archived == False, # noqa: E712
).distinct()
for r in note_q.all():
if r[0]:
if r[0] and not _is_sft_fixture_owner(r[0]):
owners.add(r[0])
return sorted(owners)
except Exception:
@@ -747,6 +760,7 @@ class TaskScheduler:
next_run = _db.query(_ST.next_run).filter(
_ST.status == "active",
_ST.next_run.isnot(None),
_background_owner_filter(_ST.owner),
).order_by(_ST.next_run.asc()).first()
if next_run and next_run[0]:
delta = (next_run[0] - _utcnow()).total_seconds()
@@ -775,6 +789,7 @@ class TaskScheduler:
due = db.query(ScheduledTask).filter(
ScheduledTask.status == "active",
ScheduledTask.next_run <= now,
_background_owner_filter(ScheduledTask.owner),
ScheduledTask.id.notin_(executing_snapshot) if executing_snapshot else True,
).all()
to_dispatch = []
@@ -1310,7 +1325,15 @@ class TaskScheduler:
# through as `command` so action_cookbook_serve can json.loads it.
elif task.action == "cookbook_serve" and task.prompt:
kwargs["command"] = task.prompt
# Model-backed actions normally use the shared Utility/Default
# chain. A task-level choice is an explicit override and must be
# available to actions such as Skills Audit as well.
if getattr(task, "model", None):
kwargs["model"] = task.model
kwargs["endpoint_url"] = getattr(task, "endpoint_url", None)
result, success = await action_fn(**kwargs)
if getattr(task, "model", None):
self._last_run_model = task.model
return result, success
except TaskNoop:
# Bubble up so _execute_task_locked can drop the run row silently.
@@ -1926,7 +1949,7 @@ class TaskScheduler:
headers = {}
try:
from core.database import SessionLocal, ModelEndpoint
from src.endpoint_resolver import normalize_base, build_headers
from src.endpoint_resolver import normalize_base, build_headers, same_endpoint_base
from src.auth_helpers import owner_filter
db2 = SessionLocal()
try:
@@ -1934,7 +1957,7 @@ class TaskScheduler:
ep_q = owner_filter(ep_q, ModelEndpoint, task.owner or None)
eps = ep_q.all()
for ep in eps:
if normalize_base(ep.base_url) in endpoint_url or endpoint_url in normalize_base(ep.base_url):
if same_endpoint_base(endpoint_url, ep.base_url):
headers = build_headers(ep.api_key, normalize_base(ep.base_url))
break
finally:
@@ -2102,7 +2125,7 @@ class TaskScheduler:
# Resolve headers
try:
from core.database import ModelEndpoint
from src.endpoint_resolver import normalize_base, build_headers
from src.endpoint_resolver import normalize_base, build_headers, same_endpoint_base
from src.auth_helpers import owner_filter
db2 = db
if not headers_from_resolver:
@@ -2110,7 +2133,7 @@ class TaskScheduler:
ep_q = owner_filter(ep_q, ModelEndpoint, task.owner or None)
eps = ep_q.all()
for ep in eps:
if normalize_base(ep.base_url) in endpoint_url or endpoint_url in normalize_base(ep.base_url):
if same_endpoint_base(endpoint_url, ep.base_url):
headers = build_headers(ep.api_key, normalize_base(ep.base_url))
break
except Exception:
+12
View File
@@ -341,6 +341,11 @@ _PRIVATE_ACTION_READS: Mapping[str, frozenset[str]] = MappingProxyType(
"manage_skills": frozenset({"list", "index", "view", "view_ref", "search"}),
"manage_tasks": frozenset({"list"}),
"manage_email_state": frozenset({"list_blocked"}),
"manage_endpoints": frozenset({"list"}),
"manage_mcp": frozenset({"list", "list_tools"}),
"manage_tokens": frozenset({"list"}),
"manage_webhooks": frozenset({"list"}),
"manage_settings": frozenset({"list", "get", "list_tools"}),
}
)
@@ -380,6 +385,13 @@ _PRIVATE_ACTION_WRITES: Mapping[str, frozenset[str]] = MappingProxyType(
"unblock_sender",
}
),
"manage_endpoints": frozenset({"add", "delete", "enable", "disable"}),
"manage_mcp": frozenset({"add", "delete", "enable", "disable", "reconnect"}),
"manage_settings": frozenset(
{"set", "delete", "reset", "disable_tool", "enable_tool"}
),
"manage_tokens": frozenset({"create", "delete"}),
"manage_webhooks": frozenset({"add", "delete", "enable", "disable"}),
}
)
+10 -2
View File
@@ -1516,13 +1516,21 @@ async def _execute_tool_block_impl(
"exit_code": 1,
"failure_kind": "turn_contract_denied",
}
if disabled_tools and not policy_names.isdisjoint(disabled_tools):
# A turn contract narrows the offered tool inventory; it is not an
# authorization grant overriding explicit execution-time restrictions.
if (
disabled_tools
and not policy_names.isdisjoint(disabled_tools)
):
desc = f"{tool}: BLOCKED"
result = {"error": f"Tool '{tool}' is disabled by user.", "exit_code": 1}
logger.info(f"Tool blocked by user: {tool}")
return desc, result
if tool_policy and any(tool_policy.blocks(name) for name in policy_names):
if (
tool_policy
and any(tool_policy.blocks(name) for name in policy_names)
):
desc = f"{tool}: BLOCKED"
result = {
"error": f"Execution of tool '{tool}' is forbade by the active guide-only policy.",
+16 -1
View File
@@ -570,7 +570,7 @@ class ToolIndex:
frozenset({"huggingface", "hugging face", "hf search",
"find a model", "search models", "search for a model",
"models for", "best model for"}):
{"search_hf_models", "list_cached_models"},
{"search_hf_models", "list_cached_models", "app_api"},
frozenset({"cached models", "list models", "my models",
"what models do i have", "is it downloaded",
"do i have", "already downloaded", "on disk"}):
@@ -627,6 +627,21 @@ class ToolIndex:
# prompts do not drag web schemas into the agent context.
if self._WEB_RE.search(query):
base.update({"web_search", "web_fetch"})
# Hardware-aware model recommendations are fulfilled by the Cookbook
# hwfit API, not by the generic endpoint/model catalog. Keep app_api in
# the caller-selected surface for natural variants such as "best model
# to run on my hardware", which do not contain the literal keyword
# phrase "best model for" above.
if (
re.search(r"\b(?:best|recommend(?:ed)?|suitable|compatible|fit)\b", ql)
and re.search(r"\bmodels?\b", ql)
and re.search(
r"\b(?:my|this|the|current)\s+(?:hardware|machine|computer|pc|server|system)\b"
r"|\b(?:gpu|vram|ram)\b",
ql,
)
):
base.add("app_api")
if re.search(r"https?://\S+(?:\.pdf\b|/pdf/)|\bPDFs?\b", query, re.I):
base.add("pdf_extract")
# Hard steering: when the query is a clear "save info about a specific
+115 -6
View File
@@ -16,6 +16,8 @@ MODEL_CHOICE_MODEL = 'odysseus-qwen3.5-tools-pre-heretic'
WEB_REFERENCE = re.compile(
r'https?://[^\s<>]+'
r'|(?<![\w@./-])(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+'
r'(?!(?:txt|md|json|csv|tsv|ya?ml|xml|log|pdf|docx?|xlsx?|pptx?|'
r'png|jpe?g|gif|webp|svg|py|js|ts|css|html?)(?![a-z]))'
r'[a-z]{2,63}(?![\w@.-])', re.I,
)
@@ -68,18 +70,125 @@ def select_experiment_inventory(inventory, routed, history, mode, *, user_text='
families.update(url_family)
research_family = {'research'} if mode == MODEL_CHOICE_MODE and has_research_hint(user_text) else set()
families.update(research_family)
families.update(recently_executed_families(
history, user_turns=6, maximum=3,
include_failed_attempts=mode == MODEL_CHOICE_MODE,
))
explicit_image_edit = (
mode == MODEL_CHOICE_MODE
and routed.capabilities == frozenset({'image_editing'})
and routed.required_read_operation is None
and any(canonical_tool(name) == 'edit_image' for name in routed.required)
)
if not explicit_image_edit:
families.update(recently_executed_families(
history, user_turns=6, maximum=3,
include_failed_attempts=mode == MODEL_CHOICE_MODE,
))
names = set().union(*(FAMILY_TOOLS.get(f, ()) for f in families))
offered = frozenset(n for n in inventory.offered
if mode == 'all' or canonical_tool(n) in names)
# The model-choice rollout was intentionally launched without lexical
# tool forcing so we could observe the fine-tuned model's own selection.
# Replays now show a narrower failure boundary: the model sometimes
# ignores an already-resolved, read-only list/search/repeat operation and
# fabricates or emits an empty lead-in. Preserve only the router's sealed
# safe-read operation in this model-specific mode. Ambiguous requests still
# have no operation and remain model-selected; mutation authority is
# unchanged.
sealed_read = (
routed.required_read_operation if mode == MODEL_CHOICE_MODE else None
)
sealed_read_required = frozenset(
name for name in offered
if sealed_read is not None
and canonical_tool(name) == canonical_tool(sealed_read.tool)
)
if sealed_read is not None and not sealed_read_required:
# Never retain an operation whose own tool was removed by permissions.
# A different required action cannot satisfy this invariant.
sealed_read = None
sealed_required = sealed_read_required
explicit_cookbook_action = (
mode == MODEL_CHOICE_MODE
and routed.capabilities == frozenset({'cookbook_admin'})
and {
canonical_tool(name) for name in routed.required
} <= {'download_model', 'serve_preset', 'stop_served_model'}
and bool(routed.required)
)
if explicit_cookbook_action:
sealed_required |= frozenset(
name for name in offered
if canonical_tool(name) in {
canonical_tool(required) for required in routed.required
}
)
if explicit_image_edit:
# An explicit supported image edit has one execution owner. Preserve
# that typed requirement so prose cannot fabricate or refuse an
# operation the user clearly requested and the backend can perform.
sealed_required |= frozenset(
name for name in offered if canonical_tool(name) == 'edit_image'
)
explicit_web_read = (
mode == MODEL_CHOICE_MODE
and routed.capabilities == frozenset({'search_browser'})
and bool(routed.required)
and {
canonical_tool(name) for name in routed.required
} <= {'web_search', 'web_fetch'}
)
if explicit_web_read:
# Search discovery and page retrieval are distinct read-only
# operations. Once the turn router resolves one exactly, retaining
# the whole warm web family lets the model substitute browser
# navigation or repeat an old search. Preserve the resolved read while
# leaving genuinely ambiguous web turns model-selected.
required_web_names = {
canonical_tool(required) for required in routed.required
}
# Keep one immutable recovery-capable set. The resolved reader still
# executes first, but a failed/empty brokered read may recover through
# page fetch or the private browser without rebuilding the contract.
# This avoids both premature abandonment and mid-turn permission
# expansion.
recovery_names = set(required_web_names) | {'private_browser'}
if 'web_search' in required_web_names:
recovery_names.add('web_fetch')
offered = frozenset(
name for name in offered
if canonical_tool(name) in recovery_names
)
sealed_required |= frozenset(
name for name in offered
if canonical_tool(name) in required_web_names
)
explicit_model_call = (
mode == MODEL_CHOICE_MODE
and bool(routed.required)
and {
canonical_tool(name) for name in routed.required
} == {'chat_with_model'}
)
if explicit_model_call:
offered = frozenset(
name for name in offered if canonical_tool(name) == 'chat_with_model'
)
sealed_required |= offered
explicit_chat_history_search = (
mode == MODEL_CHOICE_MODE
and bool(routed.required)
and {
canonical_tool(name) for name in routed.required
} == {'search_chats'}
)
if explicit_chat_history_search:
offered = frozenset(
name for name in offered if canonical_tool(name) == 'search_chats'
)
sealed_required |= offered
return replace(
inventory, offered=offered, required=frozenset(),
inventory, offered=offered, required=sealed_required,
schema_json=tuple(s for s in inventory.schema_json
if json.loads(s)['function']['name'] in offered),
required_read_operation=None, routing_experiment=mode,
required_read_operation=sealed_read, routing_experiment=mode,
# Available families are not mutation authorization. Keep the original
# request's authority; selection only changes what the model can see.
active_capabilities=routed.active_capabilities | frozenset(url_family | research_family),
+53 -25
View File
@@ -112,6 +112,16 @@ def _repair_function_arg_aliases(tool_type: str, args: dict[str, Any]) -> dict[s
args["url"] = path
args.pop("path", None)
args.pop("file_path", None)
if tool_type == "manage_documents" and "limit" not in args and "max_results" in args:
# Collection APIs use both names across the native tool surface. The
# document contract calls this integer ``limit``.
args["limit"] = args.pop("max_results")
if tool_type == "manage_tasks" and str(args.get("action") or "").casefold() == "list":
# manage_tasks has no backend result-limit argument; the canonical
# renderer applies the user's visible cap. Drop only these familiar
# collection aliases so they cannot invalidate an otherwise safe read.
args.pop("max_results", None)
args.pop("limit", None)
return args
@@ -362,11 +372,7 @@ FUNCTION_TOOL_SCHEMAS = [
"path": {"type": "string", "description": "Task-local /workspace/*.pdf path; use url for an online PDF"},
"query": {"type": "string", "description": "Required focused terms, including the target model and every requested metric/table heading"}
},
"required": ["query"],
"anyOf": [
{"required": ["url"]},
{"required": ["path"]}
]
"required": ["query"]
}
}
},
@@ -656,7 +662,11 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "object",
"properties": {
"title": {"type": "string", "description": "Document title"},
"language": {"type": "string", "description": "Programming language or format. Use richtext for formatted prose/articles the user should edit visually; use html only when the user explicitly asks for HTML source/code or a runnable HTML page (e.g. python, javascript, markdown, richtext, text, html)."},
"language": {
"type": "string",
"enum": ["python", "javascript", "typescript", "html", "css", "richtext", "markdown", "json", "yaml", "bash", "sql", "rust", "go", "java", "c", "cpp", "xml", "toml", "ini", "ruby", "php", "csv", "email", "text", "plain", "svg"],
"description": "Editor language or format. This is not a human-language code: use richtext for formatted prose/articles and markdown or text for plain prose; use html only for requested HTML source or a runnable page."
},
"content": {"type": "string", "description": "The document content"}
},
"required": ["title", "content"]
@@ -696,7 +706,7 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "function",
"function": {
"name": "suggest_document",
"description": "Suggest improvements to the active document WITHOUT editing it. Creates inline comment bubbles the user can accept or reject. Use when the user asks for suggestions, review, improvements, or feedback.",
"description": "Suggest improvements to the active document WITHOUT editing it. Creates inline comment bubbles the user can accept or reject. Use when the user asks for suggestions, review, improvements, or feedback. Every replacement must materially differ from its exact source text; never emit a no-op suggestion.",
"parameters": {
"type": "object",
"properties": {
@@ -707,7 +717,7 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "object",
"properties": {
"find": {"type": "string", "description": "Exact text in the document to suggest changing"},
"replace": {"type": "string", "description": "Suggested replacement text"},
"replace": {"type": "string", "description": "Suggested replacement text; MUST be materially different from find"},
"reason": {"type": "string", "description": "Brief explanation of why this change helps"}
},
"required": ["find", "replace", "reason"]
@@ -871,7 +881,7 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "function",
"function": {
"name": "list_models",
"description": "List all available AI models across configured endpoints. Optionally filter by keyword.",
"description": "List AI models across configured endpoints. A normal filter matches model IDs. Use filter='recommended' to detect this machine's GPU/VRAM/RAM/CPU and return ranked compatible models.",
"parameters": {
"type": "object",
"properties": {
@@ -885,13 +895,14 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "function",
"function": {
"name": "ui_control",
"description": "Control the user interface. Actions: toggle (turn tools on/off), open_panel (open a modal: documents/library, gallery, calendar/schedule, email, sessions, notes, memories/brain, skills, settings, theme, cookbook; calendar also supports `open_panel calendar month|week|year|agenda [YYYY-MM or YYYY-MM-DD]`; for 'that month/week' after a calendar listing, carry over the listed range, e.g. `open_panel calendar month 2026-09`), open_email_reply (legacy UI-only reply opener; prefer email MCP draft_email_reply for assistant-written reply drafts so a normal document-backed email draft is created), set_mode, switch_model, set_theme (built-in presets: dark, light, midnight, paper, cyberpunk, retrowave, forest, ocean, ume, copper, terminal, organs, lavender, gpt, claude, cute), create_theme (CREATE any custom theme with a name + colors object — pick distinctive, evocative hex colors that match the requested aesthetic, NOT generic defaults. The theme auto-applies after creation). When a user asks for ANY theme not in the built-in preset list, ALWAYS use create_theme.",
"description": "Control the user interface. Actions: toggle (turn tools on/off), open_panel (open a modal: documents/library, gallery, calendar/schedule, email, sessions, notes, memories/brain, skills, settings, theme, cookbook; calendar supports month/week/year/agenda plus a date; Cookbook supports models/download, launch/serve, active/running, dependencies, and settings views), open_email_reply (legacy UI-only reply opener; prefer email MCP draft_email_reply for assistant-written reply drafts so a normal document-backed email draft is created), set_mode, switch_model, set_theme (built-in presets: dark, light, midnight, paper, cyberpunk, retrowave, forest, ocean, ume, copper, terminal, organs, lavender, gpt, claude, cute), create_theme (CREATE any custom theme with a name + colors object — pick distinctive, evocative hex colors that match the requested aesthetic, NOT generic defaults. The theme auto-applies after creation), get_theme, and get_toggles. When a user asks for ANY theme not in the built-in preset list, ALWAYS use create_theme.",
"parameters": {
"type": "object",
"properties": {
"action": {"type": "string", "enum": ["toggle", "open_panel", "open_email_reply", "set_mode", "switch_model", "set_theme", "create_theme", "get_toggles"],
"action": {"type": "string", "enum": ["toggle", "open_panel", "open_email_reply", "set_mode", "switch_model", "set_theme", "create_theme", "get_theme", "get_toggles"],
"description": "The UI action. Use set_theme for presets, create_theme to build a custom theme with any hex colors"},
"name": {"type": "string", "description": "For toggle: web, bash, research, incognito, document_editor (aliases: shell, search, deepresearch, documents). For open_panel: documents, gallery, calendar/schedule, email, sessions, notes, brain/memories, skills, settings, theme/themes, cookbook. For open_email_reply: email UID. For set_theme: a preset theme name. For create_theme: the custom theme name."},
"name": {"type": "string", "description": "For toggle: web, bash, research, incognito, document_editor (aliases: shell, search, deepresearch, documents). For open_panel: documents, gallery, calendar/schedule, email, sessions, notes, brain/memories, skills, settings, theme/themes, cookbook; models and serve are Cookbook-view aliases. For open_email_reply: email UID. For set_theme: a preset theme name. For create_theme: the custom theme name."},
"view": {"type": "string", "description": "Optional open_panel subview: calendar day/week/month/year/agenda, or Cookbook models/download, launch/serve, active/running, dependencies, settings."},
"value": {"type": "string", "description": "Value: on/off for toggle, agent/chat for set_mode, model name for switch_model, theme name for set_theme, or folder for open_email_reply"},
"uid": {"type": "string", "description": "Email UID for open_email_reply"},
"folder": {"type": "string", "description": "Email folder for open_email_reply (default INBOX)"},
@@ -1443,7 +1454,7 @@ FUNCTION_TOOL_SCHEMAS = [
"type": "function",
"function": {
"name": "app_api",
"description": "Generic loopback to allowed internal Odysseus endpoints. Use this when there's no named tool for what the user wants. Hits the same routes the UI buttons hit (cookbook, gallery, library/documents, memory, notes, calendar, tasks, settings, themes, research, compare, etc.). action='endpoints' returns the OpenAPI surface (use `filter` to narrow). action='call' (default) takes method+path+body. Sensitive auth/user/admin/shell paths and host-control Cookbook mutation routes are blocked for safety. Do not use for shell commands; use named command tooling instead. Do not use for package installs, engine rebuilds, PID signalling, or email account discovery; use list_email_accounts for email accounts because /api/email/accounts is owner-filtered in tool context.",
"description": "Generic loopback to allowed internal Odysseus endpoints. Use this when there's no named tool for what the user wants. For 'best model for my hardware', call GET /api/hwfit/models with query {fit_only:true,limit:10,sort:'fit'}; it detects GPU/VRAM/RAM/CPU and returns ranked compatible models. Hits the same routes the UI buttons hit (cookbook, gallery, library/documents, memory, notes, calendar, tasks, settings, themes, research, compare, etc.). action='endpoints' returns the OpenAPI surface (use `filter` to narrow). action='call' (default) takes method+path+body. Sensitive auth/user/admin/shell paths and host-control Cookbook mutation routes are blocked for safety. Do not use for shell commands; use named command tooling instead. Do not use for package installs, engine rebuilds, PID signalling, or email account discovery; use list_email_accounts for email accounts because /api/email/accounts is owner-filtered in tool context.",
"parameters": {
"type": "object",
"properties": {
@@ -2043,18 +2054,23 @@ def function_call_to_tool_block(name: str, arguments: str) -> Optional[ToolBlock
tool_type, args = normalize_native_function_args(name, args)
if tool_type == "web_fetch" and isinstance(args.get("urls"), list):
# Some compact-model calls encode a URL as [url, ""] (an empty label
# slot) inside the batch. This is unambiguous, so normalize it without
# accepting arbitrary nested shapes.
# Some compact-model calls encode a URL as [url, label] inside the
# batch. The first value is still an explicit HTTP(S) URL and the
# second is display-only prose, so this two-string shape is
# unambiguous. Normalize it without accepting arbitrary nested data.
normalized_urls = []
for item in args["urls"]:
if (
isinstance(item, list)
and item
and len(item) in (1, 2)
and isinstance(item[0], str)
and all(not str(value or "").strip() for value in item[1:])
and item[0].strip().lower().startswith(("http://", "https://"))
and (len(item) == 1 or isinstance(item[1], str))
):
normalized_urls.append(item[0])
elif isinstance(item, list):
logger.warning("Rejecting ambiguous nested web_fetch URL item: %r", item)
return None
else:
normalized_urls.append(item)
args["urls"] = normalized_urls
@@ -2129,15 +2145,19 @@ def function_call_to_tool_block(name: str, arguments: str) -> Optional[ToolBlock
elif tool_type == "python":
content = args.get("code", "")
elif tool_type == "web_search":
# ``query`` is the canonical schema field. Some native wrappers also
# include ``command": "web_search"`` as transport metadata; treating
# that metadata as the query silently searches for the tool's name.
# Keep legacy aliases only as fallbacks when the canonical field is
# absent.
content = args.get("query", "")
queries = args.get("queries")
if isinstance(queries, list) and queries:
if not content and isinstance(queries, list) and queries:
content = str(queries[0])
elif queries:
elif not content and queries:
content = str(queries)
elif args.get("command"):
elif not content and args.get("command"):
content = args.get("command", "")
else:
content = args.get("query", "")
# Preserve the model-requested freshness filter — the web_search schema
# advertises time_filter and the executor parses {"query","time_filter"},
# but a bare query string dropped it. Mirrors the read_file JSON idiom.
@@ -2164,8 +2184,14 @@ def function_call_to_tool_block(name: str, arguments: str) -> Optional[ToolBlock
content = json.dumps(args)
elif tool_type == "create_document":
parts = [args.get("title", "Untitled")]
if args.get("language"):
parts.append(args["language"])
language = str(args.get("language") or "").strip().casefold()
# A common model slip is treating this editor-format field as a human
# language and emitting ``en``/``English``. The legacy line transport
# interpreted an unknown second line as document content, visibly
# prepending it to the user's prose. Preserve the document body and
# let the executor's content sniffer select markdown instead.
if language not in {"en", "eng", "english"} and language:
parts.append(language)
parts.append(args.get("content", ""))
content = "\n".join(parts)
elif tool_type == "edit_document":
@@ -2281,6 +2307,8 @@ def function_call_to_tool_block(name: str, arguments: str) -> Optional[ToolBlock
content = f"toggle {name} {value}"
elif action == "open_panel":
content = f"open_panel {name or value}"
if args.get("view"):
content += f" {args['view']}"
elif action == "open_email_reply":
uid = args.get("uid") or name
folder = args.get("folder") or value or "INBOX"
+90 -17
View File
@@ -17,7 +17,7 @@ from src.upload_handler import reserve_upload_references
logger = logging.getLogger(__name__)
async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
async def do_manage_calendar(content: str, owner: Optional[str] = None, *, import_event_uid: Optional[str] = None) -> Dict:
"""Handle manage_calendar tool calls: list/create/update/delete calendar events (local SQLite)."""
from core.database import SessionLocal, CalendarCal, CalendarEvent, Note
from routes.calendar_routes import (
@@ -102,6 +102,22 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
q = q.filter(CalendarCal.owner == owner)
return q
def _event_uid_candidates(raw_uid):
"""Yield exact UID first, then unambiguous UI-anchor spellings.
Calendar results render links as ``#event-<uid>``. Models sometimes
copy that href (or drop only the leading ``#``) into the UID field.
Preserve real UIDs beginning with ``event-`` by trying the exact value
first and using the stripped form only as a not-found fallback.
"""
text = str(raw_uid or "").strip()
candidates = [text]
if text.startswith("#event-"):
candidates.append(text[len("#event-"):])
elif text.startswith("event-"):
candidates.append(text[len("event-"):])
return [item for index, item in enumerate(candidates) if item and item not in candidates[:index]]
def _first_present_arg(raw_args, *names: str):
for name in names:
if name in raw_args and raw_args.get(name) is not None:
@@ -273,11 +289,11 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
"exit_code": 1,
}
if start_raw:
start_dt = _parse_dt(start_raw)
start_dt, _ = _parse_event_dt(start_raw)
else:
start_dt = datetime.utcnow().replace(hour=0, minute=0, second=0, microsecond=0)
if end_raw:
end_dt = _parse_dt(end_raw)
end_dt, _ = _parse_event_dt(end_raw)
else:
end_dt = start_dt + timedelta(days=14)
except ValueError as e:
@@ -421,13 +437,41 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
existing = (
_event_query()
.filter(
CalendarEvent.dtstart == dtstart,
CalendarEvent.status != "cancelled",
_func.lower(CalendarEvent.summary) == summary.lower(),
*([CalendarEvent.uid == import_event_uid] if import_event_uid else [
CalendarEvent.dtstart == dtstart,
CalendarEvent.status != "cancelled",
_func.lower(CalendarEvent.summary) == summary.lower(),
]),
)
.first()
)
if existing is not None:
# Repair older email-imported events whose model-generated
# location was an unrelated map URL. A concrete meeting URL
# is stronger evidence than the existing free-text location.
incoming_location = str(args.get("location") or "").strip()
changed = False
if incoming_location and re.match(
r"^https?://(?:teams\.microsoft\.com|(?:[a-z0-9-]+\.)?zoom\.us|meet\.google\.com|(?:[a-z0-9-]+\.)?webex\.com|meet\.jit\.si)/",
incoming_location,
re.IGNORECASE,
) and (
not str(existing.location or "").strip()
or not re.match(
r"^https?://(?:teams\.microsoft\.com|(?:[a-z0-9-]+\.)?zoom\.us|meet\.google\.com|(?:[a-z0-9-]+\.)?webex\.com|meet\.jit\.si)/",
str(existing.location or "").strip(),
re.IGNORECASE,
)
):
existing.location = incoming_location
changed = True
for field in ("source_email_uid", "source_email_folder", "source_email_account_id", "source_email_message_id"):
incoming = str(args.get(field) or "").strip()
if incoming and not getattr(existing, field, None):
setattr(existing, field, incoming)
changed = True
if changed:
db.commit()
reminder_note_id = None
reminder_skipped_reason = None
minutes_before = _reminder_minutes(args)
@@ -486,7 +530,7 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
"exit_code": 1,
}
uid = str(_uuid.uuid4())
uid = import_event_uid or str(_uuid.uuid4())
ev = CalendarEvent(
uid=uid, calendar_id=cal.id, summary=summary,
description=event_description,
@@ -496,6 +540,10 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
rrule=args.get("rrule", "") or "",
event_type=event_type,
importance=importance,
source_email_uid=str(args.get("source_email_uid") or "").strip() or None,
source_email_folder=str(args.get("source_email_folder") or "").strip() or None,
source_email_account_id=str(args.get("source_email_account_id") or "").strip() or None,
source_email_message_id=str(args.get("source_email_message_id") or "").strip() or None,
caldav_sync_pending="create" if cal.source == "caldav" else None,
)
db.add(ev)
@@ -545,11 +593,17 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
uid = args.get("summary")
if not uid:
return {"error": "uid is required", "exit_code": 1}
try:
base_uid = _resolve_base_uid(uid)
except ValueError as e:
return {"error": str(e), "exit_code": 1}
ev = _event_query().filter(CalendarEvent.uid == base_uid).first()
ev = None
base_uid = ""
for candidate_uid in _event_uid_candidates(uid):
try:
candidate_base_uid = _resolve_base_uid(candidate_uid)
except ValueError:
continue
ev = _event_query().filter(CalendarEvent.uid == candidate_base_uid).first()
if ev:
base_uid = candidate_base_uid
break
if not ev:
title_matches = _event_query().filter(
CalendarEvent.summary == str(uid).strip()
@@ -581,6 +635,8 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
ev.description = args["description"]
if args.get("location") is not None:
ev.location = args["location"]
previous_dtstart = ev.dtstart
previous_dtend = ev.dtend
if args.get("dtstart") is not None:
# Anchor naive/natural-language input to the USER's timezone and
# refresh is_utc, exactly like create_event. Parsing with the
@@ -594,6 +650,13 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
ev.all_day = False
ev.dtstart, _su = _parse_event_dt(args["dtstart"])
ev.is_utc = bool(_su and not _eff_all_day)
if (
args.get("dtend") is None
and previous_dtstart is not None
and previous_dtend is not None
and previous_dtend > previous_dtstart
):
ev.dtend = ev.dtstart + (previous_dtend - previous_dtstart)
if args.get("dtend") is not None:
ev.dtend, _eu = _parse_event_dt(args["dtend"])
if args.get("all_day") is None and bool(ev.all_day) and _looks_like_timed_dt(args["dtend"]):
@@ -607,6 +670,10 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
ev.event_type = _tag or None
if args.get("importance") is not None:
ev.importance = args["importance"]
for field in ("source_email_uid", "source_email_folder", "source_email_account_id", "source_email_message_id"):
incoming = str(args.get(field) or "").strip()
if incoming:
setattr(ev, field, incoming)
if args.get("rrule") is not None:
ev.rrule = args.get("rrule") or ""
elif str(args.get("repeat") or "").strip().lower() in {"none", "no", "off", "false", "single"}:
@@ -672,11 +739,17 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None) -> Dict:
return {"error": "Multiple events have that exact title; uid is required", "exit_code": 1}
if not uid:
return {"error": "uid or exact summary is required", "exit_code": 1}
try:
base_uid = _resolve_base_uid(uid)
except ValueError as e:
return {"error": str(e), "exit_code": 1}
ev = _event_query().filter(CalendarEvent.uid == base_uid).first()
ev = None
base_uid = ""
for candidate_uid in _event_uid_candidates(uid):
try:
candidate_base_uid = _resolve_base_uid(candidate_uid)
except ValueError:
continue
ev = _event_query().filter(CalendarEvent.uid == candidate_base_uid).first()
if ev:
base_uid = candidate_base_uid
break
if not ev:
return {"error": f"Event {uid} not found", "exit_code": 1}
is_caldav = ev.calendar and ev.calendar.source == "caldav" and ev.remote_href
+23 -1
View File
@@ -1829,7 +1829,11 @@ async def do_list_cached_models(content: str, owner: Optional[str] = None) -> Di
resp.raise_for_status()
data = resp.json()
if isinstance(data, dict) and data.get('error'):
raise ValueError('cache endpoint reported an error')
scan_errors.append({
'host': host_label or 'local',
'reason': str(data.get('error'))[:500],
})
return []
ms = data.get("models", []) if isinstance(data, dict) else (data or [])
for m in ms:
m["host"] = host_label or "local"
@@ -1885,7 +1889,25 @@ async def do_list_cached_models(content: str, owner: Optional[str] = None) -> Di
and (s.get("name") == raw_host or s.get("host") == host or s.get("host") == raw_host)),
{},
)
error_start = len(scan_errors)
models = await _scan_one(raw_host, host, model_dir=_dirs_for(srv))
# Friendly Cookbook names commonly double as SSH aliases. If a
# saved LAN address goes stale after a reboot/network change,
# retry the validated alias before declaring the server offline.
# This is read-only and never mutates the saved configuration.
if not models and len(scan_errors) > error_start and host != raw_host:
try:
alias = validate_remote_host(raw_host)
except Exception:
alias = None
if alias:
configured_errors = scan_errors[error_start:]
del scan_errors[error_start:]
models = await _scan_one(
raw_host, alias, model_dir=_dirs_for(srv),
)
if not models:
scan_errors[error_start:error_start] = configured_errors
else:
# Always include local. Local's saved record is the one with no host.
local_srv = next((s for s in servers if isinstance(s, dict) and not (s.get("host") or "").strip()), {})
+17 -6
View File
@@ -16,6 +16,20 @@ from src.upload_handler import reserve_upload_references
logger = logging.getLogger(__name__)
def _search_tokens(value: str) -> list[str]:
"""Normalize lightweight singular/plural variants without fuzzy matching."""
tokens = []
for token in re.findall(r"[a-z0-9]+", str(value or "").lower()):
if token in {"the", "a", "an", "note", "notes", "checklist", "list", "todo", "todos"}:
continue
if len(token) > 4 and token.endswith("ies"):
token = token[:-3] + "y"
elif len(token) > 3 and token.endswith("s") and not token.endswith("ss"):
token = token[:-1]
tokens.append(token)
return tokens
async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict:
"""Handle manage_notes tool calls: CRUD on notes and checklists."""
import uuid as _uuid
@@ -206,19 +220,16 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict:
or ""
).strip().lower()
if query:
query_terms = [
term
for term in re.findall(r"[a-z0-9]+", query)
if term not in {"the", "a", "an", "note", "notes", "checklist", "list", "todo", "todos"}
]
query_terms = _search_tokens(query)
filtered = []
for n in notes:
haystack = " ".join(
str(part or "")
for part in (n.title, n.content, n.label, n.items)
).lower()
haystack_terms = set(_search_tokens(haystack))
if query in haystack or (
query_terms and all(term in haystack for term in query_terms)
query_terms and all(term in haystack_terms for term in query_terms)
):
filtered.append(n)
notes = filtered
+9 -2
View File
@@ -8,6 +8,7 @@ tools.
tool_implementations.py and are pulled back function-locally where needed.
"""
import re
from datetime import datetime, timezone
from typing import Any, Dict, Optional
from src.constants import DEEP_RESEARCH_DIR
@@ -92,9 +93,15 @@ async def do_manage_research(content: str, owner: Optional[str] = None) -> Dict:
# the `research-` UI prefix, while action=read expects the underlying file
# stem. Exposing the exact id prevents agents from guessing or retrying
# alternate spellings after a list call.
def _completed_label(value):
try:
return datetime.fromtimestamp(float(value), timezone.utc).isoformat().replace('+00:00', 'Z')
except (TypeError, ValueError, OSError):
return 'completion time unavailable'
rows = "\n".join(
f"- [{q or '(untitled)'}](#research-{sid}) — id: {sid} — {n} sources"
for _, sid, q, n in items[:50]
f"- [{q or '(untitled)'}](#research-{sid}) — id: {sid} — completed {_completed_label(completed)} — {n} sources"
for completed, sid, q, n in items[:50]
)
return {"output": f"Research library ({len(items)} item{'s' if len(items) != 1 else ''}):\n{rows}", "exit_code": 0}
+4314 -65
View File
File diff suppressed because it is too large Load Diff