aboutsummaryrefslogtreecommitdiff
path: root/transport/local_channel.go
diff options
context:
space:
mode:
authorSangTran-127 <tranquangsang12.7@gmail.com>2026-08-08 17:24:01 +0700
committerSangTran-127 <tranquangsang12.7@gmail.com>2026-08-08 17:24:01 +0700
commit0958e46c3651f2c8b6d76bcffbd8820063647d59 (patch)
tree1e917ed409eddc0b3c29fd602e551778c96c8dc6 /transport/local_channel.go
parentacc0b6fbaace61a3befd7f726f9d80fb69391126 (diff)
downloadgoflink-0958e46c3651f2c8b6d76bcffbd8820063647d59.tar.gz
goflink-0958e46c3651f2c8b6d76bcffbd8820063647d59.zip
refactor: add context support to transport operators and runtime task management
Diffstat (limited to 'transport/local_channel.go')
-rw-r--r--transport/local_channel.go35
1 files changed, 27 insertions, 8 deletions
diff --git a/transport/local_channel.go b/transport/local_channel.go
index 4ac2c83..e3e4cba 100644
--- a/transport/local_channel.go
+++ b/transport/local_channel.go
@@ -1,6 +1,10 @@
package transport
-import "goflink/core"
+import (
+ "context"
+
+ "goflink/core"
+)
type LocalReceiveChannel[T any] struct {
recordCh <-chan core.StreamRecord[T]
@@ -12,9 +16,14 @@ func NewLocalReceiveChannel[T any](ch <-chan core.StreamRecord[T]) *LocalReceive
}
}
-func (c *LocalReceiveChannel[T]) Receive() (core.StreamRecord[T], bool) {
- val, ok := <-c.recordCh
- return val, ok
+func (c *LocalReceiveChannel[T]) Receive(ctx context.Context) (core.StreamRecord[T], bool) {
+ select {
+ case val, ok := <-c.recordCh:
+ return val, ok
+ case <-ctx.Done():
+ var zero core.StreamRecord[T]
+ return zero, false
+ }
}
func (c *LocalReceiveChannel[T]) Close() error {
@@ -26,12 +35,22 @@ type LocalEmitChannel[T any] struct {
recordCh chan<- core.StreamRecord[T]
}
-func (l LocalEmitChannel[T]) Emit(c core.StreamRecord[T]) error {
- l.recordCh <- c
- return nil
+func NewLocalEmitChannel[T any](ch chan<- core.StreamRecord[T]) *LocalEmitChannel[T] {
+ return &LocalEmitChannel[T]{
+ recordCh: ch,
+ }
+}
+
+func (l *LocalEmitChannel[T]) Emit(ctx context.Context, c core.StreamRecord[T]) error {
+ select {
+ case l.recordCh <- c:
+ return nil
+ case <-ctx.Done():
+ return ctx.Err()
+ }
}
-func (l LocalEmitChannel[T]) Close() error {
+func (l *LocalEmitChannel[T]) Close() error {
close(l.recordCh)
return nil
}