From 3151198ebdecab972aa349a508f7e3f151aaa44a Mon Sep 17 00:00:00 2001 From: Pine Date: Fri, 4 Sep 2026 20:40:19 +0800 Subject: [PATCH] =?UTF-8?q?=E8=BF=81=E7=A7=BB=E4=B8=8A=E6=B8=B8=E5=AD=90?= =?UTF-8?q?=E6=99=BA=E8=83=BD=E4=BD=93=E5=8A=9F=E8=83=BD=E7=BA=BF=EF=BC=9A?= =?UTF-8?q?ChatGroup=20=E5=88=86=E7=BB=84=E4=BD=93=E7=B3=BB=20+=20subagent?= =?UTF-8?q?=5Fmodel=20=E9=85=8D=E7=BD=AE=20+=20spawn=20=E4=BC=9A=E8=AF=9D?= =?UTF-8?q?=E5=BD=92=E5=B1=9E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - chats/models.py: SessionSource.subagent + ChatGroupKind/ChatGroup/default_chat_groups + ChatSpec 扩展 group_id/parent_session_id/root_session_id/last_finished_at + ChatsFile.groups - chats/manager.py: group CRUD 方法(list/create/update/reorder/delete) + get_or_create_chat/create_chat 的 group 归属与校验 - chats/api.py: /chats/groups 端点(GET/POST/PUT order/PUT/DELETE) - config.py+agents.py: AgentConfig.subagent_model 配置 - agent_management.py: _build_subagent_request_context 补 subagent_model→model_slot_override - console.py: _chat_registration_fields 识别 spawn_subagent 请求并注册 subagent 会话 - constant.py: DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS - 前端 chat.ts: groups API 请求函数 --- console/src/api/modules/chat.ts | 27 +++ .../agents/tools/agent_management.py | 19 +- src/pineagents/app/chats/api.py | 74 ++++++++ src/pineagents/app/chats/manager.py | 166 +++++++++++++++++- src/pineagents/app/chats/models.py | 136 +++++++++++++- src/pineagents/app/routers/agents.py | 2 + src/pineagents/app/routers/console.py | 22 +++ src/pineagents/config/config.py | 4 + src/pineagents/constant.py | 3 + 9 files changed, 447 insertions(+), 6 deletions(-) diff --git a/console/src/api/modules/chat.ts b/console/src/api/modules/chat.ts index e22ee14..970d2d6 100644 --- a/console/src/api/modules/chat.ts +++ b/console/src/api/modules/chat.ts @@ -7,6 +7,7 @@ import type { ChatDeleteResponse, ChatUpdateRequest, BatchArchiveResult, + ChatGroup, Session, } from "../types"; @@ -126,6 +127,32 @@ export const chatApi = { request(`/console/chat/stop?chat_id=${encodeURIComponent(chatId)}`, { method: "POST", }), + + listGroups: () => request("/chats/groups"), + + createGroup: (name: string) => + request("/chats/groups", { + method: "POST", + body: JSON.stringify({ name }), + }), + + updateGroup: (groupId: string, update: { name?: string; pinned?: boolean }) => + request(`/chats/groups/${encodeURIComponent(groupId)}`, { + method: "PUT", + body: JSON.stringify(update), + }), + + reorderGroups: (groupIds: string[]) => + request("/chats/groups/order", { + method: "PUT", + body: JSON.stringify({ group_ids: groupIds }), + }), + + deleteGroup: (groupId: string) => + request<{ success: boolean; group_id: string }>( + `/chats/groups/${encodeURIComponent(groupId)}`, + { method: "DELETE" }, + ), }; export const sessionApi = { diff --git a/src/pineagents/agents/tools/agent_management.py b/src/pineagents/agents/tools/agent_management.py index 3871434..825fff6 100644 --- a/src/pineagents/agents/tools/agent_management.py +++ b/src/pineagents/agents/tools/agent_management.py @@ -15,7 +15,9 @@ from agentscope.message import TextBlock from agentscope.tool import ToolChunk from agentscope.message import ToolResultState +from ...config.config import load_agent_config from ...config.utils import read_last_api +from ...constant import DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS from ...runtime.tool_registry import tool_descriptor from ...utils.http import trust_env_for_url @@ -936,7 +938,7 @@ def _coerce_bool( def _coerce_timeout( value: Any, - default: int = 600, + default: int = DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS, *, field_name: str = "timeout", ) -> int: @@ -1006,6 +1008,15 @@ def _build_subagent_request_context( ) -> dict[str, Any]: """Build request_context with approval routing + tool/skill filters.""" rc = _build_spawn_request_context(current_agent_id) + try: + agent_config = load_agent_config(current_agent_id) + subagent_model = agent_config.subagent_model + if subagent_model is not None: + rc["model_slot_override"] = subagent_model.model_dump() + except Exception: # pylint: disable=broad-exception-caught + # Subagents must remain usable when an optional per-agent model + # override cannot be loaded from a stale or synthetic test identity. + pass if extra: rc.update(extra) if allowed_tools is not None: @@ -1026,7 +1037,7 @@ async def spawn_subagent( # pylint: disable=too-many-return-statements task: str, fork: bool | str | int = False, background: bool | str | int = False, - timeout: int | float | str = 600, + timeout: int | float | str = DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS, allowed_tools: Optional[list[str] | str] = None, skills: Optional[list[str] | str] = None, batch: Optional[list[Dict[str, Any]] | str] = None, @@ -1134,7 +1145,7 @@ async def spawn_subagent( # pylint: disable=too-many-return-statements default=False, field_name="background", ) - timeout = _coerce_timeout(timeout, default=600, field_name="timeout") + timeout = _coerce_timeout(timeout, default=DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS, field_name="timeout") except ValueError as exc: return _tool_text_response(f"ERROR: {exc}") @@ -1266,7 +1277,7 @@ async def _spawn_batch( ), "timeout": _coerce_timeout( spec.get("timeout"), - default=600, + default=DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS, field_name=f"batch[{i}].timeout", ), }, diff --git a/src/pineagents/app/chats/api.py b/src/pineagents/app/chats/api.py index 755ec11..78ca80b 100644 --- a/src/pineagents/app/chats/api.py +++ b/src/pineagents/app/chats/api.py @@ -17,6 +17,10 @@ from .session import SafeJSONSession from .manager import ChatManager, MAX_BATCH_SIZE from .models import ( BatchArchiveResult, + ChatGroup, + ChatGroupCreate, + ChatGroupOrderUpdate, + ChatGroupUpdate, ChatSpec, ChatUpdate, ChatHistory, @@ -135,6 +139,76 @@ async def create_chat( return await mgr.create_chat(spec) +# ----- Chat group endpoints ----- + + +@router.get("/groups", response_model=list[ChatGroup]) +async def list_chat_groups( + mgr: ChatManager = Depends(get_chat_manager), +): + """List built-in and custom groups in display order.""" + return await mgr.list_groups() + + +@router.post("/groups", response_model=ChatGroup) +async def create_chat_group( + payload: ChatGroupCreate, + mgr: ChatManager = Depends(get_chat_manager), +): + """Create a custom chat group.""" + try: + return await mgr.create_group(payload.name) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@router.put("/groups/order", response_model=list[ChatGroup]) +async def reorder_chat_groups( + payload: ChatGroupOrderUpdate, + mgr: ChatManager = Depends(get_chat_manager), +): + """Replace the complete chat-group display order.""" + try: + return await mgr.reorder_groups(payload.group_ids) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +@router.put("/groups/{group_id}", response_model=ChatGroup) +async def update_chat_group( + group_id: str, + payload: ChatGroupUpdate, + mgr: ChatManager = Depends(get_chat_manager), +): + """Rename or pin a mutable chat group.""" + try: + group = await mgr.update_group( + group_id, + name=payload.name, + pinned=payload.pinned, + ) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + if group is None: + raise HTTPException(status_code=404, detail="Chat group not found") + return group + + +@router.delete("/groups/{group_id}", response_model=dict) +async def delete_chat_group( + group_id: str, + mgr: ChatManager = Depends(get_chat_manager), +): + """Delete a custom group and re-home its chats.""" + try: + deleted = await mgr.delete_group(group_id) + except ValueError as exc: + raise HTTPException(status_code=409, detail=str(exc)) from exc + if not deleted: + raise HTTPException(status_code=404, detail="Chat group not found") + return {"success": True, "group_id": group_id} + + @router.post("/batch-delete", response_model=dict) async def batch_delete_chats( chat_ids: list[str], diff --git a/src/pineagents/app/chats/manager.py b/src/pineagents/app/chats/manager.py index f7c558a..7266d9e 100644 --- a/src/pineagents/app/chats/manager.py +++ b/src/pineagents/app/chats/manager.py @@ -11,9 +11,12 @@ from typing import Optional from .models import ( BatchArchiveResult, BatchFailure, + ChatGroup, + ChatGroupKind, ChatSpec, ChatUpdate, SessionSource, + SOURCE_CHAT_GROUP_IDS, ) from .repo import BaseChatRepository from ..channels.schema import DEFAULT_CHANNEL @@ -24,6 +27,38 @@ logger = logging.getLogger(__name__) MAX_BATCH_SIZE = 500 +def _default_group_id(source: SessionSource) -> str: + """Return the built-in group for a chat source.""" + return SOURCE_CHAT_GROUP_IDS[source] + + +def _source_order(group: ChatGroup) -> int: + """Return a stable tail order for source-driven groups.""" + if group.kind == ChatGroupKind.cron: + return list(SessionSource).index(SessionSource.cron) + if group.kind == ChatGroupKind.subagents: + return list(SessionSource).index(SessionSource.subagent) + return -1 + + +def _is_fixed_source_group(group: ChatGroup) -> bool: + """Return whether a group represents an automated session source.""" + return group.kind in {ChatGroupKind.cron, ChatGroupKind.subagents} + + +def _ordered_groups(groups: list[ChatGroup]) -> list[ChatGroup]: + """Sort pinned groups first and keep source groups at the end.""" + return sorted( + groups, + key=lambda group: ( + _is_fixed_source_group(group), + _source_order(group) + if _is_fixed_source_group(group) + else (not group.pinned, group.order), + ), + ) + + class ChatManager: """Manages chat specifications in repository. @@ -109,6 +144,9 @@ class ChatManager: channel: str = DEFAULT_CHANNEL, name: str = "New Chat", source: str | SessionSource = SessionSource.chat, + group_id: str | None = None, + parent_session_id: str | None = None, + root_session_id: str | None = None, ) -> ChatSpec: """Get existing chat or create new one. @@ -156,6 +194,9 @@ class ChatManager: channel=channel, name=name, source=resolved_source, + group_id=group_id or _default_group_id(resolved_source), + parent_session_id=parent_session_id, + root_session_id=root_session_id, ) logger.debug(f"get_or_create_chat: created spec={spec.id}") # Call internal create without lock (already locked) @@ -175,9 +216,128 @@ class ChatManager: Chat spec """ async with self._lock: + if spec.group_id is None: + spec = spec.model_copy( + update={"group_id": _default_group_id(spec.source)}, + ) + await self._validate_group_id_locked(spec.group_id) await self._repo.upsert_chat(spec) return spec + async def _validate_group_id_locked(self, group_id: str | None) -> None: + """Validate a group ID while the manager lock is held.""" + chats_file = await self._repo.load() + if group_id not in {group.id for group in chats_file.groups}: + raise ValueError(f"Unknown chat group: {group_id}") + + async def list_groups(self) -> list[ChatGroup]: + """List persisted groups in display order.""" + async with self._lock: + chats_file = await self._repo.load() + return _ordered_groups(chats_file.groups) + + async def create_group(self, name: str) -> ChatGroup: + """Create a custom group after the current final group.""" + normalized_name = name.strip() + if not normalized_name: + raise ValueError("Group name cannot be empty") + async with self._lock: + chats_file = await self._repo.load() + next_order = ( + max( + (group.order for group in chats_file.groups), + default=-1, + ) + + 1 + ) + group = ChatGroup(name=normalized_name, order=next_order) + chats_file.groups.append(group) + await self._repo.save(chats_file) + return group + + async def update_group( + self, + group_id: str, + *, + name: str | None = None, + pinned: bool | None = None, + ) -> ChatGroup | None: + """Rename or pin a mutable group.""" + if name is None and pinned is None: + raise ValueError("At least one group field must be provided") + normalized_name = name.strip() if name is not None else None + if normalized_name == "": + raise ValueError("Group name cannot be empty") + async with self._lock: + chats_file = await self._repo.load() + for index, group in enumerate(chats_file.groups): + if group.id != group_id: + continue + if _is_fixed_source_group(group): + raise ValueError("Source groups cannot be changed") + updates = {} + if normalized_name is not None: + updates["name"] = normalized_name + if pinned is not None: + updates["pinned"] = pinned + updated = group.model_copy(update=updates) + chats_file.groups[index] = updated + await self._repo.save(chats_file) + return updated + return None + + async def reorder_groups(self, group_ids: list[str]) -> list[ChatGroup]: + """Persist a complete, duplicate-free group order.""" + async with self._lock: + chats_file = await self._repo.load() + current = {group.id: group for group in chats_file.groups} + if len(group_ids) != len(set(group_ids)): + raise ValueError("Group order contains duplicate IDs") + if set(group_ids) != set(current): + raise ValueError("Group order must contain every group ID") + fixed_ids = [ + group.id + for group in _ordered_groups(list(current.values())) + if _is_fixed_source_group(group) + ] + if group_ids[-len(fixed_ids) :] != fixed_ids: + raise ValueError("Source groups must remain at the end") + chats_file.groups = [ + current[group_id].model_copy(update={"order": index}) + for index, group_id in enumerate(group_ids) + ] + await self._repo.save(chats_file) + return _ordered_groups(chats_file.groups) + + async def delete_group(self, group_id: str) -> bool: + """Delete a custom group and return chats to their system group.""" + async with self._lock: + chats_file = await self._repo.load() + target = next( + (group for group in chats_file.groups if group.id == group_id), + None, + ) + if target is None: + return False + if target.kind != ChatGroupKind.custom: + raise ValueError("Built-in chat groups cannot be deleted") + + chats_file.groups = [ + group for group in chats_file.groups if group.id != group_id + ] + for index, chat in enumerate(chats_file.chats): + if chat.group_id != group_id: + continue + chats_file.chats[index] = chat.model_copy( + update={"group_id": _default_group_id(chat.source)}, + ) + for index, group in enumerate( + sorted(chats_file.groups, key=lambda item: item.order), + ): + group.order = index + await self._repo.save(chats_file) + return True + async def patch_chat( self, chat_id: str, @@ -223,12 +383,16 @@ class ChatManager: if existing is None: return None + if "group_id" in patch.model_fields_set: + await self._validate_group_id_locked(patch.group_id) + updates = patch.model_dump( exclude_none=True, exclude_unset=True, ) merged = existing.model_copy(update=updates) - merged.updated_at = datetime.now(timezone.utc) + if patch.model_fields_set != {"group_id"}: + merged.updated_at = datetime.now(timezone.utc) await self._repo.upsert_chat(merged) return merged diff --git a/src/pineagents/app/chats/models.py b/src/pineagents/app/chats/models.py index 0c11b0d..32dcffc 100644 --- a/src/pineagents/app/chats/models.py +++ b/src/pineagents/app/chats/models.py @@ -7,7 +7,13 @@ from enum import Enum from typing import Any, Dict, Literal, Optional from uuid import uuid4 -from pydantic import BaseModel, ConfigDict, Field, computed_field +from pydantic import ( + BaseModel, + ConfigDict, + Field, + computed_field, + model_validator, +) from pineagents.schemas import Message from ..channels.schema import DEFAULT_CHANNEL @@ -21,6 +27,71 @@ class SessionSource(str, Enum): chat = "chat" cron = "cron" + subagent = "subagent" + + +class ChatGroupKind(str, Enum): + """Distinguishes built-in and user-created chat groups.""" + + default = "default" + cron = "cron" + subagents = "subagents" + custom = "custom" + + +DEFAULT_CHAT_GROUP_ID = "default" +CRON_CHAT_GROUP_ID = "cron" +SUBAGENT_CHAT_GROUP_ID = "subagents" + +SOURCE_CHAT_GROUP_IDS = { + SessionSource.chat: DEFAULT_CHAT_GROUP_ID, + SessionSource.cron: CRON_CHAT_GROUP_ID, + SessionSource.subagent: SUBAGENT_CHAT_GROUP_ID, +} + + +class ChatGroup(BaseModel): + """One persisted group in the Console chat list.""" + + id: str = Field(default_factory=lambda: str(uuid4())) + name: str = Field(min_length=1, max_length=64) + order: int = Field(default=0, ge=0) + kind: ChatGroupKind = Field(default=ChatGroupKind.custom) + source: Optional[SessionSource] = Field( + default=None, + description="Session source represented by a built-in group", + ) + pinned: bool = Field( + default=False, + description="Whether the group is pinned above regular groups", + ) + + +def default_chat_groups() -> list[ChatGroup]: + """Return source-driven built-in groups for a chat registry.""" + return [ + ChatGroup( + id=DEFAULT_CHAT_GROUP_ID, + name="Uncategorized", + order=0, + kind=ChatGroupKind.default, + source=SessionSource.chat, + ), + ChatGroup( + id=CRON_CHAT_GROUP_ID, + name="Scheduled tasks", + order=1, + kind=ChatGroupKind.cron, + source=SessionSource.cron, + ), + ChatGroup( + id=SUBAGENT_CHAT_GROUP_ID, + name="Subagents", + order=2, + kind=ChatGroupKind.subagents, + source=SessionSource.subagent, + ), + ] class ChatSpec(BaseModel): @@ -48,6 +119,10 @@ class ChatSpec(BaseModel): default_factory=lambda: datetime.now(timezone.utc), description="Chat last update timestamp", ) + last_finished_at: Optional[datetime] = Field( + default=None, + description="When the most recent task for this chat finished", + ) meta: Dict[str, Any] = Field( default_factory=dict, description="Additional metadata", @@ -68,6 +143,18 @@ class ChatSpec(BaseModel): default=SessionSource.chat, description="What initiated this session (chat, cron, …)", ) + group_id: Optional[str] = Field( + default=None, + description="Persisted Console group identifier", + ) + parent_session_id: Optional[str] = Field( + default=None, + description="Immediate parent session for a subagent chat", + ) + root_session_id: Optional[str] = Field( + default=None, + description="Root session for a subagent chat tree", + ) @computed_field # type: ignore[misc] @property @@ -91,6 +178,42 @@ class ChatUpdate(BaseModel): default=None, description="Whether the chat is pinned to the top", ) + group_id: str | None = Field( + default=None, + description="Target Console group identifier", + ) + + +class ChatGroupCreate(BaseModel): + """Fields accepted when creating a custom chat group.""" + + model_config = ConfigDict(extra="forbid") + + name: str = Field(min_length=1, max_length=64) + + +class ChatGroupUpdate(BaseModel): + """Mutable chat-group fields.""" + + model_config = ConfigDict(extra="forbid") + + name: str | None = Field(default=None, min_length=1, max_length=64) + pinned: bool | None = None + + @model_validator(mode="after") + def require_update(self) -> "ChatGroupUpdate": + """Reject an empty group update.""" + if self.name is None and self.pinned is None: + raise ValueError("At least one group field must be provided") + return self + + +class ChatGroupOrderUpdate(BaseModel): + """Complete group order submitted by the Console.""" + + model_config = ConfigDict(extra="forbid") + + group_ids: list[str] = Field(min_length=2) class ChatHistory(BaseModel): @@ -126,3 +249,14 @@ class ChatsFile(BaseModel): version: int = 1 chats: list[ChatSpec] = Field(default_factory=list) + groups: list[ChatGroup] = Field(default_factory=default_chat_groups) + + @model_validator(mode="after") + def ensure_system_groups(self) -> "ChatsFile": + """Ensure every source-driven built-in group is present.""" + by_id = {group.id: group for group in self.groups} + defaults = default_chat_groups() + for group in defaults: + if group.id not in by_id: + self.groups.append(group) + return self diff --git a/src/pineagents/app/routers/agents.py b/src/pineagents/app/routers/agents.py index 966dace..2a0e4c9 100644 --- a/src/pineagents/app/routers/agents.py +++ b/src/pineagents/app/routers/agents.py @@ -102,6 +102,7 @@ class CreateAgentRequest(BaseModel): language: str | None = None skill_names: list[str] | None = None active_model: ModelSlotConfig | None = None + subagent_model: ModelSlotConfig | None = None backend: str = "qwenpaw" backend_settings: dict[str, Any] = Field(default_factory=dict) @@ -610,6 +611,7 @@ async def create_agent( heartbeat=HeartbeatConfig(), tools=ToolsConfig(), active_model=active_model, + subagent_model=request.subagent_model, ) _initialize_agent_workspace( diff --git a/src/pineagents/app/routers/console.py b/src/pineagents/app/routers/console.py index b09e8cd..380bc4b 100644 --- a/src/pineagents/app/routers/console.py +++ b/src/pineagents/app/routers/console.py @@ -186,6 +186,26 @@ def _is_reconnect_request(request_data: Union[AgentRequest, dict]) -> bool: return getattr(request_data, "reconnect", None) is True +def _chat_registration_fields(native_payload: dict[str, Any]) -> dict: + """Return first-class subagent fields from an internal request.""" + request_context = native_payload["meta"].get("request_context") + if not isinstance(request_context, dict): + return {} + if request_context.get("_spawn_subagent") is not True: + return {} + return { + "source": "subagent", + "parent_session_id": str( + request_context.get("parent_session_id") or "", + ) + or None, + "root_session_id": str( + request_context.get("root_session_id") or "", + ) + or None, + } + + def _empty_sse_response() -> StreamingResponse: """An SSE response that terminates immediately.""" @@ -294,6 +314,7 @@ async def post_console_chat( native_payload["sender_id"], native_payload["channel_id"], name=name, + **_chat_registration_fields(native_payload), ) tracker = workspace.task_tracker @@ -635,6 +656,7 @@ async def post_console_chat_task( # pylint: disable=too-many-statements native_payload["sender_id"], native_payload["channel_id"], name=name, + **_chat_registration_fields(native_payload), ) task_timeout: Optional[float] = None diff --git a/src/pineagents/config/config.py b/src/pineagents/config/config.py index d11cbeb..f99cc25 100644 --- a/src/pineagents/config/config.py +++ b/src/pineagents/config/config.py @@ -1731,6 +1731,10 @@ class AgentProfileConfig(BaseModel): default=None, description="Active model for this agent (provider_id + model)", ) + subagent_model: Optional["ModelSlotConfig"] = Field( + default=None, + description="Optional cheaper model used by spawned subagents", + ) language: str = Field( default="zh", description="Language setting for this agent", diff --git a/src/pineagents/constant.py b/src/pineagents/constant.py index 8d73d77..d82a1da 100644 --- a/src/pineagents/constant.py +++ b/src/pineagents/constant.py @@ -279,6 +279,9 @@ HEARTBEAT_MAX_TIMEOUT_SECONDS = 3600 HEARTBEAT_TARGET_LAST = "last" HEARTBEAT_TARGET_INBOX = "inbox" +# Parent HTTP wait for spawn_subagent foreground (/console/chat). +DEFAULT_SPAWN_FOREGROUND_TIMEOUT_SECONDS = 600 + # Debug history file for /dump_history and /load_history commands DEBUG_HISTORY_FILE = EnvVarLoader.get_str( "QWENPAW_DEBUG_HISTORY_FILE",