Fix file system access for Windows

This commit is contained in:
Frank Szendzielarz 2019-06-06 12:11:46 +02:00
parent 30c2b1b06d
commit e7feb2f869

View file

@ -122,7 +122,7 @@ func newCustomTable(path string, name string, readMeter metrics.Meter, writeMete
// compressed idx // compressed idx
idxName = fmt.Sprintf("%s.cidx", name) idxName = fmt.Sprintf("%s.cidx", name)
} }
offsets, err := os.OpenFile(filepath.Join(path, idxName), os.O_RDWR|os.O_CREATE|os.O_APPEND, 0644) offsets, err := os.OpenFile(filepath.Join(path, idxName), os.O_RDWR|os.O_CREATE, 0644)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@ -188,7 +188,7 @@ func (t *freezerTable) repair() error {
t.index.ReadAt(buffer, offsetsSize-indexEntrySize) t.index.ReadAt(buffer, offsetsSize-indexEntrySize)
lastIndex.unmarshalBinary(buffer) lastIndex.unmarshalBinary(buffer)
t.head, err = t.openFile(lastIndex.filenum, os.O_RDWR|os.O_CREATE|os.O_APPEND) t.head, err = t.openFile(lastIndex.filenum, os.O_RDWR|os.O_CREATE)
if err != nil { if err != nil {
return err return err
} }
@ -223,7 +223,7 @@ func (t *freezerTable) repair() error {
if newLastIndex.filenum != lastIndex.filenum { if newLastIndex.filenum != lastIndex.filenum {
// release earlier opened file // release earlier opened file
t.releaseFile(lastIndex.filenum) t.releaseFile(lastIndex.filenum)
t.head, err = t.openFile(newLastIndex.filenum, os.O_RDWR|os.O_CREATE|os.O_APPEND) t.head, err = t.openFile(newLastIndex.filenum, os.O_RDWR|os.O_CREATE)
if stat, err = t.head.Stat(); err != nil { if stat, err = t.head.Stat(); err != nil {
// TODO, anything more we can do here? // TODO, anything more we can do here?
// A data file has gone missing... // A data file has gone missing...
@ -269,11 +269,11 @@ func (t *freezerTable) preopen() (err error) {
} }
} }
// Open head in read/write // Open head in read/write
t.head, err = t.openFile(t.headId, os.O_RDWR|os.O_CREATE|os.O_APPEND) t.head, err = t.openFile(t.headId, os.O_RDWR|os.O_CREATE)
return err return err
} }
// truncate discards any recent data above the provided threashold number. // truncate discards any recent data above the provided threshold number.
func (t *freezerTable) truncate(items uint64) error { func (t *freezerTable) truncate(items uint64) error {
t.lock.Lock() t.lock.Lock()
defer t.lock.Unlock() defer t.lock.Unlock()
@ -299,7 +299,7 @@ func (t *freezerTable) truncate(items uint64) error {
if expected.filenum != t.headId { if expected.filenum != t.headId {
// If already open for reading, force-reopen for writing // If already open for reading, force-reopen for writing
t.releaseFile(expected.filenum) t.releaseFile(expected.filenum)
newHead, err := t.openFile(expected.filenum, os.O_RDWR|os.O_CREATE|os.O_APPEND) newHead, err := t.openFile(expected.filenum, os.O_RDWR|os.O_CREATE)
if err != nil { if err != nil {
return err return err
} }
@ -435,6 +435,10 @@ func (t *freezerTable) Append(item uint64, blob []byte) error {
defer t.lock.RUnlock() defer t.lock.RUnlock()
if _, err := t.head.Seek(0, 2); err != nil {
return err
}
if _, err := t.head.Write(blob); err != nil { if _, err := t.head.Write(blob); err != nil {
return err return err
} }
@ -444,6 +448,7 @@ func (t *freezerTable) Append(item uint64, blob []byte) error {
offset: newOffset, offset: newOffset,
} }
// Write indexEntry // Write indexEntry
t.index.Seek(0, 2)
t.index.Write(idx.marshallBinary()) t.index.Write(idx.marshallBinary())
t.writeMeter.Mark(int64(bLen + indexEntrySize)) t.writeMeter.Mark(int64(bLen + indexEntrySize))
atomic.AddUint64(&t.items, 1) atomic.AddUint64(&t.items, 1)