769 lines
34 KiB
Python
769 lines
34 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""ticket/core.py — 工单状态机与全部迁移动作。
|
||
|
||
不变量(设计文档 §4.3):
|
||
- I1 处理方唯一:status+assignee_role+assignee_id 联合表达;迁移全部乐观锁
|
||
(UPDATE ... WHERE status=前置态 + SELECT 验证,照 pipeline_service claim 范式)
|
||
- I2 human_pending ⟺ assignee_role 非空 AND assignee_id NULL
|
||
- I3 human_processing/staff_replied ⟹ assignee_id 非空
|
||
- I4 受理方每次变更必写 tk_transfers(同事务)
|
||
- I5 客户只见自己工单 + visibility='customer' 消息
|
||
- I6 终态(closed/cancelled)不接受任何迁移
|
||
- I7 状态派生待办(同一事项只发一次)
|
||
"""
|
||
|
||
import datetime
|
||
import json
|
||
import logging
|
||
import os
|
||
|
||
from appPublic.uniqueID import getID
|
||
from sqlor.dbpools import DBPools
|
||
|
||
logger = logging.getLogger("ticket.core")
|
||
|
||
MODULE_NAME = 'ticket'
|
||
DBNAME_FALLBACK = 'pipeline'
|
||
|
||
AGENT_USER = 'agent.ticket'
|
||
|
||
# 附件约束:单消息最多 9 个(前端同步拦截;client_max_size=100MB 兜底总量)
|
||
MAX_ATTACHMENTS = 9
|
||
|
||
# ── 状态集 ──
|
||
S_NEW = 'new'
|
||
S_AGENT_PROCESSING = 'agent_processing'
|
||
S_AGENT_REPLIED = 'agent_replied'
|
||
S_HUMAN_PENDING = 'human_pending'
|
||
S_HUMAN_PROCESSING = 'human_processing'
|
||
S_STAFF_REPLIED = 'staff_replied'
|
||
S_CLOSED = 'closed'
|
||
S_CANCELLED = 'cancelled'
|
||
|
||
FINAL_STATES = (S_CLOSED, S_CANCELLED)
|
||
|
||
# 错误码(设计文档 §9)
|
||
E_NOT_FOUND = ('TK_E001', '工单不存在,或您没有权限查看该工单')
|
||
|
||
|
||
def _dbname():
|
||
"""库名解析:宿主 get_module_dbname 优先,兜底 pipeline(跨宿主约定 §7.1)。"""
|
||
try:
|
||
from ahserver.serverenv import ServerEnv
|
||
fn = getattr(ServerEnv(), 'get_module_dbname', None)
|
||
if callable(fn):
|
||
n = fn(MODULE_NAME)
|
||
if n:
|
||
return n
|
||
except Exception:
|
||
pass
|
||
return DBNAME_FALLBACK
|
||
|
||
|
||
def _get_db():
|
||
db = DBPools()
|
||
if not db.databases:
|
||
from appPublic.jsonConfig import getConfig
|
||
config = getConfig()
|
||
if config.databases:
|
||
db.databases = config.databases
|
||
return db
|
||
|
||
|
||
def _now():
|
||
return datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
|
||
|
||
|
||
async def _get_param(sor, name, default):
|
||
recs = await sor.sqlExe(
|
||
"SELECT params_value FROM params WHERE params_name=${n}$ LIMIT 1", {"n": name})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if recs:
|
||
v = str(getattr(recs[0], 'params_value', '') or '').strip()
|
||
if v:
|
||
return v
|
||
return default
|
||
|
||
|
||
async def get_user_roles(sor, user_id):
|
||
"""用户角色列表({orgtypeid}.{name} 格式,与 pipeline_service 同款查询)。"""
|
||
roles = []
|
||
if not user_id:
|
||
return roles
|
||
recs = await sor.sqlExe(
|
||
"SELECT r.orgtypeid, r.name FROM userrole ur JOIN role r ON ur.roleid=r.id "
|
||
"WHERE ur.userid=${u}$", {"u": user_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
for r in (recs or []):
|
||
o = getattr(r, 'orgtypeid', '') or ''
|
||
n = getattr(r, 'name', '') or ''
|
||
if o and n and n != '*':
|
||
roles.append('%s.%s' % (o, n))
|
||
return roles
|
||
|
||
|
||
def _is_admin(roles):
|
||
return ('owner.superuser' in roles) or ('owner.admin' in roles)
|
||
|
||
|
||
# ══════════════════ 内部原语 ══════════════════
|
||
|
||
async def _next_ticket_no(sor):
|
||
"""TK{YYYYMMDD}-{4位序列}。当日计数+1;唯一索引兜底重试在调用方。"""
|
||
day = datetime.date.today().strftime('%Y%m%d')
|
||
prefix = 'TK%s-' % day
|
||
recs = await sor.sqlExe(
|
||
"SELECT COUNT(*) AS c FROM tk_tickets WHERE ticket_no LIKE ${p}$",
|
||
{"p": prefix + '%'})
|
||
await sor.sqlExe("COMMIT", {})
|
||
n = int(getattr(recs[0], 'c', 0)) if recs else 0
|
||
return '%s%04d' % (prefix, n + 1)
|
||
|
||
|
||
async def _add_message(sor, ticket_id, sender_type, sender_id, sender_role,
|
||
content, visibility='customer', attachments=None):
|
||
await sor.C('tk_messages', {
|
||
'id': getID(),
|
||
'ticket_id': ticket_id,
|
||
'sender_type': sender_type,
|
||
'sender_id': sender_id or '',
|
||
'sender_role': sender_role or '',
|
||
'content': content or '',
|
||
'attachments': json.dumps(attachments, ensure_ascii=False) if attachments else None,
|
||
'visibility': visibility,
|
||
'created_at': _now(),
|
||
})
|
||
await sor.sqlExe("COMMIT", {})
|
||
|
||
|
||
# ══════════════════ 附件(FileStorage/idfile 机制,设计文档 §8) ══════════════════
|
||
|
||
def _filestorage():
|
||
from ahserver.filestorage import FileStorage
|
||
return FileStorage()
|
||
|
||
|
||
def normalize_attachments(raw):
|
||
"""把上传附件规整为 [{id,name,webpath,size}](服务端生成,不信任客户端)。
|
||
|
||
raw 形态:str(单个 webpath 或 JSON 数组字符串)/ list[webpath|dict]
|
||
(multipart 多文件同名字段被 ahserver getPostData 聚合为 list)/ dict / None。
|
||
校验:webpath 必须是 FileStorage 落盘的相对路径('/' 开头、无 '..')、
|
||
真实存在、不越出 filesroot;并调 tfr.file_useful 把文件从临时清理队列摘除
|
||
(multipart 上传默认 1 小时后被当临时文件删掉——工单附件必须持久)。
|
||
返回 (items, err);err 为 (code,msg) 或 None。
|
||
"""
|
||
if not raw:
|
||
return [], None
|
||
if isinstance(raw, (str, dict)):
|
||
raw = [raw]
|
||
if not isinstance(raw, list):
|
||
return [], ('TK_E014', '附件参数格式无效')
|
||
# 兼容旧前端契约:JSON 数组字符串 '["/xx/a.png"]' 或 '[{"webpath":...}]'
|
||
if len(raw) == 1 and isinstance(raw[0], str) and raw[0].strip().startswith('['):
|
||
try:
|
||
parsed = json.loads(raw[0])
|
||
if isinstance(parsed, list):
|
||
raw = parsed
|
||
except Exception:
|
||
pass
|
||
if len(raw) > MAX_ATTACHMENTS:
|
||
return [], ('TK_E014', '单条消息最多上传 %d 个附件' % MAX_ATTACHMENTS)
|
||
items = []
|
||
fs = None
|
||
for it in raw:
|
||
if isinstance(it, dict):
|
||
wp = str(it.get('webpath') or it.get('path') or '').strip()
|
||
name = str(it.get('name') or '').strip()
|
||
else:
|
||
wp = str(it or '').strip()
|
||
name = ''
|
||
if not wp:
|
||
continue
|
||
if not wp.startswith('/') or '..' in wp:
|
||
return [], ('TK_E014', '附件路径无效:%s' % wp[:80])
|
||
try:
|
||
if fs is None:
|
||
fs = _filestorage()
|
||
real = fs.realPath(wp)
|
||
root = os.path.realpath(fs.root)
|
||
if not os.path.realpath(real).startswith(root + os.sep) \
|
||
or not os.path.isfile(real):
|
||
return [], ('TK_E014', '附件文件不存在或路径非法:%s' % os.path.basename(wp)[:80])
|
||
size = os.path.getsize(real)
|
||
try:
|
||
fs.tfr.file_useful(wp) # 持久化:从临时文件清理队列摘除
|
||
except Exception:
|
||
pass
|
||
except Exception as e:
|
||
return [], ('TK_E014', '附件校验失败:%s' % str(e)[:100])
|
||
if not name:
|
||
name = os.path.basename(wp)
|
||
items.append({'id': getID(), 'name': name[:200],
|
||
'webpath': wp, 'size': int(size)})
|
||
return items, None
|
||
|
||
|
||
def parse_attachments(raw):
|
||
"""tk_messages.attachments 列(JSON 字符串)→ list;容错返回 []。"""
|
||
if not raw:
|
||
return []
|
||
if isinstance(raw, list):
|
||
return raw
|
||
try:
|
||
v = json.loads(raw)
|
||
return v if isinstance(v, list) else []
|
||
except Exception:
|
||
return []
|
||
|
||
|
||
async def get_attachment_file(msg_id, att_id, user_id, org_id):
|
||
"""附件打开/下载:按消息+附件 id 取文件,I5 可见性校验。
|
||
|
||
客户只能取 visibility='customer' 消息的附件;staff/admin 全量。
|
||
返回 (True, (realpath, filename)) 或 (False, (code, msg))。
|
||
"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
recs = await sor.sqlExe(
|
||
"SELECT id, ticket_id, attachments, visibility FROM tk_messages WHERE id=${i}$",
|
||
{"i": str(msg_id or '')})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not recs:
|
||
return False, E_NOT_FOUND
|
||
m = recs[0]
|
||
t = await _load_ticket(sor, str(getattr(m, 'ticket_id', '') or ''))
|
||
if not t:
|
||
return False, E_NOT_FOUND
|
||
is_staff_view = await _check_staff_access(sor, t, user_id, roles) or _is_admin(roles)
|
||
if not is_staff_view:
|
||
if not await _check_customer_access(sor, t, user_id, org_id, roles):
|
||
return False, E_NOT_FOUND
|
||
if str(getattr(m, 'visibility', '')) != 'customer':
|
||
return False, E_NOT_FOUND
|
||
atts = parse_attachments(getattr(m, 'attachments', ''))
|
||
target = None
|
||
for a in atts:
|
||
if isinstance(a, dict) and str(a.get('id', '')) == str(att_id or ''):
|
||
target = a
|
||
break
|
||
if not target:
|
||
return False, ('TK_E015', '附件不存在')
|
||
wp = str(target.get('webpath', '') or '')
|
||
if not wp.startswith('/') or '..' in wp:
|
||
return False, ('TK_E014', '附件路径无效')
|
||
fs = _filestorage()
|
||
real = fs.realPath(wp)
|
||
root = os.path.realpath(fs.root)
|
||
if not os.path.realpath(real).startswith(root + os.sep) or not os.path.isfile(real):
|
||
return False, ('TK_E015', '附件文件已丢失,请联系平台运维')
|
||
fname = str(target.get('name', '') or os.path.basename(wp))
|
||
return True, (real, fname)
|
||
|
||
|
||
async def _add_transfer(sor, ticket_id, action, from_role='', from_user='',
|
||
to_role='', to_user='', reason='', operator_id=''):
|
||
await sor.C('tk_transfers', {
|
||
'id': getID(),
|
||
'ticket_id': ticket_id,
|
||
'action': action,
|
||
'from_role': from_role or '',
|
||
'from_user': from_user or '',
|
||
'to_role': to_role or '',
|
||
'to_user': to_user or '',
|
||
'reason': (reason or '')[:500],
|
||
'operator_id': operator_id or '',
|
||
'created_at': _now(),
|
||
})
|
||
await sor.sqlExe("COMMIT", {})
|
||
|
||
|
||
async def _cas_status(sor, ticket_id, from_status, sets, extra_where=''):
|
||
"""乐观锁状态迁移:UPDATE ... WHERE status=from_status,再 SELECT 验证。
|
||
|
||
sets: dict 列名→值(status 必含)。返回 True=迁移成功。
|
||
"""
|
||
set_sql = ', '.join('%s=${v_%s}$' % (k, k) for k in sets)
|
||
params = {'v_' + k: v for k, v in sets.items()}
|
||
params['tid'] = ticket_id
|
||
params['fs'] = from_status
|
||
await sor.sqlExe(
|
||
'UPDATE tk_tickets SET %s, updated_at=NOW() WHERE id=${tid}$ AND status=${fs}$ %s'
|
||
% (set_sql, extra_where), params)
|
||
check = await sor.sqlExe(
|
||
"SELECT id FROM tk_tickets WHERE id=${tid}$ AND status=${ns}$",
|
||
{"tid": ticket_id, "ns": sets.get('status')})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return bool(check)
|
||
|
||
|
||
async def _load_ticket(sor, ticket_id):
|
||
recs = await sor.sqlExe(
|
||
"SELECT * FROM tk_tickets WHERE id=${t}$", {"t": ticket_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return recs[0] if recs else None
|
||
|
||
|
||
def _rec_to_dict(rec):
|
||
d = {}
|
||
for k in getattr(rec, 'keys', lambda: [])():
|
||
v = rec[k]
|
||
if isinstance(v, (datetime.datetime, datetime.date)):
|
||
v = str(v)
|
||
d[k] = v
|
||
return d
|
||
|
||
|
||
# ══════════════════ 客户动作 ══════════════════
|
||
|
||
async def create_ticket(user_id, org_id, title, description, category='other',
|
||
priority='normal', attachments=None):
|
||
"""R1 客户建单。返回 (True, ticket dict) 或 (False, (code, msg))。
|
||
|
||
attachments: 客户端上传后的 FileStorage webpath(str 或 list),
|
||
服务端 normalize_attachments 校验+规整为 [{id,name,webpath,size}]。
|
||
"""
|
||
title = (title or '').strip()
|
||
description = (description or '').strip()
|
||
if not title:
|
||
return False, ('TK_E010', '请填写问题标题')
|
||
if not description:
|
||
return False, ('TK_E011', '请填写问题描述')
|
||
if not user_id:
|
||
return False, ('TK_E012', '请先登录')
|
||
atts, att_err = normalize_attachments(attachments)
|
||
if att_err:
|
||
return False, att_err
|
||
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
for _attempt in range(3): # ticket_no 唯一索引冲突重试
|
||
no = await _next_ticket_no(sor)
|
||
tid = getID()
|
||
try:
|
||
await sor.C('tk_tickets', {
|
||
'id': tid,
|
||
'ticket_no': no,
|
||
'title': title[:200],
|
||
'description': description,
|
||
'category': category or 'other',
|
||
'priority': priority or 'normal',
|
||
'status': S_NEW,
|
||
'customer_org_id': str(org_id or ''),
|
||
'customer_user_id': str(user_id),
|
||
'assignee_role': None,
|
||
'assignee_id': None,
|
||
'processing_owner': None,
|
||
'followup_count': 0,
|
||
'agent_fail_rounds': 0,
|
||
'closed_at': None,
|
||
'close_reason': None,
|
||
'created_at': _now(),
|
||
'updated_at': _now(),
|
||
})
|
||
await sor.sqlExe("COMMIT", {})
|
||
break
|
||
except Exception as e:
|
||
await sor.sqlExe("ROLLBACK", {})
|
||
if '1062' in str(getattr(e, 'args', [''])) or 'Duplicate' in str(e):
|
||
continue
|
||
logger.warning("create_ticket failed: %s", str(e)[:200])
|
||
return False, ('TK_E013', '工单创建失败:%s' % str(e)[:200])
|
||
else:
|
||
return False, ('TK_E013', '工单编号生成冲突,请重试')
|
||
|
||
await _add_message(sor, tid, 'customer', user_id, '', description,
|
||
attachments=atts or None)
|
||
logger.info("ticket created: %s (%s) by %s", tid, no, user_id)
|
||
return True, {'id': tid, 'ticket_no': no, 'status': S_NEW}
|
||
|
||
|
||
async def _check_customer_access(sor, t, user_id, org_id, roles):
|
||
"""I5 客户可见性校验。admin 豁免(代客户处理场景)。"""
|
||
if _is_admin(roles):
|
||
return True
|
||
return (str(getattr(t, 'customer_user_id', '')) == str(user_id)) or \
|
||
(org_id and str(getattr(t, 'customer_org_id', '')) == str(org_id))
|
||
|
||
|
||
async def ticket_followup(ticket_id, user_id, org_id, content, attachments=None):
|
||
"""T6/T10 客户追问。agent_replied→agent_processing(超限转人工);
|
||
staff_replied→human_processing(回当前受理人)。可带附件。"""
|
||
content = (content or '').strip()
|
||
if not content:
|
||
return False, ('TK_E010', '请填写追问内容')
|
||
atts, att_err = normalize_attachments(attachments)
|
||
if att_err:
|
||
return False, att_err
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t or not await _check_customer_access(sor, t, user_id, org_id, roles):
|
||
return False, E_NOT_FOUND
|
||
status = str(getattr(t, 'status', ''))
|
||
if status not in (S_AGENT_REPLIED, S_STAFF_REPLIED):
|
||
return False, ('TK_E002', '工单当前状态为 %s,不能追问(需状态 agent_replied/staff_replied)' % status)
|
||
|
||
await _add_message(sor, ticket_id, 'customer', user_id, '', content,
|
||
attachments=atts or None)
|
||
|
||
if status == S_STAFF_REPLIED:
|
||
# T10:回当前受理人(human_processing,assignee 不变)
|
||
ok = await _cas_status(sor, ticket_id, S_STAFF_REPLIED, {'status': S_HUMAN_PROCESSING})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
return True, {'status': S_HUMAN_PROCESSING, 'message': '已通知受理人'}
|
||
|
||
# T6:agent 追问。超限→转人工(I7 防死循环)
|
||
followup_max = int(await _get_param(sor, 'ticket_followup_max', '2'))
|
||
cur = int(getattr(t, 'followup_count', 0) or 0) + 1
|
||
if cur > followup_max:
|
||
default_role = await _get_param(sor, 'ticket_default_role', 'owner.maintainer')
|
||
ok = await _cas_status(sor, ticket_id, S_AGENT_REPLIED, {
|
||
'status': S_HUMAN_PENDING, 'assignee_role': default_role,
|
||
'assignee_id': None, 'followup_count': cur})
|
||
if ok:
|
||
await _add_message(sor, ticket_id, 'agent', AGENT_USER, '',
|
||
'客户追问超过 %d 次,自动转人工处理。' % followup_max,
|
||
visibility='internal')
|
||
await _add_transfer(sor, ticket_id, 'escalate_to_human',
|
||
from_role='', from_user=AGENT_USER,
|
||
to_role=default_role,
|
||
reason='客户追问超限(%d/%d)' % (cur, followup_max),
|
||
operator_id=AGENT_USER)
|
||
return True, {'status': S_HUMAN_PENDING,
|
||
'message': '已为您转接人工处理', 'code': 'TK_E006'}
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
|
||
ok = await _cas_status(sor, ticket_id, S_AGENT_REPLIED, {
|
||
'status': S_AGENT_PROCESSING, 'processing_owner': None,
|
||
'followup_count': cur})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
return True, {'status': S_AGENT_PROCESSING, 'message': '已收到追问,正在处理'}
|
||
|
||
|
||
async def ticket_confirm(ticket_id, user_id, org_id):
|
||
"""T5/T9 客户确认解决关闭。agent_replied/staff_replied → closed。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t or not await _check_customer_access(sor, t, user_id, org_id, roles):
|
||
return False, E_NOT_FOUND
|
||
status = str(getattr(t, 'status', ''))
|
||
if status not in (S_AGENT_REPLIED, S_STAFF_REPLIED):
|
||
return False, ('TK_E002', '工单当前状态为 %s,不能确认关闭(需已有回复)' % status)
|
||
ok = await _cas_status(sor, ticket_id, status, {
|
||
'status': S_CLOSED, 'closed_at': _now(), 'close_reason': 'resolved',
|
||
'assignee_role': None, 'assignee_id': None})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
await _add_message(sor, ticket_id, 'customer', user_id, '',
|
||
'问题已解决,工单关闭。')
|
||
logger.info("ticket closed by customer: %s", ticket_id)
|
||
return True, {'status': S_CLOSED, 'message': '工单已关闭'}
|
||
|
||
|
||
async def ticket_cancel(ticket_id, user_id, org_id):
|
||
"""T13 客户取消。agent_processing 中不可取消(避免与 poller 竞态)。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t or not await _check_customer_access(sor, t, user_id, org_id, roles):
|
||
return False, E_NOT_FOUND
|
||
status = str(getattr(t, 'status', ''))
|
||
if status in FINAL_STATES:
|
||
return False, ('TK_E002', '工单已处于终态 %s' % status)
|
||
if status == S_AGENT_PROCESSING:
|
||
return False, ('TK_E002', '智能助手正在处理中,请稍后再取消')
|
||
ok = await _cas_status(sor, ticket_id, status, {
|
||
'status': S_CANCELLED, 'closed_at': _now(), 'close_reason': 'cancelled',
|
||
'assignee_role': None, 'assignee_id': None})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
await _add_message(sor, ticket_id, 'customer', user_id, '', '客户取消工单。')
|
||
return True, {'status': S_CANCELLED, 'message': '工单已取消'}
|
||
|
||
|
||
# ══════════════════ 运维动作 ══════════════════
|
||
|
||
async def _check_staff_access(sor, t, user_id, roles):
|
||
"""运维侧访问:admin/superuser 豁免;否则须持有工单 assignee_role。"""
|
||
if _is_admin(roles):
|
||
return True
|
||
cur_role = str(getattr(t, 'assignee_role', '') or '')
|
||
return bool(cur_role) and cur_role in roles
|
||
|
||
|
||
async def ticket_claim(ticket_id, user_id):
|
||
"""T7 认领:human_pending → human_processing,assignee=自己。乐观锁防抢单。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t:
|
||
return False, E_NOT_FOUND
|
||
if not await _check_staff_access(sor, t, user_id, roles):
|
||
return False, ('TK_E001', '您不是该工单的受理角色,无法认领')
|
||
cur_role = str(getattr(t, 'assignee_role', '') or '')
|
||
ok = await _cas_status(sor, ticket_id, S_HUMAN_PENDING, {
|
||
'status': S_HUMAN_PROCESSING, 'assignee_id': user_id})
|
||
if not ok:
|
||
return False, ('TK_E005', '该工单刚被他人认领,请刷新队列')
|
||
await _add_transfer(sor, ticket_id, 'claim', from_role=cur_role,
|
||
to_role=cur_role, to_user=user_id,
|
||
operator_id=user_id)
|
||
logger.info("ticket claimed: %s by %s", ticket_id, user_id)
|
||
return True, {'status': S_HUMAN_PROCESSING, 'message': '认领成功'}
|
||
|
||
|
||
async def ticket_staff_reply(ticket_id, user_id, content, internal=False,
|
||
attachments=None):
|
||
"""T8 人工回复:human_processing → staff_replied。internal=True 只记内部备注不迁状态。
|
||
可带附件(客户可见回复与内部备注均可)。"""
|
||
content = (content or '').strip()
|
||
if not content:
|
||
return False, ('TK_E010', '请填写回复内容')
|
||
atts, att_err = normalize_attachments(attachments)
|
||
if att_err:
|
||
return False, att_err
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t:
|
||
return False, E_NOT_FOUND
|
||
cur_assignee = str(getattr(t, 'assignee_id', '') or '')
|
||
if not (_is_admin(roles) or cur_assignee == str(user_id)):
|
||
return False, ('TK_E001', '只有当前受理人可以回复该工单')
|
||
status = str(getattr(t, 'status', ''))
|
||
if status != S_HUMAN_PROCESSING:
|
||
return False, ('TK_E002', '工单当前状态为 %s,不能回复(需 human_processing)' % status)
|
||
|
||
sender_role = cur_assignee and _role_of(roles) or ''
|
||
await _add_message(sor, ticket_id, 'staff', user_id, sender_role, content,
|
||
visibility='internal' if internal else 'customer',
|
||
attachments=atts or None)
|
||
if internal:
|
||
return True, {'status': status, 'message': '内部备注已记录'}
|
||
ok = await _cas_status(sor, ticket_id, S_HUMAN_PROCESSING,
|
||
{'status': S_STAFF_REPLIED})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
logger.info("ticket replied by staff: %s", ticket_id)
|
||
return True, {'status': S_STAFF_REPLIED, 'message': '已回复客户'}
|
||
|
||
|
||
def _role_of(roles):
|
||
"""staff 消息记录用角色:优先 maintainer,否则第一个。"""
|
||
for r in roles:
|
||
if r.endswith('.maintainer'):
|
||
return r
|
||
return roles[0] if roles else ''
|
||
|
||
|
||
async def ticket_transfer(ticket_id, user_id, to_role='', to_user='', reason=''):
|
||
"""T11/T12 转派。to_user 指定→转具体人(T12);否则 to_role→转角色池(T11)。
|
||
|
||
候选校验(D4):to_role ∈ owner 组织类型角色集;to_user ∈ '0' 机构用户。
|
||
"""
|
||
reason = (reason or '').strip()
|
||
if not reason:
|
||
return False, ('TK_E010', '请填写转派原因')
|
||
if not to_role and not to_user:
|
||
return False, ('TK_E010', '请选择转派目标(角色或人员)')
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t:
|
||
return False, E_NOT_FOUND
|
||
if not await _check_staff_access(sor, t, user_id, roles):
|
||
return False, ('TK_E001', '您不是该工单的受理角色,无法转派')
|
||
status = str(getattr(t, 'status', ''))
|
||
if status != S_HUMAN_PROCESSING:
|
||
return False, ('TK_E002', '工单当前状态为 %s,不能转派(需认领后 human_processing)' % status)
|
||
cur_role = str(getattr(t, 'assignee_role', '') or '')
|
||
cur_assignee = str(getattr(t, 'assignee_id', '') or '')
|
||
|
||
if to_user:
|
||
# T12 转具体人:校验目标 ∈ '0' 机构用户
|
||
recs = await sor.sqlExe(
|
||
"SELECT id FROM users WHERE id=${u}$ AND orgid='0' AND user_status='0'",
|
||
{"u": to_user})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not recs:
|
||
return False, ('TK_E004', '目标用户不属于平台运维组织,请重新选择')
|
||
# 目标用户的角色(取 owner 类第一个作为工单角色,保持 I2/I3 语义)
|
||
troles = await get_user_roles(sor, to_user)
|
||
owner_roles = [r for r in troles if r.startswith('owner.')]
|
||
# 2026-09-10 修复:目标用户必须持有 owner 角色。原逻辑无角色时静默回退
|
||
# cur_role,产生孤儿工单:角色池无人可认领(目标不持该角色),
|
||
# 目标自己打开管理页也被权限分支挡住——待办在却查无此单。
|
||
if not owner_roles:
|
||
return False, ('TK_E004', '目标用户未持有 owner 组织类型角色,无法受理工单,请改派角色或先为其分配角色')
|
||
new_role = owner_roles[0]
|
||
ok = await _cas_status(sor, ticket_id, S_HUMAN_PROCESSING, {
|
||
'status': S_HUMAN_PROCESSING, 'assignee_role': new_role,
|
||
'assignee_id': to_user})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
await _add_transfer(sor, ticket_id, 'transfer_user',
|
||
from_role=cur_role, from_user=cur_assignee or user_id,
|
||
to_role=new_role, to_user=to_user,
|
||
reason=reason, operator_id=user_id)
|
||
logger.info("ticket transferred to user: %s -> %s", ticket_id, to_user)
|
||
return True, {'status': S_HUMAN_PROCESSING, 'message': '已转派给指定人员'}
|
||
|
||
# T11 转角色池:校验目标 ∈ owner 组织类型角色
|
||
recs = await sor.sqlExe(
|
||
"SELECT id FROM role WHERE orgtypeid='owner' AND name=${n}$",
|
||
{"n": to_role.split('.', 1)[1] if '.' in to_role else to_role})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not recs:
|
||
return False, ('TK_E003', '目标角色不在 owner 组织类型角色列表中,请重新选择')
|
||
full_role = to_role if '.' in to_role else 'owner.' + to_role
|
||
ok = await _cas_status(sor, ticket_id, S_HUMAN_PROCESSING, {
|
||
'status': S_HUMAN_PENDING, 'assignee_role': full_role,
|
||
'assignee_id': None})
|
||
if not ok:
|
||
return False, ('TK_E002', '工单状态刚发生变化,请刷新后重试')
|
||
await _add_transfer(sor, ticket_id, 'transfer_role',
|
||
from_role=cur_role, from_user=cur_assignee or user_id,
|
||
to_role=full_role, reason=reason, operator_id=user_id)
|
||
logger.info("ticket transferred to role: %s -> %s", ticket_id, full_role)
|
||
return True, {'status': S_HUMAN_PENDING, 'message': '已转派至角色池 %s' % full_role}
|
||
|
||
|
||
async def transfer_candidates():
|
||
"""D4 转派候选:owner 组织类型全部角色 + '0' 机构全部用户(动态查询)。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
rrecs = await sor.sqlExe(
|
||
"SELECT name FROM role WHERE orgtypeid='owner' AND name NOT IN ('*','customer') "
|
||
"ORDER BY name", {})
|
||
await sor.sqlExe("COMMIT", {})
|
||
urecs = await sor.sqlExe(
|
||
"SELECT DISTINCT u.id, u.username, u.nick_name FROM users u "
|
||
"JOIN userrole ur ON ur.userid=u.id JOIN role r ON r.id=ur.roleid "
|
||
"WHERE u.orgid='0' AND u.user_status='0' AND r.orgtypeid='owner' "
|
||
"AND r.name NOT IN ('*','customer') ORDER BY u.username", {})
|
||
await sor.sqlExe("COMMIT", {})
|
||
roles = [{'value': 'owner.' + str(getattr(r, 'name', '')),
|
||
'text': 'owner.' + str(getattr(r, 'name', ''))} for r in (rrecs or [])]
|
||
users = [{'value': str(getattr(u, 'id', '')),
|
||
'text': str(getattr(u, 'nick_name', '') or getattr(u, 'username', ''))
|
||
+ '(' + str(getattr(u, 'username', '')) + ')'}
|
||
for u in (urecs or [])]
|
||
return {'roles': roles, 'users': users}
|
||
|
||
|
||
# ══════════════════ 查询 ══════════════════
|
||
|
||
async def list_customer_tickets(user_id, org_id, roles=None, status='', limit=100):
|
||
"""客户工单列表(I5:本人或本机构;admin 全量)。roles=None 时自查。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
if roles is None:
|
||
roles = await get_user_roles(sor, user_id)
|
||
where = []
|
||
params = {"lim": int(limit)}
|
||
if not _is_admin(roles):
|
||
where.append("(customer_user_id=${u}$ OR customer_org_id=${o}$)")
|
||
params["u"] = str(user_id or '')
|
||
params["o"] = str(org_id or '')
|
||
if status:
|
||
where.append("status=${st}$")
|
||
params["st"] = status
|
||
wsql = (' WHERE ' + ' AND '.join(where)) if where else ''
|
||
recs = await sor.sqlExe(
|
||
"SELECT id, ticket_no, title, category, priority, status, assignee_role, "
|
||
"followup_count, closed_at, close_reason, created_at, updated_at "
|
||
"FROM tk_tickets" + wsql + " ORDER BY created_at DESC LIMIT ${lim}$", params)
|
||
await sor.sqlExe("COMMIT", {})
|
||
return [_rec_to_dict(r) for r in (recs or [])]
|
||
|
||
|
||
async def list_staff_tickets(user_id, roles=None, queue='all', limit=100):
|
||
"""运维工单列表。queue=pending(待认领:我角色池)/mine(我受理的)/all(全部,管理权限)。
|
||
roles=None 时自查。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
if roles is None:
|
||
roles = await get_user_roles(sor, user_id)
|
||
params = {"lim": int(limit)}
|
||
if queue == 'pending':
|
||
if not roles:
|
||
return []
|
||
ph = ','.join("${role%d}$" % i for i in range(len(roles)))
|
||
for i, r in enumerate(roles):
|
||
params['role%d' % i] = r
|
||
cond = ("status='human_pending' AND assignee_role IN (%s)" % ph)
|
||
elif queue == 'mine':
|
||
cond = ("status IN ('human_processing','staff_replied') "
|
||
"AND assignee_id=${u}$")
|
||
params["u"] = str(user_id or '')
|
||
else:
|
||
# all: admin/owner 角色看全量;其他账号至少能看到自己受理的工单
|
||
# (2026-09-10 修复:受理人若无 owner 角色,原逻辑直接返回[],
|
||
# 受理人看不到自己手上的工单——待办点进来却查无此单)
|
||
if _is_admin(roles) or any(r.startswith('owner.') for r in roles):
|
||
cond = ("status NOT IN ('closed','cancelled') "
|
||
"OR closed_at >= DATE_SUB(NOW(), INTERVAL 30 DAY)")
|
||
else:
|
||
cond = "assignee_id=${u}$"
|
||
params["u"] = str(user_id or '')
|
||
recs = await sor.sqlExe(
|
||
"SELECT t.id, t.ticket_no, t.title, t.category, t.priority, t.status, "
|
||
"t.assignee_role, t.assignee_id, t.customer_org_id, t.created_at, "
|
||
"t.updated_at, o.orgname AS customer_org_name, u.username AS assignee_name "
|
||
"FROM tk_tickets t "
|
||
"LEFT JOIN organization o ON t.customer_org_id=o.id "
|
||
"LEFT JOIN users u ON t.assignee_id=u.id "
|
||
"WHERE " + cond + " ORDER BY t.updated_at DESC LIMIT ${lim}$", params)
|
||
await sor.sqlExe("COMMIT", {})
|
||
return [_rec_to_dict(r) for r in (recs or [])]
|
||
|
||
|
||
async def ticket_detail(ticket_id, user_id, org_id, roles=None):
|
||
"""工单详情 + 消息流 + 转派流水(I5:客户只见 customer 消息;staff 全量)。
|
||
roles=None 时自查。"""
|
||
db = _get_db()
|
||
async with db.sqlorContext(_dbname()) as sor:
|
||
if roles is None:
|
||
roles = await get_user_roles(sor, user_id)
|
||
t = await _load_ticket(sor, ticket_id)
|
||
if not t:
|
||
return False, E_NOT_FOUND
|
||
is_customer_view = not (await _check_staff_access(sor, t, user_id, roles) or _is_admin(roles))
|
||
if is_customer_view and not await _check_customer_access(sor, t, user_id, org_id, roles):
|
||
return False, E_NOT_FOUND
|
||
|
||
vis_cond = '' if not is_customer_view else " AND visibility='customer'"
|
||
mrecs = await sor.sqlExe(
|
||
"SELECT m.id, m.sender_type, m.sender_id, m.sender_role, m.content, "
|
||
"m.attachments, m.visibility, m.created_at, u.username AS sender_name, "
|
||
"u.nick_name AS sender_nick "
|
||
"FROM tk_messages m LEFT JOIN users u ON m.sender_id=u.id "
|
||
"WHERE m.ticket_id=${t}$" + vis_cond + " ORDER BY m.created_at", {"t": ticket_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
messages = [_rec_to_dict(r) for r in (mrecs or [])]
|
||
|
||
transfers = []
|
||
if not is_customer_view:
|
||
trecs = await sor.sqlExe(
|
||
"SELECT * FROM tk_transfers WHERE ticket_id=${t}$ ORDER BY created_at",
|
||
{"t": ticket_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
transfers = [_rec_to_dict(r) for r in (trecs or [])]
|
||
|
||
d = _rec_to_dict(t)
|
||
d['messages'] = messages
|
||
d['transfers'] = transfers
|
||
d['view'] = 'customer' if is_customer_view else 'staff'
|
||
return True, d
|