From acc0b6fbaace61a3befd7f726f9d80fb69391126 Mon Sep 17 00:00:00 2001 From: SangTran-127 Date: Sat, 8 Aug 2026 12:54:43 +0700 Subject: implement core transport and processing operators with receiver, emitter, filter, and map functionalities --- transport/local_channel.go | 37 +++++++++++++++++++++++++++++++++++++ 1 file changed, 37 insertions(+) create mode 100644 transport/local_channel.go (limited to 'transport/local_channel.go') diff --git a/transport/local_channel.go b/transport/local_channel.go new file mode 100644 index 0000000..4ac2c83 --- /dev/null +++ b/transport/local_channel.go @@ -0,0 +1,37 @@ +package transport + +import "goflink/core" + +type LocalReceiveChannel[T any] struct { + recordCh <-chan core.StreamRecord[T] +} + +func NewLocalReceiveChannel[T any](ch <-chan core.StreamRecord[T]) *LocalReceiveChannel[T] { + return &LocalReceiveChannel[T]{ + recordCh: ch, + } +} + +func (c *LocalReceiveChannel[T]) Receive() (core.StreamRecord[T], bool) { + val, ok := <-c.recordCh + return val, ok +} + +func (c *LocalReceiveChannel[T]) Close() error { + // The emitter will close + return nil +} + +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 (l LocalEmitChannel[T]) Close() error { + close(l.recordCh) + return nil +} -- cgit v1.2.3