diff --git a/cmd/torrent/main.go b/cmd/torrent/main.go index a31539e76a..af4d462e43 100644 --- a/cmd/torrent/main.go +++ b/cmd/torrent/main.go @@ -1,21 +1,25 @@ package main import ( + "crypto/sha1" "fmt" + "io/ioutil" + "net/http" "os" "os/signal" - "path" "path/filepath" "syscall" "time" "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" + "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/nodeid" "yunion.io/x/onecloud/pkg/util/torrentutils" ) @@ -28,7 +32,10 @@ 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"` } func exitSignalHandlers(client *torrent.Client) { @@ -71,13 +78,15 @@ 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 { log.Fatalf("fail to save torrent file %s", err) } + rootDir = filepath.Dir(root) } else { // client mode, load mi from torrent file @@ -85,28 +94,49 @@ func main() { if err != nil { log.Fatalf("fail to open torrent file %s", err) } + 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 - if len(options.Tracker) > 0 { - // server mode - clientConfig.DataDir = path.Dir(root) - } else { - // client mode - clientConfig.DataDir = root - } + + 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, 0700) + 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 @@ -122,32 +152,50 @@ func main() { go exitSignalHandlers(client) + start := time.Now() + t, err := client.AddTorrent(mi) if err != nil { 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 { + finish := false + + for !stop { if t.BytesCompleted() == t.Info().TotalLength() { + 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) + } + } + } 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 a35fc82ea7..3b24dcb823 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -5,13 +5,16 @@ import ( "fmt" "math/rand" "net/http" + "os" "strings" "sync" + "syscall" "time" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/trace" + "yunion.io/x/pkg/util/signalutils" "yunion.io/x/pkg/utils" "yunion.io/x/onecloud/pkg/appctx" @@ -37,8 +40,8 @@ type Application struct { cors *Cors middlewares []MiddlewareFunc - // record Http server for handle shotdown - server *http.Server + isExiting bool + idleConnsClosed chan struct{} } const ( @@ -329,32 +332,73 @@ func (app *Application) initServer(addr string) *http.Server { return s } +func (app *Application) registerCleanShutdown(s *http.Server, onStop func()) { + app.idleConnsClosed = make(chan struct{}) + + // dump goroutine stack + signalutils.RegisterSignal(func() { + utils.DumpAllGoroutineStack(log.Logger().Out) + }, syscall.SIGUSR1) + + quitSignals := []os.Signal{syscall.SIGHUP, syscall.SIGINT, syscall.SIGQUIT, syscall.SIGTERM} + signalutils.RegisterSignal(func() { + if app.isExiting { + log.Infof("Quit signal received!!! clean up in progress, be patient...") + return + } + app.isExiting = true + log.Infof("Quit signal received!!! do cleanup...") + + if err := s.Shutdown(context.Background()); err != nil { + // Error from closing listeners, or context timeout: + log.Errorf("HTTP server Shutdown: %v", err) + } + if onStop != nil { + func() { + defer func() { + if r := recover(); r != nil { + log.Errorf("app exiting error: %s", r) + } + }() + onStop() + }() + } + close(app.idleConnsClosed) + }, quitSignals...) + + signalutils.StartTrap() +} + +func (app *Application) waitCleanShutdown() { + <-app.idleConnsClosed + log.Infof("Service stopped.") +} + func (app *Application) ListenAndServe(addr string) { - app.server = app.initServer(addr) - err := app.server.ListenAndServe() - if err != nil { - log.Errorf("ListAndServer: %s", err) - panic(err) - } -} - -func (app *Application) IsInServe() bool { - return app.server != nil -} - -func (app *Application) ShutDown(ctx context.Context) error { - if app.server != nil { - return app.server.Shutdown(ctx) - } - return fmt.Errorf("Not init http server ??") + app.ListenAndServeWithCleanup(addr, nil) } func (app *Application) ListenAndServeTLS(addr string, certFile, keyFile string) { + app.ListenAndServeTLSWithCleanup(addr, certFile, keyFile, nil) +} + +func (app *Application) ListenAndServeWithCleanup(addr string, onStop func()) { + app.ListenAndServeTLSWithCleanup(addr, "", "", onStop) +} + +func (app *Application) ListenAndServeTLSWithCleanup(addr string, certFile, keyFile string, onStop func()) { s := app.initServer(addr) - err := s.ListenAndServeTLS(certFile, keyFile) + app.registerCleanShutdown(s, onStop) + var err error + if len(certFile) == 0 && len(keyFile) == 0 { + err = s.ListenAndServe() + } else { + err = s.ListenAndServeTLS(certFile, keyFile) + } if err != nil && err != http.ErrServerClosed { log.Fatalf("ListAndServer fail: %s", err) } + app.waitCleanShutdown() } func isJsonContentType(r *http.Request) bool { diff --git a/pkg/baremetal/service/service.go b/pkg/baremetal/service/service.go index 13bac7f266..a0596ed205 100644 --- a/pkg/baremetal/service/service.go +++ b/pkg/baremetal/service/service.go @@ -1,7 +1,6 @@ package service import ( - "context" "os" "yunion.io/x/log" @@ -16,7 +15,6 @@ import ( type BaremetalService struct { service.SServiceBase - isExiting bool } func New() *BaremetalService { @@ -30,25 +28,9 @@ func (s *BaremetalService) StartService() { app := cloudcommon.InitApp(&o.Options.CommonOptions, false) handler.InitHandlers(app) - s.RegisterSIGUSR1() - s.RegisterQuitSignals(func() { - log.Infof("Baremetal agent quit !!!") - if s.isExiting { - return - } else { - s.isExiting = true - } - - if app.IsInServe() { - if err := app.ShutDown(context.Background()); err != nil { - log.Errorf("App shutdown err: %v", err) - } - } + cloudcommon.ServeForeverWithCleanup(app, &o.Options.CommonOptions, func() { tasks.OnStop() - os.Exit(0) }) - - cloudcommon.ServeForever(app, &o.Options.CommonOptions) } func (s *BaremetalService) startAgent() { diff --git a/pkg/cloudcommon/app.go b/pkg/cloudcommon/app.go index 5ed94b973a..6463dfbbea 100644 --- a/pkg/cloudcommon/app.go +++ b/pkg/cloudcommon/app.go @@ -24,6 +24,10 @@ func InitApp(options *CommonOptions, dbAccess bool) *appsrv.Application { } func ServeForever(app *appsrv.Application, options *CommonOptions) { + ServeForeverWithCleanup(app, options, nil) +} + +func ServeForeverWithCleanup(app *appsrv.Application, options *CommonOptions, onStop func()) { AppDBInit(app) addr := net.JoinHostPort(options.Address, strconv.Itoa(options.Port)) proto := "http" @@ -47,8 +51,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.ListenAndServeTLSWithCleanup(addr, certfile, options.SslKeyfile, onStop) } else { - app.ListenAndServe(addr) + app.ListenAndServeWithCleanup(addr, onStop) } } diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index fd8270a38c..f9fa6cc936 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 d6eaa528a9..85378b27b2 100644 --- a/pkg/cloudcommon/policy/defaults.go +++ b/pkg/cloudcommon/policy/defaults.go @@ -137,6 +137,13 @@ var ( Action: PolicyActionGet, Result: rbacutils.OwnerAllow, }, + { + Service: "image", + Resource: "images", + Action: PolicyActionPerform, + Extra: []string{"update-torrent-status"}, + Result: rbacutils.GuestAllow, + }, { Service: "log", Resource: "actions", diff --git a/pkg/cloudcommon/policy/policy.go b/pkg/cloudcommon/policy/policy.go index a9df96880f..13c5aa4c7d 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/cloudcommon/service/services.go b/pkg/cloudcommon/service/services.go index b6d7cc66a1..e6d29208d6 100644 --- a/pkg/cloudcommon/service/services.go +++ b/pkg/cloudcommon/service/services.go @@ -1,26 +1,3 @@ package service -import ( - "os" - "syscall" - - "yunion.io/x/log" - "yunion.io/x/pkg/util/signalutils" - "yunion.io/x/pkg/utils" -) - type SServiceBase struct{} - -func (s *SServiceBase) RegisterQuitSignals(quitHandler signalutils.Trap) { - quitSignals := []os.Signal{syscall.SIGHUP, syscall.SIGINT, syscall.SIGQUIT, syscall.SIGTERM} - signalutils.RegisterSignal(quitHandler, quitSignals...) - - signalutils.StartTrap() -} - -func (s *SServiceBase) RegisterSIGUSR1() { - // dump goroutine stack - signalutils.RegisterSignal(func() { - utils.DumpAllGoroutineStack(log.Logger().Out) - }, syscall.SIGUSR1) -} diff --git a/pkg/cloudir/service/service.go b/pkg/cloudir/service/service.go index 9a3c2799a7..ff07cc1334 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.ServeForeverWithCleanup(app, commonOpts, func() { + etcd.CloseDefaultEtcdClient() + }) } diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index a14fd1b72a..5b2a0a610e 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -1,7 +1,6 @@ package hostman import ( - "context" "os" "yunion.io/x/log" @@ -23,8 +22,6 @@ import ( type SHostService struct { service.SServiceBase - - isExiting bool } func (host *SHostService) StartService() { @@ -37,28 +34,6 @@ func (host *SHostService) StartService() { log.Fatalf(err.Error()) } - host.RegisterSIGUSR1() - host.RegisterQuitSignals(func() { // register quit handler - if host.isExiting { - return - } else { - host.isExiting = true - } - - if app.IsInServe() { - if err := app.ShutDown(context.Background()); err != nil { - log.Errorln(err.Error()) - } - } - - hostinfo.Stop() - storageman.Stop() - hostmetrics.Stop() - guestman.Stop() - hostutils.GetWorkManager().Stop() - os.Exit(0) - }) - if err := storageman.Init(hostInstance); err != nil { log.Fatalf(err.Error()) } @@ -82,20 +57,13 @@ func (host *SHostService) StartService() { cronManager.AddJob2( "CleanRecycleDiskFiles", 1, 3, 0, 0, storageman.CleanRecycleDiskfiles, false) - go func() { - defer func() { - if r := recover(); r != nil { - if !host.isExiting { - log.Fatalf("%s", r) - } else { - log.Errorln(r) - } - } - - }() - cloudcommon.ServeForever(app, &options.HostOptions.CommonOptions) - }() - select {} // for quit handler + cloudcommon.ServeForeverWithCleanup(app, &options.HostOptions.CommonOptions, func() { + hostinfo.Stop() + storageman.Stop() + hostmetrics.Stop() + guestman.Stop() + hostutils.GetWorkManager().Stop() + }) } func (host *SHostService) initHandlers(app *appsrv.Application) { 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 15af820b6e..80f695d656 100644 --- a/pkg/image/models/images.go +++ b/pkg/image/models/images.go @@ -899,7 +899,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) } } @@ -1046,7 +1046,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() } } @@ -1081,3 +1081,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 edc4fcf17e..7f07c220b0 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 d973cf0277..8e891f01dd 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(true) cron.AddJob1("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages) cron.Start() - defer cron.Stop() - cloudcommon.ServeForever(app, commonOpts) + cloudcommon.ServeForeverWithCleanup(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/sysutils.go b/pkg/util/sysutils/sysutils.go index 8b4d4eff88..63d9ae6731 100644 --- a/pkg/util/sysutils/sysutils.go +++ b/pkg/util/sysutils/sysutils.go @@ -3,6 +3,8 @@ package sysutils import ( "fmt" "net" + "os" + "os/exec" "strconv" "strings" @@ -231,3 +233,20 @@ func GetSerialPorts(lines []string) []string { } return ret } + +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..48914b08bd 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.ServeForeverWithCleanup(app, commonOpts, func() { + cloudcommon.CloseDB() + }) } else { log.Errorf("InitDB fail: %s", err) }