Merge "Combine startPull and sendHeartbeat threads" into rvc-dev

This commit is contained in:
Ruchir Rastogi
2020-05-27 19:18:33 +00:00
committed by Android (Google) Code Review
2 changed files with 41 additions and 68 deletions

View File

@@ -41,13 +41,8 @@ void ShellSubscriber::startNewSubscription(int in, int out, int timeoutSec) {
{ {
std::unique_lock<std::mutex> lock(mMutex); std::unique_lock<std::mutex> lock(mMutex);
if (myToken != mToken) {
// Some other subscription has already come in. Stop.
return;
}
mSubscriptionInfo = mySubscriptionInfo; mSubscriptionInfo = mySubscriptionInfo;
spawnHelperThread(myToken);
spawnHelperThreadsLocked(mySubscriptionInfo, myToken);
waitForSubscriptionToEndLocked(mySubscriptionInfo, myToken, lock, timeoutSec); waitForSubscriptionToEndLocked(mySubscriptionInfo, myToken, lock, timeoutSec);
if (mSubscriptionInfo == mySubscriptionInfo) { if (mSubscriptionInfo == mySubscriptionInfo) {
@@ -57,14 +52,9 @@ void ShellSubscriber::startNewSubscription(int in, int out, int timeoutSec) {
} }
} }
void ShellSubscriber::spawnHelperThreadsLocked(shared_ptr<SubscriptionInfo> myInfo, int myToken) { void ShellSubscriber::spawnHelperThread(int myToken) {
if (!myInfo->mPulledInfo.empty() && myInfo->mPullIntervalMin > 0) { std::thread t([this, myToken] { pullAndSendHeartbeats(myToken); });
std::thread puller([this, myToken] { startPull(myToken); }); t.detach();
puller.detach();
}
std::thread heartbeatSender([this, myToken] { sendHeartbeats(myToken); });
heartbeatSender.detach();
} }
void ShellSubscriber::waitForSubscriptionToEndLocked(shared_ptr<SubscriptionInfo> myInfo, void ShellSubscriber::waitForSubscriptionToEndLocked(shared_ptr<SubscriptionInfo> myInfo,
@@ -114,13 +104,7 @@ bool ShellSubscriber::readConfig(shared_ptr<SubscriptionInfo> subscriptionInfo)
subscriptionInfo->mPushedMatchers.push_back(pushed); subscriptionInfo->mPushedMatchers.push_back(pushed);
} }
int minInterval = -1;
for (const auto& pulled : config.pulled()) { for (const auto& pulled : config.pulled()) {
// All intervals need to be multiples of the min interval.
if (minInterval < 0 || pulled.freq_millis() < minInterval) {
minInterval = pulled.freq_millis();
}
vector<string> packages; vector<string> packages;
vector<int32_t> uids; vector<int32_t> uids;
for (const string& pkg : pulled.packages()) { for (const string& pkg : pulled.packages()) {
@@ -136,18 +120,18 @@ bool ShellSubscriber::readConfig(shared_ptr<SubscriptionInfo> subscriptionInfo)
uids); uids);
VLOG("adding matcher for pulled atom %d", pulled.matcher().atom_id()); VLOG("adding matcher for pulled atom %d", pulled.matcher().atom_id());
} }
subscriptionInfo->mPullIntervalMin = minInterval;
return true; return true;
} }
void ShellSubscriber::startPull(int myToken) { void ShellSubscriber::pullAndSendHeartbeats(int myToken) {
VLOG("ShellSubscriber: pull thread %d starting", myToken); VLOG("ShellSubscriber: helper thread %d starting", myToken);
while (true) { while (true) {
int64_t sleepTimeMs = INT_MAX;
{ {
std::lock_guard<std::mutex> lock(mMutex); std::lock_guard<std::mutex> lock(mMutex);
if (!mSubscriptionInfo || mToken != myToken) { if (!mSubscriptionInfo || mToken != myToken) {
VLOG("ShellSubscriber: pulling thread %d done!", myToken); VLOG("ShellSubscriber: helper thread %d done!", myToken);
return; return;
} }
@@ -168,11 +152,27 @@ void ShellSubscriber::startPull(int myToken) {
pullInfo.mPrevPullElapsedRealtimeMs = nowMillis; pullInfo.mPrevPullElapsedRealtimeMs = nowMillis;
} }
// Send a heartbeat, consisting of a data size of 0, if perfd hasn't recently received
// data from statsd. When it receives the data size of 0, perfd will not expect any
// atoms and recheck whether the subscription should end.
if (nowMillis - mLastWriteMs > kMsBetweenHeartbeats) {
attemptWriteToPipeLocked(/*dataSize=*/0);
} }
VLOG("ShellSubscriber: pulling thread %d sleeping for %d ms", myToken, // Determine how long to sleep before doing more work.
mSubscriptionInfo->mPullIntervalMin); for (PullInfo& pullInfo : mSubscriptionInfo->mPulledInfo) {
std::this_thread::sleep_for(std::chrono::milliseconds(mSubscriptionInfo->mPullIntervalMin)); int64_t nextPullTime = pullInfo.mPrevPullElapsedRealtimeMs + pullInfo.mInterval;
int64_t timeBeforePull = nextPullTime - nowMillis; // guaranteed to be non-negative
if (timeBeforePull < sleepTimeMs) sleepTimeMs = timeBeforePull;
}
int64_t timeBeforeHeartbeat = (mLastWriteMs + kMsBetweenHeartbeats) - nowMillis;
if (timeBeforeHeartbeat < sleepTimeMs) sleepTimeMs = timeBeforeHeartbeat;
}
VLOG("ShellSubscriber: helper thread %d sleeping for %lld ms", myToken,
(long long)sleepTimeMs);
std::this_thread::sleep_for(std::chrono::milliseconds(sleepTimeMs));
} }
} }
@@ -200,7 +200,7 @@ void ShellSubscriber::writePulledAtomsLocked(const vector<std::shared_ptr<LogEve
} }
} }
if (count > 0) attemptWriteToSocketLocked(mProto.size()); if (count > 0) attemptWriteToPipeLocked(mProto.size());
} }
void ShellSubscriber::onLogEvent(const LogEvent& event) { void ShellSubscriber::onLogEvent(const LogEvent& event) {
@@ -214,26 +214,24 @@ void ShellSubscriber::onLogEvent(const LogEvent& event) {
util::FIELD_COUNT_REPEATED | FIELD_ID_ATOM); util::FIELD_COUNT_REPEATED | FIELD_ID_ATOM);
event.ToProto(mProto); event.ToProto(mProto);
mProto.end(atomToken); mProto.end(atomToken);
attemptWriteToSocketLocked(mProto.size()); attemptWriteToPipeLocked(mProto.size());
} }
} }
} }
// Tries to write the atom encoded in mProto to the socket. If the write fails // Tries to write the atom encoded in mProto to the pipe. If the write fails
// because the read end of the pipe has closed, signals to other threads that // because the read end of the pipe has closed, signals to other threads that
// the subscription should end. // the subscription should end.
void ShellSubscriber::attemptWriteToSocketLocked(size_t dataSize) { void ShellSubscriber::attemptWriteToPipeLocked(size_t dataSize) {
// First write the payload size. // First, write the payload size.
if (!android::base::WriteFully(mSubscriptionInfo->mOutputFd, &dataSize, sizeof(dataSize))) { if (!android::base::WriteFully(mSubscriptionInfo->mOutputFd, &dataSize, sizeof(dataSize))) {
mSubscriptionInfo->mClientAlive = false; mSubscriptionInfo->mClientAlive = false;
mSubscriptionShouldEnd.notify_one(); mSubscriptionShouldEnd.notify_one();
return; return;
} }
if (dataSize == 0) return; // Then, write the payload if this is not just a heartbeat.
if (dataSize > 0 && !mProto.flush(mSubscriptionInfo->mOutputFd)) {
// Then, write the payload.
if (!mProto.flush(mSubscriptionInfo->mOutputFd)) {
mSubscriptionInfo->mClientAlive = false; mSubscriptionInfo->mClientAlive = false;
mSubscriptionShouldEnd.notify_one(); mSubscriptionShouldEnd.notify_one();
return; return;
@@ -242,28 +240,6 @@ void ShellSubscriber::attemptWriteToSocketLocked(size_t dataSize) {
mLastWriteMs = getElapsedRealtimeMillis(); mLastWriteMs = getElapsedRealtimeMillis();
} }
// Send a heartbeat, consisting solely of a data size of 0, if perfd has not
// recently received any writes from statsd. When it receives the data size of
// 0, perfd will not expect any data and recheck whether the shell command is
// still running.
void ShellSubscriber::sendHeartbeats(int myToken) {
while (true) {
{
std::lock_guard<std::mutex> lock(mMutex);
if (!mSubscriptionInfo || myToken != mToken) {
VLOG("ShellSubscriber: heartbeat thread %d done!", myToken);
return;
}
if (getElapsedRealtimeMillis() - mLastWriteMs > kMsBetweenHeartbeats) {
VLOG("ShellSubscriber: sending a heartbeat to perfd");
attemptWriteToSocketLocked(/*dataSize=*/0);
}
}
std::this_thread::sleep_for(std::chrono::milliseconds(kMsBetweenHeartbeats));
}
}
} // namespace statsd } // namespace statsd
} // namespace os } // namespace os
} // namespace android } // namespace android

View File

@@ -92,7 +92,6 @@ private:
int mOutputFd; int mOutputFd;
std::vector<SimpleAtomMatcher> mPushedMatchers; std::vector<SimpleAtomMatcher> mPushedMatchers;
std::vector<PullInfo> mPulledInfo; std::vector<PullInfo> mPulledInfo;
int mPullIntervalMin;
bool mClientAlive; bool mClientAlive;
}; };
@@ -100,27 +99,25 @@ private:
bool readConfig(std::shared_ptr<SubscriptionInfo> subscriptionInfo); bool readConfig(std::shared_ptr<SubscriptionInfo> subscriptionInfo);
void spawnHelperThreadsLocked(std::shared_ptr<SubscriptionInfo> myInfo, int myToken); void spawnHelperThread(int myToken);
void waitForSubscriptionToEndLocked(std::shared_ptr<SubscriptionInfo> myInfo, void waitForSubscriptionToEndLocked(std::shared_ptr<SubscriptionInfo> myInfo,
int myToken, int myToken,
std::unique_lock<std::mutex>& lock, std::unique_lock<std::mutex>& lock,
int timeoutSec); int timeoutSec);
void startPull(int myToken); // Helper thread that pulls atoms at a regular frequency and sends
// heartbeats to perfd if statsd hasn't recently sent any data. Statsd must
// send heartbeats for perfd to escape a blocking read call and recheck if
// the user has terminated the subscription.
void pullAndSendHeartbeats(int myToken);
void writePulledAtomsLocked(const vector<std::shared_ptr<LogEvent>>& data, void writePulledAtomsLocked(const vector<std::shared_ptr<LogEvent>>& data,
const SimpleAtomMatcher& matcher); const SimpleAtomMatcher& matcher);
void getUidsForPullAtom(vector<int32_t>* uids, const PullInfo& pullInfo); void getUidsForPullAtom(vector<int32_t>* uids, const PullInfo& pullInfo);
void attemptWriteToSocketLocked(size_t dataSize); void attemptWriteToPipeLocked(size_t dataSize);
// Send ocassional heartbeats for two reasons: (a) for statsd to detect when
// the read end of the pipe has closed and (b) for perfd to escape a
// blocking read call and recheck if the user has terminated the
// subscription.
void sendHeartbeats(int myToken);
sp<UidMap> mUidMap; sp<UidMap> mUidMap;