From d3da407cd04c97ffad75a3595798e035c081b368 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Wed, 16 Jan 2019 23:24:41 +0800 Subject: [PATCH 1/3] torrent recode --- cmd/torrent/main.go | 48 ++++++++++++++++++++++++++++++++++++++------- 1 file changed, 41 insertions(+), 7 deletions(-) diff --git a/cmd/torrent/main.go b/cmd/torrent/main.go index a31539e76a..9cb5bc4b5f 100644 --- a/cmd/torrent/main.go +++ b/cmd/torrent/main.go @@ -4,7 +4,6 @@ import ( "fmt" "os" "os/signal" - "path" "path/filepath" "syscall" "time" @@ -16,6 +15,9 @@ import ( "yunion.io/x/pkg/util/version" "yunion.io/x/structarg" + "github.com/anacrolix/torrent/storage" + "io/ioutil" + "net/http" "yunion.io/x/onecloud/pkg/util/nodeid" "yunion.io/x/onecloud/pkg/util/torrentutils" ) @@ -29,6 +31,8 @@ type Options struct { Tracker []string `help:"Tracker urls, e.g. http://10.168.222.252:6969/announce or udp://tracker.istole.it:6969"` Debug bool `help:"turn on debug"` + + CallbackURL string `help:"callback notification URL"` } func exitSignalHandlers(client *torrent.Client) { @@ -71,6 +75,7 @@ func main() { } var mi *metainfo.MetaInfo + var rootDir string if len(options.Tracker) > 0 { // server mode @@ -78,6 +83,7 @@ func main() { if err != nil { log.Fatalf("fail to save torrent file %s", err) } + rootDir = filepath.Dir(root) } else { // client mode, load mi from torrent file @@ -85,6 +91,7 @@ func main() { if err != nil { log.Fatalf("fail to open torrent file %s", err) } + rootDir = root } nodeId, err := nodeid.GetNodeId() @@ -100,13 +107,25 @@ func main() { clientConfig.Debug = options.Debug clientConfig.Seed = true clientConfig.NoUpload = false - if len(options.Tracker) > 0 { - // server mode - clientConfig.DataDir = path.Dir(root) - } else { - // client mode - clientConfig.DataDir = root + + info, err := mi.UnmarshalInfo() + if err != nil { + log.Errorf("fail to unmarshalinfo %s", err) + return } + log.Infof("To download file %s", info.Name) + tmpDir := filepath.Join(rootDir, fmt.Sprintf("%s%s", info.Name, ".tmp")) + + os.RemoveAll(tmpDir) + os.MkdirAll(tmpDir, 0755) + defer os.RemoveAll(tmpDir) + + clientConfig.DefaultStorage = storage.NewFileWithCustomPathMaker(tmpDir, + func(baseDir string, info *metainfo.Info, infoHash metainfo.Hash) string { + return filepath.Dir(baseDir) + }, + ) + clientConfig.DisableTrackers = false clientConfig.DisablePEX = true clientConfig.NoDHT = true @@ -149,6 +168,21 @@ func main() { for { if t.BytesCompleted() == t.Info().TotalLength() { fmt.Printf("\rSeeding.............") + if len(options.CallbackURL) > 0 { + for tried := 0; tried < 10; tried += 1 { + resp, err := http.Get(options.CallbackURL) + if err == nil && resp.StatusCode < 300 { + break + } + if err != nil { + log.Errorf("callback fail %s", err) + } else { + respBody, _ := ioutil.ReadAll(resp.Body) + log.Errorf("callback response error %s", string(respBody)) + } + time.Sleep(time.Duration(tried+1) * 10 * time.Second) + } + } } else { fmt.Printf("\rDownload: %.1f%%", float64(t.BytesCompleted())*100.0/float64(t.Info().TotalLength())) } From daba249249bae6990156ad48e76d6e6c4a6014c6 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Thu, 17 Jan 2019 00:08:57 +0800 Subject: [PATCH 2/3] minor updates --- cmd/torrent/main.go | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/cmd/torrent/main.go b/cmd/torrent/main.go index 9cb5bc4b5f..e573be3056 100644 --- a/cmd/torrent/main.go +++ b/cmd/torrent/main.go @@ -146,26 +146,19 @@ func main() { log.Fatalf("%s", err) } - go func() { - <-t.GotInfo() + <-t.GotInfo() + t.DownloadAll() - files := t.Info().Files - log.Debugf("Got Info, start download %d files", len(files)) - for i := 0; i < len(files); i += 1 { - log.Debugf("%d: %s", i, files[i].Path) - } - - t.DownloadAll() - }() + stop := false go func() { <-client.Closed() log.Debugf("client closed, exit!") - os.Exit(0) + stop = true }() - for { + for !stop { if t.BytesCompleted() == t.Info().TotalLength() { fmt.Printf("\rSeeding.............") if len(options.CallbackURL) > 0 { From 0e3f2e160c653c2e3e684b26d75a2fb861fb36be Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Wed, 30 Jan 2019 21:06:48 +0800 Subject: [PATCH 3/3] temp commit --- cmd/torrent/main.go | 75 +++++++++------ pkg/appsrv/appsrv.go | 45 +++++++-- pkg/cloudcommon/app.go | 6 +- pkg/cloudcommon/cronman/cronman.go | 4 +- pkg/cloudcommon/policy/defaults.go | 7 ++ pkg/cloudcommon/policy/policy.go | 9 +- pkg/cloudir/service/service.go | 5 +- pkg/image/models/image_subs.go | 43 ++++----- pkg/image/models/images.go | 21 +++- pkg/image/options/options.go | 2 + pkg/image/service/service.go | 30 +++--- pkg/image/torrent/torrent.go | 149 ++++++++++++++--------------- pkg/mcclient/auth/middleware.go | 28 ++++-- pkg/util/nodeid/nodeid.go | 8 +- pkg/util/sysutils/doc.go | 1 + pkg/util/sysutils/sysutils.go | 23 +++++ pkg/yunionconf/service/service.go | 5 +- 17 files changed, 284 insertions(+), 177 deletions(-) create mode 100644 pkg/util/sysutils/doc.go create mode 100644 pkg/util/sysutils/sysutils.go diff --git a/cmd/torrent/main.go b/cmd/torrent/main.go index e573be3056..af4d462e43 100644 --- a/cmd/torrent/main.go +++ b/cmd/torrent/main.go @@ -1,7 +1,10 @@ package main import ( + "crypto/sha1" "fmt" + "io/ioutil" + "net/http" "os" "os/signal" "path/filepath" @@ -10,14 +13,13 @@ import ( "github.com/anacrolix/torrent" "github.com/anacrolix/torrent/metainfo" + "github.com/anacrolix/torrent/storage" "yunion.io/x/log" "yunion.io/x/pkg/util/version" "yunion.io/x/structarg" - "github.com/anacrolix/torrent/storage" - "io/ioutil" - "net/http" + "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/nodeid" "yunion.io/x/onecloud/pkg/util/torrentutils" ) @@ -30,7 +32,8 @@ type Options struct { Tracker []string `help:"Tracker urls, e.g. http://10.168.222.252:6969/announce or udp://tracker.istole.it:6969"` - Debug bool `help:"turn on debug"` + Debug bool `help:"turn on debug" default:"false"` + Verbose bool `help:"verbose mode" default:"false"` CallbackURL string `help:"callback notification URL"` } @@ -77,7 +80,7 @@ func main() { var mi *metainfo.MetaInfo var rootDir string - if len(options.Tracker) > 0 { + if len(options.Tracker) > 0 && !fileutils2.Exists(options.TORRENT) { // server mode mi, err = torrentutils.GenerateTorrent(root, options.Tracker, options.TORRENT) if err != nil { @@ -94,30 +97,38 @@ func main() { rootDir = root } + info, err := mi.UnmarshalInfo() + if err != nil { + log.Errorf("fail to unmarshalinfo %s", err) + return + } + + hasher := sha1.New() + nodeId, err := nodeid.GetNodeId() if err != nil { log.Errorf("fail to generate node id: %s", err) return } - log.Infof("Set torrent server as node %s", nodeId) + hasher.Write(nodeId) + hasher.Write(info.Pieces) + + peerIdStr := fmt.Sprintf("%x", hasher.Sum(nil)) + + log.Infof("Set torrent server as node %s", peerIdStr[:20]) clientConfig := torrent.NewDefaultClientConfig() - clientConfig.PeerID = nodeId[:20] + clientConfig.PeerID = peerIdStr[:20] clientConfig.Debug = options.Debug clientConfig.Seed = true clientConfig.NoUpload = false - info, err := mi.UnmarshalInfo() - if err != nil { - log.Errorf("fail to unmarshalinfo %s", err) - return - } - log.Infof("To download file %s", info.Name) + log.Infof("To sync torrent files for %s", info.Name) tmpDir := filepath.Join(rootDir, fmt.Sprintf("%s%s", info.Name, ".tmp")) os.RemoveAll(tmpDir) - os.MkdirAll(tmpDir, 0755) + os.MkdirAll(tmpDir, 0700) defer os.RemoveAll(tmpDir) clientConfig.DefaultStorage = storage.NewFileWithCustomPathMaker(tmpDir, @@ -141,6 +152,8 @@ func main() { go exitSignalHandlers(client) + start := time.Now() + t, err := client.AddTorrent(mi) if err != nil { log.Fatalf("%s", err) @@ -158,24 +171,32 @@ func main() { stop = true }() + finish := false + for !stop { if t.BytesCompleted() == t.Info().TotalLength() { - fmt.Printf("\rSeeding.............") - if len(options.CallbackURL) > 0 { - for tried := 0; tried < 10; tried += 1 { - resp, err := http.Get(options.CallbackURL) - if err == nil && resp.StatusCode < 300 { - break + if !finish { + finish = true + fmt.Printf("Download complete, takes %d seconds\n", time.Now().Sub(start)/time.Second) + if len(options.CallbackURL) > 0 { + maxTried := 10 + for tried := 0; tried < maxTried; tried += 1 { + resp, err := http.Post(options.CallbackURL, "", nil) + if err == nil && resp.StatusCode < 300 { + break + } + if err != nil { + log.Errorf("callback fail %s", err) + } else { + defer resp.Body.Close() + respBody, _ := ioutil.ReadAll(resp.Body) + log.Errorf("callback response error %s", string(respBody)) + } + time.Sleep(time.Duration(tried+1) * 10 * time.Second) } - if err != nil { - log.Errorf("callback fail %s", err) - } else { - respBody, _ := ioutil.ReadAll(resp.Body) - log.Errorf("callback response error %s", string(respBody)) - } - time.Sleep(time.Duration(tried+1) * 10 * time.Second) } } + fmt.Printf("\rSeeding.............") } else { fmt.Printf("\rDownload: %.1f%%", float64(t.BytesCompleted())*100.0/float64(t.Info().TotalLength())) } diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index 1a1b9ada56..8e9618c0e3 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -5,8 +5,11 @@ import ( "fmt" "math/rand" "net/http" + "os" + "os/signal" "strings" "sync" + "syscall" "time" "yunion.io/x/jsonutils" @@ -36,6 +39,8 @@ type Application struct { defHandlerInfo SHandlerInfo cors *Cors middlewares []MiddlewareFunc + + idleConnsClosed chan struct{} } const ( @@ -325,20 +330,44 @@ func (app *Application) initServer(addr string) *http.Server { return s } -func (app *Application) ListenAndServe(addr string) { - s := app.initServer(addr) - err := s.ListenAndServe() - if err != nil { - log.Fatalf("ListAndServer fail: %s", err) - } +func (app *Application) registerCleanShutdown(s *http.Server, onStop func()) { + app.idleConnsClosed = make(chan struct{}) + go func() { + c := make(chan os.Signal, 1) + signal.Notify(c, syscall.SIGINT, syscall.SIGTERM) + log.Infof("Close signal received: %+v", <-c) + if err := s.Shutdown(context.Background()); err != nil { + // Error from closing listeners, or context timeout: + log.Errorf("HTTP server Shutdown: %v", err) + } + onStop() + + close(app.idleConnsClosed) + }() } -func (app *Application) ListenAndServeTLS(addr string, certFile, keyFile string) { +func (app *Application) waitCleanShutdown() { + <-app.idleConnsClosed +} + +func (app *Application) ListenAndServe(addr string, onStop func()) { s := app.initServer(addr) - err := s.ListenAndServeTLS(certFile, keyFile) + app.registerCleanShutdown(s, onStop) + err := s.ListenAndServe() if err != nil && err != http.ErrServerClosed { log.Fatalf("ListAndServer fail: %s", err) } + app.waitCleanShutdown() +} + +func (app *Application) ListenAndServeTLS(addr string, certFile, keyFile string, onStop func()) { + s := app.initServer(addr) + app.registerCleanShutdown(s, onStop) + err := s.ListenAndServeTLS(certFile, keyFile) + if err != nil && err != http.ErrServerClosed { + log.Fatalf("ListAndServerTLS fail: %s", err) + } + app.waitCleanShutdown() } func isJsonContentType(r *http.Request) bool { diff --git a/pkg/cloudcommon/app.go b/pkg/cloudcommon/app.go index 6234305c93..6e53546e49 100644 --- a/pkg/cloudcommon/app.go +++ b/pkg/cloudcommon/app.go @@ -23,7 +23,7 @@ func InitApp(options *CommonOptions, dbAccess bool) *appsrv.Application { return app } -func ServeForever(app *appsrv.Application, options *CommonOptions) { +func ServeForever(app *appsrv.Application, options *CommonOptions, onStop func()) { AppDBInit(app) addr := net.JoinHostPort(options.Address, strconv.Itoa(options.Port)) proto := "http" @@ -47,8 +47,8 @@ func ServeForever(app *appsrv.Application, options *CommonOptions) { if len(options.SslKeyfile) == 0 { log.Fatalf("Missing ssl-keyfile") } - app.ListenAndServeTLS(addr, certfile, options.SslKeyfile) + app.ListenAndServeTLS(addr, certfile, options.SslKeyfile, onStop) } else { - app.ListenAndServe(addr) + app.ListenAndServe(addr, onStop) } } diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index 9946a03a99..d2bfe0aa29 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -148,7 +148,9 @@ func (self *SCronJobManager) Start() { } func (self *SCronJobManager) Stop() { - close(self.stop) + if self.stop != nil { + close(self.stop) + } } func (self *SCronJobManager) run() { diff --git a/pkg/cloudcommon/policy/defaults.go b/pkg/cloudcommon/policy/defaults.go index 1d4a96f99f..3257705881 100644 --- a/pkg/cloudcommon/policy/defaults.go +++ b/pkg/cloudcommon/policy/defaults.go @@ -125,5 +125,12 @@ var ( Action: PolicyActionGet, Result: rbacutils.OwnerAllow, }, + { + Service: "image", + Resource: "images", + Action: PolicyActionPerform, + Extra: []string{"update-torrent-status"}, + Result: rbacutils.GuestAllow, + }, } ) diff --git a/pkg/cloudcommon/policy/policy.go b/pkg/cloudcommon/policy/policy.go index d6f2ad2e97..bafcf379f5 100644 --- a/pkg/cloudcommon/policy/policy.go +++ b/pkg/cloudcommon/policy/policy.go @@ -190,7 +190,7 @@ func queryKey(isAdmin bool, userCred mcclient.TokenCredential, service string, r } func (manager *SPolicyManager) Allow(isAdmin bool, userCred mcclient.TokenCredential, service string, resource string, action string, extra ...string) rbacutils.TRbacResult { - if manager.cache != nil { + if manager.cache != nil && userCred != nil { key := queryKey(isAdmin, userCred, service, resource, action, extra...) val := manager.cache.Get(key) if val != nil { @@ -231,7 +231,12 @@ func (manager *SPolicyManager) allowWithoutCache(isAdmin bool, userCred mcclient log.Warningf("no policies fetched") return rbacutils.Deny } - userCredJson := userCred.ToJson() + var userCredJson jsonutils.JSONObject + if userCred != nil { + userCredJson = userCred.ToJson() + } else { + userCredJson = jsonutils.NewDict() + } currentPriv := rbacutils.Deny for _, p := range policies { result := p.Allow(userCredJson, service, resource, action, extra...) diff --git a/pkg/cloudir/service/service.go b/pkg/cloudir/service/service.go index 9a3c2799a7..52d3423c03 100644 --- a/pkg/cloudir/service/service.go +++ b/pkg/cloudir/service/service.go @@ -23,11 +23,12 @@ func StartService() { log.Fatalf("init etcd fail: %s", err) return } - defer etcd.CloseDefaultEtcdClient() app := cloudcommon.InitApp(commonOpts, false) initHandlers(app) - cloudcommon.ServeForever(app, commonOpts) + cloudcommon.ServeForever(app, commonOpts, func() { + etcd.CloseDefaultEtcdClient() + }) } diff --git a/pkg/image/models/image_subs.go b/pkg/image/models/image_subs.go index c972b9740d..67ee947f6f 100644 --- a/pkg/image/models/image_subs.go +++ b/pkg/image/models/image_subs.go @@ -6,8 +6,6 @@ import ( "os" "path/filepath" - t "github.com/anacrolix/torrent" - "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -89,10 +87,12 @@ func (self *SImageSubformat) DoConvert(image *SImage) error { log.Errorf("fail to convert image torrent %s", err) return err } - err = self.seedTorrent() - if err != nil { - log.Errorf("fail to seed torrent %s", err) - return err + if options.Options.EnableTorrentService { + err = self.seedTorrent(image.Id) + if err != nil { + log.Errorf("fail to seed torrent %s", err) + return err + } } // log.Infof("Start seeding...") return nil @@ -206,10 +206,10 @@ func (self *SImageSubformat) getLocalTorrentLocation() string { return "" } -func (self *SImageSubformat) seedTorrent() error { +func (self *SImageSubformat) seedTorrent(imageId string) error { file := self.getLocalTorrentLocation() log.Debugf("add torrent %s to seed...", file) - return torrent.AddTorrent(file) + return torrent.SeedTorrent(file, imageId, self.Format) } func (self *SImageSubformat) StopTorrent() { @@ -237,14 +237,6 @@ func (self *SImageSubformat) RemoveFiles() error { return nil } -func (self *SImageSubformat) getTorrent() *t.Torrent { - torrentPath := self.getLocalTorrentLocation() - if len(torrentPath) > 0 { - return torrent.GetTorrent(torrentPath) - } - return nil -} - type SImageSubformatDetails struct { Format string @@ -258,7 +250,6 @@ type SImageSubformatDetails struct { TorrentStatus string TorrentSeeding bool - TorrentStats t.TorrentStats } func (self *SImageSubformat) GetDetails() SImageSubformatDetails { @@ -273,14 +264,11 @@ func (self *SImageSubformat) GetDetails() SImageSubformatDetails { details.TorrentChecksum = self.TorrentChecksum details.TorrentStatus = self.TorrentStatus - t := self.getTorrent() - - if t != nil { - if t.BytesMissing() == 0 { - details.TorrentSeeding = true - } - details.TorrentStats = t.Stats() + filePath := self.getLocalTorrentLocation() + if len(filePath) > 0 { + details.TorrentSeeding = torrent.GetTorrentSeeding(filePath) } + return details } @@ -342,3 +330,10 @@ func (self *SImageSubformat) checkStatus(useFast bool) { } } } + +func (self *SImageSubformat) SetStatusSeeding(seeding bool) { + filePath := self.getLocalTorrentLocation() + if len(filePath) > 0 { + torrent.SetTorrentSeeding(filePath, seeding) + } +} diff --git a/pkg/image/models/images.go b/pkg/image/models/images.go index 3b70404245..11d91c307f 100644 --- a/pkg/image/models/images.go +++ b/pkg/image/models/images.go @@ -877,7 +877,7 @@ func (self *SImage) StopTorrents() { func (self *SImage) seedTorrents() { subimgs := ImageSubformatManager.GetAllSubImages(self.Id) for i := 0; i < len(subimgs); i += 1 { - subimgs[i].seedTorrent() + subimgs[i].seedTorrent(self.Id) } } @@ -1024,7 +1024,7 @@ func (self *SImage) DoCheckStatus(ctx context.Context, userCred mcclient.TokenCr if needConvert { log.Infof("Image %s is active and need convert", self.Name) self.StartImageConvertTask(ctx, userCred, "") - } else { + } else if options.Options.EnableTorrentService { self.seedTorrents() } } @@ -1059,3 +1059,20 @@ func (self *SImage) PerformMarkPublicProtected( } return nil, nil } + +func (self *SImage) AllowPerformUpdateTorrentStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return true +} + +func (self *SImage) PerformUpdateTorrentStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + formatStr, _ := query.GetString("format") + if len(formatStr) == 0 { + return nil, httperrors.NewInputParameterError("missing parameter format") + } + subimg := ImageSubformatManager.FetchSubImage(self.Id, formatStr) + if subimg == nil { + return nil, httperrors.NewResourceNotFoundError("format %s not found", formatStr) + } + subimg.SetStatusSeeding(true) + return nil, nil +} diff --git a/pkg/image/options/options.go b/pkg/image/options/options.go index 7cf2ea593f..2a3a7e76ee 100644 --- a/pkg/image/options/options.go +++ b/pkg/image/options/options.go @@ -23,6 +23,8 @@ type SImageOptions struct { EnableTorrentService bool `help:"Enable torrent service" default:"false"` TargetImageFormats []string `help:"target image formats that the system will automatically convert to" default:"qcow2,vmdk,vhd"` + + TorrentClientPath string `help:"path to torrent executable" default:"/opt/yunion/bin/torrent"` } var ( diff --git a/pkg/image/service/service.go b/pkg/image/service/service.go index 8d6485242b..2daa6d0675 100644 --- a/pkg/image/service/service.go +++ b/pkg/image/service/service.go @@ -62,35 +62,37 @@ func StartService() { log.Infof("Auth complete!!") }) + trackers := torrent.GetTrackers() + if len(trackers) == 0 { + log.Errorf("no valid torrent-tracker") + return + } + cloudcommon.InitDB(dbOpts) - defer cloudcommon.CloseDB() app := cloudcommon.InitApp(commonOpts, true) initHandlers(app) - if opts.EnableTorrentService { - err := torrent.InitTorrentClient() - if err != nil { - log.Errorf("fail to initialize torrent client: %s", err) - return - } - torrent.InitTorrentHandler(app) - defer torrent.CloseTorrentClient() - } - if !db.CheckSync(opts.AutoSyncTable) { log.Fatalf("database schema not in sync!") } models.InitDB() - models.CheckImages() + go models.CheckImages() cron := cronman.GetCronJobManager() cron.AddJob1("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages) cron.Start() - defer cron.Stop() - cloudcommon.ServeForever(app, commonOpts) + cloudcommon.ServeForever(app, commonOpts, func() { + cloudcommon.CloseDB() + + cron.Stop() + + if options.Options.EnableTorrentService { + torrent.StopTorrents() + } + }) } diff --git a/pkg/image/torrent/torrent.go b/pkg/image/torrent/torrent.go index d41a3ccf1a..8cecbdb3b0 100644 --- a/pkg/image/torrent/torrent.go +++ b/pkg/image/torrent/torrent.go @@ -1,28 +1,52 @@ package torrent import ( - "context" "fmt" - "net/http" - - "github.com/anacrolix/torrent" - "github.com/anacrolix/torrent/metainfo" + "os" + "time" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/image/options" "yunion.io/x/onecloud/pkg/mcclient/auth" - "yunion.io/x/onecloud/pkg/util/nodeid" + "yunion.io/x/onecloud/pkg/util/sysutils" ) +type STorrentProcessState struct { + process *os.Process + seeding bool +} + var ( - torrentClient *torrent.Client - torrentTable = make(map[string]*torrent.Torrent) + torrentTable = make(map[string]*STorrentProcessState) + seedTaskWorkerMan *appsrv.SWorkerManager ) +const ( + TORRENT_TRACKER_SERVICE = "torrent-tracker" +) + +func init() { + seedTaskWorkerMan = appsrv.NewWorkerManager("seedTaskWorkerManager", 1, 1024, false) +} + +func (stat *STorrentProcessState) StopAndWait() error { + err := stat.process.Kill() + if err != nil { + log.Errorf("kill error %s", err) + return err + } + _, err = stat.process.Wait() + if err != nil { + log.Errorf("wait error %s", err) + return err + } + return nil +} + func GetTrackers() []string { - urls, err := auth.GetServiceURLs("torrent-tracker", options.Options.Region, "", "") + urls, err := auth.GetServiceURLs(TORRENT_TRACKER_SERVICE, options.Options.Region, "", "") if err != nil { log.Errorf("fail to get torrent-tracker") return nil @@ -30,93 +54,62 @@ func GetTrackers() []string { return urls } -func InitTorrentClient() error { - urls := GetTrackers() - if len(urls) == 0 { - log.Errorf("no valid torrent-tracker") - return fmt.Errorf("no valid torrent-tracker") - } - - nodeId, err := nodeid.GetNodeId() - if err != nil { - log.Errorf("fail to generate node id: %s", err) - return err - } - - log.Infof("Set torrent server as node %s", nodeId) - - clientConfig := torrent.NewDefaultClientConfig() - clientConfig.PeerID = nodeId[:20] - clientConfig.Debug = false - clientConfig.Seed = true - clientConfig.NoUpload = false - clientConfig.DataDir = options.Options.FilesystemStoreDatadir - clientConfig.DisableTrackers = false - clientConfig.DisablePEX = true - clientConfig.NoDHT = true - - client, err := torrent.NewClient(clientConfig) - if err != nil { - log.Errorf("error creating client: %s", err) - return err - } - torrentClient = client - - log.Infof("torrent client initialized") - +func SeedTorrent(torrentpath string, imageId, format string) error { + seedTaskWorkerMan.Run(func() { + log.Infof("Start seed %s ...", torrentpath) + err := seedTorrent(torrentpath, imageId, format) + if err == nil { + time.Sleep(10 * time.Second) + } + }, nil, nil) return nil } -func InitTorrentHandler(app *appsrv.Application) { - app.AddDefaultHandler("GET", "/torrent_stats", TorrentStatsHandler, "torrent_stats") -} - -func TorrentStatsHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { - if torrentClient != nil { - torrentClient.WriteStatus(w) - } -} - -func CloseTorrentClient() { - if torrentClient != nil { - torrentClient.Close() - } -} - -func AddTorrent(filepath string) error { - if torrentClient == nil { - return nil - } - - mi, err := metainfo.LoadFromFile(filepath) +func seedTorrent(torrentpath string, imageId, format string) error { + url, err := auth.GetServiceURL("image", options.Options.Region, "", "public") if err != nil { - log.Errorf("fail to open torrent file %s", err) return err } - t, err := torrentClient.AddTorrent(mi) + args := []string{ + options.Options.TorrentClientPath, + options.Options.FilesystemStoreDatadir, + torrentpath, + "--callback-url", + fmt.Sprintf("%s/images/%s/update-torrent-status?format=%s", url, imageId, format), + } + proc, err := sysutils.Start(false, args...) if err != nil { - log.Errorf("AddTorrent fail %s", err) return err } - - torrentTable[filepath] = t - - <-t.GotInfo() - t.DownloadAll() - + torrentTable[torrentpath] = &STorrentProcessState{ + process: proc, + seeding: false, + } return nil } -func GetTorrent(filepath string) *torrent.Torrent { +func SetTorrentSeeding(filepath string, seeding bool) { + if _, ok := torrentTable[filepath]; ok { + torrentTable[filepath].seeding = seeding + } +} + +func GetTorrentSeeding(filepath string) bool { if t, ok := torrentTable[filepath]; ok { - return t + return t.seeding } - return nil + return false } func RemoveTorrent(filepath string) { if t, ok := torrentTable[filepath]; ok { - t.Drop() + t.StopAndWait() delete(torrentTable, filepath) } } + +func StopTorrents() { + for k := range torrentTable { + torrentTable[k].StopAndWait() + } +} diff --git a/pkg/mcclient/auth/middleware.go b/pkg/mcclient/auth/middleware.go index 5c892a33ce..e01649696b 100644 --- a/pkg/mcclient/auth/middleware.go +++ b/pkg/mcclient/auth/middleware.go @@ -12,6 +12,12 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" ) +var ( + GuestToken = mcclient.SSimpleToken{ + User: "guest", + } +) + const ( AUTH_TOKEN = appctx.AppContextKey("X_AUTH_TOKEN") ) @@ -23,24 +29,26 @@ func Authenticate(f appsrv.FilterHandler) appsrv.FilterHandler { func AuthenticateWithDelayDecision(f appsrv.FilterHandler, delayDecision bool) appsrv.FilterHandler { return func(ctx context.Context, w http.ResponseWriter, r *http.Request) { tokenStr := r.Header.Get(mcclient.AUTH_TOKEN) + var token mcclient.TokenCredential if len(tokenStr) == 0 { log.Errorf("no auth_token found!") if !delayDecision { httperrors.UnauthorizedError(w, "Unauthorized") return } - } - token, err := Verify(tokenStr) - if err != nil { - log.Errorf("Verify token failed: %s", err) - if !delayDecision { - httperrors.UnauthorizedError(w, "InvalidToken") - return + token = &GuestToken + } else { + var err error + token, err = Verify(tokenStr) + if err != nil { + log.Errorf("Verify token failed: %s", err) + if !delayDecision { + httperrors.UnauthorizedError(w, "InvalidToken") + return + } } } - if token != nil { - ctx = context.WithValue(ctx, AUTH_TOKEN, token) - } + ctx = context.WithValue(ctx, AUTH_TOKEN, token) if taskId := r.Header.Get(mcclient.TASK_ID); taskId != "" { ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TASK_ID, taskId) diff --git a/pkg/util/nodeid/nodeid.go b/pkg/util/nodeid/nodeid.go index 75cc6e72a0..8d456598d2 100644 --- a/pkg/util/nodeid/nodeid.go +++ b/pkg/util/nodeid/nodeid.go @@ -171,20 +171,20 @@ func getLinux() ([]string, error) { return ret, nil } -func GetNodeId() (string, error) { +func GetNodeId() ([]byte, error) { var f func() ([]string, error) if runtime.GOOS == "linux" { f = getLinux } else { - return "", fmt.Errorf("Unsupported OS") + return nil, fmt.Errorf("Unsupported OS") } ret, e := f() if e != nil || len(ret) < 2 { log.Debugf("service info %s", ret) - return "", fmt.Errorf("Fail to generate Node ID") + return nil, fmt.Errorf("Fail to generate Node ID") } sn := md5.Sum([]byte(ret[0] + ret[1])) - return fmt.Sprintf("%x", sn), nil + return sn[0:], nil } diff --git a/pkg/util/sysutils/doc.go b/pkg/util/sysutils/doc.go new file mode 100644 index 0000000000..b5461ad8de --- /dev/null +++ b/pkg/util/sysutils/doc.go @@ -0,0 +1 @@ +package sysutils // import "yunion.io/x/onecloud/pkg/util/sysutils" diff --git a/pkg/util/sysutils/sysutils.go b/pkg/util/sysutils/sysutils.go new file mode 100644 index 0000000000..c286a26665 --- /dev/null +++ b/pkg/util/sysutils/sysutils.go @@ -0,0 +1,23 @@ +package sysutils + +import ( + "os" + "os/exec" +) + +func Start(closeFd bool, args ...string) (p *os.Process, err error) { + if args[0], err = exec.LookPath(args[0]); err == nil { + var procAttr os.ProcAttr + if closeFd { + procAttr.Files = []*os.File{nil, nil, nil} + } else { + procAttr.Files = []*os.File{os.Stdin, + os.Stdout, os.Stderr} + } + p, err := os.StartProcess(args[0], args, &procAttr) + if err == nil { + return p, nil + } + } + return nil, err +} diff --git a/pkg/yunionconf/service/service.go b/pkg/yunionconf/service/service.go index 2f9f2aee95..c2f25013e6 100644 --- a/pkg/yunionconf/service/service.go +++ b/pkg/yunionconf/service/service.go @@ -30,7 +30,6 @@ func StartService() { } cloudcommon.InitDB(dbOpts) - defer cloudcommon.CloseDB() app := cloudcommon.InitApp(commonOpts, true) yunionconf.InitHandlers(app) @@ -38,7 +37,9 @@ func StartService() { if db.CheckSync(opts.AutoSyncTable) { err := models.InitDB() if err == nil { - cloudcommon.ServeForever(app, commonOpts) + cloudcommon.ServeForever(app, commonOpts, func() { + cloudcommon.CloseDB() + }) } else { log.Errorf("InitDB fail: %s", err) }