Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 81 additions & 6 deletions README.en.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# EverywhereYouGo (EGo) v1.2.4
# EverywhereYouGo (EGo) v1.3.0

[中文](README.md) | English

Expand Down Expand Up @@ -53,17 +53,34 @@ Optionally set `EGO_SECRET_KEY` to customize Flask session key.

## Configuration Files

Configuration is persisted as JSON files in `config/` directory:
Configuration lives in two places, with distinct roles:

| Storage | Role |
|---------|------|
| SQLite (`ego.db`) | **Runtime source of truth** — all reads/writes go through it |
| `config/*.json` | **Export / backup medium** — for backup, versioning and migration |

| File | Content |
|------|------|
|------|---------|
| `config/parsers.json` | Parser metadata |
| `config/sources.json` | Data source definitions |
| `config/channels.json` | Push channel configurations |
| `config/templates.json` | Push templates |
| `config/bindings.json` | Channel bindings (with condition expressions) |

Can directly edit JSON and restart to take effect, or manage via WebUI. System settings (DND, log level, etc.) and runtime data (message logs) are stored in SQLite (`ego.db`).
**Load rules at startup:**

1. Database **not empty** → the database wins; JSON is not read, and the current
config is **written back** to `config/*.json` as a snapshot
2. Database **empty** and JSON present → import from JSON (first run / migration / restore)
3. Database empty and no JSON → export the initial config to JSON

So **use the WebUI for day-to-day config changes** (they take effect immediately).
Hand-editing `config/*.json` is only read on first import when the database is empty —
it is not the normal path for applying changes.

System settings (DND, log level, etc.), the message log and the queue are also stored in SQLite.
For backup/restore use **Settings → Backup**, which packages `config/*.json` + `parsers/*.py`.

## Parsers

Expand Down Expand Up @@ -101,7 +118,10 @@ Supports `and`, `or`, parentheses grouping:
Set DND time period, messages enter queue and wait, automatically flush when period ends. Urgent routes are not affected by DND.

### Message Deduplication
Channel bindings can configure `dedup_key_expr` and `dedup_window` (default 3600 seconds). Same dedup key will not be sent repeatedly within the window.
Deduplication granularity is **message × channel**: each channel binding has its own
`dedup_key_expr` and `dedup_window` (default 3600 seconds), independent of the others.
A channel that hits is skipped while the rest still send; the whole message is marked
`DISCARDED` only when **every** channel hits.

### Parallel Push
When multiple channels match, thread pool sends in parallel, total latency depends on the slowest single channel.
Expand All @@ -110,13 +130,56 @@ When multiple channels match, thread pool sends in parallel, total latency depen
Each data source automatically saves the last 20 request samples, can select samples in WebUI for test parsing and pushing.

### Message Resend
Failed messages support original resend (using parsed msg_json) or re-parse and resend.
Failed messages support original resend (using the parsed msg_json) or re-parse and resend.
By default **only the channels that failed are retried** — already-succeeded channels are
not pushed a second time. Use `scope=all` to force a full re-push.

### Import & Export
- **Backup**: Download ZIP package (`config/*.json` + `parsers/*.py`)
- **Restore**: Upload ZIP package, automatically takes effect after overwriting configuration
- **JSON Import**: Supports dry_run preview, insert/overwrite two modes, dependency check

### Channel Circuit Breaker
Automatically isolates a channel that keeps failing, so one broken third party cannot
drag down the whole send path:

- Failure ratio > 50% within a sliding window (default 60s), **or** 5 consecutive
failures → trip
- Cooldown backs off exponentially: 30 → 60 → 120 → … → 600 s (capped at 10 min)
- After cooldown it enters half-open probing; 3 consecutive successes restore it
- **4xx does not count as failure** (a business rejection is not a service outage) —
only 5xx / timeouts / connection errors are counted
- While tripped, messages **stay queued**: no retry budget is consumed and nothing is
dropped. State is persisted and survives restarts.

### Outbound Rate Limiting
Per-channel rate limit (messages per minute) to avoid getting blocked by the remote side.
Retries cannot fix a 429 — limiting has to happen **before** sending. A message that
cannot get a token waits in the queue instead of being dropped.

### Resilience UI
| Where | What you can do |
|-------|-----------------|
| Channel list → **Resilience column** | See the rate-limit badge and breaker countdown; reset a tripped channel with one click |
| Channel edit dialog | Set this channel's outbound rate limit (empty/0 = unlimited) |
| Settings → **Resilience (channel circuit breaker)** | Hand-tune sliding window, consecutive-failure threshold, cooldown base/cap, probe count, etc. |

Parameter precedence: `system_config` (settings page / direct DB edit) > environment
variable > built-in default. Changes take effect immediately, no restart needed.

### Observability
| Endpoint | Description |
|----------|-------------|
| `GET /api/metrics` | Queue depth, DLQ count, per-channel success rate, end-to-end latency, breaker & rate-limit state; `?hours=N` sets the stats window (default 24h) |
| `GET /api/resilience` | Channels currently tripped / rate-limited |
| `GET /api/queue/stats` | Queue and dead-letter counts |
| `GET /api/health` | Health check (SQLite / disk / config / queue) |

### Graceful Stop
On `SIGTERM` EGo stops accepting new messages first, then waits for in-flight tasks
(up to 30 seconds); anything unfinished is moved to the dead-letter queue — a container
restart does not lose messages.

## Internationalization

Built-in Chinese and English bilingual support, switch languages anytime via language switch button in top-right corner of navigation bar.
Expand Down Expand Up @@ -146,6 +209,18 @@ Built-in Chinese and English bilingual support, switch languages anytime via lan
| `LOG_LEVEL` | `INFO` | Log level |
| `EGO_AUTH_TOKEN` | *(empty)* | Access control Token |
| `EGO_SECRET_KEY` | *(auto)* | Flask session key |
| `EGO_INGRESS_WORKERS` | `8` | Ingress worker threads per port source |
| `EGO_INGRESS_MAX_QUEUE` | `200` | Ingress queue cap; beyond it returns 503 (backpressure) |
| `EGO_CLEANUP_INTERVAL` | `600` | Interval for purging old messages / dedup keys (s) |
| `EGO_BREAKER_WINDOW` | `60` | Circuit breaker sliding window (s) |
| `EGO_BREAKER_MIN_SAMPLES` | `5` | Min samples before the failure-ratio rule applies |
| `EGO_BREAKER_FAILURE_RATIO` | `0.5` | Failure ratio that trips the breaker |
| `EGO_BREAKER_CONSECUTIVE` | `5` | Consecutive-failure threshold (low-traffic channels) |
| `EGO_BREAKER_OPEN_BASE` | `30` | Base cooldown (s), doubles on each open |
| `EGO_BREAKER_OPEN_MAX` | `600` | Cooldown cap (s) |
| `EGO_BREAKER_HALF_OPEN_OK` | `3` | Consecutive probe successes needed to recover |
| `EGO_RATE_MAX_WAIT` | `1.0` | Max wait for a rate-limit token (s), then defer |
| `EGO_RATE_MISS_TTL` | `30` | Re-check interval for channels without a rate limit (s) |

## License

Expand Down
76 changes: 71 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# EverywhereYouGo (EGo) v1.2.4
# EverywhereYouGo (EGo) v1.3.0

[English](README.en.md) | 中文

Expand Down Expand Up @@ -51,9 +51,14 @@ EGO_AUTH_TOKEN=your-secret-token python3 main.py

可选设置 `EGO_SECRET_KEY` 自定义 Flask session 密钥。

## 配置文件
## 配置存储

配置持久化为 JSON 文件,位于 `config/` 目录:
配置有两份,角色不同:

| 存储 | 角色 |
|------|------|
| SQLite(`ego.db`) | **运行时真相源** —— 所有读写以库内数据为准 |
| `config/*.json` | **导出 / 备份介质** —— 便于备份、版本管理与迁移 |

| 文件 | 内容 |
|------|------|
Expand All @@ -63,7 +68,17 @@ EGO_AUTH_TOKEN=your-secret-token python3 main.py
| `config/templates.json` | 推送模板 |
| `config/bindings.json` | 渠道绑定(含条件表达式) |

可直接编辑 JSON 后重启生效,也可通过 WebUI 管理。系统设置(DND、日志级别等)和运行时数据(消息日志)存储在 SQLite(`ego.db`)中。
**启动时的加载规则:**

1. 数据库**已有**配置 → 以数据库为准,不读 JSON,并把当前配置**刷写**回 `config/*.json`
2. 数据库**为空**且有 JSON → 从 JSON 导入(首次启动 / 迁移 / 恢复)
3. 数据库为空且无 JSON → 导出初始配置到 JSON

因此 **日常改配置请用 WebUI**(改完即时生效)。直接编辑 `config/*.json` 只在
「数据库为空」的首次导入场景才会被读取,不是常规生效路径。

系统设置(DND、日志级别等)、消息日志与队列同样存储在 SQLite。
配置备份 / 恢复请用「系统设置 → 备份」,会打包 `config/*.json` + `parsers/*.py`。

## 解析器

Expand Down Expand Up @@ -101,7 +116,9 @@ def parse(raw_body: bytes, headers: dict, query_params: dict) -> dict:
设置免打扰时段后,消息进入队列等待,结束后自动刷新。紧急路由不受 DND 影响。

### 消息去重
渠道绑定可配置 `dedup_key_expr` 和 `dedup_window`(默认 3600 秒)。同一去重键在窗口内不重复发送。
去重粒度是 **消息 × 通道**:每个渠道绑定各自配置 `dedup_key_expr` 与 `dedup_window`
(默认 3600 秒),互不影响。命中的渠道被跳过,其余渠道照常发送;
只有**全部**渠道命中时,整条消息才标记为 `DISCARDED`。

### 并行推送
多渠道匹配时线程池并行发送,总延迟取决于最慢的单个渠道。
Expand All @@ -111,12 +128,49 @@ def parse(raw_body: bytes, headers: dict, query_params: dict) -> dict:

### 消息重发
失败消息支持原始重发(使用已解析的 msg_json)或重新解析后重发。
默认**只重发上次失败的渠道**,不会把已经成功的渠道重复推送一遍;
需要整体重推时可用 `scope=all`。

### 导入导出
- **备份**:下载 ZIP 包(`config/*.json` + `parsers/*.py`)
- **恢复**:上传 ZIP 包,覆盖配置后自动生效
- **JSON 导入**:支持 dry_run 预览、insert/overwrite 两种模式、依赖检查

### 通道熔断
第三方渠道持续故障时自动隔离,避免拖垮整条发送链路:

- 滑动窗口(默认 60s)内失败率 > 50%,**或**连续失败 ≥ 5 次 → 熔断
- 冷却时间指数退避 30 → 60 → 120 → … → 600 秒(封顶 10 分钟)
- 冷却结束后进入半开探测,连续 3 次成功才恢复
- **4xx 不计失败**(业务侧拒绝不等于服务故障),只对 5xx / 超时 / 连接类错误计数
- 熔断期间消息**留在队列等待**:不消耗重试次数、不丢弃;状态持久化,重启后仍生效

### 出站限流
按通道独立限流(条/分钟),防止发送过快被对方封禁。重试解决不了 429——
限流必须发生在发送**之前**。拿不到令牌的消息会排队等待,而不是被丢弃。

### 韧性界面
| 位置 | 能做什么 |
|------|---------|
| 通道列表 → **「韧性」列** | 查看限流徽章与熔断倒计时;熔断时可一键「手工恢复」 |
| 通道编辑弹窗 | 设置该通道的出站限流(留空/0 = 不限流) |
| 系统设置 → **韧性(通道熔断)** | 手工调整滑动窗口、连续失败阈值、冷却基数/上限、探测次数等参数 |

参数取值优先级:`system_config`(设置页 / 直接改库) > 环境变量 > 内置默认,
改完即时生效、无需重启。

### 可观测性
| 端点 | 说明 |
|------|------|
| `GET /api/metrics` | 队列深度、死信总数、各通道成功率、端到端延迟、熔断与限流状态;`?hours=N` 调整统计窗口(默认 24h) |
| `GET /api/resilience` | 当前处于熔断 / 限流状态的通道 |
| `GET /api/queue/stats` | 队列与死信计数 |
| `GET /api/health` | 健康检查(SQLite / 磁盘 / 配置 / 队列) |

### 优雅停机
收到 `SIGTERM` 时先停止接收新消息,再等待在途任务完成(最长 30 秒),
超时未完成的转入死信队列——容器重启不会丢消息。

## 国际化

内置中英文双语支持,通过导航栏右上角语言切换按钮随时切换。
Expand Down Expand Up @@ -146,6 +200,18 @@ def parse(raw_body: bytes, headers: dict, query_params: dict) -> dict:
| `LOG_LEVEL` | `INFO` | 日志等级 |
| `EGO_AUTH_TOKEN` | *(空)* | 访问控制 Token |
| `EGO_SECRET_KEY` | *(自动)* | Flask session 密钥 |
| `EGO_INGRESS_WORKERS` | `8` | 每个端口数据源的入口工作线程数 |
| `EGO_INGRESS_MAX_QUEUE` | `200` | 入口等待队列上限,超出返回 503(背压) |
| `EGO_CLEANUP_INTERVAL` | `600` | 旧消息 / 去重键的清理间隔(秒) |
| `EGO_BREAKER_WINDOW` | `60` | 熔断滑动窗口(秒) |
| `EGO_BREAKER_MIN_SAMPLES` | `5` | 窗口内触发失败率判定的最少样本数 |
| `EGO_BREAKER_FAILURE_RATIO` | `0.5` | 窗口失败率阈值(超过则熔断) |
| `EGO_BREAKER_CONSECUTIVE` | `5` | 连续失败阈值(照顾低频通道) |
| `EGO_BREAKER_OPEN_BASE` | `30` | 熔断冷却基数(秒),逐次翻倍 |
| `EGO_BREAKER_OPEN_MAX` | `600` | 熔断冷却上限(秒) |
| `EGO_BREAKER_HALF_OPEN_OK` | `3` | 恢复所需连续探测成功次数 |
| `EGO_RATE_MAX_WAIT` | `1.0` | 限流取令牌的最长等待(秒),超时改为延迟重排 |
| `EGO_RATE_MISS_TTL` | `30` | 未配置限流的通道,回查数据库的间隔(秒) |

## License

Expand Down
8 changes: 8 additions & 0 deletions api/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -136,4 +136,12 @@ def _auth_middleware():
app.source_mgr = source_mgr
app.auth_token = AUTH_TOKEN

# ── 统一的入参校验错误处理(improvement #27)──
# 业务代码里直接 raise ValidationError 即可,无需每处 try/except
from api.validation import ValidationError

@app.errorhandler(ValidationError)
def _handle_validation_error(e):
return jsonify({"status": "error", "error": str(e)}), 400

return app
87 changes: 83 additions & 4 deletions api/channels.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,13 @@
"""api/channels.py — 通道 CRUD + 插件管理"""

import os
import json
import log
import db
import channel_loader
import i18n
from flask import Blueprint, request, jsonify
from api.validation import require_name, optional_flag, ValidationError

channels_bp = Blueprint("channels", __name__)

Expand All @@ -18,22 +21,98 @@ def api_channels():
return jsonify(db.get_channels())


def _channel_config(data):
"""通道 config 必须是 JSON 对象(允许传对象或 JSON 字符串)。"""
cfg = (data or {}).get("config", {})
if isinstance(cfg, str):
try:
cfg = json.loads(cfg or "{}")
except Exception:
raise ValidationError("config must be a JSON object")
if not isinstance(cfg, dict):
raise ValidationError("config must be a JSON object")
return cfg


@channels_bp.route("/api/channels", methods=["POST"])
def api_create_channel():
data = request.json
cid = db.create_channel(data["name"], data["type"], data.get("config", {}))
data = request.json or {}
cid = db.create_channel(
require_name(data),
require_name(data, key="type", max_len=64),
_channel_config(data),
)
import config_manager
config_manager.sync_table("channels")
return jsonify({"id": cid})


@channels_bp.route("/api/channels/<int:cid>", methods=["PUT"])
def api_update_channel(cid):
data = request.json
db.update_channel(cid, **data)
data = request.json or {}
patch = {}
if "name" in data:
patch["name"] = require_name(data)
if "type" in data:
patch["type"] = require_name(data, key="type", max_len=64)
if "config" in data:
patch["config"] = _channel_config(data)
if "enabled" in data:
patch["enabled"] = optional_flag(data, "enabled")
if patch:
db.update_channel(cid, **patch)
return jsonify({"status": "ok"})


@channels_bp.route("/api/channels/<int:cid>", methods=["DELETE"])
def api_delete_channel(cid):
"""删除通道(级联清理绑定、限流、熔断、去重键;待发任务移入死信)。"""
if not db.get_channel(cid):
return jsonify({"error": i18n._("err.not_found")}), 404
db.delete_channel(cid)
import config_manager
config_manager.sync_table("channels")
return jsonify({"status": "ok"})


@channels_bp.route("/api/channels/<int:cid>/test", methods=["POST"])
def api_test_channel(cid):
"""测试通道配置是否可用。前端期望 {ok: bool, error?}。"""
ch = db.get_channel(cid)
if not ch:
return jsonify({"ok": False, "error": i18n._("err.not_found")}), 404
try:
cfg = _channel_config({"config": ch["config"]})
ok = channel_loader.create_channel(ch["type"], cfg).test()
return jsonify({"ok": bool(ok), "error": None if ok else i18n._("ch.test_fail")})
except Exception as e:
return jsonify({"ok": False, "error": str(e)[:300]})


@channels_bp.route("/api/channels/<int:cid>/duplicate", methods=["POST"])
def api_duplicate_channel(cid):
"""复制通道(默认禁用,避免复制出来就开始推送)。"""
src = db.get_channel(cid)
if not src:
return jsonify({"error": i18n._("err.not_found")}), 404

new_id = db.create_channel(src["name"] + " (copy)", src["type"],
src["config"], enabled=0)

# 出站限流存在独立表里,不显式复制就会静默丢失
try:
from rate_limiter import get_limiter
rate = get_limiter().get_rate(cid)
if rate > 0:
get_limiter().set_rate(new_id, rate)
except Exception as e:
log.logger.warning(f"Duplicate channel {cid}: rate limit not copied: {e}")

import config_manager
config_manager.sync_table("channels")
return jsonify({"id": new_id})


# ── Channel Plugins ──


Expand Down
Loading
Loading