package daemon: Handle location info request in app thread

* ... instead of queuing it for the job thread. The advantage is that
the request will be handled immediately and clients won't have to wait
for transactions (which may even require user feedback) to finish. It
complicates Volume a bit, since there are now two threads that may
access it. The shared data have been moved to a State object which is
protected by a lock.
* For commit transaction requests check whether another package request
is already pending/in progress before queuing a job. Fail immediately,
if there is.

Fixes bug #10039.
This commit is contained in:
Ingo Weinhold
2014-02-07 01:21:57 +01:00
parent 3472fc553e
commit 32cae72412
9 changed files with 408 additions and 99 deletions
+1
View File
@@ -18,6 +18,7 @@ namespace BPrivate {
enum BDaemonError { enum BDaemonError {
B_DAEMON_OK = 0, B_DAEMON_OK = 0,
B_DAEMON_INSTALLATION_LOCATION_BUSY,
B_DAEMON_CHANGE_COUNT_MISMATCH, B_DAEMON_CHANGE_COUNT_MISMATCH,
B_DAEMON_BAD_REQUEST, B_DAEMON_BAD_REQUEST,
B_DAEMON_NO_SUCH_PACKAGE, B_DAEMON_NO_SUCH_PACKAGE,
+4 -1
View File
@@ -320,6 +320,9 @@ BDaemonClient::BCommitTransactionResult::FullErrorMessage() const
const char* errorString; const char* errorString;
if (fError > 0) { if (fError > 0) {
switch ((BDaemonError)fError) { switch ((BDaemonError)fError) {
case B_DAEMON_INSTALLATION_LOCATION_BUSY:
errorString = "another package operation already in progress";
break;
case B_DAEMON_CHANGE_COUNT_MISMATCH: case B_DAEMON_CHANGE_COUNT_MISMATCH:
errorString = "transaction out of date"; errorString = "transaction out of date";
break; break;
@@ -339,7 +342,7 @@ BDaemonClient::BCommitTransactionResult::FullErrorMessage() const
} }
} else } else
errorString = strerror(fError); errorString = strerror(fError);
BString result; BString result;
if (!fErrorMessage.IsEmpty()) { if (!fErrorMessage.IsEmpty()) {
result = fErrorMessage; result = fErrorMessage;
+26 -1
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2013, Haiku, Inc. All Rights Reserved. * Copyright 2013-2014, Haiku, Inc. All Rights Reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -12,6 +12,9 @@
#include <PthreadMutexLocker.h> #include <PthreadMutexLocker.h>
// #pragma mark - JobQueue
JobQueue::JobQueue() JobQueue::JobQueue()
: :
fMutexInitialized(false), fMutexInitialized(false),
@@ -97,3 +100,25 @@ JobQueue::DequeueJob()
return NULL; return NULL;
} }
void
JobQueue::DeleteJobs(Filter* filter)
{
PthreadMutexLocker mutexLocker(fMutex);
for (JobList::Iterator it = fJobs.GetIterator(); Job* job = it.Next();) {
if (filter->FilterJob(job)) {
it.Remove();
delete job;
}
}
}
// #pragma mark - Filter
JobQueue::Filter::~Filter()
{
}
+14 -1
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2013, Haiku, Inc. All Rights Reserved. * Copyright 2013-2014, Haiku, Inc. All Rights Reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -15,6 +15,9 @@
class JobQueue { class JobQueue {
public:
class Filter;
public: public:
JobQueue(); JobQueue();
~JobQueue(); ~JobQueue();
@@ -27,6 +30,8 @@ public:
Job* DequeueJob(); Job* DequeueJob();
// returns a reference // returns a reference
void DeleteJobs(Filter* filter);
private: private:
typedef DoublyLinkedList<Job> JobList; typedef DoublyLinkedList<Job> JobList;
@@ -40,4 +45,12 @@ private:
}; };
class JobQueue::Filter {
public:
virtual ~Filter();
virtual bool FilterJob(Job* job) = 0;
};
#endif // JOB_QUEUE_H #endif // JOB_QUEUE_H
+1 -1
View File
@@ -132,7 +132,7 @@ PackageManager::InitInstalledRepository(InstalledRepository& repository)
if (Volume* volume = fRoot->GetVolume(repository.Location())) { if (Volume* volume = fRoot->GetVolume(repository.Location())) {
for (PackageFileNameHashTable::Iterator it for (PackageFileNameHashTable::Iterator it
= volume->PackagesByFileName().GetIterator(); it.HasNext();) { = volume->PackagesByFileNameIterator(); it.HasNext();) {
Package* package = it.Next(); Package* package = it.Next();
if (package->IsActive()) { if (package->IsActive()) {
BSolverPackage* solverPackage; BSolverPackage* solverPackage;
+130 -36
View File
@@ -30,13 +30,36 @@ using namespace BPackageKit::BPrivate;
using namespace BPackageKit::BManager::BPrivate; using namespace BPackageKit::BManager::BPrivate;
static const bigtime_t kCommunicationTimeout = 1000000;
// #pragma mark - AbstractVolumeJob
struct Root::AbstractVolumeJob : public Job {
AbstractVolumeJob(Volume* volume)
:
fVolume(volume)
{
}
Volume* GetVolume() const
{
return fVolume;
}
protected:
Volume* fVolume;
};
// #pragma mark - VolumeJob // #pragma mark - VolumeJob
struct Root::VolumeJob : public Job { struct Root::VolumeJob : public AbstractVolumeJob {
VolumeJob(Volume* volume, void (Root::*method)(Volume*)) VolumeJob(Volume* volume, void (Root::*method)(Volume*))
: :
fVolume(volume), AbstractVolumeJob(volume),
fMethod(method) fMethod(method)
{ {
} }
@@ -47,25 +70,49 @@ struct Root::VolumeJob : public Job {
} }
private: private:
Volume* fVolume;
void (Root::*fMethod)(Volume*); void (Root::*fMethod)(Volume*);
}; };
// #pragma mark - RequestJob // #pragma mark - ProcessNodeMonitorEventsJob
struct Root::RequestJob : public Job { struct Root::ProcessNodeMonitorEventsJob : public VolumeJob {
RequestJob(Root* root, BMessage* message) ProcessNodeMonitorEventsJob(Volume* volume, void (Root::*method)(Volume*))
: :
VolumeJob(volume, method)
{
fVolume->PackageJobPending();
}
~ProcessNodeMonitorEventsJob()
{
fVolume->PackageJobFinished();
}
};
// #pragma mark - CommitTransactionJob
struct Root::CommitTransactionJob : public AbstractVolumeJob {
CommitTransactionJob(Root* root, Volume* volume, BMessage* message)
:
AbstractVolumeJob(volume),
fRoot(root), fRoot(root),
fMessage(message) fMessage(message)
{ {
fVolume->PackageJobPending();
}
~CommitTransactionJob()
{
fVolume->PackageJobFinished();
} }
virtual void Do() virtual void Do()
{ {
fRoot->_HandleRequest(fMessage.Get()); fRoot->_CommitTransaction(fVolume, fMessage.Get());
} }
private: private:
@@ -74,6 +121,27 @@ private:
}; };
// #pragma mark - VolumeJobFilter
struct Root::VolumeJobFilter : public ::JobQueue::Filter {
VolumeJobFilter(Volume* volume)
:
fVolume(volume)
{
}
virtual bool FilterJob(Job* job)
{
AbstractVolumeJob* volumeJob = dynamic_cast<AbstractVolumeJob*>(job);
return volumeJob != NULL && volumeJob->GetVolume() == fVolume;
}
private:
Volume* fVolume;
};
// #pragma mark - Root // #pragma mark - Root
@@ -236,21 +304,64 @@ Root::GetVolume(BPackageInstallationLocation location)
void void
Root::HandleRequest(BMessage* message) Root::HandleRequest(BMessage* message)
{ {
RequestJob* job = new(std::nothrow) RequestJob(this, message); ObjectDeleter<BMessage> messageDeleter(message);
if (job == NULL) {
delete message; // get the location and the volume
int32 location;
if (message->FindInt32("location", &location) != B_OK
|| location < 0
|| location >= B_PACKAGE_INSTALLATION_LOCATION_ENUM_COUNT) {
return; return;
} }
_QueueJob(job); AutoLocker<BLocker> locker(fLock);
Volume* volume = GetVolume((BPackageInstallationLocation)location);
if (volume == NULL)
return;
switch (message->what) {
case B_MESSAGE_GET_INSTALLATION_LOCATION_INFO:
volume->HandleGetLocationInfoRequest(message);
break;
case B_MESSAGE_COMMIT_TRANSACTION:
{
// The B_MESSAGE_COMMIT_TRANSACTION request must be handled in the
// job thread. But only queue a job, if there aren't package jobs
// pending already.
if (volume->IsPackageJobPending()) {
BMessage reply(B_MESSAGE_COMMIT_TRANSACTION_REPLY);
if (reply.AddInt32("error", B_DAEMON_INSTALLATION_LOCATION_BUSY)
== B_OK) {
message->SendReply(&reply, (BHandler*)NULL,
kCommunicationTimeout);
}
return;
}
CommitTransactionJob* job = new(std::nothrow) CommitTransactionJob(
this, volume, message);
if (job == NULL)
return;
messageDeleter.Detach();
_QueueJob(job);
break;
}
default:
break;
}
} }
void void
Root::VolumeNodeMonitorEventOccurred(Volume* volume) Root::VolumeNodeMonitorEventOccurred(Volume* volume)
{ {
_QueueJob( _QueueJob(new(std::nothrow) ProcessNodeMonitorEventsJob(volume,
new(std::nothrow) VolumeJob(volume, &Root::_ProcessNodeMonitorEvents)); &Root::_ProcessNodeMonitorEvents));
} }
@@ -318,6 +429,10 @@ Root::_InitPackages(Volume* volume)
void void
Root::_DeleteVolume(Volume* volume) Root::_DeleteVolume(Volume* volume)
{ {
// delete all pending jobs for that volume
VolumeJobFilter filter(volume);
fJobQueue.DeleteJobs(&filter);
delete volume; delete volume;
} }
@@ -364,30 +479,9 @@ Root::_ProcessNodeMonitorEvents(Volume* volume)
void void
Root::_HandleRequest(BMessage* message) Root::_CommitTransaction(Volume* volume, BMessage* message)
{ {
int32 location; volume->HandleCommitTransactionRequest(message);
if (message->FindInt32("location", &location) != B_OK
|| location < 0
|| location >= B_PACKAGE_INSTALLATION_LOCATION_ENUM_COUNT) {
return;
}
// get the volume and let it handle the message
AutoLocker<BLocker> locker(fLock);
Volume* volume = GetVolume((BPackageInstallationLocation)location);
locker.Unlock();
if (volume != NULL) {
switch (message->what) {
case B_MESSAGE_GET_INSTALLATION_LOCATION_INFO:
volume->HandleGetLocationInfoRequest(message);
break;
case B_MESSAGE_COMMIT_TRANSACTION:
volume->HandleCommitTransactionRequest(message);
break;
}
}
} }
+7 -3
View File
@@ -55,10 +55,13 @@ protected:
virtual void LastReferenceReleased(); virtual void LastReferenceReleased();
private: private:
struct AbstractVolumeJob;
struct VolumeJob; struct VolumeJob;
struct RequestJob; struct ProcessNodeMonitorEventsJob;
struct CommitTransactionJob;
struct VolumeJobFilter;
friend struct RequestJob; friend struct CommitTransactionJob;
private: private:
Volume** _GetVolume(PackageFSMountType mountType); Volume** _GetVolume(PackageFSMountType mountType);
@@ -67,7 +70,8 @@ private:
void _InitPackages(Volume* volume); void _InitPackages(Volume* volume);
void _DeleteVolume(Volume* volume); void _DeleteVolume(Volume* volume);
void _ProcessNodeMonitorEvents(Volume* volume); void _ProcessNodeMonitorEvents(Volume* volume);
void _HandleRequest(BMessage* message); void _CommitTransaction(Volume* volume,
BMessage* message);
status_t _QueueJob(Job* job); status_t _QueueJob(Job* job);
+200 -48
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2013, Haiku, Inc. All Rights Reserved. * Copyright 2013-2014, Haiku, Inc. All Rights Reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -105,6 +105,146 @@ private:
}; };
// #pragma mark - State
struct Volume::State {
State()
:
fLock("volume state"),
fPackagesByFileName(),
fPackagesByNodeRef(),
fChangeCount(0),
fPendingPackageJobCount(0)
{
}
~State()
{
fPackagesByFileName.Clear();
Package* package = fPackagesByNodeRef.Clear(true);
while (package != NULL) {
Package* next = package->NodeRefHashTableNext();
delete package;
package = next;
}
}
bool Init()
{
return fLock.InitCheck() == B_OK && fPackagesByFileName.Init() == B_OK
&& fPackagesByNodeRef.Init() == B_OK;
}
bool Lock()
{
return fLock.Lock();
}
void Unlock()
{
fLock.Unlock();
}
int64 ChangeCount() const
{
return fChangeCount;
}
Package* FindPackage(const char* name) const
{
return fPackagesByFileName.Lookup(name);
}
Package* FindPackage(const node_ref& nodeRef) const
{
return fPackagesByNodeRef.Lookup(nodeRef);
}
PackageFileNameHashTable::Iterator ByFileNameIterator() const
{
return fPackagesByFileName.GetIterator();
}
PackageNodeRefHashTable::Iterator ByNodeRefIterator() const
{
return fPackagesByNodeRef.GetIterator();
}
void AddPackage(Package* package)
{
AutoLocker<BLocker> locker(fLock);
fPackagesByFileName.Insert(package);
fPackagesByNodeRef.Insert(package);
}
void RemovePackage(Package* package)
{
AutoLocker<BLocker> locker(fLock);
_RemovePackage(package);
}
void SetPackageActive(Package* package, bool active)
{
AutoLocker<BLocker> locker(fLock);
package->SetActive(active);
}
void ActivationChanged(const PackageSet& activatedPackage,
const PackageSet& deactivatePackages)
{
AutoLocker<BLocker> locker(fLock);
for (PackageSet::iterator it = activatedPackage.begin();
it != activatedPackage.end(); ++it) {
(*it)->SetActive(true);
fChangeCount++;
}
for (PackageSet::iterator it = deactivatePackages.begin();
it != deactivatePackages.end(); ++it) {
Package* package = *it;
_RemovePackage(package);
delete package;
}
}
void PackageJobPending()
{
atomic_add(&fPendingPackageJobCount, 1);
}
void PackageJobFinished()
{
atomic_add(&fPendingPackageJobCount, -1);
}
bool IsPackageJobPending() const
{
return fPendingPackageJobCount != 0;
}
private:
void _RemovePackage(Package* package)
{
fPackagesByFileName.Remove(package);
fPackagesByNodeRef.Remove(package);
fChangeCount++;
}
private:
BLocker fLock;
PackageFileNameHashTable fPackagesByFileName;
PackageNodeRefHashTable fPackagesByNodeRef;
int64 fChangeCount;
int32 fPendingPackageJobCount;
};
// #pragma mark - CommitTransactionHandler // #pragma mark - CommitTransactionHandler
@@ -159,7 +299,7 @@ struct Volume::CommitTransactionHandler {
BMessage* reply) BMessage* reply)
{ {
// check the change count // check the change count
if (transaction.ChangeCount() != fVolume->fChangeCount) if (transaction.ChangeCount() != fVolume->fState->ChangeCount())
throw Exception(B_DAEMON_CHANGE_COUNT_MISMATCH); throw Exception(B_DAEMON_CHANGE_COUNT_MISMATCH);
// collect the packages to deactivate // collect the packages to deactivate
@@ -230,7 +370,7 @@ private:
for (int32 i = 0; i < packagesToDeactivateCount; i++) { for (int32 i = 0; i < packagesToDeactivateCount; i++) {
BString packageName = packagesToDeactivate.StringAt(i); BString packageName = packagesToDeactivate.StringAt(i);
Package* package = fVolume->fPackagesByFileName.Lookup(packageName); Package* package = fVolume->fState->FindPackage(packageName);
if (package == NULL) { if (package == NULL) {
throw Exception(B_DAEMON_NO_SUCH_PACKAGE, "no such package", throw Exception(B_DAEMON_NO_SUCH_PACKAGE, "no such package",
packageName); packageName);
@@ -285,7 +425,7 @@ private:
BString packageName = packagesToActivate.StringAt(i); BString packageName = packagesToActivate.StringAt(i);
// make sure it doesn't clash with an already existing package // make sure it doesn't clash with an already existing package
Package* package = fVolume->fPackagesByFileName.Lookup(packageName); Package* package = fVolume->fState->FindPackage(packageName);
if (package != NULL) { if (package != NULL) {
if (fPackagesAlreadyAdded.find(package) if (fPackagesAlreadyAdded.find(package)
!= fPackagesAlreadyAdded.end()) { != fPackagesAlreadyAdded.end()) {
@@ -1341,14 +1481,12 @@ Volume::Volume(BLooper* looper)
fPackagesDirectoryRef(), fPackagesDirectoryRef(),
fRoot(NULL), fRoot(NULL),
fListener(NULL), fListener(NULL),
fPackagesByFileName(), fState(NULL),
fPackagesByNodeRef(),
fPendingNodeMonitorEventsLock("pending node monitor events"), fPendingNodeMonitorEventsLock("pending node monitor events"),
fPendingNodeMonitorEvents(), fPendingNodeMonitorEvents(),
fNodeMonitorEventHandleTime(0), fNodeMonitorEventHandleTime(0),
fPackagesToBeActivated(), fPackagesToBeActivated(),
fPackagesToBeDeactivated(), fPackagesToBeDeactivated(),
fChangeCount(0),
fLocationInfoReply(B_MESSAGE_GET_INSTALLATION_LOCATION_INFO_REPLY) fLocationInfoReply(B_MESSAGE_GET_INSTALLATION_LOCATION_INFO_REPLY)
{ {
looper->AddHandler(this); looper->AddHandler(this);
@@ -1360,21 +1498,15 @@ Volume::~Volume()
Unmounted(); Unmounted();
// need for error case in InitPackages() // need for error case in InitPackages()
fPackagesByFileName.Clear(); delete fState;
Package* package = fPackagesByNodeRef.Clear(true);
while (package != NULL) {
Package* next = package->NodeRefHashTableNext();
delete package;
package = next;
}
} }
status_t status_t
Volume::Init(const node_ref& rootDirectoryRef, node_ref& _packageRootRef) Volume::Init(const node_ref& rootDirectoryRef, node_ref& _packageRootRef)
{ {
if (fPackagesByFileName.Init() != B_OK || fPackagesByNodeRef.Init() != B_OK) fState = new(std::nothrow) State;
if (fState == NULL || !fState->Init())
RETURN_ERROR(B_NO_MEMORY); RETURN_ERROR(B_NO_MEMORY);
fRootDirectoryRef = rootDirectoryRef; fRootDirectoryRef = rootDirectoryRef;
@@ -1481,8 +1613,8 @@ Volume::InitPackages(Listener* listener)
status_t status_t
Volume::AddPackagesToRepository(BSolverRepository& repository, bool activeOnly) Volume::AddPackagesToRepository(BSolverRepository& repository, bool activeOnly)
{ {
for (PackageFileNameHashTable::Iterator it for (PackageFileNameHashTable::Iterator it = fState->ByFileNameIterator();
= fPackagesByFileName.GetIterator(); it.HasNext();) { it.HasNext();) {
Package* package = it.Next(); Package* package = it.Next();
if (activeOnly && !package->IsActive()) if (activeOnly && !package->IsActive())
continue; continue;
@@ -1589,10 +1721,13 @@ INFORM("Volume::InitialVerify(%p, %p)\n", nextVolume, nextNextVolume);
void void
Volume::HandleGetLocationInfoRequest(BMessage* message) Volume::HandleGetLocationInfoRequest(BMessage* message)
{ {
AutoLocker<State> stateLocker(fState);
// If the cached reply message is up-to-date, just send it. // If the cached reply message is up-to-date, just send it.
int64 changeCount; int64 changeCount;
if (fLocationInfoReply.FindInt64("change count", &changeCount) == B_OK if (fLocationInfoReply.FindInt64("change count", &changeCount) == B_OK
&& changeCount == fChangeCount) { && changeCount == fState->ChangeCount()) {
stateLocker.Unlock();
message->SendReply(&fLocationInfoReply, (BHandler*)NULL, message->SendReply(&fLocationInfoReply, (BHandler*)NULL,
kCommunicationTimeout); kCommunicationTimeout);
return; return;
@@ -1612,8 +1747,8 @@ Volume::HandleGetLocationInfoRequest(BMessage* message)
return; return;
} }
for (PackageFileNameHashTable::Iterator it for (PackageFileNameHashTable::Iterator it = fState->ByFileNameIterator();
= fPackagesByFileName.GetIterator(); it.HasNext();) { it.HasNext();) {
Package* package = it.Next(); Package* package = it.Next();
const char* fieldName = package->IsActive() const char* fieldName = package->IsActive()
? "active packages" : "inactive packages"; ? "active packages" : "inactive packages";
@@ -1625,8 +1760,12 @@ Volume::HandleGetLocationInfoRequest(BMessage* message)
} }
} }
if (fLocationInfoReply.AddInt64("change count", fChangeCount) != B_OK) if (fLocationInfoReply.AddInt64("change count", fState->ChangeCount())
!= B_OK) {
return; return;
}
stateLocker.Unlock();
message->SendReply(&fLocationInfoReply, (BHandler*)NULL, message->SendReply(&fLocationInfoReply, (BHandler*)NULL,
kCommunicationTimeout); kCommunicationTimeout);
@@ -1670,6 +1809,27 @@ Volume::HandleCommitTransactionRequest(BMessage* message)
} }
void
Volume::PackageJobPending()
{
fState->PackageJobPending();
}
void
Volume::PackageJobFinished()
{
fState->PackageJobFinished();
}
bool
Volume::IsPackageJobPending() const
{
return fState->IsPackageJobPending();
}
void void
Volume::Unmounted() Volume::Unmounted()
{ {
@@ -1738,6 +1898,13 @@ Volume::Location() const
} }
PackageFileNameHashTable::Iterator
Volume::PackagesByFileNameIterator() const
{
return fState->ByFileNameIterator();
}
int int
Volume::OpenRootDirectory() const Volume::OpenRootDirectory() const
{ {
@@ -1852,7 +2019,7 @@ Volume::CreateTransaction(BPackageInstallationLocation location,
} }
// init the transaction // init the transaction
error = _transaction.SetTo(location, fChangeCount, directoryName); error = _transaction.SetTo(location, fState->ChangeCount(), directoryName);
if (error != B_OK) { if (error != B_OK) {
BEntry entry; BEntry entry;
_transactionDirectory.GetEntry(&entry); _transactionDirectory.GetEntry(&entry);
@@ -1979,7 +2146,7 @@ Volume::_PackagesEntryCreated(const char* name)
{ {
INFORM("Volume::_PackagesEntryCreated(\"%s\")\n", name); INFORM("Volume::_PackagesEntryCreated(\"%s\")\n", name);
// Ignore the event, if the package is already known. // Ignore the event, if the package is already known.
Package* package = fPackagesByFileName.Lookup(name); Package* package = fState->FindPackage(name);
if (package != NULL) { if (package != NULL) {
if (package->EntryCreatedIgnoreLevel() > 0) { if (package->EntryCreatedIgnoreLevel() > 0) {
package->DecrementEntryCreatedIgnoreLevel(); package->DecrementEntryCreatedIgnoreLevel();
@@ -2029,7 +2196,7 @@ void
Volume::_PackagesEntryRemoved(const char* name) Volume::_PackagesEntryRemoved(const char* name)
{ {
INFORM("Volume::_PackagesEntryRemoved(\"%s\")\n", name); INFORM("Volume::_PackagesEntryRemoved(\"%s\")\n", name);
Package* package = fPackagesByFileName.Lookup(name); Package* package = fState->FindPackage(name);
if (package == NULL) if (package == NULL)
return; return;
@@ -2081,16 +2248,13 @@ Volume::_FillInActivationChangeItem(PackageFSActivationChangeItem* item,
void void
Volume::_AddPackage(Package* package) Volume::_AddPackage(Package* package)
{ {
fPackagesByFileName.Insert(package); fState->AddPackage(package);
fPackagesByNodeRef.Insert(package);
} }
void void
Volume::_RemovePackage(Package* package) Volume::_RemovePackage(Package* package)
{ {
fPackagesByFileName.Remove(package); fState->RemovePackage(package);
fPackagesByNodeRef.Remove(package);
fChangeCount++;
} }
@@ -2157,7 +2321,7 @@ Volume::_GetActivePackages(int fd)
// mark the returned packages active // mark the returned packages active
for (uint32 i = 0; i < request->packageCount; i++) { for (uint32 i = 0; i < request->packageCount; i++) {
Package* package = fPackagesByNodeRef.Lookup( Package* package = fState->FindPackage(
node_ref(request->infos[i].packageDeviceID, node_ref(request->infos[i].packageDeviceID,
request->infos[i].packageNodeID)); request->infos[i].packageNodeID));
if (package == NULL) { if (package == NULL) {
@@ -2169,18 +2333,17 @@ Volume::_GetActivePackages(int fd)
continue; continue;
} }
package->SetActive(true); fState->SetPackageActive(package, true);
INFORM("active package: \"%s\"\n", package->FileName().String()); INFORM("active package: \"%s\"\n", package->FileName().String());
} }
for (PackageNodeRefHashTable::Iterator it = fPackagesByNodeRef.GetIterator(); for (PackageNodeRefHashTable::Iterator it = fState->ByNodeRefIterator();
it.HasNext();) { it.HasNext();) {
Package* package = it.Next(); Package* package = it.Next();
if (!package->IsActive()) if (!package->IsActive())
INFORM("inactive package: \"%s\"\n", package->FileName().String()); INFORM("inactive package: \"%s\"\n", package->FileName().String());
} }
PackageNodeRefHashTable fPackagesByNodeRef;
// INFORM("%" B_PRIu32 " active packages:\n", request->packageCount); // INFORM("%" B_PRIu32 " active packages:\n", request->packageCount);
// for (uint32 i = 0; i < request->packageCount; i++) { // for (uint32 i = 0; i < request->packageCount; i++) {
// INFORM(" dev: %" B_PRIdDEV ", node: %" B_PRIdINO "\n", // INFORM(" dev: %" B_PRIdDEV ", node: %" B_PRIdINO "\n",
@@ -2280,8 +2443,8 @@ Volume::_CreateActivationFileContent(const PackageSet& toActivate,
const PackageSet& toDeactivate, BString& _content) const PackageSet& toDeactivate, BString& _content)
{ {
BString activationFileContent; BString activationFileContent;
for (PackageFileNameHashTable::Iterator it for (PackageFileNameHashTable::Iterator it = fState->ByFileNameIterator();
= fPackagesByFileName.GetIterator(); it.HasNext();) { it.HasNext();) {
Package* package = it.Next(); Package* package = it.Next();
if (package->IsActive() if (package->IsActive()
&& toDeactivate.find(package) == toDeactivate.end()) { && toDeactivate.find(package) == toDeactivate.end()) {
@@ -2446,16 +2609,5 @@ packagesToActivate.size(), packagesToDeactivate.size());
// Update our state, i.e. remove deactivated packages and mark activated // Update our state, i.e. remove deactivated packages and mark activated
// packages accordingly. // packages accordingly.
for (PackageSet::iterator it = packagesToActivate.begin(); fState->ActivationChanged(packagesToActivate, packagesToDeactivate);
it != packagesToActivate.end(); ++it) {
(*it)->SetActive(true);
fChangeCount++;
}
for (PackageSet::iterator it = packagesToDeactivate.begin();
it != packagesToDeactivate.end(); ++it) {
Package* package = *it;
_RemovePackage(package);
delete package;
}
} }
+25 -8
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2013, Haiku, Inc. All Rights Reserved. * Copyright 2013-2014, Haiku, Inc. All Rights Reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -23,6 +23,21 @@
#include "Package.h" #include "Package.h"
// Locking Policy
// ==============
//
// A Volume object is accessed by two threads:
// 1. The application thread: initially (c'tor and Init()) and when handling a
// location info request (HandleGetLocationInfoRequest()).
// 2. The corresponding Root object's job thread (any other operation).
//
// The only thread synchronization needed is for the status information accessed
// by HandleGetLocationInfoRequest() and modified by the job thread. The data
// are encapsulated in a Volume::State object which contains a lock. The lock
// must be held by the app thread when accessing the data (it reads only) and
// by the job thread when modifying the data (not needed when reading).
using BPackageKit::BPrivate::BActivationTransaction; using BPackageKit::BPrivate::BActivationTransaction;
using BPackageKit::BPrivate::BDaemonClient; using BPackageKit::BPrivate::BDaemonClient;
@@ -61,6 +76,10 @@ public:
void HandleCommitTransactionRequest( void HandleCommitTransactionRequest(
BMessage* message); BMessage* message);
void PackageJobPending();
void PackageJobFinished();
bool IsPackageJobPending() const;
void Unmounted(); void Unmounted();
virtual void MessageReceived(BMessage* message); virtual void MessageReceived(BMessage* message);
@@ -90,10 +109,8 @@ public:
void SetRoot(Root* root) void SetRoot(Root* root)
{ fRoot = root; } { fRoot = root; }
const PackageFileNameHashTable& PackagesByFileName() const PackageFileNameHashTable::Iterator PackagesByFileNameIterator()
{ return fPackagesByFileName; } const;
const PackageNodeRefHashTable& PackagesByNodeRef() const
{ return fPackagesByNodeRef; }
int OpenRootDirectory() const; int OpenRootDirectory() const;
@@ -120,6 +137,7 @@ public:
private: private:
struct NodeMonitorEvent; struct NodeMonitorEvent;
struct State;
struct CommitTransactionHandler; struct CommitTransactionHandler;
friend struct CommitTransactionHandler; friend struct CommitTransactionHandler;
@@ -196,15 +214,14 @@ private:
node_ref fPackagesDirectoryRef; node_ref fPackagesDirectoryRef;
Root* fRoot; Root* fRoot;
Listener* fListener; Listener* fListener;
PackageFileNameHashTable fPackagesByFileName; State* fState;
PackageNodeRefHashTable fPackagesByNodeRef;
BLocker fPendingNodeMonitorEventsLock; BLocker fPendingNodeMonitorEventsLock;
NodeMonitorEventList fPendingNodeMonitorEvents; NodeMonitorEventList fPendingNodeMonitorEvents;
bigtime_t fNodeMonitorEventHandleTime; bigtime_t fNodeMonitorEventHandleTime;
PackageSet fPackagesToBeActivated; PackageSet fPackagesToBeActivated;
PackageSet fPackagesToBeDeactivated; PackageSet fPackagesToBeDeactivated;
int64 fChangeCount;
BMessage fLocationInfoReply; BMessage fLocationInfoReply;
// only accessed in the application thread
}; };