Compare commits
2 Commits
395475628c
...
39b7df2ce0
| Author | SHA1 | Date | |
|---|---|---|---|
| 39b7df2ce0 | |||
| 6b9fb39807 |
@ -6,7 +6,7 @@ import json
|
|||||||
from yuxi.channel.domain.middleware.configurable import Configurable
|
from yuxi.channel.domain.middleware.configurable import Configurable
|
||||||
from yuxi.channel.domain.port.config_reload_port import ConfigReloadPort
|
from yuxi.channel.domain.port.config_reload_port import ConfigReloadPort
|
||||||
from yuxi.channel.domain.service.pipeline import Pipeline
|
from yuxi.channel.domain.service.pipeline import Pipeline
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class ConfigService:
|
class ConfigService:
|
||||||
|
|||||||
@ -9,7 +9,7 @@ from yuxi.channel.domain.model.message.dispatch_result import SendResult
|
|||||||
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
||||||
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
||||||
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@ -57,7 +57,7 @@ class FeishuAdapter:
|
|||||||
return self._ws
|
return self._ws
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def route_contributor(self) -> ChannelRouteContributor | None:
|
def route_contributor(self) -> ChannelRouteContributorPort | None:
|
||||||
return _FeishuRouteContributor()
|
return _FeishuRouteContributor()
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
|
|||||||
@ -8,7 +8,7 @@ from yuxi.channel.domain.model.message.dispatch_result import SendResult
|
|||||||
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
||||||
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
||||||
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@ -55,7 +55,7 @@ class HooksAdapter:
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def route_contributor(self) -> ChannelRouteContributor | None:
|
def route_contributor(self) -> ChannelRouteContributorPort | None:
|
||||||
return _HooksRouteContributor()
|
return _HooksRouteContributor()
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
|
|||||||
@ -8,7 +8,7 @@ from yuxi.channel.domain.model.message.dispatch_result import SendResult
|
|||||||
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
from yuxi.channel.domain.model.message.unified_message import UnifiedMessage
|
||||||
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
from yuxi.channel.domain.model.shared.channel_capabilities import ChannelCapabilities
|
||||||
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
from yuxi.channel.domain.model.shared.channel_type import ChannelType
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
from yuxi.channel.domain.port.sse_push_port import SsePushPort
|
from yuxi.channel.domain.port.sse_push_port import SsePushPort
|
||||||
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
||||||
|
|
||||||
@ -41,7 +41,7 @@ class WebAdapter:
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def route_contributor(self) -> ChannelRouteContributor | None:
|
def route_contributor(self) -> ChannelRouteContributorPort | None:
|
||||||
return _WebRouteContributor()
|
return _WebRouteContributor()
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
|
|||||||
@ -37,12 +37,12 @@ from yuxi.channel.domain.repository.outbox_repository import OutboxRepositoryPor
|
|||||||
from yuxi.channel.domain.repository.session_repository import SessionRepositoryPort
|
from yuxi.channel.domain.repository.session_repository import SessionRepositoryPort
|
||||||
from yuxi.channel.domain.service.pipeline import Pipeline
|
from yuxi.channel.domain.service.pipeline import Pipeline
|
||||||
from yuxi.channel.infrastructure.agent.agent_adapter import AgentAdapter
|
from yuxi.channel.infrastructure.agent.agent_adapter import AgentAdapter
|
||||||
from yuxi.channel.infrastructure.cache.redis_bot_loop_guard import RedisBotLoopGuard
|
from yuxi.channel.infrastructure.cache_infra.redis_bot_loop_guard import RedisBotLoopGuard
|
||||||
from yuxi.channel.infrastructure.cache.redis_cache import RedisCache
|
from yuxi.channel.infrastructure.cache_infra.redis_cache import RedisCache
|
||||||
from yuxi.channel.infrastructure.cache.redis_circuit_breaker import RedisCircuitBreaker
|
from yuxi.channel.infrastructure.cache_infra.redis_circuit_breaker import RedisCircuitBreaker
|
||||||
from yuxi.channel.infrastructure.cache.redis_rate_limiter import RedisRateLimiter
|
from yuxi.channel.infrastructure.cache_infra.redis_rate_limiter import RedisRateLimiter
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
from yuxi.channel.infrastructure.config.redis_config_reload import RedisConfigReload
|
from yuxi.channel.infrastructure.configuration.redis_config_reload import RedisConfigReload
|
||||||
from yuxi.channel.infrastructure.content_filter.composite_content_filter import CompositeContentFilter
|
from yuxi.channel.infrastructure.content_filter.composite_content_filter import CompositeContentFilter
|
||||||
from yuxi.channel.infrastructure.content_filter.llm_content_filter import LlmContentFilter
|
from yuxi.channel.infrastructure.content_filter.llm_content_filter import LlmContentFilter
|
||||||
from yuxi.channel.infrastructure.content_filter.redis_content_filter import RedisContentFilter
|
from yuxi.channel.infrastructure.content_filter.redis_content_filter import RedisContentFilter
|
||||||
|
|||||||
@ -29,7 +29,7 @@ _LAZY_IMPORTS = {
|
|||||||
"KeywordMatcherPort": "yuxi.channel.domain.port",
|
"KeywordMatcherPort": "yuxi.channel.domain.port",
|
||||||
"CachePort": "yuxi.channel.domain.port",
|
"CachePort": "yuxi.channel.domain.port",
|
||||||
"ChannelAdapterPort": "yuxi.channel.domain.port",
|
"ChannelAdapterPort": "yuxi.channel.domain.port",
|
||||||
"ChannelRouteContributor": "yuxi.channel.domain.port",
|
"ChannelRouteContributorPort": "yuxi.channel.domain.port",
|
||||||
"CircuitBreakerPort": "yuxi.channel.domain.port",
|
"CircuitBreakerPort": "yuxi.channel.domain.port",
|
||||||
"ConfigReloadPort": "yuxi.channel.domain.port",
|
"ConfigReloadPort": "yuxi.channel.domain.port",
|
||||||
"ContentFilterPort": "yuxi.channel.domain.port",
|
"ContentFilterPort": "yuxi.channel.domain.port",
|
||||||
|
|||||||
@ -1,16 +1,22 @@
|
|||||||
from yuxi.channel.domain.port.agent_port import AgentPort
|
from yuxi.channel.domain.port.external import (
|
||||||
from yuxi.channel.domain.port.bot_loop_guard_port import BotLoopGuardPort
|
ChannelAdapterPort,
|
||||||
from yuxi.channel.domain.port.cache_port import CachePort
|
ChannelRouteContributorPort,
|
||||||
from yuxi.channel.domain.port.channel_adapter_port import ChannelAdapterPort
|
SignatureVerifyPort,
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
WsConnectionPort,
|
||||||
from yuxi.channel.domain.port.circuit_breaker_port import CircuitBreakerPort
|
)
|
||||||
from yuxi.channel.domain.port.config_reload_port import ConfigReloadPort
|
from yuxi.channel.domain.port.internal import (
|
||||||
from yuxi.channel.domain.port.content_filter_port import ContentFilterPort, FilterResult
|
AgentPort,
|
||||||
from yuxi.channel.domain.port.event_publisher_port import DomainEvent, EventPublisherPort
|
BotLoopGuardPort,
|
||||||
from yuxi.channel.domain.port.keyword_matcher_port import KeywordMatcherPort
|
CachePort,
|
||||||
from yuxi.channel.domain.port.metrics_port import MetricsPort
|
CircuitBreakerPort,
|
||||||
from yuxi.channel.domain.port.queue_port import QueuePort
|
ConfigReloadPort,
|
||||||
from yuxi.channel.domain.port.rate_limit_port import RateLimitPort
|
ContentFilterPort,
|
||||||
from yuxi.channel.domain.port.signature_verify_port import SignatureVerifyPort
|
DomainEvent,
|
||||||
from yuxi.channel.domain.port.sse_push_port import SsePushPort
|
EventPublisherPort,
|
||||||
from yuxi.channel.domain.port.ws_connection_port import WsConnectionPort
|
FilterResult,
|
||||||
|
KeywordMatcherPort,
|
||||||
|
MetricsPort,
|
||||||
|
QueuePort,
|
||||||
|
RateLimitPort,
|
||||||
|
SsePushPort,
|
||||||
|
)
|
||||||
|
|||||||
4
backend/package/yuxi/channel/domain/port/external/__init__.py
vendored
Normal file
4
backend/package/yuxi/channel/domain/port/external/__init__.py
vendored
Normal file
@ -0,0 +1,4 @@
|
|||||||
|
from yuxi.channel.domain.port.external.channel_adapter_port import ChannelAdapterPort
|
||||||
|
from yuxi.channel.domain.port.external.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
|
from yuxi.channel.domain.port.external.signature_verify_port import SignatureVerifyPort
|
||||||
|
from yuxi.channel.domain.port.external.ws_connection_port import WsConnectionPort
|
||||||
@ -4,6 +4,6 @@ from typing import Protocol, runtime_checkable
|
|||||||
|
|
||||||
|
|
||||||
@runtime_checkable
|
@runtime_checkable
|
||||||
class ChannelRouteContributor(Protocol):
|
class ChannelRouteContributorPort(Protocol):
|
||||||
@property
|
@property
|
||||||
def router(self) -> object: ...
|
def router(self) -> object: ...
|
||||||
@ -0,0 +1,12 @@
|
|||||||
|
from yuxi.channel.domain.port.internal.agent_port import AgentPort
|
||||||
|
from yuxi.channel.domain.port.internal.bot_loop_guard_port import BotLoopGuardPort
|
||||||
|
from yuxi.channel.domain.port.internal.cache_port import CachePort
|
||||||
|
from yuxi.channel.domain.port.internal.circuit_breaker_port import CircuitBreakerPort
|
||||||
|
from yuxi.channel.domain.port.internal.config_reload_port import ConfigReloadPort
|
||||||
|
from yuxi.channel.domain.port.internal.content_filter_port import ContentFilterPort, FilterResult
|
||||||
|
from yuxi.channel.domain.port.internal.event_publisher_port import DomainEvent, EventPublisherPort
|
||||||
|
from yuxi.channel.domain.port.internal.keyword_matcher_port import KeywordMatcherPort
|
||||||
|
from yuxi.channel.domain.port.internal.metrics_port import MetricsPort
|
||||||
|
from yuxi.channel.domain.port.internal.queue_port import QueuePort
|
||||||
|
from yuxi.channel.domain.port.internal.rate_limit_port import RateLimitPort
|
||||||
|
from yuxi.channel.domain.port.internal.sse_push_port import SsePushPort
|
||||||
@ -0,0 +1,61 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
|
from yuxi.channel.domain.port.bot_loop_guard_port import BotLoopGuardPort
|
||||||
|
from yuxi.channel.domain.port.cache_port import CachePort
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class RedisBotLoopGuard:
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
cache_port: CachePort,
|
||||||
|
*,
|
||||||
|
dm_budget: int = 10,
|
||||||
|
group_budget: int = 20,
|
||||||
|
window_seconds: int = 60,
|
||||||
|
cooldown_seconds: int = 300,
|
||||||
|
) -> None:
|
||||||
|
self._cache = cache_port
|
||||||
|
self._dm_budget = dm_budget
|
||||||
|
self._group_budget = group_budget
|
||||||
|
self._window = window_seconds
|
||||||
|
self._cooldown = cooldown_seconds
|
||||||
|
|
||||||
|
async def check(
|
||||||
|
self, session_id: str, *, sender_id: str = "", is_group: bool = False
|
||||||
|
) -> bool:
|
||||||
|
if is_group and not sender_id:
|
||||||
|
logger.warning(
|
||||||
|
"bot loop guard: group message without sender_id rejected, session=%s",
|
||||||
|
session_id,
|
||||||
|
)
|
||||||
|
return False
|
||||||
|
|
||||||
|
if is_group and sender_id:
|
||||||
|
key = f"channel:bot_loop:group:{session_id}:{sender_id}"
|
||||||
|
budget = self._group_budget
|
||||||
|
else:
|
||||||
|
key = f"channel:bot_loop:dm:{session_id}"
|
||||||
|
budget = self._dm_budget
|
||||||
|
|
||||||
|
count = await self._cache.incr(key)
|
||||||
|
if count == 1:
|
||||||
|
await self._cache.expire(key, self._window)
|
||||||
|
|
||||||
|
if count > budget:
|
||||||
|
await self._cache.expire(key, self._cooldown)
|
||||||
|
return False
|
||||||
|
|
||||||
|
return True
|
||||||
|
|
||||||
|
async def reset(
|
||||||
|
self, session_id: str, *, sender_id: str = "", is_group: bool = False
|
||||||
|
) -> None:
|
||||||
|
if is_group and sender_id:
|
||||||
|
key = f"channel:bot_loop:group:{session_id}:{sender_id}"
|
||||||
|
else:
|
||||||
|
key = f"channel:bot_loop:dm:{session_id}"
|
||||||
|
await self._cache.delete(key)
|
||||||
@ -0,0 +1,50 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
|
import redis.asyncio as aioredis
|
||||||
|
|
||||||
|
from yuxi.channel.domain.port.cache_port import CachePort
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class RedisCache(CachePort):
|
||||||
|
def __init__(self, redis: aioredis.Redis) -> None:
|
||||||
|
self._redis = redis
|
||||||
|
|
||||||
|
async def get(self, key: str) -> str | None:
|
||||||
|
raw = await self._redis.get(key)
|
||||||
|
if raw is None:
|
||||||
|
return None
|
||||||
|
return raw.decode() if isinstance(raw, bytes) else raw
|
||||||
|
|
||||||
|
async def set(
|
||||||
|
self, key: str, value: str, *, ex: int | None = None, nx: bool = False
|
||||||
|
) -> bool:
|
||||||
|
result = await self._redis.set(key, value, ex=ex, nx=nx)
|
||||||
|
if result is None:
|
||||||
|
return False
|
||||||
|
if isinstance(result, bool):
|
||||||
|
return result
|
||||||
|
if isinstance(result, bytes):
|
||||||
|
return result == b"OK"
|
||||||
|
return str(result) == "OK"
|
||||||
|
|
||||||
|
async def delete(self, key: str) -> None:
|
||||||
|
await self._redis.delete(key)
|
||||||
|
|
||||||
|
async def incr(self, key: str) -> int:
|
||||||
|
return await self._redis.incr(key)
|
||||||
|
|
||||||
|
async def expire(self, key: str, seconds: int) -> None:
|
||||||
|
await self._redis.expire(key, seconds)
|
||||||
|
|
||||||
|
async def ttl(self, key: str) -> int:
|
||||||
|
return await self._redis.ttl(key)
|
||||||
|
|
||||||
|
async def eval(self, script: str, keys: list[str], args: list[str | int]) -> tuple:
|
||||||
|
return await self._redis.eval(script, len(keys), *keys, *args)
|
||||||
|
|
||||||
|
async def publish(self, channel: str, message: str) -> None:
|
||||||
|
await self._redis.publish(channel, message)
|
||||||
@ -0,0 +1,107 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from yuxi.channel.domain.port.cache_port import CachePort
|
||||||
|
from yuxi.channel.domain.port.circuit_breaker_port import CircuitBreakerPort
|
||||||
|
|
||||||
|
|
||||||
|
_IS_AVAILABLE_SCRIPT = """
|
||||||
|
local state_key = KEYS[1]
|
||||||
|
local probe_key = KEYS[2]
|
||||||
|
local recovery_timeout = tonumber(ARGV[1])
|
||||||
|
local half_open_max = tonumber(ARGV[2])
|
||||||
|
|
||||||
|
local state = redis.call('GET', state_key)
|
||||||
|
if not state then
|
||||||
|
return {1, 'closed'}
|
||||||
|
end
|
||||||
|
|
||||||
|
if state == 'open' then
|
||||||
|
local ttl = redis.call('TTL', state_key)
|
||||||
|
if ttl > 0 then
|
||||||
|
return {0, 'open'}
|
||||||
|
end
|
||||||
|
redis.call('SET', state_key, 'half_open', 'EX', recovery_timeout)
|
||||||
|
redis.call('DEL', probe_key)
|
||||||
|
return {1, 'half_open'}
|
||||||
|
end
|
||||||
|
|
||||||
|
if state == 'half_open' then
|
||||||
|
local count = redis.call('INCR', probe_key)
|
||||||
|
if count == 1 then
|
||||||
|
redis.call('EXPIRE', probe_key, recovery_timeout)
|
||||||
|
end
|
||||||
|
if count <= half_open_max then
|
||||||
|
return {1, 'half_open'}
|
||||||
|
end
|
||||||
|
return {0, 'half_open'}
|
||||||
|
end
|
||||||
|
|
||||||
|
return {1, 'closed'}
|
||||||
|
"""
|
||||||
|
|
||||||
|
_RECORD_FAILURE_SCRIPT = """
|
||||||
|
local failure_key = KEYS[1]
|
||||||
|
local state_key = KEYS[2]
|
||||||
|
local probe_key = KEYS[3]
|
||||||
|
local recovery_timeout = tonumber(ARGV[1])
|
||||||
|
local failure_threshold = tonumber(ARGV[2])
|
||||||
|
|
||||||
|
local count = redis.call('INCR', failure_key)
|
||||||
|
if count == 1 then
|
||||||
|
redis.call('EXPIRE', failure_key, recovery_timeout * 2)
|
||||||
|
end
|
||||||
|
|
||||||
|
if count >= failure_threshold then
|
||||||
|
redis.call('SET', state_key, 'open', 'EX', recovery_timeout)
|
||||||
|
redis.call('DEL', probe_key)
|
||||||
|
return count
|
||||||
|
end
|
||||||
|
|
||||||
|
return count
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
class RedisCircuitBreaker:
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
cache_port: CachePort,
|
||||||
|
*,
|
||||||
|
failure_threshold: int = 5,
|
||||||
|
recovery_timeout: int = 30,
|
||||||
|
half_open_max: int = 1,
|
||||||
|
) -> None:
|
||||||
|
self._cache = cache_port
|
||||||
|
self._failure_threshold = failure_threshold
|
||||||
|
self._recovery_timeout = recovery_timeout
|
||||||
|
self._half_open_max = half_open_max
|
||||||
|
|
||||||
|
async def is_available(self, agent_config_id: int) -> bool:
|
||||||
|
state_key = f"channel:circuit:state:{agent_config_id}"
|
||||||
|
probe_key = f"channel:circuit:probe:{agent_config_id}"
|
||||||
|
|
||||||
|
result = await self._cache.eval(
|
||||||
|
_IS_AVAILABLE_SCRIPT,
|
||||||
|
keys=[state_key, probe_key],
|
||||||
|
args=[str(self._recovery_timeout), str(self._half_open_max)],
|
||||||
|
)
|
||||||
|
allowed = result[0]
|
||||||
|
return bool(allowed)
|
||||||
|
|
||||||
|
async def record_success(self, agent_config_id: int) -> None:
|
||||||
|
state_key = f"channel:circuit:state:{agent_config_id}"
|
||||||
|
failure_key = f"channel:circuit:failures:{agent_config_id}"
|
||||||
|
probe_key = f"channel:circuit:probe:{agent_config_id}"
|
||||||
|
await self._cache.delete(failure_key)
|
||||||
|
await self._cache.delete(probe_key)
|
||||||
|
await self._cache.delete(state_key)
|
||||||
|
|
||||||
|
async def record_failure(self, agent_config_id: int) -> None:
|
||||||
|
failure_key = f"channel:circuit:failures:{agent_config_id}"
|
||||||
|
state_key = f"channel:circuit:state:{agent_config_id}"
|
||||||
|
probe_key = f"channel:circuit:probe:{agent_config_id}"
|
||||||
|
|
||||||
|
await self._cache.eval(
|
||||||
|
_RECORD_FAILURE_SCRIPT,
|
||||||
|
keys=[failure_key, state_key, probe_key],
|
||||||
|
args=[str(self._recovery_timeout), str(self._failure_threshold)],
|
||||||
|
)
|
||||||
@ -0,0 +1,34 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
|
import redis.asyncio as aioredis
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class RedisRateLimiter:
|
||||||
|
def __init__(self, redis: aioredis.Redis) -> None:
|
||||||
|
self._redis = redis
|
||||||
|
|
||||||
|
async def check_and_incr(
|
||||||
|
self, key: str, *, max_attempts: int, window_seconds: int, lockout_seconds: int = 0
|
||||||
|
) -> bool:
|
||||||
|
count = await self._redis.incr(key)
|
||||||
|
if count == 1:
|
||||||
|
await self._redis.expire(key, window_seconds)
|
||||||
|
if count <= max_attempts:
|
||||||
|
return True
|
||||||
|
if lockout_seconds > 0:
|
||||||
|
lockout_key = key.replace(":attempts:", ":lockout:")
|
||||||
|
await self._redis.set(lockout_key, "1", ex=lockout_seconds)
|
||||||
|
return False
|
||||||
|
|
||||||
|
async def is_locked(self, key: str) -> tuple[bool, int]:
|
||||||
|
ttl = await self._redis.ttl(key)
|
||||||
|
if ttl is None or ttl < 0:
|
||||||
|
return False, 0
|
||||||
|
return True, ttl
|
||||||
|
|
||||||
|
async def reset(self, key: str) -> None:
|
||||||
|
await self._redis.delete(key)
|
||||||
@ -5,7 +5,7 @@ import logging
|
|||||||
from fastapi import FastAPI
|
from fastapi import FastAPI
|
||||||
|
|
||||||
from yuxi.channel.domain.port.channel_adapter_port import ChannelAdapterPort
|
from yuxi.channel.domain.port.channel_adapter_port import ChannelAdapterPort
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
@ -30,11 +30,11 @@ def register_all_channel_routes(
|
|||||||
return registered
|
return registered
|
||||||
|
|
||||||
|
|
||||||
def _get_route_contributor(adapter: ChannelAdapterPort) -> ChannelRouteContributor | None:
|
def _get_route_contributor(adapter: ChannelAdapterPort) -> ChannelRouteContributorPort | None:
|
||||||
contributor = getattr(adapter, "route_contributor", None)
|
contributor = getattr(adapter, "route_contributor", None)
|
||||||
if contributor is None:
|
if contributor is None:
|
||||||
return None
|
return None
|
||||||
if isinstance(contributor, ChannelRouteContributor):
|
if isinstance(contributor, ChannelRouteContributorPort):
|
||||||
return contributor
|
return contributor
|
||||||
if hasattr(contributor, "router"):
|
if hasattr(contributor, "router"):
|
||||||
return contributor
|
return contributor
|
||||||
|
|||||||
@ -5,7 +5,7 @@ from unittest.mock import AsyncMock, MagicMock, patch
|
|||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildAdapters:
|
class TestBuildAdapters:
|
||||||
|
|||||||
@ -6,7 +6,7 @@ import pytest
|
|||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
||||||
from yuxi.channel.application.service.auth_service import AuthService
|
from yuxi.channel.application.service.auth_service import AuthService
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildAuthService:
|
class TestBuildAuthService:
|
||||||
|
|||||||
@ -5,7 +5,7 @@ from unittest.mock import patch
|
|||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory
|
from yuxi.channel.container import ChannelContainerFactory
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildConfig:
|
class TestBuildConfig:
|
||||||
|
|||||||
@ -5,7 +5,7 @@ from unittest.mock import MagicMock
|
|||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildInfra:
|
class TestBuildInfra:
|
||||||
|
|||||||
@ -7,7 +7,7 @@ import pytest
|
|||||||
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle
|
||||||
from yuxi.channel.application.service.auth_service import AuthService
|
from yuxi.channel.application.service.auth_service import AuthService
|
||||||
from yuxi.channel.domain.service.pipeline import Pipeline
|
from yuxi.channel.domain.service.pipeline import Pipeline
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildPipeline:
|
class TestBuildPipeline:
|
||||||
|
|||||||
@ -6,7 +6,7 @@ import pytest
|
|||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle, _WorkerBundle
|
from yuxi.channel.container import ChannelContainerFactory, _InfraBundle, _WorkerBundle
|
||||||
from yuxi.channel.domain.service.pipeline import Pipeline
|
from yuxi.channel.domain.service.pipeline import Pipeline
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestBuildWorkers:
|
class TestBuildWorkers:
|
||||||
|
|||||||
@ -1,10 +1,10 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
|
|
||||||
|
|
||||||
def test_channel_route_contributor_is_protocol() -> None:
|
def test_channel_route_contributor_port_is_protocol() -> None:
|
||||||
assert hasattr(ChannelRouteContributor, "router")
|
assert hasattr(ChannelRouteContributorPort, "router")
|
||||||
|
|
||||||
|
|
||||||
class _FakeContributor:
|
class _FakeContributor:
|
||||||
@ -15,4 +15,4 @@ class _FakeContributor:
|
|||||||
|
|
||||||
def test_fake_contributor_is_instance() -> None:
|
def test_fake_contributor_is_instance() -> None:
|
||||||
obj = _FakeContributor()
|
obj = _FakeContributor()
|
||||||
assert isinstance(obj, ChannelRouteContributor)
|
assert isinstance(obj, ChannelRouteContributorPort)
|
||||||
|
|||||||
@ -0,0 +1 @@
|
|||||||
|
|
||||||
@ -5,7 +5,7 @@ from pathlib import Path
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
@ -5,7 +5,7 @@ from unittest.mock import AsyncMock, MagicMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.config.redis_config_reload import RedisConfigReload
|
from yuxi.channel.infrastructure.configuration.redis_config_reload import RedisConfigReload
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
@ -4,7 +4,7 @@ import pytest
|
|||||||
from fastapi import FastAPI, APIRouter
|
from fastapi import FastAPI, APIRouter
|
||||||
|
|
||||||
from yuxi.channel.interfaces.rest.router.registry import register_all_channel_routes, _get_route_contributor
|
from yuxi.channel.interfaces.rest.router.registry import register_all_channel_routes, _get_route_contributor
|
||||||
from yuxi.channel.domain.port.channel_route_contributor import ChannelRouteContributor
|
from yuxi.channel.domain.port.channel_route_contributor_port import ChannelRouteContributorPort
|
||||||
|
|
||||||
|
|
||||||
class _FakeContributor:
|
class _FakeContributor:
|
||||||
|
|||||||
@ -2,7 +2,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestChannelConfig:
|
class TestChannelConfig:
|
||||||
|
|||||||
@ -7,7 +7,7 @@ import pytest
|
|||||||
import redis.asyncio as aioredis
|
import redis.asyncio as aioredis
|
||||||
|
|
||||||
from yuxi.channel.container import ChannelContainerFactory
|
from yuxi.channel.container import ChannelContainerFactory
|
||||||
from yuxi.channel.infrastructure.config.channel_config import ChannelConfig
|
from yuxi.channel.infrastructure.configuration.channel_config import ChannelConfig
|
||||||
|
|
||||||
|
|
||||||
class TestChannelContainerFactory:
|
class TestChannelContainerFactory:
|
||||||
|
|||||||
@ -4,7 +4,7 @@ from unittest.mock import AsyncMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.cache.redis_bot_loop_guard import RedisBotLoopGuard
|
from yuxi.channel.infrastructure.cache_infra.redis_bot_loop_guard import RedisBotLoopGuard
|
||||||
|
|
||||||
|
|
||||||
class TestRedisBotLoopGuard:
|
class TestRedisBotLoopGuard:
|
||||||
|
|||||||
@ -4,7 +4,7 @@ from unittest.mock import AsyncMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.cache.redis_cache import RedisCache
|
from yuxi.channel.infrastructure.cache_infra.redis_cache import RedisCache
|
||||||
|
|
||||||
|
|
||||||
class TestRedisCache:
|
class TestRedisCache:
|
||||||
|
|||||||
@ -4,7 +4,7 @@ from unittest.mock import AsyncMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.cache.redis_circuit_breaker import RedisCircuitBreaker
|
from yuxi.channel.infrastructure.cache_infra.redis_circuit_breaker import RedisCircuitBreaker
|
||||||
|
|
||||||
|
|
||||||
class TestRedisCircuitBreaker:
|
class TestRedisCircuitBreaker:
|
||||||
|
|||||||
@ -4,7 +4,7 @@ from unittest.mock import AsyncMock, MagicMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.config.redis_config_reload import RedisConfigReload
|
from yuxi.channel.infrastructure.configuration.redis_config_reload import RedisConfigReload
|
||||||
|
|
||||||
|
|
||||||
class TestRedisConfigReload:
|
class TestRedisConfigReload:
|
||||||
|
|||||||
@ -4,7 +4,7 @@ from unittest.mock import AsyncMock
|
|||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from yuxi.channel.infrastructure.cache.redis_rate_limiter import RedisRateLimiter
|
from yuxi.channel.infrastructure.cache_infra.redis_rate_limiter import RedisRateLimiter
|
||||||
|
|
||||||
|
|
||||||
class TestRedisRateLimiter:
|
class TestRedisRateLimiter:
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user