Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Implemented connect-websocket, disconnect-websocket, and update-user lambda
- Loading branch information
Showing
8 changed files
with
189 additions
and
16 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
55 changes: 55 additions & 0 deletions
55
packages/laconia-acceptance-test/src/DynamoDbWebSocketServer.js
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,55 @@ | ||
const AWS = require("aws-sdk"); | ||
const ensureApiGatewayManagementApi = require("aws-apigatewaymanagementapi"); | ||
|
||
ensureApiGatewayManagementApi(AWS); | ||
|
||
module.exports = class DynamoDbWebSocketServer { | ||
constructor(connectionTableName, endpoint) { | ||
this.connectionTableName = connectionTableName; | ||
this.DDB = new AWS.DynamoDB(); | ||
this.documentClient = new AWS.DynamoDB.DocumentClient(); | ||
this.endpoint = endpoint; | ||
} | ||
|
||
addConnection(connectionId) { | ||
const putParams = { | ||
TableName: this.connectionTableName, | ||
Item: { connectionId } | ||
}; | ||
|
||
return this.documentClient.put(putParams).promise(); | ||
} | ||
|
||
removeConnection(connectionId) { | ||
const deleteParams = { | ||
TableName: this.connectionTableName, | ||
Key: { connectionId } | ||
}; | ||
|
||
return this.documentClient.delete(deleteParams).promise(); | ||
} | ||
|
||
async broadcast(message) { | ||
const client = new AWS.ApiGatewayManagementApi({ | ||
apiVersion: "2018-11-29", | ||
endpoint: this.endpoint | ||
}); | ||
|
||
const data = await this.documentClient | ||
.scan({ | ||
TableName: this.connectionTableName | ||
}) | ||
.promise(); | ||
|
||
return Promise.all( | ||
data.Items.map(item => | ||
client | ||
.postToConnection({ | ||
ConnectionId: item.connectionId, | ||
Data: JSON.stringify(message) | ||
}) | ||
.promise() | ||
) | ||
); | ||
} | ||
}; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,18 @@ | ||
const laconia = require("@laconia/core"); | ||
const { res } = require("@laconia/event").apigateway; | ||
const DynamoDbWebSocketServer = require("./DynamoDbWebSocketServer"); | ||
|
||
const app = async (connectionId, { webSocketServer }) => { | ||
return webSocketServer.addConnection(connectionId); | ||
}; | ||
|
||
const adapter = app => async (event, laconiaContext) => { | ||
await app(event.requestContext.connectionId, laconiaContext); | ||
return res("success", 200); | ||
}; | ||
|
||
exports.handler = laconia(adapter(app)).register(() => ({ | ||
webSocketServer: new DynamoDbWebSocketServer( | ||
process.env.CONNECTION_TABLE_NAME | ||
) | ||
})); |
18 changes: 18 additions & 0 deletions
18
packages/laconia-acceptance-test/src/disconnect-websocket.js
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,18 @@ | ||
const laconia = require("@laconia/core"); | ||
const { res } = require("@laconia/event").apigateway; | ||
const DynamoDbWebSocketServer = require("./DynamoDbWebSocketServer"); | ||
|
||
const app = async (connectionId, { webSocketServer }) => { | ||
return webSocketServer.removeConnection(connectionId); | ||
}; | ||
|
||
const adapter = app => async (event, laconiaContext) => { | ||
await app(event.requestContext.connectionId, laconiaContext); | ||
return res("success", 200); | ||
}; | ||
|
||
exports.handler = laconia(adapter(app)).register(() => ({ | ||
webSocketServer: new DynamoDbWebSocketServer( | ||
process.env.CONNECTION_TABLE_NAME | ||
) | ||
})); |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,18 @@ | ||
const laconia = require("@laconia/core"); | ||
const kinesis = require("@laconia/adapter").kinesis(); | ||
const DynamoDbWebSocketServer = require("./DynamoDbWebSocketServer"); | ||
|
||
const app = async (orderEvents, { webSocketServer, event }) => { | ||
const acceptedEvents = orderEvents.filter(o => o.eventType === "accepted"); | ||
|
||
if (acceptedEvents.length > 0) { | ||
return webSocketServer.broadcast({ message: "order accepted" }); | ||
} | ||
}; | ||
|
||
exports.handler = laconia(kinesis(app)).register(() => ({ | ||
webSocketServer: new DynamoDbWebSocketServer( | ||
process.env.CONNECTION_TABLE_NAME, | ||
process.env.WEBSOCKET_ENDPOINT | ||
) | ||
})); |