model_puller.go 4.5 KB

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