From fabbe340d6c640da29f0c2f0dc819e30feb51859 Mon Sep 17 00:00:00 2001 From: Abiba Date: Thu, 11 Jun 2026 00:57:29 +0000 Subject: [PATCH] feat(router): Phase 2 - Atomic session token tracking via Redis Lua script --- router/router.py | 20 ++++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/router/router.py b/router/router.py index 9c53957..2dbeceb 100644 --- a/router/router.py +++ b/router/router.py @@ -2,6 +2,19 @@ import os, json, time, logging, traceback, threading, queue, statistics, math import requests, redis from flask import Flask, request, jsonify, Response, stream_with_context + + +# Phase 2: Atomic session token update Redis Lua script +SESSION_LUA_SCRIPT = """ +local key = KEYS[1] +local new_val = tonumber(ARGV[1]) +local current = tonumber(redis.call('GET', key) or 0) +local max_val = math.max(current, new_val) +redis.call('SET', key, max_val, 'EX', 86400) +return max_val +""" + + REDIS_URL = os.environ.get("REDIS_URL", "redis://redis:6379") GPU_MOE_URL = os.environ.get("GPU_MOE_URL", "http://192.168.68.15:8080/v1") GPU_DENSE_URL = os.environ.get("GPU_DENSE_URL", "http://192.168.68.8:8080/v1") @@ -435,15 +448,14 @@ def chat(): # Allow agent to override queue timeout via header q_timeout = int(request.headers.get("X-Queue-Timeout", str(QUEUE_TIMEOUT))) - # Cross-turn context tracking: accumulate tokens per session + # Cross-turn context tracking: accumulate tokens per session (Phase 2: atomic Lua) session_id = request.headers.get("X-Session-Id", "") session_tokens = 0 if session_id and r: try: - prev = int(r.get("session:" + session_id) or 0) current = estimate_tokens(rd.get("messages",[])) - session_tokens = max(prev, current) # context only grows - r.set("session:" + session_id, session_tokens, ex=86400) # TTL 24h + # Atomic GET/MAX/SET via Lua script prevents race conditions + session_tokens = r.eval(SESSION_LUA_SCRIPT, 1, "session:" + session_id, current) except Exception: pass d = route(rd, tier, agent)