From 60fe104da23dff91eae72f4f7556a4c3f8106d96 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Fri, 25 Jan 2019 17:18:20 +0800 Subject: [PATCH] snapshot recycle, disk recycle, dhcp relay use raw socket --- pkg/cloudcommon/cronman/cronman.go | 16 +-- pkg/cloudcommon/dhcp/conn.go | 80 ++++++++++- pkg/cloudcommon/dhcp/conn_linux.go | 2 +- pkg/cloudcommon/dhcp/server.go | 2 +- pkg/compute/service/service.go | 2 +- pkg/hostman/guestfs/localfs.go | 1 - pkg/hostman/guestman/guestman.go | 2 - pkg/hostman/guestman/guesttasks.go | 24 +++- pkg/hostman/guestman/qemu-kvm.go | 6 +- pkg/hostman/host_services.go | 11 +- pkg/hostman/hostinfo/hostdhcp/dhcprelay.go | 22 +-- pkg/hostman/hostinfo/hostdhcp/dhcpserver.go | 12 +- pkg/hostman/hostinfo/hostinfo.go | 4 +- pkg/hostman/hostinfo/hostinfohelper.go | 7 +- pkg/hostman/hostutils/hostutils.go | 2 +- pkg/hostman/monitor/qmp.go | 3 +- pkg/hostman/options/options.go | 4 +- pkg/hostman/storageman/core.go | 43 ++++++ pkg/hostman/storageman/storagelocal.go | 146 ++++++++++++++++++-- pkg/image/service/service.go | 2 +- pkg/util/cgrouputils/cgrouputils.go | 3 +- pkg/util/cgrouputils/cgrouputils_test.go | 19 +++ 22 files changed, 344 insertions(+), 69 deletions(-) create mode 100644 pkg/util/cgrouputils/cgrouputils_test.go diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index 9946a03a99..fd8270a38c 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -17,13 +17,6 @@ type TCronJobFunction func(ctx context.Context, userCred mcclient.TokenCredentia var manager *SCronJobManager -func init() { - manager = &SCronJobManager{ - jobs: make([]*SCronJob, 0), - workers: appsrv.NewWorkerManager("CronJobWorkers", 4, 1024, true), - } -} - type ICronTimer interface { Next(time.Time) time.Time } @@ -93,7 +86,14 @@ type SCronJobManager struct { workers *appsrv.SWorkerManager } -func GetCronJobManager() *SCronJobManager { +func GetCronJobManager(idDbWorker bool) *SCronJobManager { + if manager == nil { + manager = &SCronJobManager{ + jobs: make([]*SCronJob, 0), + workers: appsrv.NewWorkerManager("CronJobWorkers", 4, 1024, idDbWorker), + } + } + return manager } diff --git a/pkg/cloudcommon/dhcp/conn.go b/pkg/cloudcommon/dhcp/conn.go index 0b1001fa37..fec2d22420 100644 --- a/pkg/cloudcommon/dhcp/conn.go +++ b/pkg/cloudcommon/dhcp/conn.go @@ -15,11 +15,12 @@ package dhcp import ( - //"errors" + "errors" "fmt" "io" "net" "strconv" + "syscall" "time" "golang.org/x/net/ipv4" @@ -73,7 +74,11 @@ func NewConn(addr string) (*Conn, error) { return newConn(addr, newPortableConn) } -func newConn(addr string, n func(int) (conn, error)) (*Conn, error) { +func NewSocketConn(addr string) (*Conn, error) { + return newConn(addr, newSocketConn) +} + +func newConn(addr string, n func(net.IP, int) (conn, error)) (*Conn, error) { if addr == "" { addr = "0.0.0.0:67" } @@ -95,7 +100,7 @@ func newConn(addr string, n func(int) (conn, error)) (*Conn, error) { } } - c, err := n(udpAddr.Port) + c, err := n(udpAddr.IP, udpAddr.Port) if err != nil { return nil, err } @@ -221,7 +226,7 @@ type portableConn struct { conn *ipv4.PacketConn } -func newPortableConn(port int) (conn, error) { +func newPortableConn(_ net.IP, port int) (conn, error) { c, err := net.ListenPacket("udp4", fmt.Sprintf(":%d", port)) if err != nil { return nil, err @@ -265,3 +270,70 @@ func (c *portableConn) SetReadDeadline(t time.Time) error { func (c *portableConn) SetWriteDeadline(t time.Time) error { return c.conn.SetWriteDeadline(t) } + +type socketConn struct { + sock int +} + +func newSocketConn(addr net.IP, port int) (conn, error) { + sock, err := syscall.Socket(syscall.AF_INET, syscall.SOCK_DGRAM, 0) + if err != nil { + return nil, err + } + if err = syscall.SetsockoptInt(sock, syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1); err != nil { + return nil, err + } + byteAddr := [4]byte{} + copy(byteAddr[:], addr.To4()[:4]) + lsa := &syscall.SockaddrInet4{ + Port: port, + Addr: byteAddr, + } + if err = syscall.Bind(sock, lsa); err != nil { + return nil, err + } + if err = syscall.SetNonblock(sock, false); err != nil { + return nil, err + } + return &socketConn{sock}, nil +} + +func (s *socketConn) Close() error { + return syscall.Close(s.sock) +} + +func (s *socketConn) Recv(b []byte) (rb []byte, addr *net.UDPAddr, ifidx int, err error) { + n, a, err := syscall.Recvfrom(s.sock, b, 0) + if err != nil { + return nil, nil, 0, err + } + if addr, ok := a.(*syscall.SockaddrInet4); !ok { + return nil, nil, 0, errors.New("Recvfrom recevice address is not famliy Inet4") + } else { + ip := net.IP{addr.Addr[0], addr.Addr[1], addr.Addr[2], addr.Addr[3]} + udpAddr := &net.UDPAddr{ + IP: ip, + Port: addr.Port, + } + // there is no interface index info + return b[:n], udpAddr, 0, nil + } +} + +func (s *socketConn) Send(b []byte, addr *net.UDPAddr, ifidx int) error { + destIp := [4]byte{} + copy(destIp[:], addr.IP.To4()[:4]) + destAddr := &syscall.SockaddrInet4{ + Addr: destIp, + Port: addr.Port, + } + return syscall.Sendto(s.sock, b, 0, destAddr) +} + +func (s *socketConn) SetReadDeadline(t time.Time) error { + return errors.New("Not Implement") +} + +func (s *socketConn) SetWriteDeadline(t time.Time) error { + return errors.New("Not Implement") +} diff --git a/pkg/cloudcommon/dhcp/conn_linux.go b/pkg/cloudcommon/dhcp/conn_linux.go index 9bad9ce2b6..23da403aef 100644 --- a/pkg/cloudcommon/dhcp/conn_linux.go +++ b/pkg/cloudcommon/dhcp/conn_linux.go @@ -40,7 +40,7 @@ func NewSnooperConn(addr string) (*Conn, error) { return newConn(addr, newLinuxConn) } -func newLinuxConn(port int) (conn, error) { +func newLinuxConn(_ net.IP, port int) (conn, error) { if port == 0 { return nil, errors.New("must specify a listen port") } diff --git a/pkg/cloudcommon/dhcp/server.go b/pkg/cloudcommon/dhcp/server.go index 3abdcc4500..a5a0ebae5c 100644 --- a/pkg/cloudcommon/dhcp/server.go +++ b/pkg/cloudcommon/dhcp/server.go @@ -22,7 +22,7 @@ func NewDHCPServer(address string, port int) *DHCPServer { } func NewDHCPServer2(address string, port int) (*DHCPServer, *Conn, error) { - conn, err := NewConn(fmt.Sprintf("%s:%d", address, port)) + conn, err := NewSocketConn(fmt.Sprintf("%s:%d", address, port)) if err != nil { return nil, nil, err } diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 048f5ef7ed..bec2b86325 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -59,7 +59,7 @@ func StartService() { log.Errorf("InitDB fail: %s", err) } - cron := cronman.GetCronJobManager() + cron := cronman.GetCronJobManager(true) cron.AddJob1("CleanPendingDeleteServers", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) cron.AddJob1("CleanPendingDeleteDisks", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks) cron.AddJob1("CleanPendingDeleteLoadbalancers", time.Duration(opts.LoadbalancerPendingDeleteCheckInterval)*time.Second, models.LoadbalancerAgentManager.CleanPendingDeleteLoadbalancers) diff --git a/pkg/hostman/guestfs/localfs.go b/pkg/hostman/guestfs/localfs.go index 368587dc6a..7409fed896 100644 --- a/pkg/hostman/guestfs/localfs.go +++ b/pkg/hostman/guestfs/localfs.go @@ -111,7 +111,6 @@ func (f *SLocalGuestFS) Zerofiles(dir string, caseInsensitive bool) error { return fmt.Errorf("No such file %s", sPath) } -// TODO: test func (f *SLocalGuestFS) Passwd(account, password string, caseInsensitive bool) error { var proc = exec.Command("chroot", f.mountPath, "passwd", account) stdin, err := proc.StdinPipe() diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 67f9ad4bb5..e98c8c5d64 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -637,8 +637,6 @@ func (m *SGuestManager) ExitGuestCleanup() { guest.ExitCleanup(false) } - // TODO - // m.StopCpusetBalancer() cgrouputils.CgroupCleanAll() // TODO // hostmetrics? diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 57d9f41b76..6d90697a84 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -48,10 +48,9 @@ func NewGuestStopTask(guest *SKVMGuestInstance, ctx context.Context, timeout int func (s *SGuestStopTask) Start() { if s.IsRunning() && s.IsMonitorAlive() { - // Do Powerdown, s.Monitor.SimpleCommand("system_powerdown", s.onPowerdownGuest) } else { - s.CheckGuestRunningLater() + s.checkGuestRunning() } } @@ -62,8 +61,7 @@ func (s *SGuestStopTask) onPowerdownGuest(results string) { } func (s *SGuestStopTask) checkGuestRunning() { - if !s.IsRunning() || - time.Now().Sub(s.startPowerdown) > time.Duration(s.timeout)*time.Second { + if !s.IsRunning() || time.Now().Sub(s.startPowerdown) > time.Duration(s.timeout)*time.Second { s.Stop() // force stop hostutils.TaskComplete(s.ctx, nil) } else { @@ -447,7 +445,7 @@ func (s *SGuestResumeTask) onResumeSucc(res string) { } func (s *SGuestResumeTask) onStartRunning() { - s.removeStatefile() + // s.removeStatefile() XXX 可能不用了,先注释了 if s.ctx != nil && len(appctx.AppContextTaskId(s.ctx)) > 0 { hostutils.TaskComplete(s.ctx, nil) } @@ -659,9 +657,21 @@ func NewGuestReloadDiskTask( } } -func (s *SGuestReloadDiskTask) WaitSnapshotReplaced(callback func()) { - // TODO +func (s *SGuestReloadDiskTask) WaitSnapshotReplaced(callback func()) error { + var retry = 0 + for { + retry += 1 + if retry == 300 { + return fmt.Errorf( + "SnapshotDeleteJob.deleting_disk_snapshot always has %s", s.disk.GetId()) + } + if _, ok := storageman.DELETEING_SNAPSHOTS[s.disk.GetId()]; ok { + time.Sleep(time.Second * 1) + } + } + callback() + return nil } func (s *SGuestReloadDiskTask) Start() { diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 32fc80d1dd..30fd0cc65e 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -142,7 +142,6 @@ func (s *SKVMGuestInstance) isSelfQemuPid(pid, uuid string) bool { cmdlineFile := fmt.Sprintf("/proc/%s/cmdline", pid) fi, err := os.Stat(cmdlineFile) if err != nil { - log.Warningf("IsSelfQemuPid Stat File %s error %s", cmdlineFile, err) return false } if !fi.Mode().IsRegular() { @@ -641,8 +640,6 @@ func (s *SKVMGuestInstance) delTmpDisks(ctx context.Context, migrated bool) erro } func (s *SKVMGuestInstance) Delete(ctx context.Context, migrated bool) error { - // self._del_bw_limit() - // self._del_netmon_nic() ?? 需要开发? if err := s.delTmpDisks(ctx, migrated); err != nil { return err } @@ -1076,8 +1073,7 @@ func (s *SKVMGuestInstance) ExecReloadDiskTask(ctx context.Context, disk storage if s.IsRunning() { if s.isLiveSnapshotEnabled() { task := NewGuestReloadDiskTask(ctx, s, disk) - task.WaitSnapshotReplaced(task.Start) - return nil, nil + return nil, task.WaitSnapshotReplaced(task.Start) } else { return nil, fmt.Errorf("Guest dosen't support reload disk") } diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index 9e715f3268..793ab3fa8a 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -8,6 +8,7 @@ import ( "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/cloudcommon" + "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/service" "yunion.io/x/onecloud/pkg/hostman/downloader" "yunion.io/x/onecloud/pkg/hostman/guestman" @@ -42,8 +43,7 @@ func (host *SHostService) StartService() { } if app.IsInServe() { - err := app.ShowDown(context.Background()) - if err != nil { + if err := app.ShowDown(context.Background()); err != nil { log.Errorln(err.Error()) } } @@ -52,7 +52,6 @@ func (host *SHostService) StartService() { storageman.Stop() guestman.Stop() hostutils.GetWorkManager().Stop() - os.Exit(0) }) @@ -67,8 +66,12 @@ func (host *SHostService) StartService() { hostInstance.StartRegister(2, guestman.GetGuestManager().Bootstrap) }) host.initHandlers(app) - <-hostinfo.Instance().IsRegistered // wait host and guest init + + cronManager := cronman.GetCronJobManager(false) + cronManager.AddJob2( + "CleanRecycleDiskFiles", 1, 3, 0, 0, storageman.CleanRecycleDiskfiles, false) + cloudcommon.ServeForever(app, &options.HostOptions.CommonOptions) } diff --git a/pkg/hostman/hostinfo/hostdhcp/dhcprelay.go b/pkg/hostman/hostinfo/hostdhcp/dhcprelay.go index 49f39083cb..2e0e7d8412 100644 --- a/pkg/hostman/hostinfo/hostdhcp/dhcprelay.go +++ b/pkg/hostman/hostinfo/hostdhcp/dhcprelay.go @@ -10,7 +10,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/dhcp" ) -const DEFAULT_DHCP_RELAY_PORT = 168 +const DEFAULT_DHCP_RELAY_PORT = 68 type recvFunc func(pkt *dhcp.Packet) @@ -33,7 +33,6 @@ type SDHCPRelay struct { destport int cache sync.Map - // cache map[string]*SRelayCache } func NewDHCPRelay(addrs []string) (*SDHCPRelay, error) { @@ -47,15 +46,7 @@ func NewDHCPRelay(addrs []string) (*SDHCPRelay, error) { relay.destaddr = net.ParseIP(addr) relay.destport = port relay.cache = sync.Map{} - // relay.cache = make(map[string]*SRelayCache, 0) - log.Infof("DHCP Relay Bind addr %s port %d", - DEFAULT_DHCP_LISTEN_ADDR, DEFAULT_DHCP_RELAY_PORT) - relay.server, relay.conn, err = dhcp.NewDHCPServer2(DEFAULT_DHCP_LISTEN_ADDR, DEFAULT_DHCP_RELAY_PORT) - if err != nil { - log.Errorln(err) - return nil, err - } return relay, nil } @@ -69,8 +60,17 @@ func (r *SDHCPRelay) Start() { }() } -func (r *SDHCPRelay) Setup(addr string) { +func (r *SDHCPRelay) Setup(addr string) error { + var err error r.srcaddr = addr + log.Infof("DHCP Relay Bind addr %s port %d", r.srcaddr, DEFAULT_DHCP_RELAY_PORT) + r.server, r.conn, err = dhcp.NewDHCPServer2(r.srcaddr, DEFAULT_DHCP_RELAY_PORT) + if err != nil { + log.Errorln(err) + return err + } + r.Start() + return nil } func (r *SDHCPRelay) ServeDHCP(pkt dhcp.Packet, addr *net.UDPAddr, intf *net.Interface) (dhcp.Packet, error) { diff --git a/pkg/hostman/hostinfo/hostdhcp/dhcpserver.go b/pkg/hostman/hostinfo/hostdhcp/dhcpserver.go index ef41a6e44f..1694f4cde7 100644 --- a/pkg/hostman/hostinfo/hostdhcp/dhcpserver.go +++ b/pkg/hostman/hostinfo/hostdhcp/dhcpserver.go @@ -16,7 +16,7 @@ import ( "yunion.io/x/onecloud/pkg/util/netutils2" ) -var DEFAULT_DHCP_LISTEN_ADDR = "0.0.0.0" +var DEFAULT_DHCP_BIND_ADDR = "0.0.0.0" type SGuestDHCPServer struct { server *dhcp.DHCPServer @@ -37,16 +37,13 @@ func NewGuestDHCPServer(iface string, relay []string) (*SGuestDHCPServer, error) return nil, err } } - guestdhcp.server = dhcp.NewDHCPServer(DEFAULT_DHCP_LISTEN_ADDR, options.HostOptions.DhcpServerPort) + guestdhcp.server = dhcp.NewDHCPServer(DEFAULT_DHCP_BIND_ADDR, options.HostOptions.DhcpServerPort) guestdhcp.iface = iface return guestdhcp, nil } func (s *SGuestDHCPServer) Start() { log.Infof("SGuestDHCPServer starting ...") - if s.relay != nil { - s.relay.Start() - } go func() { err := s.server.ListenAndServe(s) if err != nil { @@ -55,10 +52,11 @@ func (s *SGuestDHCPServer) Start() { }() } -func (s *SGuestDHCPServer) RelaySetup(addr string) { +func (s *SGuestDHCPServer) RelaySetup(addr string) error { if s.relay != nil { - s.relay.Setup(addr) + return s.relay.Setup(addr) } + return nil } func (s *SGuestDHCPServer) getGuestConfig(guestDesc, guestNic jsonutils.JSONObject) *dhcp.ResponseConfig { diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index ca3e304259..1c7ecc987d 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -148,7 +148,9 @@ func (h *SHostInfo) parseConfig() error { h.Nics = append(h.Nics, nic) } for i := 0; i < len(h.Nics); i++ { - h.Nics[i].SetupDhcpRelay() + if err := h.Nics[i].SetupDhcpRelay(); err != nil { + return err + } } if man, err := isolated_device.NewManager(h); err != nil { diff --git a/pkg/hostman/hostinfo/hostinfohelper.go b/pkg/hostman/hostinfo/hostinfohelper.go index aea174035e..cc2c67ef6d 100644 --- a/pkg/hostman/hostinfo/hostinfohelper.go +++ b/pkg/hostman/hostinfo/hostinfohelper.go @@ -191,10 +191,13 @@ func (n *SNIC) EnableDHCPRelay() bool { } } -func (n *SNIC) SetupDhcpRelay() { +func (n *SNIC) SetupDhcpRelay() error { if n.EnableDHCPRelay() { - n.dhcpServer.RelaySetup(n.Ip) + if err := n.dhcpServer.RelaySetup(n.Ip); err != nil { + return err + } } + return nil } func (n *SNIC) SetWireId(wire, wireId string, bandwidth int64) { diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 30c1eae71f..0a260eb1b9 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -46,7 +46,7 @@ func TaskFailed(ctx context.Context, reason string) { if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil { modules.ComputeTasks.TaskFailed2(GetComputeSession(ctx), taskId.(string), reason) } else { - log.Errorln("Reqeuest task failed missing task id, with reason(%v)", reason) + log.Errorf("Reqeuest task failed missing task id, with reason(%s)", reason) } } diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 058b083a68..32ee2edfd4 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -178,7 +178,6 @@ func (m *QmpMonitor) read(r io.Reader) { // remove reader timeout m.rwc.SetReadDeadline(time.Time{}) - log.Infof("Qmp Connected") m.connected = true m.timeout = false go m.query() @@ -189,7 +188,7 @@ func (m *QmpMonitor) read(r io.Reader) { log.Infof("Scan over ...") err := scanner.Err() if err != nil { - log.Errorln(err) + log.Infof("QMP Disconnected: %s", err) } if m.timeout { m.OnMonitorTimeout(err) diff --git a/pkg/hostman/options/options.go b/pkg/hostman/options/options.go index db0c8d2657..e0fc80ce7e 100644 --- a/pkg/hostman/options/options.go +++ b/pkg/hostman/options/options.go @@ -82,7 +82,9 @@ 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"` + + SnapshotDirSuffix string `help:"Snapshot dir name equal diskId concat snapshot dir suffix" default:"_snap"` + SnapshotRecycleDay int `default:"1" help:"Snapshot Recycle delete Duration day"` } var HostOptions SHostOptions diff --git a/pkg/hostman/storageman/core.go b/pkg/hostman/storageman/core.go index 85694c3a7a..70c337efbc 100644 --- a/pkg/hostman/storageman/core.go +++ b/pkg/hostman/storageman/core.go @@ -1,14 +1,22 @@ package storageman import ( + "context" "fmt" + "io/ioutil" "path" "strings" + "time" + + "yunion.io/x/log" + "yunion.io/x/pkg/util/timeutils" "yunion.io/x/onecloud/pkg/cloudcommon/storagetypes" "yunion.io/x/onecloud/pkg/hostman/hostutils" "yunion.io/x/onecloud/pkg/hostman/options" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/fileutils2" + "yunion.io/x/onecloud/pkg/util/procutils" ) const MINIMAL_FREE_SPACE = 128 @@ -279,3 +287,38 @@ func Init(host hostutils.IHost) error { func Stop() { // pass do nothing } + +func cleanDailyFiles(storagePath, subDir string, keepDay int) { + recycleDir := path.Join(storagePath, subDir) + if !fileutils2.Exists(recycleDir) { + return + } + + // before mark should be deleted + markTime := timeutils.UtcNow().Add(time.Hour * 24 * -1 * time.Duration(keepDay)) + files, err := ioutil.ReadDir(recycleDir) + if err != nil { + log.Errorln(err) + return + } + + for _, file := range files { + date, err := timeutils.ParseTimeStr(file.Name()) + if err != nil { + log.Errorln(err) + continue + } + if date.Before(markTime) { + log.Infof("Real delete %s", file) + subDirPath := path.Join(recycleDir, file.Name()) + procutils.NewCommand("rm", "-rf", subDirPath) + } + } +} + +func CleanRecycleDiskfiles(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + for _, d := range options.HostOptions.LocalImagePath { + cleanDailyFiles(d, _RECYCLE_BIN_, options.HostOptions.RecycleDiskfileKeepDays) + cleanDailyFiles(d, _IMGSAVE_BACKUPS_, options.HostOptions.RecycleDiskfileKeepDays) + } +} diff --git a/pkg/hostman/storageman/storagelocal.go b/pkg/hostman/storageman/storagelocal.go index 131701179d..ad7f8bec78 100644 --- a/pkg/hostman/storageman/storagelocal.go +++ b/pkg/hostman/storageman/storagelocal.go @@ -3,28 +3,34 @@ package storageman import ( "context" "fmt" + "io/ioutil" "os" "path" + "regexp" "time" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/util/timeutils" + + "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/storagetypes" "yunion.io/x/onecloud/pkg/hostman/guestfs/fsdriver" "yunion.io/x/onecloud/pkg/hostman/hostutils" "yunion.io/x/onecloud/pkg/hostman/options" "yunion.io/x/onecloud/pkg/hostman/storageman/remotefile" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modules" "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/procutils" "yunion.io/x/onecloud/pkg/util/qemuimg" - "yunion.io/x/pkg/util/timeutils" ) var ( - _FUSE_MOUNT_PATH_ = "fusemnt" - _FUSE_TMP_PATH_ = "fusetmp" - _SNAPSHOT_PATH_ = "snapshots" + _FUSE_MOUNT_PATH_ = "fusemnt" + _FUSE_TMP_PATH_ = "fusetmp" + _SNAPSHOT_PATH_ = "snapshots" + DELETEING_SNAPSHOTS = map[string]bool{} ) type SLocalStorage struct { @@ -117,10 +123,6 @@ func (s *SLocalStorage) CreateDisk(diskId string) IDisk { return disk } -func (s *SLocalStorage) StartSnapshotRecycle() { - //TODO -} - func (s *SLocalStorage) Accessible() bool { if !fileutils2.Exists(s.Path) { if _, err := procutils.NewCommand("mkdir", "-p", s.Path).Run(); err != nil { @@ -324,3 +326,131 @@ func (s *SLocalStorage) DeleteSnapshots(ctx context.Context, params interface{}) } return nil, nil } + +/*************************Background delete snapshot job****************************/ + +func (s *SLocalStorage) StartSnapshotRecycle() { + log.Infof("Snapshot recyle job started") + if !fileutils2.Exists(s.GetSnapshotDir()) { + procutils.NewCommand("mkdir", "-p", s.GetSnapshotDir()).Run() + } + cronman.GetCronJobManager(false).AddJob2("SnapshotRecycle", options.HostOptions.SnapshotRecycleDay, 2, 0, 0, s.snapshotRecycle, true) +} + +func (s *SLocalStorage) snapshotRecycle(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + res, err := modules.Snapshots.GetById(hostutils.GetComputeSession(ctx), "max-count", nil) + if err != nil { + log.Errorln(err) + return + } + maxSnapshotCount, err := res.Int("max_count") + if err != nil { + log.Errorln("Request region get snapshot max count failed") + return + } + files, err := ioutil.ReadDir(s.GetSnapshotDir()) + if err != nil { + log.Errorln(err) + return + } + for _, file := range files { + s.checkSnapshots(file.Name(), int(maxSnapshotCount)) + } +} + +func (s *SLocalStorage) checkSnapshots(snapshotDir string, maxSnapshotCount int) { + re := regexp.MustCompile(`^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}_snap$`) + if !re.MatchString(snapshotDir) { + log.Warningf("snapshot_dir got unexcept file %s", snapshotDir) + return + } + + diskId := snapshotDir[:len(snapshotDir)-len(options.HostOptions.SnapshotDirSuffix)] + snapshotPath := path.Join(s.GetSnapshotDir(), snapshotDir) + + // If disk is Deleted, request delete this disk all snapshots + if !fileutils2.Exists(path.Join(s.Path, diskId)) && fileutils2.Exists(snapshotPath) { + params := jsonutils.NewDict() + params.Set("disk_id", jsonutils.NewString(diskId)) + _, err := modules.Snapshots.PerformClassAction( + hostutils.GetComputeSession(context.Background()), + "delete-disk-snapshots", params) + if err != nil { + log.Infof("Request delele disk %s snapshots failed %s", diskId, err) + } + return + } + + snapshots, err := ioutil.ReadDir(snapshotPath) + if err != nil { + log.Errorln(err) + return + } + + // if snapshot count greater than maxsnapshot count, do convert + if len(snapshots) >= maxSnapshotCount { + s.requestConvertSnapshot(snapshotPath, diskId) + } +} + +func (s *SLocalStorage) requestConvertSnapshot(snapshotPath, diskId string) { + res, err := modules.Disks.GetSpecific( + hostutils.GetComputeSession(context.Background()), diskId, "convert-snapshot", nil) + if err != nil { + log.Errorln(err) + return + } + + var ( + deleteSnapshot, _ = res.GetString("delete_snapshot") + convertSnapshot, _ = res.GetString("convert_snapshot") + pendingDelete, _ = res.Bool("pending_delete") + ) + log.Infof("start convert disk(%s) snapshot(%s), delete_snapshot is %s", + diskId, convertSnapshot, deleteSnapshot) + convertSnapshotPath := path.Join(snapshotPath, convertSnapshot) + outfile := convertSnapshotPath + ".tmp" + img, err := qemuimg.NewQemuImage(convertSnapshot) + if err != nil { + log.Errorln(err) + return + } + err = img.Convert2Qcow2To(outfile, true) + if err != nil { + log.Errorln(err) + return + } + s.requestDeleteSnapshot( + diskId, snapshotPath, deleteSnapshot, convertSnapshotPath, outfile, pendingDelete) +} + +func (s *SLocalStorage) requestDeleteSnapshot( + diskId, snapshotPath, deleteSnapshot, convertSnapshotPath, + outfile string, pendingDelete bool, +) { + deleteSnapshotPath := path.Join(snapshotPath, deleteSnapshot) + DELETEING_SNAPSHOTS[diskId] = true + defer delete(DELETEING_SNAPSHOTS, diskId) + _, err := modules.Snapshots.PerformAction(hostutils.GetComputeSession(context.Background()), + deleteSnapshot, "deleted", nil) + if err != nil { + log.Errorln(err) + return + } + if out, err := procutils.NewCommand("rm", "-f", convertSnapshotPath).Run(); err != nil { + log.Errorf("%s", out) + return + } + if out, err := procutils.NewCommand("mv", "-f", outfile, convertSnapshotPath).Run(); err != nil { + log.Errorf("%s", out) + return + } + if !pendingDelete { + if out, err := procutils.NewCommand("rm", "-f", deleteSnapshotPath).Run(); err != nil { + log.Errorf("%s", out) + return + } + } +} + +/******************************* END *****************************/ diff --git a/pkg/image/service/service.go b/pkg/image/service/service.go index 8d6485242b..d973cf0277 100644 --- a/pkg/image/service/service.go +++ b/pkg/image/service/service.go @@ -86,7 +86,7 @@ func StartService() { models.CheckImages() - cron := cronman.GetCronJobManager() + cron := cronman.GetCronJobManager(true) cron.AddJob1("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages) cron.Start() diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgrouputils.go index e774e86961..1a71767992 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgrouputils.go @@ -228,8 +228,9 @@ func (c *CGroupTask) MoveTasksToRoot() { func (c *CGroupTask) RemoveTask() bool { if c.taskIsExist() { c.MoveTasksToRoot() + log.Infof("Remove task path %s", c.TaskPath()) if err := os.Remove(c.TaskPath()); err != nil { - log.Errorf("Remove task %s", err) + log.Errorf("Remove task path failed %s", err) return false } } diff --git a/pkg/util/cgrouputils/cgrouputils_test.go b/pkg/util/cgrouputils/cgrouputils_test.go new file mode 100644 index 0000000000..9295a4471e --- /dev/null +++ b/pkg/util/cgrouputils/cgrouputils_test.go @@ -0,0 +1,19 @@ +package cgrouputils + +import ( + "bufio" + "fmt" + "os" + "strings" + "testing" +) + +func TestCgroupSet(t *testing.T) { + reader := bufio.NewReader(os.Stdin) + fmt.Print("Enter pid: ") + pid, _ := reader.ReadString('\n') + pid = strings.TrimSpace(pid) + t.Logf("Start %s cgroup set", pid) + CgroupSet(pid, 1) + CgroupCleanAll() +}