Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 36 additions & 11 deletions backend/app/api/channels.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@
ChannelInboundEvent,
ChatSession,
Message,
Team,
User,
utc_now,
)
Expand Down Expand Up @@ -275,27 +276,48 @@ def create_channel_binding(
ensure_current_user_tenant(request.tenant_id, current_user)
if request.channel not in SUPPORTED_CHANNELS:
raise HTTPException(status_code=400, detail=f"v1 仅支持渠道: {sorted(SUPPORTED_CHANNELS)}")
ensure_agent_scope_manager(db, request.tenant_id, request.agent_id, current_user)
# 挂员工集或绑团队二选一:都给/都不给均拒绝
if bool(request.agent_id) == bool(request.team_id):
raise HTTPException(status_code=400, detail="agent_id 与 team_id 必须且只能提供一个")
if request.team_id:
team = db.get(Team, request.team_id)
if team is None or team.tenant_id != request.tenant_id:
raise HTTPException(status_code=404, detail="Team not found")
from app.teams.service import get_team_leader

leader = get_team_leader(db, team.id)
if leader is None:
raise HTTPException(status_code=400, detail="团队暂未设置 TL,请先设置 TL 后再绑定渠道")
# 复用员工绑定同款守卫:创建者须能管理现任 TL 员工
ensure_agent_scope_manager(db, request.tenant_id, leader.agent_id, current_user)
# agent_id 为非空遗留列(列表过滤/挂载回退仍在用):团队绑定回写现任 TL,
# 入站路由始终以 binding.team_id 解析的现任 TL 为准,换帅自动跟随
target_agent_id = leader.agent_id
else:
ensure_agent_scope_manager(db, request.tenant_id, request.agent_id, current_user)
target_agent_id = request.agent_id
# 同一员工同一渠道允许多个绑定实例,总是新建
binding = ChannelBinding(
tenant_id=request.tenant_id,
agent_id=request.agent_id,
agent_id=target_agent_id,
channel=request.channel,
status="pending",
created_by_user_id=current_user.id,
team_id=request.team_id,
)
db.add(binding)
db.flush()
# 新绑定自动挂载默认员工
db.add(
ChannelBindingAgent(
tenant_id=request.tenant_id,
binding_id=binding.id,
agent_id=request.agent_id,
is_default=True,
sort_order=0,
if not request.team_id:
# 新绑定自动挂载默认员工;团队绑定走 TL 直路由,不写挂载行
db.add(
ChannelBindingAgent(
tenant_id=request.tenant_id,
binding_id=binding.id,
agent_id=target_agent_id,
is_default=True,
sort_order=0,
)
)
)
db.commit()
db.refresh(binding)
return channel_binding_read(db, binding)
Expand Down Expand Up @@ -470,6 +492,9 @@ def update_channel_binding_agents(
_ensure_binding_manager(db, tenant_id, binding, current_user)
if request.agents is None and request.auto_route is None:
raise HTTPException(status_code=400, detail="无有效更新内容")
if request.agents is not None and binding.team_id:
# 团队绑定的接待员工由团队现任 TL 决定,不允许整表替换员工挂载
raise HTTPException(status_code=400, detail="团队绑定的渠道不支持修改员工挂载")
default_agent_id: str | None = None
if request.agents is not None:
if not request.agents:
Expand Down
143 changes: 128 additions & 15 deletions backend/app/api/chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@
from app.agents.branching import model_for_agent, visible_published_skills
from app.channels.service_outbox import stage_channel_delivery
from app.core import AgentLoop
from app.core.capability_manifest import CapabilityManifestBuilder
from app.core.cancellation import cancel_chat_turn
from app.core.capability_manifest import CapabilityManifestBuilder
from app.core.harness_session_cleanup import (
harness_task_workspace_path,
remove_harness_session_workspace,
Expand All @@ -43,11 +43,17 @@
ScheduledTaskRun,
Skill,
SkillFeedback,
Team,
User,
new_id,
utc_now,
)
from app.feedback import enqueue_feedback_analysis
from app.harness import (
HarnessArtifactAccessError,
normalize_harness_artifact_path,
open_harness_artifact,
)
from app.knowledge.citations import CITATION_EXCERPT_CHAR_LIMIT, compact_knowledge_citation_labels
from app.llm import LLMClient, LLMError
from app.observability.spans import (
Expand All @@ -56,16 +62,11 @@
reset_span_sink,
set_span_sink,
)
from app.scheduled_tasks.schema import ScheduledTaskDraftRead
from app.scheduled_tasks.service import DEFAULT_TASK_TIME, detect_scheduled_task_draft
from app.security.auth import get_current_user
from app.security.permissions import agent_owned_by_user, is_admin_user
from app.security.tenant import ensure_tenant
from app.harness import (
HarnessArtifactAccessError,
normalize_harness_artifact_path,
open_harness_artifact,
)
from app.scheduled_tasks.schema import ScheduledTaskDraftRead
from app.scheduled_tasks.service import DEFAULT_TASK_TIME, detect_scheduled_task_draft
from app.session.attachments import (
parse_chat_attachment,
validate_chat_turn_attachments,
Expand All @@ -82,6 +83,8 @@
MessageFeedbackRequest,
MessageRead,
)
from app.teams.service import get_team_leader
from app.teams.wakeup import build_tl_chat_message, process_tl_reply

router = APIRouter(prefix="/api/chat", tags=["chat"])
logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -175,7 +178,9 @@ class HumanHandoffReplyRequest(BaseModel):
reply: str


def session_read(row: ChatSession, *, is_scheduled: bool = False) -> ChatSessionRead:
def session_read(
row: ChatSession, *, is_scheduled: bool = False, team_name: str | None = None
) -> ChatSessionRead:
return ChatSessionRead(
id=row.id,
tenant_id=row.tenant_id,
Expand All @@ -188,6 +193,8 @@ def session_read(row: ChatSession, *, is_scheduled: bool = False) -> ChatSession
summary=row.summary,
last_agent_question=row.last_agent_question,
is_scheduled=is_scheduled,
team_id=row.team_id,
team_name=team_name,
created_at=row.created_at.isoformat(),
updated_at=row.updated_at.isoformat(),
)
Expand Down Expand Up @@ -985,14 +992,26 @@ def chat_turn(
_ensure_request_tenant(request.tenant_id, current_user)
request = request.model_copy(update={"user_id": current_user.id})
request = _validate_chat_turn_attachments(request)
team_tl_team: Team | None = None
if request.session_id:
chat_session = _ensure_chat_session_available(db, request.tenant_id, current_user.id, request.session_id)
_ensure_team_session_human_writable(chat_session)
request = _bind_request_to_session_agent(db, request, chat_session, current_user)
team_tl_team = _team_tl_session_team(db, chat_session)
else:
_ensure_chat_agent_available(db, request.tenant_id, request.agent_id, current_user)
ensure_tenant(db, request.tenant_id)
if not request.message.strip() and not request.attachments:
raise HTTPException(status_code=400, detail="Message cannot be empty")
original_message = request.message
if team_tl_team is not None:
# 团队 TL 会话:注入团队上下文(花名册/未闭环任务/黑板/派任务格式)后再走正常引擎
request = request.model_copy(
update={
"message": build_tl_chat_message(db, team_tl_team, original_message),
"interaction_mode": "team_tl",
}
)
if request.session_id:
scheduled_response = _maybe_handle_scheduled_task_request(db, request, chat_session)
if scheduled_response:
Expand All @@ -1001,6 +1020,21 @@ def chat_turn(
return response
response = AgentLoop(db).handle_turn(request)
_schedule_session_title_summary(request.tenant_id, request.user_id, response.session_id, request.agent_id)
if team_tl_team is not None:
# TL 回复后处理:解析派任务块并创建任务(与 tl_chat 端点同语义);
# 后处理失败不影响本轮回复
try:
process_tl_reply(
db,
team=team_tl_team,
session=chat_session,
user=current_user,
user_message=original_message,
reply=response.reply or "",
client_turn_id=request.client_turn_id,
)
except Exception:
logger.exception("team TL reply post-processing failed")
if request.interaction_mode == "scheduled_task" and request.agent_id:
draft = detect_scheduled_task_draft(
db,
Expand All @@ -1026,13 +1060,26 @@ def chat_stream(
request = request.model_copy(update={"user_id": current_user.id})
request = _validate_chat_turn_attachments(request)
ensure_tenant(db, request.tenant_id)
team_tl_team_id: str | None = None
if request.session_id:
chat_session = _ensure_chat_session_available(db, request.tenant_id, current_user.id, request.session_id)
_ensure_team_session_human_writable(chat_session)
request = _bind_request_to_session_agent(db, request, chat_session, current_user)
team_tl_team = _team_tl_session_team(db, chat_session)
team_tl_team_id = team_tl_team.id if team_tl_team is not None else None
else:
_ensure_chat_agent_available(db, request.tenant_id, request.agent_id, current_user)
if not request.message.strip() and not request.attachments:
raise HTTPException(status_code=400, detail="Message cannot be empty")
original_message = request.message
if team_tl_team_id is not None:
# 团队 TL 会话:注入团队上下文(花名册/未闭环任务/黑板/派任务格式)后再走正常引擎
request = request.model_copy(
update={
"message": build_tl_chat_message(db, db.get(Team, team_tl_team_id), original_message),
"interaction_mode": "team_tl",
}
)

relay_ready = threading.Event()
worker_done = threading.Event()
Expand Down Expand Up @@ -1199,6 +1246,27 @@ def persist_span(event_type: str, payload: dict[str, object]) -> None:
event_source_session_id,
request.agent_id,
)
if team_tl_team_id is not None:
# 团队 TL 会话:complete 后做派任务后处理(与 tl_chat 端点同语义);
# 后处理失败不影响本轮回复
try:
tl_team = worker_db.get(Team, team_tl_team_id)
tl_session = worker_db.get(
ChatSession, event_source_session_id or request.session_id or ""
)
tl_user = worker_db.get(User, request.user_id) if request.user_id else None
if tl_team is not None and tl_session is not None and tl_user is not None:
process_tl_reply(
worker_db,
team=tl_team,
session=tl_session,
user=tl_user,
user_message=original_message,
reply=str(data.get("reply") or ""),
client_turn_id=request.client_turn_id,
)
except Exception:
logger.exception("team TL reply post-processing failed")
if event_source_session_id:
summary_payload = _session_title_summary_payload(worker_db, request.tenant_id, event_source_session_id)
if summary_payload:
Expand Down Expand Up @@ -1834,7 +1902,19 @@ def list_chat_sessions(
).all()
if session_id
}
return [session_read(row, is_scheduled=row.id in scheduled_session_ids) for row in rows]
team_ids = {row.team_id for row in rows if row.team_id}
team_names = {
team.id: team.name
for team in db.exec(select(Team).where(Team.id.in_(team_ids))).all()
} if team_ids else {}
return [
session_read(
row,
is_scheduled=row.id in scheduled_session_ids,
team_name=team_names.get(row.team_id) if row.team_id else None,
)
for row in rows
]


@router.put("/sessions/{session_id}", response_model=ChatSessionRead)
Expand All @@ -1858,7 +1938,8 @@ def rename_chat_session(
ScheduledTaskRun.session_id == row.id,
)
).first() is not None
return session_read(row, is_scheduled=is_scheduled)
team = db.get(Team, row.team_id) if row.team_id else None
return session_read(row, is_scheduled=is_scheduled, team_name=team.name if team else None)


@router.delete("/sessions/{session_id}")
Expand Down Expand Up @@ -2413,11 +2494,45 @@ def _bind_request_to_session_agent(
def _ensure_chat_session_available(db: Session, tenant_id: str, user_id: str, session_id: str) -> ChatSession:
ensure_tenant(db, tenant_id)
row = db.get(ChatSession, session_id)
if not row or row.tenant_id != tenant_id or row.user_id != user_id:
if not row or row.tenant_id != tenant_id:
raise HTTPException(status_code=404, detail="Session not found")
# 团队会话(team_id 非空)对本租户成员开放发言(如 TL 工作台聊天室);
# 普通会话仍仅创建者可见
if row.user_id != user_id and not row.team_id:
raise HTTPException(status_code=404, detail="Session not found")
return row


def _ensure_team_session_human_writable(chat_session: ChatSession) -> None:
"""团队内部会话(任务执行/竞标/验收)仅可查看,不允许人工 /turn、/stream 写入。

判据与 _team_tl_session_team 一致:只有「TL 对话」标题的团队会话才对人类开放发言,
其余团队会话由唤醒机制自主驱动,人工写入会污染任务历史并绕过 Agent 权限校验。
"""
if chat_session.team_id and "TL 对话" not in (chat_session.title or ""):
raise HTTPException(status_code=403, detail="Team execution sessions are read-only")


def _team_tl_session_team(db: Session, chat_session: ChatSession) -> Team | None:
"""识别团队 TL 会话:session 挂 team_id 且绑定 agent 是该团队现任 TL。

非 TL 的团队会话(任务执行/竞标等)不对人直接聊,返回 None(不注入、不后处理)。
"""
if not chat_session.team_id or not chat_session.agent_id:
return None
# 团队会话全量绑定 team_id 后,任务验收/竞标打分等会话同样挂在 TL 名下;
# 只有「TL 对话」标题的会话才按人对 TL 聊天处理(与 team-threads 列表同判据)
if "TL 对话" not in (chat_session.title or ""):
return None
team = db.get(Team, chat_session.team_id)
if team is None or team.tenant_id != chat_session.tenant_id or team.status != "active":
return None
leader = get_team_leader(db, team.id)
if leader is None or leader.agent_id != chat_session.agent_id:
return None
return team


def _get_feedback_target_message(db: Session, tenant_id: str, user_id: str, message_id: str) -> Message:
ensure_tenant(db, tenant_id)
row = db.get(Message, message_id)
Expand Down Expand Up @@ -3347,9 +3462,7 @@ def _event_trace_line(
label = "等待SOP"
elif runtime_decision in {"start_skill", "start_new_task"}:
label = "选择SOP"
elif runtime_decision == "suspend_current_and_start_new_skill":
label = "切换SOP"
elif (
elif runtime_decision == "suspend_current_and_start_new_skill" or (
runtime_decision
in {"answer_related_question_then_resume", "answer_chitchat_then_resume"}
and from_skill_id
Expand Down
Loading