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
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,13 +88,13 @@ cp fail2ban/filter.d/* /etc/fail2ban/filter.d/
| 字段 | 类型 | 默认 | 含义 | 公共 mirror 推荐起点 |
| --- | --- | --- | --- | --- |
| `relay_idle_timeout` | int 秒 | 0 | relay 阶段双向无 I/O 多久后关闭连接。语义同 rsyncd `timeout`。 | `600` |
| `relay_max_duration` | int 秒 | 0 | relay 阶段总时长上限。超时关闭,rsync 客户端通常会重连续传。 | `14400`(4h) |
| `relay_max_duration` | int 秒 | 0 | relay 阶段总时长上限。超时关闭;rsync 本身不会自动重连,是否续传取决于调用方(脚本/cron)。 | `14400`(4h) |
| `tcp_keepalive` | int 秒 | 0 | 客户端连接和上游连接的 TCP keepalive 周期,0 沿用 OS 默认(通常 ~2h)。 | `120` |
| `per_ip_max_active_connections` | int | 0 | 单 IP 对单上游最多并发 relay 连接数,proxy-wide 默认值。NAT/校园出口 IP 时取值需放宽。 | `4` |
| `dial_timeout` | int 秒 | 0 | 拨号上游的超时;0 沿用内核 SYN 重试(~75s)。 | `5`(LAN) |
| `min_throughput_bytes` | int64 字节 | 0 | relay 阶段最近 `min_throughput_window` 秒内累计收发须 ≥ 此值,否则视作慢速吸血并关闭。0 关闭整组检查。 | `1048576` |
| `min_throughput_window` | int 秒 | 60 | 上述滑动窗口长度。 | `60` |
| `min_throughput_grace` | int 秒 | = window | 连接刚开始的豁免期,避免误杀大 module 的 file-list 阶段。 | `600` |
| `min_throughput_window` | int 秒 | 60 | 上述滑动窗口长度。有效下限 = `min_throughput_bytes` / 本值;窗口越短反应越快但等效速率越高。 | `300`(≈3.4 KiB/s) |
| `min_throughput_grace` | int 秒 | = window | 连接刚开始的豁免期;豁免期内完全不判定,豁免结束后才开始第一个采样窗口,避免误杀大 module 的 file-list 阶段。 | `600` |

各项触发的事件都有对应的 Prometheus counter(`rsync_proxy_relay_idle_timeout_terminated_total`、`rsync_proxy_relay_max_duration_terminated_total`、`rsync_proxy_throughput_floor_terminated_total`、`rsync_proxy_per_ip_rejected_total`、`rsync_proxy_upstream_dial_errors_total`),便于先观察再调参。

Expand Down
42 changes: 29 additions & 13 deletions assets/config.example.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,11 @@ motd = "Served by rsync-proxy (https://github.com/ustclug/rsync-proxy)"

# Hard upper bound (seconds) on the total wall-clock duration of the
# bidirectional relay phase. When exceeded the proxy closes the
# connection regardless of activity. rsync clients typically resume
# on reconnect. 0 (the default) disables this cap. For public mirrors,
# 14400 (4h) bounds the worst case while still accommodating large
# initial syncs.
# connection regardless of activity. rsync does not automatically
# reconnect; retrying or resuming depends on how the caller invokes
# it (e.g. a script or cron job). 0 (the default) disables this cap.
# For public mirrors, 14400 (4h) bounds the worst case while still
# accommodating large initial syncs.
#relay_max_duration = 0

# TCP keepalive period (seconds) applied to accepted client
Expand Down Expand Up @@ -59,16 +60,31 @@ motd = "Served by rsync-proxy (https://github.com/ustclug/rsync-proxy)"
# Throughput floor enforced during the relay phase. A connection
# that transfers fewer than min_throughput_bytes (counting both
# directions) over the most recent min_throughput_window seconds is
# treated as slow-leeching and terminated. min_throughput_grace
# gives a fresh connection this many seconds of ramp-up before the
# floor is enforced; this matters because rsync's initial file-list
# phase can take many minutes for very large modules (millions of
# files), during which the apparent throughput is low. Defaults to
# min_throughput_window when unset. Set min_throughput_bytes to 0
# (the default) to disable the check. Typical starting values:
# bytes = 1048576 (1 MiB), window = 60, grace = 600.
# treated as slow-leeching and terminated. Set min_throughput_bytes
# to 0 (the default) to disable the check.
#
# Choosing values: the effective floor is
# min_throughput_bytes / min_throughput_window. A short window reacts
# quickly but also raises the effective rate and is easily tripped by
# a client that merely stalls for one window; a longer window smooths
# transient stalls and lowers the rate. For a public mirror,
# 1048576 (1 MiB) over 300 (5 minutes) works out to ~3.4 KiB/s
# averaged over 5 minutes, which still terminates genuinely stalled
# clients (most offenders transfer only a few hundred bytes per
# minute) while sparing slow-but-legitimate ones. Avoid combining
# 1 MiB with the 60s default, which is ~17.5 KiB/s and has been
# observed to kill real transfers.
#
# min_throughput_grace exempts a fresh connection from the floor for
# this many seconds; no measurement is taken during the grace period,
# and the first sampling window starts only after it ends. This
# matters because rsync's initial file-list phase can take many
# minutes for very large modules (millions of files), during which
# the apparent throughput is low. Defaults to min_throughput_window
# when unset. Typical starting values:
# bytes = 1048576 (1 MiB), window = 300, grace = 600.
#min_throughput_bytes = 0
#min_throughput_window = 60
#min_throughput_window = 300
#min_throughput_grace = 600

[upstreams.u1]
Expand Down
5 changes: 3 additions & 2 deletions pkg/server/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,9 @@ type ProxySettings struct {
// RelayMaxDurationSecs is the maximum total wall-clock duration
// (in seconds) of the bidirectional relay phase. When exceeded
// the proxy closes the connection regardless of activity. 0 (the
// default) disables this hard cap. rsync clients will typically
// reconnect and resume on the next run.
// default) disables this hard cap. rsync does not automatically
// reconnect; retrying or resuming depends on how the caller invokes
// it (e.g. a script or cron job).
RelayMaxDurationSecs int `toml:"relay_max_duration"`
// TCPKeepAliveSecs enables TCP keepalive on accepted client
// connections and on dialed upstream connections. The value is
Expand Down
52 changes: 35 additions & 17 deletions pkg/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -1112,13 +1112,19 @@ func (s *Server) relay(ctx context.Context, index uint32, downConn net.Conn) (re
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
// Sliding-window throughput state. lastSampleTime moves
// forward each time we reach a full window, with
// lastSampleBytes capturing the cumulative byte count at
// that moment. delta over the next window must meet
// minBytes or the connection is terminated.
// Sliding-window throughput state. The first sampling
// window is anchored at the moment the grace period ends,
// so the ramp-up phase (e.g. rsync's initial file-list
// exchange) is never counted towards the first measurement.
// Afterwards lastSampleTime moves forward each time a full
// window is reached, with lastSampleBytes capturing the
// cumulative byte count at that moment. delta over the next
// window must meet minBytes or the connection is terminated.
// When minGrace is 0 the window starts at relayStartedAt,
// preserving the original behaviour.
lastSampleTime := relayStartedAt
lastSampleBytes := int64(0)
graceEnded := minGrace <= 0
for {
select {
case <-sentClosed:
Expand All @@ -1145,20 +1151,32 @@ func (s *Server) relay(ctx context.Context, index uint32, downConn net.Conn) (re
_ = downConn.Close()
return
}
if throughputEnabled && now.Sub(relayStartedAt) >= minGrace && now.Sub(lastSampleTime) >= minWindow {
if throughputEnabled {
curBytes := info.SentBytes.Load() + info.ReceivedBytes.Load()
delta := curBytes - lastSampleBytes
if delta < minBytes {
idleTimedOut.Store(true)
s.getUpstreamCounters(upstreamName).throughputFloorTerminated.Add(1)
s.accessLog.F("client %s for module %s below throughput floor (%d bytes < %d bytes in %s), closing",
ip, moduleName, delta, minBytes, minWindow)
_ = upConn.Close()
_ = downConn.Close()
return
if !graceEnded {
if now.Sub(relayStartedAt) >= minGrace {
// Grace just ended: anchor a fresh
// window here so the ramp-up is
// excluded from the first
// measurement.
graceEnded = true
lastSampleTime = now
lastSampleBytes = curBytes
}
} else if now.Sub(lastSampleTime) >= minWindow {
delta := curBytes - lastSampleBytes
if delta < minBytes {
idleTimedOut.Store(true)
s.getUpstreamCounters(upstreamName).throughputFloorTerminated.Add(1)
s.accessLog.F("client %s for module %s below throughput floor (%d bytes < %d bytes in %s), closing",
ip, moduleName, delta, minBytes, minWindow)
_ = upConn.Close()
_ = downConn.Close()
return
}
lastSampleTime = now
lastSampleBytes = curBytes
}
lastSampleTime = now
lastSampleBytes = curBytes
}
}
}
Expand Down
74 changes: 74 additions & 0 deletions pkg/server/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1807,6 +1807,80 @@ func TestThroughputFloorTerminatesSlowConnection(t *testing.T) {
assert.Contains(t, string(logData), "for module fake")
}

// TestThroughputFloorExcludesGracePeriodFromFirstWindow verifies that
// bytes transferred during the grace period are excluded from the first
// throughput measurement window. A relay that is quiet during the
// ramp-up (as during rsync's initial file-list exchange) but then
// transfers well above the floor must not be torn down by the first
// evaluation, which previously averaged over the whole grace period.
func TestThroughputFloorExcludesGracePeriodFromFirstWindow(t *testing.T) {
srv := startServer(t)
defer srv.Close()
srv.RelayIdleTimeout = 0
srv.RelayMaxDuration = 0
// The grace period (2s) spans several ticks of the 1s-minimum
// ticker and the window (1s) is one tick, so the bytes sent during
// the ramp-up would dominate a window anchored at the start of the
// relay.
srv.MinThroughputBytes = 1000
srv.MinThroughputWindow = time.Second
srv.MinThroughputGrace = 2 * time.Second
accessLogPath := setupAccessLog(t, srv)

r := require.New(t)

fakeRsync := rsync.NewServer(func(conn *rsync.Conn) {
defer conn.Close()

if _, _, err := doServerHandshake(conn, RsyncdServerVersion); err != nil {
return
}
_, _ = io.ReadAll(conn)
})
fakeRsync.Start()
defer fakeRsync.Close()

srv.modules = map[string][]Target{
"fake": {{Upstream: "u1", Addr: fakeRsync.Listener.Addr().String()}},
}
srv.upstreams = []upstreamConfig{{Name: "u1"}}
srv.upstreamQueues = map[string]*queue.Queue{"u1": queue.New(0, 0)}

rawConn, err := net.Dial("tcp", srv.TCPListener.Addr().String())
r.NoError(err)
conn := rsync.NewConn(rawConn)
defer conn.Close()

_, err = doClientHandshake(conn, RsyncdServerVersion, "fake")
r.NoError(err)

// Ramp-up: trickle a few bytes, far below the floor. These are
// attributed to the grace period and must not count towards the
// first measured window.
_, err = conn.Write([]byte("hello"))
r.NoError(err)

// Wait until the grace period has elapsed.
time.Sleep(2100 * time.Millisecond)

// Transfer well above the floor for longer than one window plus one
// tick, so the first post-grace window is evaluated and passes.
deadline := time.Now().Add(2200 * time.Millisecond)
payload := make([]byte, 4000)
for time.Now().Before(deadline) {
_, err = conn.Write(payload)
r.NoError(err, "connection must remain usable after the grace period")
time.Sleep(250 * time.Millisecond)
}

assert.Zero(t, srv.getUpstreamCounters("u1").throughputFloorTerminated.Load(),
"ramp-up bytes must not trigger the throughput floor")

logData, err := os.ReadFile(accessLogPath)
r.NoError(err)
assert.NotContains(t, string(logData), "below throughput floor")
}

// TestLoadConfigPropagatesDialTimeoutAndThroughputSettings verifies
// that dial_timeout and the min_throughput_* settings are parsed,
// validated, and propagated into the dialer and Server fields. It also
Expand Down
Loading