Compare commits
2
Commits
f665ab17bc
...
304cf79826
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
304cf79826 | ||
|
|
d2f219ca19 |
@@ -83,6 +83,12 @@ VOLUME /data
|
|||||||
# The default public port; PORT changes it.
|
# The default public port; PORT changes it.
|
||||||
EXPOSE 8080
|
EXPOSE 8080
|
||||||
|
|
||||||
|
# Requests the backend's health check through nginx, on the port from
|
||||||
|
# PORT, so it fails unless both answer. upaas reads the result 60
|
||||||
|
# seconds after a deploy and fails the deploy unless it is healthy.
|
||||||
|
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
|
||||||
|
CMD wget -q -O /dev/null "http://127.0.0.1:${PORT:-8080}/.well-known/healthcheck"
|
||||||
|
|
||||||
# The nginx image stops its container with SIGQUIT; the entrypoint
|
# The nginx image stops its container with SIGQUIT; the entrypoint
|
||||||
# acts on TERM and INT.
|
# acts on TERM and INT.
|
||||||
STOPSIGNAL SIGTERM
|
STOPSIGNAL SIGTERM
|
||||||
|
|||||||
@@ -194,6 +194,46 @@ only inside the container, on `127.0.0.1:8081`. The image:
|
|||||||
- Writes buffered reports to disk on `docker stop`, and exits non-zero if nginx
|
- Writes buffered reports to disk on `docker stop`, and exits non-zero if nginx
|
||||||
or the backend exits on its own, so the platform restarts it
|
or the backend exits on its own, so the platform restarts it
|
||||||
|
|
||||||
|
## Running under upaas
|
||||||
|
|
||||||
|
What the [upaas](https://git.eeqj.de/sneak/upaas) app for netwatch needs:
|
||||||
|
|
||||||
|
- **Port:** container port `8080`.
|
||||||
|
- **Volume:** container path `/data`; the reports are kept in `/data/reports`.
|
||||||
|
- **First run:** upaas bind-mounts the host directory it is given and does not
|
||||||
|
create it, and the backend, which runs as uid 1000, does not start unless it
|
||||||
|
can write there. Create the directory, owned by uid 1000, before the first
|
||||||
|
deploy:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
mkdir -p /path/to/data
|
||||||
|
chown 1000:1000 /path/to/data
|
||||||
|
```
|
||||||
|
|
||||||
|
- **Environment variables:** none is required. An empty one counts as unset, and
|
||||||
|
one set to a value netwatch cannot use stops the container at start, with the
|
||||||
|
reason in its log.
|
||||||
|
- `PORT`, default `8080`: the container port, from 1 to 65535. `8081` cannot
|
||||||
|
be used: the backend listens on it inside the container
|
||||||
|
- `REPORTS_PER_MINUTE`, default `60`: reports each client address may send a
|
||||||
|
minute
|
||||||
|
- `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the
|
||||||
|
report files may take
|
||||||
|
- `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call
|
||||||
|
the API
|
||||||
|
- `DEBUG`, default `false`: debug logging
|
||||||
|
- `DATA_DIR`, default `/data/reports`: leave unset; reports kept outside
|
||||||
|
`/data` do not survive a redeploy
|
||||||
|
- `TRUSTED_PROXIES`, default loopback and RFC1918: leave unset. The
|
||||||
|
backend's only client is nginx, on loopback, which passes on the client
|
||||||
|
address; nginx takes it from `X-Forwarded-For` only from RFC1918
|
||||||
|
addresses.
|
||||||
|
- **Health check:** the image's `HEALTHCHECK` requests
|
||||||
|
`/.well-known/healthcheck` through nginx every 30 seconds, so it fails unless
|
||||||
|
both nginx and the backend answer. upaas reads the container's health 60
|
||||||
|
seconds after a deploy and fails the deploy unless it is `healthy`. The
|
||||||
|
container also stops when either process exits.
|
||||||
|
|
||||||
## Browser Compatibility
|
## Browser Compatibility
|
||||||
|
|
||||||
Requires a modern browser with ES modules, Fetch API, Canvas API, and CSS custom
|
Requires a modern browser with ES modules, Fetch API, Canvas API, and CSS custom
|
||||||
|
|||||||
@@ -23,6 +23,22 @@ latest run passes.
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
- 2026-09-29: report file names can no longer collide (issue #61): each is
|
||||||
|
`reports-<timestamp>-<number>.jsonl.zst`, where the number goes up by one for
|
||||||
|
each file the server starts to write, so two flushes in the same millisecond,
|
||||||
|
such as a flush for size and the final flush at shutdown, each get a file of
|
||||||
|
their own instead of the second one failing. A failed write leaves a gap in
|
||||||
|
the numbers.
|
||||||
|
- 2026-09-29: ready to run under upaas (issue #59): the image has a
|
||||||
|
`HEALTHCHECK` that requests `/.well-known/healthcheck` through nginx on the
|
||||||
|
port from `PORT`. The backend no longer reads a bad `PORT` as 0 or a bad
|
||||||
|
`DEBUG` as false: those, and a `BIND_ADDRESS` that is not an IP address, stop
|
||||||
|
it from starting with an error naming the variable, as the limits,
|
||||||
|
`CORS_ALLOWED_ORIGINS` and, now by name, `TRUSTED_PROXIES` already did.
|
||||||
|
`bin/entrypoint.sh` also refuses a `PORT` outside 1 to 65535, and `8081`,
|
||||||
|
where the backend listens inside the container, naming `PORT`. `README.md` has
|
||||||
|
a "Running under upaas" section, whose first-run steps create the host
|
||||||
|
directory for `/data` owned by uid 1000; the image does not change its owner
|
||||||
- 2026-09-29: nginx listens on `PORT` (issue #26), 8080 when unset or empty: the
|
- 2026-09-29: nginx listens on `PORT` (issue #26), 8080 when unset or empty: the
|
||||||
nginx image renders `nginx.conf` as a template at container start, filling in
|
nginx image renders `nginx.conf` as a template at container start, filling in
|
||||||
`PORT` and no other variable. `bin/entrypoint.sh` refuses to start when `PORT`
|
`PORT` and no other variable. `bin/entrypoint.sh` refuses to start when `PORT`
|
||||||
|
|||||||
+11
-3
@@ -91,6 +91,10 @@ The loopback entries cover the reverse proxy that shares the container; the
|
|||||||
RFC1918 ranges match `nginx.conf`. A request whose direct peer is outside this
|
RFC1918 ranges match `nginx.conf`. A request whose direct peer is outside this
|
||||||
set has its forwarded headers ignored, and the direct peer is logged instead.
|
set has its forwarded headers ignored, and the direct peer is logged instead.
|
||||||
|
|
||||||
|
A variable set to a value the server cannot use, such as `PORT=abc`,
|
||||||
|
`DEBUG=maybe` or a `BIND_ADDRESS` that is not an IP address, stops it from
|
||||||
|
starting, with an error naming the variable. An empty variable counts as unset.
|
||||||
|
|
||||||
### Container image
|
### Container image
|
||||||
|
|
||||||
The root `Dockerfile` builds one image in which nginx listens on the public port
|
The root `Dockerfile` builds one image in which nginx listens on the public port
|
||||||
@@ -102,9 +106,13 @@ which `netwatch` owns.
|
|||||||
|
|
||||||
### Report storage
|
### Report storage
|
||||||
|
|
||||||
Reports are written as `reports-<timestamp>.jsonl.zst` files in `DATA_DIR`.
|
Reports are written as `reports-<timestamp>-<number>.jsonl.zst` files in
|
||||||
Each file contains one JSON object per line, compressed with zstd. Files are
|
`DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by
|
||||||
created with `O_EXCL` to prevent overwrites.
|
time. The number starts at 1 when the server starts and goes up by one for each
|
||||||
|
file the server starts to write, so two files written in the same millisecond
|
||||||
|
still get different names, and a failed write leaves a gap in the numbers. Each
|
||||||
|
file contains one JSON object per line, compressed with zstd. Files are created
|
||||||
|
with `O_EXCL` to prevent overwrites.
|
||||||
|
|
||||||
### Report limits
|
### Report limits
|
||||||
|
|
||||||
|
|||||||
@@ -6,7 +6,10 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"math"
|
||||||
|
"net/netip"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"sneak.berlin/go/netwatch/internal/globals"
|
"sneak.berlin/go/netwatch/internal/globals"
|
||||||
@@ -37,6 +40,9 @@ var (
|
|||||||
errNotOrigin = errors.New(
|
errNotOrigin = errors.New(
|
||||||
"must be an origin, scheme://host with an optional port",
|
"must be an origin, scheme://host with an optional port",
|
||||||
)
|
)
|
||||||
|
errNotPort = errors.New("must be a port number, 1 to 65535")
|
||||||
|
errNotBool = errors.New("must be true or false")
|
||||||
|
errNotIP = errors.New("must be an IP address, or empty")
|
||||||
)
|
)
|
||||||
|
|
||||||
// Params defines the dependencies for Config.
|
// Params defines the dependencies for Config.
|
||||||
@@ -65,7 +71,8 @@ type Config struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// New loads configuration from env, .env files, and config
|
// New loads configuration from env, .env files, and config
|
||||||
// files, returning a fully resolved Config.
|
// files, returning a fully resolved Config. It fails, with an error
|
||||||
|
// naming the setting, on a value the server cannot use.
|
||||||
func New(
|
func New(
|
||||||
_ fx.Lifecycle,
|
_ fx.Lifecycle,
|
||||||
params Params,
|
params Params,
|
||||||
@@ -103,15 +110,29 @@ func New(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Read with strconv: viper's GetInt and GetBool would read a value
|
||||||
|
// they cannot parse as 0 or false instead of failing.
|
||||||
|
port, err := strconv.Atoi(viper.GetString("PORT"))
|
||||||
|
if err != nil || port < 1 || port > math.MaxUint16 {
|
||||||
|
return nil, fmt.Errorf("PORT %q: %w",
|
||||||
|
viper.GetString("PORT"), errNotPort)
|
||||||
|
}
|
||||||
|
|
||||||
|
debug, err := strconv.ParseBool(viper.GetString("DEBUG"))
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("DEBUG %q: %w",
|
||||||
|
viper.GetString("DEBUG"), errNotBool)
|
||||||
|
}
|
||||||
|
|
||||||
s := &Config{
|
s := &Config{
|
||||||
BindAddress: viper.GetString("BIND_ADDRESS"),
|
BindAddress: viper.GetString("BIND_ADDRESS"),
|
||||||
CORSAllowedOrigins: splitList(viper.GetString("CORS_ALLOWED_ORIGINS")),
|
CORSAllowedOrigins: splitList(viper.GetString("CORS_ALLOWED_ORIGINS")),
|
||||||
DataDir: viper.GetString("DATA_DIR"),
|
DataDir: viper.GetString("DATA_DIR"),
|
||||||
DataDirMaxBytes: viper.GetInt64("DATA_DIR_MAX_BYTES"),
|
DataDirMaxBytes: viper.GetInt64("DATA_DIR_MAX_BYTES"),
|
||||||
Debug: viper.GetBool("DEBUG"),
|
Debug: debug,
|
||||||
MetricsPassword: viper.GetString("METRICS_PASSWORD"),
|
MetricsPassword: viper.GetString("METRICS_PASSWORD"),
|
||||||
MetricsUsername: viper.GetString("METRICS_USERNAME"),
|
MetricsUsername: viper.GetString("METRICS_USERNAME"),
|
||||||
Port: viper.GetInt("PORT"),
|
Port: port,
|
||||||
ReportsPerMinute: viper.GetInt("REPORTS_PER_MINUTE"),
|
ReportsPerMinute: viper.GetInt("REPORTS_PER_MINUTE"),
|
||||||
SentryDSN: viper.GetString("SENTRY_DSN"),
|
SentryDSN: viper.GetString("SENTRY_DSN"),
|
||||||
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
|
TrustedProxies: splitList(viper.GetString("TRUSTED_PROXIES")),
|
||||||
@@ -119,19 +140,7 @@ func New(
|
|||||||
params: ¶ms,
|
params: ¶ms,
|
||||||
}
|
}
|
||||||
|
|
||||||
// viper reads a value that is not a number as 0, so this also
|
err = s.check()
|
||||||
// catches a mistyped setting.
|
|
||||||
if s.ReportsPerMinute <= 0 {
|
|
||||||
return nil, fmt.Errorf("REPORTS_PER_MINUTE %q: %w",
|
|
||||||
viper.GetString("REPORTS_PER_MINUTE"), errNotPositive)
|
|
||||||
}
|
|
||||||
|
|
||||||
if s.DataDirMaxBytes <= 0 {
|
|
||||||
return nil, fmt.Errorf("DATA_DIR_MAX_BYTES %q: %w",
|
|
||||||
viper.GetString("DATA_DIR_MAX_BYTES"), errNotPositive)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = checkOrigins(s.CORSAllowedOrigins)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -144,6 +153,32 @@ func New(
|
|||||||
return s, nil
|
return s, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// check fails with an error naming the first setting here whose value
|
||||||
|
// the server cannot use. New checks PORT and DEBUG as it reads them,
|
||||||
|
// and the middleware checks TRUSTED_PROXIES as it parses it.
|
||||||
|
func (s *Config) check() error {
|
||||||
|
// viper reads a value that is not a number as 0, so this also
|
||||||
|
// catches a mistyped setting.
|
||||||
|
if s.ReportsPerMinute <= 0 {
|
||||||
|
return fmt.Errorf("REPORTS_PER_MINUTE %q: %w",
|
||||||
|
viper.GetString("REPORTS_PER_MINUTE"), errNotPositive)
|
||||||
|
}
|
||||||
|
|
||||||
|
if s.DataDirMaxBytes <= 0 {
|
||||||
|
return fmt.Errorf("DATA_DIR_MAX_BYTES %q: %w",
|
||||||
|
viper.GetString("DATA_DIR_MAX_BYTES"), errNotPositive)
|
||||||
|
}
|
||||||
|
|
||||||
|
if s.BindAddress != "" {
|
||||||
|
_, err := netip.ParseAddr(s.BindAddress)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("BIND_ADDRESS %q: %w", s.BindAddress, errNotIP)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return checkOrigins(s.CORSAllowedOrigins)
|
||||||
|
}
|
||||||
|
|
||||||
// checkOrigins fails on the first CORS_ALLOWED_ORIGINS entry that is
|
// checkOrigins fails on the first CORS_ALLOWED_ORIGINS entry that is
|
||||||
// not a plain origin, scheme://host with an optional port, as browsers
|
// not a plain origin, scheme://host with an optional port, as browsers
|
||||||
// send it; anything more, such as a trailing "/", would match no page.
|
// send it; anything more, such as a trailing "/", would match no page.
|
||||||
|
|||||||
@@ -29,6 +29,63 @@ func requireConfigError(t *testing.T, setting string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestSettingsLoadAsGiven: valid values pass the checks and are used
|
||||||
|
// as given. bin/entrypoint.sh starts the server with these
|
||||||
|
// BIND_ADDRESS and PORT values.
|
||||||
|
func TestSettingsLoadAsGiven(t *testing.T) {
|
||||||
|
t.Setenv("BIND_ADDRESS", "127.0.0.1")
|
||||||
|
t.Setenv("PORT", "8081")
|
||||||
|
t.Setenv("DEBUG", "true")
|
||||||
|
|
||||||
|
var cfg *config.Config
|
||||||
|
|
||||||
|
app := fx.New(
|
||||||
|
fx.NopLogger,
|
||||||
|
fx.Provide(globals.New, logger.New, config.New),
|
||||||
|
fx.Populate(&cfg),
|
||||||
|
)
|
||||||
|
|
||||||
|
err := app.Err()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("config error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if cfg.BindAddress != "127.0.0.1" || cfg.Port != 8081 || !cfg.Debug {
|
||||||
|
t.Fatalf("BindAddress, Port, Debug = %q, %d, %t; "+
|
||||||
|
"want \"127.0.0.1\", 8081, true",
|
||||||
|
cfg.BindAddress, cfg.Port, cfg.Debug)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestPortMustBeAPortNumber: viper reads a value that is not a number
|
||||||
|
// as 0, on which the server would listen on a random port.
|
||||||
|
func TestPortMustBeAPortNumber(t *testing.T) {
|
||||||
|
for _, value := range []string{"abc", "0", "65536", "8080.5"} {
|
||||||
|
t.Run(value, func(t *testing.T) {
|
||||||
|
t.Setenv("PORT", value)
|
||||||
|
|
||||||
|
requireConfigError(t, "PORT")
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestDebugMustBeTrueOrFalse: viper reads any other value, such as
|
||||||
|
// "yes", as false.
|
||||||
|
func TestDebugMustBeTrueOrFalse(t *testing.T) {
|
||||||
|
t.Setenv("DEBUG", "yes")
|
||||||
|
|
||||||
|
requireConfigError(t, "DEBUG")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestBindAddressMustBeAnIPAddress: a host name would be looked up
|
||||||
|
// only once the server starts listening, and a mistyped one would stop
|
||||||
|
// it then with an error that does not name the setting.
|
||||||
|
func TestBindAddressMustBeAnIPAddress(t *testing.T) {
|
||||||
|
t.Setenv("BIND_ADDRESS", "localhost")
|
||||||
|
|
||||||
|
requireConfigError(t, "BIND_ADDRESS")
|
||||||
|
}
|
||||||
|
|
||||||
// TestReportsPerMinuteMustBePositive: unchecked, zero would panic
|
// TestReportsPerMinuteMustBePositive: unchecked, zero would panic
|
||||||
// when the routes are built, and a negative rate would lift the
|
// when the routes are built, and a negative rate would lift the
|
||||||
// limit.
|
// limit.
|
||||||
|
|||||||
@@ -77,8 +77,8 @@ func New(
|
|||||||
return s, nil
|
return s, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// parseTrustedProxies converts CIDR strings into prefixes,
|
// parseTrustedProxies converts the TRUSTED_PROXIES entries into
|
||||||
// failing fast on any malformed entry.
|
// prefixes, failing fast on any malformed entry.
|
||||||
func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
|
func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
|
||||||
prefixes := make([]netip.Prefix, 0, len(cidrs))
|
prefixes := make([]netip.Prefix, 0, len(cidrs))
|
||||||
|
|
||||||
@@ -86,7 +86,7 @@ func parseTrustedProxies(cidrs []string) ([]netip.Prefix, error) {
|
|||||||
prefix, err := netip.ParsePrefix(cidr)
|
prefix, err := netip.ParsePrefix(cidr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf(
|
return nil, fmt.Errorf(
|
||||||
"trusted proxy %q: %w", cidr, err,
|
"TRUSTED_PROXIES %q: %w", cidr, err,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -40,8 +40,8 @@ func TestParseTrustedProxiesRejectsMalformed(t *testing.T) {
|
|||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
_, err := middleware.ParseTrustedProxies([]string{"not-a-cidr"})
|
_, err := middleware.ParseTrustedProxies([]string{"not-a-cidr"})
|
||||||
if err == nil {
|
if err == nil || !strings.Contains(err.Error(), "TRUSTED_PROXIES") {
|
||||||
t.Fatal("expected error for malformed CIDR, got nil")
|
t.Fatalf("error = %v, want one naming TRUSTED_PROXIES", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,7 +1,15 @@
|
|||||||
package reportbuf
|
package reportbuf
|
||||||
|
|
||||||
|
import "time"
|
||||||
|
|
||||||
// Flush writes the buffered reports to a file now, as the periodic
|
// Flush writes the buffered reports to a file now, as the periodic
|
||||||
// flush does, so tests need not wait a minute for it.
|
// flush does, so tests need not wait a minute for it.
|
||||||
func (b *Buffer) Flush() error {
|
func (b *Buffer) Flush() error {
|
||||||
return b.flushLocked()
|
return b.flushLocked()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// StopClock makes every report file the buffer writes from now on
|
||||||
|
// carry the timestamp at, as if all were written in one millisecond.
|
||||||
|
func (b *Buffer) StopClock(at time.Time) {
|
||||||
|
b.now = func() time.Time { return at }
|
||||||
|
}
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"sneak.berlin/go/netwatch/internal/config"
|
"sneak.berlin/go/netwatch/internal/config"
|
||||||
@@ -30,7 +31,8 @@ const (
|
|||||||
dirPerms fs.FileMode = 0o750
|
dirPerms fs.FileMode = 0o750
|
||||||
filePerms fs.FileMode = 0o640
|
filePerms fs.FileMode = 0o640
|
||||||
|
|
||||||
// Report files are named filePrefix + timestamp + fileSuffix.
|
// Report files are named filePrefix + timestamp + "-" + number +
|
||||||
|
// fileSuffix; see writeFile.
|
||||||
filePrefix = "reports-"
|
filePrefix = "reports-"
|
||||||
fileSuffix = ".jsonl.zst"
|
fileSuffix = ".jsonl.zst"
|
||||||
)
|
)
|
||||||
@@ -56,6 +58,12 @@ type Buffer struct {
|
|||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
maxBytes int64
|
maxBytes int64
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
|
// now is the clock report files are named by: time.Now, except
|
||||||
|
// in tests that need two flushes to share a timestamp.
|
||||||
|
now func() time.Time
|
||||||
|
// seq numbers the report files, so that two named in the same
|
||||||
|
// millisecond still get different names.
|
||||||
|
seq atomic.Uint64
|
||||||
stopOnce sync.Once
|
stopOnce sync.Once
|
||||||
// usedBytes is what Append checks against maxBytes: the size
|
// usedBytes is what Append checks against maxBytes: the size
|
||||||
// of the report files in dataDir, plus the reports not yet
|
// of the report files in dataDir, plus the reports not yet
|
||||||
@@ -79,6 +87,7 @@ func New(
|
|||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
log: params.Logger.Get(),
|
log: params.Logger.Get(),
|
||||||
maxBytes: params.Config.DataDirMaxBytes,
|
maxBytes: params.Config.DataDirMaxBytes,
|
||||||
|
now: time.Now,
|
||||||
}
|
}
|
||||||
|
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
@@ -211,11 +220,14 @@ func (b *Buffer) drainBuf() []byte {
|
|||||||
// writeFile creates a timestamped zstd-compressed JSONL file
|
// writeFile creates a timestamped zstd-compressed JSONL file
|
||||||
// in the data directory.
|
// in the data directory.
|
||||||
func (b *Buffer) writeFile(data []byte) error {
|
func (b *Buffer) writeFile(data []byte) error {
|
||||||
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
|
// The timestamp comes first, so the names sort by time; the number
|
||||||
path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix)
|
// after it tells apart files named in the same millisecond.
|
||||||
|
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z")
|
||||||
|
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
|
||||||
|
path := filepath.Join(b.dataDir, name)
|
||||||
|
|
||||||
// path is built from the operator-supplied dataDir plus a
|
// path is built from the operator-supplied dataDir plus a
|
||||||
// generated timestamp, so it carries no external input.
|
// generated timestamp and number, so it carries no external input.
|
||||||
f, err := os.OpenFile( //nolint:gosec // see comment above
|
f, err := os.OpenFile( //nolint:gosec // see comment above
|
||||||
path,
|
path,
|
||||||
os.O_WRONLY|os.O_CREATE|os.O_EXCL,
|
os.O_WRONLY|os.O_CREATE|os.O_EXCL,
|
||||||
|
|||||||
@@ -3,9 +3,11 @@ package reportbuf_test
|
|||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"io/fs"
|
"io/fs"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"slices"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
@@ -18,6 +20,7 @@ import (
|
|||||||
"sneak.berlin/go/netwatch/internal/logger"
|
"sneak.berlin/go/netwatch/internal/logger"
|
||||||
"sneak.berlin/go/netwatch/internal/reportbuf"
|
"sneak.berlin/go/netwatch/internal/reportbuf"
|
||||||
|
|
||||||
|
"github.com/klauspost/compress/zstd"
|
||||||
"go.uber.org/fx"
|
"go.uber.org/fx"
|
||||||
"go.uber.org/fx/fxtest"
|
"go.uber.org/fx/fxtest"
|
||||||
)
|
)
|
||||||
@@ -210,10 +213,6 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
|
|||||||
t.Fatalf("flush: %v", err)
|
t.Fatalf("flush: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// The second report is written at shutdown, and must not land in
|
|
||||||
// the first file's millisecond (see TestWrittenReportsKeepCounting).
|
|
||||||
time.Sleep(time.Millisecond)
|
|
||||||
|
|
||||||
err = buf.Append(report)
|
err = buf.Append(report)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("second report, after the first was written: %v", err)
|
t.Fatalf("second report, after the first was written: %v", err)
|
||||||
@@ -254,10 +253,6 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
|
|||||||
t.Fatalf("with %d bytes of report files: %v", used, err)
|
t.Fatalf("with %d bytes of report files: %v", used, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Report files are named to the millisecond; two in the same
|
|
||||||
// one collide (https://git.eeqj.de/sneak/netwatch/issues/61).
|
|
||||||
time.Sleep(time.Millisecond)
|
|
||||||
|
|
||||||
err = buf.Flush()
|
err = buf.Flush()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("flush: %v", err)
|
t.Fatalf("flush: %v", err)
|
||||||
@@ -315,6 +310,43 @@ func TestConcurrentAppendsStopAtCap(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestTwoFlushesInOneMillisecond flushes twice within one millisecond,
|
||||||
|
// as a flush for size and the final flush at shutdown can: each flush
|
||||||
|
// must write a file of its own, and the files must hold every report.
|
||||||
|
func TestTwoFlushesInOneMillisecond(t *testing.T) {
|
||||||
|
const flushes = 2
|
||||||
|
|
||||||
|
dir := t.TempDir()
|
||||||
|
t.Setenv("DATA_DIR", dir)
|
||||||
|
|
||||||
|
buf := startBuffer(t)
|
||||||
|
buf.StopClock(time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC))
|
||||||
|
|
||||||
|
for id := 1; id <= flushes; id++ {
|
||||||
|
err := buf.Append(map[string]int{"id": id})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("append report %d: %v", id, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = buf.Flush()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("flush %d: %v", id, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
files := readReportFiles(t, dir)
|
||||||
|
if len(files) != flushes {
|
||||||
|
t.Fatalf("%d report files after %d flushes", len(files), flushes)
|
||||||
|
}
|
||||||
|
|
||||||
|
for id := 1; id <= flushes; id++ {
|
||||||
|
want := fmt.Sprintf(`{"id":%d}`+"\n", id)
|
||||||
|
if !slices.Contains(files, want) {
|
||||||
|
t.Fatalf("no report file holds report %d alone", id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// reportFilesBytes returns the total size of the report files in dir.
|
// reportFilesBytes returns the total size of the report files in dir.
|
||||||
func reportFilesBytes(t *testing.T, dir string) int64 {
|
func reportFilesBytes(t *testing.T, dir string) int64 {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
@@ -338,6 +370,43 @@ func reportFilesBytes(t *testing.T, dir string) int64 {
|
|||||||
return total
|
return total
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// readReportFiles returns the decompressed contents of each report
|
||||||
|
// file in dir.
|
||||||
|
func readReportFiles(t *testing.T, dir string) []string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
files := os.DirFS(dir)
|
||||||
|
|
||||||
|
names, err := fs.Glob(files, "reports-*.jsonl.zst")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list report files: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
dec, err := zstd.NewReader(nil)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("create zstd decoder: %v", err)
|
||||||
|
}
|
||||||
|
defer dec.Close()
|
||||||
|
|
||||||
|
contents := make([]string, 0, len(names))
|
||||||
|
|
||||||
|
for _, name := range names {
|
||||||
|
compressed, readErr := fs.ReadFile(files, name)
|
||||||
|
if readErr != nil {
|
||||||
|
t.Fatalf("read %s: %v", name, readErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
data, decErr := dec.DecodeAll(compressed, nil)
|
||||||
|
if decErr != nil {
|
||||||
|
t.Fatalf("decompress %s: %v", name, decErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
contents = append(contents, string(data))
|
||||||
|
}
|
||||||
|
|
||||||
|
return contents
|
||||||
|
}
|
||||||
|
|
||||||
func writeBytes(t *testing.T, path string, n int) {
|
func writeBytes(t *testing.T, path string, n int) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
|||||||
+14
-2
@@ -10,8 +10,9 @@ set -u
|
|||||||
|
|
||||||
# PORT is the public port nginx listens on, 8080 when unset or empty.
|
# PORT is the public port nginx listens on, 8080 when unset or empty.
|
||||||
# nginx would take a value such as localhost or unix:/tmp/x.sock as an
|
# nginx would take a value such as localhost or unix:/tmp/x.sock as an
|
||||||
# address and start anyway, so anything but digits stops the container
|
# address and start anyway, and reports a bad port without naming
|
||||||
# here, before either process starts.
|
# PORT, so a value that is not a usable port stops the container here,
|
||||||
|
# before either process starts.
|
||||||
export PORT="${PORT:-8080}"
|
export PORT="${PORT:-8080}"
|
||||||
case "$PORT" in
|
case "$PORT" in
|
||||||
*[!0-9]*)
|
*[!0-9]*)
|
||||||
@@ -19,6 +20,17 @@ case "$PORT" in
|
|||||||
exit 1
|
exit 1
|
||||||
;;
|
;;
|
||||||
esac
|
esac
|
||||||
|
# The length is checked first because, for a number too big for it,
|
||||||
|
# the shell's test prints an error and is false, so the range checks
|
||||||
|
# alone would let it through.
|
||||||
|
if [ "${#PORT}" -gt 5 ] || [ "$PORT" -lt 1 ] || [ "$PORT" -gt 65535 ]; then
|
||||||
|
echo "entrypoint: PORT must be from 1 to 65535, not '$PORT'" >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
if [ "$PORT" -eq 8081 ]; then
|
||||||
|
echo "entrypoint: PORT cannot be 8081, netwatch-server listens there" >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
# A stop signal is only noted here; the loop below acts on it.
|
# A stop signal is only noted here; the loop below acts on it.
|
||||||
stop_requested=""
|
stop_requested=""
|
||||||
|
|||||||
Reference in New Issue
Block a user