Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ee959fc33a |
+2
-3
@@ -65,9 +65,8 @@ RUN yarn install --frozen-lockfile
|
|||||||
RUN apk add --no-cache git make
|
RUN apk add --no-cache git make
|
||||||
COPY . .
|
COPY . .
|
||||||
# make frontend-check is the frontend half of make check (test + lint +
|
# make frontend-check is the frontend half of make check (test + lint +
|
||||||
# fmt-check); its test step runs the unit tests, then the production
|
# fmt-check); its test step is the production yarn build, so this both
|
||||||
# yarn build, so this both produces dist/ and gates the image on
|
# produces dist/ and gates the image on lint/fmt-check/test regressions.
|
||||||
# lint/fmt-check/test regressions.
|
|
||||||
# This node stage has neither Go nor Docker; the lint and builder stages
|
# This node stage has neither Go nor Docker; the lint and builder stages
|
||||||
# above gate the backend half.
|
# above gate the backend half.
|
||||||
RUN make frontend-check
|
RUN make frontend-check
|
||||||
|
|||||||
@@ -52,8 +52,8 @@ halves, so the root `make check` fails if either one is broken. We provide:
|
|||||||
- `script/fmt` — format all files (writes): prettier, then gofmt over `backend/`
|
- `script/fmt` — format all files (writes): prettier, then gofmt over `backend/`
|
||||||
- `script/fmt-check` — check formatting (read-only): prettier, then gofmt
|
- `script/fmt-check` — check formatting (read-only): prettier, then gofmt
|
||||||
- `script/check` — run test, lint, and fmt-check
|
- `script/check` — run test, lint, and fmt-check
|
||||||
- `script/frontend-test` — run the unit tests in `test/unit/` with Node's
|
- `script/frontend-test` — run the production build as the frontend's test (no
|
||||||
built-in test runner, then the production build
|
unit tests yet)
|
||||||
- `script/frontend-lint` — run prettier in check mode
|
- `script/frontend-lint` — run prettier in check mode
|
||||||
- `script/frontend-fmt` — format everything prettier understands (writes)
|
- `script/frontend-fmt` — format everything prettier understands (writes)
|
||||||
- `script/frontend-fmt-check` — check prettier formatting (read-only)
|
- `script/frontend-fmt-check` — check prettier formatting (read-only)
|
||||||
@@ -136,14 +136,8 @@ Local hosts are tracked separately from WAN stats.
|
|||||||
### Latency measurement
|
### Latency measurement
|
||||||
|
|
||||||
HEAD requests with `mode: 'no-cors'` and `cache: 'no-store'`, timed with
|
HEAD requests with `mode: 'no-cors'` and `cache: 'no-store'`, timed with
|
||||||
`performance.now()`. Each check times out after 80% of the refresh interval (24
|
`performance.now()`. 1-second timeout; anything over 1000ms is clamped to
|
||||||
seconds at 30 seconds) and is then recorded as a timeout, so a round's checks
|
unreachable. IPv4 only.
|
||||||
have all finished before the next round is due. When no WAN host answers, a
|
|
||||||
recovery probe checks 4 random WAN hosts every half second, giving up the checks
|
|
||||||
it started half a second before. As soon as one answers, a new round starts at
|
|
||||||
once, as it does after an interval change. A round started early gives up the
|
|
||||||
last round's checks if they are still waiting, and that round records nothing,
|
|
||||||
so rounds never overlap. IPv4 only.
|
|
||||||
|
|
||||||
### Color coding
|
### Color coding
|
||||||
|
|
||||||
|
|||||||
@@ -27,19 +27,10 @@ latest run passes.
|
|||||||
(issue #54): when a report would take them past it, the oldest report files
|
(issue #54): when a report would take them past it, the oldest report files
|
||||||
are deleted to make room, each deletion logged, and at start files already
|
are deleted to make room, each deletion logged, and at start files already
|
||||||
past it are deleted the same way. A file still being written is never deleted.
|
past it are deleted the same way. A file still being written is never deleted.
|
||||||
A report is refused with 507 only when the reports waiting to be written fill
|
A report is refused with 507 only when the reports not yet written leave no
|
||||||
the cap on their own, and then no file is deleted. The reports of a failed
|
room for it on their own, and then no file is deleted. A file that cannot be
|
||||||
write stop counting, and the part of its file written is removed. A file that
|
deleted still counts until the next start; one already deleted by hand counts
|
||||||
cannot be deleted still counts until the next start; one already deleted by
|
as freed
|
||||||
hand counts as freed
|
|
||||||
- 2026-10-03: each target check times out after 80% of the refresh interval
|
|
||||||
(issue #78), 24 seconds at 30 seconds, where it was capped at 3 seconds. A
|
|
||||||
round started early, after an interval change or when the recovery probe finds
|
|
||||||
a target answering, gives up the last round's checks if they are still
|
|
||||||
waiting, so rounds never overlap; the recovery probe gives up its own checks
|
|
||||||
after half a second. The frontend has its first unit tests, run by
|
|
||||||
`script/frontend-test` with Node's built-in test runner; for them,
|
|
||||||
`index.html` now links `src/styles.css`, which `src/main.js` used to import
|
|
||||||
- 2026-09-29: the container sets up its own data directory (issue #75):
|
- 2026-09-29: the container sets up its own data directory (issue #75):
|
||||||
`bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it
|
`bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it
|
||||||
and `/data` to the `netwatch` user with mode 750 before starting the backend
|
and `/data` to the `netwatch` user with mode 750 before starting the backend
|
||||||
|
|||||||
+9
-11
@@ -148,17 +148,15 @@ credentials, so it is bounded instead. Both refusals below answer with the same
|
|||||||
`X-RateLimit-Reset` headers.
|
`X-RateLimit-Reset` headers.
|
||||||
- **Size cap.** The report files in `DATA_DIR` may total at most
|
- **Size cap.** The report files in `DATA_DIR` may total at most
|
||||||
`DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports
|
`DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports
|
||||||
waiting in memory count at their uncompressed size until they are written;
|
waiting in memory count at their uncompressed size until they are written.
|
||||||
those lost to a failed write stop counting, and the part of its file written
|
When a report would take the total past the cap, the oldest report files are
|
||||||
is removed. When a report would take the total past the cap, the oldest report
|
deleted to make room, and each deletion is logged with the file's name and
|
||||||
files are deleted to make room, and each deletion is logged with the file's
|
size; a file still being written is never deleted. A report is refused with
|
||||||
name and size; a file still being written is never deleted. A report is
|
507, and nothing of it is stored, only when the reports not yet written leave
|
||||||
refused with 507, and nothing of it is stored, only when the reports waiting
|
no room for it on their own, and then no file is deleted. At start, report
|
||||||
to be written fill the cap on their own, and then no file is deleted. At
|
files past the cap, as after lowering it, are deleted the same way. So the cap
|
||||||
start, report files past the cap, as after lowering it, are deleted the same
|
is how much of the newest reports is kept: the default of 1 GiB is small
|
||||||
way. So the cap is how much of the newest reports is kept: the default of 1
|
enough for any host; set it to the space you can give `DATA_DIR`.
|
||||||
GiB is small enough for any host; set it to the space you can give
|
|
||||||
`DATA_DIR`.
|
|
||||||
|
|
||||||
### CORS
|
### CORS
|
||||||
|
|
||||||
|
|||||||
@@ -86,12 +86,11 @@ func (s *Handlers) decodeErrorStatus(err error) int {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// appendErrorStatus logs a failure to store a report and returns
|
// appendErrorStatus logs a failure to store a report and returns
|
||||||
// the status to send: 507 when the reports waiting to be written fill
|
// the status to send: 507 when the report files are at their size
|
||||||
// the size cap, otherwise 500.
|
// cap, otherwise 500.
|
||||||
func (s *Handlers) appendErrorStatus(err error) int {
|
func (s *Handlers) appendErrorStatus(err error) int {
|
||||||
if errors.Is(err, reportbuf.ErrFull) {
|
if errors.Is(err, reportbuf.ErrFull) {
|
||||||
s.log.Warn("report refused: " +
|
s.log.Warn("report refused: report files at their size cap")
|
||||||
"reports waiting to be written fill the size cap")
|
|
||||||
|
|
||||||
return http.StatusInsufficientStorage
|
return http.StatusInsufficientStorage
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -68,9 +68,9 @@ func TestHandleReportStorageFailureIsNon2xx(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestHandleReportFullIs507 checks the answer when the reports waiting
|
// TestHandleReportFullIs507 checks the answer when the report files
|
||||||
// to be written fill the size cap: 507 and the usual error body, which
|
// are at their size cap: 507 and the usual error body, which tells
|
||||||
// tells the client nothing more.
|
// the client nothing more.
|
||||||
func TestHandleReportFullIs507(t *testing.T) {
|
func TestHandleReportFullIs507(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -1,9 +1,6 @@
|
|||||||
package reportbuf
|
package reportbuf
|
||||||
|
|
||||||
import (
|
import "time"
|
||||||
"os"
|
|
||||||
"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.
|
||||||
@@ -16,10 +13,3 @@ func (b *Buffer) Flush() error {
|
|||||||
func (b *Buffer) StopClock(at time.Time) {
|
func (b *Buffer) StopClock(at time.Time) {
|
||||||
b.now = func() time.Time { return at }
|
b.now = func() time.Time { return at }
|
||||||
}
|
}
|
||||||
|
|
||||||
// OnFileCreated makes the buffer call fn with each report file it
|
|
||||||
// writes from now on, once the file is created and before anything is
|
|
||||||
// written to it.
|
|
||||||
func (b *Buffer) OnFileCreated(fn func(f *os.File)) {
|
|
||||||
b.fileCreated = fn
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -38,10 +38,10 @@ const (
|
|||||||
fileSuffix = ".jsonl.zst"
|
fileSuffix = ".jsonl.zst"
|
||||||
)
|
)
|
||||||
|
|
||||||
// ErrFull is returned by Append when the reports waiting to be
|
// ErrFull is returned by Append when the reports not yet written
|
||||||
// written leave no room for the report under the configured maximum
|
// leave no room for the report under the configured maximum size,
|
||||||
// size, however many report files are deleted.
|
// however many report files are deleted.
|
||||||
var ErrFull = errors.New("reports waiting to be written fill the size cap")
|
var ErrFull = errors.New("report files at their size cap")
|
||||||
|
|
||||||
// Params defines the dependencies for Buffer.
|
// Params defines the dependencies for Buffer.
|
||||||
type Params struct {
|
type Params struct {
|
||||||
@@ -64,10 +64,6 @@ type Buffer struct {
|
|||||||
buf bytes.Buffer
|
buf bytes.Buffer
|
||||||
dataDir string
|
dataDir string
|
||||||
done chan struct{}
|
done chan struct{}
|
||||||
// fileCreated is called with each report file once it is
|
|
||||||
// created, before anything is written to it: it does nothing,
|
|
||||||
// except in tests that hold the write open or make it fail.
|
|
||||||
fileCreated func(f *os.File)
|
|
||||||
// files are the report files that may be deleted to make room,
|
// files are the report files that may be deleted to make room,
|
||||||
// oldest first: those in dataDir at start, then each one this
|
// oldest first: those in dataDir at start, then each one this
|
||||||
// buffer writes, once it is complete. A file still being written
|
// buffer writes, once it is complete. A file still being written
|
||||||
@@ -85,8 +81,8 @@ type Buffer struct {
|
|||||||
seq atomic.Uint64
|
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 waiting to
|
// of the report files in dataDir, plus the reports not yet
|
||||||
// be written to one at their uncompressed size.
|
// written to one at their uncompressed size.
|
||||||
usedBytes int64
|
usedBytes int64
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -102,12 +98,11 @@ func New(
|
|||||||
}
|
}
|
||||||
|
|
||||||
b := &Buffer{
|
b := &Buffer{
|
||||||
dataDir: dir,
|
dataDir: dir,
|
||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
fileCreated: func(*os.File) {},
|
log: params.Logger.Get(),
|
||||||
log: params.Logger.Get(),
|
maxBytes: params.Config.DataDirMaxBytes,
|
||||||
maxBytes: params.Config.DataDirMaxBytes,
|
now: time.Now,
|
||||||
now: time.Now,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
lc.Append(fx.Hook{
|
lc.Append(fx.Hook{
|
||||||
@@ -211,9 +206,9 @@ func (b *Buffer) Append(v any) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// deleteOldestFiles deletes report files, oldest first, until n more
|
// deleteOldestFiles deletes report files, oldest first, until n more
|
||||||
// bytes fit under maxBytes. It deletes none when the reports waiting
|
// bytes fit under maxBytes. It deletes none when the reports not yet
|
||||||
// to be written leave no room for n even with every file gone, since
|
// written leave no room for n even with every file gone, since that
|
||||||
// that would lose the files for nothing. The caller must hold b.mu.
|
// would lose the files for nothing. The caller must hold b.mu.
|
||||||
func (b *Buffer) deleteOldestFiles(n int64) {
|
func (b *Buffer) deleteOldestFiles(n int64) {
|
||||||
for b.usedBytes+n > b.maxBytes && len(b.files) > 0 {
|
for b.usedBytes+n > b.maxBytes && len(b.files) > 0 {
|
||||||
if b.usedBytes-b.filesBytes+n > b.maxBytes {
|
if b.usedBytes-b.filesBytes+n > b.maxBytes {
|
||||||
@@ -287,41 +282,15 @@ func (b *Buffer) drainBuf() []byte {
|
|||||||
return data
|
return data
|
||||||
}
|
}
|
||||||
|
|
||||||
// writeFile writes data, reports drained from the buffer, to a new
|
// writeFile creates a timestamped zstd-compressed JSONL file
|
||||||
// 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
|
// The timestamp comes first, so the names sort by time; the number
|
||||||
// after it tells apart files named in the same millisecond.
|
// after it tells apart files named in the same millisecond.
|
||||||
ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z")
|
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)
|
name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix)
|
||||||
|
path := filepath.Join(b.dataDir, name)
|
||||||
|
|
||||||
size, err := b.createFile(filepath.Join(b.dataDir, name), data)
|
|
||||||
|
|
||||||
// The reports no longer wait to be written, so they stop counting
|
|
||||||
// at their uncompressed size. If the write failed they are lost;
|
|
||||||
// otherwise they count as the file, which from here on may be
|
|
||||||
// deleted to make room.
|
|
||||||
b.mu.Lock()
|
|
||||||
defer b.mu.Unlock()
|
|
||||||
|
|
||||||
b.usedBytes -= int64(len(data))
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
b.usedBytes += size
|
|
||||||
b.files = append(b.files, reportFile{name: name, size: size})
|
|
||||||
b.filesBytes += size
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// createFile creates the file at path holding data compressed with
|
|
||||||
// zstd, and returns its size. If the write fails once the file is
|
|
||||||
// created, it removes the file, so that a failed write leaves nothing
|
|
||||||
// behind to take room.
|
|
||||||
func (b *Buffer) createFile(path string, data []byte) (int64, error) {
|
|
||||||
// path is built from the operator-supplied dataDir plus a
|
// path is built from the operator-supplied dataDir plus a
|
||||||
// generated timestamp and number, 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
|
||||||
@@ -330,54 +299,53 @@ func (b *Buffer) createFile(path string, data []byte) (int64, error) {
|
|||||||
filePerms,
|
filePerms,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, fmt.Errorf("create report file: %w", err)
|
return fmt.Errorf("create report file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
b.fileCreated(f)
|
// Closes the file on the early returns below. The success
|
||||||
|
// path closes it explicitly to check the error; closing it
|
||||||
|
// a second time here is harmless.
|
||||||
|
defer func() { _ = f.Close() }()
|
||||||
|
|
||||||
size, err := writeCompressed(f, data)
|
|
||||||
if err != nil {
|
|
||||||
_ = f.Close()
|
|
||||||
|
|
||||||
return 0, errors.Join(err, os.Remove(path))
|
|
||||||
}
|
|
||||||
|
|
||||||
err = f.Close()
|
|
||||||
if err != nil {
|
|
||||||
err = fmt.Errorf("close report file: %w", err)
|
|
||||||
|
|
||||||
return 0, errors.Join(err, os.Remove(path))
|
|
||||||
}
|
|
||||||
|
|
||||||
return size, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// writeCompressed writes data to f compressed with zstd, and returns
|
|
||||||
// the size of f.
|
|
||||||
func writeCompressed(f *os.File, data []byte) (int64, error) {
|
|
||||||
enc, err := zstd.NewWriter(f)
|
enc, err := zstd.NewWriter(f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, fmt.Errorf("create zstd encoder: %w", err)
|
return fmt.Errorf("create zstd encoder: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
_, err = enc.Write(data)
|
_, err = enc.Write(data)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = enc.Close()
|
_ = enc.Close()
|
||||||
|
|
||||||
return 0, fmt.Errorf("write compressed data: %w", err)
|
return fmt.Errorf("write compressed data: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = enc.Close()
|
err = enc.Close()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, fmt.Errorf("close zstd encoder: %w", err)
|
return fmt.Errorf("close zstd encoder: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
info, err := f.Stat()
|
info, err := f.Stat()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, fmt.Errorf("stat report file: %w", err)
|
return fmt.Errorf("stat report file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return info.Size(), nil
|
err = f.Close()
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("close report file: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The reports counted at their uncompressed size while they
|
||||||
|
// waited; now they count as the file, which from here on may be
|
||||||
|
// deleted to make room. After a failed write they stay counted as
|
||||||
|
// they were, which errs toward refusing reports early rather than
|
||||||
|
// letting the files pass the cap.
|
||||||
|
b.mu.Lock()
|
||||||
|
b.usedBytes += info.Size() - int64(len(data))
|
||||||
|
b.files = append(b.files, reportFile{name: name, size: info.Size()})
|
||||||
|
b.filesBytes += info.Size()
|
||||||
|
b.mu.Unlock()
|
||||||
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// reportFiles returns the report files in dir, oldest first:
|
// reportFiles returns the report files in dir, oldest first:
|
||||||
|
|||||||
@@ -424,111 +424,36 @@ func TestFileDeletedByHandFreesRoom(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestFileBeingWrittenIsNeverDeleted holds a write open, its file
|
// TestUncountedReportFileIsNeverDeleted: only the report files counted,
|
||||||
// created but not complete, while a report needs room. The file may
|
// those in DATA_DIR at start and those the buffer has finished
|
||||||
// not be deleted to make it, so the report is refused. Once the write
|
// writing, are deleted to make room. A file still being written is not
|
||||||
// is complete, the file is deleted when room is needed.
|
// among them, so it is kept even when that means refusing a report. A
|
||||||
func TestFileBeingWrittenIsNeverDeleted(t *testing.T) {
|
// write cannot be paused here, so a report file put in DATA_DIR after
|
||||||
report := map[string]string{"id": "writing"}
|
// start, which is not counted either, stands in for one being written.
|
||||||
|
func TestUncountedReportFileIsNeverDeleted(t *testing.T) {
|
||||||
|
report := map[string]string{"id": "uncounted"}
|
||||||
|
dir := t.TempDir()
|
||||||
|
uncounted := reportFilePath(dir, 1)
|
||||||
|
|
||||||
t.Setenv("DATA_DIR", t.TempDir())
|
t.Setenv("DATA_DIR", dir)
|
||||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
||||||
|
|
||||||
buf := startBuffer(t)
|
buf := startBuffer(t)
|
||||||
|
|
||||||
created := make(chan string)
|
|
||||||
release := make(chan struct{})
|
|
||||||
|
|
||||||
buf.OnFileCreated(func(f *os.File) {
|
|
||||||
created <- f.Name()
|
|
||||||
|
|
||||||
<-release
|
|
||||||
})
|
|
||||||
|
|
||||||
err := buf.Append(report)
|
err := buf.Append(report)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("report that fills the cap exactly: %v", err)
|
t.Fatalf("report that fills the cap exactly: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
flushed := make(chan error)
|
writeBytes(t, uncounted, 100)
|
||||||
|
|
||||||
go func() { flushed <- buf.Flush() }()
|
|
||||||
|
|
||||||
writing := <-created
|
|
||||||
|
|
||||||
// Errorf, not Fatalf, until the write is released, so that a
|
|
||||||
// failure here does not leave it held.
|
|
||||||
err = buf.Append(report)
|
err = buf.Append(report)
|
||||||
if !errors.Is(err, reportbuf.ErrFull) {
|
if !errors.Is(err, reportbuf.ErrFull) {
|
||||||
t.Errorf("report while the first was being written: "+
|
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
|
||||||
"error = %v, want ErrFull", err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if !exists(t, writing) {
|
if !exists(t, uncounted) {
|
||||||
t.Error("report file deleted while it was being written")
|
t.Fatal("report file deleted, though the buffer had not counted it")
|
||||||
}
|
|
||||||
|
|
||||||
close(release)
|
|
||||||
|
|
||||||
err = <-flushed
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("flush: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Writes from here on, the final one at stop included, go
|
|
||||||
// through unheld.
|
|
||||||
buf.OnFileCreated(func(*os.File) {})
|
|
||||||
|
|
||||||
err = buf.Append(report)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("report after the write was complete: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if exists(t, writing) {
|
|
||||||
t.Fatal("complete report file kept when room was needed")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFailedWriteStopsCounting makes a write fail once its file is
|
|
||||||
// created. Its reports are lost, so they stop counting, and the part of
|
|
||||||
// the file written is removed, so it takes no room.
|
|
||||||
func TestFailedWriteStopsCounting(t *testing.T) {
|
|
||||||
report := map[string]string{"id": "failed"}
|
|
||||||
|
|
||||||
t.Setenv("DATA_DIR", t.TempDir())
|
|
||||||
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
|
||||||
|
|
||||||
buf := startBuffer(t)
|
|
||||||
|
|
||||||
var failed string
|
|
||||||
|
|
||||||
// Closing the file under the write makes the write fail.
|
|
||||||
buf.OnFileCreated(func(f *os.File) {
|
|
||||||
failed = f.Name()
|
|
||||||
|
|
||||||
_ = f.Close()
|
|
||||||
})
|
|
||||||
|
|
||||||
err := buf.Append(report)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("report that fills the cap exactly: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = buf.Flush()
|
|
||||||
if err == nil {
|
|
||||||
t.Fatal("flush succeeded, though its file was closed under it")
|
|
||||||
}
|
|
||||||
|
|
||||||
if exists(t, failed) {
|
|
||||||
t.Fatal("file of the failed write kept")
|
|
||||||
}
|
|
||||||
|
|
||||||
// Writes from here on, the final one at stop included, succeed.
|
|
||||||
buf.OnFileCreated(func(*os.File) {})
|
|
||||||
|
|
||||||
err = buf.Append(report)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("report after the failed write: %v", err)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -9,9 +9,6 @@
|
|||||||
type="image/svg+xml"
|
type="image/svg+xml"
|
||||||
href="data:image/svg+xml,<svg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 100 100'><text y='.9em' font-size='90'>📡</text></svg>"
|
href="data:image/svg+xml,<svg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 100 100'><text y='.9em' font-size='90'>📡</text></svg>"
|
||||||
/>
|
/>
|
||||||
<!-- Linked here, not imported by src/main.js, so the unit tests can
|
|
||||||
import that module in Node, which cannot import CSS. -->
|
|
||||||
<link rel="stylesheet" href="/src/styles.css" />
|
|
||||||
</head>
|
</head>
|
||||||
<body class="bg-gray-900 text-white min-h-screen">
|
<body class="bg-gray-900 text-white min-h-screen">
|
||||||
<div id="app"></div>
|
<div id="app"></div>
|
||||||
|
|||||||
@@ -1,14 +1,13 @@
|
|||||||
#!/bin/sh
|
#!/bin/sh
|
||||||
# script/frontend-test: run the frontend test suite: the unit tests in
|
# script/frontend-test: run the frontend test suite. The frontend has no
|
||||||
# test/unit/ with Node's built-in test runner, then the production
|
# unit tests; the production build serves as the test (fails on broken
|
||||||
# build, which fails on broken code.
|
# code).
|
||||||
set -eu
|
set -eu
|
||||||
|
|
||||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||||
|
|
||||||
main() {
|
main() {
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
timeout 30 node --test test/unit/*.test.js
|
|
||||||
timeout 30 yarn build
|
timeout 30 yarn build
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+29
-54
@@ -1,14 +1,14 @@
|
|||||||
|
import "./styles.css";
|
||||||
|
|
||||||
// --- Configuration -----------------------------------------------------------
|
// --- Configuration -----------------------------------------------------------
|
||||||
|
|
||||||
// Timing, axis labels, and display constants. A target check times out
|
// Timing, axis labels, and display constants. Latency above maxLatency is
|
||||||
// after requestTimeout, 80% of updateInterval, so a round's checks have
|
// clamped to "unreachable". The sparkline Y-axis is capped at
|
||||||
// all finished before the next round is due; latency above maxLatency is
|
|
||||||
// recorded as a timeout. The sparkline Y-axis is capped at
|
|
||||||
// graphMaxLatency — values above it pin to the top of the chart but still
|
// graphMaxLatency — values above it pin to the top of the chart but still
|
||||||
// display their real value in the latency figure. The history buffer holds
|
// display their real value in the latency figure. The history buffer holds
|
||||||
// maxHistoryPoints samples (historyDuration / updateInterval).
|
// maxHistoryPoints samples (historyDuration / updateInterval).
|
||||||
// reportInterval is how often collected samples are POSTed to the backend.
|
// reportInterval is how often collected samples are POSTed to the backend.
|
||||||
export const CONFIG = {
|
const CONFIG = {
|
||||||
updateInterval: 3000,
|
updateInterval: 3000,
|
||||||
maxHistoryPoints: 100,
|
maxHistoryPoints: 100,
|
||||||
reportInterval: 60000,
|
reportInterval: 60000,
|
||||||
@@ -16,7 +16,7 @@ export const CONFIG = {
|
|||||||
return (this.maxHistoryPoints * this.updateInterval) / 1000;
|
return (this.maxHistoryPoints * this.updateInterval) / 1000;
|
||||||
},
|
},
|
||||||
get requestTimeout() {
|
get requestTimeout() {
|
||||||
return this.updateInterval * 0.8;
|
return Math.min(this.updateInterval - 100, 3000);
|
||||||
},
|
},
|
||||||
get maxLatency() {
|
get maxLatency() {
|
||||||
return this.requestTimeout;
|
return this.requestTimeout;
|
||||||
@@ -503,16 +503,12 @@ class Reporter {
|
|||||||
|
|
||||||
// --- Latency Measurement -----------------------------------------------------
|
// --- Latency Measurement -----------------------------------------------------
|
||||||
|
|
||||||
// Checks one target. The check times out after CONFIG.requestTimeout; the
|
async function measureLatency(url) {
|
||||||
// caller can give it up sooner through the optional signal, which also ends
|
|
||||||
// it as a timeout.
|
|
||||||
export async function measureLatency(url, signal) {
|
|
||||||
const controller = new AbortController();
|
const controller = new AbortController();
|
||||||
const timeoutId = setTimeout(
|
const timeoutId = setTimeout(
|
||||||
() => controller.abort(),
|
() => controller.abort(),
|
||||||
CONFIG.requestTimeout,
|
CONFIG.requestTimeout,
|
||||||
);
|
);
|
||||||
signal?.addEventListener("abort", () => controller.abort());
|
|
||||||
|
|
||||||
const targetUrl = new URL(url);
|
const targetUrl = new URL(url);
|
||||||
targetUrl.searchParams.set("_cb", Date.now().toString());
|
targetUrl.searchParams.set("_cb", Date.now().toString());
|
||||||
@@ -1105,7 +1101,7 @@ function sortAndRebuildWAN(state) {
|
|||||||
|
|
||||||
// --- Main Loop ---------------------------------------------------------------
|
// --- Main Loop ---------------------------------------------------------------
|
||||||
|
|
||||||
async function tick(state, signal, onOffline) {
|
async function tick(state, onOffline) {
|
||||||
const ts = Date.now();
|
const ts = Date.now();
|
||||||
|
|
||||||
if (state.paused) {
|
if (state.paused) {
|
||||||
@@ -1127,12 +1123,11 @@ async function tick(state, signal, onOffline) {
|
|||||||
log.debug(`Tick #${state.tickCount + 1} started`);
|
log.debug(`Tick #${state.tickCount + 1} started`);
|
||||||
|
|
||||||
const results = await Promise.all(
|
const results = await Promise.all(
|
||||||
state.allHosts.map((h) => measureLatency(h.url, signal)),
|
state.allHosts.map((h) => measureLatency(h.url)),
|
||||||
);
|
);
|
||||||
|
|
||||||
// User may have paused, or the next round may have given up this
|
// User may have paused while awaiting results — discard them
|
||||||
// one's checks, while awaiting results — discard them
|
if (state.paused) return;
|
||||||
if (state.paused || signal.aborted) return;
|
|
||||||
|
|
||||||
state.tickCount++;
|
state.tickCount++;
|
||||||
|
|
||||||
@@ -1173,10 +1168,9 @@ async function tick(state, signal, onOffline) {
|
|||||||
|
|
||||||
// --- Recovery Probe ----------------------------------------------------------
|
// --- Recovery Probe ----------------------------------------------------------
|
||||||
|
|
||||||
// When offline, check 4 random WAN hosts every 500ms, giving up the checks
|
// When offline, rapidly poll 4 random WAN hosts every 500ms. As soon as any
|
||||||
// started 500ms before, so at most 4 are ever waiting. As soon as one
|
// responds, stop probing and fire a normal tick to refresh all hosts.
|
||||||
// answers, stop probing and start a new round at once.
|
function startRecoveryProbe(state, triggerTick) {
|
||||||
function startRecoveryProbe(state, startRounds) {
|
|
||||||
if (state._recoveryProbeId) return; // already running
|
if (state._recoveryProbeId) return; // already running
|
||||||
const candidates = [...state.wan];
|
const candidates = [...state.wan];
|
||||||
for (let i = candidates.length - 1; i > 0; i--) {
|
for (let i = candidates.length - 1; i > 0; i--) {
|
||||||
@@ -1187,18 +1181,15 @@ function startRecoveryProbe(state, startRounds) {
|
|||||||
log.notice(
|
log.notice(
|
||||||
`Recovery probe started (${canaries.map((h) => h.name).join(", ")})`,
|
`Recovery probe started (${canaries.map((h) => h.name).join(", ")})`,
|
||||||
);
|
);
|
||||||
state._recoveryProbeId = setInterval(() => {
|
state._recoveryProbeId = setInterval(async () => {
|
||||||
if (state.paused) return;
|
if (state.paused) return;
|
||||||
state._recoveryProbeChecks?.abort();
|
const results = await Promise.all(
|
||||||
const checks = new AbortController();
|
canaries.map((h) => measureLatency(h.url)),
|
||||||
state._recoveryProbeChecks = checks;
|
);
|
||||||
for (const host of canaries) {
|
if (results.some((r) => r.error === null)) {
|
||||||
measureLatency(host.url, checks.signal).then((r) => {
|
log.notice("Recovery probe: connectivity detected");
|
||||||
if (r.error !== null || checks.signal.aborted) return;
|
stopRecoveryProbe(state);
|
||||||
log.notice("Recovery probe: connectivity detected");
|
triggerTick();
|
||||||
stopRecoveryProbe(state);
|
|
||||||
startRounds();
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
}, 500);
|
}, 500);
|
||||||
}
|
}
|
||||||
@@ -1207,7 +1198,6 @@ function stopRecoveryProbe(state) {
|
|||||||
if (state._recoveryProbeId) {
|
if (state._recoveryProbeId) {
|
||||||
clearInterval(state._recoveryProbeId);
|
clearInterval(state._recoveryProbeId);
|
||||||
state._recoveryProbeId = null;
|
state._recoveryProbeId = null;
|
||||||
state._recoveryProbeChecks?.abort();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1384,34 +1374,18 @@ async function init() {
|
|||||||
updateClocks();
|
updateClocks();
|
||||||
setInterval(updateClocks, 1000);
|
setInterval(updateClocks, 1000);
|
||||||
|
|
||||||
// Rounds never overlap: a round first gives up the last round's checks
|
|
||||||
// if they are still waiting, and the last round then records nothing.
|
|
||||||
// At a steady interval they never are, as they time out at 80% of it;
|
|
||||||
// they can be when a round starts early, after an interval change or
|
|
||||||
// when the recovery probe finds a target answering.
|
|
||||||
let roundChecks = new AbortController();
|
|
||||||
function doTick() {
|
function doTick() {
|
||||||
roundChecks.abort();
|
tick(state, () => startRecoveryProbe(state, doTick));
|
||||||
roundChecks = new AbortController();
|
|
||||||
tick(state, roundChecks.signal, () =>
|
|
||||||
startRecoveryProbe(state, startRounds),
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Starts a round now and then one every CONFIG.updateInterval.
|
doTick();
|
||||||
let tickIntervalId;
|
let tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
||||||
function startRounds() {
|
|
||||||
clearInterval(tickIntervalId);
|
|
||||||
doTick();
|
|
||||||
tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
|
||||||
}
|
|
||||||
|
|
||||||
startRounds();
|
|
||||||
|
|
||||||
document
|
document
|
||||||
.getElementById("interval-select")
|
.getElementById("interval-select")
|
||||||
.addEventListener("change", (e) => {
|
.addEventListener("change", (e) => {
|
||||||
const newInterval = parseInt(e.target.value, 10);
|
const newInterval = parseInt(e.target.value, 10);
|
||||||
|
clearInterval(tickIntervalId);
|
||||||
CONFIG.updateInterval = newInterval;
|
CONFIG.updateInterval = newInterval;
|
||||||
log.notice(
|
log.notice(
|
||||||
`Interval changed to ${humanDuration(newInterval / 1000)}, history reset`,
|
`Interval changed to ${humanDuration(newInterval / 1000)}, history reset`,
|
||||||
@@ -1460,7 +1434,8 @@ async function init() {
|
|||||||
|
|
||||||
// Start immediately with new interval
|
// Start immediately with new interval
|
||||||
stopRecoveryProbe(state);
|
stopRecoveryProbe(state);
|
||||||
startRounds();
|
doTick();
|
||||||
|
tickIntervalId = setInterval(doTick, CONFIG.updateInterval);
|
||||||
});
|
});
|
||||||
|
|
||||||
window.addEventListener("resize", () => handleResize(state));
|
window.addEventListener("resize", () => handleResize(state));
|
||||||
@@ -1469,7 +1444,7 @@ async function init() {
|
|||||||
|
|
||||||
// Bootstrap only when loaded as the page: a real DOM containing the #app
|
// Bootstrap only when loaded as the page: a real DOM containing the #app
|
||||||
// mount point this module renders into. Importing the module in a unit test
|
// mount point this module renders into. Importing the module in a unit test
|
||||||
// (which has no #app) runs nothing, so its exports can be tested in isolation.
|
// (which has no #app) runs nothing, so buildReport can be tested in isolation.
|
||||||
if (typeof document !== "undefined" && document.getElementById("app")) {
|
if (typeof document !== "undefined" && document.getElementById("app")) {
|
||||||
if (document.readyState === "loading") {
|
if (document.readyState === "loading") {
|
||||||
document.addEventListener("DOMContentLoaded", init);
|
document.addEventListener("DOMContentLoaded", init);
|
||||||
|
|||||||
@@ -1,66 +0,0 @@
|
|||||||
// Unit tests for src/main.js, run by script/frontend-test with Node's
|
|
||||||
// built-in test runner. Importing the module does not start the page.
|
|
||||||
|
|
||||||
import { test } from "node:test";
|
|
||||||
import assert from "node:assert/strict";
|
|
||||||
import { CONFIG, measureLatency } from "../../src/main.js";
|
|
||||||
|
|
||||||
// measureLatency writes timeouts to the debug log, which looks for its
|
|
||||||
// panel in the page. There is no page here.
|
|
||||||
globalThis.document = { getElementById: () => null };
|
|
||||||
|
|
||||||
// Mocks the clock for test t, so that a check lasting seconds takes no real
|
|
||||||
// time, and replaces fetch with a target that answers after answerAfter
|
|
||||||
// milliseconds of that clock, or never when answerAfter is Infinity. Both
|
|
||||||
// are restored when the test ends.
|
|
||||||
function mockTarget(t, answerAfter) {
|
|
||||||
t.mock.timers.enable({ apis: ["setTimeout", "Date"] });
|
|
||||||
t.mock.method(performance, "now", () => Date.now());
|
|
||||||
t.mock.method(
|
|
||||||
globalThis,
|
|
||||||
"fetch",
|
|
||||||
(url, { signal }) =>
|
|
||||||
new Promise((resolve, reject) => {
|
|
||||||
if (answerAfter !== Infinity) setTimeout(resolve, answerAfter);
|
|
||||||
signal.addEventListener("abort", () => reject(signal.reason));
|
|
||||||
}),
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
// The result of check if it has ended, otherwise "still waiting".
|
|
||||||
function settled(check) {
|
|
||||||
return Promise.race([
|
|
||||||
check,
|
|
||||||
new Promise((resolve) => setImmediate(resolve, "still waiting")),
|
|
||||||
]);
|
|
||||||
}
|
|
||||||
|
|
||||||
for (const interval of [10000, 30000]) {
|
|
||||||
const timeout = interval * 0.8;
|
|
||||||
// Over 3 seconds, which the timeout used to be capped at.
|
|
||||||
const slowAnswer = timeout - 1000;
|
|
||||||
|
|
||||||
test(`at a ${interval}ms interval, an answer after ${slowAnswer}ms is recorded with its real time`, async (t) => {
|
|
||||||
CONFIG.updateInterval = interval;
|
|
||||||
mockTarget(t, slowAnswer);
|
|
||||||
const check = measureLatency("https://target.test");
|
|
||||||
t.mock.timers.tick(slowAnswer);
|
|
||||||
assert.deepEqual(await settled(check), {
|
|
||||||
latency: slowAnswer,
|
|
||||||
error: null,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
test(`at a ${interval}ms interval, a target that never answers is recorded as a timeout after ${timeout}ms`, async (t) => {
|
|
||||||
CONFIG.updateInterval = interval;
|
|
||||||
mockTarget(t, Infinity);
|
|
||||||
const check = measureLatency("https://target.test");
|
|
||||||
t.mock.timers.tick(timeout - 1);
|
|
||||||
assert.equal(await settled(check), "still waiting");
|
|
||||||
t.mock.timers.tick(1);
|
|
||||||
assert.deepEqual(await settled(check), {
|
|
||||||
latency: null,
|
|
||||||
error: "timeout",
|
|
||||||
});
|
|
||||||
});
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user