Skip to content
Merged
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
1 change: 1 addition & 0 deletions app/api/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
# API 层:路由、请求校验与统一响应封装
18 changes: 18 additions & 0 deletions app/api/routes/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
# 路由聚合:main.py 以 /api 前缀统一挂载
#
# 当前开发阶段仅挂载核心数据查询与管理接口:
# 会话/消息(sessions)、客户(customers)、订单(orders)、工单(tickets)、健康检查(health)。
# 已下线(代码保留,未挂载):
# 风险队列 app/api/routes/risk.py;Agent 副驾与 SSE app/api/routes/deprecated_agent.py。
# 恢复方式:在下方重新 include_router 对应 router 即可。

from fastapi import APIRouter

from app.api.routes import customers, health, orders, sessions, tickets

api_router = APIRouter()
api_router.include_router(health.router)
api_router.include_router(sessions.router)
api_router.include_router(customers.router)
api_router.include_router(orders.router)
api_router.include_router(tickets.router)
43 changes: 43 additions & 0 deletions app/api/routes/customers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
# 消费者路由:基础信息查询与跨会话/订单/工单的统一时间线

from typing import Annotated, Literal

from fastapi import APIRouter, Depends, Query, Request
from sqlalchemy.orm import Session as DbSession

from app.core.deps import get_db
from app.core.envelope import ok
from app.services import consumer_service, timeline_service

router = APIRouter(prefix="/customers", tags=["customers"])


@router.get("/{customer_id}")
def get_customer(request: Request, customer_id: str, db: Annotated[DbSession, Depends(get_db)]) -> dict:
"""客户基础信息与统计(会话数、订单数、工单数、未闭环工单数)。"""
return ok(request, consumer_service.get_customer(db, customer_id))


@router.get("/{customer_id}/timeline")
def get_timeline(
request: Request,
customer_id: str,
db: Annotated[DbSession, Depends(get_db)],
from_time: str | None = Query(None, alias="from", description="RFC3339 起始时间"),
to_time: str | None = Query(None, alias="to", description="RFC3339 结束时间"),
event_type: str | None = Query(None, max_length=200, description="事件类型,逗号分隔多值"),
order: Literal["asc", "desc"] = Query("desc", description="默认倒序(主管看板)"),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
) -> dict:
items, total = timeline_service.get_timeline(
db,
customer_id,
from_time=from_time,
to_time=to_time,
event_type=event_type,
order=order,
page=page,
page_size=page_size,
)
return ok(request, {"items": items}, page=page, page_size=page_size, total=total)
87 changes: 87 additions & 0 deletions app/api/routes/deprecated_agent.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
# 【已下线】Agent 副驾结论与 SSE 事件流路由——当前开发阶段仅保留核心数据查询与管理 API。
#
# 下线原因:本阶段聚焦客户/会话/消息/订单/工单基础接口,副驾结论与实时推送不参与联调。
# 恢复方式:在 app/api/routes/__init__.py 重新 include_router(deprecated_agent.router)。
# 依赖模块(app/agents/copilot.py、app/core/events.py)均保留完好,可直接恢复挂载。

import asyncio
import json
from collections.abc import AsyncIterator
from typing import Annotated, Any, Literal

from fastapi import APIRouter, Depends, Header, Query, Request
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session as DbSession

from app.agents import copilot
from app.core import config
from app.core.deps import get_db
from app.core.envelope import ok, request_id_of
from app.core.errors import session_not_found
from app.core.events import event_bus
from app.models import ServiceSession

router = APIRouter(prefix="/sessions", tags=["deprecated-agent"])


@router.get("/{session_id}/copilot")
def get_session_copilot(
request: Request,
session_id: str,
db: Annotated[DbSession, Depends(get_db)],
mode: Literal["auto", "fast", "reasoning", "vision", "mock"] = Query("auto"),
refresh: bool = Query(False, description="true 时强制重新生成,忽略缓存"),
) -> dict:
"""触发或读取当前会话的 Agent 结论(不改变业务状态)。"""
outcome = copilot.generate_copilot(db, session_id, mode=mode, refresh=refresh, request_id=request_id_of(request))
return ok(
request,
{
"insight": outcome.insight,
"draft_reply": outcome.draft_reply,
"generated_at": outcome.generated_at,
"analysis_id": outcome.analysis_id,
"evidence": outcome.insight.get("evidence", []),
"model_route": outcome.model_route,
"degraded": outcome.degraded,
},
)


@router.get("/{session_id}/stream")
async def stream_session_events(
request: Request,
session_id: str,
db: Annotated[DbSession, Depends(get_db)],
last_event_id: str | None = Header(None, alias="Last-Event-ID"),
) -> StreamingResponse:
"""SSE 推送 Agent 生成、动作执行与承诺状态变化;断线可用 Last-Event-ID 重连。"""
if db.get(ServiceSession, session_id) is None:
raise session_not_found(session_id)
return StreamingResponse(
_event_stream(session_id, last_event_id),
media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no"},
)


async def _event_stream(session_id: str, last_event_id: str | None) -> AsyncIterator[str]:
yield ": connected\n\n"
for record in event_bus.replay_after(session_id, last_event_id):
yield _format_sse(record)
subscription = event_bus.register(session_id)
try:
while True:
try:
record = await asyncio.wait_for(subscription.queue.get(), timeout=config.SSE_KEEPALIVE_SECONDS)
except TimeoutError:
yield ": keep-alive\n\n"
continue
yield _format_sse(record)
finally:
event_bus.unregister(session_id, subscription)


def _format_sse(record: dict[str, Any]) -> str:
data = json.dumps(record["data"], ensure_ascii=False)
return f"id: {record['id']}\nevent: {record['event']}\ndata: {data}\n\n"
42 changes: 42 additions & 0 deletions app/api/routes/health.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
# 健康检查:服务、数据库与模型网关状态(GET /api/health)
#
# 数据库不可用时返回 HTTP 503 DATABASE_UNAVAILABLE(见接口文档 4.1)。

from typing import Annotated

from fastapi import APIRouter, Depends, Request
from sqlalchemy import text
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session as DbSession

from app.core import config
from app.core.deps import get_db
from app.core.envelope import ok
from app.core.errors import ApiError, ErrorCode
from app.core.time_utils import to_api_time, utc_now

router = APIRouter(tags=["health"])


@router.get("/health")
def health(request: Request, db: Annotated[DbSession, Depends(get_db)]) -> dict:
try:
db.execute(text("SELECT 1"))
except SQLAlchemyError as exc:
raise ApiError(
503,
ErrorCode.DATABASE_UNAVAILABLE,
"数据库不可用",
{"reason": type(exc).__name__},
) from exc

return ok(
request,
{
"status": "ok",
"database": "ok",
"model_provider": config.model_provider_name(),
"version": config.APP_VERSION,
"server_time": to_api_time(utc_now()),
},
)
34 changes: 34 additions & 0 deletions app/api/routes/orders.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
# 订单路由:订单事实查询与条件检索

from typing import Annotated

from fastapi import APIRouter, Depends, Query, Request
from sqlalchemy.orm import Session as DbSession

from app.core.deps import get_db
from app.core.envelope import ok
from app.services import order_service

router = APIRouter(prefix="/orders", tags=["orders"])


@router.get("")
def list_orders(
request: Request,
db: Annotated[DbSession, Depends(get_db)],
order_no: str | None = Query(None, max_length=64, description="订单号精确匹配"),
customer_id: str | None = Query(None, max_length=64, description="按消费者筛选"),
session_id: str | None = Query(None, max_length=64, description="按关联会话筛选"),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
) -> dict:
"""按订单号或关联条件(消费者/会话)检索订单列表。"""
items, total = order_service.list_orders(
db, order_no=order_no, customer_id=customer_id, session_id=session_id, page=page, page_size=page_size
)
return ok(request, {"items": items}, page=page, page_size=page_size, total=total)


@router.get("/{order_id}")
def get_order(request: Request, order_id: str, db: Annotated[DbSession, Depends(get_db)]) -> dict:
return ok(request, order_service.get_order(db, order_id))
31 changes: 31 additions & 0 deletions app/api/routes/risk.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# 风险队列路由:主管视角的会话、工单与承诺聚合
#
# 【当前阶段已下线】未在 app/api/routes/__init__.py 挂载(聚焦核心数据查询与管理 API)。
# 恢复方式:重新 include_router(risk.router);依赖 app/services/risk_service.py 保留完好。

from typing import Annotated, Literal

from fastapi import APIRouter, Depends, Query, Request
from sqlalchemy.orm import Session as DbSession

from app.core.deps import get_db
from app.core.envelope import ok
from app.services import risk_service

router = APIRouter(prefix="/risk-queue", tags=["risk"])


@router.get("")
def get_risk_queue(
request: Request,
db: Annotated[DbSession, Depends(get_db)],
risk_level: Literal["L0", "L1", "L2", "L3"] | None = Query(None),
ticket_type: str | None = Query(None, max_length=50),
status: Literal["open", "closed", "pending"] | None = Query(None),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
) -> dict:
items, total = risk_service.build_risk_queue(
db, risk_level=risk_level, ticket_type=ticket_type, status=status, page=page, page_size=page_size
)
return ok(request, {"items": items}, page=page, page_size=page_size, total=total)
80 changes: 80 additions & 0 deletions app/api/routes/sessions.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
# 会话路由:队列列表、聚合详情与聊天消息检索
#
# 当前开发阶段已下线:Agent 副驾结论(GET /sessions/{id}/copilot)与 SSE 事件流
# (GET /sessions/{id}/stream),代码迁移至 app/api/routes/deprecated_agent.py 保留。

from typing import Annotated, Literal

from fastapi import APIRouter, Depends, Header, Query, Request
from sqlalchemy.orm import Session as DbSession

from app.core.deps import get_db
from app.core.envelope import ok
from app.schemas.api import SessionMessageCreateRequest
from app.services import session_service

router = APIRouter(prefix="/sessions", tags=["sessions"])


@router.get("")
def list_sessions(
request: Request,
db: Annotated[DbSession, Depends(get_db)],
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
q: str | None = Query(None, max_length=100, description="脱敏昵称、会话 ID、场景关键词"),
risk_level: Literal["L0", "L1", "L2", "L3"] | None = Query(None),
status: Literal["open", "closed", "pending"] | None = Query(None),
scene_major: str | None = Query(None, max_length=50),
customer_id: str | None = Query(None, max_length=64, description="按消费者筛选"),
sort: Literal["last_message_at", "risk"] | None = Query(None, description="默认风险优先视图"),
) -> dict:
items, total = session_service.list_sessions(
db,
q=q,
risk_level=risk_level,
status=status,
scene_major=scene_major,
customer_id=customer_id,
sort=sort or "risk",
page=page,
page_size=page_size,
)
return ok(request, {"items": items}, page=page, page_size=page_size, total=total)


@router.get("/{session_id}")
def get_session_detail(
request: Request,
session_id: str,
db: Annotated[DbSession, Depends(get_db)],
include: str | None = Query(None, description="可选:events,orders,tickets"),
) -> dict:
return ok(request, session_service.get_session_detail(db, session_id, include=include))


@router.get("/{session_id}/messages")
def get_session_messages(
request: Request,
session_id: str,
db: Annotated[DbSession, Depends(get_db)],
order: Literal["asc", "desc"] = Query("asc", description="默认按 seq_no 升序(聊天读取顺序)"),
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
) -> dict:
"""按会话 ID 分页获取聊天消息记录。"""
items, total = session_service.list_messages(db, session_id, order=order, page=page, page_size=page_size)
return ok(request, {"items": items}, page=page, page_size=page_size, total=total)


@router.post("/{session_id}/messages")
def create_session_message(
request: Request,
session_id: str,
payload: SessionMessageCreateRequest,
db: Annotated[DbSession, Depends(get_db)],
x_operator_id: str | None = Header(None, alias="X-Operator-ID", description="操作人标识,写入审计事件"),
) -> dict:
"""模拟客服发送消息(send=true)或保存草稿(send=false)。"""
result = session_service.append_message(db, session_id, payload=payload, operator_id=x_operator_id)
return ok(request, result)
36 changes: 36 additions & 0 deletions app/api/routes/tickets.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
# 工单路由:统一核心字段 + 白名单 detail + 字段式更新(PATCH)

from typing import Annotated

from fastapi import APIRouter, Depends, Header, Query, Request
from sqlalchemy.orm import Session as DbSession

from app.core.deps import get_db
from app.core.envelope import ok
from app.schemas.api import TicketUpdateRequest
from app.services import ticket_service

router = APIRouter(prefix="/tickets", tags=["tickets"])


@router.get("/{ticket_id}")
def get_ticket(
request: Request,
ticket_id: str,
db: Annotated[DbSession, Depends(get_db)],
include: str | None = Query(None, description="可选:events"),
) -> dict:
include_events = bool(include and "events" in {part.strip() for part in include.split(",")})
return ok(request, ticket_service.get_ticket(db, ticket_id, include_events=include_events))


@router.patch("/{ticket_id}")
def update_ticket(
request: Request,
ticket_id: str,
payload: TicketUpdateRequest,
db: Annotated[DbSession, Depends(get_db)],
x_operator_id: str | None = Header(None, alias="X-Operator-ID", description="操作人标识,写入审计事件"),
) -> dict:
"""更新工单状态/优先级/处理人/备注;状态变化自动写入 service_event 审计与时间线。"""
return ok(request, ticket_service.update_ticket(db, ticket_id, payload=payload, operator_id=x_operator_id))
Loading
Loading