-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathindex.js
More file actions
126 lines (96 loc) · 3.43 KB
/
Copy pathindex.js
File metadata and controls
126 lines (96 loc) · 3.43 KB
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
import { Client, Server } from 'quic'
import ArbitraryPromise from 'arbitrary-promise'
// simple convenience helper
const rejectPromise = (promise, err, message) => {
promise.reject(Object.assign(err, { class: message }))
}
// the underlying quic library can only handle strings or buffers.
// This is used to convert to one of them.
const convertToSendType = (data) => {
// buffers should be sendable
if (Buffer.isBuffer(data)) return data
// objects must be stringified
if (typeof data === 'object') return JSON.stringify(data)
// all else should be strings
return data
}
class Quic {
listen(port, address = 'localhost') {
const promise = new ArbitraryPromise([['resolve', 'then'], ['reject', 'onError'], ['handleData', 'onData']])
if (!port) return promise.reject('must supply port argument!')
this._server = new Server()
this._server
.on('error', err => rejectPromise(promise, err, 'server error'))
.on('session', session => {
session
.on('error', err => rejectPromise(promise, err, 'server session error'))
.on('stream', stream => {
let message = ''
let buffer
stream
.on('error', err => rejectPromise(promise, err, 'server stream error'))
.on('data', data => {
message += data.toString()
if (buffer) buffer = Buffer.concat([buffer, data])
else buffer = data
})
.on('end', () => {
const oldWrite = stream.write.bind(stream)
stream.write = data => {
const convertedData = convertToSendType(data)
oldWrite(convertedData)
stream.end()
}
promise.handleData(message, stream, buffer)
})
})
})
this._server.listen(port, address).then(promise.resolve).catch(promise.reject)
return promise
}
async stopListening() {
// TODO returned promise to return onError instead of catch
this._server && await this._server.close()
delete this._server
}
getServer() {
return this._server
}
getAddress() {
const defaul = { port: 0, family: '', address: '' }
return this._server && this._server.address() || defaul
}
send(port, address, data) {
const promise = new ArbitraryPromise([['resolve', 'then'], ['reject', 'onError'], ['handleData', 'onData']])
if (!port || !address || !data) return promise.reject('must supply three parameters')
const convertedData = convertToSendType(data)
const client = new Client()
client.on('error', err => {
rejectPromise(promise, err, 'client error')
})
// These clients are ephemeral so we'll nuke em when they're done
client.on('close', () => client.destroy())
client.connect(port, address).then(() => {
const stream = client.request()
let message = ''
let buffer
stream
.on('error', err => rejectPromise(promise, err, 'client stream error'))
.on('data', data => {
message += data.toString()
if (buffer) buffer = Buffer.concat([buffer, data])
else buffer = data
})
.on('end', () => {
client.close()
promise.handleData(message, buffer)
})
stream.write(convertedData, () => {
promise.resolve()
stream.end()
})
})
return promise
}
}
module.exports = new Quic()