aboutsummaryrefslogtreecommitdiffstats
path: root/messagebus/src/vespa/messagebus/sequencer.cpp
blob: 59e3164e11d6cb28b1bb3bba617c8db928a3f4da (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
// Copyright Yahoo. Licensed under the terms of the Apache 2.0 license. See LICENSE in the project root.
#include "sequencer.h"
#include "tracelevel.h"
#include <vespa/vespalib/util/stringfmt.h>
#include <cassert>

using vespalib::make_string;

namespace mbus {

Sequencer::Sequencer(IMessageHandler &sender) :
    _lock(),
    _sender(sender),
    _seqMap()
{
    // empty
}

Sequencer::~Sequencer()
{
    for (auto & entry : _seqMap) {
        MessageQueue *queue = entry.second;
        if (queue != nullptr) {
            while (queue->size() > 0) {
                Message *msg = queue->front();
                queue->pop();
                msg->discard();
                delete msg;
            }
            delete queue;
        }
    }
}

Message::UP
Sequencer::filter(Message::UP msg)
{
    uint64_t seqId = msg->getSequenceId();
    msg->setContext(Context(seqId));
    {
        std::lock_guard guard(_lock);
        auto it = _seqMap.find(seqId);
        if (it != _seqMap.end()) {
            if (it->second == nullptr) {
                it->second = new MessageQueue();
            }
            msg->getTrace().trace(TraceLevel::COMPONENT,
                                  make_string("Sequencer queued message with sequence id '%" PRIu64 "'.", seqId));
            it->second->push(msg.get());
            msg.release();
            return {};
        }
        _seqMap[seqId] = nullptr; // insert empty queue
    }
    return msg;
}

void
Sequencer::sequencedSend(Message::UP msg)
{
    msg->getTrace().trace(TraceLevel::COMPONENT,
                          make_string("Sequencer sending message with sequence id '%" PRIu64 "'.",
                                      msg->getContext().value.UINT64));
    msg->pushHandler(*this);
    _sender.handleMessage(std::move(msg));
}

void
Sequencer::handleMessage(Message::UP msg)
{
    if (msg->hasSequenceId()) {
        msg = filter(std::move(msg));
        if (msg.get() != nullptr) {
            sequencedSend(std::move(msg));
        }
    } else {
        _sender.handleMessage(std::move(msg)); // unsequenced
    }
}

void
Sequencer::handleReply(Reply::UP reply)
{
    uint64_t seq = reply->getContext().value.UINT64;
    reply->getTrace().trace(TraceLevel::COMPONENT,
                            make_string("Sequencer received reply with sequence id '%" PRIu64 "'.", seq));
    Message::UP msg;
    {
        std::lock_guard guard(_lock);
        auto it = _seqMap.find(seq);
        MessageQueue *que = it->second;
        assert(it != _seqMap.end());
        if (que == nullptr || que->size() == 0) {
            if (que != nullptr) {
                delete que;
            }
            _seqMap.erase(it);
        } else {
            msg.reset(que->front());
            que->pop();
        }
    }
    if (msg) {
        sequencedSend(std::move(msg));
    }
    IReplyHandler &handler = reply->getCallStack().pop(*reply);
    handler.handleReply(std::move(reply));
}

} // namespace mbus