replicate is a long lived process that replicates keys to the peer.
()
| 275 | |
| 276 | // replicate is a long lived process that replicates keys to the peer. |
| 277 | func (p *peerSync) replicate() { |
| 278 | for { |
| 279 | select { |
| 280 | case <-p.endSync: |
| 281 | return |
| 282 | case <-p.r.quit: |
| 283 | return |
| 284 | case key := <-p.queue: |
| 285 | p.lock.Lock() |
| 286 | log, ok := p.tasks[key] |
| 287 | delete(p.tasks, key) |
| 288 | p.lock.Unlock() |
| 289 | if !ok { |
| 290 | continue |
| 291 | } |
| 292 | |
| 293 | task, err := collapseOpLog(log) |
| 294 | if err != nil { |
| 295 | p.r.n.logger.Printf("[ERROR] onecache.replicator: key %v tasks invalid: %v\n", key, err) |
| 296 | continue |
| 297 | } |
| 298 | |
| 299 | // Ensure we are still the owner of the key. |
| 300 | if keys := p.r.n.ring.getOwnedKeys([]string{key}); len(keys) == 0 { |
| 301 | continue |
| 302 | } |
| 303 | |
| 304 | p.lock.Lock() |
| 305 | p.replicating = key |
| 306 | p.lock.Unlock() |
| 307 | |
| 308 | switch task { |
| 309 | case task_BACKFILL: |
| 310 | p.replicateBackfill(key, p.name) |
| 311 | case task_DIRTY_VALUE: |
| 312 | p.replicateDirtyValue(key, p.name) |
| 313 | case task_TOUCH: |
| 314 | p.replicateTouch(key, p.name) |
| 315 | case task_DELETE: |
| 316 | p.replicateDelete(key, p.name) |
| 317 | } |
| 318 | |
| 319 | p.lock.Lock() |
| 320 | p.replicating = "" |
| 321 | p.lock.Unlock() |
| 322 | } |
| 323 | } |
| 324 | } |
| 325 | |
| 326 | // replicateDirtyValue replicates a key due to a dirty value. |
| 327 | func (p *peerSync) replicateDirtyValue(key, peer string) error { |
no test coverage detected