2026-09-19 12:22:21 +08:00

271 lines
14 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单事务事件+状态写入、按月分区、提交后广播、3s 轮询兜底)是**权威实现**
# 同名契约在此覆盖上方 M11a 版本,避免两套逻辑漂移;返回结构向后兼容(只增字段)。
# 主链路m11b_api → tx_write.apply_runtime_event一个事务写 事件+实体状态+快照
# → 提交后才广播 → 分段计时裁决/广播/端到端并与 SLA 预算比对)。
from .m11b_api import ( # noqa: E402 契约层覆盖(放最后,导入即生效)
contracts as m11b_contracts,
pbl_entity_state_apply,
pbl_entity_state_get,
pbl_runtime_broadcast_pull,
pbl_runtime_event_append,
pbl_runtime_broadcast_pull,
pbl_runtime_broadcast_stats,
pbl_runtime_event_poll,
pbl_world_broadcast,
)
CONTRACTS = {
"pbl_runtime_event_append": pbl_runtime_event_append,
"pbl_runtime_event_poll": pbl_runtime_event_poll,
"pbl_runtime_broadcast_pull": pbl_runtime_broadcast_pull,
"pbl_runtime_broadcast_stats": pbl_runtime_broadcast_stats,
"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,
}