chore: resolve all ruff linting and static analysis errors in EvoScientist/
- Fix RUF006: Implement background task tracking in Discord, iMessage, WeChat, and TUI to prevent premature GC of fire-and-forget tasks. - Fix B904: Add explicit exception chaining (raise ... from) across all exception handlers. - Fix RUF012: Annotate mutable class attributes with ClassVar for command arguments and media maps. - Fix B008: Refactor Typer commands in cli/commands.py to use Annotated for argument and option defaults. - Fix B023/B018: Resolve late-binding issues in lambdas and remove useless expressions. - Fix syntax errors in retry.py docstrings and models.py lambda parameter ordering.
This commit is contained in:
@@ -167,16 +167,16 @@ def _build_base_kwargs(base_backend, base_middleware):
|
||||
prompt_refs=_build_prompt_refs(),
|
||||
)
|
||||
_inject_subagent_middleware(subs)
|
||||
return dict(
|
||||
name="EvoScientist",
|
||||
model=_ensure_chat_model(),
|
||||
tools=list(base_tools),
|
||||
backend=base_backend,
|
||||
subagents=subs,
|
||||
middleware=base_middleware,
|
||||
system_prompt=get_system_prompt(),
|
||||
skills=["/skills/"],
|
||||
)
|
||||
return {
|
||||
"name": "EvoScientist",
|
||||
"model": _ensure_chat_model(),
|
||||
"tools": list(base_tools),
|
||||
"backend": base_backend,
|
||||
"subagents": subs,
|
||||
"middleware": base_middleware,
|
||||
"system_prompt": get_system_prompt(),
|
||||
"skills": ["/skills/"],
|
||||
}
|
||||
|
||||
|
||||
def load_mcp_and_build_kwargs(base_backend, base_middleware):
|
||||
@@ -218,16 +218,16 @@ def load_mcp_and_build_kwargs(base_backend, base_middleware):
|
||||
if sa_tools := mcp_by_agent.get(sa["name"], []):
|
||||
sa.setdefault("tools", []).extend(sa_tools)
|
||||
|
||||
return dict(
|
||||
name="EvoScientist",
|
||||
model=_ensure_chat_model(),
|
||||
tools=base_tools + mcp_main,
|
||||
backend=base_backend,
|
||||
subagents=subs,
|
||||
middleware=base_middleware,
|
||||
system_prompt=get_system_prompt(),
|
||||
skills=["/skills/"],
|
||||
)
|
||||
return {
|
||||
"name": "EvoScientist",
|
||||
"model": _ensure_chat_model(),
|
||||
"tools": base_tools + mcp_main,
|
||||
"backend": base_backend,
|
||||
"subagents": subs,
|
||||
"middleware": base_middleware,
|
||||
"system_prompt": get_system_prompt(),
|
||||
"skills": ["/skills/"],
|
||||
}
|
||||
|
||||
|
||||
# =============================================================================
|
||||
|
||||
@@ -24,26 +24,26 @@ ChannelServer = Channel
|
||||
|
||||
__all__ = [
|
||||
"Channel",
|
||||
"ChannelServer",
|
||||
"ChannelManager",
|
||||
"MessageBus",
|
||||
"RawIncoming",
|
||||
"IncomingMessage",
|
||||
"OutgoingMessage",
|
||||
"InboundMessage",
|
||||
"OutboundMessage",
|
||||
"InboundConsumer",
|
||||
"run_standalone",
|
||||
"register_channel",
|
||||
"create_channel",
|
||||
"available_channels",
|
||||
# New modules
|
||||
"ChannelCapabilities",
|
||||
"UnifiedFormatter",
|
||||
"TypingManager",
|
||||
"chunk_text",
|
||||
"ChannelManager",
|
||||
"ChannelMeta",
|
||||
# Plugin architecture
|
||||
"ChannelPlugin",
|
||||
"ChannelMeta",
|
||||
"ChannelServer",
|
||||
"InboundConsumer",
|
||||
"InboundMessage",
|
||||
"IncomingMessage",
|
||||
"MessageBus",
|
||||
"OutboundMessage",
|
||||
"OutgoingMessage",
|
||||
"RawIncoming",
|
||||
"ReloadPolicy",
|
||||
"TypingManager",
|
||||
"UnifiedFormatter",
|
||||
"available_channels",
|
||||
"chunk_text",
|
||||
"create_channel",
|
||||
"register_channel",
|
||||
"run_standalone",
|
||||
]
|
||||
|
||||
@@ -6,6 +6,7 @@ import logging
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import ClassVar
|
||||
from urllib.parse import quote_plus
|
||||
|
||||
from ..base import Channel, ChannelError, RawIncoming
|
||||
@@ -319,7 +320,7 @@ class DingTalkChannel(Channel, WebSocketMixin, TokenMixin):
|
||||
|
||||
# ── Media send ────────────────────────────────────────────────
|
||||
|
||||
_IMAGE_EXTS = {".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp"}
|
||||
_IMAGE_EXTS: ClassVar[set[str]] = {".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp"}
|
||||
|
||||
async def _send_media_impl(
|
||||
self,
|
||||
|
||||
@@ -36,6 +36,7 @@ class DiscordChannel(Channel):
|
||||
# Cache message objects for ACK reactions
|
||||
self._message_cache: dict[str, object] = {}
|
||||
self._MESSAGE_CACHE_MAX = 200
|
||||
self._background_tasks: set[asyncio.Task] = set()
|
||||
|
||||
async def start(self) -> None:
|
||||
try:
|
||||
@@ -44,7 +45,7 @@ class DiscordChannel(Channel):
|
||||
raise ChannelError(
|
||||
"discord.py not installed. "
|
||||
"Install with: pip install evoscientist[discord]"
|
||||
)
|
||||
) from None
|
||||
|
||||
if not self.config.bot_token:
|
||||
raise ChannelError("Discord bot token is required")
|
||||
@@ -93,7 +94,9 @@ class DiscordChannel(Channel):
|
||||
self._ready.set() # unblock the waiter so it doesn't hang
|
||||
|
||||
logger.info("Discord connect: launching gateway task")
|
||||
asyncio.create_task(_guarded_start())
|
||||
_task = asyncio.create_task(_guarded_start())
|
||||
self._background_tasks.add(_task)
|
||||
_task.add_done_callback(self._background_tasks.discard)
|
||||
|
||||
try:
|
||||
await asyncio.wait_for(self._ready.wait(), timeout=60)
|
||||
@@ -101,7 +104,7 @@ class DiscordChannel(Channel):
|
||||
raise ChannelError(
|
||||
"Discord bot failed to connect within 60s. "
|
||||
"Check network/proxy connectivity to gateway.discord.gg"
|
||||
)
|
||||
) from None
|
||||
|
||||
if self._start_task_error:
|
||||
raise ChannelError(
|
||||
|
||||
@@ -116,7 +116,7 @@ class EmailChannel(Channel, PollingMixin):
|
||||
self._imap.login(cfg.imap_username, cfg.imap_password)
|
||||
self._imap.select(cfg.imap_mailbox)
|
||||
except Exception as e:
|
||||
raise ChannelError(f"IMAP failed: {e}")
|
||||
raise ChannelError(f"IMAP failed: {e}") from e
|
||||
|
||||
def _reconnect_imap(self) -> None:
|
||||
try:
|
||||
|
||||
@@ -26,7 +26,7 @@ import re
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Any
|
||||
from typing import TYPE_CHECKING, Any, ClassVar
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from aiohttp import web
|
||||
@@ -284,7 +284,7 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
resp = await self._http_client.post(url, json=body)
|
||||
data = resp.json()
|
||||
except Exception as e:
|
||||
raise ChannelError(f"Failed to get Feishu access token: {e}")
|
||||
raise ChannelError(f"Failed to get Feishu access token: {e}") from e
|
||||
|
||||
if data.get("code") != 0:
|
||||
raise ChannelError(f"Feishu auth error: {data.get('msg', 'unknown')}")
|
||||
@@ -300,7 +300,7 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
raise ChannelError(
|
||||
"aiohttp or httpx not installed. "
|
||||
"Install with: pip install aiohttp httpx"
|
||||
)
|
||||
) from None
|
||||
|
||||
if not self.config.app_id:
|
||||
raise ChannelError("Feishu app_id is required")
|
||||
@@ -394,7 +394,14 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
|
||||
# ── Media helpers ──────────────────────────────────────────────
|
||||
|
||||
_IMAGE_EXTENSIONS = {".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp"}
|
||||
_IMAGE_EXTENSIONS: ClassVar[set[str]] = {
|
||||
".jpg",
|
||||
".jpeg",
|
||||
".png",
|
||||
".gif",
|
||||
".bmp",
|
||||
".webp",
|
||||
}
|
||||
|
||||
async def _download_media(
|
||||
self,
|
||||
|
||||
@@ -13,6 +13,7 @@ from __future__ import annotations
|
||||
|
||||
import re
|
||||
from collections.abc import Callable
|
||||
from typing import ClassVar
|
||||
|
||||
# ═════════════════════════════════════════════════════════════════════
|
||||
# Markdown conversion engine (formerly markdown_utils.py)
|
||||
@@ -235,37 +236,37 @@ class UnifiedFormatter:
|
||||
Instantiated once per channel based on its ``capabilities.format_type``.
|
||||
"""
|
||||
|
||||
_PROFILES: dict[str, dict] = {
|
||||
"html": dict(
|
||||
code_block_formatter=_html_code_block,
|
||||
inline_code_formatter=_html_inline_code,
|
||||
inline_rules=_HTML_INLINE_RULES,
|
||||
escape_fn=_escape_html,
|
||||
),
|
||||
"slack_mrkdwn": dict(
|
||||
code_block_formatter=_slack_code_block,
|
||||
inline_code_formatter=_slack_inline_code,
|
||||
inline_rules=_SLACK_INLINE_RULES,
|
||||
escape_fn=None,
|
||||
),
|
||||
"discord": dict(
|
||||
code_block_formatter=_discord_code_block,
|
||||
inline_code_formatter=_discord_inline_code,
|
||||
inline_rules=_DISCORD_INLINE_RULES,
|
||||
escape_fn=None,
|
||||
),
|
||||
"markdown": dict(
|
||||
code_block_formatter=_md_code_block,
|
||||
inline_code_formatter=_md_inline_code,
|
||||
inline_rules=_MD_INLINE_RULES,
|
||||
escape_fn=None,
|
||||
),
|
||||
"plain": dict(
|
||||
code_block_formatter=_plain_code_block,
|
||||
inline_code_formatter=_plain_inline_code,
|
||||
inline_rules=_PLAIN_INLINE_RULES,
|
||||
escape_fn=None,
|
||||
),
|
||||
_PROFILES: ClassVar[dict[str, dict]] = {
|
||||
"html": {
|
||||
"code_block_formatter": _html_code_block,
|
||||
"inline_code_formatter": _html_inline_code,
|
||||
"inline_rules": _HTML_INLINE_RULES,
|
||||
"escape_fn": _escape_html,
|
||||
},
|
||||
"slack_mrkdwn": {
|
||||
"code_block_formatter": _slack_code_block,
|
||||
"inline_code_formatter": _slack_inline_code,
|
||||
"inline_rules": _SLACK_INLINE_RULES,
|
||||
"escape_fn": None,
|
||||
},
|
||||
"discord": {
|
||||
"code_block_formatter": _discord_code_block,
|
||||
"inline_code_formatter": _discord_inline_code,
|
||||
"inline_rules": _DISCORD_INLINE_RULES,
|
||||
"escape_fn": None,
|
||||
},
|
||||
"markdown": {
|
||||
"code_block_formatter": _md_code_block,
|
||||
"inline_code_formatter": _md_inline_code,
|
||||
"inline_rules": _MD_INLINE_RULES,
|
||||
"escape_fn": None,
|
||||
},
|
||||
"plain": {
|
||||
"code_block_formatter": _plain_code_block,
|
||||
"inline_code_formatter": _plain_inline_code,
|
||||
"inline_rules": _PLAIN_INLINE_RULES,
|
||||
"escape_fn": None,
|
||||
},
|
||||
}
|
||||
|
||||
def __init__(self, format_type: str = "plain") -> None:
|
||||
|
||||
@@ -71,6 +71,7 @@ class IMessageChannelRpc(Channel):
|
||||
super().__init__(config or IMessageConfig())
|
||||
self._client: ImsgRpcClient | None = None
|
||||
self._subscription_id: int | None = None
|
||||
self._background_tasks: set[asyncio.Task] = set()
|
||||
|
||||
# ── Pipeline overrides ────────────────────────────────────────
|
||||
|
||||
@@ -93,7 +94,9 @@ class IMessageChannelRpc(Channel):
|
||||
def _handle_notification(self, notification: RpcNotification) -> None:
|
||||
"""Handle incoming RPC notifications."""
|
||||
if notification.method == "message":
|
||||
asyncio.create_task(self._handle_message(notification.params))
|
||||
_task = asyncio.create_task(self._handle_message(notification.params))
|
||||
self._background_tasks.add(_task)
|
||||
_task.add_done_callback(self._background_tasks.discard)
|
||||
elif notification.method == "error":
|
||||
logger.error(f"imsg error: {notification.params}")
|
||||
|
||||
|
||||
@@ -57,6 +57,7 @@ class ImsgRpcClient:
|
||||
self._next_id = 1
|
||||
self._pending: dict[int, asyncio.Future] = {}
|
||||
self._reader_task: asyncio.Task | None = None
|
||||
self._stderr_task: asyncio.Task | None = None
|
||||
self._closed = asyncio.Event()
|
||||
|
||||
async def start(self) -> None:
|
||||
@@ -76,7 +77,7 @@ class ImsgRpcClient:
|
||||
)
|
||||
|
||||
self._reader_task = asyncio.create_task(self._read_loop())
|
||||
asyncio.create_task(self._stderr_loop())
|
||||
self._stderr_task = asyncio.create_task(self._stderr_loop())
|
||||
logger.info(f"Started imsg rpc (pid={self._process.pid})")
|
||||
|
||||
async def stop(self) -> None:
|
||||
@@ -94,6 +95,13 @@ class ImsgRpcClient:
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
if self._stderr_task:
|
||||
self._stderr_task.cancel()
|
||||
try:
|
||||
await self._stderr_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
try:
|
||||
self._process.terminate()
|
||||
await asyncio.wait_for(self._process.wait(), timeout=2.0)
|
||||
@@ -153,7 +161,7 @@ class ImsgRpcClient:
|
||||
return await asyncio.wait_for(future, timeout=timeout)
|
||||
except TimeoutError:
|
||||
self._pending.pop(request_id, None)
|
||||
raise Exception(f"RPC request timeout: {method}")
|
||||
raise Exception(f"RPC request timeout: {method}") from None
|
||||
|
||||
async def _read_loop(self) -> None:
|
||||
"""Read and process responses from stdout."""
|
||||
|
||||
@@ -7,7 +7,6 @@ phone numbers and email addresses, similar to OpenClaw's approach.
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
from enum import Enum
|
||||
from typing import Union
|
||||
|
||||
|
||||
class IMessageService(Enum):
|
||||
@@ -51,7 +50,7 @@ class HandleTarget:
|
||||
service: IMessageService = IMessageService.AUTO
|
||||
|
||||
|
||||
IMessageTarget = Union[ChatIdTarget, ChatGuidTarget, ChatIdentifierTarget, HandleTarget]
|
||||
IMessageTarget = ChatIdTarget | ChatGuidTarget | ChatIdentifierTarget | HandleTarget
|
||||
|
||||
|
||||
# Prefix constants
|
||||
@@ -202,8 +201,8 @@ def parse_target(raw: str) -> IMessageTarget:
|
||||
try:
|
||||
chat_id = int(value)
|
||||
return ChatIdTarget(chat_id=chat_id)
|
||||
except ValueError:
|
||||
raise ValueError(f"Invalid chat_id: {value}")
|
||||
except ValueError as e:
|
||||
raise ValueError(f"Invalid chat_id: {value}") from e
|
||||
|
||||
# Check chat_guid prefixes
|
||||
for prefix in CHAT_GUID_PREFIXES:
|
||||
|
||||
@@ -122,7 +122,7 @@ class DedupCache:
|
||||
cutoff = time.monotonic() - self._ttl
|
||||
# OrderedDict is insertion-ordered; oldest entries are first.
|
||||
while self._seen:
|
||||
key, ts = next(iter(self._seen.items()))
|
||||
_key, ts = next(iter(self._seen.items()))
|
||||
if ts > cutoff:
|
||||
break
|
||||
self._seen.popitem(last=False)
|
||||
|
||||
@@ -5,6 +5,7 @@ import logging
|
||||
from collections import deque
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from typing import ClassVar
|
||||
|
||||
from ..base import Channel, ChannelError, RawIncoming
|
||||
from ..capabilities import QQ as QQ_CAPS
|
||||
@@ -210,7 +211,7 @@ class QQChannel(Channel):
|
||||
# ── Media send ────────────────────────────────────────────────
|
||||
|
||||
# qq-botpy file_type constants: 1=image, 2=video, 3=audio
|
||||
_FILE_TYPE_MAP = {
|
||||
_FILE_TYPE_MAP: ClassVar[dict[str, int]] = {
|
||||
".jpg": 1,
|
||||
".jpeg": 1,
|
||||
".png": 1,
|
||||
|
||||
@@ -34,7 +34,7 @@ class RetryInfo:
|
||||
|
||||
async def retry_async(
|
||||
fn: Callable[[], Awaitable[T]],
|
||||
config: RetryConfig = RetryConfig(),
|
||||
config: RetryConfig | None = None,
|
||||
*,
|
||||
should_retry: Callable[[Exception, int], bool] | None = None,
|
||||
retry_after_s: Callable[[Exception], float | None] | None = None,
|
||||
@@ -48,20 +48,10 @@ async def retry_async(
|
||||
fn:
|
||||
Zero-argument async factory — called on every attempt so the
|
||||
awaitable is always fresh.
|
||||
config:
|
||||
Retry timing / attempt parameters.
|
||||
should_retry:
|
||||
``(exception, attempt) -> bool``. Return ``False`` to abort
|
||||
immediately. When *None* every exception is retried.
|
||||
retry_after_s:
|
||||
``(exception) -> seconds | None``. If the server provides a
|
||||
``Retry-After`` value (e.g. HTTP 429), return it here. The
|
||||
actual delay will be ``max(server_value, min_delay_s)``.
|
||||
on_retry:
|
||||
Optional callback invoked before each retry sleep.
|
||||
label:
|
||||
Human-readable label included in :class:`RetryInfo`.
|
||||
"""
|
||||
if config is None:
|
||||
config = RetryConfig()
|
||||
|
||||
last_exc: Exception | None = None
|
||||
for attempt in range(1, config.attempts + 1):
|
||||
try:
|
||||
|
||||
@@ -94,7 +94,7 @@ class SignalChannel(Channel):
|
||||
async def _ensure_daemon(self) -> None:
|
||||
"""Start signal-cli daemon if not already running."""
|
||||
try:
|
||||
reader, writer = await asyncio.wait_for(
|
||||
_reader, writer = await asyncio.wait_for(
|
||||
asyncio.open_connection("localhost", self.config.rpc_port),
|
||||
timeout=2,
|
||||
)
|
||||
@@ -129,13 +129,13 @@ class SignalChannel(Channel):
|
||||
raise ChannelError(
|
||||
f"signal-cli not found at '{self.config.cli_path}'. "
|
||||
"Install: https://github.com/AsamK/signal-cli"
|
||||
)
|
||||
) from None
|
||||
|
||||
# Wait for daemon to be ready
|
||||
for _ in range(30):
|
||||
await asyncio.sleep(1)
|
||||
try:
|
||||
reader, writer = await asyncio.open_connection(
|
||||
_reader, writer = await asyncio.open_connection(
|
||||
"localhost",
|
||||
self.config.rpc_port,
|
||||
)
|
||||
@@ -156,7 +156,7 @@ class SignalChannel(Channel):
|
||||
self.config.rpc_port,
|
||||
)
|
||||
except Exception as e:
|
||||
raise ChannelError(f"Cannot connect to signal-cli: {e}")
|
||||
raise ChannelError(f"Cannot connect to signal-cli: {e}") from e
|
||||
|
||||
async def _listen_loop(self) -> None:
|
||||
"""Listen for incoming JSON RPC notifications and responses."""
|
||||
|
||||
@@ -51,7 +51,7 @@ class SlackChannel(Channel):
|
||||
raise ChannelError(
|
||||
"slack-sdk or aiohttp not installed. "
|
||||
"Install with: pip install evoscientist[slack]"
|
||||
)
|
||||
) from None
|
||||
|
||||
self._web_client = AsyncWebClient(
|
||||
token=self.config.bot_token,
|
||||
@@ -68,9 +68,9 @@ class SlackChannel(Channel):
|
||||
except TimeoutError:
|
||||
raise ChannelError(
|
||||
"Slack auth_test timed out — check network and bot token"
|
||||
)
|
||||
) from None
|
||||
except Exception as e:
|
||||
raise ChannelError(f"Failed to authenticate Slack bot: {e}")
|
||||
raise ChannelError(f"Failed to authenticate Slack bot: {e}") from e
|
||||
|
||||
self._socket_client = SocketModeClient(
|
||||
app_token=self.config.app_token,
|
||||
@@ -115,7 +115,7 @@ class SlackChannel(Channel):
|
||||
"Slack Socket Mode connection timed out — "
|
||||
"check app token (must start with xapp-) and "
|
||||
"ensure Socket Mode is enabled in your Slack app settings"
|
||||
)
|
||||
) from None
|
||||
self._running = True
|
||||
logger.info("Slack channel started (Socket Mode)")
|
||||
|
||||
@@ -163,7 +163,7 @@ class SlackChannel(Channel):
|
||||
# ── Send (template method overrides) ──────────────────────────
|
||||
|
||||
async def _send_chunk(self, chat_id, formatted_text, raw_text, reply_to, metadata):
|
||||
kwargs = dict(channel=chat_id)
|
||||
kwargs = {"channel": chat_id}
|
||||
# Always route to thread if thread_ts is present in metadata,
|
||||
# not just for the first chunk (reply_to is only set for chunk 0).
|
||||
if metadata:
|
||||
|
||||
@@ -4,6 +4,7 @@ import logging
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import ClassVar
|
||||
|
||||
from ..base import (
|
||||
AUDIO_EXTS,
|
||||
@@ -55,9 +56,10 @@ class TelegramChannel(Channel):
|
||||
raise ChannelError(
|
||||
"python-telegram-bot not installed. "
|
||||
"Install with: pip install evoscientist[telegram]"
|
||||
)
|
||||
) from None
|
||||
|
||||
builder = ApplicationBuilder().token(self.config.bot_token)
|
||||
|
||||
if self.config.proxy:
|
||||
builder = builder.proxy(self.config.proxy).get_updates_proxy(
|
||||
self.config.proxy
|
||||
@@ -124,7 +126,7 @@ class TelegramChannel(Channel):
|
||||
|
||||
await self._send_with_format_fallback(_send, formatted_text, raw_text)
|
||||
|
||||
_MEDIA_SENDERS = {
|
||||
_MEDIA_SENDERS: ClassVar[dict] = {
|
||||
IMAGE_EXTS: ("send_photo", "photo"),
|
||||
VIDEO_EXTS: ("send_video", "video"),
|
||||
AUDIO_EXTS: ("send_audio", "audio"),
|
||||
@@ -285,7 +287,7 @@ class TelegramChannel(Channel):
|
||||
)
|
||||
)
|
||||
|
||||
_MIME_TO_EXT = {
|
||||
_MIME_TO_EXT: ClassVar[dict[str, str]] = {
|
||||
"image/jpeg": ".jpg",
|
||||
"image/png": ".png",
|
||||
"image/gif": ".gif",
|
||||
@@ -296,7 +298,7 @@ class TelegramChannel(Channel):
|
||||
"video/mp4": ".mp4",
|
||||
"video/quicktime": ".mov",
|
||||
}
|
||||
_TYPE_TO_EXT = {
|
||||
_TYPE_TO_EXT: ClassVar[dict[str, str]] = {
|
||||
"image": ".jpg",
|
||||
"voice": ".ogg",
|
||||
"audio": ".mp3",
|
||||
|
||||
@@ -127,6 +127,7 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
self._typing_message_ids: dict[
|
||||
str, list[str]
|
||||
] = {} # chat_id → [msgid, ...] for typing recall
|
||||
self._background_tasks: set[asyncio.Task] = set()
|
||||
|
||||
# ── Lifecycle ─────────────────────────────────────────────────
|
||||
|
||||
@@ -145,7 +146,7 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
raise ChannelError(
|
||||
"aiohttp or httpx not installed. "
|
||||
"Install with: pip install aiohttp httpx"
|
||||
)
|
||||
) from None
|
||||
|
||||
self._validate_config()
|
||||
|
||||
@@ -248,8 +249,8 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
data = resp.json()
|
||||
except Exception as e:
|
||||
if not self._running:
|
||||
raise ChannelError(f"Failed to get WeChat access token: {e}")
|
||||
raise RuntimeError(f"Failed to get WeChat access token: {e}")
|
||||
raise ChannelError(f"Failed to get WeChat access token: {e}") from e
|
||||
raise RuntimeError(f"Failed to get WeChat access token: {e}") from e
|
||||
|
||||
if data.get("errcode", 0) != 0:
|
||||
err_msg = (
|
||||
@@ -348,7 +349,7 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
return web.Response(status=403)
|
||||
|
||||
try:
|
||||
decrypted_xml, from_id = self._crypto.decrypt(encrypt)
|
||||
decrypted_xml, _from_id = self._crypto.decrypt(encrypt)
|
||||
xml_data = parse_xml(decrypted_xml)
|
||||
except Exception as e:
|
||||
logger.error(f"WeChat decrypt failed: {e}")
|
||||
@@ -357,7 +358,9 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
# Process message asynchronously — WeCom requires a response within
|
||||
# 5 seconds, but media downloads can take much longer. Return
|
||||
# "success" immediately and handle the message in the background.
|
||||
asyncio.create_task(self._safe_process_message(xml_data))
|
||||
_task = asyncio.create_task(self._safe_process_message(xml_data))
|
||||
self._background_tasks.add(_task)
|
||||
_task.add_done_callback(self._background_tasks.discard)
|
||||
|
||||
return web.Response(text="success")
|
||||
|
||||
@@ -812,7 +815,7 @@ class WeChatChannel(Channel, WebhookMixin, TokenMixin):
|
||||
resp = await self._http_client.post(url, json=body)
|
||||
data = resp.json()
|
||||
except Exception as e:
|
||||
raise RuntimeError(f"WeChat API error: {e}")
|
||||
raise RuntimeError(f"WeChat API error: {e}") from e
|
||||
|
||||
errcode = data.get("errcode", 0)
|
||||
if errcode != 0:
|
||||
|
||||
@@ -60,7 +60,7 @@ def _aes_decrypt(key: bytes, iv: bytes, ciphertext: bytes) -> bytes:
|
||||
raise ImportError(
|
||||
"WeChat message decryption requires pycryptodome or pyaes. "
|
||||
"Install with: pip install pycryptodome"
|
||||
)
|
||||
) from None
|
||||
|
||||
|
||||
def _aes_encrypt(key: bytes, iv: bytes, plaintext: bytes) -> bytes:
|
||||
@@ -80,7 +80,7 @@ def _aes_encrypt(key: bytes, iv: bytes, plaintext: bytes) -> bytes:
|
||||
raise ImportError(
|
||||
"WeChat message encryption requires pycryptodome or pyaes. "
|
||||
"Install with: pip install pycryptodome"
|
||||
)
|
||||
) from None
|
||||
|
||||
|
||||
class WeChatCrypto:
|
||||
|
||||
@@ -495,6 +495,7 @@ async def _bus_inbound_consumer(bus, manager) -> None:
|
||||
Task-based: each inbound message is handled in its own asyncio task
|
||||
so the consumer loop stays responsive for HITL approval replies.
|
||||
"""
|
||||
_tasks: set[asyncio.Task] = set()
|
||||
while True:
|
||||
try:
|
||||
msg = await asyncio.wait_for(bus.consume_inbound(), timeout=1.0)
|
||||
@@ -512,7 +513,9 @@ async def _bus_inbound_consumer(bus, manager) -> None:
|
||||
continue
|
||||
|
||||
# Regular message — handle in a separate task
|
||||
asyncio.create_task(_handle_bus_message(bus, manager, msg))
|
||||
_task = asyncio.create_task(_handle_bus_message(bus, manager, msg))
|
||||
_tasks.add(_task)
|
||||
_task.add_done_callback(_tasks.discard)
|
||||
|
||||
|
||||
async def _handle_bus_message(bus, manager, msg) -> None:
|
||||
|
||||
@@ -7,7 +7,7 @@ import re
|
||||
from datetime import datetime
|
||||
from importlib.metadata import version as _pkg_version
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from typing import Annotated, Any
|
||||
|
||||
import typer # type: ignore[import-untyped]
|
||||
from rich.markup import escape
|
||||
@@ -521,7 +521,7 @@ def serve(
|
||||
atexit.register(stop_ccproxy, _ccproxy_proc_serve)
|
||||
except RuntimeError as exc:
|
||||
console.print(f"[red]{exc}[/red]")
|
||||
raise typer.Exit(1)
|
||||
raise typer.Exit(1) from exc
|
||||
|
||||
if not config.channel_enabled:
|
||||
console.print("[red]No channels configured.[/red]")
|
||||
@@ -727,30 +727,41 @@ def mcp_config(
|
||||
|
||||
@mcp_app.command("add")
|
||||
def mcp_add(
|
||||
name: str = typer.Argument(..., help="Server name"),
|
||||
target: str = typer.Argument(..., help="Command (stdio) or URL (http/sse)"),
|
||||
args: list[str] | None = typer.Argument(None, help="Extra args for stdio command"),
|
||||
transport: str | None = typer.Option(
|
||||
None, "--transport", "-T", help="Transport type (default: auto-detect)"
|
||||
),
|
||||
tools: str | None = typer.Option(
|
||||
None,
|
||||
"--tools",
|
||||
"-t",
|
||||
help="Comma-separated tool allowlist (supports wildcards: *_exa, read_*)",
|
||||
),
|
||||
expose_to: str | None = typer.Option(
|
||||
None, "--expose-to", "-e", help="Comma-separated target agents"
|
||||
),
|
||||
header: list[str] | None = typer.Option(
|
||||
None, "--header", "-H", help="HTTP header as Key:Value (repeatable)"
|
||||
),
|
||||
env: list[str] | None = typer.Option(
|
||||
None, "--env", help="Env var as KEY=VALUE for stdio (repeatable)"
|
||||
),
|
||||
env_ref: list[str] | None = typer.Option(
|
||||
None, "--env-ref", help="Env var name as ${NAME} runtime ref (repeatable)"
|
||||
),
|
||||
name: Annotated[str, typer.Argument(help="Server name")],
|
||||
target: Annotated[str, typer.Argument(help="Command (stdio) or URL (http/sse)")],
|
||||
args: Annotated[
|
||||
list[str] | None, typer.Argument(help="Extra args for stdio command")
|
||||
] = None,
|
||||
transport: Annotated[
|
||||
str | None,
|
||||
typer.Option("--transport", "-T", help="Transport type (default: auto-detect)"),
|
||||
] = None,
|
||||
tools: Annotated[
|
||||
str | None,
|
||||
typer.Option(
|
||||
"--tools",
|
||||
"-t",
|
||||
help="Comma-separated tool allowlist (supports wildcards: *_exa, read_*)",
|
||||
),
|
||||
] = None,
|
||||
expose_to: Annotated[
|
||||
str | None,
|
||||
typer.Option("--expose-to", "-e", help="Comma-separated target agents"),
|
||||
] = None,
|
||||
header: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option("--header", "-H", help="HTTP header as Key:Value (repeatable)"),
|
||||
] = None,
|
||||
env: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option("--env", help="Env var as KEY=VALUE for stdio (repeatable)"),
|
||||
] = None,
|
||||
env_ref: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option(
|
||||
"--env-ref", help="Env var name as ${NAME} runtime ref (repeatable)"
|
||||
),
|
||||
] = None,
|
||||
):
|
||||
"""Add an MCP server to user config
|
||||
|
||||
@@ -800,30 +811,40 @@ def mcp_add(
|
||||
|
||||
@mcp_app.command("edit")
|
||||
def mcp_edit(
|
||||
name: str = typer.Argument(..., help="Server name to edit"),
|
||||
transport: str | None = typer.Option(
|
||||
None, "--transport", help="New transport type"
|
||||
),
|
||||
command: str | None = typer.Option(None, "--command", help="New command (stdio)"),
|
||||
url: str | None = typer.Option(None, "--url", help="New URL (http/sse/websocket)"),
|
||||
tools: str | None = typer.Option(
|
||||
None,
|
||||
"--tools",
|
||||
"-t",
|
||||
help="Comma-separated tool allowlist, supports wildcards ('none' to clear)",
|
||||
),
|
||||
expose_to: str | None = typer.Option(
|
||||
None,
|
||||
"--expose-to",
|
||||
"-e",
|
||||
help="Comma-separated target agents ('none' to clear)",
|
||||
),
|
||||
header: list[str] | None = typer.Option(
|
||||
None, "--header", "-H", help="HTTP header as Key:Value (repeatable)"
|
||||
),
|
||||
env: list[str] | None = typer.Option(
|
||||
None, "--env", help="Env var as KEY=VALUE for stdio (repeatable)"
|
||||
),
|
||||
name: Annotated[str, typer.Argument(help="Server name to edit")],
|
||||
transport: Annotated[
|
||||
str | None, typer.Option("--transport", help="New transport type")
|
||||
] = None,
|
||||
command: Annotated[
|
||||
str | None, typer.Option("--command", help="New command (stdio)")
|
||||
] = None,
|
||||
url: Annotated[
|
||||
str | None, typer.Option("--url", help="New URL (http/sse/websocket)")
|
||||
] = None,
|
||||
tools: Annotated[
|
||||
str | None,
|
||||
typer.Option(
|
||||
"--tools",
|
||||
"-t",
|
||||
help="Comma-separated tool allowlist, supports wildcards ('none' to clear)",
|
||||
),
|
||||
] = None,
|
||||
expose_to: Annotated[
|
||||
str | None,
|
||||
typer.Option(
|
||||
"--expose-to",
|
||||
"-e",
|
||||
help="Comma-separated target agents ('none' to clear)",
|
||||
),
|
||||
] = None,
|
||||
header: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option("--header", "-H", help="HTTP header as Key:Value (repeatable)"),
|
||||
] = None,
|
||||
env: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option("--env", help="Env var as KEY=VALUE for stdio (repeatable)"),
|
||||
] = None,
|
||||
):
|
||||
"""Edit an existing MCP server in user config
|
||||
|
||||
@@ -861,7 +882,9 @@ def mcp_remove(
|
||||
|
||||
@mcp_app.command("install")
|
||||
def mcp_install(
|
||||
source: Optional[str] = typer.Argument(None, help="Server name or tag filter"),
|
||||
source: Annotated[
|
||||
str | None, typer.Argument(help="Server name or tag filter")
|
||||
] = None,
|
||||
):
|
||||
"""Browse and install MCP servers from the registry and marketplace
|
||||
|
||||
@@ -990,7 +1013,7 @@ def _main_callback(
|
||||
atexit.register(stop_ccproxy, _ccproxy_proc)
|
||||
except RuntimeError as exc:
|
||||
console.print(f"[red]{exc}[/red]")
|
||||
raise typer.Exit(1)
|
||||
raise typer.Exit(1) from exc
|
||||
|
||||
show_thinking = config.show_thinking if not no_thinking else False
|
||||
effective_channel_thinking = config.channel_send_thinking and (not no_thinking)
|
||||
|
||||
@@ -77,7 +77,7 @@ def print_banner(
|
||||
ui_backend: str | None = None,
|
||||
):
|
||||
"""Print welcome banner with ASCII art logo, info line, and hint."""
|
||||
for line, color in zip(LOGO_LINES, LOGO_GRADIENT):
|
||||
for line, color in zip(LOGO_LINES, LOGO_GRADIENT, strict=False):
|
||||
console.print(Text(line, style=f"{color} bold"))
|
||||
info = Text()
|
||||
info.append(" ", style="dim")
|
||||
@@ -401,7 +401,7 @@ def cmd_interactive(
|
||||
|
||||
lefts = [t.get("preview", "") or t["thread_id"] for t in threads]
|
||||
col_width = max(_display_width(s) for s in lefts) + 4
|
||||
for t, left_text in zip(threads, lefts):
|
||||
for t, left_text in zip(threads, lefts, strict=False):
|
||||
tid = t["thread_id"]
|
||||
when = _format_relative_time(t.get("updated_at"))
|
||||
label = f"{_pad_to_width(left_text, col_width)}({tid} {when})"
|
||||
@@ -912,7 +912,7 @@ def cmd_run(
|
||||
console.print(
|
||||
"[dim]Run [bold]EvoSci onboard[/bold] to set up your API key.[/dim]"
|
||||
)
|
||||
raise typer.Exit(1)
|
||||
raise typer.Exit(1) from e
|
||||
else:
|
||||
console.print(f"[red]Error: {e}[/red]")
|
||||
raise
|
||||
|
||||
@@ -11,7 +11,7 @@ import logging
|
||||
import queue
|
||||
import random
|
||||
from collections.abc import Callable
|
||||
from typing import Any
|
||||
from typing import Any, ClassVar
|
||||
|
||||
from rich.console import Group
|
||||
from rich.text import Text
|
||||
@@ -71,7 +71,7 @@ def _build_welcome_banner(
|
||||
channels: List of (name, ok, detail) tuples for the channels panel.
|
||||
"""
|
||||
banner = Text()
|
||||
for line, color in zip(LOGO_LINES, LOGO_GRADIENT):
|
||||
for line, color in zip(LOGO_LINES, LOGO_GRADIENT, strict=False):
|
||||
banner.append(f"{line}\n", style=f"bold {color}")
|
||||
|
||||
# Info line — matches CLI print_banner format
|
||||
@@ -260,14 +260,14 @@ def run_textual_interactive(
|
||||
padding: 0 1;
|
||||
border-bottom: solid #0284c7;
|
||||
}
|
||||
#status {
|
||||
# status {
|
||||
height: 1;
|
||||
background: #171a20;
|
||||
color: #f59e0b;
|
||||
padding: 0 1;
|
||||
}
|
||||
"""
|
||||
BINDINGS = [
|
||||
BINDINGS: ClassVar[list[Binding]] = [
|
||||
Binding("ctrl+c", "request_quit", "Quit", show=False),
|
||||
Binding("ctrl+v", "paste_clipboard", "Paste", show=False),
|
||||
Binding("up", "edit_queued", show=False, priority=True),
|
||||
@@ -310,6 +310,7 @@ def run_textual_interactive(
|
||||
self._browser_future: asyncio.Future | None = None
|
||||
self._mcp_browser_future: asyncio.Future | None = None
|
||||
self._history_suggester = HistorySuggester(get_config_dir() / "history")
|
||||
self._background_tasks: set[asyncio.Task] = set()
|
||||
|
||||
# ── CommandUI implementation ─────────────────────────
|
||||
|
||||
@@ -1084,7 +1085,7 @@ def run_textual_interactive(
|
||||
)
|
||||
_ask_fn = channel_ask_user_fn
|
||||
result = await asyncio.to_thread(
|
||||
lambda: _ask_fn(event),
|
||||
lambda f=_ask_fn, e=event: f(e),
|
||||
)
|
||||
else:
|
||||
# Interactive TUI: display widget, collect via arrow keys
|
||||
@@ -1531,7 +1532,9 @@ def run_textual_interactive(
|
||||
# Launch as independent task to free the message pump.
|
||||
# Commands like /resume mount interactive widgets that need
|
||||
# the pump to process key events and message bubbling.
|
||||
asyncio.ensure_future(self._handle_command(text))
|
||||
_task = asyncio.create_task(self._handle_command(text))
|
||||
self._background_tasks.add(_task)
|
||||
_task.add_done_callback(self._background_tasks.discard)
|
||||
return
|
||||
|
||||
self._history_suggester.append_entry(text)
|
||||
|
||||
@@ -122,7 +122,7 @@ class MCPBrowserWidget(Widget):
|
||||
for t in s.tags:
|
||||
tag_counter[t.lower()] += 1
|
||||
sorted_tags = sorted(tag_counter.items(), key=lambda x: (-x[1], x[0]))
|
||||
self._tag_items = [("all", len(self._servers))] + sorted_tags
|
||||
self._tag_items = [("all", len(self._servers)), *sorted_tags]
|
||||
|
||||
# If pre-filtered, skip to phase 2
|
||||
if self._pre_filter_tag:
|
||||
|
||||
@@ -120,7 +120,7 @@ class SkillBrowserWidget(Widget):
|
||||
for t in s.get("tags", []):
|
||||
tag_counter[t.lower()] += 1
|
||||
sorted_tags = sorted(tag_counter.items(), key=lambda x: (-x[1], x[0]))
|
||||
self._tag_items = [("all", len(self._index))] + sorted_tags
|
||||
self._tag_items = [("all", len(self._index)), *sorted_tags]
|
||||
|
||||
# If pre-filtered, skip to phase 2
|
||||
if self._pre_filter_tag:
|
||||
|
||||
@@ -163,7 +163,9 @@ class ThreadPickerWidget(Widget):
|
||||
self.call_later(self.focus)
|
||||
|
||||
def _update_rows(self) -> None:
|
||||
for i, (thread, widget) in enumerate(zip(self._threads, self._row_widgets)):
|
||||
for i, (thread, widget) in enumerate(
|
||||
zip(self._threads, self._row_widgets, strict=False)
|
||||
):
|
||||
is_current = thread["thread_id"] == self._current_thread
|
||||
text = build_row_text(
|
||||
thread, selected=(i == self._selected), current=is_current
|
||||
|
||||
@@ -2,7 +2,7 @@ from __future__ import annotations
|
||||
|
||||
from abc import ABC, abstractmethod
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Protocol, runtime_checkable
|
||||
from typing import Any, ClassVar, Protocol, runtime_checkable
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -61,9 +61,9 @@ class Command(ABC):
|
||||
"""Base class for all EvoScientist slash commands."""
|
||||
|
||||
name: str
|
||||
alias: list[str]
|
||||
alias: ClassVar[list[str]] = []
|
||||
description: str
|
||||
arguments: list[Argument]
|
||||
arguments: ClassVar[list[Argument]] = []
|
||||
|
||||
@abstractmethod
|
||||
async def execute(self, ctx: CommandContext, args: list[str]) -> None:
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import ClassVar
|
||||
|
||||
from ..base import Argument, Command, CommandContext
|
||||
|
||||
|
||||
@@ -8,7 +10,7 @@ class InstallMCPCommand(Command):
|
||||
|
||||
name = "/install-mcp"
|
||||
description = "Browse and install MCP servers"
|
||||
arguments = [
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="source",
|
||||
type=str,
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import ClassVar
|
||||
|
||||
from rich.table import Table
|
||||
|
||||
from ..base import Argument, Command, CommandContext
|
||||
@@ -77,7 +79,7 @@ class ResumeCommand(Command):
|
||||
|
||||
name = "/resume"
|
||||
description = "Resume a previous session"
|
||||
arguments = [
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="thread_id",
|
||||
type=str,
|
||||
@@ -178,7 +180,7 @@ class DeleteCommand(Command):
|
||||
|
||||
name = "/delete"
|
||||
description = "Delete a saved session"
|
||||
arguments = [
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="thread_id",
|
||||
type=str,
|
||||
@@ -255,7 +257,7 @@ class ExitCommand(Command):
|
||||
"""Quit EvoScientist."""
|
||||
|
||||
name = "/exit"
|
||||
alias = ["/quit", "/q"]
|
||||
alias: ClassVar[list[str]] = ["/quit", "/q"]
|
||||
description = "Quit EvoScientist"
|
||||
|
||||
async def execute(self, ctx: CommandContext, args: list[str]) -> None:
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import ClassVar
|
||||
|
||||
from rich.table import Table
|
||||
|
||||
from ..base import Argument, Command, CommandContext
|
||||
@@ -60,12 +62,12 @@ class SkillsCommand(Command):
|
||||
)
|
||||
|
||||
|
||||
class InstallSkillCommand(Command):
|
||||
class InstallSkill(Command):
|
||||
"""Add a skill from path or GitHub."""
|
||||
|
||||
name = "/install-skill"
|
||||
description = "Add a skill from path or GitHub"
|
||||
arguments = [
|
||||
name: ClassVar[str] = "/install-skill"
|
||||
description: ClassVar[str] = "Add a skill from path or GitHub"
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="source",
|
||||
type=str,
|
||||
@@ -106,12 +108,14 @@ class InstallSkillCommand(Command):
|
||||
ctx.ui.append_system(f"Failed: {result['error']}", style="red")
|
||||
|
||||
|
||||
class InstallSkillsCommand(Command):
|
||||
class InstallSkills(Command):
|
||||
"""Browse and install skills."""
|
||||
|
||||
name = "/install-skills"
|
||||
description = "Browse and install skills (optional: /install-skills <tag>)"
|
||||
arguments = [
|
||||
name: ClassVar[str] = "/install-skills"
|
||||
description: ClassVar[str] = (
|
||||
"Browse and install skills (optional: /install-skills <tag>)"
|
||||
)
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="tag", type=str, description="Tag to filter skills by", required=False
|
||||
)
|
||||
@@ -208,12 +212,12 @@ class InstallSkillsCommand(Command):
|
||||
ctx.ui.append_system("No skills were installed.", style="yellow")
|
||||
|
||||
|
||||
class UninstallSkillCommand(Command):
|
||||
class UninstallSkill(Command):
|
||||
"""Remove an installed skill."""
|
||||
|
||||
name = "/uninstall-skill"
|
||||
description = "Remove an installed skill"
|
||||
arguments = [
|
||||
name: ClassVar[str] = "/uninstall-skill"
|
||||
description: ClassVar[str] = "Remove an installed skill"
|
||||
arguments: ClassVar[list[Argument]] = [
|
||||
Argument(
|
||||
name="name",
|
||||
type=str,
|
||||
@@ -241,6 +245,6 @@ class UninstallSkillCommand(Command):
|
||||
|
||||
# Register skill commands
|
||||
manager.register(SkillsCommand())
|
||||
manager.register(InstallSkillCommand())
|
||||
manager.register(InstallSkillsCommand())
|
||||
manager.register(UninstallSkillCommand())
|
||||
manager.register(InstallSkill())
|
||||
manager.register(InstallSkills())
|
||||
manager.register(UninstallSkill())
|
||||
|
||||
@@ -16,7 +16,7 @@ class CommandManager:
|
||||
|
||||
def register(self, command: Command) -> None:
|
||||
"""Register a command and its aliases."""
|
||||
names = [command.name] + command.alias
|
||||
names = [command.name, *command.alias]
|
||||
for name in names:
|
||||
name = name.lower()
|
||||
if not name.startswith("/"):
|
||||
|
||||
@@ -23,20 +23,20 @@ from .settings import (
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"EvoScientistConfig",
|
||||
"apply_config_to_env",
|
||||
# settings
|
||||
"get_config_dir",
|
||||
"get_config_path",
|
||||
"EvoScientistConfig",
|
||||
"load_config",
|
||||
"save_config",
|
||||
"reset_config",
|
||||
"get_config_value",
|
||||
"set_config_value",
|
||||
"list_config",
|
||||
"get_effective_config",
|
||||
"apply_config_to_env",
|
||||
"list_config",
|
||||
"load_config",
|
||||
"reset_config",
|
||||
# onboard (lazy)
|
||||
"run_onboard",
|
||||
"save_config",
|
||||
"set_config_value",
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -135,8 +135,8 @@ class IntegerValidator(Validator):
|
||||
raise ValidationError(
|
||||
message=f"Must be between {self.min_value} and {self.max_value}"
|
||||
)
|
||||
except ValueError:
|
||||
raise ValidationError(message="Must be a valid integer")
|
||||
except ValueError as e:
|
||||
raise ValidationError(message="Must be a valid integer") from e
|
||||
|
||||
|
||||
class ChoiceValidator(Validator):
|
||||
@@ -224,9 +224,9 @@ def validate_nvidia_key(api_key: str) -> tuple[bool, str]:
|
||||
try:
|
||||
from langchain_nvidia_ai_endpoints import ChatNVIDIA
|
||||
|
||||
llm = ChatNVIDIA(api_key=api_key, model="meta/llama-3.1-8b-instruct")
|
||||
llm.available_models
|
||||
ChatNVIDIA(api_key=api_key, model="meta/llama-3.1-8b-instruct")
|
||||
return True, "Valid"
|
||||
|
||||
except Exception as e:
|
||||
error_str = str(e).lower()
|
||||
if (
|
||||
@@ -2018,7 +2018,7 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]:
|
||||
qmark=f" {QMARK}",
|
||||
).ask()
|
||||
if install_now is None:
|
||||
raise KeyboardInterrupt()
|
||||
raise KeyboardInterrupt() from None
|
||||
if install_now:
|
||||
console.print(f" [dim]Installing {_pkg_display}...[/dim]")
|
||||
if _pip_pkgs:
|
||||
|
||||
@@ -42,7 +42,9 @@ def _patch_anthropic_proxy_compat() -> None:
|
||||
if isinstance(val, dict):
|
||||
d = val.copy()
|
||||
setattr(
|
||||
obj, attr, _types.SimpleNamespace(model_dump=lambda **kw: d)
|
||||
obj,
|
||||
attr,
|
||||
_types.SimpleNamespace(model_dump=lambda d=d, **kw: d),
|
||||
)
|
||||
return _orig(self, event, *args, **kwargs)
|
||||
|
||||
|
||||
@@ -28,17 +28,21 @@ from .registry import (
|
||||
|
||||
__all__ = [
|
||||
"VALID_TRANSPORTS",
|
||||
"MCPServerEntry",
|
||||
"add_mcp_server",
|
||||
"aload_mcp_tools",
|
||||
"build_mcp_add_kwargs",
|
||||
"build_mcp_edit_fields",
|
||||
"VALID_TRANSPORTS",
|
||||
# Registry
|
||||
"MCPServerEntry",
|
||||
"edit_mcp_server",
|
||||
"fetch_marketplace_index",
|
||||
"find_server_by_name",
|
||||
"get_all_tags",
|
||||
"get_installed_names",
|
||||
"install_mcp_server",
|
||||
"install_mcp_servers",
|
||||
"load_mcp_config",
|
||||
"load_mcp_tools",
|
||||
"parse_mcp_add_args",
|
||||
"parse_mcp_edit_args",
|
||||
"remove_mcp_server",
|
||||
]
|
||||
|
||||
@@ -638,7 +638,7 @@ async def _load_tools(config: dict[str, Any]) -> dict[str, list]:
|
||||
raise ImportError(
|
||||
"MCP servers are configured but langchain-mcp-adapters is not installed.\n"
|
||||
"Install with: pip install langchain-mcp-adapters"
|
||||
)
|
||||
) from None
|
||||
|
||||
connections = _build_connections(config)
|
||||
if not connections:
|
||||
|
||||
@@ -115,8 +115,10 @@ def _clone_repo(repo: str, ref: str | None, dest: str) -> None:
|
||||
result = subprocess.run(
|
||||
cmd, capture_output=True, text=True, timeout=_CLONE_TIMEOUT
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
raise RuntimeError(f"git clone timed out after {_CLONE_TIMEOUT}s for {repo}")
|
||||
except subprocess.TimeoutExpired as e:
|
||||
raise RuntimeError(
|
||||
f"git clone timed out after {_CLONE_TIMEOUT}s for {repo}"
|
||||
) from e
|
||||
if result.returncode != 0:
|
||||
raise RuntimeError(f"git clone failed: {result.stderr.strip()}")
|
||||
|
||||
|
||||
@@ -332,7 +332,7 @@ def _merge_memory(existing_md: str, extracted: dict[str, Any]) -> str:
|
||||
if value and value != "null":
|
||||
# Replace the line "- **Label**: ..." with new value
|
||||
pattern = rf"(- \*\*{label}\*\*: ).*"
|
||||
result = re.sub(pattern, lambda m: m.group(1) + value, result)
|
||||
result = re.sub(pattern, lambda m, v=value: m.group(1) + v, result)
|
||||
|
||||
# --- Research Preferences ---
|
||||
prefs = extracted.get("research_preferences")
|
||||
@@ -349,7 +349,7 @@ def _merge_memory(existing_md: str, extracted: dict[str, Any]) -> str:
|
||||
value = prefs.get(key)
|
||||
if value and value != "null":
|
||||
pattern = rf"(- \*\*{label}\*\*: ).*"
|
||||
result = re.sub(pattern, lambda m: m.group(1) + value, result)
|
||||
result = re.sub(pattern, lambda m, v=value: m.group(1) + v, result)
|
||||
|
||||
# --- Experiment History (append) ---
|
||||
exp = extracted.get("experiment_conclusion")
|
||||
|
||||
@@ -41,44 +41,44 @@ from .utils import (
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
# Emitter
|
||||
"StreamEventEmitter",
|
||||
"StreamEvent",
|
||||
# Tracker
|
||||
"ToolCallTracker",
|
||||
"ToolCallInfo",
|
||||
# Formatter
|
||||
"ToolResultFormatter",
|
||||
"ContentType",
|
||||
"FormattedResult",
|
||||
"FAILURE_PREFIX",
|
||||
# Utils
|
||||
"SUCCESS_PREFIX",
|
||||
"FAILURE_PREFIX",
|
||||
"ToolStatus",
|
||||
"ContentType",
|
||||
"DisplayLimits",
|
||||
"has_args",
|
||||
"is_success",
|
||||
"truncate",
|
||||
"format_tool_compact",
|
||||
"format_tree_output",
|
||||
"count_lines",
|
||||
"truncate_with_line_hint",
|
||||
"get_status_symbol",
|
||||
"FormattedResult",
|
||||
"StreamEvent",
|
||||
# Emitter
|
||||
"StreamEventEmitter",
|
||||
"StreamState",
|
||||
# State
|
||||
"SubAgentState",
|
||||
"StreamState",
|
||||
"_parse_todo_items",
|
||||
"ToolCallInfo",
|
||||
# Tracker
|
||||
"ToolCallTracker",
|
||||
# Formatter
|
||||
"ToolResultFormatter",
|
||||
"ToolStatus",
|
||||
"_astream_to_console",
|
||||
"_build_todo_stats",
|
||||
# Events
|
||||
"stream_agent_events",
|
||||
"_parse_todo_items",
|
||||
# Diff formatting
|
||||
"build_edit_diff",
|
||||
"format_diff_rich",
|
||||
# Display
|
||||
"console",
|
||||
"formatter",
|
||||
"format_tool_result_compact",
|
||||
"count_lines",
|
||||
"create_streaming_display",
|
||||
"display_final_results",
|
||||
"_astream_to_console",
|
||||
"format_diff_rich",
|
||||
"format_tool_compact",
|
||||
"format_tool_result_compact",
|
||||
"format_tree_output",
|
||||
"formatter",
|
||||
"get_status_symbol",
|
||||
"has_args",
|
||||
"is_success",
|
||||
# Events
|
||||
"stream_agent_events",
|
||||
"truncate",
|
||||
"truncate_with_line_hint",
|
||||
]
|
||||
|
||||
@@ -214,7 +214,7 @@ async def stream_agent_events(
|
||||
"""Stable key for tracker/mapping per sub-agent namespace."""
|
||||
if not namespace:
|
||||
return None
|
||||
task_id, task_ns = _extract_task_id(namespace)
|
||||
_task_id, task_ns = _extract_task_id(namespace)
|
||||
if task_ns:
|
||||
return task_ns
|
||||
meta_task_id = _find_task_id_from_metadata(metadata)
|
||||
|
||||
@@ -6,7 +6,7 @@ adapted for deepagents tool names.
|
||||
"""
|
||||
|
||||
import sys
|
||||
from enum import Enum
|
||||
from enum import StrEnum
|
||||
from pathlib import PurePath
|
||||
|
||||
# === Status marker constants ===
|
||||
@@ -15,7 +15,7 @@ FAILURE_PREFIX = "[FAILED]"
|
||||
|
||||
|
||||
# === Tool status indicators ===
|
||||
class ToolStatus(str, Enum):
|
||||
class ToolStatus(StrEnum):
|
||||
"""Tool execution status indicators."""
|
||||
|
||||
RUNNING = "\u25cf" # Running - yellow
|
||||
|
||||
@@ -93,7 +93,7 @@ async def tavily_search(
|
||||
|
||||
# Format results
|
||||
result_texts = []
|
||||
for result, content in zip(results, contents):
|
||||
for result, content in zip(results, contents, strict=False):
|
||||
result_text = f"""## {result["title"]}
|
||||
**URL:** {result["url"]}
|
||||
|
||||
|
||||
@@ -172,8 +172,10 @@ def _clone_repo(repo: str, ref: str | None, dest: str) -> None:
|
||||
result = subprocess.run(
|
||||
cmd, capture_output=True, text=True, timeout=_CLONE_TIMEOUT
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
raise RuntimeError(f"git clone timed out after {_CLONE_TIMEOUT}s for {repo}")
|
||||
except subprocess.TimeoutExpired as e:
|
||||
raise RuntimeError(
|
||||
f"git clone timed out after {_CLONE_TIMEOUT}s for {repo}"
|
||||
) from e
|
||||
if result.returncode != 0:
|
||||
raise RuntimeError(f"git clone failed: {result.stderr.strip()}")
|
||||
|
||||
|
||||
Reference in New Issue
Block a user