diff --git a/pkg/hostman/downloader/downloader.go b/pkg/hostman/downloader/downloader.go new file mode 100644 index 0000000000..ab718c7862 --- /dev/null +++ b/pkg/hostman/downloader/downloader.go @@ -0,0 +1,109 @@ +package downloader + +import ( + "compress/zlib" + "io" + "net/http" + "os" + "time" + + "yunion.io/x/log" +) + +const ( + CHUNK_SIZE = 1024 * 8 + DEFAULT_RATE_LIMIT = 50 + COMPRESS_LEVEL = 1 +) + +type SDownloadProvider struct { + w http.ResponseWriter + rateLimit int + compress bool +} + +func NewDownloadProvider(w http.ResponseWriter, compress bool, rateLimit int) *SDownloadProvider { + if rateLimit <= 0 { + rateLimit = DEFAULT_RATE_LIMIT + } + return &SDownloadProvider{w, rateLimit, compress} +} + +func (d *SDownloadProvider) Start( + prepareDownload func() error, onDownloadComplete func(), + downloadFilePath string, headers http.Header, +) error { + if prepareDownload != nil { + if err := prepareDownload(); err != nil { + log.Errorln(err) + return err + } + } + if headers.Get("Content-Type") == "" { + headers.Set("Content-Type", "application/octet-stream") + } + for k, _ := range headers { + d.w.Header().Add(k, headers.Get(k)) + } + + fi, err := os.Open(downloadFilePath) + if err != nil { + log.Errorln(err) + return err + } + defer fi.Close() + + var ( + end = false + chunk = make([]byte, CHUNK_SIZE) + writer io.Writer = d.w + startTime = time.Now() + sendBytes = 0 + ) + + if d.compress { + zw, err := zlib.NewWriterLevel(d.w, COMPRESS_LEVEL) + if err != nil { + log.Errorln(err) + return err + } + writer = zw + defer zw.Flush() // it's cool + defer zw.Close() + } + + for !end { + if _, err := fi.Read(chunk); err == io.EOF { + end = true + } else if err != nil && err != io.EOF { + log.Errorln(err) + return err + } + + if size, err := writer.Write(chunk); err != nil { + log.Errorln(err) + return err + } else { + sendBytes += size + timeDur := time.Now().Sub(startTime) + exceptDur := float64(sendBytes) / 1000.0 / 1000.0 / float64(d.rateLimit) + if exceptDur > timeDur.Seconds() { + time.Sleep(time.Duration(exceptDur-timeDur.Seconds()) * time.Second) + } + } + } + + // if d.compress { + // zw := writer.(*zlib.Writer) + // zw.Flush() + // } + + sendMb := float64(sendBytes) / 1000.0 / 1000.0 + timeDur := time.Now().Sub(startTime) + log.Infof("Send data: %fMB rate: %fMB/sec", sendMb/timeDur.Seconds()) + + if onDownloadComplete != nil { + onDownloadComplete() + } + return nil +} diff --git a/pkg/hostman/downloader/downloadhandler.go b/pkg/hostman/downloader/downloadhandler.go new file mode 100644 index 0000000000..38261153a1 --- /dev/null +++ b/pkg/hostman/downloader/downloadhandler.go @@ -0,0 +1,183 @@ +package downloader + +import ( + "context" + "fmt" + "net/http" + "time" + + "yunion.io/x/onecloud/pkg/appsrv" + "yunion.io/x/onecloud/pkg/hostman/hostutils" + "yunion.io/x/onecloud/pkg/hostman/options" + "yunion.io/x/onecloud/pkg/hostman/storageman" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/util/fileutils2" +) + +var ( + keyWords = []string{"download"} + streamingWorkerMan *appsrv.SWorkerManager +) + +func init() { + streamingWorkerMan = appsrv.NewWorkerManager("streaming_worker", 20, 1024, false) +} + +func AddDownloadHandler(prefix string, app *appsrv.Application) { + for _, kerword := range keyWords { + hi := app.AddHandler2("GET", fmt.Sprintf("%s/%s//", prefix, kerword), + auth.Authenticate(download), nil, "download", nil) + customizeHandlerInfo(hi) + + hi = app.AddHandler2("GET", fmt.Sprintf("%s/%s/disks//", + prefix, kerword), auth.Authenticate(diskDownload), nil, "disk_download", nil) + customizeHandlerInfo(hi) + + hi = app.AddHandler2("GET", fmt.Sprintf( + "%s/%s/snapshots///", + prefix, kerword), auth.Authenticate(snapshotDownload), + nil, "snapshot_download", nil) + customizeHandlerInfo(hi) + + app.AddHandler("HEAD", fmt.Sprintf("%s/%s/disks//", + prefix, kerword), auth.Authenticate(diskHead)) + app.AddHandler("HEAD", + fmt.Sprintf("%s/%s/snapshots///", + prefix, kerword), auth.Authenticate(snapshotHead)) + } +} + +func customizeHandlerInfo(info *appsrv.SHandlerInfo) { + switch info.GetName(nil) { + case "disk_download", "download", "snapshot_download": + info.SetProcessTimeout(time.Minute * 30).SetWorkerManager(streamingWorkerMan) + } +} + +func isCompress(r *http.Request) bool { + return r.Header.Get("X-Compress-Content") == "zlib" +} + +func download(ctx context.Context, w http.ResponseWriter, r *http.Request) { + var ( + params, _, _ = appsrv.FetchEnv(ctx, w, r) + id = params[""] + action = params[""] + rateLimit = options.HostOptions.BandwidthLimit + compress = isCompress(r) + ) + + switch action { + case "images": + hand := NewImageCacheDownloadProvider(w, compress, rateLimit, id) + if !fileutils2.Exists(hand.downloadFilePath()) { + httperrors.NotFoundError(w, "Image cache %s not found", id) + } else { + hand.Start() + } + case "servers": + hand := NewGuestDownloadProvider(w, compress, rateLimit, id) + if !fileutils2.Exists(hand.fullPath()) { + httperrors.NotFoundError(w, "Guest %s not found", id) + } else { + if err := hand.Start(); err != nil { + hostutils.Response(ctx, w, err) + } + } + default: + hostutils.Response(ctx, w, httperrors.NewNotFoundError("%s Not found", action)) + } +} + +func diskPrecheck( + ctx context.Context, w http.ResponseWriter, r *http.Request, +) (storageman.IDisk, error) { + var ( + params, _, _ = appsrv.FetchEnv(ctx, w, r) + storageId = params[""] + diskId = params[""] + ) + storage := storageman.GetManager().GetStorage(storageId) + if storage == nil { + return nil, httperrors.NewNotFoundError("Storage %s not found", storageId) + } + disk := storage.GetDiskById(diskId) + if disk == nil { + return nil, httperrors.NewNotFoundError("Disk %s not found", diskId) + } + return disk, nil +} + +func diskDownload(ctx context.Context, w http.ResponseWriter, r *http.Request) { + disk, err := diskPrecheck(ctx, w, r) + if err != nil { + hostutils.Response(ctx, w, err) + } else { + var compress = isCompress(r) + hand := NewImageDownloadProvider(w, + compress, options.HostOptions.BandwidthLimit, disk, "") + if err := hand.Start(); err != nil { + hostutils.Response(ctx, w, err) + } + } +} + +func diskHead(ctx context.Context, w http.ResponseWriter, r *http.Request) { + disk, err := diskPrecheck(ctx, w, r) + if err != nil { + hostutils.Response(ctx, w, err) + } else { + var compress = isCompress(r) + hand := NewImageDownloadProvider(w, + compress, options.HostOptions.BandwidthLimit, disk, "") + if err := hand.HandlerHead(); err != nil { + hostutils.Response(ctx, w, err) + } + } +} + +func snapshotPrecheck( + ctx context.Context, w http.ResponseWriter, r *http.Request, +) (string, error) { + var ( + params, _, _ = appsrv.FetchEnv(ctx, w, r) + storageId = params[""] + diskId = params[""] + snapshotId = params["snapshotId"] + ) + + storage := storageman.GetManager().GetStorage(storageId) + if storage == nil { + return "", httperrors.NewNotFoundError("Storage %s not found", storageId) + } + return storage.GetSnapshotPathByIds(diskId, snapshotId), nil +} + +func snapshotDownload(ctx context.Context, w http.ResponseWriter, r *http.Request) { + snapshotPath, err := snapshotPrecheck(ctx, w, r) + if err != nil { + hostutils.Response(ctx, w, err) + } else { + var compress = isCompress(r) + hand := NewSnapshotDownloadProvider(w, + compress, options.HostOptions.BandwidthLimit, snapshotPath) + if err := hand.Start(); err != nil { + hostutils.Response(ctx, w, err) + } + } +} + +func snapshotHead(ctx context.Context, w http.ResponseWriter, r *http.Request) { + snapshotPath, err := snapshotPrecheck(ctx, w, r) + if err != nil { + hostutils.Response(ctx, w, err) + } else { + var compress = isCompress(r) + hand := NewSnapshotDownloadProvider(w, + compress, options.HostOptions.BandwidthLimit, snapshotPath) + if err := hand.HandlerHead(); err != nil { + hostutils.Response(ctx, w, err) + } + } +} diff --git a/pkg/hostman/downloader/guest_downloader.go b/pkg/hostman/downloader/guest_downloader.go new file mode 100644 index 0000000000..ac071bd724 --- /dev/null +++ b/pkg/hostman/downloader/guest_downloader.go @@ -0,0 +1,56 @@ +package downloader + +import ( + "net/http" + "os" + "path" + + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/hostman/options" + "yunion.io/x/onecloud/pkg/util/fileutils2" + "yunion.io/x/onecloud/pkg/util/tarutils" +) + +type SGuestDownloadProvider struct { + *SDownloadProvider + serverId string +} + +func NewGuestDownloadProvider( + w http.ResponseWriter, compress bool, rateLimit int, sid string, +) *SGuestDownloadProvider { + return &SGuestDownloadProvider{ + SDownloadProvider: NewDownloadProvider(w, compress, rateLimit), + serverId: sid, + } +} + +func (s *SGuestDownloadProvider) fullPath() string { + return path.Join(options.HostOptions.ServersPath, s.serverId) +} + +func (s *SGuestDownloadProvider) getHeaders() http.Header { + hdrs := http.Header{} + hdrs.Set("X-Image-Meta-Disk_format", "tar") + return hdrs +} + +func (i *SGuestDownloadProvider) onDownloadComplete() { + if fileutils2.Exists(i.downloadFilePath()) { + os.Remove(i.downloadFilePath()) + } +} + +func (s *SGuestDownloadProvider) downloadFilePath() string { + return s.fullPath() + ".tar" +} + +func (s *SGuestDownloadProvider) prepareDownload() error { + log.Infof("Compress %s to %s", s.fullPath(), s.downloadFilePath()) + return tarutils.TarSparseFile(s.fullPath(), s.downloadFilePath()) +} + +func (s *SGuestDownloadProvider) Start() error { + return s.SDownloadProvider.Start(s.prepareDownload, + s.onDownloadComplete, s.downloadFilePath(), s.getHeaders()) +} diff --git a/pkg/hostman/downloader/image_downloader.go b/pkg/hostman/downloader/image_downloader.go new file mode 100644 index 0000000000..4f2b069072 --- /dev/null +++ b/pkg/hostman/downloader/image_downloader.go @@ -0,0 +1,94 @@ +package downloader + +import ( + "fmt" + "net/http" + "os" + + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/hostman/storageman" + "yunion.io/x/onecloud/pkg/util/fileutils2" + "yunion.io/x/onecloud/pkg/util/qemuimg" + "yunion.io/x/onecloud/pkg/util/tarutils" + "yunion.io/x/pkg/utils" +) + +type SImageDownloadProvider struct { + *SDownloadProvider + disk storageman.IDisk + compressFormat string +} + +func NewImageDownloadProvider(w http.ResponseWriter, compress bool, rateLimit int, disk storageman.IDisk, compressFormat string) *SImageDownloadProvider { + return &SImageDownloadProvider{ + SDownloadProvider: NewDownloadProvider(w, compress, rateLimit), + disk: disk, + compressFormat: compressFormat, + } +} + +func (i *SImageDownloadProvider) fullPath() string { + return i.disk.GetPath() +} + +func (i *SImageDownloadProvider) downloadFilePath() string { + if utils.IsInStringArray(i.compressFormat, []string{"qcow2", "tar"}) { + return i.fullPath() + "." + i.compressFormat + } else { + return i.fullPath() + } +} + +func (i *SImageDownloadProvider) prepareDownload() error { + log.Infof(fmt.Sprintf("Compress %s to %s", i.fullPath(), i.downloadFilePath())) + switch i.compressFormat { + case "qcow2": + img, err := qemuimg.NewQemuImage(i.fullPath()) + if err != nil { + return err + } + _, err = img.CloneQcow2(i.downloadFilePath(), true) + return err + case "tar": + return tarutils.TarSparseFile(i.fullPath(), i.downloadFilePath()) + default: + return nil + } +} + +func (i *SImageDownloadProvider) onDownloadComplete() { + if i.downloadFilePath() != i.fullPath() && fileutils2.Exists(i.downloadFilePath()) { + os.Remove(i.downloadFilePath()) + } +} + +func (i *SImageDownloadProvider) getHeaders() http.Header { + hdrs := http.Header{} + if utils.IsInStringArray(i.compressFormat, []string{"qcow2", "tar"}) { + hdrs.Set("X-Image-Meta-Disk_format", i.compressFormat) + } + return hdrs +} + +func (i *SImageDownloadProvider) Start() error { + return i.SDownloadProvider.Start(i.prepareDownload, i.onDownloadComplete, + i.downloadFilePath(), i.getHeaders()) +} + +func (i *SImageDownloadProvider) HandlerHead() error { + headers := i.getHeaders() + if len(i.compressFormat) > 0 { + headers.Set("X-Image-Meta-Checksum", "error") + } else { + checksum, err := fileutils2.MD5(i.fullPath()) + if err != nil { + return err + } + headers.Set("X-Image-Meta-Checksum", checksum) + } + for k, _ := range headers { + i.w.Header().Add(k, headers.Get(k)) + } + i.w.WriteHeader(200) + return nil +} diff --git a/pkg/hostman/downloader/imagecache_downloader.go b/pkg/hostman/downloader/imagecache_downloader.go new file mode 100644 index 0000000000..5afd3238ce --- /dev/null +++ b/pkg/hostman/downloader/imagecache_downloader.go @@ -0,0 +1,35 @@ +package downloader + +import ( + "net/http" + "path" + + "yunion.io/x/onecloud/pkg/hostman/storageman" +) + +type SImageCacheDownloadProvider struct { + *SDownloadProvider + imageId string +} + +func NewImageCacheDownloadProvider( + w http.ResponseWriter, compress bool, rateLimit int, imageId string, +) *SImageCacheDownloadProvider { + return &SImageCacheDownloadProvider{ + SDownloadProvider: NewDownloadProvider(w, compress, rateLimit), + imageId: imageId, + } +} + +func (s *SImageCacheDownloadProvider) getHeaders() http.Header { + return http.Header{} +} + +func (s *SImageCacheDownloadProvider) downloadFilePath() string { + return path.Join( + storageman.GetManager().LocalStorageImagecacheManager.GetPath(), s.imageId) +} + +func (s *SImageCacheDownloadProvider) Start() error { + return s.SDownloadProvider.Start(nil, nil, s.downloadFilePath(), s.getHeaders()) +} diff --git a/pkg/hostman/downloader/snapshot_downloader.go b/pkg/hostman/downloader/snapshot_downloader.go new file mode 100644 index 0000000000..30409bb27c --- /dev/null +++ b/pkg/hostman/downloader/snapshot_downloader.go @@ -0,0 +1,48 @@ +package downloader + +import ( + "net/http" + + "yunion.io/x/onecloud/pkg/util/fileutils2" +) + +type SSnapshotDownloadProvider struct { + *SDownloadProvider + snapshotPath string +} + +func NewSnapshotDownloadProvider( + w http.ResponseWriter, compress bool, rateLimit int, snapshotPath string, +) *SSnapshotDownloadProvider { + return &SSnapshotDownloadProvider{ + SDownloadProvider: NewDownloadProvider(w, compress, rateLimit), + snapshotPath: snapshotPath, + } +} + +func (s *SSnapshotDownloadProvider) getHeaders() http.Header { + hdrs := http.Header{} + hdrs.Set("X-Image-Meta-Disk_format", "") + return hdrs +} + +func (s *SSnapshotDownloadProvider) HandlerHead() error { + headers := s.getHeaders() + if fileutils2.Exists(s.snapshotPath) { + chksum, err := fileutils2.MD5(s.snapshotPath) + if err != nil { + return err + } + headers.Set("X-Image-Meta-Checksum", chksum) + } + s.w.WriteHeader(200) + return nil +} + +func (s *SSnapshotDownloadProvider) downloadFilePath() string { + return s.snapshotPath +} + +func (s *SSnapshotDownloadProvider) Start() error { + return s.SDownloadProvider.Start(nil, nil, s.downloadFilePath(), s.getHeaders()) +} diff --git a/pkg/hostman/downloadhandler/downloadhandler.go b/pkg/hostman/downloadhandler/downloadhandler.go deleted file mode 100644 index c13d5f3dc0..0000000000 --- a/pkg/hostman/downloadhandler/downloadhandler.go +++ /dev/null @@ -1,87 +0,0 @@ -package downloadhandler - -import ( - "context" - "fmt" - "net/http" - "time" - - "yunion.io/x/onecloud/pkg/appsrv" - "yunion.io/x/onecloud/pkg/hostman/hostutils" - "yunion.io/x/onecloud/pkg/httperrors" - "yunion.io/x/onecloud/pkg/mcclient/auth" -) - -var ( - keyWords = []string{"download"} - streamingWorkerMan *appsrv.SWorkerManager -) - -func init() { - streamingWorkerMan = appsrv.NewWorkerManager("streaming_worker", 20, 1024, false) -} - -func AddDownloadHandler(prefix string, app *appsrv.Application) { - for _, kerword := range keyWords { - hi := app.AddHandler2("GET", fmt.Sprintf("%s/%s//", prefix, kerword), - auth.Authenticate(download), nil, "download", nil) - customizeHandlerInfo(hi) - - hi = app.AddHandler2("GET", fmt.Sprintf("%s/%s/disks//", - prefix, kerword), auth.Authenticate(diskDownload), nil, "disk_download", nil) - customizeHandlerInfo(hi) - - hi = app.AddHandler2("GET", fmt.Sprintf( - "%s/%s/snapshots///", - prefix, kerword), auth.Authenticate(snapshotDownload), - nil, "snapshot_download", nil) - customizeHandlerInfo(hi) - - app.AddHandler("HEAD", fmt.Sprintf("%s/%s/disks//", - prefix, kerword), auth.Authenticate(diskHead)) - app.AddHandler("HEAD", - fmt.Sprintf("%s/%s/snapshots///", - prefix, kerword), auth.Authenticate(snapshotHead)) - } -} - -func customizeHandlerInfo(info *appsrv.SHandlerInfo) { - switch info.GetName(nil) { - case "disk_download", "download", "snapshot_download": - info.SetProcessTimeout(time.Minute * 30).SetWorkerManager(streamingWorkerMan) - } -} - -func download(ctx context.Context, w http.ResponseWriter, r *http.Request) { - var ( - params, _, _ = appsrv.FetchEnv(ctx, w, r) - sid = params[""] - action = params[""] - zlib = r.Header.Get("X-Compress-Content") - compress bool - ) - - if zlib == "zlib" { - compress = true - } - - switch action { - case "images": - // ImagecacheDownloadProvider() - case "servers": - default: - hostutils.Response(ctx, w, httperrors.NewNotFoundError("%s Not found", action)) - } -} - -func diskDownload(ctx context.Context, w http.ResponseWriter, r *http.Request) { -} - -func snapshotDownload(ctx context.Context, w http.ResponseWriter, r *http.Request) { -} - -func diskHead(ctx context.Context, w http.ResponseWriter, r *http.Request) { -} - -func snapshotHead(ctx context.Context, w http.ResponseWriter, r *http.Request) { -} diff --git a/pkg/hostman/options/options.go b/pkg/hostman/options/options.go index 1899149c1c..6e43f41aab 100644 --- a/pkg/hostman/options/options.go +++ b/pkg/hostman/options/options.go @@ -78,6 +78,7 @@ type SHostOptions struct { ManageNtpConfiguration bool `default:"true"` LogSystemdUnits []string `help:"Systemd units log collected by fluent-bit"` BandwidthLimit int `default:"50" help:"Bandwidth upper bound when migrating disk image in MB/sec"` + SnapshotDirSuffix string `help:"Snapshot dir name equal diskId concat snapshot dir suffix" default:"_snap"` } var HostOptions SHostOptions diff --git a/pkg/hostman/storageman/disklocal.go b/pkg/hostman/storageman/disklocal.go index 7d16451a48..48b468c16a 100644 --- a/pkg/hostman/storageman/disklocal.go +++ b/pkg/hostman/storageman/disklocal.go @@ -55,7 +55,7 @@ func (d *SLocalDisk) GetPath() string { } func (d *SLocalDisk) GetSnapshotDir() string { - return path.Join(d.Storage.GetSnapshotDir(), d.Id+"_snap") + return path.Join(d.Storage.GetSnapshotDir(), d.Id+options.HostOptions.SnapshotDirSuffix) } func (d *SLocalDisk) Probe() error { diff --git a/pkg/hostman/storageman/storagebase.go b/pkg/hostman/storageman/storagebase.go index a9524b5698..be55226e3d 100644 --- a/pkg/hostman/storageman/storagebase.go +++ b/pkg/hostman/storageman/storagebase.go @@ -27,6 +27,7 @@ type IStorage interface { SetPath(string) GetPath() string GetSnapshotDir() string + GetSnapshotPathByIds(diskId, snapshotId string) string GetFreeSizeMb() int GetCapacity() int diff --git a/pkg/hostman/storageman/storagelocal.go b/pkg/hostman/storageman/storagelocal.go index 0740ad9b19..478747b111 100644 --- a/pkg/hostman/storageman/storagelocal.go +++ b/pkg/hostman/storageman/storagelocal.go @@ -54,6 +54,11 @@ func (s *SLocalStorage) GetSnapshotDir() string { return path.Join(s.Path, _SNAPSHOT_PATH_) } +func (s *SLocalStorage) GetSnapshotPathByIds(diskId, snapshotId string) string { + return path.Join(s.GetSnapshotDir(), + diskId+options.HostOptions.SnapshotDirSuffix, snapshotId) +} + func (s *SLocalStorage) SyncStorageInfo() (jsonutils.JSONObject, error) { content := jsonutils.NewDict() content.Set("name", jsonutils.NewString(s.StorageName)) diff --git a/pkg/util/tarutils/tarutils.go b/pkg/util/tarutils/tarutils.go new file mode 100644 index 0000000000..67da0ada42 --- /dev/null +++ b/pkg/util/tarutils/tarutils.go @@ -0,0 +1,25 @@ +package tarutils + +import ( + "os" + "path/filepath" + + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/util/procutils" +) + +func TarSparseFile(origin, tar string) error { + origin, _ = filepath.Abs(origin) + tar, _ = filepath.Abs(tar) + workDir := filepath.Dir(origin) + originFile := filepath.Base(origin) + if err := os.Chdir(workDir); err != nil { + log.Errorln(err) + return err + } + _, err := procutils.NewCommand("tar", "-Scf", tar, originFile).Run() + if err != nil { + log.Errorln("Tar sparse file error: %s", err) + } + return nil +}