从0开始的agent框架 第4章 定时任务
·
从 0 开始的 Agent 框架 第 4 章 定时任务
定时任务完全可以从 Agent 框架里剥离出来,但是这里讲讲我的设计。
众所周知,假设我们将定时任务理解为一个 xx:xx:xx 会到期的任务,而不是 3 秒后到期的任务,那么肯定是时间越近,到期越快。
所以,我们可以维护一个数据结构,支持随机插入,顺序读取,和删去第一个元素。
所以我们考虑维护一个最小堆,堆顶始终是最近到期的任务。
和 Chromium 的设计一致,只不过我们使用 Python。
这样一来就好办了,我们只需要单开一个线程负责等待,等待堆顶需要等待的时间,弹出堆顶任务,然后继续等待就好了。
插入时会麻烦一点,我们大约需要 O(n) 的时间调整堆,但是我们考虑到,不会频繁添加等待任务,所以这也是可接受的。
插入时会打断等待线程的等待,并且会使等待线程重新计算sleep的时间,但是这也只是O(1)时间的等待,所以也可接受。
所以,我们开始写代码,工具调用希望是这样的:
# 在 timed_task_tools.py 中
@registry.tool(name="create_timed_task", description="Create a timed task that triggers after a delay (seconds).")
def create_timed_task(delay_seconds: float, message: str, ) -> str:
payload = {
"type": "reminder",
"message": message,
}
print("create_timed_task", payload)
task_id = task_manager.add_task(delay_seconds, payload)
return f"Task {task_id} created"
def add_task(self, delay_seconds: float, payload: dict, task_id: Optional[str] = None) -> str:
"""
添加一个定时任务
:param delay_seconds: 延迟秒数(可以为浮点数)
:param payload: 任务触发时放入队列的数据(如 {"type": "reminder", "message": "喝水"})
:param task_id: 可选,不提供则自动生成
:return: task_id
"""
if task_id is None:
task_id = str(uuid.uuid4())
trigger_time = time.time() + delay_seconds
item = (trigger_time, task_id, payload)
print("item", item)
with self._heap_lock:
print("in heap_lock", self._heap_lock)
heapq.heappush(self._heap, item)
# print("heap push", heapq.heappop(self._heap))
self._condition.notify() # 唤醒 worker 重新检查堆顶
print("condition notify", self._condition)
# print("")
logger.info(f"Task {task_id} scheduled in {delay_seconds}s")
return task_id
删除时可以选择标记删除或者直接删除,均可。
下一期我准备讲讲流式输出,改善用户体验。
更多推荐


所有评论(0)