1 |
|
2 |
|
3 |
|
4 |
|
5 | 'use strict';
|
6 |
|
7 | const MongooseError = require('../error/mongooseError');
|
8 | const Readable = require('stream').Readable;
|
9 | const promiseOrCallback = require('../helpers/promiseOrCallback');
|
10 | const eachAsync = require('../helpers/cursor/eachAsync');
|
11 | const util = require('util');
|
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 | function AggregationCursor(agg) {
|
38 | Readable.call(this, { objectMode: true });
|
39 |
|
40 | this.cursor = null;
|
41 | this.agg = agg;
|
42 | this._transforms = [];
|
43 | const model = agg._model;
|
44 | delete agg.options.cursor.useMongooseAggCursor;
|
45 | this._mongooseOptions = {};
|
46 |
|
47 | _init(model, this, agg);
|
48 | }
|
49 |
|
50 | util.inherits(AggregationCursor, Readable);
|
51 |
|
52 |
|
53 |
|
54 |
|
55 |
|
56 | function _init(model, c, agg) {
|
57 | if (!model.collection.buffer) {
|
58 | model.hooks.execPre('aggregate', agg, function() {
|
59 | c.cursor = model.collection.aggregate(agg._pipeline, agg.options || {});
|
60 | c.emit('cursor', c.cursor);
|
61 | });
|
62 | } else {
|
63 | model.collection.emitter.once('queue', function() {
|
64 | model.hooks.execPre('aggregate', agg, function() {
|
65 | c.cursor = model.collection.aggregate(agg._pipeline, agg.options || {});
|
66 | c.emit('cursor', c.cursor);
|
67 | });
|
68 | });
|
69 | }
|
70 | }
|
71 |
|
72 |
|
73 |
|
74 |
|
75 |
|
76 | AggregationCursor.prototype._read = function() {
|
77 | const _this = this;
|
78 | _next(this, function(error, doc) {
|
79 | if (error) {
|
80 | return _this.emit('error', error);
|
81 | }
|
82 | if (!doc) {
|
83 | _this.push(null);
|
84 | _this.cursor.close(function(error) {
|
85 | if (error) {
|
86 | return _this.emit('error', error);
|
87 | }
|
88 | setTimeout(function() {
|
89 |
|
90 | const isNotClosedAutomatically = !_this.destroyed;
|
91 | if (isNotClosedAutomatically) {
|
92 | _this.emit('close');
|
93 | }
|
94 | }, 0);
|
95 | });
|
96 | return;
|
97 | }
|
98 | _this.push(doc);
|
99 | });
|
100 | };
|
101 |
|
102 | if (Symbol.asyncIterator != null) {
|
103 | const msg = 'Mongoose does not support using async iterators with an ' +
|
104 | 'existing aggregation cursor. See http://bit.ly/mongoose-async-iterate-aggregation';
|
105 |
|
106 | AggregationCursor.prototype[Symbol.asyncIterator] = function() {
|
107 | throw new MongooseError(msg);
|
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 |
|
142 |
|
143 |
|
144 | AggregationCursor.prototype.map = function(fn) {
|
145 | this._transforms.push(fn);
|
146 | return this;
|
147 | };
|
148 |
|
149 |
|
150 |
|
151 |
|
152 |
|
153 | AggregationCursor.prototype._markError = function(error) {
|
154 | this._error = error;
|
155 | return this;
|
156 | };
|
157 |
|
158 |
|
159 |
|
160 |
|
161 |
|
162 |
|
163 |
|
164 |
|
165 |
|
166 |
|
167 |
|
168 |
|
169 |
|
170 | AggregationCursor.prototype.close = function(callback) {
|
171 | return promiseOrCallback(callback, cb => {
|
172 | this.cursor.close(error => {
|
173 | if (error) {
|
174 | cb(error);
|
175 | return this.listeners('error').length > 0 && this.emit('error', error);
|
176 | }
|
177 | this.emit('close');
|
178 | cb(null);
|
179 | });
|
180 | });
|
181 | };
|
182 |
|
183 |
|
184 |
|
185 |
|
186 |
|
187 |
|
188 |
|
189 |
|
190 |
|
191 |
|
192 |
|
193 | AggregationCursor.prototype.next = function(callback) {
|
194 | return promiseOrCallback(callback, cb => {
|
195 | _next(this, cb);
|
196 | });
|
197 | };
|
198 |
|
199 |
|
200 |
|
201 |
|
202 |
|
203 |
|
204 |
|
205 |
|
206 |
|
207 |
|
208 |
|
209 |
|
210 |
|
211 |
|
212 |
|
213 | AggregationCursor.prototype.eachAsync = function(fn, opts, callback) {
|
214 | const _this = this;
|
215 | if (typeof opts === 'function') {
|
216 | callback = opts;
|
217 | opts = {};
|
218 | }
|
219 | opts = opts || {};
|
220 |
|
221 | return eachAsync(function(cb) { return _next(_this, cb); }, fn, opts, callback);
|
222 | };
|
223 |
|
224 |
|
225 |
|
226 |
|
227 |
|
228 | AggregationCursor.prototype.transformNull = function(val) {
|
229 | if (arguments.length === 0) {
|
230 | val = true;
|
231 | }
|
232 | this._mongooseOptions.transformNull = val;
|
233 | return this;
|
234 | };
|
235 |
|
236 |
|
237 |
|
238 |
|
239 |
|
240 |
|
241 |
|
242 |
|
243 |
|
244 |
|
245 |
|
246 |
|
247 | AggregationCursor.prototype.addCursorFlag = function(flag, value) {
|
248 | const _this = this;
|
249 | _waitForCursor(this, function() {
|
250 | _this.cursor.addCursorFlag(flag, value);
|
251 | });
|
252 | return this;
|
253 | };
|
254 |
|
255 |
|
256 |
|
257 |
|
258 |
|
259 | function _waitForCursor(ctx, cb) {
|
260 | if (ctx.cursor) {
|
261 | return cb();
|
262 | }
|
263 | ctx.once('cursor', function() {
|
264 | cb();
|
265 | });
|
266 | }
|
267 |
|
268 |
|
269 |
|
270 |
|
271 |
|
272 |
|
273 | function _next(ctx, cb) {
|
274 | let callback = cb;
|
275 | if (ctx._transforms.length) {
|
276 | callback = function(err, doc) {
|
277 | if (err || (doc === null && !ctx._mongooseOptions.transformNull)) {
|
278 | return cb(err, doc);
|
279 | }
|
280 | cb(err, ctx._transforms.reduce(function(doc, fn) {
|
281 | return fn(doc);
|
282 | }, doc));
|
283 | };
|
284 | }
|
285 |
|
286 | if (ctx._error) {
|
287 | return process.nextTick(function() {
|
288 | callback(ctx._error);
|
289 | });
|
290 | }
|
291 |
|
292 | if (ctx.cursor) {
|
293 | return ctx.cursor.next(function(error, doc) {
|
294 | if (error) {
|
295 | return callback(error);
|
296 | }
|
297 | if (!doc) {
|
298 | return callback(null, null);
|
299 | }
|
300 |
|
301 | callback(null, doc);
|
302 | });
|
303 | } else {
|
304 | ctx.once('cursor', function() {
|
305 | _next(ctx, cb);
|
306 | });
|
307 | }
|
308 | }
|
309 |
|
310 | module.exports = AggregationCursor;
|