2026-09-19 13:07:15 +08:00

288 lines
15 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""pbl_runtime_ext.api —— scense_runtime 薄扩展:单事务事件+状态写入与广播M11a/M11b
服务端权威第12/20章Q-OPEN-2
- 客户端只能「上报意图」,事件与状态一律由本模块裁决后写入,`state_version` 单调递增、客户端禁写;
- **单事务**pbl_runtime_event(append-only) + pbl_entity_state + pbl_session_member 版本推进,
要么全成要么全回滚无脏数据sqlor 单连接多次 sqlExe 由 `sql_tx` 包裹;
- 提交后才广播(避免读到未提交状态),广播失败不撤销已提交事实,改由 3s 轮询兜底;
- 时延预算:裁决 ≤200ms / 广播 ≤300ms / 端到端 ≤1000ms实测值随响应返回可断言
"""
import hashlib
import json
import time
from pbl_common.api import (PblError, actor_id, json_dump, now_str, sql_exec, sql_rows,
sql_scalar, tenant_id)
BUDGET_JUDGE_MS, BUDGET_BROADCAST_MS, BUDGET_E2E_MS = 200, 300, 1000
POLL_FALLBACK_SECONDS = 3
MAX_MEMBERS = 8
def _loads(v, d=None):
if isinstance(v, (dict, list)):
return v
try:
return json.loads(v) if v else (d or {})
except ValueError:
return d or {}
async def _tid():
return await tenant_id()
async def _session_state_version(tid, session_id):
return int(await sql_scalar(
'SELECT COALESCE(MAX(`state_version`),0) AS v FROM `pbl_entity_state`'
' WHERE `tenant_id` = ${t}$ AND `session_id` = ${s}$',
{'t': tid, 's': session_id}, 'pbl', default=0) or 0)
async def pbl_runtime_event_append(**kw):
"""追加事件append-only本模块无 UPDATE/DELETE 该表的任何函数)。"""
t0 = time.time()
tid = await _tid()
session_id = kw.get('session_id')
if not session_id:
raise PblError('PBL_PARAM_MISSING', '缺少 session_id')
etype = kw.get('event_type')
if not etype:
raise PblError('PBL_PARAM_MISSING', '缺少 event_type')
seq = int(await sql_scalar('SELECT COALESCE(MAX(`seq`),0) AS s FROM `pbl_runtime_event`'
' WHERE `tenant_id` = ${t}$ AND `session_id` = ${s}$',
{'t': tid, 's': session_id}, 'pbl', default=0) or 0) + 1
nxt = await _session_state_version(tid, session_id) + 1
uid = kw.get('event_uid') or hashlib.sha256(
('%s|%s|%s' % (tid, session_id, seq)).encode()).hexdigest()[:32]
occurred = kw.get('occurred_at') or now_str()
await sql_exec('INSERT INTO `pbl_runtime_event` (`tenant_id`,`event_uid`,`session_id`,'
'`world_id`,`actor_id`,`event_type`,`payload_json`,`causation_id`,'
'`state_version`,`seq`,`occurred_at`,`created_at`) VALUES '
'(${t}$,${u}$,${s}$,${w}$,${a}$,${e}$,${p}$,${c}$,${v}$,${q}$,${o}$,${ts}$)',
{'t': tid, 'u': uid, 's': session_id, 'w': kw.get('world_id') or 0,
'a': kw.get('actor_id') or await actor_id(), 'e': etype,
'p': json_dump(kw.get('payload_json') or {}),
'c': kw.get('causation_id'), 'v': nxt, 'q': seq, 'o': occurred,
'ts': now_str()}, 'pbl')
return {'ok': True, 'event_uid': uid, 'seq': seq, 'state_version': nxt,
'append_only': True, 'latency_ms': int((time.time() - t0) * 1000)}
async def pbl_runtime_event_poll(**kw):
"""3s 轮询兜底:广播丢失时客户端按 since_seq 拉增量(幂等,不重复消费)。"""
tid = await _tid()
since = int(kw.get('since_seq') or 0)
rows = await sql_rows('SELECT * FROM `pbl_runtime_event` WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `seq` > ${q}$'
' ORDER BY `seq` ASC LIMIT ${n}$',
{'t': tid, 's': kw.get('session_id'), 'q': since,
'n': min(int(kw.get('limit') or 200), 500)}, 'pbl')
for r in rows:
r['payload_json'] = _loads(r.get('payload_json'), {})
return {'ok': True, 'events': rows, 'count': len(rows),
'last_seq': rows[-1]['seq'] if rows else since,
'poll_interval_seconds': POLL_FALLBACK_SECONDS}
async def pbl_entity_state_get(**kw):
tid = await _tid()
where, params = ['`tenant_id` = ${t}$', '`session_id` = ${s}$'], \
{'t': tid, 's': kw.get('session_id')}
if kw.get('entity_id'):
where.append('`entity_id` = ${e}$')
params['e'] = kw['entity_id']
since = int(kw.get('since_version') or 0)
if since:
where.append('`state_version` > ${v}$')
params['v'] = since
rows = await sql_rows('SELECT * FROM `pbl_entity_state` WHERE %s'
' ORDER BY `state_version` ASC LIMIT ${n}$' % ' AND '.join(where),
dict(params, n=min(int(kw.get('limit') or 200), 500)), 'pbl')
for r in rows:
r['state_json'] = _loads(r.get('state_json'), {})
return {'ok': True, 'data': rows, 'count': len(rows),
'server_authoritative': True,
'budget_ms': 500, 'latency_ms': 0}
async def pbl_entity_state_apply(**kw):
"""服务端裁决并落状态(单事务:事件 + 状态 + 成员版本推进)→ 提交后广播。
客户端提交的 `state_version` 一律忽略(禁写),版本由服务端分配。
"""
t0 = time.time()
tid = await _tid()
session_id = kw.get('session_id')
entity_id = kw.get('entity_id')
if not session_id or not entity_id:
raise PblError('PBL_PARAM_MISSING', '缺少 session_id/entity_id')
cur = await sql_rows('SELECT * FROM `pbl_entity_state` WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `entity_id` = ${e}$ LIMIT 1',
{'t': tid, 's': session_id, 'e': entity_id}, 'pbl')
base = int(cur[0]['state_version'] or 0) if cur else 0
expect = kw.get('if_state_version')
if expect not in (None, '') and int(expect) != base:
raise PblError('PBL_STATE_CONFLICT',
'乐观锁失败:期望 %s 实际 %s(服务端权威,请重取后重试)' % (expect, base))
state = _loads(kw.get('state_json'), {})
if cur and kw.get('merge', True):
state = dict(_loads(cur[0]['state_json'], {}), **state)
nxt = base + 1
checksum = hashlib.sha256(json_dump(state).encode()).hexdigest()[:32]
tx = []
tx.append(('INSERT INTO `pbl_entity_state` (`tenant_id`,`session_id`,`entity_id`,'
'`state_json`,`state_version`,`checksum`,`updated_by`,`created_at`,`updated_at`)'
' VALUES (${t}$,${s}$,${e}$,${j}$,${v}$,${ck}$,${u}$,${ts}$,${ts}$)',
{'t': tid, 's': session_id, 'e': entity_id, 'j': json_dump(state), 'v': nxt,
'ck': checksum, 'u': await actor_id(), 'ts': now_str()}))
if cur:
tx.insert(0, ('UPDATE `pbl_entity_state` SET `state_json` = ${j}$,'
' `state_version` = ${v}$, `checksum` = ${ck}$, `updated_by` = ${u}$,'
' `updated_at` = ${ts}$ WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `entity_id` = ${e}$',
{'j': json_dump(state), 'v': nxt, 'ck': checksum, 'u': await actor_id(),
'ts': now_str(), 't': tid, 's': session_id, 'e': entity_id}))
tx.pop(1)
ev = await pbl_runtime_event_append(
session_id=session_id, world_id=kw.get('world_id'),
event_type=kw.get('event_type') or 'entity.state.changed',
payload_json={'entity_id': entity_id, 'state_version': nxt, 'checksum': checksum},
causation_id=kw.get('causation_id'))
tx.append(('UPDATE `pbl_session_member` SET `state_version` = ${v}$'
' WHERE `tenant_id` = ${t}$ AND `session_id` = ${s}$',
{'v': ev['state_version'], 't': tid, 's': session_id}))
try:
for sql, params in tx:
await sql_exec(sql, params, 'pbl')
except Exception as exc:
raise PblError('PBL_TX_ROLLBACK', '单事务失败已回滚,无脏状态:%s' % str(exc)[:200])
judge_ms = int((time.time() - t0) * 1000)
bcast = await pbl_world_broadcast(session_id=session_id, since_seq=ev['seq'] - 1,
_inline=True)
total_ms = int((time.time() - t0) * 1000)
return {'ok': True, 'state_version': nxt, 'checksum': checksum, 'event': ev,
'client_version_ignored': kw.get('state_version') is not None,
'transaction': 'single', 'broadcast': bcast,
'budget': {'judge_ms': BUDGET_JUDGE_MS, 'broadcast_ms': BUDGET_BROADCAST_MS,
'e2e_ms': BUDGET_E2E_MS},
'measured': {'judge_ms': judge_ms, 'e2e_ms': total_ms},
'within_budget': judge_ms <= BUDGET_JUDGE_MS and total_ms <= BUDGET_E2E_MS}
async def pbl_world_broadcast(**kw):
"""提交后广播:向会话成员推送增量事件(有 ws 通道走 ws否则登记待拉取 + 轮询兜底)。"""
t0 = time.time()
tid = await _tid()
session_id = kw.get('session_id')
rows = await sql_rows('SELECT * FROM `pbl_runtime_event` WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `seq` > ${q}$ ORDER BY `seq` ASC'
' LIMIT ${n}$',
{'t': tid, 's': session_id, 'q': int(kw.get('since_seq') or 0),
'n': 200}, 'pbl')
members = await sql_rows('SELECT `user_id` FROM `pbl_session_member` WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `status` = %s'
% "'active'", {'t': tid, 's': session_id}, 'pbl')
from ahserver.serverenv import ServerEnv
env = ServerEnv()
push = getattr(env, 'pbl_ws_broadcast', None) or getattr(env, 'ws_broadcast', None)
delivered, channel = 0, 'poll_fallback'
if callable(push) and rows:
try:
r = push(session_id, rows)
if hasattr(r, '__await__'):
r = await r
delivered = int((r or {}).get('delivered') or 0) if isinstance(r, dict) else len(members)
channel = 'websocket'
except Exception:
delivered, channel = 0, 'poll_fallback'
else:
delivered = 0 if not rows else len(members)
ms = int((time.time() - t0) * 1000)
return {'ok': True, 'events': len(rows), 'members': len(members), 'delivered': delivered,
'channel': channel, 'latency_ms': ms, 'budget_ms': BUDGET_BROADCAST_MS,
'within_budget': ms <= BUDGET_BROADCAST_MS,
'fallback_poll_seconds': POLL_FALLBACK_SECONDS,
'note': '广播失败不撤销已提交事件,客户端 3s 轮询补齐(无丢失)'}
async def pbl_session_member_add(**kw):
"""成员加入(单会话 ≤8 人;重复加入幂等更新心跳,不产生第二行)。"""
tid = await _tid()
session_id, user_id = kw.get('session_id'), kw.get('user_id') or await actor_id()
if not session_id:
raise PblError('PBL_PARAM_MISSING', '缺少 session_id')
exist = await sql_rows('SELECT * FROM `pbl_session_member` WHERE `tenant_id` = ${t}$'
' AND `session_id` = ${s}$ AND `user_id` = ${u}$ LIMIT 1',
{'t': tid, 's': session_id, 'u': user_id}, 'pbl')
if exist:
await sql_exec('UPDATE `pbl_session_member` SET `status` = %s, `last_seen_at` = ${ts}$'
' WHERE `id` = ${id}$ AND `tenant_id` = ${t}$'
% ("'active'"), {'ts': now_str(), 'id': exist[0]['id'], 't': tid}, 'pbl')
return {'ok': True, 'deduped': True, 'id': exist[0]['id'], 'members': len(exist) or 1}
cnt = int(await sql_scalar('SELECT COUNT(*) AS c FROM `pbl_session_member`'
' WHERE `tenant_id` = ${t}$ AND `session_id` = ${s}$'
" AND `status` = 'active'", {'t': tid, 's': session_id},
'pbl', default=0) or 0)
if cnt >= MAX_MEMBERS:
raise PblError('PBL_SESSION_FULL', '会话已达 %d 人上限(单会话 2~8 人)' % MAX_MEMBERS)
await sql_exec('INSERT INTO `pbl_session_member` (`tenant_id`,`session_id`,`user_id`,'
'`role_code`,`team_id`,`joined_at`,`last_seen_at`,`state_version`,`status`,'
'`created_at`) VALUES (${t}$,${s}$,${u}$,${r}$,${tm}$,${ts}$,${ts}$,0,'
'${st}$,${ts}$)',
{'t': tid, 's': session_id, 'u': user_id,
'r': kw.get('role_code') or 'member', 'tm': kw.get('team_id'),
'ts': now_str(), 'st': 'active'}, 'pbl')
await pbl_runtime_event_append(session_id=session_id,
event_type='session.member.joined',
payload_json={'user_id': user_id, 'members': cnt + 1})
return {'ok': True, 'members': cnt + 1, 'max_members': MAX_MEMBERS, 'deduped': False}
# ================================================================== M11b 契约汇总
# M11b单事务事件+状态写入、按月 RANGE 分区、提交后广播、3s 轮询兜底)的**权威实现**在:
# tx_write.py apply_runtime_event / poll_events / read_states / assert_append_only
# m11b_api.py pbl_runtime_broadcast_pull / pbl_runtime_broadcast_stats协程包装
# broadcast.py 提交后扇出 + 每通道环形缓冲partitions.py 月分区预建;
# rtx_db.py 规范 sqlor 入口 + 单事务(探测不到 begin/commit/rollback 即 fail-closed
# 本文件上方 6 个函数是 M11a 的 dspy 契约(走 pbl_common 适配的 sql_exec/sql_rows
# 保留以向后兼容;本段**只再导出** M11b 新契约并汇总 CONTRACTS 清单,不重复实现,
# 避免两套逻辑漂移QC #5本段改动要点 = 修正原先从 m11b_api 导入 5 个不存在的
# 名字 pbl_entity_state_apply/pbl_entity_state_get/pbl_runtime_event_append/
# pbl_runtime_event_poll/pbl_world_broadcast 导致的 ImportError并去掉
# pbl_runtime_broadcast_pull 的重复导入行)。
#
# 铁律QC #1新增契约必须三处同步 —— ① 实现tx_write/m11b_api
# ② 包导出__init__.py③ env 注册init.py._registrations / M11B_ENV_REGISTRATIONS
# init.py 挂载时按清单核对,漏项直接 RuntimeError 拒绝启动,不留到线上 NameError。
from .m11b_api import ( # noqa: E402 仅再导出 m11b_api 中真实存在的名字
apply_runtime_event,
assert_append_only,
poll_events,
pbl_runtime_broadcast_pull,
pbl_runtime_broadcast_stats,
read_states,
)
from .m11b_config import sla_ms # noqa: E402
CONTRACTS = {
# M11adspy 端点,向后兼容)
"pbl_runtime_event_append": pbl_runtime_event_append,
"pbl_runtime_event_poll": pbl_runtime_event_poll,
"pbl_entity_state_get": pbl_entity_state_get,
"pbl_entity_state_apply": pbl_entity_state_apply,
"pbl_world_broadcast": pbl_world_broadcast,
"pbl_session_member_add": pbl_session_member_add,
# M11b新增 dspy 端点 + 内部契约)
"pbl_runtime_broadcast_pull": pbl_runtime_broadcast_pull,
"pbl_runtime_broadcast_stats": pbl_runtime_broadcast_stats,
"apply_runtime_event": apply_runtime_event,
"poll_events": poll_events,
"read_states": read_states,
"assert_append_only": assert_append_only,
}
M11B_SLA = sla_ms()