feat(router): Phase 2 - Atomic session token tracking via Redis Lua script
This commit is contained in:
+16
-4
@@ -2,6 +2,19 @@ import os, json, time, logging, traceback, threading, queue, statistics, math
|
|||||||
import requests, redis
|
import requests, redis
|
||||||
from flask import Flask, request, jsonify, Response, stream_with_context
|
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")
|
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_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")
|
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
|
# Allow agent to override queue timeout via header
|
||||||
q_timeout = int(request.headers.get("X-Queue-Timeout", str(QUEUE_TIMEOUT)))
|
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_id = request.headers.get("X-Session-Id", "")
|
||||||
session_tokens = 0
|
session_tokens = 0
|
||||||
if session_id and r:
|
if session_id and r:
|
||||||
try:
|
try:
|
||||||
prev = int(r.get("session:" + session_id) or 0)
|
|
||||||
current = estimate_tokens(rd.get("messages",[]))
|
current = estimate_tokens(rd.get("messages",[]))
|
||||||
session_tokens = max(prev, current) # context only grows
|
# Atomic GET/MAX/SET via Lua script prevents race conditions
|
||||||
r.set("session:" + session_id, session_tokens, ex=86400) # TTL 24h
|
session_tokens = r.eval(SESSION_LUA_SCRIPT, 1, "session:" + session_id, current)
|
||||||
except Exception: pass
|
except Exception: pass
|
||||||
|
|
||||||
d = route(rd, tier, agent)
|
d = route(rd, tier, agent)
|
||||||
|
|||||||
Reference in New Issue
Block a user