Compare commits
2
Commits
c88c5da063
...
0987817aae
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0987817aae | ||
|
|
ced1956b06 |
+4
-1
@@ -68,8 +68,10 @@ FROM nginx@sha256:15e96e59aa3b0aada3a121296e3bce117721f42d88f5f64217ef4b18f458c6
|
|||||||
RUN addgroup -g 1000 -S netwatch && \
|
RUN addgroup -g 1000 -S netwatch && \
|
||||||
adduser -u 1000 -S netwatch -G netwatch
|
adduser -u 1000 -S netwatch -G netwatch
|
||||||
|
|
||||||
|
# At start-up the nginx image renders every template here into
|
||||||
|
# conf.d; bin/entrypoint.sh says how.
|
||||||
RUN rm /etc/nginx/conf.d/default.conf
|
RUN rm /etc/nginx/conf.d/default.conf
|
||||||
COPY nginx.conf /etc/nginx/conf.d/netwatch.conf
|
COPY nginx.conf /etc/nginx/templates/netwatch.conf.template
|
||||||
COPY --from=frontend /app/dist /usr/share/nginx/html
|
COPY --from=frontend /app/dist /usr/share/nginx/html
|
||||||
COPY --from=builder /src/netwatch-server /usr/local/bin/netwatch-server
|
COPY --from=builder /src/netwatch-server /usr/local/bin/netwatch-server
|
||||||
COPY bin/entrypoint.sh /usr/local/bin/entrypoint.sh
|
COPY bin/entrypoint.sh /usr/local/bin/entrypoint.sh
|
||||||
@@ -78,6 +80,7 @@ ENV DATA_DIR=/data/reports
|
|||||||
RUN mkdir -p /data/reports && chown -R netwatch:netwatch /data
|
RUN mkdir -p /data/reports && chown -R netwatch:netwatch /data
|
||||||
VOLUME /data
|
VOLUME /data
|
||||||
|
|
||||||
|
# The default public port; PORT changes it.
|
||||||
EXPOSE 8080
|
EXPOSE 8080
|
||||||
|
|
||||||
# The nginx image stops its container with SIGQUIT; the entrypoint
|
# The nginx image stops its container with SIGQUIT; the entrypoint
|
||||||
|
|||||||
@@ -23,6 +23,17 @@ 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 counts the files
|
||||||
|
written since the server started, 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
|
||||||
|
- 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
|
||||||
|
`PORT` and no other variable. `bin/entrypoint.sh` refuses to start when `PORT`
|
||||||
|
is not digits only. `server_tokens off` keeps the nginx version out of
|
||||||
|
responses. `script/frontend-viewport-test` renders the template the same way.
|
||||||
|
Gzip and a `50x.html` error page are not added
|
||||||
- 2026-09-29: bounded the report endpoint (issue #20): `POST /api/v1/reports`
|
- 2026-09-29: bounded the report endpoint (issue #20): `POST /api/v1/reports`
|
||||||
still needs no credentials, but each client address, as resolved through
|
still needs no credentials, but each client address, as resolved through
|
||||||
`TRUSTED_PROXIES`, may send `REPORTS_PER_MINUTE` (default 60) reports a
|
`TRUSTED_PROXIES`, may send `REPORTS_PER_MINUTE` (default 60) reports a
|
||||||
|
|||||||
+6
-3
@@ -102,9 +102,12 @@ 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 counts the files the server has written since it started, so
|
||||||
|
two files written in the same millisecond still get different names. Each file
|
||||||
|
contains one JSON object per line, compressed with zstd. Files are created with
|
||||||
|
`O_EXCL` to prevent overwrites.
|
||||||
|
|
||||||
### Report limits
|
### Report limits
|
||||||
|
|
||||||
|
|||||||
@@ -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,9 @@ type Buffer struct {
|
|||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
maxBytes int64
|
maxBytes int64
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
|
// 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
|
||||||
@@ -211,11 +216,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 {
|
||||||
|
// The timestamp comes first, so the names sort by time; the number
|
||||||
|
// after it tells apart files named in the same millisecond.
|
||||||
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
|
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
|
||||||
path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix)
|
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,60 @@ 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) {
|
||||||
|
// A pair of flushes may straddle a millisecond, which proves
|
||||||
|
// nothing, so pairs are flushed until one falls within one.
|
||||||
|
const maxPairs = 1000
|
||||||
|
|
||||||
|
dir := t.TempDir()
|
||||||
|
t.Setenv("DATA_DIR", dir)
|
||||||
|
|
||||||
|
buf := startBuffer(t)
|
||||||
|
flushes := 0
|
||||||
|
|
||||||
|
for pair := 1; ; pair++ {
|
||||||
|
start := time.Now().Truncate(time.Millisecond)
|
||||||
|
|
||||||
|
for range 2 {
|
||||||
|
flushes++
|
||||||
|
|
||||||
|
err := buf.Append(map[string]int{"id": flushes})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("append report %d: %v", flushes, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = buf.Flush()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("flush %d: %v", flushes, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if time.Now().Truncate(time.Millisecond).Equal(start) {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
|
if pair == maxPairs {
|
||||||
|
t.Fatalf("no pair of flushes fell within one millisecond "+
|
||||||
|
"in %d tries", maxPairs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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 +387,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()
|
||||||
|
|
||||||
|
|||||||
+19
-2
@@ -8,6 +8,18 @@
|
|||||||
# No set -e: kill and wait return non-zero here in normal operation.
|
# No set -e: kill and wait return non-zero here in normal operation.
|
||||||
set -u
|
set -u
|
||||||
|
|
||||||
|
# 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
|
||||||
|
# address and start anyway, so anything but digits stops the container
|
||||||
|
# here, before either process starts.
|
||||||
|
export PORT="${PORT:-8080}"
|
||||||
|
case "$PORT" in
|
||||||
|
*[!0-9]*)
|
||||||
|
echo "entrypoint: PORT must be a port number, not '$PORT'" >&2
|
||||||
|
exit 1
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
|
||||||
# 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=""
|
||||||
trap 'stop_requested=yes' TERM INT
|
trap 'stop_requested=yes' TERM INT
|
||||||
@@ -23,8 +35,13 @@ backend=$!
|
|||||||
|
|
||||||
# nginx starts through the nginx image's own entrypoint, which applies
|
# nginx starts through the nginx image's own entrypoint, which applies
|
||||||
# the image's start-up configuration and then replaces itself with
|
# the image's start-up configuration and then replaces itself with
|
||||||
# nginx.
|
# nginx. Part of that start-up configuration renders nginx.conf into
|
||||||
/docker-entrypoint.sh nginx -g 'daemon off;' &
|
# conf.d with nginx listening on PORT. NGINX_ENVSUBST_FILTER limits
|
||||||
|
# that rendering to PORT: a variable nginx itself uses, such as $uri,
|
||||||
|
# would otherwise be replaced by an environment variable of the same
|
||||||
|
# name.
|
||||||
|
NGINX_ENVSUBST_FILTER='^PORT$' \
|
||||||
|
/docker-entrypoint.sh nginx -g 'daemon off;' &
|
||||||
nginx=$!
|
nginx=$!
|
||||||
|
|
||||||
running() {
|
running() {
|
||||||
|
|||||||
+7
-1
@@ -1,7 +1,13 @@
|
|||||||
|
# A template: the nginx image renders it into conf.d at container start,
|
||||||
|
# filling in PORT and nothing else. bin/entrypoint.sh sets PORT and that
|
||||||
|
# limit.
|
||||||
server {
|
server {
|
||||||
listen 8080;
|
listen ${PORT};
|
||||||
server_name _;
|
server_name _;
|
||||||
|
|
||||||
|
# Keep the nginx version out of the Server header and error pages.
|
||||||
|
server_tokens off;
|
||||||
|
|
||||||
root /usr/share/nginx/html;
|
root /usr/share/nginx/html;
|
||||||
index index.html;
|
index index.html;
|
||||||
|
|
||||||
|
|||||||
@@ -61,10 +61,13 @@ main() {
|
|||||||
# host.
|
# host.
|
||||||
docker network create --internal "$NETWORK" > /dev/null
|
docker network create --internal "$NETWORK" > /dev/null
|
||||||
|
|
||||||
|
# nginx.conf is a template: the image renders it over its own
|
||||||
|
# default.conf, with the same port and limit bin/entrypoint.sh uses.
|
||||||
docker run -d --rm --name "$SERVER" \
|
docker run -d --rm --name "$SERVER" \
|
||||||
--network "$NETWORK" --network-alias netwatch \
|
--network "$NETWORK" --network-alias netwatch \
|
||||||
|
-e PORT=8080 -e NGINX_ENVSUBST_FILTER='^PORT$' \
|
||||||
-v "$ROOT/dist:/usr/share/nginx/html:ro" \
|
-v "$ROOT/dist:/usr/share/nginx/html:ro" \
|
||||||
-v "$ROOT/nginx.conf:/etc/nginx/conf.d/default.conf:ro" \
|
-v "$ROOT/nginx.conf:/etc/nginx/templates/default.conf.template:ro" \
|
||||||
"$SERVER_IMAGE" > /dev/null
|
"$SERVER_IMAGE" > /dev/null
|
||||||
|
|
||||||
# The image's own entrypoint already exposes CDP on 9222 and passes
|
# The image's own entrypoint already exposes CDP on 9222 and passes
|
||||||
|
|||||||
Reference in New Issue
Block a user