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 packages/orchestrator/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ require (
github.com/go-git/go-billy/v5 v5.9.0
github.com/go-openapi/runtime v0.29.2
github.com/go-openapi/strfmt v0.26.1
github.com/gofrs/flock v0.13.0
github.com/google/go-containerregistry v0.21.7
github.com/google/nftables v0.3.0
github.com/google/uuid v1.6.0
Expand Down
2 changes: 2 additions & 0 deletions packages/orchestrator/go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

91 changes: 72 additions & 19 deletions packages/orchestrator/pkg/factories/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,12 @@ import (
"os"
"os/signal"
"slices"
"strconv"
"strings"
"syscall"
"time"

"github.com/gofrs/flock"
"github.com/google/uuid"
"github.com/soheilhy/cmux"
"go.opentelemetry.io/otel/attribute"
Expand Down Expand Up @@ -170,6 +172,61 @@ func ensureDirs(c cfg.Config) error {
return nil
}

func acquireOrchestratorLock(path string) (*flock.Flock, error) {
fileLock := flock.New(path, flock.SetPermissions(0o644))
locked, err := fileLock.TryLock()
if err != nil {
return nil, fmt.Errorf("lock file: %w", err)
}
if !locked {
// flock(2) is released by the kernel on crash, so reaching here means a
// live process holds the lock. Surface the PID it recorded, if any.
if pid, perr := readLockHolderPID(path); perr == nil && pid > 0 {
return nil, fmt.Errorf("another instance is running with pid %d", pid)
}

return nil, errors.New("another instance is running")
}

// Record our PID so a future conflicting instance can report it. We hold the
// exclusive advisory lock, so no other process writes this file concurrently.
if err := writeLockHolderPID(path); err != nil {
_ = fileLock.Unlock()

return nil, fmt.Errorf("write lock holder pid: %w", err)
}

return fileLock, nil
}

func writeLockHolderPID(path string) error {
f, err := os.OpenFile(path, os.O_WRONLY|os.O_TRUNC, 0o644)
if err != nil {
return fmt.Errorf("open lock file: %w", err)
}
defer f.Close()

if _, err := fmt.Fprintf(f, "%d\n", os.Getpid()); err != nil {
return fmt.Errorf("write pid: %w", err)
}

return nil
}

func readLockHolderPID(path string) (int, error) {
data, err := os.ReadFile(path)
if err != nil {
return 0, fmt.Errorf("read lock file: %w", err)
}

pid, err := strconv.Atoi(strings.TrimSpace(string(data)))
if err != nil {
return 0, fmt.Errorf("parse pid: %w", err)
}

return pid, nil
}

func run(config cfg.Config, opts Options) (success bool) {
success = true

Expand All @@ -178,31 +235,27 @@ func run(config cfg.Config, opts Options) (success bool) {

services := cfg.GetServices(config)

// Check if the orchestrator crashed and restarted
// Skip this check in development mode
// We don't want to lock if the service is running with force stop; the subsequent start would fail.
if !env.IsDevelopment() && !config.ForceStop && services.RunsOrchestrator() {
fileLockName := config.OrchestratorLockPath
info, err := os.Stat(fileLockName)
if err == nil {
log.Fatalf("Orchestrator was already started at %s, exiting", info.ModTime())
}
usesSandboxRuntime := services.UsesSandboxRuntime()

Comment thread
wj-e2b marked this conversation as resolved.
f, err := os.Create(fileLockName)
// Enforce a single host-level sandbox runtime instance.
// Skip this check in development mode.
if !env.IsDevelopment() && usesSandboxRuntime {
f, err := acquireOrchestratorLock(config.OrchestratorLockPath)
if err != nil {
log.Fatalf("Failed to create lock file %s: %v", fileLockName, err)
log.Fatalf("Failed to acquire orchestrator lock %s: %v", config.OrchestratorLockPath, err)
}
defer func() {
fileErr := f.Close()
if fileErr != nil {
log.Printf("Failed to close lock file %s: %v", fileLockName, fileErr)
log.Printf("Failed to close lock file %s: %v", config.OrchestratorLockPath, fileErr)
}

// Remove the lock file on graceful shutdown
if success == true {
if fileErr = os.Remove(fileLockName); fileErr != nil {
log.Printf("Failed to remove lock file %s: %v", fileLockName, fileErr)
}
// Remove the lock file on clean shutdown so a rollback to the older
// stat-based release can start: that guard exits whenever the lock
// file exists and cannot tell that this process is already gone.
// TODO: Remove this os.Remove once all hosts run a flock-based
// release and rollback to the stat-based guard is no longer possible.
if rmErr := os.Remove(config.OrchestratorLockPath); rmErr != nil && !os.IsNotExist(rmErr) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Avoid unlinking a successor's lock file

When a replacement process starts after f.Close() releases the flock but before this os.Remove, it can acquire and write the existing lock path; the exiting process then unlinks that successor's locked inode. A later third process will create a fresh /orchestrator.lock and acquire an independent flock, so rolling/restart timing can still allow two sandbox-runtime instances on the same host.

Useful? React with 👍 / 👎.

log.Printf("Failed to remove lock file %s: %v", config.OrchestratorLockPath, rmErr)
}
}()
Comment thread
wj-e2b marked this conversation as resolved.
}
Expand Down Expand Up @@ -626,7 +679,7 @@ func run(config cfg.Config, opts Options) (success bool) {
// Sandbox-runtime reclaim must run before newStorage below: reclaim deletes
// leaked ns-* from /run/netns, and NewStorageLocal snapshots the remaining
// namespaces as foreign at construction.
if services.UsesSandboxRuntime() && !config.DisableStartupReclaim {
if usesSandboxRuntime && !config.DisableStartupReclaim {
startupreclaim.Run(ctx, startupreclaim.Config{
NetworkConfig: config.NetworkConfig,
EgressProxy: egressSetup.Proxy,
Expand Down
Loading