I don't know how else to phrase this, but I think it's easier if you look at the code:
func (self *fromChannel[T]) Subscribe(cx context.Context, observer Observer[T]) {
// Receive from self.src and observer.OnNext it
// If canceled, observer.OnComplete and return
// If self.src is closed, observer.OnComplete and return
// Current implementation (I don't know if it's correct)
loop:
for {
select {
case value := <-self.src:
observer.OnNext(value)
case <-cx.Done():
observer.OnComplete()
break loop
}
observer.OnComplete()
}
}
Observer:
type Observer[T any] interface {
OnNext(value T)
OnError(err error)
OnComplete()
}