package gitx import ( "bufio" "container/list" "context" "fmt" "io" "os/exec" "strconv" "strings" "sync" "time" ) // idlePerDir bounds how many readers one repository keeps, so a hot repo does not serialise. const idlePerDir = 4 // idleTotal bounds the whole server, because a barerepo holds more repositories than processes. const idleTotal = 64 // IdleLife is how long an unused reader is kept, and Sweep closes what is older. const IdleLife = 10 * time.Minute // reader is one cat-file --batch kept open, because forking git costs more than the read does. type reader struct { dir string cmd *exec.Cmd in io.WriteCloser out *bufio.Reader idle time.Time spot *list.Element } var ( poolMu sync.Mutex // free holds readers by directory, and order holds the same readers oldest first. free = map[string][]*reader{} order = list.New() ) // take returns an idle reader for dir, or nothing if none is waiting. func take(dir string) *reader { poolMu.Lock() defer poolMu.Unlock() have := free[dir] if len(have) == 0 { return nil } r := have[len(have)-1] free[dir] = have[:len(have)-1] if len(free[dir]) == 0 { delete(free, dir) } order.Remove(r.spot) r.spot = nil return r } // put keeps a reader for the next request, or closes it when the pool is full. func put(r *reader) { poolMu.Lock() if len(free[r.dir]) >= idlePerDir { poolMu.Unlock() r.close() return } r.idle = time.Now() free[r.dir] = append(free[r.dir], r) r.spot = order.PushBack(r) var evict *reader if order.Len() > idleTotal { evict = drop(order.Front()) } poolMu.Unlock() if evict != nil { evict.close() } } // drop removes one reader from both structures, and the caller closes it outside the lock. func drop(e *list.Element) *reader { if e == nil { return nil } r := e.Value.(*reader) order.Remove(e) r.spot = nil have := free[r.dir] for i, other := range have { if other == r { free[r.dir] = append(have[:i], have[i+1:]...) break } } if len(free[r.dir]) == 0 { delete(free, r.dir) } return r } // CloseIdleReaders closes every reader unused for longer than age, and returns how many. func CloseIdleReaders(age time.Duration) int { cutoff := time.Now().Add(-age) var stale []*reader poolMu.Lock() for e := order.Front(); e != nil; { next := e.Next() if e.Value.(*reader).idle.After(cutoff) { break } stale = append(stale, drop(e)) e = next } poolMu.Unlock() for _, r := range stale { r.close() } return len(stale) } // start opens a new reader, which is the only place a cat-file process is created. func start(dir string) (*reader, error) { cmd := exec.Command(Bin, "cat-file", "--batch") cmd.Dir = dir cmd.Env = env() in, err := cmd.StdinPipe() if err != nil { return nil, err } out, err := cmd.StdoutPipe() if err != nil { in.Close() return nil, err } if err := cmd.Start(); err != nil { in.Close() return nil, err } return &reader{dir: dir, cmd: cmd, in: in, out: bufio.NewReaderSize(out, 64<<10)}, nil } func (r *reader) close() { r.in.Close() if r.cmd.Process != nil { r.cmd.Process.Kill() } r.cmd.Wait() } // ask writes every spec and reads every answer, and any surprise means the stream is out of step. func (r *reader) ask(specs []string) (map[string]*Object, error) { // Written from another goroutine, because a big batch fills the pipe before git has answered. sent := make(chan error, 1) go func() { _, err := io.WriteString(r.in, strings.Join(specs, "\n")+"\n") sent <- err }() out := make(map[string]*Object, len(specs)) for _, spec := range specs { header, err := r.out.ReadString('\n') if err != nil { <-sent return nil, err } fields := strings.Fields(header) // " missing" is git's answer, and it rereads the pack directory before saying it. if len(fields) < 3 { continue } size, err := strconv.ParseInt(fields[2], 10, 64) if err != nil || size < 0 { <-sent return nil, fmt.Errorf("cat-file said %q", strings.TrimSpace(header)) } body := make([]byte, size+1) if _, err := io.ReadFull(r.out, body); err != nil { <-sent return nil, err } out[spec] = &Object{SHA: fields[0], Type: fields[1], Size: size, Body: string(body[:size])} } if err := <-sent; err != nil { return nil, err } return out, nil } // Batch reads many objects without starting a process, which is what chapter 25's budget needs. func Batch(ctx context.Context, dir string, specs []string) (map[string]*Object, error) { if len(specs) == 0 { return map[string]*Object{}, nil } if err := ctx.Err(); err != nil { return nil, err } if r := take(dir); r != nil { out, err := r.ask(specs) if err == nil { put(r) return out, nil } r.close() } r, err := start(dir) if err != nil { return nil, err } out, err := r.ask(specs) if err != nil { r.close() return nil, err } put(r) return out, nil }