mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
Revert "event: move type fixation logic into Feed.init (#27249)"
This reverts commit b2a95a215a.
This commit is contained in:
parent
7cc5d3c1bc
commit
7410a9c069
1 changed files with 22 additions and 12 deletions
|
|
@ -57,8 +57,7 @@ func (e feedTypeError) Error() string {
|
||||||
return "event: wrong type in " + e.op + " got " + e.got.String() + ", want " + e.want.String()
|
return "event: wrong type in " + e.op + " got " + e.got.String() + ", want " + e.want.String()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *Feed) init(etype reflect.Type) {
|
func (f *Feed) init() {
|
||||||
f.etype = etype
|
|
||||||
f.removeSub = make(chan interface{})
|
f.removeSub = make(chan interface{})
|
||||||
f.sendLock = make(chan struct{}, 1)
|
f.sendLock = make(chan struct{}, 1)
|
||||||
f.sendLock <- struct{}{}
|
f.sendLock <- struct{}{}
|
||||||
|
|
@ -71,6 +70,8 @@ func (f *Feed) init(etype reflect.Type) {
|
||||||
// The channel should have ample buffer space to avoid blocking other subscribers.
|
// The channel should have ample buffer space to avoid blocking other subscribers.
|
||||||
// Slow subscribers are not dropped.
|
// Slow subscribers are not dropped.
|
||||||
func (f *Feed) Subscribe(channel interface{}) Subscription {
|
func (f *Feed) Subscribe(channel interface{}) Subscription {
|
||||||
|
f.once.Do(f.init)
|
||||||
|
|
||||||
chanval := reflect.ValueOf(channel)
|
chanval := reflect.ValueOf(channel)
|
||||||
chantyp := chanval.Type()
|
chantyp := chanval.Type()
|
||||||
if chantyp.Kind() != reflect.Chan || chantyp.ChanDir()&reflect.SendDir == 0 {
|
if chantyp.Kind() != reflect.Chan || chantyp.ChanDir()&reflect.SendDir == 0 {
|
||||||
|
|
@ -78,13 +79,11 @@ func (f *Feed) Subscribe(channel interface{}) Subscription {
|
||||||
}
|
}
|
||||||
sub := &feedSub{feed: f, channel: chanval, err: make(chan error, 1)}
|
sub := &feedSub{feed: f, channel: chanval, err: make(chan error, 1)}
|
||||||
|
|
||||||
f.once.Do(func() { f.init(chantyp.Elem()) })
|
|
||||||
if f.etype != chantyp.Elem() {
|
|
||||||
panic(feedTypeError{op: "Subscribe", got: chantyp, want: reflect.ChanOf(reflect.SendDir, f.etype)})
|
|
||||||
}
|
|
||||||
|
|
||||||
f.mu.Lock()
|
f.mu.Lock()
|
||||||
defer f.mu.Unlock()
|
defer f.mu.Unlock()
|
||||||
|
if !f.typecheck(chantyp.Elem()) {
|
||||||
|
panic(feedTypeError{op: "Subscribe", got: chantyp, want: reflect.ChanOf(reflect.SendDir, f.etype)})
|
||||||
|
}
|
||||||
// Add the select case to the inbox.
|
// Add the select case to the inbox.
|
||||||
// The next Send will add it to f.sendCases.
|
// The next Send will add it to f.sendCases.
|
||||||
cas := reflect.SelectCase{Dir: reflect.SelectSend, Chan: chanval}
|
cas := reflect.SelectCase{Dir: reflect.SelectSend, Chan: chanval}
|
||||||
|
|
@ -92,6 +91,15 @@ func (f *Feed) Subscribe(channel interface{}) Subscription {
|
||||||
return sub
|
return sub
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// note: callers must hold f.mu
|
||||||
|
func (f *Feed) typecheck(typ reflect.Type) bool {
|
||||||
|
if f.etype == nil {
|
||||||
|
f.etype = typ
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return f.etype == typ
|
||||||
|
}
|
||||||
|
|
||||||
func (f *Feed) remove(sub *feedSub) {
|
func (f *Feed) remove(sub *feedSub) {
|
||||||
// Delete from inbox first, which covers channels
|
// Delete from inbox first, which covers channels
|
||||||
// that have not been added to f.sendCases yet.
|
// that have not been added to f.sendCases yet.
|
||||||
|
|
@ -120,17 +128,19 @@ func (f *Feed) remove(sub *feedSub) {
|
||||||
func (f *Feed) Send(value interface{}) (nsent int) {
|
func (f *Feed) Send(value interface{}) (nsent int) {
|
||||||
rvalue := reflect.ValueOf(value)
|
rvalue := reflect.ValueOf(value)
|
||||||
|
|
||||||
f.once.Do(func() { f.init(rvalue.Type()) })
|
f.once.Do(f.init)
|
||||||
if f.etype != rvalue.Type() {
|
|
||||||
panic(feedTypeError{op: "Send", got: rvalue.Type(), want: f.etype})
|
|
||||||
}
|
|
||||||
|
|
||||||
<-f.sendLock
|
<-f.sendLock
|
||||||
|
|
||||||
// Add new cases from the inbox after taking the send lock.
|
// Add new cases from the inbox after taking the send lock.
|
||||||
f.mu.Lock()
|
f.mu.Lock()
|
||||||
f.sendCases = append(f.sendCases, f.inbox...)
|
f.sendCases = append(f.sendCases, f.inbox...)
|
||||||
f.inbox = nil
|
f.inbox = nil
|
||||||
|
|
||||||
|
if !f.typecheck(rvalue.Type()) {
|
||||||
|
f.sendLock <- struct{}{}
|
||||||
|
f.mu.Unlock()
|
||||||
|
panic(feedTypeError{op: "Send", got: rvalue.Type(), want: f.etype})
|
||||||
|
}
|
||||||
f.mu.Unlock()
|
f.mu.Unlock()
|
||||||
|
|
||||||
// Set the sent value on all channels.
|
// Set the sent value on all channels.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue