1 | 'use strict'
|
2 |
|
3 |
|
4 |
|
5 |
|
6 | var xtend = require('xtend')
|
7 |
|
8 | var Readable = require('readable-stream').Readable
|
9 | var streamsOpts = { objectMode: true }
|
10 | var defaultStoreOptions = {
|
11 | clean: true
|
12 | }
|
13 |
|
14 |
|
15 |
|
16 |
|
17 |
|
18 |
|
19 |
|
20 |
|
21 |
|
22 |
|
23 |
|
24 | var Map = require('es6-map')
|
25 |
|
26 |
|
27 |
|
28 |
|
29 |
|
30 |
|
31 |
|
32 | function Store (options) {
|
33 | if (!(this instanceof Store)) {
|
34 | return new Store(options)
|
35 | }
|
36 |
|
37 | this.options = options || {}
|
38 |
|
39 |
|
40 | this.options = xtend(defaultStoreOptions, options)
|
41 |
|
42 | this._inflights = new Map()
|
43 | }
|
44 |
|
45 |
|
46 |
|
47 |
|
48 |
|
49 |
|
50 | Store.prototype.put = function (packet, cb) {
|
51 | this._inflights.set(packet.messageId, packet)
|
52 |
|
53 | if (cb) {
|
54 | cb()
|
55 | }
|
56 |
|
57 | return this
|
58 | }
|
59 |
|
60 |
|
61 |
|
62 |
|
63 |
|
64 | Store.prototype.createStream = function () {
|
65 | var stream = new Readable(streamsOpts)
|
66 | var destroyed = false
|
67 | var values = []
|
68 | var i = 0
|
69 |
|
70 | this._inflights.forEach(function (value, key) {
|
71 | values.push(value)
|
72 | })
|
73 |
|
74 | stream._read = function () {
|
75 | if (!destroyed && i < values.length) {
|
76 | this.push(values[i++])
|
77 | } else {
|
78 | this.push(null)
|
79 | }
|
80 | }
|
81 |
|
82 | stream.destroy = function () {
|
83 | if (destroyed) {
|
84 | return
|
85 | }
|
86 |
|
87 | var self = this
|
88 |
|
89 | destroyed = true
|
90 |
|
91 | process.nextTick(function () {
|
92 | self.emit('close')
|
93 | })
|
94 | }
|
95 |
|
96 | return stream
|
97 | }
|
98 |
|
99 |
|
100 |
|
101 |
|
102 | Store.prototype.del = function (packet, cb) {
|
103 | packet = this._inflights.get(packet.messageId)
|
104 | if (packet) {
|
105 | this._inflights.delete(packet.messageId)
|
106 | cb(null, packet)
|
107 | } else if (cb) {
|
108 | cb(new Error('missing packet'))
|
109 | }
|
110 |
|
111 | return this
|
112 | }
|
113 |
|
114 |
|
115 |
|
116 |
|
117 | Store.prototype.get = function (packet, cb) {
|
118 | packet = this._inflights.get(packet.messageId)
|
119 | if (packet) {
|
120 | cb(null, packet)
|
121 | } else if (cb) {
|
122 | cb(new Error('missing packet'))
|
123 | }
|
124 |
|
125 | return this
|
126 | }
|
127 |
|
128 |
|
129 |
|
130 |
|
131 | Store.prototype.close = function (cb) {
|
132 | if (this.options.clean) {
|
133 | this._inflights = null
|
134 | }
|
135 | if (cb) {
|
136 | cb()
|
137 | }
|
138 | }
|
139 |
|
140 | module.exports = Store
|