Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions cmd/relayfile-cli/commandspec.go
Original file line number Diff line number Diff line change
Expand Up @@ -936,6 +936,7 @@ func mountOptions() []cliOptionSpec {
{Flags: "--mode <mode>", Description: "mount mode: poll (recommended) or fuse"},
{Flags: "--interval <duration>", Description: "sync interval"},
{Flags: "--interval-jitter <ratio>", Description: "sync interval jitter ratio (0.0-1.0)"},
{Flags: "--startup-jitter <duration>", Description: "maximum random delay before the first sync cycle, so mounts started together (scheduled sandboxes, scoped siblings) do not bootstrap in lockstep (0 disables)"},
{Flags: "--timeout <duration>", Description: "per-sync timeout"},
{Flags: "--bootstrap-timeout <duration>", Description: "hard cap for the one-time/full-tree bootstrap pull (0 = unbounded while making progress)"},
{Flags: "--bootstrap-max-files-per-cycle <n>", Description: "maximum files materialized per resumable tree-bootstrap cycle (-1 = legacy unbounded tree behavior)"},
Expand Down
60 changes: 59 additions & 1 deletion cmd/relayfile-cli/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"fmt"
"io"
"log"
"math"
mathrand "math/rand/v2"
"mime"
"net/http"
Expand Down Expand Up @@ -7357,6 +7358,7 @@ func runMount(args []string) error {
mode := fs.String("mode", envOrDefault("RELAYFILE_MOUNT_MODE", defaultMountMode), "mount mode: poll (recommended) or fuse")
interval := fs.Duration("interval", durationEnv("RELAYFILE_MOUNT_INTERVAL", defaultMountInterval), "sync interval")
intervalJitter := fs.Float64("interval-jitter", floatEnv("RELAYFILE_MOUNT_INTERVAL_JITTER", 0.2), "sync interval jitter ratio (0.0-1.0)")
startupJitter := fs.Duration("startup-jitter", durationEnv("RELAYFILE_MOUNT_STARTUP_JITTER", defaultStartupJitter), "maximum random delay before the first sync cycle, so mounts started together (scheduled sandboxes, scoped siblings) do not bootstrap in lockstep (0 disables)")
timeout := fs.Duration("timeout", durationEnv("RELAYFILE_MOUNT_TIMEOUT", defaultMountTimeout), "per-sync timeout")
bootstrapTimeout := fs.Duration("bootstrap-timeout", durationEnv("RELAYFILE_BOOTSTRAP_TIMEOUT", 0), "hard cap for the one-time/full-tree bootstrap pull (0 = unbounded while making progress)")
bootstrapMaxFiles := fs.Int("bootstrap-max-files-per-cycle", intEnv("RELAYFILE_BOOTSTRAP_MAX_FILES_PER_CYCLE", 2000), "maximum files materialized per resumable tree-bootstrap cycle (-1 = legacy unbounded tree behavior)")
Expand Down Expand Up @@ -7820,6 +7822,7 @@ func runMount(args []string) error {
*timeout = defaultMountTimeout
}
*intervalJitter = clampJitterRatio(*intervalJitter)
*startupJitter = clampStartupJitter(*startupJitter)

registerPID := shouldRegisterMountPID(*daemonized, *once)
if *daemonized {
Expand Down Expand Up @@ -7935,6 +7938,7 @@ func runMount(args []string) error {
*timeout,
*interval,
*intervalJitter,
*startupJitter,
*websocketEnabled,
*once,
*daemonized,
Expand Down Expand Up @@ -14095,6 +14099,42 @@ func boolPtr(value bool) *bool {
return &value
}

// defaultStartupJitter spreads the first (possibly full-tree) sync of mounts
// that start in the same instant — scheduled sandboxes at a cron boundary, and
// every scoped runner of one mount — so they do not all hit the
// single-threaded workspace Durable Object in the same second. Mirrors
// cmd/relayfile-mount.
const (
defaultStartupJitter = 5 * time.Second
maxStartupJitter = 5 * time.Minute
)

// startupSplaySample draws the uniform sample for the startup splay; tests pin
// it so a wait can be asserted without depending on a random draw.
var startupSplaySample = mathrand.Float64

func clampStartupJitter(value time.Duration) time.Duration {
if value < 0 {
return 0
}
if value > maxStartupJitter {
return maxStartupJitter
}
return value
}

// startupSplayDelay maps a uniform sample in [0,1) onto [0, max).
func startupSplayDelay(max time.Duration, sample float64) time.Duration {
max = clampStartupJitter(max)
if max == 0 || sample <= 0 {
return 0
}
if sample >= 1 {
sample = math.Nextafter(1, 0)
}
return time.Duration(sample * float64(max))
}

func clampJitterRatio(value float64) float64 {
if value < 0 {
return 0
Expand Down Expand Up @@ -14343,6 +14383,7 @@ func runMountLoop(rootCtx context.Context, syncer *mountsync.Syncer, localDir, w
timeout,
interval,
intervalJitter,
0,
websocketEnabled,
once,
daemonized,
Expand All @@ -14352,7 +14393,7 @@ func runMountLoop(rootCtx context.Context, syncer *mountsync.Syncer, localDir, w
)
}

func runMountLoopWithAuthLock(rootCtx context.Context, syncer *mountsync.Syncer, localDir, workspaceID, serverURL, delegatedCredsFile string, timeout, interval time.Duration, intervalJitter float64, websocketEnabled, once, daemonized bool, pidFile, logFile string, authMu *sync.Mutex) error {
func runMountLoopWithAuthLock(rootCtx context.Context, syncer *mountsync.Syncer, localDir, workspaceID, serverURL, delegatedCredsFile string, timeout, interval time.Duration, intervalJitter float64, startupJitter time.Duration, websocketEnabled, once, daemonized bool, pidFile, logFile string, authMu *sync.Mutex) error {
interval = enforcePollIntervalFloor(interval)
httpClient, _ := syncerClient(syncer)
record, _ := workspaceRecordByID(workspaceID)
Expand Down Expand Up @@ -14719,6 +14760,23 @@ func runMountLoopWithAuthLock(rootCtx context.Context, syncer *mountsync.Syncer,
}

log.Print(mountStartBanner(localDir, interval, intervalJitter))
// Each runner (one per scope) draws its own splay. The watcher is already
// admitting local edits, so only the first remote cycle waits.
if delay := startupSplayDelay(startupJitter, startupSplaySample()); delay > 0 {
log.Printf("mount startup splay: waiting %s before first sync", delay.Round(time.Millisecond))
splay := time.NewTimer(delay)
select {
case <-splay.C:
case <-rootCtx.Done():
splay.Stop()
if once {
return fmt.Errorf("initial sync cancelled before first cycle: %w", rootCtx.Err())
}
log.Printf("mount sync stopping: %v", rootCtx.Err())
writeSnapshot()
return nil
}
}
initialErr := runCycle(true)
logStuckEventSummary(syncer, initialErr)
if mountsync.IsBootstrapTerminalError(initialErr) {
Expand Down
77 changes: 77 additions & 0 deletions cmd/relayfile-cli/startup_splay_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package main

import (
"context"
"io"
"log"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

"github.com/agentworkforce/relayfile/internal/mountsync"
)

// `relayfile mount` runs its own loop, not relayfile-mount's. Its first
// (possibly full-tree) cycle must also wait out the startup splay, or every
// scheduled sandbox and every scoped runner bootstraps in the same second.
func TestPublicMountLoopWaitsStartupSplayBeforeFirstRequest(t *testing.T) {
t.Setenv("HOME", t.TempDir())
clearRelayfileEnv(t)
var requests atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
requests.Add(1)
switch {
case strings.Contains(r.URL.Path, "/fs/tree"):
_, _ = w.Write([]byte(`{"path":"/","entries":[]}`))
case strings.Contains(r.URL.Path, "/fs/events"):
_, _ = w.Write([]byte(`{"events":[]}`))
default:
http.NotFound(w, r)
}
}))
defer server.Close()
previous := startupSplaySample
startupSplaySample = func() float64 { return 0.5 } // 2.5m of a 5m window
t.Cleanup(func() { startupSplaySample = previous })
prevLog := log.Writer()
log.SetOutput(io.Discard)
defer log.SetOutput(prevLog)

localDir := filepath.Join(t.TempDir(), "mount")
ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
defer cancel()
disableWebSocket := false
syncer, err := mountsync.NewSyncer(mountsync.NewHTTPClient(server.URL, "test-token", server.Client()), mountsync.SyncerOptions{
WorkspaceID: "ws_public_startup_splay", RemoteRoot: "/", LocalRoot: localDir,
Interval: time.Hour, WebSocket: &disableWebSocket, RootCtx: ctx, Logger: log.New(io.Discard, "", 0),
})
if err != nil {
t.Fatalf("NewSyncer: %v", err)
}
err = runMountLoopWithAuthLock(ctx, syncer, localDir, "ws_public_startup_splay", server.URL, "",
time.Second, time.Hour, 0, maxStartupJitter, false, true, false,
mountPIDFile(localDir), mountLogFile(localDir), &sync.Mutex{})
if err == nil {
t.Fatal("one-shot public mount returned success while its startup splay outlasted the deadline")
}
if got := requests.Load(); got != 0 {
t.Fatalf("public mount made %d request(s) during the startup splay, want 0", got)
}
}

func TestPublicMountStartupSplayDelayBounds(t *testing.T) {
if got := startupSplayDelay(0, 0.9); got != 0 {
t.Fatalf("disabled splay = %s, want 0", got)
}
if got := startupSplayDelay(10*time.Second, 0.5); got != 5*time.Second {
t.Fatalf("half sample = %s, want 5s", got)
}
if got := startupSplayDelay(time.Hour, 0.999); got >= maxStartupJitter {
t.Fatalf("oversized splay = %s, want < %s", got, maxStartupJitter)
}
}
91 changes: 80 additions & 11 deletions cmd/relayfile-mount/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"flag"
"fmt"
"log"
"math"
"math/rand"
"net/http"
"os"
Expand Down Expand Up @@ -66,6 +67,7 @@ type mountConfig struct {
syncMode string
interval time.Duration
intervalJitter float64
startupJitter time.Duration
timeout time.Duration
bootstrapTimeout time.Duration
bootstrapMaxFiles int
Expand Down Expand Up @@ -122,6 +124,7 @@ func main() {
syncModeFlag := flag.String("sync-mode", envOrDefault("RELAYFILE_MOUNT_SYNC_MODE", syncModeMirror), "sync behavior: mirror (pull and push), pull-only (poll mode only; mirror remote changes without writeback), or write-only (push local changes without mirroring provider history)")
interval := flag.Duration("interval", durationEnv("RELAYFILE_MOUNT_INTERVAL", 30*time.Second), "sync interval")
intervalJitter := flag.Float64("interval-jitter", floatEnv("RELAYFILE_MOUNT_INTERVAL_JITTER", 0.2), "sync interval jitter ratio (0.0-1.0)")
startupJitter := flag.Duration("startup-jitter", durationEnv("RELAYFILE_MOUNT_STARTUP_JITTER", defaultStartupJitter), "maximum random delay before the first sync cycle, so mounts started together (scheduled sandboxes, scoped siblings) do not bootstrap in lockstep (0 disables)")
timeout := flag.Duration("timeout", durationEnv("RELAYFILE_MOUNT_TIMEOUT", 15*time.Second), "per-sync timeout")
bootstrapTimeout := flag.Duration("bootstrap-timeout", durationEnv("RELAYFILE_BOOTSTRAP_TIMEOUT", 0), "hard cap for the one-time/full-tree bootstrap pull (0 = unbounded while making progress)")
bootstrapMaxFiles := flag.Int("bootstrap-max-files-per-cycle", intEnv("RELAYFILE_BOOTSTRAP_MAX_FILES_PER_CYCLE", 2000), "maximum files materialized per resumable tree-bootstrap cycle (-1 = legacy unbounded tree behavior)")
Expand Down Expand Up @@ -193,6 +196,7 @@ func main() {
}
allRemotePaths := append(remotePaths.Values(), fileRemotePaths...)
*intervalJitter = clampJitterRatio(*intervalJitter)
*startupJitter = clampStartupJitter(*startupJitter)
resolvedMode, err := resolveMountMode(*mode, *fuse)
if err != nil {
log.Fatalf("invalid mount mode: %v", err)
Expand Down Expand Up @@ -233,6 +237,7 @@ func main() {
syncMode: resolvedSyncMode,
interval: *interval,
intervalJitter: *intervalJitter,
startupJitter: *startupJitter,
timeout: *timeout,
bootstrapTimeout: *bootstrapTimeout,
bootstrapMaxFiles: *bootstrapMaxFiles,
Expand Down Expand Up @@ -727,7 +732,23 @@ func runSinglePollingMount(rootCtx context.Context, cfg mountConfig) error {
// "just finished as part of this attempt."
priorBootstrapComplete := bootstrapAlreadyComplete(cfg.localDir)

if err := run(true); err != nil {
flushedDuringSplay, err := waitStartupSplay(rootCtx, startupSplayDelay(cfg.startupJitter, startupSplaySample()), cfg.flushReq)
if err != nil {
if cfg.once {
return fmt.Errorf("initial sync cancelled before first cycle: %w", err)
}
log.Printf("mount sync stopping: %v", err)
return nil
}
if flushedDuringSplay {
// An explicit flush ends the splay: the operator asked for a sync
// now, and the notifier stops waiting for an ack long before a
// full-length splay would elapse. The kicked reconcile is the first
// cycle.
if err := serviceFlushRequest(rootCtx, cfg, syncer); err != nil {
return err
}
} else if err := run(true); err != nil {
return err
}
if cfg.once {
Expand Down Expand Up @@ -791,16 +812,8 @@ func runSinglePollingMount(rootCtx context.Context, cfg mountConfig) error {
log.Printf("mount sync stopping: %v", rootCtx.Err())
return nil
case <-cfg.flushReq:
kickErr := kickReconcile(rootCtx, cfg, syncer)
if recErr := recordFlushAck(cfg, kickErr); recErr != nil {
log.Printf("mount flush ack failed: %v", recErr)
} else if kickErr != nil {
log.Printf("mount flush requested via SIGUSR1; failed: %v", kickErr)
} else {
log.Printf("mount flush requested via SIGUSR1; ack recorded")
}
if kickErr != nil && mountsync.IsBootstrapTerminalError(kickErr) {
return kickErr
if err := serviceFlushRequest(rootCtx, cfg, syncer); err != nil {
return err
}
case <-wsTicker.C:
if mountWebSocketEnabled(cfg) {
Expand Down Expand Up @@ -1532,6 +1545,62 @@ func mountReconcileUsesWebSocketCadence(cfg mountConfig, watcherActive bool) boo
return mountWebSocketEnabled(cfg) && (cfg.syncMode == syncModePullOnly || watcherActive)
}

// defaultStartupJitter spreads the first (possibly full-tree) sync of mounts
// that start in the same instant. Scheduled sandboxes launch together at cron
// boundaries and every scoped sibling starts its own Syncer at once; without a
// splay they all hit the single-threaded workspace Durable Object in the same
// second. Five seconds is small against the bootstrap readiness budget.
const (
defaultStartupJitter = 5 * time.Second
maxStartupJitter = 5 * time.Minute
)

// startupSplaySample draws the uniform sample for the startup splay; tests pin
// it so a wait can be asserted without depending on a random draw.
var startupSplaySample = rand.Float64

func clampStartupJitter(value time.Duration) time.Duration {
if value < 0 {
return 0
}
if value > maxStartupJitter {
return maxStartupJitter
}
return value
}

// startupSplayDelay maps a uniform sample in [0,1) onto [0, max).
func startupSplayDelay(max time.Duration, sample float64) time.Duration {
max = clampStartupJitter(max)
if max == 0 || sample <= 0 {
return 0
}
if sample >= 1 {
sample = math.Nextafter(1, 0)
}
return time.Duration(sample * float64(max))
}

// waitStartupSplay waits out the splay unless ctx ends or a flush request
// arrives first; flushed reports the latter so the caller services it.
func waitStartupSplay(ctx context.Context, delay time.Duration, flushReq <-chan struct{}) (flushed bool, err error) {
if delay <= 0 {
return false, ctx.Err()
}
log.Printf("mount startup splay: waiting %s before first sync", delay.Round(time.Millisecond))
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-timer.C:
return false, nil
case <-flushReq:
log.Printf("mount startup splay: flush requested; starting first sync now")
return true, nil
case <-ctx.Done():
return false, ctx.Err()
}
}

func clampJitterRatio(value float64) float64 {
if value < 0 {
return 0
Expand Down
19 changes: 19 additions & 0 deletions cmd/relayfile-mount/notify_flush.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"context"
"log"
"os"
"time"

Expand All @@ -15,6 +16,24 @@ func kickReconcile(rootCtx context.Context, cfg mountConfig, syncer *mountsync.S
return syncer.Reconcile(ctx)
}

// serviceFlushRequest runs the reconcile a SIGUSR1 flush asked for and
// records its ack. Only a terminal bootstrap error is returned; any other
// failure is reported through the ack.
func serviceFlushRequest(rootCtx context.Context, cfg mountConfig, syncer *mountsync.Syncer) error {
kickErr := kickReconcile(rootCtx, cfg, syncer)
if recErr := recordFlushAck(cfg, kickErr); recErr != nil {
log.Printf("mount flush ack failed: %v", recErr)
} else if kickErr != nil {
log.Printf("mount flush requested via SIGUSR1; failed: %v", kickErr)
} else {
log.Printf("mount flush requested via SIGUSR1; ack recorded")
}
if kickErr != nil && mountsync.IsBootstrapTerminalError(kickErr) {
return kickErr
}
return nil
}

func recordFlushAck(cfg mountConfig, kickErr error) error {
prev, err := mountlease.ReadFlushAck(cfg.baseURL, cfg.workspaceID, cfg.localDir)
if err != nil {
Expand Down
Loading
Loading