/* * Copyright 2005, Ingo Weinhold, bonefish@users.sf.net. All rights reserved. * Distributed under the terms of the MIT License. */ #include #include #include #include #include #include #include #include #include #include #include #include #include "Debug.h" #include "MessageDeliverer.h" #include "Referenceable.h" // sDeliverer -- the singleton instance MessageDeliverer *MessageDeliverer::sDeliverer = NULL; static const bigtime_t kRetryDelay = 100000; // 100 ms // per port sanity limits static const int32 kMaxMessagesPerPort = 10000; static const int32 kMaxDataPerPort = 50 * 1024 * 1024; // 50 MB // Message class MessageDeliverer::Message : public Referenceable { public: Message(void *data, int32 dataSize, bigtime_t timeout) : Referenceable(true), fData(data), fDataSize(dataSize), fCreationTime(system_time()), fBusy(false) { if (B_INFINITE_TIMEOUT - fCreationTime <= timeout) fTimeoutTime = B_INFINITE_TIMEOUT; else if (timeout <= 0) fTimeoutTime = fCreationTime; else fTimeoutTime = fCreationTime + timeout; } ~Message() { free(fData); } void *Data() const { return fData; } int32 DataSize() const { return fDataSize; } bigtime_t CreationTime() const { return fCreationTime; } bigtime_t TimeoutTime() const { return fTimeoutTime; } bool HasTimeout() const { return (fTimeoutTime < B_INFINITE_TIMEOUT); } void SetBusy(bool busy) { fBusy = busy; } bool IsBusy() const { return fBusy; } private: void *fData; int32 fDataSize; bigtime_t fCreationTime; bigtime_t fTimeoutTime; bool fBusy; }; // TargetMessage class MessageDeliverer::TargetMessage : public DoublyLinkedListLinkImpl { public: TargetMessage(Message *message, int32 token) : fMessage(message), fToken(token) { if (fMessage) fMessage->AddReference(); } ~TargetMessage() { if (fMessage) fMessage->RemoveReference(); } Message *GetMessage() const { return fMessage; } int32 Token() const { return fToken; } private: Message *fMessage; int32 fToken; }; // TargetMessageHandle class MessageDeliverer::TargetMessageHandle { public: TargetMessageHandle(TargetMessage *message) : fMessage(message) { } TargetMessageHandle(const TargetMessageHandle &other) : fMessage(other.fMessage) { } TargetMessage *GetMessage() const { return fMessage; } TargetMessageHandle &operator=(const TargetMessageHandle &other) { fMessage = other.fMessage; return *this; } bool operator==(const TargetMessageHandle &other) const { return (fMessage == other.fMessage); } bool operator!=(const TargetMessageHandle &other) const { return (fMessage != other.fMessage); } bool operator<(const TargetMessageHandle &other) const { bigtime_t timeout = fMessage->GetMessage()->TimeoutTime(); bigtime_t otherTimeout = other.fMessage->GetMessage()->TimeoutTime(); if (timeout < otherTimeout) return true; if (timeout > otherTimeout) return false; return (fMessage < other.fMessage); } private: TargetMessage *fMessage; }; // TargetPort class MessageDeliverer::TargetPort { public: TargetPort(port_id portID) : fPortID(portID), fMessages(), fMessageCount(0), fMessageSize(0) { } ~TargetPort() { while (!fMessages.IsEmpty()) PopMessage(); } port_id PortID() const { return fPortID; } status_t PushMessage(Message *message, int32 token) { PRINT(("MessageDeliverer::TargetPort::PushMessage(port: %ld, %p, %ld)\n", fPortID, message, token)); // create a target message TargetMessage *targetMessage = new(nothrow) TargetMessage(message, token); if (!targetMessage) return B_NO_MEMORY; // push it fMessages.Insert(targetMessage); fMessageCount++; fMessageSize += targetMessage->GetMessage()->DataSize(); // add it to the timeoutable messages, if it has a timeout if (message->HasTimeout()) fTimeoutableMessages.insert(targetMessage); _EnforceLimits(); return B_OK; } Message *PeekMessage(int32 &token) const { if (!fMessages.Head()) return NULL; token = fMessages.Head()->Token(); return fMessages.Head()->GetMessage(); } void PopMessage() { if (fMessages.Head()) { PRINT(("MessageDeliverer::TargetPort::PopMessage(): port: %ld, %p\n", fPortID, fMessages.Head()->GetMessage())); _RemoveMessage(fMessages.Head()); } } void DropTimedOutMessages() { bigtime_t now = system_time(); while (fTimeoutableMessages.begin() != fTimeoutableMessages.end()) { TargetMessage *message = fTimeoutableMessages.begin()->GetMessage(); if (message->GetMessage()->TimeoutTime() > now) break; PRINT(("MessageDeliverer::TargetPort::DropTimedOutMessages(): port: %ld: " "message %p timed out\n", fPortID, message->GetMessage())); _RemoveMessage(message); } } bool IsEmpty() const { return fMessages.IsEmpty(); } private: void _RemoveMessage(TargetMessage *message) { fMessages.Remove(message); fMessageCount--; fMessageSize -= message->GetMessage()->DataSize(); if (message->GetMessage()->HasTimeout()) fTimeoutableMessages.erase(message); delete message; } void _EnforceLimits() { // message count while (fMessageCount > kMaxMessagesPerPort) { PRINT(("MessageDeliverer::TargetPort::_EnforceLimits(): port: %ld: hit maximum " "message count limit.\n", fPortID)); PopMessage(); } // message size while (fMessageSize > kMaxDataPerPort) { PRINT(("MessageDeliverer::TargetPort::_EnforceLimits(): port: %ld: hit maximum " "message size limit.\n", fPortID)); PopMessage(); } } typedef DoublyLinkedList MessageList; port_id fPortID; MessageList fMessages; int32 fMessageCount; int32 fMessageSize; set fTimeoutableMessages; }; // TargetPortMap struct MessageDeliverer::TargetPortMap : public map { }; // #pragma mark - // constructor MessageDeliverer::MessageDeliverer() : fLock("message deliverer"), fTargetPorts(NULL), fDelivererThread(-1), fTerminating(false) { } // destructor MessageDeliverer::~MessageDeliverer() { fTerminating = true; if (fDelivererThread >= 0) { int32 result; wait_for_thread(fDelivererThread, &result); } delete fTargetPorts; } // Init status_t MessageDeliverer::Init() { // create the target port map fTargetPorts = new(nothrow) TargetPortMap; if (!fTargetPorts) return B_NO_MEMORY; // spawn the deliverer thread fDelivererThread = spawn_thread(MessageDeliverer::_DelivererThreadEntry, "message deliverer", B_NORMAL_PRIORITY, this); if (fDelivererThread < 0) return fDelivererThread; // resume the deliverer thread resume_thread(fDelivererThread); return B_OK; } // CreateDefault status_t MessageDeliverer::CreateDefault() { if (sDeliverer) return B_OK; // create the deliverer MessageDeliverer *deliverer = new(nothrow) MessageDeliverer; if (!deliverer) return B_NO_MEMORY; // init it status_t error = deliverer->Init(); if (error != B_OK) { delete deliverer; return error; } sDeliverer = deliverer; return B_OK; } // DeleteDefault void MessageDeliverer::DeleteDefault() { if (sDeliverer) { delete sDeliverer; sDeliverer = NULL; } } // Default MessageDeliverer * MessageDeliverer::Default() { return sDeliverer; } // DeliverMessage status_t MessageDeliverer::DeliverMessage(BMessage *message, BMessenger target, bigtime_t timeout) { BMessenger::Private messengerPrivate(target); return DeliverMessage(message, messengerPrivate.Port(), messengerPrivate.IsPreferredTarget() ? B_PREFERRED_TOKEN : messengerPrivate.Token(), timeout); } // DeliverMessage status_t MessageDeliverer::DeliverMessage(BMessage *message, port_id port, int32 token, bigtime_t timeout) { if (!message) return B_BAD_VALUE; // Set the token now, so that the header contains room for it. // It will be set when sending the message anyway, but if it is not set // before flattening, the header will not contain room for it, and it // will not possible to send the message flattened later. BMessage::Private(message).SetTarget(token, (token < 0)); // flatten the message BMallocIO mallocIO; status_t error = message->Flatten(&mallocIO); if (error != B_OK) return error; return DeliverMessage(mallocIO.Buffer(), mallocIO.BufferLength(), port, token, timeout); } // DeliverMessage status_t MessageDeliverer::DeliverMessage(const void *message, int32 messageSize, BMessenger target, bigtime_t timeout) { BMessenger::Private messengerPrivate(target); return DeliverMessage(message, messageSize, messengerPrivate.Port(), messengerPrivate.IsPreferredTarget() ? B_PREFERRED_TOKEN : messengerPrivate.Token(), timeout); } // DeliverMessage status_t MessageDeliverer::DeliverMessage(const void *message, int32 messageSize, port_id port, int32 token, bigtime_t timeout) { messaging_target target; target.port = port; target.token = token; return DeliverMessage(message, messageSize, &target, 1, timeout); } // DeliverMessage status_t MessageDeliverer::DeliverMessage(BMessage *message, const BMessenger *targets, int32 targetCount, bigtime_t timeout) { if (!message || targetCount < 0 || !targets) return B_BAD_VALUE; // convert the reply targets messaging_target *messagingTargets = new(nothrow) messaging_target[targetCount]; if (!messagingTargets) return B_NO_MEMORY; ArrayDeleter _(messagingTargets); for (int i = 0; i < targetCount; i++) { BMessenger messenger(targets[i]); BMessenger::Private messengerPrivate(messenger); messaging_target &target = messagingTargets[i]; target.port = messengerPrivate.Port(); target.token = messengerPrivate.IsPreferredTarget() ? B_PREFERRED_TOKEN : messengerPrivate.Token(); } // Set a dummy token now, so that the header contains room for it. // It will be set when sending the message anyway, but if it is not set // before flattening, the header will not contain room for it, and it // will not possible to send the message flattened later. BMessage::Private(message).SetTarget(0, false); // flatten the message BMallocIO mallocIO; status_t error = message->Flatten(&mallocIO); if (error != B_OK) return error; return DeliverMessage(mallocIO.Buffer(), mallocIO.BufferLength(), messagingTargets, targetCount, timeout); } // DeliverMessage status_t MessageDeliverer::DeliverMessage(const void *messageData, int32 messageSize, const messaging_target *targets, int32 targetCount, bigtime_t timeout) { if (!messageData || messageSize <= 0) return B_BAD_VALUE; // clone the buffer void *data = malloc(messageSize); if (!data) return B_NO_MEMORY; memcpy(data, messageData, messageSize); // create a Message Message *message = new(nothrow) Message(data, messageSize, timeout); if (!message) { free(data); return B_NO_MEMORY; } Reference _(message, true); // add the message to the respective target ports BAutolock locker(fLock); for (int32 i = 0; i < targetCount; i++) { // get the target port TargetPort *port = _GetTargetPort(targets[i].port, true); if (!port) return B_NO_MEMORY; // try sending the message, if there are no queued messages yet if (port->IsEmpty()) { status_t error = _SendMessage(message, targets[i].port, targets[i].token); // if the message was delivered OK, we're done with the target if (error == B_OK) { _PutTargetPort(port); continue; } // if the port is not full, but an error occurred, we skip this target if (error != B_WOULD_BLOCK) { _PutTargetPort(port); if (targetCount == 1) return error; continue; } } // add the message status_t error = port->PushMessage(message, targets[i].token); _PutTargetPort(port); if (error != B_OK) return error; } return B_OK; } // _GetTargetPort MessageDeliverer::TargetPort * MessageDeliverer::_GetTargetPort(port_id portID, bool create) { // get the port from the map TargetPortMap::iterator it = fTargetPorts->find(portID); if (it != fTargetPorts->end()) return it->second; if (!create) return NULL; // create a port TargetPort *port = new(nothrow) TargetPort(portID); if (!port) return NULL; (*fTargetPorts)[portID] = port; return port; } // _PutTargetPort void MessageDeliverer::_PutTargetPort(TargetPort *port) { if (!port) return; if (port->IsEmpty()) { fTargetPorts->erase(port->PortID()); delete port; } } // _SendMessage status_t MessageDeliverer::_SendMessage(Message *message, port_id portID, int32 token) { status_t error = BMessage::Private::SendFlattenedMessage(message->Data(), message->DataSize(), portID, token, (token < 0), 0); //PRINT(("MessageDeliverer::_SendMessage(%p, port: %ld, token: %ld): %lx\n", //message, portID, token, error)); return error; } // _DelivererThreadEntry int32 MessageDeliverer::_DelivererThreadEntry(void *data) { return ((MessageDeliverer*)data)->_DelivererThread(); } // _DelivererThread int32 MessageDeliverer::_DelivererThread() { while (!fTerminating) { snooze(kRetryDelay); if (fTerminating) break; // iterate through all target ports and try sending the messages BAutolock _(fLock); for (TargetPortMap::iterator it = fTargetPorts->begin(); it != fTargetPorts->end();) { TargetPort *port = it->second; bool portError = false; port->DropTimedOutMessages(); // try sending all messages int32 token; while (Message *message = port->PeekMessage(token)) { status_t error = B_OK; // if (message->TimeoutTime() > system_time()) { error = _SendMessage(message, port->PortID(), token); // } else { // // timeout, drop message // PRINT(("MessageDeliverer::_DelivererThread(): port %ld, " // "message %p timed out\n", port->PortID(), message)); // } if (error == B_OK) { port->PopMessage(); } else if (error == B_WOULD_BLOCK) { // no luck yet -- port is still full break; } else { // unexpected error -- probably the port is gone portError = true; } } // next port if (portError || port->IsEmpty()) { TargetPortMap::iterator oldIt = it; ++it; delete port; fTargetPorts->erase(oldIt); } else ++it; } } return 0; }