Compare commits
No commits in common. "main" and "74f0116" have entirely different histories.
10 changed files with 60 additions and 328 deletions
|
|
@ -1,13 +0,0 @@
|
|||
/target
|
||||
/result
|
||||
/result-*
|
||||
/cache
|
||||
/config
|
||||
/fumi.rns
|
||||
.git
|
||||
.jj
|
||||
.env
|
||||
*.db
|
||||
*.db-wal
|
||||
*.db-shm
|
||||
server.key
|
||||
|
|
@ -1,22 +0,0 @@
|
|||
# Default (non-rns) build only: microReticulum's cmake/C++ build isn't set
|
||||
# up for the musl toolchain alpine gives us, mirroring packages.static in
|
||||
# flake.nix. Use the Nix flake if you need the rns feature.
|
||||
FROM rust:1-alpine AS builder
|
||||
RUN apk add --no-cache musl-dev gcc
|
||||
WORKDIR /usr/src/bunshin
|
||||
COPY Cargo.toml Cargo.lock build.rs ./
|
||||
COPY src ./src
|
||||
RUN cargo build --release --locked
|
||||
|
||||
FROM alpine:3.20
|
||||
RUN apk add --no-cache ca-certificates \
|
||||
&& adduser -D -h /data -u 10000 bunshin
|
||||
COPY --from=builder /usr/src/bunshin/target/release/bunshin /usr/local/bin/bunshin
|
||||
COPY docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh
|
||||
RUN chmod +x /usr/local/bin/docker-entrypoint.sh
|
||||
USER bunshin
|
||||
WORKDIR /data
|
||||
VOLUME /data
|
||||
EXPOSE 1961/tcp
|
||||
ENTRYPOINT ["docker-entrypoint.sh"]
|
||||
CMD ["serve", "--key", "/data/server.key", "--db", "/data/mail.db", "--host", "0.0.0.0", "--port", "1961"]
|
||||
26
README.md
26
README.md
|
|
@ -6,29 +6,11 @@ A Rust implementation of the [Smol Mail](https://code.randogoth.com/randogoth/sm
|
|||
|
||||
The server never sees plaintext, sender identities or any private key. It learns only which mailbox an envelope is for, its size, and when it arrived.
|
||||
|
||||
The quickest way to run it is the container image below. For a NixOS host, the flake also provides turnkey deployment as a special case: import `nixosModules.default`, point it at a key, and `nixos-rebuild switch`.
|
||||
|
||||
## Deploying with a container
|
||||
|
||||
Pull the published image and bring up a server in one command:
|
||||
|
||||
```
|
||||
podman run -d --name bunshin -p 1961:1961 -v bunshin-data:/data code.randogoth.com/randogoth/bunshin
|
||||
```
|
||||
|
||||
`docker` works the same way — the image is a standard OCI image either way. The entrypoint generates `/data/server.key` on first run if it's missing, then runs `serve` against `/data/server.key` and `/data/mail.db` on `0.0.0.0:1961`; `podman logs bunshin` prints the generated public key to publish to clients. A key already in the volume is left alone, so restarts and upgrades keep the same identity.
|
||||
|
||||
Pass your own arguments to run `keygen` or a customized `serve` instead of the default — they take the same flags as a bare-metal install, e.g. `podman run --rm -v bunshin-data:/data code.randogoth.com/randogoth/bunshin keygen --key /data/server.key --force`.
|
||||
|
||||
To build the image locally instead of pulling (same non-`rns` build as the release image, see the `Containerfile` header comment):
|
||||
|
||||
```
|
||||
podman build -t bunshin -f Containerfile .
|
||||
```
|
||||
The flake's main purpose is turnkey deployment on a NixOS host: import `nixosModules.default`, point it at a key, and `nixos-rebuild switch`.
|
||||
|
||||
## Deploying on NixOS
|
||||
|
||||
For a NixOS host that already manages the rest of its config with Nix, import the module instead of running the container:
|
||||
Add bunshin as a flake input and import the module:
|
||||
|
||||
```nix
|
||||
{
|
||||
|
|
@ -52,7 +34,7 @@ For a NixOS host that already manages the rest of its config with Nix, import th
|
|||
}
|
||||
```
|
||||
|
||||
`services.bunshin` also takes `host`, `port`, `dataDir`, `maxEnvelope`, `quota`, `requestsQuota`, `retentionDays`, `requestsRetentionDays`, `maxTokens`, `rateConnections`, `rateSends`, `rateTokens`, `maxConnections` and `domains`; see `flake.nix` for defaults. The module renders a `systemd` unit that runs `bunshin serve` under `DynamicUser`; unlike the container entrypoint, it does not generate a key.
|
||||
`services.bunshin` also takes `host`, `port`, `dataDir`, `maxEnvelope`, `quota`, `requestsQuota`, `retentionDays`, `requestsRetentionDays`, `maxTokens`, `rateConnections`, `rateSends`, `rateTokens` and `domains`; see `flake.nix` for defaults. The module renders a `systemd` unit that runs `bunshin serve` under `DynamicUser`; it does not generate a key.
|
||||
|
||||
For the RNS carrier, build `packages.rns`, set `services.bunshin.package` to it, and enable `services.bunshin.rns` with its `keyFile`; see [RNS.md](RNS.md) for the protocol and the remaining options.
|
||||
|
||||
|
|
@ -133,7 +115,7 @@ bunshin keygen --key server.key
|
|||
bunshin serve --key server.key --db mail.db --host 0.0.0.0 --port 1961
|
||||
```
|
||||
|
||||
`serve` accepts `--max-envelope`, `--quota`, `--requests-quota`, `--retention-days`, `--requests-retention-days`, `--max-tokens`, `--invite-token`, `--rate-connections`, `--rate-sends`, `--rate-tokens` and `--max-connections` to control size limits, the mailbox's two quota tiers, their retention, the accept-token cap, registration gating, abuse control and the concurrent-connection cap. Rate limits are sliding-window: hits age out 60 s after they happen, so a burst straddling a window boundary cannot exceed the configured rate. Connections past `--max-connections` (default 100, 0 = unlimited) are refused, mirroring the RNS carrier's link cap. Run `bunshin serve --help` for defaults.
|
||||
`serve` accepts `--max-envelope`, `--quota`, `--requests-quota`, `--retention-days`, `--requests-retention-days`, `--max-tokens`, `--invite-token`, `--rate-connections`, `--rate-sends` and `--rate-tokens` to control size limits, the mailbox's two quota tiers, their retention, the accept-token cap, registration gating and abuse control. Run `bunshin serve --help` for defaults.
|
||||
|
||||
`--verbose` (or `services.bunshin.verbose` in the NixOS module) raises logging to debug level, which adds server-side metrics on both carriers: per-operation timing, request and response sizes, handshake duration, per-session summaries, and RNS link events. The default info level stays quiet on success.
|
||||
|
||||
|
|
|
|||
2
RNS.md
2
RNS.md
|
|
@ -16,7 +16,7 @@ Implemented behind the `rns` cargo feature (protocol 1.2; the default build stil
|
|||
| Ops, status codes, envelope | §6, §12 | identical |
|
||||
| AUTH / REGISTER binding | Noise handshake hash, server static key | derived, see §2 |
|
||||
| Announce | — | at startup and every 2 h; interfaces rate-limit to ≈1/hour, which is the ceiling |
|
||||
| Abuse control | per-IP sliding-window rate limits, concurrent-connection cap | per-link request + byte limits, concurrent-link cap |
|
||||
| Abuse control | per-IP rate limits | per-link request + byte limits, concurrent-link cap |
|
||||
|
||||
One wire-format detail the upstream spec leaves implicit: microReticulum splices the request payload and the response into their msgpack envelopes verbatim, so both directions carry the smolmail payload as a msgpack binary. The shim unpacks on the way in and packs on the way out; the Rust side only ever sees `op u8 || body` and `status u8 || payload`.
|
||||
|
||||
|
|
|
|||
|
|
@ -1,12 +0,0 @@
|
|||
#!/bin/sh
|
||||
# Generates the server key on first run so a bare `docker run` against an
|
||||
# empty volume works; a key already at BUNSHIN_KEY is left untouched.
|
||||
set -e
|
||||
|
||||
: "${BUNSHIN_KEY:=/data/server.key}"
|
||||
|
||||
if [ "$1" = "serve" ] && [ ! -f "$BUNSHIN_KEY" ]; then
|
||||
bunshin keygen --key "$BUNSHIN_KEY"
|
||||
fi
|
||||
|
||||
exec bunshin "$@"
|
||||
81
flake.nix
81
flake.nix
|
|
@ -109,17 +109,15 @@
|
|||
name = "bunshin";
|
||||
};
|
||||
|
||||
# Builds packages.static and publishes it as a Forgejo release on
|
||||
# this repo, tagged by short commit hash (creating the release if
|
||||
# it doesn't exist yet, replacing the asset if it does — safe to
|
||||
# rerun for the same commit). Needs FORGEJO_TOKEN (a token with
|
||||
# write:repository scope) in the environment; run from a checkout
|
||||
# so `git rev-parse` and the `.` flake ref resolve to the right
|
||||
# place.
|
||||
# Builds packages.static and uploads it to this repo owner's
|
||||
# Forgejo generic package registry, tagged by short commit hash.
|
||||
# Needs FORGEJO_TOKEN (a token with write:package scope) in the
|
||||
# environment; run from a checkout so `git rev-parse` and the `.`
|
||||
# flake ref resolve to the right place.
|
||||
apps.release-static = flake-utils.lib.mkApp {
|
||||
drv = pkgs.writeShellApplication {
|
||||
name = "bunshin-release-static";
|
||||
runtimeInputs = [ pkgs.nix pkgs.curl pkgs.git pkgs.jq ];
|
||||
runtimeInputs = [ pkgs.nix pkgs.curl pkgs.git ];
|
||||
text = ''
|
||||
if [ -f .env ]; then
|
||||
set -a
|
||||
|
|
@ -127,63 +125,12 @@
|
|||
. ./.env
|
||||
set +a
|
||||
fi
|
||||
: "''${FORGEJO_TOKEN:?set FORGEJO_TOKEN (env or .env) to a Forgejo token with write:repository scope}"
|
||||
|
||||
rev=$(git rev-parse --short HEAD)
|
||||
sha=$(git rev-parse HEAD)
|
||||
out=$(nix build .#static --no-link --print-out-paths)
|
||||
bin="$out/bin/bunshin"
|
||||
|
||||
api="https://code.randogoth.com/api/v1/repos/randogoth/bunshin"
|
||||
auth=(-H "Authorization: token ''${FORGEJO_TOKEN}")
|
||||
|
||||
release_id=$(curl -sS "''${auth[@]}" "$api/releases/tags/$rev" | jq -r '.id // empty')
|
||||
if [ -z "$release_id" ]; then
|
||||
release_id=$(curl -sSf "''${auth[@]}" -H "Content-Type: application/json" \
|
||||
-d "$(jq -n --arg tag "$rev" --arg sha "$sha" \
|
||||
'{tag_name:$tag, target_commitish:$sha, name:$tag, body:"Static musl build.", draft:false, prerelease:false}')" \
|
||||
"$api/releases" | jq -r '.id')
|
||||
fi
|
||||
|
||||
asset_id=$(curl -sSf "''${auth[@]}" "$api/releases/$release_id/assets" | jq -r '.[] | select(.name=="bunshin") | .id' | head -1)
|
||||
if [ -n "$asset_id" ]; then
|
||||
curl -sSf -X DELETE "''${auth[@]}" "$api/releases/$release_id/assets/$asset_id" >/dev/null
|
||||
fi
|
||||
curl -sSf "''${auth[@]}" -F "attachment=@$bin;filename=bunshin" "$api/releases/$release_id/assets?name=bunshin" >/dev/null
|
||||
|
||||
echo "released: https://code.randogoth.com/randogoth/bunshin/releases/tag/$rev"
|
||||
'';
|
||||
};
|
||||
};
|
||||
|
||||
# Builds the Containerfile image (same non-rns build the release
|
||||
# binary is, see its header comment) and pushes it to this repo's
|
||||
# Forgejo container registry, tagged by short commit hash and
|
||||
# `latest`. Needs FORGEJO_TOKEN (a token with write:package scope)
|
||||
# in the environment; run from a checkout so `git rev-parse` and
|
||||
# the Containerfile build context resolve to the right place.
|
||||
apps.release-container = flake-utils.lib.mkApp {
|
||||
drv = pkgs.writeShellApplication {
|
||||
name = "bunshin-release-container";
|
||||
runtimeInputs = [ pkgs.podman pkgs.git ];
|
||||
text = ''
|
||||
if [ -f .env ]; then
|
||||
set -a
|
||||
# shellcheck disable=SC1091
|
||||
. ./.env
|
||||
set +a
|
||||
fi
|
||||
: "''${FORGEJO_TOKEN:?set FORGEJO_TOKEN (env or .env) to a Forgejo token with write:package scope}"
|
||||
|
||||
rev=$(git rev-parse --short HEAD)
|
||||
registry="code.randogoth.com/randogoth/bunshin"
|
||||
|
||||
echo "''${FORGEJO_TOKEN}" | podman login code.randogoth.com -u randogoth --password-stdin
|
||||
podman build -f Containerfile -t "$registry:$rev" -t "$registry:latest" .
|
||||
podman push "$registry:$rev"
|
||||
podman push "$registry:latest"
|
||||
|
||||
echo "pushed: https://code.randogoth.com/randogoth/-/packages/container/bunshin"
|
||||
url="https://code.randogoth.com/api/packages/randogoth/generic/bunshin/$rev/bunshin"
|
||||
curl -sSf -H "Authorization: token ''${FORGEJO_TOKEN}" --upload-file "$out/bin/bunshin" "$url"
|
||||
echo "uploaded: $url"
|
||||
'';
|
||||
};
|
||||
};
|
||||
|
|
@ -204,7 +151,7 @@
|
|||
|
||||
package = mkOption {
|
||||
type = types.package;
|
||||
default = self.packages.${pkgs.stdenv.hostPlatform.system}.default;
|
||||
default = self.packages.${pkgs.system}.default;
|
||||
description = ''
|
||||
bunshin package to run. Use `packages.rns` when the RNS
|
||||
carrier is enabled: the default package is built without it.
|
||||
|
|
@ -355,12 +302,6 @@
|
|||
description = "Max SEND operations per minute, per accept token.";
|
||||
};
|
||||
|
||||
maxConnections = mkOption {
|
||||
type = types.ints.unsigned;
|
||||
default = 100;
|
||||
description = "Max concurrently served TCP connections; past the cap connections are refused. 0 means unlimited.";
|
||||
};
|
||||
|
||||
inviteToken = mkOption {
|
||||
type = types.nullOr types.str;
|
||||
default = null;
|
||||
|
|
@ -514,7 +455,6 @@
|
|||
rate_connections = cfg.rateConnections;
|
||||
rate_sends = cfg.rateSends;
|
||||
rate_tokens = cfg.rateTokens;
|
||||
max_connections = cfg.maxConnections;
|
||||
}
|
||||
// lib.optionalAttrs (cfg.inviteToken != null) { invite_token = cfg.inviteToken; }
|
||||
// lib.optionalAttrs (cfg.inviteTokenFile != null) {
|
||||
|
|
@ -562,7 +502,6 @@
|
|||
--rate-connections ${toString cfg.rateConnections}
|
||||
--rate-sends ${toString cfg.rateSends}
|
||||
--rate-tokens ${toString cfg.rateTokens}
|
||||
--max-connections ${toString cfg.maxConnections}
|
||||
)
|
||||
${lib.optionalString cfg.rns.enable ''
|
||||
args+=(
|
||||
|
|
|
|||
|
|
@ -25,7 +25,6 @@ const DEFAULT_MAX_TOKENS: u16 = 1024;
|
|||
const DEFAULT_RATE_CONNECTIONS: u32 = 120;
|
||||
const DEFAULT_RATE_SENDS: u32 = 60;
|
||||
const DEFAULT_RATE_TOKENS: u32 = 30;
|
||||
const DEFAULT_MAX_CONNECTIONS: usize = 100;
|
||||
|
||||
/// One resolved domain, ready to serve.
|
||||
pub struct DomainConfig {
|
||||
|
|
@ -44,7 +43,6 @@ pub struct DomainConfig {
|
|||
pub rate_connections: u32,
|
||||
pub rate_sends: u32,
|
||||
pub rate_tokens: u32,
|
||||
pub max_connections: usize,
|
||||
}
|
||||
|
||||
/// The union of every setting either table accepts; `deny_unknown_fields`
|
||||
|
|
@ -69,7 +67,6 @@ struct Table {
|
|||
rate_connections: Option<u32>,
|
||||
rate_sends: Option<u32>,
|
||||
rate_tokens: Option<u32>,
|
||||
max_connections: Option<usize>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
|
|
@ -190,10 +187,6 @@ pub fn load(path: &str) -> anyhow::Result<Vec<DomainConfig>> {
|
|||
.rate_tokens
|
||||
.or(file.defaults.rate_tokens)
|
||||
.unwrap_or(DEFAULT_RATE_TOKENS),
|
||||
max_connections: table
|
||||
.max_connections
|
||||
.or(file.defaults.max_connections)
|
||||
.unwrap_or(DEFAULT_MAX_CONNECTIONS),
|
||||
});
|
||||
}
|
||||
Ok(domains)
|
||||
|
|
@ -225,7 +218,6 @@ mod tests {
|
|||
assert_eq!(d.max_envelope, 768 << 10);
|
||||
assert_eq!(d.quota, 64 << 20);
|
||||
assert_eq!(d.rate_tokens, 30);
|
||||
assert_eq!(d.max_connections, 100);
|
||||
assert!(d.invite_token.is_none());
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ enum Command {
|
|||
Serve {
|
||||
/// Serve one domain per [domains.<name>] table; every domain has its
|
||||
/// own key, port and mailbox database
|
||||
#[arg(long, conflicts_with_all = ["key", "db", "host", "port", "max_envelope", "quota", "requests_quota", "retention_days", "requests_retention_days", "max_tokens", "invite_token", "rate_connections", "rate_sends", "rate_tokens", "max_connections"])]
|
||||
#[arg(long, conflicts_with_all = ["key", "db", "host", "port", "max_envelope", "quota", "requests_quota", "retention_days", "requests_retention_days", "max_tokens", "invite_token", "rate_connections", "rate_sends", "rate_tokens"])]
|
||||
#[cfg_attr(feature = "rns", arg(conflicts_with_all = ["enabled", "rns_key", "rns_max_envelope", "rns_fetch_budget", "rns_max_links", "rns_rate_link_requests", "rns_rate_link_bytes", "rns_link_idle", "rns_udp", "rns_udp_forward"]))]
|
||||
config: Option<String>,
|
||||
#[arg(long, default_value = "server.key")]
|
||||
|
|
@ -76,9 +76,6 @@ enum Command {
|
|||
rate_sends: u32,
|
||||
#[arg(long = "rate-tokens", default_value_t = 30)]
|
||||
rate_tokens: u32,
|
||||
/// Refuse connections past this many concurrent ones (0 = unlimited)
|
||||
#[arg(long = "max-connections", default_value_t = 100)]
|
||||
max_connections: usize,
|
||||
#[cfg(feature = "rns")]
|
||||
#[command(flatten)]
|
||||
rns: RnsServeArgs,
|
||||
|
|
@ -156,7 +153,6 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
max_connections,
|
||||
} => {
|
||||
if let Some(path) = config {
|
||||
return server::serve_domains(config::load(&path)?);
|
||||
|
|
@ -176,7 +172,6 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
max_connections,
|
||||
})
|
||||
}
|
||||
#[cfg(feature = "rns")]
|
||||
|
|
@ -196,7 +191,6 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
max_connections,
|
||||
rns,
|
||||
} => {
|
||||
if let Some(path) = config {
|
||||
|
|
@ -254,7 +248,6 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
max_connections,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
124
src/ratelimit.rs
124
src/ratelimit.rs
|
|
@ -1,18 +1,20 @@
|
|||
//! Sliding-window per-IP counters, the whole of the server's abuse control.
|
||||
//! Fixed-window per-IP counter, the whole of the server's abuse control.
|
||||
//!
|
||||
//! A server cannot see senders, so quotas, size caps and this are all it has.
|
||||
//! Hits are logged per key and pruned once they age out of the window, so a
|
||||
//! burst straddling a window boundary cannot exceed the configured rate the
|
||||
//! way a fixed-window counter would.
|
||||
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Mutex;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
struct Window {
|
||||
start: Instant,
|
||||
count: u32,
|
||||
}
|
||||
|
||||
pub struct RateLimiter {
|
||||
limit: u32,
|
||||
window: Duration,
|
||||
hits: Mutex<HashMap<String, VecDeque<Instant>>>,
|
||||
hits: Mutex<HashMap<String, Window>>,
|
||||
}
|
||||
|
||||
impl RateLimiter {
|
||||
|
|
@ -25,44 +27,44 @@ impl RateLimiter {
|
|||
}
|
||||
|
||||
pub fn allow(&self, ip: &str) -> bool {
|
||||
self.allow_at(ip, Instant::now())
|
||||
}
|
||||
|
||||
fn allow_at(&self, key: &str, now: Instant) -> bool {
|
||||
if self.limit == 0 {
|
||||
return true;
|
||||
}
|
||||
let now = Instant::now();
|
||||
let mut hits = self.hits.lock().unwrap();
|
||||
let entry = hits.entry(key.to_string()).or_default();
|
||||
prune(entry, now, self.window);
|
||||
if entry.len() >= self.limit as usize {
|
||||
let entry = hits.entry(ip.to_string()).or_insert(Window {
|
||||
start: now,
|
||||
count: 0,
|
||||
});
|
||||
if now.duration_since(entry.start) >= self.window {
|
||||
entry.start = now;
|
||||
entry.count = 0;
|
||||
}
|
||||
if entry.count >= self.limit {
|
||||
return false;
|
||||
}
|
||||
entry.push_back(now);
|
||||
entry.count += 1;
|
||||
if hits.len() > 4096 {
|
||||
let window = self.window;
|
||||
hits.retain(|_, e| !e.is_empty() && now.duration_since(*e.back().unwrap()) <= window);
|
||||
hits.retain(|_, w| now.duration_since(w.start) < window);
|
||||
}
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
fn prune(entry: &mut VecDeque<Instant>, now: Instant, window: Duration) {
|
||||
while entry
|
||||
.front()
|
||||
.is_some_and(|t| now.duration_since(*t) > window)
|
||||
{
|
||||
entry.pop_front();
|
||||
}
|
||||
}
|
||||
|
||||
/// Sliding-window byte counter for per-link transfer budgets: `RateLimiter`
|
||||
/// Fixed-window byte counter for per-link transfer budgets: `RateLimiter`
|
||||
/// counts events, this counts bytes, so it gets its own small type.
|
||||
#[cfg(feature = "rns")]
|
||||
pub struct ByteRateLimiter {
|
||||
limit: u64,
|
||||
window: Duration,
|
||||
hits: Mutex<HashMap<String, VecDeque<(Instant, u64)>>>,
|
||||
hits: Mutex<HashMap<String, ByteWindow>>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "rns")]
|
||||
struct ByteWindow {
|
||||
start: Instant,
|
||||
bytes: u64,
|
||||
}
|
||||
|
||||
#[cfg(feature = "rns")]
|
||||
|
|
@ -77,29 +79,26 @@ impl ByteRateLimiter {
|
|||
|
||||
/// Records `bytes` against `key` if the window still has room for them.
|
||||
pub fn allow(&self, key: &str, bytes: usize) -> bool {
|
||||
self.allow_at(key, bytes, Instant::now())
|
||||
}
|
||||
|
||||
fn allow_at(&self, key: &str, bytes: usize, now: Instant) -> bool {
|
||||
if self.limit == 0 {
|
||||
return true;
|
||||
}
|
||||
let now = Instant::now();
|
||||
let mut hits = self.hits.lock().unwrap();
|
||||
let entry = hits.entry(key.to_string()).or_default();
|
||||
while entry
|
||||
.front()
|
||||
.is_some_and(|(t, _)| now.duration_since(*t) > self.window)
|
||||
{
|
||||
entry.pop_front();
|
||||
let entry = hits.entry(key.to_string()).or_insert(ByteWindow {
|
||||
start: now,
|
||||
bytes: 0,
|
||||
});
|
||||
if now.duration_since(entry.start) >= self.window {
|
||||
entry.start = now;
|
||||
entry.bytes = 0;
|
||||
}
|
||||
let used: u64 = entry.iter().map(|(_, n)| *n).sum();
|
||||
if used + bytes as u64 > self.limit {
|
||||
if entry.bytes + bytes as u64 > self.limit {
|
||||
return false;
|
||||
}
|
||||
entry.push_back((now, bytes as u64));
|
||||
entry.bytes += bytes as u64;
|
||||
if hits.len() > 4096 {
|
||||
let window = self.window;
|
||||
hits.retain(|_, e| !e.is_empty() && now.duration_since(e.back().unwrap().0) <= window);
|
||||
hits.retain(|_, w| now.duration_since(w.start) < window);
|
||||
}
|
||||
true
|
||||
}
|
||||
|
|
@ -112,10 +111,9 @@ mod tests {
|
|||
#[test]
|
||||
fn allows_up_to_limit_then_blocks() {
|
||||
let rl = RateLimiter::new(2);
|
||||
let now = Instant::now();
|
||||
assert!(rl.allow_at("1.2.3.4", now));
|
||||
assert!(rl.allow_at("1.2.3.4", now));
|
||||
assert!(!rl.allow_at("1.2.3.4", now));
|
||||
assert!(rl.allow("1.2.3.4"));
|
||||
assert!(rl.allow("1.2.3.4"));
|
||||
assert!(!rl.allow("1.2.3.4"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -133,42 +131,4 @@ mod tests {
|
|||
assert!(rl.allow("2.2.2.2"));
|
||||
assert!(!rl.allow("1.1.1.1"));
|
||||
}
|
||||
|
||||
/// The regression the sliding window exists for: a burst straddling a
|
||||
/// window boundary must not be handed a fresh window, only the slots
|
||||
/// its own hits have aged out. A hit counts for the whole window,
|
||||
/// including the instant it is exactly window old.
|
||||
#[test]
|
||||
fn boundary_burst_cannot_double_the_rate() {
|
||||
let rl = RateLimiter::new(2);
|
||||
let t0 = Instant::now();
|
||||
let window = Duration::from_secs(60);
|
||||
assert!(rl.allow_at("1.2.3.4", t0));
|
||||
assert!(rl.allow_at("1.2.3.4", t0 + Duration::from_secs(1)));
|
||||
// At the boundary a fixed-window counter would reset and admit
|
||||
// two more; both hits are still window-old or fresher.
|
||||
assert!(!rl.allow_at("1.2.3.4", t0 + window));
|
||||
// Half a second later the t0 hit has aged out but the +1s hit
|
||||
// has not, so exactly one slot is free.
|
||||
let later = t0 + window + Duration::from_millis(500);
|
||||
assert!(rl.allow_at("1.2.3.4", later));
|
||||
assert!(!rl.allow_at("1.2.3.4", later));
|
||||
}
|
||||
|
||||
#[cfg(feature = "rns")]
|
||||
#[test]
|
||||
fn byte_boundary_burst_cannot_double_the_budget() {
|
||||
let bl = ByteRateLimiter::new(100);
|
||||
let t0 = Instant::now();
|
||||
let window = Duration::from_secs(60);
|
||||
assert!(bl.allow_at("link", 60, t0));
|
||||
assert!(bl.allow_at("link", 40, t0 + Duration::from_secs(1)));
|
||||
// At the boundary both hits still count: 100 bytes used, no room.
|
||||
assert!(!bl.allow_at("link", 1, t0 + window));
|
||||
// Half a second later only the 60-byte hit has aged out, freeing
|
||||
// exactly its 60 bytes.
|
||||
let later = t0 + window + Duration::from_millis(500);
|
||||
assert!(bl.allow_at("link", 60, later));
|
||||
assert!(!bl.allow_at("link", 1, later));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
//! TCP accept loop, per-connection handling, and the background purge loop.
|
||||
|
||||
use std::net::{TcpListener, TcpStream};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::bind::TransportBindValues;
|
||||
|
|
@ -28,7 +28,6 @@ pub struct ServeArgs {
|
|||
pub rate_connections: u32,
|
||||
pub rate_sends: u32,
|
||||
pub rate_tokens: u32,
|
||||
pub max_connections: usize,
|
||||
}
|
||||
|
||||
/// The single-domain CLI path: one anonymous domain.
|
||||
|
|
@ -49,7 +48,6 @@ pub fn run(args: ServeArgs) -> anyhow::Result<()> {
|
|||
rate_connections: args.rate_connections,
|
||||
rate_sends: args.rate_sends,
|
||||
rate_tokens: args.rate_tokens,
|
||||
max_connections: args.max_connections,
|
||||
}])
|
||||
}
|
||||
|
||||
|
|
@ -96,9 +94,8 @@ pub fn serve_domains(domains: Vec<DomainConfig>) -> anyhow::Result<()> {
|
|||
purge_loop(purge_db_path, main_retention_secs, requests_retention_secs)
|
||||
});
|
||||
|
||||
let gate = Arc::new(ConnGate::new(domain.max_connections));
|
||||
let db_path = domain.db_path.clone();
|
||||
std::thread::spawn(move || accept_loop(listener, config, static_key, db_path, gate));
|
||||
std::thread::spawn(move || accept_loop(listener, config, static_key, db_path));
|
||||
}
|
||||
|
||||
// Each domain's accept loop runs in its own thread; nothing fails here.
|
||||
|
|
@ -123,7 +120,6 @@ fn accept_loop(
|
|||
config: Arc<ServerConfig>,
|
||||
static_key: [u8; KEY_LEN],
|
||||
db_path: String,
|
||||
gate: Arc<ConnGate>,
|
||||
) {
|
||||
for incoming in listener.incoming() {
|
||||
let stream = match incoming {
|
||||
|
|
@ -143,19 +139,9 @@ fn accept_loop(
|
|||
continue;
|
||||
}
|
||||
|
||||
let permit = match gate.try_acquire() {
|
||||
Some(p) => p,
|
||||
None => {
|
||||
log::warn!("connection limit reached, refusing {peer_ip}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let config = Arc::clone(&config);
|
||||
let db_path = db_path.clone();
|
||||
std::thread::spawn(move || {
|
||||
// The permit is dropped with the connection, freeing its slot.
|
||||
let _permit = permit;
|
||||
if let Err(e) = handle_connection(stream, &config, &static_key, &db_path, &peer_ip) {
|
||||
log::info!("connection error from {peer_ip}: {e}");
|
||||
}
|
||||
|
|
@ -163,55 +149,6 @@ fn accept_loop(
|
|||
}
|
||||
}
|
||||
|
||||
/// Cap on concurrently served TCP connections, the mirror of the RNS
|
||||
/// carrier's link cap (SPEC.md sec 13.8): without one, a connection flood
|
||||
/// would exhaust threads one unbounded spawn at a time. `max == 0` disables
|
||||
/// the cap; past it, connections are refused, never queued.
|
||||
struct ConnGate {
|
||||
max: usize,
|
||||
active: Mutex<usize>,
|
||||
}
|
||||
|
||||
impl ConnGate {
|
||||
fn new(max: usize) -> Self {
|
||||
ConnGate {
|
||||
max,
|
||||
active: Mutex::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
fn try_acquire(self: &Arc<Self>) -> Option<ConnPermit> {
|
||||
if self.max == 0 {
|
||||
return Some(ConnPermit {
|
||||
gate: Arc::clone(self),
|
||||
});
|
||||
}
|
||||
let mut active = self.active.lock().unwrap();
|
||||
if *active >= self.max {
|
||||
return None;
|
||||
}
|
||||
*active += 1;
|
||||
Some(ConnPermit {
|
||||
gate: Arc::clone(self),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// Dropping releases the slot, so the count stays accurate whatever return
|
||||
/// path or panic closes the connection.
|
||||
struct ConnPermit {
|
||||
gate: Arc<ConnGate>,
|
||||
}
|
||||
|
||||
impl Drop for ConnPermit {
|
||||
fn drop(&mut self) {
|
||||
if self.gate.max == 0 {
|
||||
return;
|
||||
}
|
||||
*self.gate.active.lock().unwrap() -= 1;
|
||||
}
|
||||
}
|
||||
|
||||
/// Also the hostile harness's entry point: it drives real connections
|
||||
/// through the same accept/handshake/session path the TCP carrier serves.
|
||||
pub(crate) fn handle_connection(
|
||||
|
|
@ -328,27 +265,3 @@ fn purge_loop(db_path: String, main_retention_secs: i64, requests_retention_secs
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn gate_refuses_past_cap_and_frees_on_drop() {
|
||||
let gate = Arc::new(ConnGate::new(2));
|
||||
let a = gate.try_acquire().unwrap();
|
||||
let b = gate.try_acquire().unwrap();
|
||||
assert!(gate.try_acquire().is_none());
|
||||
drop(b);
|
||||
assert!(gate.try_acquire().is_some());
|
||||
drop(a);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zero_cap_means_unlimited() {
|
||||
let gate = Arc::new(ConnGate::new(0));
|
||||
for _ in 0..100 {
|
||||
assert!(gate.try_acquire().is_some());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue