model_puller.go 4.2 KB

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