From 8e22ce50198bef6ce182cb1426b012400a348ef3 Mon Sep 17 00:00:00 2001 From: Jamie Liu Date: Wed, 19 Jan 2022 13:44:18 -0800 Subject: [PATCH] Consistently order Pipe.mu before other file mutexes and MM.activeMu. PiperOrigin-RevId: 422894869 --- pkg/sentry/kernel/pipe/pipe_util.go | 46 ++++++++------ test/syscalls/linux/BUILD | 2 + test/syscalls/linux/splice.cc | 96 +++++++++++++++++++++++++++++ 3 files changed, 126 insertions(+), 18 deletions(-) diff --git a/pkg/sentry/kernel/pipe/pipe_util.go b/pkg/sentry/kernel/pipe/pipe_util.go index c20ce7325..2c3a0fff0 100644 --- a/pkg/sentry/kernel/pipe/pipe_util.go +++ b/pkg/sentry/kernel/pipe/pipe_util.go @@ -43,21 +43,26 @@ func (p *Pipe) Release(context.Context) { // Read reads from the Pipe into dst. func (p *Pipe) Read(ctx context.Context, dst usermem.IOSequence) (int64, error) { - n, err := dst.CopyOutFrom(ctx, p) + n, err := p.read(dst.NumBytes(), func(srcs safemem.BlockSeq) (uint64, error) { + var done uint64 + for !srcs.IsEmpty() { + src := srcs.Head() + n, err := dst.CopyOut(ctx, src.ToSlice()) + done += uint64(n) + if err != nil { + return done, err + } + dst = dst.DropFirst(n) + srcs = srcs.Tail() + } + return done, nil + }, true /* removeFromSrc */) if n > 0 { p.queue.Notify(waiter.WritableEvents) } return n, err } -// ReadToBlocks implements safemem.Reader.ReadToBlocks for Pipe.Read. -func (p *Pipe) ReadToBlocks(dsts safemem.BlockSeq) (uint64, error) { - n, err := p.read(int64(dsts.NumBytes()), func(srcs safemem.BlockSeq) (uint64, error) { - return safemem.CopySeq(dsts, srcs) - }, true /* removeFromSrc */) - return uint64(n), err -} - func (p *Pipe) read(count int64, f func(srcs safemem.BlockSeq) (uint64, error), removeFromSrc bool) (int64, error) { p.mu.Lock() defer p.mu.Unlock() @@ -81,7 +86,20 @@ func (p *Pipe) WriteTo(ctx context.Context, w io.Writer, count int64, dup bool) // Write writes to the Pipe from src. func (p *Pipe) Write(ctx context.Context, src usermem.IOSequence) (int64, error) { - n, err := src.CopyInTo(ctx, p) + n, err := p.write(src.NumBytes(), func(dsts safemem.BlockSeq) (uint64, error) { + var done uint64 + for !dsts.IsEmpty() { + dst := dsts.Head() + n, err := src.CopyIn(ctx, dst.ToSlice()) + done += uint64(n) + if err != nil { + return done, err + } + src = src.DropFirst(n) + dsts = dsts.Tail() + } + return done, nil + }) if n > 0 { p.queue.Notify(waiter.ReadableEvents) } @@ -94,14 +112,6 @@ func (p *Pipe) Write(ctx context.Context, src usermem.IOSequence) (int64, error) return n, err } -// WriteFromBlocks implements safemem.Writer.WriteFromBlocks for Pipe.Write. -func (p *Pipe) WriteFromBlocks(srcs safemem.BlockSeq) (uint64, error) { - n, err := p.write(int64(srcs.NumBytes()), func(dsts safemem.BlockSeq) (uint64, error) { - return safemem.CopySeq(dsts, srcs) - }) - return uint64(n), err -} - func (p *Pipe) write(count int64, f func(safemem.BlockSeq) (uint64, error)) (int64, error) { p.mu.Lock() defer p.mu.Unlock() diff --git a/test/syscalls/linux/BUILD b/test/syscalls/linux/BUILD index ed9e3b80c..8afa2b410 100644 --- a/test/syscalls/linux/BUILD +++ b/test/syscalls/linux/BUILD @@ -2310,9 +2310,11 @@ cc_binary( linkstatic = 1, deps = [ "//test/util:file_descriptor", + "@com_google_absl//absl/cleanup", "@com_google_absl//absl/strings", "@com_google_absl//absl/time", gtest, + "//test/util:memory_util", "//test/util:signal_util", "//test/util:temp_path", "//test/util:test_main", diff --git a/test/syscalls/linux/splice.cc b/test/syscalls/linux/splice.cc index 4a10ae8d2..c9ec489ba 100644 --- a/test/syscalls/linux/splice.cc +++ b/test/syscalls/linux/splice.cc @@ -22,10 +22,12 @@ #include "gmock/gmock.h" #include "gtest/gtest.h" +#include "absl/cleanup/cleanup.h" #include "absl/strings/string_view.h" #include "absl/time/clock.h" #include "absl/time/time.h" #include "test/util/file_descriptor.h" +#include "test/util/memory_util.h" #include "test/util/signal_util.h" #include "test/util/temp_path.h" #include "test/util/test_util.h" @@ -830,6 +832,100 @@ TEST(SpliceTest, ToPipeWithSmallCapacityDoesNotSpin) { EXPECT_EQ(signaled, 1); } +// Regression test for b/208679047. +TEST(SpliceTest, FromPipeWithConcurrentIo) { + // Create a file containing two copies of the same byte. Two bytes are + // necessary because both the read() and splice() loops below advance the file + // offset by one byte before lseek(); use of the file offset is required since + // the mutex protecting the file offset is implicated in the circular lock + // ordering that this test attempts to reproduce. + // + // This can't use memfd_create() because, in Linux, memfd_create(2) creates a + // struct file using alloc_file_pseudo() without going through + // do_dentry_open(), so FMODE_ATOMIC_POS is not set despite the created file + // having type S_IFREG ("regular file"). + const TempPath file = ASSERT_NO_ERRNO_AND_VALUE(TempPath::CreateFile()); + const FileDescriptor fd = + ASSERT_NO_ERRNO_AND_VALUE(Open(file.path(), O_RDWR)); + constexpr char kSplicedByte = 0x01; + for (int i = 0; i < 2; i++) { + ASSERT_THAT(WriteFd(fd.get(), &kSplicedByte, 1), + SyscallSucceedsWithValue(1)); + } + + // Create a pipe. + int pipe_fds[2]; + ASSERT_THAT(pipe(pipe_fds), SyscallSucceeds()); + const FileDescriptor rfd(pipe_fds[0]); + FileDescriptor wfd(pipe_fds[1]); + + DisableSave ds; + std::atomic done(false); + + // Create a thread that reads from fd until the end of the test. + ScopedThread memfd_reader([&] { + char file_buf; + while (!done.load()) { + ASSERT_THAT(lseek(fd.get(), 0, SEEK_SET), SyscallSucceeds()); + int n = ReadFd(fd.get(), &file_buf, 1); + if (n == 0) { + // fd was at offset 2 (EOF). In Linux, this is possible even after + // lseek(0) because splice() doesn't attempt atomicity with respect to + // concurrent lseek(), so the effect of lseek() may be lost. + continue; + } + ASSERT_THAT(n, SyscallSucceedsWithValue(1)); + ASSERT_EQ(file_buf, kSplicedByte); + } + }); + + // Create a thread that reads from the pipe until the end of the test. + ScopedThread pipe_reader([&] { + char pipe_buf; + while (!done.load()) { + int n = ReadFd(rfd.get(), &pipe_buf, 1); + if (n == 0) { + // This should only happen due to cleanup_threads (below) closing wfd. + EXPECT_TRUE(done.load()); + return; + } + ASSERT_THAT(n, SyscallSucceedsWithValue(1)); + ASSERT_EQ(pipe_buf, kSplicedByte); + } + }); + + // Create a thread that repeatedly invokes madvise(MADV_DONTNEED) on the same + // page of memory. (Having a thread attempt to lock MM.activeMu for writing is + // necessary to create a deadlock from the circular lock ordering, since + // otherwise both uses of MM.activeMu are for reading and may proceed + // concurrently.) + ScopedThread mm_locker([&] { + const Mapping m = ASSERT_NO_ERRNO_AND_VALUE( + MmapAnon(kPageSize, PROT_READ | PROT_WRITE, MAP_PRIVATE)); + while (!done.load()) { + madvise(m.ptr(), kPageSize, MADV_DONTNEED); + } + }); + + // This must come after the ScopedThreads since its destructor must run before + // theirs. + const absl::Cleanup cleanup_threads = [&] { + done.store(true); + // Ensure that pipe_reader is unblocked after setting done, so that it will + // be able to observe done being true. + wfd.reset(); + }; + + // Repeatedly splice from memfd to the pipe. The test passes if this does not + // deadlock. + const int kIterations = 5000; + for (int i = 0; i < kIterations; i++) { + ASSERT_THAT(lseek(fd.get(), 0, SEEK_SET), SyscallSucceeds()); + ASSERT_THAT(splice(fd.get(), nullptr, wfd.get(), nullptr, 1, 0), + SyscallSucceedsWithValue(1)); + } +} + } // namespace } // namespace testing