Compare commits
8 Commits
v1.1.0
...
528bdef433
| Author | SHA1 | Date | |
|---|---|---|---|
| 528bdef433 | |||
| dfe082629f | |||
| 39aff83837 | |||
| 27be28ee9c | |||
| a8a6fdfab5 | |||
| e9df802478 | |||
| 6c6b962e24 | |||
| e64075d5d7 |
@@ -103,15 +103,14 @@ func connectOnce(ctx context.Context, cfg Config, handle Handler) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
dialCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
dialCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||||
conn, res, err := websocket.Dial(dialCtx, wsURL, dialOpts)
|
conn, _, err := websocket.Dial(dialCtx, wsURL, dialOpts) //nolint:bodyclose // successful upgrades have a nil response body owned by conn
|
||||||
cancel()
|
cancel()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("dial: %w", err)
|
return fmt.Errorf("dial: %w", err)
|
||||||
}
|
}
|
||||||
// websocket.Dial returns the upgrade response separately from the
|
// On a successful upgrade coder/websocket transfers ownership of the
|
||||||
// conn. Body is empty on a successful upgrade but Go's net/http
|
// response stream to conn and deliberately sets res.Body to nil. Closing
|
||||||
// still expects it closed to release the connection.
|
// the connection below releases that stream.
|
||||||
defer func() { _ = res.Body.Close() }()
|
|
||||||
defer conn.CloseNow() //nolint:errcheck
|
defer conn.CloseNow() //nolint:errcheck
|
||||||
|
|
||||||
// Send hello.
|
// Send hello.
|
||||||
|
|||||||
@@ -0,0 +1,46 @@
|
|||||||
|
package wsclient
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/coder/websocket"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
||||||
|
serverErr := make(chan error, 1)
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
conn, err := websocket.Accept(w, r, nil)
|
||||||
|
if err != nil {
|
||||||
|
serverErr <- err
|
||||||
|
return
|
||||||
|
}
|
||||||
|
defer conn.CloseNow() //nolint:errcheck
|
||||||
|
|
||||||
|
// Wait for the agent hello so Dial and the first client write have both
|
||||||
|
// completed before ending the connection normally.
|
||||||
|
if _, _, err := conn.Read(r.Context()); err != nil {
|
||||||
|
serverErr <- err
|
||||||
|
return
|
||||||
|
}
|
||||||
|
serverErr <- conn.Close(websocket.StatusNormalClosure, "test complete")
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
err := connectOnce(ctx, Config{
|
||||||
|
ServerURL: srv.URL,
|
||||||
|
AgentToken: "test-token",
|
||||||
|
HeartbeatPeriod: time.Hour,
|
||||||
|
}, nil)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("connectOnce returned nil after server disconnected")
|
||||||
|
}
|
||||||
|
if err := <-serverErr; err != nil {
|
||||||
|
t.Fatalf("server websocket: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -512,11 +512,27 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
|
|||||||
// Connect the agent so DrainPending can dispatch.
|
// Connect the agent so DrainPending can dispatch.
|
||||||
c := agentDial(t, srv, ts, hostID, token)
|
c := agentDial(t, srv, ts, hostID, token)
|
||||||
sendHello(t, c, "serialise-host")
|
sendHello(t, c, "serialise-host")
|
||||||
// Drain the on-hello goroutine's pass first (no pending rows yet),
|
// Wait for the on-hello push to settle.
|
||||||
// then wait for the schedule.set so the connection is fully settled.
|
|
||||||
_ = drainUntil(t, c, api.MsgScheduleSet)
|
_ = drainUntil(t, c, api.MsgScheduleSet)
|
||||||
|
|
||||||
// Insert 5 pending rows now that the on-hello drain has already run.
|
// A real agent is always in a read loop. Keep this test client
|
||||||
|
// reading in the background for the rest of the test: without an
|
||||||
|
// active reader the server-side conn can be dropped under parallel
|
||||||
|
// load, which unregisters it from the hub and makes DrainPending
|
||||||
|
// no-op (conn == nil) — the historical source of this test's
|
||||||
|
// flakiness (it would observe 0 or a partial drain). The reader also
|
||||||
|
// consumes the command.run envelopes our drains emit.
|
||||||
|
readerCtx, stopReader := context.WithCancel(context.Background())
|
||||||
|
defer stopReader()
|
||||||
|
go func() {
|
||||||
|
for {
|
||||||
|
if _, _, err := c.Read(readerCtx); err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Insert 5 due pending rows.
|
||||||
now := time.Now().UTC()
|
now := time.Now().UTC()
|
||||||
for i := range 5 {
|
for i := range 5 {
|
||||||
pid := ulid.Make().String()
|
pid := ulid.Make().String()
|
||||||
@@ -533,7 +549,8 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Spawn 10 goroutines all calling DrainPending concurrently.
|
// Fire 10 concurrent DrainPending calls. The per-host mutex must
|
||||||
|
// ensure each row is dispatched at most once (no double-dispatch).
|
||||||
var wg sync.WaitGroup
|
var wg sync.WaitGroup
|
||||||
for range 10 {
|
for range 10 {
|
||||||
wg.Add(1)
|
wg.Add(1)
|
||||||
@@ -544,24 +561,26 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
|
|||||||
}
|
}
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
// Drain any envelopes the agent received so we don't block below.
|
// Drain to completion. The fire-and-forget on-hello DrainPending
|
||||||
// We read with short timeouts and stop when the connection goes quiet.
|
// shares the same per-host mutex and can hold it during the burst,
|
||||||
drainDeadline := time.Now().Add(500 * time.Millisecond)
|
// leaving rows for a later pass — exactly how production drains
|
||||||
for time.Now().Before(drainDeadline) {
|
// (repeatedly, via the 30s tick / on reconnect). Re-drain until the
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
// queue is empty; because every drain is still serialised, each row
|
||||||
_, _, err := c.Read(ctx)
|
// is dispatched at most once, so the exactly-5 job count below proves
|
||||||
cancel()
|
// there was no double-dispatch.
|
||||||
if err != nil {
|
deadline := time.Now().Add(5 * time.Second)
|
||||||
break
|
for countPendingForHost(t, st, hostID) > 0 && time.Now().Before(deadline) {
|
||||||
}
|
srv.DrainPending(context.Background(), hostID)
|
||||||
|
time.Sleep(10 * time.Millisecond)
|
||||||
}
|
}
|
||||||
|
|
||||||
// All 5 pending rows must be gone.
|
// All 5 pending rows must be drained.
|
||||||
if n := countPendingForHost(t, st, hostID); n != 0 {
|
if n := countPendingForHost(t, st, hostID); n != 0 {
|
||||||
t.Errorf("pending rows after concurrent drain: got %d, want 0", n)
|
t.Errorf("pending rows after drain-to-completion: got %d, want 0", n)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Exactly 5 backup job rows (one per pending row), not 10+ from a race.
|
// Exactly 5 backup job rows (one per pending row) — never more, which
|
||||||
|
// would mean the per-host mutex failed to prevent double-dispatch.
|
||||||
var n int
|
var n int
|
||||||
_ = st.DB().QueryRow(
|
_ = st.DB().QueryRow(
|
||||||
`SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'backup' AND actor_kind = 'schedule'`,
|
`SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'backup' AND actor_kind = 'schedule'`,
|
||||||
|
|||||||
@@ -82,10 +82,14 @@
|
|||||||
<div class="text-[12px] text-ok mb-3 mono">✓ saved</div>
|
<div class="text-[12px] text-ok mb-3 mono">✓ saved</div>
|
||||||
{{end}}
|
{{end}}
|
||||||
<p class="text-[12.5px] text-ink-mid leading-[1.6] mb-4 max-w-[640px]">
|
<p class="text-[12.5px] text-ink-mid leading-[1.6] mb-4 max-w-[640px]">
|
||||||
Only needed for rest-server repos that distinguish an append-only
|
Required for prune. On rest-server repos this is the
|
||||||
user (everyday backups) from a delete-capable user (prune /
|
delete-capable user, as distinct from the append-only user used
|
||||||
forget). For S3 / B2 / SFTP / local, leave this blank — the
|
for everyday backups. Note that <strong>forget</strong> always
|
||||||
everyday repo credentials handle prune too.
|
runs with the everyday repo credentials, so those must have
|
||||||
|
delete authority. For S3 / B2 / SFTP / local, enter the same
|
||||||
|
delete-capable repository credentials here if you want prune
|
||||||
|
enabled. <strong>Prune is skipped when admin credentials are
|
||||||
|
unset</strong>, on any backend.
|
||||||
</p>
|
</p>
|
||||||
<div class="grid grid-cols-2 gap-4">
|
<div class="grid grid-cols-2 gap-4">
|
||||||
<div>
|
<div>
|
||||||
|
|||||||
Reference in New Issue
Block a user