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: 0 additions & 1 deletion core/repomanager/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ go_library(
"//core/git",
"//core/workspace",
"//tangopb",
"@com_github_gofrs_flock//:flock",
"@org_uber_go_zap//:zap",
],
)
Expand Down
184 changes: 138 additions & 46 deletions core/repomanager/repo_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,17 @@ import (
"os"
"path/filepath"
"strings"
"time"
"sync"

"github.com/uber/tango/core/git"
"github.com/uber/tango/core/workspace"
"github.com/uber/tango/tangopb"
"github.com/gofrs/flock"
"go.uber.org/zap"
)

const (
lockFileName = ".tango.lease.lock"

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

not needed anymore. We only need a mutex per repo, and workers concurrency is managed by avail chan

lockTimeout = 5 * time.Minute

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

instead of timeout, we just wait until a workspace is available. If client doesn't wanna wait, it can just cancel the request

lockRetryDelay = 15 * time.Second
)
const defaultPoolSize = 3

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

will do next diff, make it configurable


// RepoManager manages repository workspaces with a pool of workers per repo.
type RepoManager interface {
Lease(ctx context.Context, desc tangopb.BuildDescription) (workspace.Workspace, error)
}
Expand All @@ -29,66 +25,162 @@ type repoManager struct {
git git.Interface
rootWorkspace string
logger *zap.SugaredLogger
poolSize int

mu sync.Mutex
pools map[string]*workerPool
}

// workerPool manages a fixed set of worker slots for a single repo.
// The origin directory holds the initial clone; workers are cheap local copies.
type workerPool struct {
originDir string
originMu sync.Mutex // one lock per repo for orginal clone
cloned bool

avail chan *workerSlot // available slots; = pool capacity
}

// workerSlot is a pre-allocated workspace directory that may or may not
// have been cloned yet. Lazy creation on first use.
type workerSlot struct {
dir string
created bool
}

// Params for creating a RepoManager.
type Params struct {
Git git.Interface
Logger *zap.SugaredLogger
RootWorkspace string
PoolSize int // number of worker workspaces per repo; 0 uses default (3)
}

// NewRepoManager creates a new repo manager with the given git interface and root workspace.
// NewRepoManager creates a new repo manager with pooled worker workspaces.
func NewRepoManager(p Params) RepoManager {
return &repoManager{git: p.Git, rootWorkspace: p.RootWorkspace, logger: p.Logger}
size := p.PoolSize
if size <= 0 {
size = defaultPoolSize
}
return &repoManager{
git: p.Git,
rootWorkspace: p.RootWorkspace,
logger: p.Logger,
poolSize: size,
pools: make(map[string]*workerPool),
}
}

// Lease tries to take an exclusive lease on the repo’s workspace.
// If another process holds it, wait with a timeout.
func (r *repoManager) Lease(ctx context.Context, desc tangopb.BuildDescription) (workspace.Workspace, error) {
repo := toShortRemote(desc.Remote)
repoDir := filepath.Join(r.rootWorkspace, repo)
// Lock file must be in the root workspace as git clone cannot be executed on a non-empty directory.
lockPath := filepath.Join(r.rootWorkspace, lockFileName)
if err := os.MkdirAll(r.rootWorkspace, 0o755); err != nil {
return nil, fmt.Errorf("mkdir repo dir: %w", err)
func (r *repoManager) poolFor(repo string) *workerPool {
r.mu.Lock()
defer r.mu.Unlock()

if pool, ok := r.pools[repo]; ok {
return pool
}

fl := flock.New(lockPath)
pool := &workerPool{
originDir: filepath.Join(r.rootWorkspace, repo),
avail: make(chan *workerSlot, r.poolSize),
}

waitCtx, cancel := context.WithTimeout(ctx, lockTimeout)
defer cancel()
// Pre-allocate fixed worker slots. Existing directories from a previous
// run are detected and reused without re-cloning.
workersDir := filepath.Join(r.rootWorkspace, ".workers", repo)
for i := 1; i <= r.poolSize; i++ {
dir := filepath.Join(workersDir, fmt.Sprintf("worker-%d", i))
slot := &workerSlot{dir: dir}
if _, err := os.Stat(filepath.Join(dir, ".git")); err == nil {
slot.created = true
r.logger.Debugf("discovered existing worker: %s", dir)
}
pool.avail <- slot
}

r.pools[repo] = pool
return pool
}

// Lease borrows a worker workspace from the pool.
// If all workers are leased, it blocks until one is returned or ctx is cancelled.
func (r *repoManager) Lease(ctx context.Context, desc tangopb.BuildDescription) (workspace.Workspace, error) {
repo := toShortRemote(desc.Remote)
pool := r.poolFor(repo)

// This will block until:
// 1) the lock is released by another process, OR
// 2) waitCtx expires/canceled
locked, err := fl.TryLockContext(waitCtx, lockRetryDelay)
if err != nil {
return nil, fmt.Errorf("failed to acquire lock for %s: %w", r.rootWorkspace, err)
if err := pool.ensureOrigin(ctx, r.git, desc.Remote); err != nil {
return nil, err
}
if !locked {
return nil, fmt.Errorf("lock timeout after %s for %s", lockTimeout, r.rootWorkspace)

// Acquire a worker slot (blocks if all slots are leased)
var slot *workerSlot
select {
case slot = <-pool.avail:
case <-ctx.Done():
return nil, ctx.Err()
}
if _, err := os.Stat(filepath.Join(repoDir, ".git")); os.IsNotExist(err) {
if err := r.git.Clone(ctx, desc.Remote, repoDir); err != nil {
flockErr := fl.Unlock()
if flockErr != nil {
r.logger.Errorf("unlock failed: %w", flockErr)
}
removeErr := os.RemoveAll(repoDir)
if removeErr != nil {
r.logger.Errorf("remove repo dir failed: %w", removeErr)
}
return nil, fmt.Errorf("clone failed: %w", err)

// Lazily create the worker clone on first use
if !slot.created {
if err := r.createWorker(ctx, pool.originDir, slot.dir); err != nil {
pool.avail <- slot // return slot so others can retry
return nil, fmt.Errorf("create worker: %w", err)
}
slot.created = true
}
// Use a Git interface rooted at the repo directory so commands run in the correct working directory.
repoGit := git.New(repoDir)
return workspace.NewWorkspace(workspace.WorkspaceParams{
Path: repoDir,
Lock: fl,

repoGit := git.New(slot.dir)
ws := workspace.NewWorkspace(workspace.WorkspaceParams{
Path: slot.dir,
Git: repoGit,
Logger: r.logger,
}), nil
})
return &pooledWorkspace{Workspace: ws, pool: pool, slot: slot}, nil
}

// ensureOrigin clones the origin repository if it doesn't exist yet.
func (p *workerPool) ensureOrigin(ctx context.Context, g git.Interface, remote string) error {
p.originMu.Lock()
defer p.originMu.Unlock()

if p.cloned {
return nil
}
if _, err := os.Stat(filepath.Join(p.originDir, ".git")); err == nil {
p.cloned = true
return nil
}

if err := os.MkdirAll(filepath.Dir(p.originDir), 0o755); err != nil {
return fmt.Errorf("mkdir origin dir: %w", err)
}
if err := g.Clone(ctx, remote, p.originDir); err != nil {
os.RemoveAll(p.originDir)
return fmt.Errorf("clone origin: %w", err)
}
p.cloned = true
return nil
}

// createWorker creates a worker by cloning the origin with --local
// (fast and space-efficient).
func (r *repoManager) createWorker(ctx context.Context, originDir, workerDir string) error {
os.RemoveAll(workerDir) // clean up any partial/corrupted previous state
if err := os.MkdirAll(filepath.Dir(workerDir), 0o755); err != nil {
return err
}
return r.git.Clone(ctx, originDir, workerDir, "--local")
}

// pooledWorkspace wraps a workspace and returns its slot to the pool on release.
type pooledWorkspace struct {
workspace.Workspace
pool *workerPool
slot *workerSlot
}

func (pw *pooledWorkspace) Release() error {
pw.pool.avail <- pw.slot
return nil
}

func toShortRemote(remote string) string {
Expand Down
Loading