Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

feat: use new sync service #17

Merged
merged 8 commits into from
Oct 15, 2021
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21,125 changes: 18,335 additions & 2,790 deletions package-lock.json

Large diffs are not rendered by default.

10 changes: 6 additions & 4 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -39,13 +39,15 @@
"test:browser": "aegir test --target browser"
},
"dependencies": {
"ioredis": "^4.26.0",
"ipaddr.js": "^2.0.0",
"emittery": "^0.10.0",
"ipaddr.js": "^2.0.1",
"isomorphic-ws": "^4.0.1",
"lodash.flatten": "^4.4.0",
"winston": "^3.3.3"
"winston": "^3.3.3",
"ws": "^8.2.3"
},
"devDependencies": {
"@types/ioredis": "^4.26.0",
"@types/ws": "^8.2.0",
"aegir": "^33.1.0"
},
"engines": {
Expand Down
17 changes: 10 additions & 7 deletions src/sync/index.js
Original file line number Diff line number Diff line change
@@ -1,14 +1,18 @@
'use strict'

const { redisClient } = require('./redis')
const { createSocket } = require('./socket')
const { createState } = require('./state')
const { createTopic } = require('./topic')
const { createSugar } = require('./sugar')
const { createPubSub } = require('./pubsub')

/** @typedef {import('winston').Logger} Logger */
/** @typedef {import('events').EventEmitter} EventEmitter */
/** @typedef {import('../runtime').RunEnv} RunEnv */
/** @typedef {import('../runtime').RunParams} RunParams */
/** @typedef {import('./types').SyncClient} SyncClient */
/** @typedef {import('./types').Request} Request */
/** @typedef {import('./types').Response} Response */

/**
* Returns a new sync client that is bound to the provided runEnv. All the operations
Expand All @@ -28,19 +32,18 @@ function newBoundClient (runenv) {
* @returns {Promise<SyncClient>}
*/
async function newClient (logger, extractor) {
const redis = await redisClient(logger)
const socket = await createSocket(logger)
const pubsub = createPubSub(logger, socket)

const base = {
...createState(logger, extractor, redis),
...createTopic(logger, extractor, redis)
...createState(logger, extractor, pubsub, socket),
...createTopic(logger, extractor, pubsub, socket)
}

return {
...base,
...createSugar(base),
close: () => {
redis.disconnect()
}
close: socket.close
}
}

Expand Down
55 changes: 55 additions & 0 deletions src/sync/pubsub.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
'use strict'

/** @typedef {import('../runtime').RunParams} RunParams */
/** @typedef {import('./types').PubSub} PubSub */
/** @typedef {import('./types').Response} Response */
/** @typedef {import('./types').Request} Request */
/** @typedef {import('./types').Socket} Socket */
/** @typedef {import('events').EventEmitter} EventEmitter */

/**
* @param {import('winston').Logger} logger
* @param {Socket} socket
* @returns {PubSub}
*/
function createPubSub (logger, socket) {
return {
publish: async (topic, payload) => {
const res = await socket.requestOnce({
publish: {
topic: topic,
payload: payload
}
})

if (res.error) {
throw res.error
}

return res.publish.seq
},
subscribe: async (key) => {
const { cancel, wait: waitSocket } = socket.request({
subscribe: {
topic: key
}
})

const wait = (async function * () {
for await (const res of waitSocket) {
if (res.error) {
throw new Error(res.error)
}

yield JSON.parse(res.subscribe)
}
})()

return { cancel, wait }
}
}
}

module.exports = {
createPubSub
}
43 changes: 0 additions & 43 deletions src/sync/redis.js

This file was deleted.

116 changes: 116 additions & 0 deletions src/sync/socket.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
'use strict'

const Emittery = require('emittery')
const WebSocket = require('isomorphic-ws')

const ENV_SYNC_SERVICE_HOST = 'SYNC_SERVICE_HOST'
const ENV_SYNC_SERVICE_PORT = 'SYNC_SERVICE_PORT'

/** @typedef {import('winston').Logger} Logger */
/** @typedef {import('events').EventEmitter} EventEmitter */
/** @typedef {import('../runtime').RunEnv} RunEnv */
/** @typedef {import('../runtime').RunParams} RunParams */
/** @typedef {import('./types').SyncClient} SyncClient */
/** @typedef {import('./types').Request} Request */
/** @typedef {import('./types').Response} Response */
/** @typedef {import('./types').ResponseIterator} ResponseIterator */
/** @typedef {import('./types').Socket} Socket */

/**
* @param {import('winston').Logger} logger
* @returns {Promise<Socket>}
*/
function createSocket (logger) {
const address = socketAddress()
const ws = new WebSocket(address)
const emitter = new Emittery()
let next = 0

return new Promise((resolve, reject) => {
ws.onopen = function open () {
resolve({
request,
requestOnce,
close: () => {
ws.close()
}
})
}

ws.onclose = function close () {
logger.info('connection to sync server closed')
}

ws.onmessage = function incoming (event) {
const res = /** @type Response */(JSON.parse(event.data.toString()))
emitter.emit(res.id, res)
}

/**
* @param {Request} req
* @returns {Promise<Response>}
*/
const requestOnce = async function (req) {
const id = (next++).toString()
const promise = emitter.once(id)

req.id = id
ws.send(JSON.stringify(req))

const data = await promise
return data
}

/**
* @param {Request} req
* @returns {ResponseIterator}
*/
const request = function (req) {
const id = (next++).toString()
const it = emitter.events(id)
let run = true

req.id = id
ws.send(JSON.stringify(req))

const cancel = () => {
run = false
emitter.clearListeners(id)
}

const wait = (async function * () {
try {
for await (const data of it) {
yield data
}
} catch (e) {
if (run) throw e
}
})()

return {
cancel,
wait
}
}
})
}

function socketAddress () {
let host = process.env[ENV_SYNC_SERVICE_HOST]
let port = process.env[ENV_SYNC_SERVICE_PORT]

if (!port) {
port = '5050'
}

if (!host) {
host = 'testground-sync-service'
}

return `ws://${host}:${port}`
}

module.exports = {
createSocket
}
Loading