-
-
Notifications
You must be signed in to change notification settings - Fork 398
/
Copy pathutility_grpc_streaming.go
120 lines (99 loc) · 3.3 KB
/
utility_grpc_streaming.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
// This file is part of arduino-cli.
//
// Copyright 2024 ARDUINO SA (http://www.arduino.cc/)
//
// This software is released under the GNU General Public License version 3,
// which covers the main part of arduino-cli.
// The terms of this license can be found at:
// https://www.gnu.org/licenses/gpl-3.0.en.html
//
// You can be released from the requirements of the above licenses by purchasing
// a commercial license. Buying such a license is mandatory if you want to
// modify or otherwise use the software for commercial activities involving the
// Arduino software without disclosing the source code of your own applications.
// To purchase a commercial license, send an email to license@arduino.cc.
package commands
import (
"context"
"errors"
"sync"
"google.golang.org/grpc/metadata"
)
type streamingResponseProxyToChan[T any] struct {
ctx context.Context
respChan chan<- *T
respLock sync.Mutex
}
func streamResponseToChan[T any](ctx context.Context) (*streamingResponseProxyToChan[T], <-chan *T) {
respChan := make(chan *T, 1)
w := &streamingResponseProxyToChan[T]{
ctx: ctx,
respChan: respChan,
}
go func() {
<-ctx.Done()
w.respLock.Lock()
close(w.respChan)
w.respChan = nil
w.respLock.Unlock()
}()
return w, respChan
}
func (w *streamingResponseProxyToChan[T]) Send(resp *T) error {
w.respLock.Lock()
if w.respChan != nil {
w.respChan <- resp
}
w.respLock.Unlock()
return nil
}
func (w *streamingResponseProxyToChan[T]) Context() context.Context {
return w.ctx
}
func (w *streamingResponseProxyToChan[T]) RecvMsg(m any) error {
return errors.New("RecvMsg not implemented")
}
func (w *streamingResponseProxyToChan[T]) SendHeader(metadata.MD) error {
return errors.New("SendHeader not implemented")
}
func (w *streamingResponseProxyToChan[T]) SendMsg(m any) error {
return errors.New("SendMsg not implemented")
}
func (w *streamingResponseProxyToChan[T]) SetHeader(metadata.MD) error {
return errors.New("SetHeader not implemented")
}
func (w *streamingResponseProxyToChan[T]) SetTrailer(tr metadata.MD) {
}
// streamingResponseProxyToCallback is a streaming response proxy that
// forwards the responses to a callback function
type streamingResponseProxyToCallback[T any] struct {
ctx context.Context
cb func(*T) error
}
// creates a streaming response proxy that forwards the responses to a callback function
func streamResponseToCallback[T any](ctx context.Context, cb func(*T) error) *streamingResponseProxyToCallback[T] {
if cb == nil {
cb = func(*T) error { return nil }
}
return &streamingResponseProxyToCallback[T]{ctx: ctx, cb: cb}
}
func (w *streamingResponseProxyToCallback[T]) Send(resp *T) error {
return w.cb(resp)
}
func (w *streamingResponseProxyToCallback[T]) Context() context.Context {
return w.ctx
}
func (w *streamingResponseProxyToCallback[T]) RecvMsg(m any) error {
return errors.New("RecvMsg not implemented")
}
func (w *streamingResponseProxyToCallback[T]) SendHeader(metadata.MD) error {
return errors.New("SendHeader not implemented")
}
func (w *streamingResponseProxyToCallback[T]) SendMsg(m any) error {
return errors.New("SendMsg not implemented")
}
func (w *streamingResponseProxyToCallback[T]) SetHeader(metadata.MD) error {
return errors.New("SetHeader not implemented")
}
func (w *streamingResponseProxyToCallback[T]) SetTrailer(tr metadata.MD) {
}