Skip to content
Draft
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
99 changes: 91 additions & 8 deletions go/boss/boss.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package boss

import (
"context"
"encoding/json"
"fmt"
"io"
Expand All @@ -11,10 +12,12 @@ import (
"strconv"
"strings"
"syscall"
"time"

"github.com/open-lambda/open-lambda/go/boss/autoscaling"
"github.com/open-lambda/open-lambda/go/boss/cloudvm"
"github.com/open-lambda/open-lambda/go/boss/config"
"github.com/open-lambda/open-lambda/go/boss/etcd"
"github.com/open-lambda/open-lambda/go/boss/lambdastore"
)

Expand Down Expand Up @@ -144,15 +147,81 @@ func (b *Boss) RegistryHandler(w http.ResponseWriter, r *http.Request) {
}
}

// leaderOrRedirect wraps a handler so that follower bosses redirect mutating
// requests to the current leader. read-only GET requests are served locally on every replica.
func leaderOrRedirect(le *etcd.LeaderElection, next http.HandlerFunc) http.HandlerFunc {
if le == nil {
return next
}
return func(w http.ResponseWriter, r *http.Request) {
if le.IsLeader() || r.Method == http.MethodGet {
next(w, r)
return
}
ctx, cancel := context.WithTimeout(r.Context(), 3*time.Second)
defer cancel()
addr, err := le.LeaderAddr(ctx)
if err != nil {
http.Error(w, "no leader available", http.StatusServiceUnavailable)
return
}
http.Redirect(w, r, fmt.Sprintf("http://%s%s", addr, r.URL.RequestURI()), http.StatusTemporaryRedirect)
}
}

// BossMain is the main function for the boss.
func BossMain() (err error) {
fmt.Printf("WARNING! Boss incomplete (only use this as part of development process).\n")

pool, err := cloudvm.NewWorkerPool(config.BossConf.Platform, config.BossConf.Worker_Cap)
var etcdClient *etcd.Client
if len(config.BossConf.Etcd.Endpoints) > 0 {
timeout := time.Duration(config.BossConf.Etcd.Dial_timeout_sec) * time.Second
etcdClient, err = etcd.NewClient(config.BossConf.Etcd.Endpoints, config.BossConf.Etcd.Prefix, timeout)
if err != nil {
return fmt.Errorf("etcd connect failed: %w", err)
}
defer etcdClient.Close()
slog.Info("etcd connected", "endpoints", config.BossConf.Etcd.Endpoints)
}

var le *etcd.LeaderElection
var selfAddr string
if etcdClient != nil && config.BossConf.Etcd.Enable_HA {
le, err = etcdClient.NewLeaderElection(context.Background())
if err != nil {
return fmt.Errorf("leader election init failed: %w", err)
}

selfHost := config.BossConf.Boss_host
if selfHost == "" {
if h, e := os.Hostname(); e == nil {
selfHost = h
} else {
selfHost = "localhost"
}
}
selfAddr = fmt.Sprintf("%s:%s", selfHost, config.BossConf.Boss_port)
}

pool, err := cloudvm.NewWorkerPool(config.BossConf.Platform, config.BossConf.Worker_Cap, etcdClient)
if err != nil {
return err
}

if le != nil {
go func() {
slog.Info("campaigning for leadership", "self", selfAddr)
if err := le.Campaign(context.Background(), selfAddr); err != nil {
slog.Error("leader election campaign failed", "err", err)
return
}
if err := pool.ResyncFromEtcd(); err != nil {
slog.Error("etcd resync on election win failed", "err", err)
}
slog.Info("elected leader", "addr", selfAddr)
}()
}

store, err := lambdastore.NewLambdaStore(config.BossConf.GetLambdaStoreURL(), pool)
if err != nil {
return err
Expand All @@ -168,16 +237,24 @@ func BossMain() (err error) {
boss.autoScaler.Launch(boss.workerPool)
}

// Launch 1 worker by default when boss starts
slog.Info("Launching 1 worker by default")
boss.workerPool.SetTarget(1)
// restore a prior target, else default to 1.
if pool.GetTarget() == 0 {
slog.Info("Launching 1 worker by default")
boss.workerPool.SetTarget(1)
} else {
slog.Info("resuming from etcd", "target", pool.GetTarget())
boss.workerPool.SetTarget(pool.GetTarget())
}

wrap := func(h http.HandlerFunc) http.HandlerFunc {
return leaderOrRedirect(le, h)
}

http.HandleFunc(BOSS_STATUS_PATH, boss.BossStatus)
http.HandleFunc(SCALING_PATH, boss.ScalingWorker)
http.HandleFunc(RUN_PATH, boss.workerPool.RunLambda)
http.HandleFunc(SCALING_PATH, wrap(boss.ScalingWorker))
http.HandleFunc(RUN_PATH, wrap(boss.workerPool.RunLambda))
http.HandleFunc(SHUTDOWN_PATH, boss.Close)

http.HandleFunc(REGISTRY_BASE_PATH, boss.RegistryHandler)
http.HandleFunc(REGISTRY_BASE_PATH, wrap(boss.RegistryHandler))

// clean up if signal hits us
c := make(chan os.Signal, 1)
Expand All @@ -186,6 +263,12 @@ func BossMain() (err error) {
go func() {
<-c
slog.Info("received kill signal, cleaning up")
if le != nil {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
le.Resign(ctx)
le.Close()
cancel()
}
boss.Close(nil, nil)
os.Exit(0)
}()
Expand Down
4 changes: 4 additions & 0 deletions go/boss/cloudvm/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import (
"net/http"
"os"
"sync"

"github.com/open-lambda/open-lambda/go/boss/etcd"
)

// The WorkerState integer
Expand Down Expand Up @@ -51,6 +53,8 @@ type WorkerPool struct {
totalTask int32
sumLatency int64
nLatency int64

etcd *etcd.Client
}

/*
Expand Down
Loading
Loading