| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1725,12 +1725,12 @@ func (c *grpcStorageClient) OpenWriter(params *openWriterParams, opts ...storage | |||
| 1725 | 1725 | ||
| 1726 | 1726 | var o *storagepb.Object | |
| 1727 | 1727 | uploadBuff := func(ctx context.Context) error { | |
| 1728 | - obj, err := gw.uploadBuffer(recvd, offset, doneReading) | ||
| 1728 | + obj, err := gw.uploadBuffer(ctx, recvd, offset, doneReading) | ||
| 1729 | 1729 | o = obj | |
| 1730 | 1730 | return err | |
| 1731 | 1731 | } | |
| 1732 | 1732 | ||
| 1733 | - err = run(gw.ctx, uploadBuff, gw.settings.retry, s.idempotent) | ||
| 1733 | + err = run(bucketContext(gw.ctx, gw.bucket), uploadBuff, gw.settings.retry, s.idempotent) | ||
| 1734 | 1734 | if err != nil { | |
| 1735 | 1735 | return err | |
| 1736 | 1736 | } | |
@@ -2666,11 +2666,10 @@ type gRPCBidiWriteBufferSender interface { | |||
| 2666 | 2666 | // If flush is true, implementations must not return until the data in buf is | |
| 2667 | 2667 | // stable. If finishWrite is true, implementations must return the object on | |
| 2668 | 2668 | // success. | |
| 2669 | - sendBuffer(buf []byte, offset int64, flush, finishWrite bool) (*storagepb.Object, error) | ||
| 2669 | + sendBuffer(ctx context.Context, buf []byte, offset int64, flush, finishWrite bool) (*storagepb.Object, error) | ||
| 2670 | 2670 | } | |
| 2671 | 2671 | ||
| 2672 | 2672 | type gRPCOneshotBidiWriteBufferSender struct { | |
| 2673 | - ctx context.Context | ||
| 2674 | 2673 | firstMessage *storagepb.BidiWriteObjectRequest | |
| 2675 | 2674 | raw *gapic.Client | |
| 2676 | 2675 | stream storagepb.Storage_BidiWriteObjectClient | |
@@ -2691,17 +2690,16 @@ func (w *gRPCWriter) newGRPCOneshotBidiWriteBufferSender() (*gRPCOneshotBidiWrit | |||
| 2691 | 2690 | } | |
| 2692 | 2691 | ||
| 2693 | 2692 | return &gRPCOneshotBidiWriteBufferSender{ | |
| 2694 | - ctx: bucketContext(w.ctx, w.bucket), | ||
| 2695 | 2693 | firstMessage: firstMessage, | |
| 2696 | 2694 | raw: w.c.raw, | |
| 2697 | 2695 | settings: w.settings, | |
| 2698 | 2696 | }, nil | |
| 2699 | 2697 | } | |
| 2700 | 2698 | ||
| 2701 | - func (s *gRPCOneshotBidiWriteBufferSender) sendBuffer(buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 2699 | + func (s *gRPCOneshotBidiWriteBufferSender) sendBuffer(ctx context.Context, buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 2702 | 2700 | var firstMessage *storagepb.BidiWriteObjectRequest | |
| 2703 | 2701 | if s.stream == nil { | |
| 2704 | - s.stream, err = s.raw.BidiWriteObject(s.ctx, s.settings.gax...) | ||
| 2702 | + s.stream, err = s.raw.BidiWriteObject(ctx, s.settings.gax...) | ||
| 2705 | 2703 | if err != nil { | |
| 2706 | 2704 | return | |
| 2707 | 2705 | } | |
@@ -2737,7 +2735,6 @@ func (s *gRPCOneshotBidiWriteBufferSender) sendBuffer(buf []byte, offset int64, | |||
| 2737 | 2735 | } | |
| 2738 | 2736 | ||
| 2739 | 2737 | type gRPCResumableBidiWriteBufferSender struct { | |
| 2740 | - ctx context.Context | ||
| 2741 | 2738 | queryRetry *retryConfig | |
| 2742 | 2739 | upid string | |
| 2743 | 2740 | progress func(int64) | |
@@ -2748,7 +2745,7 @@ type gRPCResumableBidiWriteBufferSender struct { | |||
| 2748 | 2745 | settings *settings | |
| 2749 | 2746 | } | |
| 2750 | 2747 | ||
| 2751 | - func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender() (*gRPCResumableBidiWriteBufferSender, error) { | ||
| 2748 | + func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender(ctx context.Context) (*gRPCResumableBidiWriteBufferSender, error) { | ||
| 2752 | 2749 | req := &storagepb.StartResumableWriteRequest{ | |
| 2753 | 2750 | WriteObjectSpec: w.spec, | |
| 2754 | 2751 | CommonObjectRequestParams: toProtoCommonObjectRequestParams(w.encryptionKey), | |
@@ -2758,7 +2755,6 @@ func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender() (*gRPCResumableBidi | |||
| 2758 | 2755 | ObjectChecksums: toProtoChecksums(w.sendCRC32C, w.attrs), | |
| 2759 | 2756 | } | |
| 2760 | 2757 | ||
| 2761 | - ctx := bucketContext(w.ctx, w.bucket) | ||
| 2762 | 2758 | var upid string | |
| 2763 | 2759 | err := run(ctx, func(ctx context.Context) error { | |
| 2764 | 2760 | upres, err := w.c.raw.StartResumableWrite(ctx, req, w.settings.gax...) | |
@@ -2778,7 +2774,6 @@ func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender() (*gRPCResumableBidi | |||
| 2778 | 2774 | } | |
| 2779 | 2775 | ||
| 2780 | 2776 | return &gRPCResumableBidiWriteBufferSender{ | |
| 2781 | - ctx: ctx, | ||
| 2782 | 2777 | queryRetry: w.settings.retry, | |
| 2783 | 2778 | upid: upid, | |
| 2784 | 2779 | progress: w.progress, | |
@@ -2791,9 +2786,9 @@ func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender() (*gRPCResumableBidi | |||
| 2791 | 2786 | ||
| 2792 | 2787 | // queryProgress is a helper that queries the status of the resumable upload | |
| 2793 | 2788 | // associated with the given upload ID. | |
| 2794 | - func (s *gRPCResumableBidiWriteBufferSender) queryProgress() (int64, error) { | ||
| 2789 | + func (s *gRPCResumableBidiWriteBufferSender) queryProgress(ctx context.Context) (int64, error) { | ||
| 2795 | 2790 | var persistedSize int64 | |
| 2796 | - err := run(s.ctx, func(ctx context.Context) error { | ||
| 2791 | + err := run(ctx, func(ctx context.Context) error { | ||
| 2797 | 2792 | q, err := s.raw.QueryWriteStatus(ctx, &storagepb.QueryWriteStatusRequest{ | |
| 2798 | 2793 | UploadId: s.upid, | |
| 2799 | 2794 | }, s.settings.gax...) | |
@@ -2805,15 +2800,15 @@ func (s *gRPCResumableBidiWriteBufferSender) queryProgress() (int64, error) { | |||
| 2805 | 2800 | return persistedSize, err | |
| 2806 | 2801 | } | |
| 2807 | 2802 | ||
| 2808 | - func (s *gRPCResumableBidiWriteBufferSender) sendBuffer(buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 2803 | + func (s *gRPCResumableBidiWriteBufferSender) sendBuffer(ctx context.Context, buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 2809 | 2804 | reconnected := false | |
| 2810 | 2805 | if s.stream == nil { | |
| 2811 | 2806 | // Determine offset and reconnect | |
| 2812 | - s.flushOffset, err = s.queryProgress() | ||
| 2807 | + s.flushOffset, err = s.queryProgress(ctx) | ||
| 2813 | 2808 | if err != nil { | |
| 2814 | 2809 | return | |
| 2815 | 2810 | } | |
| 2816 | - s.stream, err = s.raw.BidiWriteObject(s.ctx, s.settings.gax...) | ||
| 2811 | + s.stream, err = s.raw.BidiWriteObject(ctx, s.settings.gax...) | ||
| 2817 | 2812 | if err != nil { | |
| 2818 | 2813 | return | |
| 2819 | 2814 | } | |
@@ -2885,7 +2880,7 @@ func (s *gRPCResumableBidiWriteBufferSender) sendBuffer(buf []byte, offset int64 | |||
| 2885 | 2880 | // The final Object is returned on success if doneReading is true. | |
| 2886 | 2881 | // | |
| 2887 | 2882 | // Returns object and any error that is not retriable. | |
| 2888 | - func (w *gRPCWriter) uploadBuffer(recvd int, start int64, doneReading bool) (obj *storagepb.Object, err error) { | ||
| 2883 | + func (w *gRPCWriter) uploadBuffer(ctx context.Context, recvd int, start int64, doneReading bool) (obj *storagepb.Object, err error) { | ||
| 2889 | 2884 | if w.streamSender == nil { | |
| 2890 | 2885 | if w.append { | |
| 2891 | 2886 | // Appendable object semantics | |
@@ -2895,7 +2890,7 @@ func (w *gRPCWriter) uploadBuffer(recvd int, start int64, doneReading bool) (obj | |||
| 2895 | 2890 | w.streamSender, err = w.newGRPCOneshotBidiWriteBufferSender() | |
| 2896 | 2891 | } else { | |
| 2897 | 2892 | // Resumable write semantics | |
| 2898 | - w.streamSender, err = w.newGRPCResumableBidiWriteBufferSender() | ||
| 2893 | + w.streamSender, err = w.newGRPCResumableBidiWriteBufferSender(ctx) | ||
| 2899 | 2894 | } | |
| 2900 | 2895 | if err != nil { | |
| 2901 | 2896 | return | |
@@ -2915,7 +2910,7 @@ func (w *gRPCWriter) uploadBuffer(recvd int, start int64, doneReading bool) (obj | |||
| 2915 | 2910 | l = len(data) | |
| 2916 | 2911 | flush = true | |
| 2917 | 2912 | } | |
| 2918 | - obj, err = w.streamSender.sendBuffer(data[:l], offset, flush, flush && doneReading) | ||
| 2913 | + obj, err = w.streamSender.sendBuffer(ctx, data[:l], offset, flush, flush && doneReading) | ||
| 2919 | 2914 | if err != nil { | |
| 2920 | 2915 | return nil, err | |
| 2921 | 2916 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -29,7 +29,6 @@ import ( | |||
| 29 | 29 | ) | |
| 30 | 30 | ||
| 31 | 31 | type gRPCAppendBidiWriteBufferSender struct { | |
| 32 | - ctx context.Context | ||
| 33 | 32 | bucket string | |
| 34 | 33 | routingToken *string | |
| 35 | 34 | raw *gapic.Client | |
@@ -51,7 +50,6 @@ type gRPCAppendBidiWriteBufferSender struct { | |||
| 51 | 50 | ||
| 52 | 51 | func (w *gRPCWriter) newGRPCAppendBidiWriteBufferSender() (*gRPCAppendBidiWriteBufferSender, error) { | |
| 53 | 52 | s := &gRPCAppendBidiWriteBufferSender{ | |
| 54 | - ctx: w.ctx, | ||
| 55 | 53 | bucket: w.spec.GetResource().GetBucket(), | |
| 56 | 54 | raw: w.c.raw, | |
| 57 | 55 | settings: w.c.settings, | |
@@ -68,7 +66,7 @@ func (w *gRPCWriter) newGRPCAppendBidiWriteBufferSender() (*gRPCAppendBidiWriteB | |||
| 68 | 66 | return s, nil | |
| 69 | 67 | } | |
| 70 | 68 | ||
| 71 | - func (s *gRPCAppendBidiWriteBufferSender) connect() (err error) { | ||
| 69 | + func (s *gRPCAppendBidiWriteBufferSender) connect(ctx context.Context) (err error) { | ||
| 72 | 70 | err = func() error { | |
| 73 | 71 | // If this is a forced first message, we've already determined it's safe to | |
| 74 | 72 | // send. | |
@@ -107,19 +105,19 @@ func (s *gRPCAppendBidiWriteBufferSender) connect() (err error) { | |||
| 107 | 105 | return err | |
| 108 | 106 | } | |
| 109 | 107 | ||
| 110 | - return s.startReceiver() | ||
| 108 | + return s.startReceiver(ctx) | ||
| 111 | 109 | } | |
| 112 | 110 | ||
| 113 | 111 | func (s *gRPCAppendBidiWriteBufferSender) withRequestParams(ctx context.Context) context.Context { | |
| 114 | 112 | param := fmt.Sprintf("appendable=true&bucket=%s", s.bucket) | |
| 115 | 113 | if s.routingToken != nil { | |
| 116 | 114 | param = param + fmt.Sprintf("&routing_token=%s", *s.routingToken) | |
| 117 | 115 | } | |
| 118 | - return gax.InsertMetadataIntoOutgoingContext(s.ctx, "x-goog-request-params", param) | ||
| 116 | + return gax.InsertMetadataIntoOutgoingContext(ctx, "x-goog-request-params", param) | ||
| 119 | 117 | } | |
| 120 | 118 | ||
| 121 | - func (s *gRPCAppendBidiWriteBufferSender) startReceiver() (err error) { | ||
| 122 | - s.stream, err = s.raw.BidiWriteObject(s.withRequestParams(s.ctx), s.settings.gax...) | ||
| 119 | + func (s *gRPCAppendBidiWriteBufferSender) startReceiver(ctx context.Context) (err error) { | ||
| 120 | + s.stream, err = s.raw.BidiWriteObject(s.withRequestParams(ctx), s.settings.gax...) | ||
| 123 | 121 | if err != nil { | |
| 124 | 122 | return | |
| 125 | 123 | } | |
@@ -282,12 +280,12 @@ func (s *gRPCAppendBidiWriteBufferSender) sendOnConnectedStream(buf []byte, offs | |||
| 282 | 280 | return | |
| 283 | 281 | } | |
| 284 | 282 | ||
| 285 | - func (s *gRPCAppendBidiWriteBufferSender) sendBuffer(buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 283 | + func (s *gRPCAppendBidiWriteBufferSender) sendBuffer(ctx context.Context, buf []byte, offset int64, flush, finishWrite bool) (obj *storagepb.Object, err error) { | ||
| 286 | 284 | for { | |
| 287 | 285 | sendFirstMessage := false | |
| 288 | 286 | if s.stream == nil { | |
| 289 | 287 | sendFirstMessage = true | |
| 290 | - if err = s.connect(); err != nil { | ||
| 288 | + if err = s.connect(ctx); err != nil { | ||
| 291 | 289 | return | |
| 292 | 290 | } | |
| 293 | 291 | } | |
| Back | FazBrowse Home | New Git URL |
0 commit comments