* Added a timeout to the delivery functions. This is the time after which
the message will be finally dropped. Makes sense for periodic message runners for instance. * Set the target of a BMessage before flattening it. Thus there will be space in the flattened header for it. git-svn-id: file:///srv/svn/repos/haiku/trunk/current@11071 a95241bf-73f2-0310-859d-f6bbb57e9c96
This commit is contained in:
@@ -26,13 +26,19 @@ static const bigtime_t kRetryDelay = 20000; // 20 ms
|
|||||||
// Message
|
// Message
|
||||||
class MessageDeliverer::Message : public Referenceable {
|
class MessageDeliverer::Message : public Referenceable {
|
||||||
public:
|
public:
|
||||||
Message(void *data, int32 dataSize)
|
Message(void *data, int32 dataSize, bigtime_t timeout)
|
||||||
: Referenceable(true),
|
: Referenceable(true),
|
||||||
fData(data),
|
fData(data),
|
||||||
fDataSize(dataSize),
|
fDataSize(dataSize),
|
||||||
fCreationTime(system_time()),
|
fCreationTime(system_time()),
|
||||||
fBusy(false)
|
fBusy(false)
|
||||||
{
|
{
|
||||||
|
if (B_INFINITE_TIMEOUT - fCreationTime <= timeout)
|
||||||
|
fTimeoutTime = B_INFINITE_TIMEOUT;
|
||||||
|
else if (timeout <= 0)
|
||||||
|
fTimeoutTime = fCreationTime;
|
||||||
|
else
|
||||||
|
fTimeoutTime = fCreationTime + timeout;
|
||||||
}
|
}
|
||||||
|
|
||||||
~Message()
|
~Message()
|
||||||
@@ -50,11 +56,16 @@ public:
|
|||||||
return fDataSize;
|
return fDataSize;
|
||||||
}
|
}
|
||||||
|
|
||||||
bigtime_t CreationTime()
|
bigtime_t CreationTime() const
|
||||||
{
|
{
|
||||||
return fCreationTime;
|
return fCreationTime;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
bigtime_t TimeoutTime() const
|
||||||
|
{
|
||||||
|
return fTimeoutTime;
|
||||||
|
}
|
||||||
|
|
||||||
void SetBusy(bool busy)
|
void SetBusy(bool busy)
|
||||||
{
|
{
|
||||||
fBusy = busy;
|
fBusy = busy;
|
||||||
@@ -69,6 +80,7 @@ private:
|
|||||||
void *fData;
|
void *fData;
|
||||||
int32 fDataSize;
|
int32 fDataSize;
|
||||||
bigtime_t fCreationTime;
|
bigtime_t fCreationTime;
|
||||||
|
bigtime_t fTimeoutTime;
|
||||||
bool fBusy;
|
bool fBusy;
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -274,21 +286,30 @@ MessageDeliverer::Default()
|
|||||||
|
|
||||||
// DeliverMessage
|
// DeliverMessage
|
||||||
status_t
|
status_t
|
||||||
MessageDeliverer::DeliverMessage(BMessage *message, BMessenger target)
|
MessageDeliverer::DeliverMessage(BMessage *message, BMessenger target,
|
||||||
|
bigtime_t timeout)
|
||||||
{
|
{
|
||||||
BMessenger::Private messengerPrivate(target);
|
BMessenger::Private messengerPrivate(target);
|
||||||
return DeliverMessage(message, messengerPrivate.Port(),
|
return DeliverMessage(message, messengerPrivate.Port(),
|
||||||
messengerPrivate.IsPreferredTarget()
|
messengerPrivate.IsPreferredTarget()
|
||||||
? B_PREFERRED_TOKEN : messengerPrivate.Token());
|
? B_PREFERRED_TOKEN : messengerPrivate.Token(),
|
||||||
|
timeout);
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeliverMessage
|
// DeliverMessage
|
||||||
status_t
|
status_t
|
||||||
MessageDeliverer::DeliverMessage(BMessage *message, port_id port, int32 token)
|
MessageDeliverer::DeliverMessage(BMessage *message, port_id port, int32 token,
|
||||||
|
bigtime_t timeout)
|
||||||
{
|
{
|
||||||
if (!message)
|
if (!message)
|
||||||
return B_BAD_VALUE;
|
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
|
// flatten the message
|
||||||
BMallocIO mallocIO;
|
BMallocIO mallocIO;
|
||||||
status_t error = message->Flatten(&mallocIO);
|
status_t error = message->Flatten(&mallocIO);
|
||||||
@@ -296,35 +317,36 @@ MessageDeliverer::DeliverMessage(BMessage *message, port_id port, int32 token)
|
|||||||
return error;
|
return error;
|
||||||
|
|
||||||
return DeliverMessage(mallocIO.Buffer(), mallocIO.BufferLength(), port,
|
return DeliverMessage(mallocIO.Buffer(), mallocIO.BufferLength(), port,
|
||||||
token);
|
token, timeout);
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeliverMessage
|
// DeliverMessage
|
||||||
status_t
|
status_t
|
||||||
MessageDeliverer::DeliverMessage(const void *message, int32 messageSize,
|
MessageDeliverer::DeliverMessage(const void *message, int32 messageSize,
|
||||||
BMessenger target)
|
BMessenger target, bigtime_t timeout)
|
||||||
{
|
{
|
||||||
BMessenger::Private messengerPrivate(target);
|
BMessenger::Private messengerPrivate(target);
|
||||||
return DeliverMessage(message, messageSize, messengerPrivate.Port(),
|
return DeliverMessage(message, messageSize, messengerPrivate.Port(),
|
||||||
messengerPrivate.IsPreferredTarget()
|
messengerPrivate.IsPreferredTarget()
|
||||||
? B_PREFERRED_TOKEN : messengerPrivate.Token());
|
? B_PREFERRED_TOKEN : messengerPrivate.Token(),
|
||||||
|
timeout);
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeliverMessage
|
// DeliverMessage
|
||||||
status_t
|
status_t
|
||||||
MessageDeliverer::DeliverMessage(const void *message, int32 messageSize,
|
MessageDeliverer::DeliverMessage(const void *message, int32 messageSize,
|
||||||
port_id port, int32 token)
|
port_id port, int32 token, bigtime_t timeout)
|
||||||
{
|
{
|
||||||
messaging_target target;
|
messaging_target target;
|
||||||
target.port = port;
|
target.port = port;
|
||||||
target.token = token;
|
target.token = token;
|
||||||
return DeliverMessage(message, messageSize, &target, 1);
|
return DeliverMessage(message, messageSize, &target, 1, timeout);
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeliverMessage
|
// DeliverMessage
|
||||||
status_t
|
status_t
|
||||||
MessageDeliverer::DeliverMessage(const void *messageData, int32 messageSize,
|
MessageDeliverer::DeliverMessage(const void *messageData, int32 messageSize,
|
||||||
const messaging_target *targets, int32 targetCount)
|
const messaging_target *targets, int32 targetCount, bigtime_t timeout)
|
||||||
{
|
{
|
||||||
if (!messageData || messageSize <= 0)
|
if (!messageData || messageSize <= 0)
|
||||||
return B_BAD_VALUE;
|
return B_BAD_VALUE;
|
||||||
@@ -336,7 +358,7 @@ MessageDeliverer::DeliverMessage(const void *messageData, int32 messageSize,
|
|||||||
memcpy(data, messageData, messageSize);
|
memcpy(data, messageData, messageSize);
|
||||||
|
|
||||||
// create a Message
|
// create a Message
|
||||||
Message *message = new(nothrow) Message(data, messageSize);
|
Message *message = new(nothrow) Message(data, messageSize, timeout);
|
||||||
if (!message) {
|
if (!message) {
|
||||||
free(data);
|
free(data);
|
||||||
return B_NO_MEMORY;
|
return B_NO_MEMORY;
|
||||||
@@ -448,7 +470,13 @@ MessageDeliverer::_DelivererThread()
|
|||||||
// try sending all messages
|
// try sending all messages
|
||||||
int32 token;
|
int32 token;
|
||||||
while (Message *message = port->PeekMessage(token)) {
|
while (Message *message = port->PeekMessage(token)) {
|
||||||
status_t error = _SendMessage(message, port->PortID(), token);
|
status_t error = B_OK;
|
||||||
|
if (message->TimeoutTime() > system_time()) {
|
||||||
|
error = _SendMessage(message, port->PortID(), token);
|
||||||
|
} else {
|
||||||
|
// timeout, drop message
|
||||||
|
}
|
||||||
|
|
||||||
if (error == B_OK) {
|
if (error == B_OK) {
|
||||||
port->PopMessage();
|
port->PopMessage();
|
||||||
} else if (error == B_WOULD_BLOCK) {
|
} else if (error == B_WOULD_BLOCK) {
|
||||||
|
|||||||
@@ -23,14 +23,17 @@ public:
|
|||||||
static void DeleteDefault();
|
static void DeleteDefault();
|
||||||
static MessageDeliverer *Default();
|
static MessageDeliverer *Default();
|
||||||
|
|
||||||
status_t DeliverMessage(BMessage *message, BMessenger target);
|
status_t DeliverMessage(BMessage *message, BMessenger target,
|
||||||
status_t DeliverMessage(BMessage *message, port_id port, int32 token);
|
bigtime_t timeout = B_INFINITE_TIMEOUT);
|
||||||
|
status_t DeliverMessage(BMessage *message, port_id port, int32 token,
|
||||||
|
bigtime_t timeout = B_INFINITE_TIMEOUT);
|
||||||
status_t DeliverMessage(const void *message, int32 messageSize,
|
status_t DeliverMessage(const void *message, int32 messageSize,
|
||||||
BMessenger target);
|
BMessenger target, bigtime_t timeout = B_INFINITE_TIMEOUT);
|
||||||
status_t DeliverMessage(const void *message, int32 messageSize,
|
status_t DeliverMessage(const void *message, int32 messageSize,
|
||||||
port_id port, int32 token);
|
port_id port, int32 token, bigtime_t timeout = B_INFINITE_TIMEOUT);
|
||||||
status_t DeliverMessage(const void *message, int32 messageSize,
|
status_t DeliverMessage(const void *message, int32 messageSize,
|
||||||
const messaging_target *targets, int32 targetCount);
|
const messaging_target *targets, int32 targetCount,
|
||||||
|
bigtime_t timeout = B_INFINITE_TIMEOUT);
|
||||||
|
|
||||||
private:
|
private:
|
||||||
class Message;
|
class Message;
|
||||||
|
|||||||
Reference in New Issue
Block a user