model_puller.go 4.1 KB

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