You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Python为LiveKit房间设置结束时间并提前1分钟发警告

语音面试会话开发疑问解答

背景

我使用Python结合LiveKit与多模态代理开发语音面试会话,需实现以下功能:

  • 指定时间自动断开代理连接;
  • 结束前1分钟发送「收尾提示」消息,让代理自然结束会话。

当前实现可正常运行,但存在以下疑问:

  1. 在LiveKit中使用Python调度定时事件是否有更优方案?
  2. 我对关闭与清理的处理是否正确?
  3. 用户断开通话时shutdown hook无法工作,该如何修复?

当前实现代码

import asyncio
import datetime
from datetime import timedelta
import json
import logging

logger = logging.getLogger(__name__)
transcript = []  # 补充定义,避免未定义报错

async def end_interview_at_time(ctx, session, end_time, warning_time, my_shutdown_hook):
    """处理面试会话的定时提醒与关闭逻辑"""
    warning_sent = False

    while True:
        now = datetime.datetime.now()

        # 结束前1分钟发送提醒
        if not warning_sent and now >= warning_time:
            logger.info("发送收尾提示给代理")
            await send_message(ctx, "你还有1分钟时间,请结束会话。")
            warning_sent = True

        # 结束会话
        if now >= end_time:
            logger.info("面试时间结束,开始清理")
            ctx.add_shutdown_callback(my_shutdown_hook)
            await room_manager.cleanup_room(ctx.room.name)
            ctx.shutdown(reason="面试结束")
            break

        await asyncio.sleep(10)

async def entrypoint(ctx):
    metadata = json.loads(ctx.job.metadata)

    async def my_shutdown_hook():
        print("=== 保存对话记录 ===")
        try:
            interview_id = metadata.get("interview_id", 1)
            store_conversation(interview_id, transcript)
        except Exception as e:
            logger.error(f"保存对话记录出错: {e}")

    try:
        await ctx.connect(auto_subscribe=AutoSubscribe.SUBSCRIBE_ALL)
        await ctx.wait_for_participant()

        interview_start = datetime.datetime.now()
        duration = metadata.get("duration_minutes", 10)
        end_time = interview_start + timedelta(minutes=duration + 0.25)
        warning_time = end_time - timedelta(minutes=1)

        # 初始化多模态代理
        model = openai.realtime.RealtimeModel(
            instructions="欢迎语...",
            voice="shimmer",
            modalities=["audio", "text"]
        )
        agent = MultimodalAgent(model=model)
        agent.start(ctx.room)
        session = model.sessions[0]

        # 发送欢迎消息
        await send_message(ctx, "欢迎来到面试,现在开始吧!")

        # 监听代理消息事件
        @agent.on("user_speech_committed")
        def on_user_speech(msg):
            transcript.append({"role": "应聘者", "content": msg})

        @agent.on("agent_speech_committed")
        def on_agent_speech(msg):
            transcript.append({"role": "面试官", "content": msg})

        # 启动面试定时器
        timer_task = asyncio.create_task(
            end_interview_at_time(ctx, session, end_time, warning_time, my_shutdown_hook)
        )
        await timer_task

    except Exception as e:
        logger.error(f"面试出错: {e}")
        ctx.add_shutdown_callback(my_shutdown_hook)
        await room_manager.cleanup_room(ctx.room.name)
        ctx.shutdown(reason="出错")

疑问解答

1. LiveKit中Python定时事件的更优方案

当前轮询(while True+sleep(10))的方式效率较低,推荐使用asyncio原生定时逻辑,通过计算时间差直接sleep到目标时间,避免无意义的轮询:

async def end_interview_at_time(ctx, session, end_time, warning_time, my_shutdown_hook):
    """优化后的定时处理逻辑"""
    now = datetime.datetime.now()
    
    # 等待到提醒时间(如果还没到)
    warning_delay = (warning_time - now).total_seconds()
    if warning_delay > 0:
        await asyncio.sleep(warning_delay)
        logger.info("发送收尾提示给代理")
        await send_message(ctx, "你还有1分钟时间,请结束会话。")
    
    # 等待到结束时间(如果还没到)
    end_delay = (end_time - datetime.datetime.now()).total_seconds()
    if end_delay > 0:
        await asyncio.sleep(end_delay)
    
    # 执行结束清理
    logger.info("面试时间结束,开始清理")
    await agent.stop()  # 先停止代理会话
    await room_manager.cleanup_room(ctx.room.name)
    ctx.shutdown(reason="面试结束")

如果需要更复杂的定时逻辑(如重复任务),可以引入APScheduler库,但对于单会话的定时需求,asyncio原生方案足够轻量高效。

2. 关闭与清理的处理优化

当前实现存在几个可以改进的点:

  • 提前注册shutdown hook:应该在会话初始化阶段就注册hook,而不是到结束时才添加,确保异常退出、用户断开等场景都能触发。比如在ctx.connect后立即添加:
    await ctx.connect(auto_subscribe=AutoSubscribe.SUBSCRIBE_ALL)
    ctx.add_shutdown_callback(my_shutdown_hook)  # 提前注册hook
    
  • 清理顺序调整:正确的清理顺序应为:停止代理会话 → 清理房间 → 触发ctx shutdown,避免资源泄漏。
  • 避免重复注册hook:异常块中无需再次添加hook,只要提前注册过,ctx.shutdown()会自动触发。

3. 用户断开通话时shutdown hook失效的修复

用户断开通话时,LiveKit会触发participant_disconnected事件,需要监听该事件并主动执行清理与shutdown逻辑,同时取消原定时任务:

async def entrypoint(ctx):
    metadata = json.loads(ctx.job.metadata)
    timer_task = None  # 保存定时任务引用

    def my_shutdown_hook():
        print("=== 保存对话记录 ===")
        try:
            interview_id = metadata.get("interview_id", 1)
            store_conversation(interview_id, transcript)
        except Exception as e:
            logger.error(f"保存对话记录出错: {e}")

    try:
        await ctx.connect(auto_subscribe=AutoSubscribe.SUBSCRIBE_ALL)
        ctx.add_shutdown_callback(my_shutdown_hook)  # 提前注册hook

        # 监听用户断开事件
        @ctx.room.on("participant_disconnected")
        async def on_participant_disconnected(participant):
            nonlocal timer_task
            logger.info(f"用户 {participant.identity} 断开连接")
            if timer_task:
                timer_task.cancel()  # 取消定时任务,避免重复执行
            await agent.stop()
            await room_manager.cleanup_room(ctx.room.name)
            ctx.shutdown(reason="用户断开连接")

        await ctx.wait_for_participant()
        # ... 其他初始化逻辑 ...

        # 启动定时器
        timer_task = asyncio.create_task(
            end_interview_at_time(ctx, session, end_time, warning_time, my_shutdown_hook)
        )
        await timer_task

    except Exception as e:
        logger.error(f"面试出错: {e}")
        await agent.stop()
        await room_manager.cleanup_room(ctx.room.name)
        ctx.shutdown(reason="出错")

另外,注意将my_shutdown_hook改为同步函数(如果LiveKit的shutdown callback不支持异步),或者确认SDK会自动await异步回调,否则hook可能无法执行。


内容的提问来源于stack exchange,提问作者Fatima Sayeed Amani

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 11:15:55