Files
buildkit/session/upload/uploadprovider/provider.go
Tonis Tiigi 3294c68b99 uploadprovider: allow closing used sources
Helps detecting when uploadprovider has finished
with the source. This avoids deadlock, should the
same reader be shared by multiple uploads that sync
with each other.

Signed-off-by: Tonis Tiigi <tonistiigi@gmail.com>
2024-08-14 13:16:58 +03:00

86 lines
1.6 KiB
Go

package uploadprovider
import (
"io"
"path"
"sync"
"github.com/moby/buildkit/identity"
"github.com/moby/buildkit/session/upload"
"github.com/pkg/errors"
"google.golang.org/grpc"
"google.golang.org/grpc/metadata"
)
func New() *Uploader {
return &Uploader{m: map[string]io.ReadCloser{}}
}
type Uploader struct {
mu sync.Mutex
m map[string]io.ReadCloser
}
func (hp *Uploader) Add(r io.ReadCloser) string {
id := identity.NewID()
hp.m[id] = r
return "http://buildkit-session/" + id
}
func (hp *Uploader) Register(server *grpc.Server) {
upload.RegisterUploadServer(server, hp)
}
func (hp *Uploader) Pull(stream upload.Upload_PullServer) error {
opts, _ := metadata.FromIncomingContext(stream.Context()) // if no metadata continue with empty object
var p string
urls, ok := opts["urlpath"]
if ok && len(urls) > 0 {
p = urls[0]
}
p = path.Base(p)
hp.mu.Lock()
r, ok := hp.m[p]
if !ok {
hp.mu.Unlock()
return errors.Errorf("no http response from session for %s", p)
}
delete(hp.m, p)
hp.mu.Unlock()
_, err := io.Copy(&writer{stream}, r)
if err1 := r.Close(); err == nil {
err = err1
}
return err
}
type writer struct {
grpc.ServerStream
}
func (w *writer) Write(dt []byte) (int, error) {
// avoid sending too big messages on grpc stream
const maxChunkSize = 3 * 1024 * 1024
if len(dt) > maxChunkSize {
n1, err := w.Write(dt[:maxChunkSize])
if err != nil {
return n1, err
}
dt = dt[maxChunkSize:]
var n2 int
if n2, err := w.Write(dt); err != nil {
return n1 + n2, err
}
return n1 + n2, nil
}
if err := w.SendMsg(&upload.BytesMessage{Data: dt}); err != nil {
return 0, err
}
return len(dt), nil
}