| 435 | } |
| 436 | |
| 437 | func (dc *DCPClient) openStream(vbID uint16, maxRetries uint32) error { |
| 438 | |
| 439 | var openStreamErr error |
| 440 | var attempts uint32 |
| 441 | for { |
| 442 | // Cancel open for stopped client |
| 443 | select { |
| 444 | |
| 445 | case <-dc.terminator: |
| 446 | return nil |
| 447 | default: |
| 448 | } |
| 449 | |
| 450 | openStreamErr = dc.openStreamRequest(vbID) |
| 451 | if openStreamErr == nil { |
| 452 | return nil |
| 453 | } |
| 454 | |
| 455 | var rollbackErr gocbcore.DCPRollbackError |
| 456 | switch { |
| 457 | case errors.As(openStreamErr, &rollbackErr): |
| 458 | if dc.failOnRollback { |
| 459 | InfofCtx(dc.ctx, KeyDCP, "Open stream for vbID %d failed due to rollback or range error, closing client based on failOnRollback=true", vbID) |
| 460 | return fmt.Errorf("%w, failOnRollback requested", openStreamErr) |
| 461 | } |
| 462 | InfofCtx(dc.ctx, KeyDCP, "Open stream for vbID %d failed due to rollback or range error, will roll back metadata and retry: %v", vbID, openStreamErr) |
| 463 | |
| 464 | dc.rollback(dc.ctx, vbID, rollbackErr.SeqNo) |
| 465 | case errors.Is(openStreamErr, gocbcore.ErrMemdRangeError): |
| 466 | err := fmt.Errorf("Invalid metadata out of range for vbID %d, err: %v metadata %+v, shutting down agent", vbID, openStreamErr, dc.metadata.GetMeta(vbID)) |
| 467 | WarnfCtx(dc.ctx, "%s", err) |
| 468 | return err |
| 469 | case errors.Is(openStreamErr, ErrVbUUIDMismatch): |
| 470 | WarnfCtx(dc.ctx, "Closing Stream for vbID: %d, %s", vbID, openStreamErr) |
| 471 | return openStreamErr |
| 472 | case errors.Is(openStreamErr, gocbcore.ErrShutdown): |
| 473 | WarnfCtx(dc.ctx, "Closing stream for vbID %d, agent has been shut down", vbID) |
| 474 | return openStreamErr |
| 475 | case errors.Is(openStreamErr, ErrTimeout): |
| 476 | InfofCtx(dc.ctx, KeyDCP, "Timeout attempting to open stream for vb %d, will retry", vbID) |
| 477 | default: |
| 478 | WarnfCtx(dc.ctx, "Unknown error opening stream for vbID %d: %v", vbID, openStreamErr) |
| 479 | } |
| 480 | if maxRetries == infiniteOpenStreamRetries { |
| 481 | continue |
| 482 | } else if attempts > maxRetries { |
| 483 | break |
| 484 | } |
| 485 | attempts++ |
| 486 | } |
| 487 | |
| 488 | return fmt.Errorf("openStream failed to complete after %d attempts, last error: %w", attempts, openStreamErr) |
| 489 | } |
| 490 | |
| 491 | func (dc *DCPClient) rollback(ctx context.Context, vbID uint16, seqNo gocbcore.SeqNo) { |
| 492 | if dc.dbStats != nil { |