mirror of
https://github.com/netbirdio/gvisor.git
synced 2026-05-22 17:12:49 -07:00
Add POLLRDNORM/POLLWRNORM support.
On Linux these are meant to be equivalent to POLLIN/POLLOUT. Rather than hack these on in sys_poll etc it felt cleaner to just cleanup the call sites to notify for both events. This is what linux does as well. Fixes #5544 PiperOrigin-RevId: 364859977
This commit is contained in:
committed by
gVisor bot
parent
72ff6a1cac
commit
e7ca2a51a8
@@ -248,7 +248,7 @@ func (c *ConnectedEndpoint) Writable() bool {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
|
||||
return fdnotifier.NonBlockingPoll(int32(c.file.FD()), waiter.EventOut)&waiter.EventOut != 0
|
||||
return fdnotifier.NonBlockingPoll(int32(c.file.FD()), waiter.WritableEvents)&waiter.WritableEvents != 0
|
||||
}
|
||||
|
||||
// Passcred implements transport.ConnectedEndpoint.Passcred.
|
||||
@@ -345,7 +345,7 @@ func (c *ConnectedEndpoint) Readable() bool {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
|
||||
return fdnotifier.NonBlockingPoll(int32(c.file.FD()), waiter.EventIn)&waiter.EventIn != 0
|
||||
return fdnotifier.NonBlockingPoll(int32(c.file.FD()), waiter.ReadableEvents)&waiter.ReadableEvents != 0
|
||||
}
|
||||
|
||||
// SendQueuedSize implements transport.Receiver.SendQueuedSize.
|
||||
|
||||
@@ -41,13 +41,13 @@ func TestWait(t *testing.T) {
|
||||
|
||||
defer file.DecRef(ctx)
|
||||
|
||||
r := file.Readiness(waiter.EventIn)
|
||||
r := file.Readiness(waiter.ReadableEvents)
|
||||
if r != 0 {
|
||||
t.Fatalf("File is ready for read when it shouldn't be.")
|
||||
}
|
||||
|
||||
e, ch := waiter.NewChannelEntry(nil)
|
||||
file.EventRegister(&e, waiter.EventIn)
|
||||
file.EventRegister(&e, waiter.ReadableEvents)
|
||||
defer file.EventUnregister(&e)
|
||||
|
||||
// Check that there are no notifications yet.
|
||||
|
||||
@@ -107,7 +107,7 @@ func (i *Inotify) Readiness(mask waiter.EventMask) waiter.EventMask {
|
||||
defer i.evMu.Unlock()
|
||||
|
||||
if !i.events.Empty() {
|
||||
ready |= waiter.EventIn
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
|
||||
return mask & ready
|
||||
@@ -246,7 +246,7 @@ func (i *Inotify) queueEvent(ev *Event) {
|
||||
// can do.
|
||||
i.evMu.Unlock()
|
||||
|
||||
i.Queue.Notify(waiter.EventIn)
|
||||
i.Queue.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
|
||||
// newWatchLocked creates and adds a new watch to target.
|
||||
|
||||
@@ -101,7 +101,7 @@ func (t *TimerOperations) SetTime(s ktime.Setting) (ktime.Time, ktime.Setting) {
|
||||
func (t *TimerOperations) Readiness(mask waiter.EventMask) waiter.EventMask {
|
||||
var ready waiter.EventMask
|
||||
if atomic.LoadUint64(&t.val) != 0 {
|
||||
ready |= waiter.EventIn
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
return ready
|
||||
}
|
||||
@@ -143,7 +143,7 @@ func (t *TimerOperations) Write(context.Context, *fs.File, usermem.IOSequence, i
|
||||
// Notify implements ktime.TimerListener.Notify.
|
||||
func (t *TimerOperations) Notify(exp uint64, setting ktime.Setting) (ktime.Setting, bool) {
|
||||
atomic.AddUint64(&t.val, exp)
|
||||
t.events.Notify(waiter.EventIn)
|
||||
t.events.Notify(waiter.ReadableEvents)
|
||||
return ktime.Setting{}, false
|
||||
}
|
||||
|
||||
|
||||
@@ -143,7 +143,7 @@ func (l *lineDiscipline) setTermios(task *kernel.Task, args arch.SyscallArgument
|
||||
l.inQueue.pushWaitBufLocked(l)
|
||||
l.inQueue.readable = true
|
||||
l.inQueue.mu.Unlock()
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
|
||||
return 0, err
|
||||
@@ -187,9 +187,9 @@ func (l *lineDiscipline) inputQueueRead(ctx context.Context, dst usermem.IOSeque
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.masterWaiter.Notify(waiter.EventOut)
|
||||
l.masterWaiter.Notify(waiter.WritableEvents)
|
||||
if pushed {
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
@@ -204,7 +204,7 @@ func (l *lineDiscipline) inputQueueWrite(ctx context.Context, src usermem.IOSequ
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
return n, nil
|
||||
}
|
||||
return 0, syserror.ErrWouldBlock
|
||||
@@ -222,9 +222,9 @@ func (l *lineDiscipline) outputQueueRead(ctx context.Context, dst usermem.IOSequ
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.replicaWaiter.Notify(waiter.EventOut)
|
||||
l.replicaWaiter.Notify(waiter.WritableEvents)
|
||||
if pushed {
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
@@ -239,7 +239,7 @@ func (l *lineDiscipline) outputQueueWrite(ctx context.Context, src usermem.IOSeq
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
return n, nil
|
||||
}
|
||||
return 0, syserror.ErrWouldBlock
|
||||
@@ -399,7 +399,7 @@ func (*inputQueueTransformer) transform(l *lineDiscipline, q *queue, buf []byte)
|
||||
// Anything written to the readBuf will have to be echoed.
|
||||
if l.termios.LEnabled(linux.ECHO) {
|
||||
l.outQueue.writeBytes(cBytes, l)
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
|
||||
// If we finish a line, make it available for reading.
|
||||
|
||||
@@ -71,7 +71,7 @@ func (q *queue) readReadiness(t *linux.KernelTermios) waiter.EventMask {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if len(q.readBuf) > 0 && q.readable {
|
||||
return waiter.EventIn
|
||||
return waiter.ReadableEvents
|
||||
}
|
||||
return waiter.EventMask(0)
|
||||
}
|
||||
@@ -81,7 +81,7 @@ func (q *queue) writeReadiness(t *linux.KernelTermios) waiter.EventMask {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if q.waitBufLen < waitBufMaxBytes {
|
||||
return waiter.EventOut
|
||||
return waiter.WritableEvents
|
||||
}
|
||||
return waiter.EventMask(0)
|
||||
}
|
||||
|
||||
@@ -141,7 +141,7 @@ func (l *lineDiscipline) setTermios(task *kernel.Task, args arch.SyscallArgument
|
||||
l.inQueue.pushWaitBufLocked(l)
|
||||
l.inQueue.readable = true
|
||||
l.inQueue.mu.Unlock()
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
|
||||
return 0, err
|
||||
@@ -185,9 +185,9 @@ func (l *lineDiscipline) inputQueueRead(ctx context.Context, dst usermem.IOSeque
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.masterWaiter.Notify(waiter.EventOut)
|
||||
l.masterWaiter.Notify(waiter.WritableEvents)
|
||||
if pushed {
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
@@ -202,7 +202,7 @@ func (l *lineDiscipline) inputQueueWrite(ctx context.Context, src usermem.IOSequ
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.replicaWaiter.Notify(waiter.EventIn)
|
||||
l.replicaWaiter.Notify(waiter.ReadableEvents)
|
||||
return n, nil
|
||||
}
|
||||
return 0, syserror.ErrWouldBlock
|
||||
@@ -220,9 +220,9 @@ func (l *lineDiscipline) outputQueueRead(ctx context.Context, dst usermem.IOSequ
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.replicaWaiter.Notify(waiter.EventOut)
|
||||
l.replicaWaiter.Notify(waiter.WritableEvents)
|
||||
if pushed {
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
@@ -237,7 +237,7 @@ func (l *lineDiscipline) outputQueueWrite(ctx context.Context, src usermem.IOSeq
|
||||
return 0, err
|
||||
}
|
||||
if n > 0 {
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
return n, nil
|
||||
}
|
||||
return 0, syserror.ErrWouldBlock
|
||||
@@ -397,7 +397,7 @@ func (*inputQueueTransformer) transform(l *lineDiscipline, q *queue, buf []byte)
|
||||
// Anything written to the readBuf will have to be echoed.
|
||||
if l.termios.LEnabled(linux.ECHO) {
|
||||
l.outQueue.writeBytes(cBytes, l)
|
||||
l.masterWaiter.Notify(waiter.EventIn)
|
||||
l.masterWaiter.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
|
||||
// If we finish a line, make it available for reading.
|
||||
|
||||
@@ -69,7 +69,7 @@ func (q *queue) readReadiness(t *linux.KernelTermios) waiter.EventMask {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if len(q.readBuf) > 0 && q.readable {
|
||||
return waiter.EventIn
|
||||
return waiter.ReadableEvents
|
||||
}
|
||||
return waiter.EventMask(0)
|
||||
}
|
||||
@@ -79,7 +79,7 @@ func (q *queue) writeReadiness(t *linux.KernelTermios) waiter.EventMask {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if q.waitBufLen < waitBufMaxBytes {
|
||||
return waiter.EventOut
|
||||
return waiter.WritableEvents
|
||||
}
|
||||
return waiter.EventMask(0)
|
||||
}
|
||||
|
||||
@@ -185,7 +185,7 @@ func (efd *EventFileDescription) read(ctx context.Context, dst usermem.IOSequenc
|
||||
// Notify writers. We do this even if we were already writable because
|
||||
// it is possible that a writer is waiting to write the maximum value
|
||||
// to the event.
|
||||
efd.queue.Notify(waiter.EventOut)
|
||||
efd.queue.Notify(waiter.WritableEvents)
|
||||
|
||||
var buf [8]byte
|
||||
usermem.ByteOrder.PutUint64(buf[:], val)
|
||||
@@ -238,7 +238,7 @@ func (efd *EventFileDescription) Signal(val uint64) error {
|
||||
efd.mu.Unlock()
|
||||
|
||||
// Always trigger a notification.
|
||||
efd.queue.Notify(waiter.EventIn)
|
||||
efd.queue.Notify(waiter.ReadableEvents)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -254,11 +254,11 @@ func (efd *EventFileDescription) Readiness(mask waiter.EventMask) waiter.EventMa
|
||||
|
||||
ready := waiter.EventMask(0)
|
||||
if efd.val > 0 {
|
||||
ready |= waiter.EventIn
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
|
||||
if efd.val < math.MaxUint64-1 {
|
||||
ready |= waiter.EventOut
|
||||
ready |= waiter.WritableEvents
|
||||
}
|
||||
|
||||
return mask & ready
|
||||
|
||||
@@ -49,7 +49,7 @@ func TestEventFD(t *testing.T) {
|
||||
|
||||
// Register a callback for a write event.
|
||||
w, ch := waiter.NewChannelEntry(nil)
|
||||
eventfd.EventRegister(&w, waiter.EventIn)
|
||||
eventfd.EventRegister(&w, waiter.ReadableEvents)
|
||||
defer eventfd.EventUnregister(&w)
|
||||
|
||||
data := []byte("00000124")
|
||||
|
||||
@@ -316,7 +316,7 @@ func (conn *connection) callFutureLocked(t *kernel.Task, r *Request) (*futureRes
|
||||
conn.fd.completions[r.id] = fut
|
||||
|
||||
// Signal the readers that there is something to read.
|
||||
conn.fd.waitQueue.Notify(waiter.EventIn)
|
||||
conn.fd.waitQueue.Notify(waiter.ReadableEvents)
|
||||
|
||||
return fut, nil
|
||||
}
|
||||
|
||||
@@ -368,10 +368,10 @@ func (fd *DeviceFD) readinessLocked(mask waiter.EventMask) waiter.EventMask {
|
||||
}
|
||||
|
||||
// FD is always writable.
|
||||
ready |= waiter.EventOut
|
||||
ready |= waiter.WritableEvents
|
||||
if !fd.queue.Empty() {
|
||||
// Have reqs available, FD is readable.
|
||||
ready |= waiter.EventIn
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
|
||||
return ready & mask
|
||||
|
||||
@@ -180,7 +180,7 @@ func ReadTest(serverTask *kernel.Task, fd *vfs.FileDescription, inIOseq usermem.
|
||||
|
||||
// Register for notifications.
|
||||
w, ch := waiter.NewChannelEntry(nil)
|
||||
dev.EventRegister(&w, waiter.EventIn)
|
||||
dev.EventRegister(&w, waiter.ReadableEvents)
|
||||
for {
|
||||
// Issue the request and break out if it completes with anything other than
|
||||
// "would block".
|
||||
|
||||
@@ -286,7 +286,7 @@ func (fs *filesystem) Release(ctx context.Context) {
|
||||
fs.umounted = true
|
||||
fs.conn.Abort(ctx)
|
||||
// Notify all the waiters on this fd.
|
||||
fs.conn.fd.waitQueue.Notify(waiter.EventIn)
|
||||
fs.conn.fd.waitQueue.Notify(waiter.ReadableEvents)
|
||||
|
||||
fs.conn.fd.mu.Unlock()
|
||||
|
||||
|
||||
@@ -192,7 +192,7 @@ func (c *ConnectedEndpoint) Writable() bool {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
|
||||
return fdnotifier.NonBlockingPoll(int32(c.fd), waiter.EventOut)&waiter.EventOut != 0
|
||||
return fdnotifier.NonBlockingPoll(int32(c.fd), waiter.WritableEvents)&waiter.WritableEvents != 0
|
||||
}
|
||||
|
||||
// Passcred implements transport.ConnectedEndpoint.Passcred.
|
||||
@@ -282,7 +282,7 @@ func (c *ConnectedEndpoint) Readable() bool {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
|
||||
return fdnotifier.NonBlockingPoll(int32(c.fd), waiter.EventIn)&waiter.EventIn != 0
|
||||
return fdnotifier.NonBlockingPoll(int32(c.fd), waiter.ReadableEvents)&waiter.ReadableEvents != 0
|
||||
}
|
||||
|
||||
// SendQueuedSize implements transport.Receiver.SendQueuedSize.
|
||||
|
||||
@@ -117,8 +117,8 @@ func (sfd *SignalFileDescription) Read(ctx context.Context, dst usermem.IOSequen
|
||||
func (sfd *SignalFileDescription) Readiness(mask waiter.EventMask) waiter.EventMask {
|
||||
sfd.mu.Lock()
|
||||
defer sfd.mu.Unlock()
|
||||
if mask&waiter.EventIn != 0 && sfd.target.PendingSignals()&sfd.mask != 0 {
|
||||
return waiter.EventIn // Pending signals.
|
||||
if mask&waiter.ReadableEvents != 0 && sfd.target.PendingSignals()&sfd.mask != 0 {
|
||||
return waiter.ReadableEvents // Pending signals.
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
@@ -105,7 +105,7 @@ func (tfd *TimerFileDescription) SetTime(s ktime.Setting) (ktime.Time, ktime.Set
|
||||
func (tfd *TimerFileDescription) Readiness(mask waiter.EventMask) waiter.EventMask {
|
||||
var ready waiter.EventMask
|
||||
if atomic.LoadUint64(&tfd.val) != 0 {
|
||||
ready |= waiter.EventIn
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
return ready
|
||||
}
|
||||
@@ -138,7 +138,7 @@ func (tfd *TimerFileDescription) Release(context.Context) {
|
||||
// Notify implements ktime.TimerListener.Notify.
|
||||
func (tfd *TimerFileDescription) Notify(exp uint64, setting ktime.Setting) (ktime.Setting, bool) {
|
||||
atomic.AddUint64(&tfd.val, exp)
|
||||
tfd.events.Notify(waiter.EventIn)
|
||||
tfd.events.Notify(waiter.ReadableEvents)
|
||||
return ktime.Setting{}, false
|
||||
}
|
||||
|
||||
|
||||
@@ -213,8 +213,8 @@ func (e *EventPoll) eventsAvailable() bool {
|
||||
func (e *EventPoll) Readiness(mask waiter.EventMask) waiter.EventMask {
|
||||
ready := waiter.EventMask(0)
|
||||
|
||||
if (mask&waiter.EventIn) != 0 && e.eventsAvailable() {
|
||||
ready |= waiter.EventIn
|
||||
if (mask&waiter.ReadableEvents) != 0 && e.eventsAvailable() {
|
||||
ready |= waiter.ReadableEvents
|
||||
}
|
||||
|
||||
return ready
|
||||
@@ -290,7 +290,7 @@ func (p *pollEntry) Callback(*waiter.Entry, waiter.EventMask) {
|
||||
p.curList = &e.readyList
|
||||
e.listsMu.Unlock()
|
||||
|
||||
e.Notify(waiter.EventIn)
|
||||
e.Notify(waiter.ReadableEvents)
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -45,7 +45,7 @@ func (e *EventPoll) afterLoad() {
|
||||
e.waitingList.Remove(entry)
|
||||
e.readyList.PushBack(entry)
|
||||
entry.curList = &e.readyList
|
||||
e.Notify(waiter.EventIn)
|
||||
e.Notify(waiter.ReadableEvents)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -29,7 +29,7 @@ func TestFileDestroyed(t *testing.T) {
|
||||
ctx := contexttest.Context(t)
|
||||
efile := NewEventPoll(ctx)
|
||||
e := efile.FileOperations.(*EventPoll)
|
||||
if err := e.AddEntry(id, 0, waiter.EventIn, [2]int32{}); err != nil {
|
||||
if err := e.AddEntry(id, 0, waiter.ReadableEvents, [2]int32{}); err != nil {
|
||||
t.Fatalf("addEntry failed: %v", err)
|
||||
}
|
||||
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user