model_puller.go 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243
  1. package main
  2. /*
  3. Locking
  4. =======
  5. These methods are never called from the outside so don't follow the locking
  6. policy in model.go.
  7. TODO(jb): Refactor this into smaller and cleaner pieces.
  8. TODO(jb): Increase performance by taking apparent peer bandwidth into account.
  9. */
  10. import (
  11. "bytes"
  12. "errors"
  13. "fmt"
  14. "io"
  15. "os"
  16. "path"
  17. "sync"
  18. "time"
  19. "github.com/calmh/syncthing/buffers"
  20. )
  21. func (m *Model) pullFile(name string) error {
  22. m.RLock()
  23. var localFile = m.local[name]
  24. var globalFile = m.global[name]
  25. var nodeIDs = m.whoHas(name)
  26. m.RUnlock()
  27. if len(nodeIDs) == 0 {
  28. return fmt.Errorf("%s: no connected nodes with file available", name)
  29. }
  30. filename := path.Join(m.dir, name)
  31. sdir := path.Dir(filename)
  32. _, err := os.Stat(sdir)
  33. if err != nil && os.IsNotExist(err) {
  34. os.MkdirAll(sdir, 0777)
  35. }
  36. tmpFilename := tempName(filename, globalFile.Modified)
  37. tmpFile, err := os.Create(tmpFilename)
  38. if err != nil {
  39. return err
  40. }
  41. defer tmpFile.Close()
  42. contentChan := make(chan content, 32)
  43. var applyDone sync.WaitGroup
  44. applyDone.Add(1)
  45. go func() {
  46. applyContent(contentChan, tmpFile)
  47. applyDone.Done()
  48. }()
  49. local, remote := localFile.Blocks.To(globalFile.Blocks)
  50. var fetchDone sync.WaitGroup
  51. // One local copy routine
  52. fetchDone.Add(1)
  53. go func() {
  54. for _, block := range local {
  55. data, err := m.Request("<local>", name, block.Offset, block.Length, block.Hash)
  56. if err != nil {
  57. break
  58. }
  59. contentChan <- content{
  60. offset: int64(block.Offset),
  61. data: data,
  62. }
  63. }
  64. fetchDone.Done()
  65. }()
  66. // N remote copy routines
  67. var remoteBlocks = blockIterator{blocks: remote}
  68. for i := 0; i < opts.Advanced.RequestsInFlight; i++ {
  69. curNode := nodeIDs[i%len(nodeIDs)]
  70. fetchDone.Add(1)
  71. go func(nodeID string) {
  72. for {
  73. block, ok := remoteBlocks.Next()
  74. if !ok {
  75. break
  76. }
  77. data, err := m.RequestGlobal(nodeID, name, block.Offset, block.Length, block.Hash)
  78. if err != nil {
  79. break
  80. }
  81. contentChan <- content{
  82. offset: int64(block.Offset),
  83. data: data,
  84. }
  85. }
  86. fetchDone.Done()
  87. }(curNode)
  88. }
  89. fetchDone.Wait()
  90. close(contentChan)
  91. applyDone.Wait()
  92. err = hashCheck(tmpFilename, globalFile.Blocks)
  93. if err != nil {
  94. return fmt.Errorf("%s: %s (deleting)", path.Base(name), err.Error())
  95. }
  96. err = os.Chtimes(tmpFilename, time.Unix(globalFile.Modified, 0), time.Unix(globalFile.Modified, 0))
  97. if err != nil {
  98. return err
  99. }
  100. err = os.Rename(tmpFilename, filename)
  101. if err != nil {
  102. return err
  103. }
  104. return nil
  105. }
  106. func (m *Model) puller() {
  107. for {
  108. time.Sleep(time.Second)
  109. var ns []string
  110. m.RLock()
  111. for n := range m.need {
  112. ns = append(ns, n)
  113. }
  114. m.RUnlock()
  115. if len(ns) == 0 {
  116. continue
  117. }
  118. var limiter = make(chan bool, opts.Advanced.FilesInFlight)
  119. for _, n := range ns {
  120. limiter <- true
  121. go func(n string) {
  122. defer func() {
  123. <-limiter
  124. }()
  125. f, ok := m.GlobalFile(n)
  126. if !ok {
  127. return
  128. }
  129. var err error
  130. if f.Flags&FlagDeleted == 0 {
  131. if opts.Debug.TraceFile {
  132. debugf("FILE: Pull %q", n)
  133. }
  134. err = m.pullFile(n)
  135. } else {
  136. if opts.Debug.TraceFile {
  137. debugf("FILE: Remove %q", n)
  138. }
  139. // Cheerfully ignore errors here
  140. _ = os.Remove(path.Join(m.dir, n))
  141. }
  142. if err == nil {
  143. m.UpdateLocal(f)
  144. } else {
  145. warnln(err)
  146. }
  147. }(n)
  148. }
  149. }
  150. }
  151. type content struct {
  152. offset int64
  153. data []byte
  154. }
  155. func applyContent(cc <-chan content, dst io.WriterAt) error {
  156. var err error
  157. for c := range cc {
  158. _, err = dst.WriteAt(c.data, c.offset)
  159. if err != nil {
  160. return err
  161. }
  162. buffers.Put(c.data)
  163. }
  164. return nil
  165. }
  166. func hashCheck(name string, correct []Block) error {
  167. rf, err := os.Open(name)
  168. if err != nil {
  169. return err
  170. }
  171. defer rf.Close()
  172. current, err := Blocks(rf, BlockSize)
  173. if err != nil {
  174. return err
  175. }
  176. if len(current) != len(correct) {
  177. return errors.New("incorrect number of blocks")
  178. }
  179. for i := range current {
  180. if bytes.Compare(current[i].Hash, correct[i].Hash) != 0 {
  181. return fmt.Errorf("hash mismatch: %x != %x", current[i], correct[i])
  182. }
  183. }
  184. return nil
  185. }
  186. type blockIterator struct {
  187. sync.Mutex
  188. blocks []Block
  189. }
  190. func (i *blockIterator) Next() (b Block, ok bool) {
  191. i.Lock()
  192. defer i.Unlock()
  193. if len(i.blocks) == 0 {
  194. return
  195. }
  196. b, i.blocks = i.blocks[0], i.blocks[1:]
  197. ok = true
  198. return
  199. }