tail_volume.go 2.0 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485
  1. package operation
  2. import (
  3. "context"
  4. "fmt"
  5. "io"
  6. "google.golang.org/grpc"
  7. "github.com/chrislusf/seaweedfs/weed/pb/volume_server_pb"
  8. "github.com/chrislusf/seaweedfs/weed/storage/needle"
  9. )
  10. func TailVolume(masterFn GetMasterFn, grpcDialOption grpc.DialOption, vid needle.VolumeId, sinceNs uint64, timeoutSeconds int, fn func(n *needle.Needle) error) error {
  11. // find volume location, replication, ttl info
  12. lookup, err := Lookup(masterFn, vid.String())
  13. if err != nil {
  14. return fmt.Errorf("look up volume %d: %v", vid, err)
  15. }
  16. if len(lookup.Locations) == 0 {
  17. return fmt.Errorf("unable to locate volume %d", vid)
  18. }
  19. volumeServer := lookup.Locations[0].Url
  20. return TailVolumeFromSource(volumeServer, grpcDialOption, vid, sinceNs, timeoutSeconds, fn)
  21. }
  22. func TailVolumeFromSource(volumeServer string, grpcDialOption grpc.DialOption, vid needle.VolumeId, sinceNs uint64, idleTimeoutSeconds int, fn func(n *needle.Needle) error) error {
  23. return WithVolumeServerClient(volumeServer, grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
  24. ctx, cancel := context.WithCancel(context.Background())
  25. defer cancel()
  26. stream, err := client.VolumeTailSender(ctx, &volume_server_pb.VolumeTailSenderRequest{
  27. VolumeId: uint32(vid),
  28. SinceNs: sinceNs,
  29. IdleTimeoutSeconds: uint32(idleTimeoutSeconds),
  30. })
  31. if err != nil {
  32. return err
  33. }
  34. for {
  35. resp, recvErr := stream.Recv()
  36. if recvErr != nil {
  37. if recvErr == io.EOF {
  38. break
  39. } else {
  40. return recvErr
  41. }
  42. }
  43. needleHeader := resp.NeedleHeader
  44. needleBody := resp.NeedleBody
  45. if len(needleHeader) == 0 {
  46. continue
  47. }
  48. for !resp.IsLastChunk {
  49. resp, recvErr = stream.Recv()
  50. if recvErr != nil {
  51. if recvErr == io.EOF {
  52. break
  53. } else {
  54. return recvErr
  55. }
  56. }
  57. needleBody = append(needleBody, resp.NeedleBody...)
  58. }
  59. n := new(needle.Needle)
  60. n.ParseNeedleHeader(needleHeader)
  61. n.ReadNeedleBodyBytes(needleBody, needle.CurrentVersion)
  62. err = fn(n)
  63. if err != nil {
  64. return err
  65. }
  66. }
  67. return nil
  68. })
  69. }