# -*- coding: utf-8 -*- """ demo —— W-07 实时演示模块业务逻辑。 数据模型: demo_session 演示会话(world_id/status/camera_mode/viewer_count/extra_json) demo_session_state 实体状态覆盖层(base_json=演示开始基线,demo_json=演示中当前值) demo_hot_script 演示中保存即热加载的脚本 demo_presence 多人预览在线用户 核心原则(状态隔离): 演示期间对实体的改动只写入覆盖层 demo_session_state,绝不写正式 entity/scene/world 表; 结束演示(stop_demo)即删除覆盖层与热加载脚本/在线表 —— 正式数据保持不变。 """ import json import time import random import datetime try: from appPublic.sqlor import DBPools except ImportError: # pragma: no cover from sqlor import DBPools try: from appPublic.log import debug, info, warning, exception except ImportError: # pragma: no cover from appPublic.log import debug, info, warning, exception try: from ahserver.server_env import ServerEnv except ImportError: # pragma: no cover try: from ahserver import ServerEnv except ImportError: ServerEnv = None class DemoParamError(ValueError): """参数非法 → 400""" class DemoNotFoundError(LookupError): """对象不存在 → 404""" class DemoDuplicateError(Exception): """重复操作 → 409""" class DemoStateError(Exception): """会话状态不允许该操作 → 400""" DEMO_CAMERA_MODES = ('orbit', 'topdown', 'firstperson') DEMO_STATUS_RUNNING = 'preview' DEMO_STATUS_PAUSED = 'paused' DEMO_STATUS_ENDED = 'ended' def _now(): return datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S') def _new_id(prefix): return '%s_%d_%d' % (prefix, int(time.time() * 1000), random.randint(1000, 9999)) def _ctx(): if ServerEnv is None: raise RuntimeError('ServerEnv 不可用,demo 模块需在 ahserver 宿主中加载') env = ServerEnv() dbname = env.get_module_dbname('demo') return DBPools().sqlorContext(dbname) def _loads(s): try: return json.loads(s or '{}') except Exception: return {} async def _get_session(sor, session_id): recs = await sor.R('demo_session', {'id': session_id}) return recs[0] if recs else None async def _count_viewers(sor, session_id): rows = await sor.sqlExe( "SELECT COUNT(*) AS cnt FROM demo_presence WHERE session_id=${session_id}$", {'session_id': session_id}) if rows: return int(getattr(rows[0], 'cnt', 0) or 0) return 0 async def _check_world(world_id): """跨模块校验世界是否存在;world 模块未挂载/表缺失时跳过校验(防御式)。""" try: env = ServerEnv() wdb = env.get_module_dbname('world') async with DBPools().sqlorContext(wdb) as sor: recs = await sor.R('world', {'id': world_id}) return bool(recs), True except Exception as e: warning('demo: 世界存在性校验不可用(跳过): %s' % e) return True, False # ---------------- W-07a 开始演示(进入预览态) ---------------- async def start_demo(world_id, session_name='', entities=None): if not world_id or not isinstance(world_id, str): raise DemoParamError('world_id 不能为空') if isinstance(entities, str): try: entities = json.loads(entities) except Exception: raise DemoParamError('entities 不是合法 JSON') if entities is not None and not isinstance(entities, list): raise DemoParamError('entities 必须是数组') ent_items = [] for ent in (entities or []): if not isinstance(ent, dict) or not ent.get('entity_id'): raise DemoParamError('entities 每项必须包含 entity_id') f = ent.get('fields') if isinstance(ent.get('fields'), dict) else {} ent_items.append({ 'entity_id': str(ent['entity_id']), 'entity_name': str(ent.get('entity_name') or ent['entity_id'])[:128], 'fields': f, }) exists, checked = await _check_world(world_id) if checked and not exists: raise DemoNotFoundError('世界不存在: %s' % world_id) sid = _new_id('ds') now = _now() name = (session_name or '')[:128] or ('演示-%s' % world_id) async with _ctx() as sor: recs = await sor.sqlExe( "SELECT id FROM demo_session WHERE world_id=${world_id}$ " "AND status IN ('preview','paused') LIMIT 1", {'world_id': world_id}) if recs: raise DemoDuplicateError('该世界已有进行中的演示会话: %s' % recs[0].id) await sor.C('demo_session', { 'id': sid, 'world_id': world_id, 'session_name': name, 'status': DEMO_STATUS_RUNNING, 'camera_mode': 'orbit', 'owner_id': '0', 'viewer_count': 0, 'extra_json': '{}', 'created_at': now, 'updated_at': now, }) entities_out = [] for item in ent_items: base = json.dumps(item['fields'], ensure_ascii=False) await sor.C('demo_session_state', { 'id': _new_id('st'), 'session_id': sid, 'entity_id': item['entity_id'], 'entity_name': item['entity_name'], 'base_json': base, 'demo_json': base, 'updated_at': now, }) entities_out.append({'entity_id': item['entity_id'], 'entity_name': item['entity_name']}) return { 'session_id': sid, 'world_id': world_id, 'session_name': name, 'status': DEMO_STATUS_RUNNING, 'camera_mode': 'orbit', 'entities': entities_out, } # ---------------- W-07b 结束演示(含状态隔离清理) ---------------- async def stop_demo(session_id): if not session_id: raise DemoParamError('session_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束') now = _now() await sor.U('demo_session', {'id': session_id, 'status': DEMO_STATUS_ENDED, 'updated_at': now}) # 状态隔离:删除覆盖层/热加载脚本/在线表 —— 正式 entity/scene/world 数据从未被修改 await sor.D('demo_session_state', {'session_id': session_id}) await sor.D('demo_hot_script', {'session_id': session_id}) await sor.D('demo_presence', {'session_id': session_id}) return {'session_id': session_id, 'status': DEMO_STATUS_ENDED, 'overlay_cleared': True} # ---------------- W-07e 视角切换 ---------------- async def switch_camera(session_id, camera_mode): if not session_id or not camera_mode: raise DemoParamError('session_id 与 camera_mode 不能为空') if camera_mode not in DEMO_CAMERA_MODES: raise DemoParamError('非法视角模式: %s,可选 %s' % (camera_mode, ','.join(DEMO_CAMERA_MODES))) async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束,无法切换视角') await sor.U('demo_session', {'id': session_id, 'camera_mode': camera_mode, 'updated_at': _now()}) return {'session_id': session_id, 'camera_mode': camera_mode} # ---------------- W-07f 暂停/继续 ---------------- async def pause_demo(session_id, paused): if not session_id: raise DemoParamError('session_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束') new_status = DEMO_STATUS_PAUSED if paused else DEMO_STATUS_RUNNING await sor.U('demo_session', {'id': session_id, 'status': new_status, 'updated_at': _now()}) return {'session_id': session_id, 'status': new_status, 'paused': paused} # ---------------- W-07c 编辑实体(演示覆盖层,<1s 实时反映) ---------------- async def update_demo_entity(session_id, entity_id, entity_name='', fields=None): if not session_id or not entity_id: raise DemoParamError('session_id/entity_id 不能为空') if fields is None or not isinstance(fields, dict): raise DemoParamError('fields 必须是对象') if isinstance(entity_id, str) is False: raise DemoParamError('entity_id 必须是字符串') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束,无法编辑实体') now = _now() demo_json = json.dumps(fields, ensure_ascii=False) sts = await sor.R('demo_session_state', {'session_id': session_id, 'entity_id': str(entity_id)}) if sts: await sor.U('demo_session_state', { 'session_id': session_id, 'entity_id': str(entity_id), 'entity_name': (entity_name or '')[:128] or sts[0].entity_name, 'demo_json': demo_json, 'updated_at': now, }) else: await sor.C('demo_session_state', { 'id': _new_id('st'), 'session_id': session_id, 'entity_id': str(entity_id), 'entity_name': (entity_name or str(entity_id))[:128], 'base_json': demo_json, 'demo_json': demo_json, 'updated_at': now, }) return {'session_id': session_id, 'entity_id': str(entity_id), 'fields': fields, 'updated_at': now} # ---------------- W-07i 查看实体属性 ---------------- async def get_entity_attrs(session_id, entity_id): if not session_id or not entity_id: raise DemoParamError('session_id/entity_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) sts = await sor.R('demo_session_state', {'session_id': session_id, 'entity_id': str(entity_id)}) if not sts: raise DemoNotFoundError('实体不存在于演示会话: %s' % entity_id) st = sts[0] base = _loads(st.base_json) demo = _loads(st.demo_json) merged = dict(base) merged.update(demo) return { 'entity_id': str(entity_id), 'entity_name': st.entity_name or str(entity_id), 'base': base, 'demo': demo, 'merged': merged, } # ---------------- W-07d 脚本保存热加载 ---------------- async def hotload_script(session_id, script_id, script_name='', script_type='python', script_content=''): if not session_id or not script_id: raise DemoParamError('session_id/script_id 不能为空') if not script_content or not str(script_content).strip(): raise DemoParamError('脚本内容不能为空') if script_type not in ('python', 'expr'): raise DemoParamError('非法脚本类型: %s' % script_type) try: compile(str(script_content), '' % script_id, 'exec') except SyntaxError as e: raise DemoParamError('脚本语法错误: %s (第 %s 行)' % (getattr(e, 'msg', str(e)), getattr(e, 'lineno', '?'))) async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束,无法热加载脚本') now = _now() sts = await sor.R('demo_hot_script', {'session_id': session_id, 'script_id': str(script_id)}) if sts: await sor.U('demo_hot_script', { 'session_id': session_id, 'script_id': str(script_id), 'script_name': (script_name or '')[:128] or sts[0].script_name, 'script_type': script_type, 'script_content': str(script_content), 'updated_at': now, }) else: await sor.C('demo_hot_script', { 'id': _new_id('hs'), 'session_id': session_id, 'script_id': str(script_id), 'script_name': (script_name or str(script_id))[:128], 'script_type': script_type, 'script_content': str(script_content), 'updated_at': now, }) return {'session_id': session_id, 'script_id': str(script_id), 'status': 'hot_loaded', 'updated_at': now} # ---------------- W-07j 多人预览 ---------------- async def presence_join(session_id, user_id='', user_name=''): if not session_id: raise DemoParamError('session_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) if rec.status == DEMO_STATUS_ENDED: raise DemoStateError('演示会话已结束,无法加入') uid = user_id or ('anon_%s' % _new_id('u')) now = _now() sts = await sor.R('demo_presence', {'session_id': session_id, 'user_id': uid}) if sts: await sor.U('demo_presence', {'session_id': session_id, 'user_id': uid, 'last_seen_at': now}) else: await sor.C('demo_presence', { 'id': _new_id('pr'), 'session_id': session_id, 'user_id': uid, 'user_name': (user_name or '')[:64] or uid, 'joined_at': now, 'last_seen_at': now, }) cnt = await _count_viewers(sor, session_id) await sor.U('demo_session', {'id': session_id, 'viewer_count': cnt, 'updated_at': now}) return {'session_id': session_id, 'user_id': uid, 'viewer_count': cnt} async def presence_leave(session_id, user_id=''): if not session_id or not user_id: raise DemoParamError('session_id 与 user_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) await sor.D('demo_presence', {'session_id': session_id, 'user_id': user_id}) cnt = await _count_viewers(sor, session_id) await sor.U('demo_session', {'id': session_id, 'viewer_count': cnt, 'updated_at': _now()}) return {'session_id': session_id, 'user_id': user_id, 'viewer_count': cnt} # ---------------- W-07g 全屏(前端行为,会话级记录) ---------------- async def set_fullscreen(session_id, enabled): if not session_id: raise DemoParamError('session_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) extra = _loads(rec.extra_json) extra['fullscreen'] = bool(enabled) await sor.U('demo_session', { 'id': session_id, 'extra_json': json.dumps(extra, ensure_ascii=False), 'updated_at': _now(), }) return {'session_id': session_id, 'fullscreen': bool(enabled)} # ---------------- 状态快照(前端 500ms 轮询,<1s 实时反映) ---------------- async def get_demo_state(session_id): if not session_id: raise DemoParamError('session_id 不能为空') async with _ctx() as sor: rec = await _get_session(sor, session_id) if not rec: raise DemoNotFoundError('演示会话不存在: %s' % session_id) sts = await sor.R('demo_session_state', {'session_id': session_id}) entities = [] for st in sts: entities.append({ 'entity_id': st.entity_id, 'entity_name': st.entity_name, 'state': _loads(st.demo_json), }) scr = await sor.R('demo_hot_script', {'session_id': session_id}) scripts = [{ 'script_id': s.script_id, 'script_name': s.script_name, 'script_type': s.script_type, 'updated_at': s.updated_at, } for s in scr] pre = await sor.R('demo_presence', {'session_id': session_id}) return { 'session_id': rec.id, 'world_id': rec.world_id, 'session_name': rec.session_name, 'status': rec.status, 'camera_mode': rec.camera_mode, 'viewer_count': len(pre), 'entities': entities, 'scripts': scripts, 'updated_at': rec.updated_at, } # ---------------- 进行中会话列表 ---------------- async def list_active_sessions(): async with _ctx() as sor: ns = {'page': 1, 'rows': 100, 'sort': 'created_at', 'order': 'desc'} return await sor.sqlPaging( "SELECT id, world_id, session_name, status, camera_mode, viewer_count, created_at, updated_at " "FROM demo_session WHERE status IN ('preview','paused')", ns)