You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

215 lines
6.6 KiB

  1. package storage
  2. import (
  3. "fmt"
  4. "github.com/seaweedfs/seaweedfs/weed/util/mem"
  5. "io"
  6. "time"
  7. "github.com/seaweedfs/seaweedfs/weed/glog"
  8. "github.com/seaweedfs/seaweedfs/weed/storage/backend"
  9. "github.com/seaweedfs/seaweedfs/weed/storage/needle"
  10. "github.com/seaweedfs/seaweedfs/weed/storage/super_block"
  11. . "github.com/seaweedfs/seaweedfs/weed/storage/types"
  12. )
  13. const PagedReadLimit = 1024 * 1024
  14. // read fills in Needle content by looking up n.Id from NeedleMapper
  15. func (v *Volume) readNeedle(n *needle.Needle, readOption *ReadOption, onReadSizeFn func(size Size)) (count int, err error) {
  16. v.dataFileAccessLock.RLock()
  17. defer v.dataFileAccessLock.RUnlock()
  18. nv, ok := v.nm.Get(n.Id)
  19. if !ok || nv.Offset.IsZero() {
  20. return -1, ErrorNotFound
  21. }
  22. readSize := nv.Size
  23. if readSize.IsDeleted() {
  24. if readOption != nil && readOption.ReadDeleted && readSize != TombstoneFileSize {
  25. glog.V(3).Infof("reading deleted %s", n.String())
  26. readSize = -readSize
  27. } else {
  28. return -1, ErrorDeleted
  29. }
  30. }
  31. if readSize == 0 {
  32. return 0, nil
  33. }
  34. if onReadSizeFn != nil {
  35. onReadSizeFn(readSize)
  36. }
  37. if readOption != nil && readOption.AttemptMetaOnly && readSize > PagedReadLimit {
  38. readOption.VolumeRevision = v.SuperBlock.CompactionRevision
  39. err = n.ReadNeedleMeta(v.DataBackend, nv.Offset.ToActualOffset(), readSize, v.Version())
  40. if err == needle.ErrorSizeMismatch && OffsetSize == 4 {
  41. readOption.IsOutOfRange = true
  42. err = n.ReadNeedleMeta(v.DataBackend, nv.Offset.ToActualOffset()+int64(MaxPossibleVolumeSize), readSize, v.Version())
  43. }
  44. if err != nil {
  45. return 0, err
  46. }
  47. if !n.IsCompressed() && !n.IsChunkedManifest() {
  48. readOption.IsMetaOnly = true
  49. }
  50. }
  51. if readOption == nil || !readOption.IsMetaOnly {
  52. err = n.ReadData(v.DataBackend, nv.Offset.ToActualOffset(), readSize, v.Version())
  53. if err == needle.ErrorSizeMismatch && OffsetSize == 4 {
  54. err = n.ReadData(v.DataBackend, nv.Offset.ToActualOffset()+int64(MaxPossibleVolumeSize), readSize, v.Version())
  55. }
  56. v.checkReadWriteError(err)
  57. if err != nil {
  58. return 0, err
  59. }
  60. }
  61. count = int(n.DataSize)
  62. if !n.HasTtl() {
  63. return
  64. }
  65. ttlMinutes := n.Ttl.Minutes()
  66. if ttlMinutes == 0 {
  67. return
  68. }
  69. if !n.HasLastModifiedDate() {
  70. return
  71. }
  72. if time.Now().Before(time.Unix(0, int64(n.AppendAtNs)).Add(time.Duration(ttlMinutes) * time.Minute)) {
  73. return
  74. }
  75. return -1, ErrorNotFound
  76. }
  77. // read needle at a specific offset
  78. func (v *Volume) readNeedleMetaAt(n *needle.Needle, offset int64, size int32) (err error) {
  79. v.dataFileAccessLock.RLock()
  80. defer v.dataFileAccessLock.RUnlock()
  81. // read deleted meta data
  82. if size < 0 {
  83. size = -size
  84. }
  85. err = n.ReadNeedleMeta(v.DataBackend, offset, Size(size), v.Version())
  86. if err == needle.ErrorSizeMismatch && OffsetSize == 4 {
  87. err = n.ReadNeedleMeta(v.DataBackend, offset+int64(MaxPossibleVolumeSize), Size(size), v.Version())
  88. }
  89. if err != nil {
  90. return err
  91. }
  92. return nil
  93. }
  94. // read fills in Needle content by looking up n.Id from NeedleMapper
  95. func (v *Volume) readNeedleDataInto(n *needle.Needle, readOption *ReadOption, writer io.Writer, offset int64, size int64) (err error) {
  96. v.dataFileAccessLock.RLock()
  97. defer v.dataFileAccessLock.RUnlock()
  98. nv, ok := v.nm.Get(n.Id)
  99. if !ok || nv.Offset.IsZero() {
  100. return ErrorNotFound
  101. }
  102. readSize := nv.Size
  103. if readSize.IsDeleted() {
  104. if readOption != nil && readOption.ReadDeleted && readSize != TombstoneFileSize {
  105. glog.V(3).Infof("reading deleted %s", n.String())
  106. readSize = -readSize
  107. } else {
  108. return ErrorDeleted
  109. }
  110. }
  111. if readSize == 0 {
  112. return nil
  113. }
  114. if readOption.VolumeRevision != v.SuperBlock.CompactionRevision {
  115. // the volume is compacted
  116. readOption.IsOutOfRange = false
  117. err = n.ReadNeedleMeta(v.DataBackend, nv.Offset.ToActualOffset(), readSize, v.Version())
  118. }
  119. buf := mem.Allocate(min(1024*1024, int(size)))
  120. defer mem.Free(buf)
  121. actualOffset := nv.Offset.ToActualOffset()
  122. if readOption.IsOutOfRange {
  123. actualOffset += int64(MaxPossibleVolumeSize)
  124. }
  125. return n.ReadNeedleDataInto(v.DataBackend, actualOffset, buf, writer, offset, size)
  126. }
  127. func min(x, y int) int {
  128. if x < y {
  129. return x
  130. }
  131. return y
  132. }
  133. // read fills in Needle content by looking up n.Id from NeedleMapper
  134. func (v *Volume) ReadNeedleBlob(offset int64, size Size) ([]byte, error) {
  135. v.dataFileAccessLock.RLock()
  136. defer v.dataFileAccessLock.RUnlock()
  137. return needle.ReadNeedleBlob(v.DataBackend, offset, size, v.Version())
  138. }
  139. type VolumeFileScanner interface {
  140. VisitSuperBlock(super_block.SuperBlock) error
  141. ReadNeedleBody() bool
  142. VisitNeedle(n *needle.Needle, offset int64, needleHeader, needleBody []byte) error
  143. }
  144. func ScanVolumeFile(dirname string, collection string, id needle.VolumeId,
  145. needleMapKind NeedleMapKind,
  146. volumeFileScanner VolumeFileScanner) (err error) {
  147. var v *Volume
  148. if v, err = loadVolumeWithoutIndex(dirname, collection, id, needleMapKind); err != nil {
  149. return fmt.Errorf("failed to load volume %d: %v", id, err)
  150. }
  151. if err = volumeFileScanner.VisitSuperBlock(v.SuperBlock); err != nil {
  152. return fmt.Errorf("failed to process volume %d super block: %v", id, err)
  153. }
  154. defer v.Close()
  155. version := v.Version()
  156. offset := int64(v.SuperBlock.BlockSize())
  157. return ScanVolumeFileFrom(version, v.DataBackend, offset, volumeFileScanner)
  158. }
  159. func ScanVolumeFileFrom(version needle.Version, datBackend backend.BackendStorageFile, offset int64, volumeFileScanner VolumeFileScanner) (err error) {
  160. n, nh, rest, e := needle.ReadNeedleHeader(datBackend, version, offset)
  161. if e != nil {
  162. if e == io.EOF {
  163. return nil
  164. }
  165. return fmt.Errorf("cannot read %s at offset %d: %v", datBackend.Name(), offset, e)
  166. }
  167. for n != nil {
  168. var needleBody []byte
  169. if volumeFileScanner.ReadNeedleBody() {
  170. // println("needle", n.Id.String(), "offset", offset, "size", n.Size, "rest", rest)
  171. if needleBody, err = n.ReadNeedleBody(datBackend, version, offset+NeedleHeaderSize, rest); err != nil {
  172. glog.V(0).Infof("cannot read needle head [%d, %d) body [%d, %d) body length %d: %v", offset, offset+NeedleHeaderSize, offset+NeedleHeaderSize, offset+NeedleHeaderSize+rest, rest, err)
  173. // err = fmt.Errorf("cannot read needle body: %v", err)
  174. // return
  175. }
  176. }
  177. err := volumeFileScanner.VisitNeedle(n, offset, nh, needleBody)
  178. if err == io.EOF {
  179. return nil
  180. }
  181. if err != nil {
  182. glog.V(0).Infof("visit needle error: %v", err)
  183. return fmt.Errorf("visit needle error: %v", err)
  184. }
  185. offset += NeedleHeaderSize + rest
  186. glog.V(4).Infof("==> new entry offset %d", offset)
  187. if n, nh, rest, err = needle.ReadNeedleHeader(datBackend, version, offset); err != nil {
  188. if err == io.EOF {
  189. return nil
  190. }
  191. return fmt.Errorf("cannot read needle header at offset %d: %v", offset, err)
  192. }
  193. glog.V(4).Infof("new entry needle size:%d rest:%d", n.Size, rest)
  194. }
  195. return nil
  196. }