aboutsummaryrefslogtreecommitdiff
path: root/transport/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/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/channel.go')
-rw-r--r--transport/channel.go16
1 files changed, 13 insertions, 3 deletions
diff --git a/transport/channel.go b/transport/channel.go
index 011f8ae..1db6ea8 100644
--- a/transport/channel.go
+++ b/transport/channel.go
@@ -1,13 +1,23 @@
package transport
-import "goflink/core"
+import (
+ "context"
+
+ "goflink/core"
+)
type Receiver[T any] interface {
- Receive() (core.StreamRecord[T], bool)
+ // Receive returns the next record, or ok=false once the stream is done or
+ // ctx is cancelled.
+ Receive(ctx context.Context) (core.StreamRecord[T], bool)
Close() error
}
+// Emitter mirrors core.Emitter plus Close. It is declared separately because
+// core cannot import transport without an import cycle.
type Emitter[T any] interface {
- Emit(core.StreamRecord[T]) error
+ // Emit blocks until the record is handed off, or returns ctx.Err() if ctx
+ // is cancelled first.
+ Emit(ctx context.Context, record core.StreamRecord[T]) error
Close() error
}