summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/document_source_merge_cursors.cpp
blob: 1d4ed9347b3ec5fd0530914f57f470a8a132df25 (plain)
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
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
/**
 * Copyright 2013 (c) 10gen Inc.
 *
 * This program is free software: you can redistribute it and/or  modify
 * it under the terms of the GNU Affero General Public License, version 3,
 * as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful,
 * but WITHOUT ANY WARRANTY; without even the implied warranty of
 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
 * GNU Affero General Public License for more details.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program.  If not, see <http://www.gnu.org/licenses/>.
 */

#include "mongo/pch.h"

#include "mongo/db/pipeline/document_source.h"

#include "mongo/db/cmdline.h"

namespace mongo {

    const char DocumentSourceMergeCursors::name[] = "$mergeCursors";

    const char* DocumentSourceMergeCursors::getSourceName() const {
        return name;
    }

    void DocumentSourceMergeCursors::setSource(DocumentSource *pSource) {
        /* this doesn't take a source */
        verify(false);
    }

    DocumentSourceMergeCursors::DocumentSourceMergeCursors(
            const CursorIds& cursorIds,
            const intrusive_ptr<ExpressionContext> &pExpCtx)
        : DocumentSource(pExpCtx)
        , _cursorIds(cursorIds)
        , _unstarted(true)
    {}

    intrusive_ptr<DocumentSource> DocumentSourceMergeCursors::create(
            const CursorIds& cursorIds,
            const intrusive_ptr<ExpressionContext> &pExpCtx) {
        return new DocumentSourceMergeCursors(cursorIds, pExpCtx);
    }

    intrusive_ptr<DocumentSource> DocumentSourceMergeCursors::createFromBson(
            BSONElement* pBsonElement,
            const intrusive_ptr<ExpressionContext>& pExpCtx) {

        massert(17026, string("Expected an Array, but got a ") + typeName(pBsonElement->type()),
                pBsonElement->type() == Array);

        CursorIds cursorIds;
        BSONObj array = pBsonElement->embeddedObject();
        BSONForEach(cursor, array) {
            massert(17027, string("Expected an Object, but got a ") + typeName(cursor.type()),
                    cursor.type() == Object);

            cursorIds.push_back(make_pair(ConnectionString(cursor["host"].String()),
                                          cursor["id"].Long()));
        }
        
        return new DocumentSourceMergeCursors(cursorIds, pExpCtx);
    }

    Value DocumentSourceMergeCursors::serialize(bool explain) const {
        vector<Value> cursors;
        for (size_t i = 0; i < _cursorIds.size(); i++) {
            cursors.push_back(Value(DOC("host" << Value(_cursorIds[i].first.toString())
                                     << "id" << _cursorIds[i].second)));
        }
        return Value(DOC(getSourceName() << Value(cursors)));
    }

    DocumentSourceMergeCursors::CursorAndConnection::CursorAndConnection(
            ConnectionString host,
            NamespaceString ns,
            CursorId id)
        : connection(host)
        , cursor(connection.get(), ns, id, 0, 0)
    {}

    boost::optional<Document> DocumentSourceMergeCursors::getNext() {
        if (_unstarted) {
            _unstarted = false;

            // open each cursor and send message asking for a batch
            for (CursorIds::const_iterator it = _cursorIds.begin(); it !=_cursorIds.end(); ++it) {
                _cursors.push_back(boost::make_shared<CursorAndConnection>(
                            it->first, pExpCtx->ns, it->second));
                verify(_cursors.back()->connection->lazySupported());
                _cursors.back()->cursor.initLazy(); // shouldn't block
            }

            // wait for all cursors to return a batch
            // TODO need a way to keep cursors alive if some take longer than 10 minutes.
            for (Cursors::const_iterator it = _cursors.begin(); it !=_cursors.end(); ++it) {
                bool retry = false;
                bool ok = (*it)->cursor.initLazyFinish(retry); // blocks here for first batch

                uassert(17028,
                        "error reading response from " + _cursors.back()->connection->toString(),
                        ok);
                verify(!retry);
            }

            _currentCursor = _cursors.begin();
        }

        // purge eof cursors and release their connections
        while (!_cursors.empty() && !(*_currentCursor)->cursor.more()) {
            (*_currentCursor)->connection.done();
            _cursors.erase(_currentCursor);
            _currentCursor = _cursors.begin();
        }

        if (_cursors.empty())
            return boost::none;

        BSONObj next = (*_currentCursor)->cursor.next();
        uassert(17029, str::stream() << "Received error in response from "
                                     << (*_currentCursor)->connection->toString()
                                     << ": " << next,
                !next.hasField("$err"));

        // advance _currentCursor, wrapping if needed
        if (++_currentCursor == _cursors.end())
            _currentCursor = _cursors.begin();

        return Document(next);
    }

    void DocumentSourceMergeCursors::dispose() {
        _cursors.clear();
        _currentCursor = _cursors.end();
    }
}