+51
![autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>](/assets/img/avatar_default.png)

![gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>](/assets/img/avatar_default.png)






qiuqiua
GitHub
QuantumGhost
盐粒 Yanli
wangxiaolei
Stephen Zhou
gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
Cursx
autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
lif
非法操作
Asuka Minato
fenglin
qiaofenglin
-LAN-
TomoOkuyama
Tomo Okuyama
dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
zyssyz123
hj24
Coding On Star
CodingOnStar
yyh
Xiangxuan Qu
fghpdf
coopercoder
zhaiguangpeng
Junyan Qin
E.G
GlobalStar117
Claude Haiku 4.5
CodingOnStar
crazywoola
heyszt
NeatGuyCoding
Yeuoly
zxhlyh
moonpanda
warlocgao
github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com>
claude[bot] <41898282+claude[bot]@users.noreply.github.com>
KVOJJJin
eux
bangjiehan
FFXN
Jyong
Nie Ronghua
JQSevenMiao
jiasiqi
Seokrin Taron Sung
CrabSAMA
Copilot
yihong
Joel
Wu Tianwei
yessenia
Jax
niveshdandyan
OSS Contributor
niveshdandyan
Sean Kenneth Doherty
9ef6b90843
Signed-off-by: majiayu000 <[email protected]> Signed-off-by: dependabot[bot] <[email protected]> Signed-off-by: NeatGuyCoding <[email protected]> Signed-off-by: -LAN- <[email protected]> Signed-off-by: yihong0618 <[email protected]> Co-authored-by: QuantumGhost <[email protected]> Co-authored-by: 盐粒 Yanli <[email protected]> Co-authored-by: wangxiaolei <[email protected]> Co-authored-by: Stephen Zhou <[email protected]> Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> Co-authored-by: Cursx <[email protected]> Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> Co-authored-by: lif <[email protected]> Co-authored-by: 非法操作 <[email protected]> Co-authored-by: Asuka Minato <[email protected]> Co-authored-by: fenglin <[email protected]> Co-authored-by: qiaofenglin <[email protected]> Co-authored-by: -LAN- <[email protected]> Co-authored-by: TomoOkuyama <[email protected]> Co-authored-by: Tomo Okuyama <[email protected]> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: zyssyz123 <[email protected]> Co-authored-by: hj24 <[email protected]> Co-authored-by: Coding On Star <[email protected]> Co-authored-by: CodingOnStar <[email protected]> Co-authored-by: yyh <[email protected]> Co-authored-by: Xiangxuan Qu <[email protected]> Co-authored-by: fghpdf <[email protected]> Co-authored-by: coopercoder <[email protected]> Co-authored-by: zhaiguangpeng <[email protected]> Co-authored-by: Junyan Qin (Chin) <[email protected]> Co-authored-by: E.G <[email protected]> Co-authored-by: GlobalStar117 <[email protected]> Co-authored-by: Claude Haiku 4.5 <[email protected]> Co-authored-by: CodingOnStar <[email protected]> Co-authored-by: crazywoola <[email protected]> Co-authored-by: heyszt <[email protected]> Co-authored-by: NeatGuyCoding <[email protected]> Co-authored-by: Yeuoly <[email protected]> Co-authored-by: zxhlyh <[email protected]> Co-authored-by: moonpanda <[email protected]> Co-authored-by: warlocgao <[email protected]> Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: claude[bot] <41898282+claude[bot]@users.noreply.github.com> Co-authored-by: KVOJJJin <[email protected]> Co-authored-by: eux <[email protected]> Co-authored-by: bangjiehan <[email protected]> Co-authored-by: FFXN <[email protected]> Co-authored-by: Jyong <[email protected]> Co-authored-by: Nie Ronghua <[email protected]> Co-authored-by: JQSevenMiao <[email protected]> Co-authored-by: jiasiqi <[email protected]> Co-authored-by: Seokrin Taron Sung <[email protected]> Co-authored-by: CrabSAMA <[email protected]> Co-authored-by: Copilot <[email protected]> Co-authored-by: yihong <[email protected]> Co-authored-by: Joel <[email protected]> Co-authored-by: Wu Tianwei <[email protected]> Co-authored-by: yessenia <[email protected]> Co-authored-by: Jax <[email protected]> Co-authored-by: niveshdandyan <[email protected]> Co-authored-by: OSS Contributor <[email protected]> Co-authored-by: niveshdandyan <[email protected]> Co-authored-by: Sean Kenneth Doherty <[email protected]>
125 lines
3.8 KiB
Python
125 lines
3.8 KiB
Python
import hashlib
|
|
import logging
|
|
from typing import TypeVar
|
|
|
|
from redis import RedisError
|
|
|
|
from core.trigger.debug.events import BaseDebugEvent
|
|
from extensions.ext_redis import redis_client
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
TRIGGER_DEBUG_EVENT_TTL = 300
|
|
|
|
TTriggerDebugEvent = TypeVar("TTriggerDebugEvent", bound="BaseDebugEvent")
|
|
|
|
|
|
class TriggerDebugEventBus:
|
|
"""
|
|
Unified Redis-based trigger debug service with polling support.
|
|
|
|
Uses {tenant_id} hash tags for Redis Cluster compatibility.
|
|
Supports multiple event types through a generic dispatch/poll interface.
|
|
"""
|
|
|
|
# LUA_SELECT: Atomic poll or register for event
|
|
# KEYS[1] = trigger_debug_inbox:{<tenant_id>}:<address_id>
|
|
# KEYS[2] = trigger_debug_waiting_pool:{<tenant_id>}:...
|
|
# ARGV[1] = address_id
|
|
LUA_SELECT = (
|
|
"local v=redis.call('GET',KEYS[1]);"
|
|
"if v then redis.call('DEL',KEYS[1]);return v end;"
|
|
"redis.call('SADD',KEYS[2],ARGV[1]);"
|
|
f"redis.call('EXPIRE',KEYS[2],{TRIGGER_DEBUG_EVENT_TTL});"
|
|
"return false"
|
|
)
|
|
|
|
# LUA_DISPATCH: Dispatch event to all waiting addresses
|
|
# KEYS[1] = trigger_debug_waiting_pool:{<tenant_id>}:...
|
|
# ARGV[1] = tenant_id
|
|
# ARGV[2] = event_json
|
|
LUA_DISPATCH = (
|
|
"local a=redis.call('SMEMBERS',KEYS[1]);"
|
|
"if #a==0 then return 0 end;"
|
|
"redis.call('DEL',KEYS[1]);"
|
|
"for i=1,#a do "
|
|
f"redis.call('SET','trigger_debug_inbox:{{'..ARGV[1]..'}}'..':'..a[i],ARGV[2],'EX',{TRIGGER_DEBUG_EVENT_TTL});"
|
|
"end;"
|
|
"return #a"
|
|
)
|
|
|
|
@classmethod
|
|
def dispatch(
|
|
cls,
|
|
tenant_id: str,
|
|
event: BaseDebugEvent,
|
|
pool_key: str,
|
|
) -> int:
|
|
"""
|
|
Dispatch event to all waiting addresses in the pool.
|
|
|
|
Args:
|
|
tenant_id: Tenant ID for hash tag
|
|
event: Event object to dispatch
|
|
pool_key: Pool key (generate using build_{?}_pool_key(...))
|
|
|
|
Returns:
|
|
Number of addresses the event was dispatched to
|
|
"""
|
|
event_data = event.model_dump_json()
|
|
try:
|
|
result = redis_client.eval(
|
|
cls.LUA_DISPATCH,
|
|
1,
|
|
pool_key,
|
|
tenant_id,
|
|
event_data,
|
|
)
|
|
return int(result)
|
|
except RedisError:
|
|
logger.exception("Failed to dispatch event to pool: %s", pool_key)
|
|
return 0
|
|
|
|
@classmethod
|
|
def poll(
|
|
cls,
|
|
event_type: type[TTriggerDebugEvent],
|
|
pool_key: str,
|
|
tenant_id: str,
|
|
user_id: str,
|
|
app_id: str,
|
|
node_id: str,
|
|
) -> TTriggerDebugEvent | None:
|
|
"""
|
|
Poll for an event or register to the waiting pool.
|
|
|
|
If an event is available in the inbox, return it immediately.
|
|
Otherwise, register the address to the waiting pool for future dispatch.
|
|
|
|
Args:
|
|
event_class: Event class for deserialization and type safety
|
|
pool_key: Pool key (generate using build_{?}_pool_key(...))
|
|
tenant_id: Tenant ID
|
|
user_id: User ID for address calculation
|
|
app_id: App ID for address calculation
|
|
node_id: Node ID for address calculation
|
|
|
|
Returns:
|
|
Event object if available, None otherwise
|
|
"""
|
|
address_id: str = hashlib.sha256(f"{user_id}|{app_id}|{node_id}".encode()).hexdigest()
|
|
address: str = f"trigger_debug_inbox:{{{tenant_id}}}:{address_id}"
|
|
|
|
try:
|
|
event_data = redis_client.eval(
|
|
cls.LUA_SELECT,
|
|
2,
|
|
address,
|
|
pool_key,
|
|
address_id,
|
|
)
|
|
return event_type.model_validate_json(json_data=event_data) if event_data else None
|
|
except RedisError:
|
|
logger.exception("Failed to poll event from pool: %s", pool_key)
|
|
return None
|