mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
client:
1.try other hosts for network errors 2.use single instance for File and Dir object metanode: 1.use tcp packet notify all raft followers to free inodes datanode: 1.remove read check appliedID 2.change extentcache fp to 100 3.fix always repair has delete extents repl: 1.modify funcName from reciveFromFollower to checkLocalResultAndReceiveFromFollower
This commit is contained in:
parent
879c55295a
commit
493af438b0
@ -54,7 +54,7 @@ var (
|
||||
)
|
||||
|
||||
// NewDir returns a new directory.
|
||||
func NewDir(s *Super, i *Inode) *Dir {
|
||||
func NewDir(s *Super, i *Inode) fs.Node {
|
||||
return &Dir{
|
||||
super: s,
|
||||
inode: i,
|
||||
@ -88,6 +88,10 @@ func (d *Dir) Create(ctx context.Context, req *fuse.CreateRequest, resp *fuse.Cr
|
||||
child := NewFile(d.super, inode)
|
||||
d.super.ec.OpenStream(inode.ino)
|
||||
|
||||
d.super.fslock.Lock()
|
||||
d.super.nodeCache[inode.ino] = child
|
||||
d.super.fslock.Unlock()
|
||||
|
||||
elapsed := time.Since(start)
|
||||
log.LogDebugf("TRACE Create: parent(%v) req(%v) resp(%v) ino(%v) (%v)ns", d.inode.ino, req, resp, inode.ino, elapsed.Nanoseconds())
|
||||
return child, child, nil
|
||||
@ -99,6 +103,10 @@ func (d *Dir) Forget() {
|
||||
defer func() {
|
||||
log.LogDebugf("TRACE Forget: ino(%v)", ino)
|
||||
}()
|
||||
|
||||
d.super.fslock.Lock()
|
||||
delete(d.super.nodeCache, ino)
|
||||
d.super.fslock.Unlock()
|
||||
}
|
||||
|
||||
// Mkdir handles the mkdir request.
|
||||
@ -114,6 +122,10 @@ func (d *Dir) Mkdir(ctx context.Context, req *fuse.MkdirRequest) (fs.Node, error
|
||||
d.super.ic.Put(inode)
|
||||
child := NewDir(d.super, inode)
|
||||
|
||||
d.super.fslock.Lock()
|
||||
d.super.nodeCache[inode.ino] = child
|
||||
d.super.fslock.Unlock()
|
||||
|
||||
elapsed := time.Since(start)
|
||||
log.LogDebugf("TRACE Mkdir: parent(%v) req(%v) ino(%v) (%v)ns", d.inode.ino, req, inode.ino, elapsed.Nanoseconds())
|
||||
return child, nil
|
||||
@ -171,12 +183,17 @@ func (d *Dir) Lookup(ctx context.Context, req *fuse.LookupRequest, resp *fuse.Lo
|
||||
}
|
||||
mode := inode.mode
|
||||
|
||||
var child fs.Node
|
||||
if mode.IsDir() {
|
||||
child = NewDir(d.super, inode)
|
||||
} else {
|
||||
child = NewFile(d.super, inode)
|
||||
d.super.fslock.Lock()
|
||||
child, ok := d.super.nodeCache[ino]
|
||||
if !ok {
|
||||
if mode.IsDir() {
|
||||
child = NewDir(d.super, inode)
|
||||
} else {
|
||||
child = NewFile(d.super, inode)
|
||||
}
|
||||
d.super.nodeCache[ino] = child
|
||||
}
|
||||
d.super.fslock.Unlock()
|
||||
|
||||
resp.EntryValid = LookupValidDuration
|
||||
return child, nil
|
||||
@ -276,6 +293,10 @@ func (d *Dir) Symlink(ctx context.Context, req *fuse.SymlinkRequest) (fs.Node, e
|
||||
d.super.ic.Put(inode)
|
||||
child := NewFile(d.super, inode)
|
||||
|
||||
d.super.fslock.Lock()
|
||||
d.super.nodeCache[inode.ino] = child
|
||||
d.super.fslock.Unlock()
|
||||
|
||||
elapsed := time.Since(start)
|
||||
log.LogDebugf("TRACE Symlink: parent(%v) req(%v) ino(%v) (%v)ns", parentIno, req, inode.ino, elapsed.Nanoseconds())
|
||||
return child, nil
|
||||
@ -306,7 +327,14 @@ func (d *Dir) Link(ctx context.Context, req *fuse.LinkRequest, old fs.Node) (fs.
|
||||
|
||||
newInode := NewInode(info)
|
||||
d.super.ic.Put(newInode)
|
||||
newFile := NewFile(d.super, newInode)
|
||||
|
||||
d.super.fslock.Lock()
|
||||
newFile, ok := d.super.nodeCache[newInode.ino]
|
||||
if !ok {
|
||||
newFile = NewFile(d.super, newInode)
|
||||
d.super.nodeCache[newInode.ino] = newFile
|
||||
}
|
||||
d.super.fslock.Unlock()
|
||||
|
||||
elapsed := time.Since(start)
|
||||
log.LogDebugf("TRACE Link: parent(%v) name(%v) ino(%v) (%v)ns", d.inode.ino, req.NewName, newInode.ino, elapsed.Nanoseconds())
|
||||
|
||||
@ -55,7 +55,7 @@ var (
|
||||
)
|
||||
|
||||
// NewFile returns a new file.
|
||||
func NewFile(s *Super, i *Inode) *File {
|
||||
func NewFile(s *Super, i *Inode) fs.Node {
|
||||
return &File{super: s, inode: i}
|
||||
}
|
||||
|
||||
@ -86,6 +86,10 @@ func (f *File) Forget() {
|
||||
log.LogDebugf("TRACE Forget: ino(%v)", ino)
|
||||
}()
|
||||
|
||||
f.super.fslock.Lock()
|
||||
delete(f.super.nodeCache, ino)
|
||||
f.super.fslock.Unlock()
|
||||
|
||||
if err := f.super.ec.EvictStream(ino); err != nil {
|
||||
log.LogWarnf("Forget: stream not ready to evict, ino(%v) err(%v)", ino, err)
|
||||
return
|
||||
|
||||
@ -18,6 +18,7 @@ import (
|
||||
"fmt"
|
||||
"github.com/juju/errors"
|
||||
"golang.org/x/net/context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"bazil.org/fuse"
|
||||
@ -37,6 +38,9 @@ type Super struct {
|
||||
ec *stream.ExtentClient
|
||||
orphan *OrphanInodeList
|
||||
enSyncWrite bool
|
||||
|
||||
nodeCache map[uint64]fs.Node
|
||||
fslock sync.Mutex
|
||||
}
|
||||
|
||||
// Functions that Super needs to implement
|
||||
@ -76,6 +80,7 @@ func NewSuper(volname, owner, master string, icacheTimeout, lookupValid, attrVal
|
||||
}
|
||||
s.ic = NewInodeCache(inodeExpiration, MaxInodeCache)
|
||||
s.orphan = NewOrphanInodeList()
|
||||
s.nodeCache = make(map[uint64]fs.Node)
|
||||
log.LogInfof("NewSuper: cluster(%v) volname(%v) icacheExpiration(%v) LookupValidDuration(%v) AttrValidDuration(%v)", s.cluster, s.volname, inodeExpiration, LookupValidDuration, AttrValidDuration)
|
||||
return s, nil
|
||||
}
|
||||
|
||||
@ -1,12 +1,13 @@
|
||||
{
|
||||
"mountpoint": "/mnt/fuse",
|
||||
"volname": "intest",
|
||||
"master": "10.196.31.173:80,10.196.31.141:80,10.196.30.200:80",
|
||||
"logpath": "/export/Logs/baudstorage",
|
||||
"loglvl": "info",
|
||||
"mountPoint": "/mnt/fuse",
|
||||
"volName": "test",
|
||||
"owner": "cfs",
|
||||
"masterAddr": "10.196.31.173:80,10.196.31.141:80,10.196.30.200:80",
|
||||
"logDir": "/export/Logs/cfs",
|
||||
"logLevel": "info",
|
||||
"consulAddr": "http://cbconsul-cfs01.cbmonitor.svc.ht7.n.jd.local",
|
||||
"exporterPort": 9513,
|
||||
"profport": "10094"
|
||||
"profPort": "10094"
|
||||
}
|
||||
|
||||
|
||||
|
||||
@ -165,14 +165,6 @@ func (dp *DataPartition) CheckLeader(request *repl.Packet, connect net.Conn) (er
|
||||
return
|
||||
}
|
||||
|
||||
if dp.appliedID < dp.maxAppliedID {
|
||||
err = storage.TryAgainError
|
||||
logContent := fmt.Sprintf("action[ReadCheck] %v localID=%v maxID=%v.",
|
||||
request.LogMessage(request.GetOpMsg(), connect.RemoteAddr().String(), request.StartT, nil), dp.appliedID, dp.maxAppliedID)
|
||||
log.LogErrorf(logContent)
|
||||
return
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@ -176,14 +176,6 @@ func (s *DataNode) parseConfig(cfg *config.Config) (err error) {
|
||||
LocalIP = cfg.GetString(ConfigKeyLocalIP)
|
||||
port = cfg.GetString(ConfigKeyPort)
|
||||
if regexpPort, err = regexp.Compile("^(\\d)+$"); err != nil {
|
||||
return
|
||||
}
|
||||
if !regexpPort.MatchString(port) {
|
||||
err = ErrBadConfFile
|
||||
return
|
||||
}
|
||||
s.port = port
|
||||
if len(cfg.GetArray(ConfigKeyMasterAddr)) == 0 {
|
||||
return fmt.Errorf("Err:no port")
|
||||
}
|
||||
if !regexpPort.MatchString(port) {
|
||||
@ -316,6 +308,7 @@ func (s *DataNode) registerHandler() {
|
||||
http.HandleFunc("/extent", s.getExtentAPI)
|
||||
http.HandleFunc("/block", s.getBlockCrcAPI)
|
||||
http.HandleFunc("/stats", s.getStatAPI)
|
||||
http.HandleFunc("/raftStatus", s.getRaftStatus)
|
||||
}
|
||||
|
||||
func (s *DataNode) startTCPService() (err error) {
|
||||
|
||||
@ -67,6 +67,25 @@ func (s *DataNode) getStatAPI(w http.ResponseWriter, r *http.Request) {
|
||||
s.buildSuccessResp(w, response)
|
||||
}
|
||||
|
||||
func (s *DataNode) getRaftStatus(w http.ResponseWriter, r *http.Request) {
|
||||
const (
|
||||
paramRaftID = "raftID"
|
||||
)
|
||||
if err := r.ParseForm(); err != nil {
|
||||
err = fmt.Errorf("parse form fail: %v", err)
|
||||
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
|
||||
return
|
||||
}
|
||||
raftID, err := strconv.ParseUint(r.FormValue(paramRaftID), 10, 64)
|
||||
if err != nil {
|
||||
err = fmt.Errorf("parse param %v fail: %v", paramRaftID, err)
|
||||
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
|
||||
return
|
||||
}
|
||||
raftStatus := s.raftStore.RaftStatus(raftID)
|
||||
s.buildSuccessResp(w, raftStatus)
|
||||
}
|
||||
|
||||
func (s *DataNode) getPartitionsAPI(w http.ResponseWriter, r *http.Request) {
|
||||
partitions := make([]interface{}, 0)
|
||||
s.space.RangePartitions(func(dp *DataPartition) bool {
|
||||
|
||||
@ -74,6 +74,8 @@ func (m *metadataManager) HandleMetadataOperation(conn net.Conn, p *Packet,
|
||||
err = m.opCreateInode(conn, p, remoteAddr)
|
||||
case proto.OpMetaLinkInode:
|
||||
err = m.opMetaLinkInode(conn, p, remoteAddr)
|
||||
case proto.OpMetaFreeInodesOnRaftFollower:
|
||||
err = m.opFreeInodeOnRaftFollower(conn, p, remoteAddr)
|
||||
case proto.OpMetaUnlinkInode:
|
||||
err = m.opMetaUnlinkInode(conn, p, remoteAddr)
|
||||
case proto.OpMetaInodeGet:
|
||||
|
||||
@ -169,6 +169,22 @@ func (m *metadataManager) opMetaLinkInode(conn net.Conn, p *Packet,
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpCreate
|
||||
func (m *metadataManager) opFreeInodeOnRaftFollower(conn net.Conn, p *Packet,
|
||||
remoteAddr string) (err error) {
|
||||
mp, err := m.getPartition(p.PartitionID)
|
||||
if err != nil {
|
||||
p.PacketErrorWithBody(proto.OpErr, []byte(err.Error()))
|
||||
m.respondToClient(conn, p)
|
||||
return
|
||||
}
|
||||
mp.(*metaPartition).internalDelete(p.Data[:p.Size])
|
||||
p.PacketOkReply()
|
||||
m.respondToClient(conn, p)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpCreate
|
||||
func (m *metadataManager) opCreateDentry(conn net.Conn, p *Packet,
|
||||
remoteAddr string) (err error) {
|
||||
|
||||
@ -44,3 +44,18 @@ func NewPacketToDeleteExtent(dp *DataPartition, ext *proto.ExtentKey) *Packet {
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
// NewPacketToDeleteExtent returns a new packet to delete the extent.
|
||||
func NewPacketToFreeInodeOnRaftFollower(partitionID uint64, freeInodes []byte) *Packet {
|
||||
p := new(Packet)
|
||||
p.Magic = proto.ProtoMagic
|
||||
p.Opcode = proto.OpMetaFreeInodesOnRaftFollower
|
||||
p.PartitionID = partitionID
|
||||
p.ExtentType = proto.NormalExtentType
|
||||
p.ReqID = proto.GenerateRequestID()
|
||||
p.Data = make([]byte, len(freeInodes))
|
||||
copy(p.Data, freeInodes)
|
||||
p.Size = uint32(len(p.Data))
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
@ -360,6 +360,17 @@ func (mp *metaPartition) IsLeader() (leaderAddr string, ok bool) {
|
||||
return
|
||||
}
|
||||
|
||||
func (mp *metaPartition) GetPeers() (peers []string) {
|
||||
peers = make([]string, 0)
|
||||
for _, peer := range mp.config.Peers {
|
||||
if mp.config.NodeId == peer.ID {
|
||||
continue
|
||||
}
|
||||
peers = append(peers, peer.Addr)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// GetCursor returns the cursor stored in the config.
|
||||
func (mp *metaPartition) GetCursor() uint64 {
|
||||
return mp.config.Cursor
|
||||
|
||||
@ -20,6 +20,7 @@ import (
|
||||
"github.com/chubaofs/cfs/proto"
|
||||
"github.com/chubaofs/cfs/util/log"
|
||||
"github.com/juju/errors"
|
||||
"net"
|
||||
"os"
|
||||
"path"
|
||||
"runtime"
|
||||
@ -204,13 +205,12 @@ func (mp *metaPartition) deleteMarkedInodes(inoSlice []*Inode) {
|
||||
for _, ino := range shouldCommit {
|
||||
bufSlice = append(bufSlice, ino.MarshalKey()...)
|
||||
}
|
||||
// raft Commit
|
||||
_, err := mp.Put(opFSMInternalDeleteInode, bufSlice)
|
||||
err := mp.syncToRaftFollowersFreeInode(bufSlice)
|
||||
if err != nil {
|
||||
for _, ino := range shouldCommit {
|
||||
mp.freeList.Push(ino)
|
||||
}
|
||||
log.LogWarnf("[deleteInodeTree] raft commit inode list: %v, "+
|
||||
log.LogWarnf("[deleteInodeTreeOnRaftPeers] syncToRaftFollowersFreeInode inode list: %v, "+
|
||||
"response %s", shouldCommit, err.Error())
|
||||
}
|
||||
log.LogDebugf("[deleteInodeTree] inode list: %v", shouldCommit)
|
||||
@ -218,6 +218,56 @@ func (mp *metaPartition) deleteMarkedInodes(inoSlice []*Inode) {
|
||||
|
||||
}
|
||||
|
||||
func (mp *metaPartition) syncToRaftFollowersFreeInode(hasDeleteInodes []byte) (err error) {
|
||||
raftPeers := mp.GetPeers()
|
||||
raftPeersError := make([]error, len(raftPeers))
|
||||
wg := new(sync.WaitGroup)
|
||||
for index, target := range raftPeers {
|
||||
wg.Add(1)
|
||||
raftPeersError[index] = mp.notifyRaftFollowerToFreeInodes(wg, target, hasDeleteInodes)
|
||||
}
|
||||
wg.Wait()
|
||||
for index := 0; index < len(raftPeersError); index++ {
|
||||
if raftPeersError[index] != nil {
|
||||
err = raftPeersError[index]
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (mp *metaPartition) notifyRaftFollowerToFreeInodes(wg *sync.WaitGroup, target string, hasDeleteInodes []byte) (err error) {
|
||||
var conn *net.TCPConn
|
||||
conn, err = mp.config.ConnPool.GetConnect(target)
|
||||
defer func() {
|
||||
wg.Done()
|
||||
if err != nil {
|
||||
log.LogWarnf(err.Error())
|
||||
mp.config.ConnPool.PutConnect(conn, ForceClosedConnect)
|
||||
} else {
|
||||
mp.config.ConnPool.PutConnect(conn, NoClosedConnect)
|
||||
}
|
||||
}()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
request := NewPacketToFreeInodeOnRaftFollower(mp.config.PartitionId, hasDeleteInodes)
|
||||
if err = request.WriteToConn(conn); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
if err = request.ReadFromConn(conn, proto.NoReadDeadlineTime); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
if request.ResultCode != proto.OpOk {
|
||||
err = fmt.Errorf("request(%v) error(%v)", request.GetUniqueLogId(), string(request.Data[request.Size]))
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (mp *metaPartition) doDeleteMarkedInodes(ext *proto.ExtentKey) (err error) {
|
||||
// get the data node view
|
||||
dp := mp.vol.GetPartition(ext.PartitionId)
|
||||
@ -228,7 +278,14 @@ func (mp *metaPartition) doDeleteMarkedInodes(ext *proto.ExtentKey) (err error)
|
||||
}
|
||||
// delete the data node
|
||||
conn, err := mp.config.ConnPool.GetConnect(dp.Hosts[0])
|
||||
defer mp.config.ConnPool.PutConnect(conn, ForceClosedConnect)
|
||||
|
||||
defer func() {
|
||||
if err != nil {
|
||||
mp.config.ConnPool.PutConnect(conn, ForceClosedConnect)
|
||||
} else {
|
||||
mp.config.ConnPool.PutConnect(conn, NoClosedConnect)
|
||||
}
|
||||
}()
|
||||
|
||||
if err != nil {
|
||||
err = errors.Errorf("get conn from pool %s, "+
|
||||
|
||||
@ -18,6 +18,7 @@ import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"github.com/chubaofs/cfs/proto"
|
||||
"github.com/chubaofs/cfs/util/log"
|
||||
"io"
|
||||
)
|
||||
|
||||
@ -149,6 +150,7 @@ func (mp *metaPartition) internalDelete(val []byte) (err error) {
|
||||
}
|
||||
return
|
||||
}
|
||||
log.LogDebugf("recive raftLeader free inode(%v)", ino.Inode)
|
||||
mp.internalDeleteInode(ino)
|
||||
}
|
||||
}
|
||||
|
||||
@ -83,6 +83,9 @@ const (
|
||||
OpMetaSetattr uint8 = 0x30
|
||||
OpMetaReleaseOpen uint8 = 0x31
|
||||
|
||||
//Operations: MetaNode Leader -> MetaNode Follower
|
||||
OpMetaFreeInodesOnRaftFollower uint8 = 0x32
|
||||
|
||||
// Operations: Master -> MetaNode
|
||||
OpCreateMetaPartition uint8 = 0x40
|
||||
OpMetaNodeHeartbeat uint8 = 0x41
|
||||
|
||||
@ -32,6 +32,7 @@ type RaftStore interface {
|
||||
CreatePartition(cfg *PartitionConfig) (Partition, error)
|
||||
Stop()
|
||||
RaftConfig() *raft.Config
|
||||
RaftStatus(raftID uint64) (raftStatus *raft.Status)
|
||||
NodeManager
|
||||
}
|
||||
|
||||
@ -48,6 +49,10 @@ func (s *raftStore) RaftConfig() *raft.Config {
|
||||
return s.raftConfig
|
||||
}
|
||||
|
||||
func (s *raftStore) RaftStatus(raftID uint64) (raftStatus *raft.Status) {
|
||||
return s.raftServer.Status(raftID)
|
||||
}
|
||||
|
||||
// AddNode adds a new node to the raft store.
|
||||
func (s *raftStore) AddNode(nodeID uint64, addr string) {
|
||||
if s.resolver != nil {
|
||||
|
||||
@ -256,7 +256,7 @@ func (rp *ReplProtocol) reciveAllFollowerResponse() {
|
||||
rp.deletePacket(request, e)
|
||||
}()
|
||||
for index := 0; index < len(request.followersAddrs); index++ {
|
||||
err := rp.receiveFromFollower(request, index)
|
||||
err := rp.checkLocalResultAndReceiveFromFollower(request, index)
|
||||
if err != nil {
|
||||
rp.setReplProtocolError(request, index)
|
||||
request.PackErrorBody(ActionReceiveFromFollower, err.Error())
|
||||
@ -268,7 +268,7 @@ func (rp *ReplProtocol) reciveAllFollowerResponse() {
|
||||
}
|
||||
|
||||
// Read the response from the follower
|
||||
func (rp *ReplProtocol) receiveFromFollower(request *Packet, index int) (err error) {
|
||||
func (rp *ReplProtocol) checkLocalResultAndReceiveFromFollower(request *Packet, index int) (err error) {
|
||||
// Receive p response from one member
|
||||
if request.followerConns[index] == nil {
|
||||
err = fmt.Errorf(ConnIsNullErr)
|
||||
|
||||
@ -65,7 +65,9 @@ func (reader *ExtentReader) Read(req *ExtentRequest) (readBytes int, err error)
|
||||
replyPacket.Data = req.Data[readBytes : readBytes+bufSize]
|
||||
e := replyPacket.readFromConn(conn, proto.ReadDeadlineTime)
|
||||
if e != nil {
|
||||
return errors.Annotatef(e, "Extent Reader Read: failed to read from connect, readBytes(%v)", readBytes), false
|
||||
log.LogErrorf("Extent Reader Read: failed to read from connect, readBytes(%v) err(%v)", readBytes, e)
|
||||
// Upon receiving NotALeaderError, other hosts will be retried.
|
||||
return NotALeaderError, false
|
||||
}
|
||||
|
||||
//log.LogDebugf("ExtentReader Read: ResultCode(%v) req(%v) reply(%v) readBytes(%v)", replyPacket.GetResultMsg(), reqPacket, replyPacket, readBytes)
|
||||
|
||||
@ -332,7 +332,9 @@ func (s *Streamer) doOverwrite(req *ExtentRequest, direct bool) (total int, err
|
||||
err = sc.Send(reqPacket, func(conn *net.TCPConn) (error, bool) {
|
||||
e := replyPacket.ReadFromConn(conn, proto.ReadDeadlineTime)
|
||||
if e != nil {
|
||||
return errors.Annotatef(e, "Stream Writer doOverwrite: ino(%v) failed to read from connect", s.inode), false
|
||||
log.LogErrorf("Stream Writer doOverwrite: ino(%v) failed to read from connect, err(%v)", s.inode, e)
|
||||
// Upon receiving NotALeaderError, other hosts will be retried.
|
||||
return NotALeaderError, false
|
||||
}
|
||||
|
||||
if replyPacket.ResultCode == proto.OpAgain {
|
||||
|
||||
@ -39,18 +39,17 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
ExtCrcHeaderFileName = "EXTENT_CRC"
|
||||
ExtBaseExtentIDFileName = "EXTENT_META"
|
||||
TinyDeleteFileOpt = os.O_CREATE | os.O_RDWR
|
||||
TinyExtDeletedFileName = "TINYEXTENT_DELETE"
|
||||
NormalExtDeletedFileName = "NORMALEXTENT_DELETE"
|
||||
MaxExtentCount = 20000
|
||||
TinyExtentCount = 64
|
||||
TinyExtentStartID = 1
|
||||
MinExtentID = 1024
|
||||
EveryTinyDeleteRecordSize = 24
|
||||
UpdateCrcInterval = 600
|
||||
PerAllocExtentIDCntOnVerifyFile = 1000
|
||||
ExtCrcHeaderFileName = "EXTENT_CRC"
|
||||
ExtBaseExtentIDFileName = "EXTENT_META"
|
||||
TinyDeleteFileOpt = os.O_CREATE | os.O_RDWR
|
||||
TinyExtDeletedFileName = "TINYEXTENT_DELETE"
|
||||
NormalExtDeletedFileName = "NORMALEXTENT_DELETE"
|
||||
MaxExtentCount = 20000
|
||||
TinyExtentCount = 64
|
||||
TinyExtentStartID = 1
|
||||
MinExtentID = 1024
|
||||
EveryTinyDeleteRecordSize = 24
|
||||
UpdateCrcInterval = 600
|
||||
)
|
||||
|
||||
var (
|
||||
@ -120,7 +119,7 @@ type ExtentStore struct {
|
||||
blockSize int
|
||||
partitionID uint64
|
||||
verifyExtentFp *os.File
|
||||
HasAllocSpaceExtentIDOnVerfiyFile uint64
|
||||
hasAllocSpaceExtentIDOnVerfiyFile uint64
|
||||
}
|
||||
|
||||
func MkdirAll(name string) (err error) {
|
||||
@ -153,12 +152,12 @@ func NewExtentStore(dataDir string, partitionID uint64, storeSize int) (s *Exten
|
||||
}
|
||||
|
||||
s.extentInfoMap = make(map[uint64]*ExtentInfo, 0)
|
||||
s.cache = NewExtentCache(20)
|
||||
s.cache = NewExtentCache(100)
|
||||
if err = s.initBaseFileID(); err != nil {
|
||||
err = fmt.Errorf("init base field ID: %v", err)
|
||||
return
|
||||
}
|
||||
s.HasAllocSpaceExtentIDOnVerfiyFile = s.GetPreAllocSpaceExtentIDOnVerfiyFile()
|
||||
s.hasAllocSpaceExtentIDOnVerfiyFile = s.GetPreAllocSpaceExtentIDOnVerfiyFile()
|
||||
s.storeSize = storeSize
|
||||
s.closeC = make(chan bool, 1)
|
||||
s.closed = false
|
||||
@ -287,18 +286,13 @@ func (s *ExtentStore) initBaseFileID() (err error) {
|
||||
// Write writes the given extent to the disk.
|
||||
func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, crc uint32, isUpdateSize bool, isSync bool) (err error) {
|
||||
var (
|
||||
has bool
|
||||
e *Extent
|
||||
ei *ExtentInfo
|
||||
e *Extent
|
||||
ei *ExtentInfo
|
||||
)
|
||||
s.eiMutex.RLock()
|
||||
ei, has = s.extentInfoMap[extentID]
|
||||
ei, _ = s.extentInfoMap[extentID]
|
||||
s.eiMutex.RUnlock()
|
||||
if !has {
|
||||
err = ExtentNotFoundError
|
||||
return
|
||||
}
|
||||
e, err = s.extentWithHeader(extentID)
|
||||
e, err = s.extentWithHeader(ei)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@ -339,7 +333,10 @@ func IsTinyExtent(extentID uint64) bool {
|
||||
// Read reads the extent based on the given id.
|
||||
func (s *ExtentStore) Read(extentID uint64, offset, size int64, nbuf []byte, isRepairRead bool) (crc uint32, err error) {
|
||||
var e *Extent
|
||||
if e, err = s.extentWithHeader(extentID); err != nil {
|
||||
s.eiMutex.RLock()
|
||||
ei := s.extentInfoMap[extentID]
|
||||
s.eiMutex.RUnlock()
|
||||
if e, err = s.extentWithHeader(ei); err != nil {
|
||||
return
|
||||
}
|
||||
if err = s.checkOffsetAndSize(extentID, offset, size); err != nil {
|
||||
@ -366,19 +363,14 @@ func (s *ExtentStore) tinyDelete(e *Extent, offset, size, tinyDeleteFileOffset i
|
||||
// MarkDelete marks the given extent as deleted.
|
||||
func (s *ExtentStore) MarkDelete(extentID uint64, offset, size, tinyDeleteFileOffset int64) (err error) {
|
||||
var (
|
||||
e *Extent
|
||||
ei *ExtentInfo
|
||||
has bool
|
||||
e *Extent
|
||||
ei *ExtentInfo
|
||||
)
|
||||
|
||||
s.eiMutex.RLock()
|
||||
ei, has = s.extentInfoMap[extentID]
|
||||
ei = s.extentInfoMap[extentID]
|
||||
s.eiMutex.RUnlock()
|
||||
if !has {
|
||||
return
|
||||
}
|
||||
|
||||
if e, err = s.extentWithHeader(extentID); err != nil {
|
||||
if e, err = s.extentWithHeader(ei); err != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
@ -393,10 +385,8 @@ func (s *ExtentStore) MarkDelete(extentID uint64, offset, size, tinyDeleteFileOf
|
||||
}
|
||||
s.PersistenceHasDeleteExtent(extentID)
|
||||
ei.IsDeleted = true
|
||||
ei.ModifyTime = time.Now().Unix()
|
||||
s.cache.Del(e.extentID)
|
||||
s.eiMutex.Lock()
|
||||
delete(s.extentInfoMap, extentID)
|
||||
s.eiMutex.Unlock()
|
||||
s.DeleteBlockCrc(extentID)
|
||||
|
||||
return
|
||||
@ -678,7 +668,22 @@ func (s *ExtentStore) extent(extentID uint64) (e *Extent, err error) {
|
||||
return
|
||||
}
|
||||
|
||||
func (s *ExtentStore) extentWithHeader(extentID uint64) (e *Extent, err error) {
|
||||
func (s *ExtentStore) extentWithHeader(ei *ExtentInfo) (e *Extent, err error) {
|
||||
var ok bool
|
||||
if ei == nil || ei.IsDeleted {
|
||||
err = ExtentNotFoundError
|
||||
return
|
||||
}
|
||||
if e, ok = s.cache.Get(ei.FileID); !ok {
|
||||
if e, err = s.loadExtentFromDisk(ei.FileID, true); err != nil {
|
||||
err = fmt.Errorf("load %v from disk: %v", s.getExtentKey(ei.FileID), err)
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *ExtentStore) extentWithHeaderByExtentID(extentID uint64) (e *Extent, err error) {
|
||||
var ok bool
|
||||
if e, ok = s.cache.Get(extentID); !ok {
|
||||
if e, err = s.loadExtentFromDisk(extentID, true); err != nil {
|
||||
@ -743,9 +748,10 @@ func (s *ExtentStore) loopAutoComputeExtentCrc() {
|
||||
func (s *ExtentStore) ScanBlocks(extentID uint64) (bcs []*BlockCrc, err error) {
|
||||
var blockCnt int
|
||||
bcs = make([]*BlockCrc, 0)
|
||||
e, err := s.extentWithHeader(extentID)
|
||||
ei := s.extentInfoMap[extentID]
|
||||
e, err := s.extentWithHeader(ei)
|
||||
if err != nil {
|
||||
return
|
||||
return bcs, err
|
||||
}
|
||||
blockCnt = int(e.Size() / util.BlockSize)
|
||||
if e.Size()%util.BlockSize != 0 {
|
||||
@ -768,12 +774,24 @@ func (arr ExtentInfoArr) Swap(i, j int) { arr[i], arr[j] = arr[j], arr[i] }
|
||||
|
||||
func (s *ExtentStore) autoComputeExtentCrc() {
|
||||
extentInfos := make([]*ExtentInfo, 0)
|
||||
deleteExtents := make([]*ExtentInfo, 0)
|
||||
s.eiMutex.RLock()
|
||||
for _, ei := range s.extentInfoMap {
|
||||
extentInfos = append(extentInfos, ei)
|
||||
if ei.IsDeleted && time.Now().Unix()-ei.ModifyTime > UpdateCrcInterval {
|
||||
deleteExtents = append(deleteExtents, ei)
|
||||
}
|
||||
}
|
||||
s.eiMutex.RUnlock()
|
||||
|
||||
if len(deleteExtents) > 0 {
|
||||
s.eiMutex.Lock()
|
||||
for _, ei := range deleteExtents {
|
||||
delete(s.extentInfoMap, ei.FileID)
|
||||
}
|
||||
s.eiMutex.Unlock()
|
||||
}
|
||||
|
||||
sort.Sort(ExtentInfoArr(extentInfos))
|
||||
|
||||
for _, ei := range extentInfos {
|
||||
@ -784,7 +802,7 @@ func (s *ExtentStore) autoComputeExtentCrc() {
|
||||
if atomic.LoadUint32(&ei.Crc) == 0 {
|
||||
continue
|
||||
}
|
||||
e, err := s.extentWithHeader(ei.FileID)
|
||||
e, err := s.extentWithHeader(ei)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
@ -62,8 +62,8 @@ func (s *ExtentStore) GetPreAllocSpaceExtentIDOnVerfiyFile() (extentID uint64) {
|
||||
}
|
||||
|
||||
func (s *ExtentStore) PreAllocSpaceOnVerfiyFile(currExtentID uint64) {
|
||||
if currExtentID > atomic.LoadUint64(&s.HasAllocSpaceExtentIDOnVerfiyFile) {
|
||||
prevAllocSpaceExtentID := int64(atomic.LoadUint64(&s.HasAllocSpaceExtentIDOnVerfiyFile))
|
||||
if currExtentID > atomic.LoadUint64(&s.hasAllocSpaceExtentIDOnVerfiyFile) {
|
||||
prevAllocSpaceExtentID := int64(atomic.LoadUint64(&s.hasAllocSpaceExtentIDOnVerfiyFile))
|
||||
endAllocSpaceExtentID := int64(prevAllocSpaceExtentID + 1000)
|
||||
size := int64(1000 * util.BlockHeaderSize)
|
||||
err := syscall.Fallocate(int(s.verifyExtentFp.Fd()), 1, prevAllocSpaceExtentID*util.BlockHeaderSize, size)
|
||||
@ -75,7 +75,7 @@ func (s *ExtentStore) PreAllocSpaceOnVerfiyFile(currExtentID uint64) {
|
||||
if _, err = s.metadataFp.WriteAt(data, 8); err != nil {
|
||||
return
|
||||
}
|
||||
atomic.StoreUint64(&s.HasAllocSpaceExtentIDOnVerfiyFile, uint64(endAllocSpaceExtentID))
|
||||
atomic.StoreUint64(&s.hasAllocSpaceExtentIDOnVerfiyFile, uint64(endAllocSpaceExtentID))
|
||||
log.LogInfof("Action(PreAllocSpaceOnVerfiyFile) PartitionID(%v) currentExtent(%v)"+
|
||||
"PrevAllocSpaceExtentIDOnVerifyFile(%v) EndAllocSpaceExtentIDOnVerifyFile(%v)"+
|
||||
" has allocSpaceOnVerifyFile to (%v)", s.partitionID, currExtentID, prevAllocSpaceExtentID, endAllocSpaceExtentID,
|
||||
|
||||
6
vendor/github.com/tiglabs/raft/raft.go
generated
vendored
6
vendor/github.com/tiglabs/raft/raft.go
generated
vendored
@ -538,6 +538,12 @@ func (s *raft) sendMessage(m *proto.Message) {
|
||||
|
||||
func (s *raft) maybeChange(respErr bool) {
|
||||
updated := false
|
||||
if s.prevSoftSt.term != s.raftFsm.term && s.raftFsm.leader == s.config.NodeID &&
|
||||
s.curApplied.Get() < s.raftFsm.raftLog.committed {
|
||||
logger.Warn("raft:[%v] changed leader wait applied. curApplied %v committed %v at term %d.",
|
||||
s.raftFsm.id, s.curApplied.Get(), s.raftFsm.raftLog.committed, s.raftFsm.term)
|
||||
return
|
||||
}
|
||||
if s.prevSoftSt.term != s.raftFsm.term {
|
||||
updated = true
|
||||
s.prevSoftSt.term = s.raftFsm.term
|
||||
|
||||
Loading…
Reference in New Issue
Block a user