Merge "[dataloader] fix fd leak" into rvc-dev am: 8caaf33214

Change-Id: Ia1390f0a745e8cdf6bfc0a1a0c5f4709a93399e5
This commit is contained in:
Automerger Merge Worker
2020-03-16 05:52:53 +00:00

View File

@@ -364,13 +364,8 @@ private:
} }
} }
void onDestroy() final { void onDestroy() final {
ALOGE("Sending EXIT to server.");
sendRequest(mOutFd, EXIT);
// Make sure the receiver thread stopped. // Make sure the receiver thread stopped.
CHECK(!mReceiverThread.joinable()); CHECK(!mReceiverThread.joinable());
mInFd.reset();
mOutFd.reset();
} }
// Installation. // Installation.
@@ -568,13 +563,6 @@ private:
// Streaming. // Streaming.
bool initStreaming(unique_fd inout) { bool initStreaming(unique_fd inout) {
mInFd.reset(dup(inout));
mOutFd.reset(dup(inout));
if (mInFd < 0 || mOutFd < 0) {
ALOGE("Failed to dup FDs.");
return false;
}
mEventFd.reset(eventfd(0, EFD_CLOEXEC)); mEventFd.reset(eventfd(0, EFD_CLOEXEC));
if (mEventFd < 0) { if (mEventFd < 0) {
ALOGE("Failed to create eventfd."); ALOGE("Failed to create eventfd.");
@@ -583,7 +571,7 @@ private:
// Awaiting adb handshake. // Awaiting adb handshake.
char okay_buf[OKAY.size()]; char okay_buf[OKAY.size()];
if (!android::base::ReadFully(mInFd, okay_buf, OKAY.size())) { if (!android::base::ReadFully(inout, okay_buf, OKAY.size())) {
ALOGE("Failed to receive OKAY. Abort."); ALOGE("Failed to receive OKAY. Abort.");
return false; return false;
} }
@@ -593,13 +581,23 @@ private:
return false; return false;
} }
mReceiverThread = std::thread([this]() { receiver(); }); {
std::lock_guard lock{mOutFdLock};
mOutFd.reset(::dup(inout));
if (mOutFd < 0) {
ALOGE("Failed to create streaming fd.");
}
}
mReceiverThread =
std::thread([this, io = std::move(inout)]() mutable { receiver(std::move(io)); });
ALOGI("Started streaming..."); ALOGI("Started streaming...");
return true; return true;
} }
// IFS callbacks. // IFS callbacks.
void onPendingReads(dataloader::PendingReads pendingReads) final { void onPendingReads(dataloader::PendingReads pendingReads) final {
std::lock_guard lock{mOutFdLock};
CHECK(mIfs); CHECK(mIfs);
for (auto&& pendingRead : pendingReads) { for (auto&& pendingRead : pendingReads) {
const android::dataloader::FileId& fileId = pendingRead.id; const android::dataloader::FileId& fileId = pendingRead.id;
@@ -625,12 +623,12 @@ private:
} }
} }
void receiver() { void receiver(unique_fd inout) {
std::vector<uint8_t> data; std::vector<uint8_t> data;
std::vector<IncFsDataBlock> instructions; std::vector<IncFsDataBlock> instructions;
std::unordered_map<FileIdx, unique_fd> writeFds; std::unordered_map<FileIdx, unique_fd> writeFds;
while (!mStopReceiving) { while (!mStopReceiving) {
const int res = waitForDataOrSignal(mInFd, mEventFd); const int res = waitForDataOrSignal(inout, mEventFd);
if (res == 0) { if (res == 0) {
continue; continue;
} }
@@ -640,10 +638,11 @@ private:
break; break;
} }
if (res == mEventFd) { if (res == mEventFd) {
ALOGE("Received stop signal. Exit."); ALOGE("Received stop signal. Sending EXIT to server.");
sendRequest(inout, EXIT);
break; break;
} }
if (!readChunk(mInFd, data)) { if (!readChunk(inout, data)) {
ALOGE("Failed to read a message. Abort."); ALOGE("Failed to read a message. Abort.");
mStatusListener->reportStatus(DATA_LOADER_NO_CONNECTION); mStatusListener->reportStatus(DATA_LOADER_NO_CONNECTION);
break; break;
@@ -656,7 +655,7 @@ private:
ALOGI("Stop signal received. Sending exit command (remaining bytes: %d).", ALOGI("Stop signal received. Sending exit command (remaining bytes: %d).",
int(remainingData.size())); int(remainingData.size()));
sendRequest(mOutFd, EXIT); sendRequest(inout, EXIT);
mStopReceiving = true; mStopReceiving = true;
break; break;
} }
@@ -699,6 +698,11 @@ private:
writeInstructions(instructions); writeInstructions(instructions);
} }
writeInstructions(instructions); writeInstructions(instructions);
{
std::lock_guard lock{mOutFdLock};
mOutFd.reset();
}
} }
void writeInstructions(std::vector<IncFsDataBlock>& instructions) { void writeInstructions(std::vector<IncFsDataBlock>& instructions) {
@@ -742,7 +746,7 @@ private:
std::string mArgs; std::string mArgs;
android::dataloader::FilesystemConnectorPtr mIfs = nullptr; android::dataloader::FilesystemConnectorPtr mIfs = nullptr;
android::dataloader::StatusListenerPtr mStatusListener = nullptr; android::dataloader::StatusListenerPtr mStatusListener = nullptr;
android::base::unique_fd mInFd; std::mutex mOutFdLock;
android::base::unique_fd mOutFd; android::base::unique_fd mOutFd;
android::base::unique_fd mEventFd; android::base::unique_fd mEventFd;
std::thread mReceiverThread; std::thread mReceiverThread;