INTEXT-39 Remove Experimental WebSocket Code

JIRA: https://jira.spring.io/browse/INTEXT-39
This commit is contained in:
Gary Russell
2014-11-19 17:50:51 -05:00
parent 05b7f39f39
commit 3318a4ddeb
25 changed files with 4 additions and 2301 deletions

View File

@@ -3,6 +3,9 @@ Spring Integration IP Extensions
Welcome to the Spring Integration IP Extensions project. It is intended to supplement the spring-integration-ip module with, for example, custom serializers/deserializers.
__There are currently no extensions in this project. The previous experimental WebSocket support has been superceded by the support in spring-messaging.__
# Building
If you encounter out of memory errors during the build, increase available heap and permgen for Gradle:
@@ -48,20 +51,3 @@ To generate IDEA metadata (.iml and .ipr files), do the following:
For more information, please visit the Spring Integration website at:
[http://www.springsource.org/spring-integration](http://www.springsource.org/spring-integration)
# WebSocket Server Demo
This demonstrates how to use the TCP adapters to provide a very lightweight websocket server.
Run WebSocketServerTests as a Java Application (main) and open
`file:///.../spring-integration-extensions/spring-integration-ip-extensions/src/test/java/org/springframework/integration/ip/extensions/sockjs/ws.html`
in a browser. Opening the page opens the WebSocket.
Sending 'start' begins sending an incrementing # once per second. 'stop' stops the stream (leaving the socket open), 'start' resumes again. Multiple browser instances get their own sequence #.
# Bitcoin Sample
The [bitcoin-rt project](https://github.com/cbeams/bitcoin-rt) provides a sample using the Spring Integration IP extensions:
[https://github.com/cbeams/bitcoin-rt/tree/master/java-spring-integration](https://github.com/cbeams/bitcoin-rt/tree/master/java-spring-integration)

View File

@@ -5,7 +5,7 @@ buildscript {
maven { url 'https://repo.springsource.org/plugins-snapshot' }
}
dependencies {
classpath 'org.springframework.build.gradle:docbook-reference-plugin:0.1.5'
// classpath 'org.springframework.build.gradle:docbook-reference-plugin:0.1.5'
}
}

View File

@@ -1,13 +0,0 @@
<html>
<body>
This document is the API specification for Spring Integration
<hr />
<div id="overviewBody">
<p>
Spring Integration IP Extensions project is intended to supplement
the spring-integration-ip module with, for example, custom
serializers/deserializers.
</p>
</div>
</body>
</html>

View File

@@ -1,201 +0,0 @@
Apache License
Version 2.0, January 2004
http://www.apache.org/licenses/
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
1. Definitions.
"License" shall mean the terms and conditions for use, reproduction,
and distribution as defined by Sections 1 through 9 of this document.
"Licensor" shall mean the copyright owner or entity authorized by
the copyright owner that is granting the License.
"Legal Entity" shall mean the union of the acting entity and all
other entities that control, are controlled by, or are under common
control with that entity. For the purposes of this definition,
"control" means (i) the power, direct or indirect, to cause the
direction or management of such entity, whether by contract or
otherwise, or (ii) ownership of fifty percent (50%) or more of the
outstanding shares, or (iii) beneficial ownership of such entity.
"You" (or "Your") shall mean an individual or Legal Entity
exercising permissions granted by this License.
"Source" form shall mean the preferred form for making modifications,
including but not limited to software source code, documentation
source, and configuration files.
"Object" form shall mean any form resulting from mechanical
transformation or translation of a Source form, including but
not limited to compiled object code, generated documentation,
and conversions to other media types.
"Work" shall mean the work of authorship, whether in Source or
Object form, made available under the License, as indicated by a
copyright notice that is included in or attached to the work
(an example is provided in the Appendix below).
"Derivative Works" shall mean any work, whether in Source or Object
form, that is based on (or derived from) the Work and for which the
editorial revisions, annotations, elaborations, or other modifications
represent, as a whole, an original work of authorship. For the purposes
of this License, Derivative Works shall not include works that remain
separable from, or merely link (or bind by name) to the interfaces of,
the Work and Derivative Works thereof.
"Contribution" shall mean any work of authorship, including
the original version of the Work and any modifications or additions
to that Work or Derivative Works thereof, that is intentionally
submitted to Licensor for inclusion in the Work by the copyright owner
or by an individual or Legal Entity authorized to submit on behalf of
the copyright owner. For the purposes of this definition, "submitted"
means any form of electronic, verbal, or written communication sent
to the Licensor or its representatives, including but not limited to
communication on electronic mailing lists, source code control systems,
and issue tracking systems that are managed by, or on behalf of, the
Licensor for the purpose of discussing and improving the Work, but
excluding communication that is conspicuously marked or otherwise
designated in writing by the copyright owner as "Not a Contribution."
"Contributor" shall mean Licensor and any individual or Legal Entity
on behalf of whom a Contribution has been received by Licensor and
subsequently incorporated within the Work.
2. Grant of Copyright License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
copyright license to reproduce, prepare Derivative Works of,
publicly display, publicly perform, sublicense, and distribute the
Work and such Derivative Works in Source or Object form.
3. Grant of Patent License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
(except as stated in this section) patent license to make, have made,
use, offer to sell, sell, import, and otherwise transfer the Work,
where such license applies only to those patent claims licensable
by such Contributor that are necessarily infringed by their
Contribution(s) alone or by combination of their Contribution(s)
with the Work to which such Contribution(s) was submitted. If You
institute patent litigation against any entity (including a
cross-claim or counterclaim in a lawsuit) alleging that the Work
or a Contribution incorporated within the Work constitutes direct
or contributory patent infringement, then any patent licenses
granted to You under this License for that Work shall terminate
as of the date such litigation is filed.
4. Redistribution. You may reproduce and distribute copies of the
Work or Derivative Works thereof in any medium, with or without
modifications, and in Source or Object form, provided that You
meet the following conditions:
(a) You must give any other recipients of the Work or
Derivative Works a copy of this License; and
(b) You must cause any modified files to carry prominent notices
stating that You changed the files; and
(c) You must retain, in the Source form of any Derivative Works
that You distribute, all copyright, patent, trademark, and
attribution notices from the Source form of the Work,
excluding those notices that do not pertain to any part of
the Derivative Works; and
(d) If the Work includes a "NOTICE" text file as part of its
distribution, then any Derivative Works that You distribute must
include a readable copy of the attribution notices contained
within such NOTICE file, excluding those notices that do not
pertain to any part of the Derivative Works, in at least one
of the following places: within a NOTICE text file distributed
as part of the Derivative Works; within the Source form or
documentation, if provided along with the Derivative Works; or,
within a display generated by the Derivative Works, if and
wherever such third-party notices normally appear. The contents
of the NOTICE file are for informational purposes only and
do not modify the License. You may add Your own attribution
notices within Derivative Works that You distribute, alongside
or as an addendum to the NOTICE text from the Work, provided
that such additional attribution notices cannot be construed
as modifying the License.
You may add Your own copyright statement to Your modifications and
may provide additional or different license terms and conditions
for use, reproduction, or distribution of Your modifications, or
for any such Derivative Works as a whole, provided Your use,
reproduction, and distribution of the Work otherwise complies with
the conditions stated in this License.
5. Submission of Contributions. Unless You explicitly state otherwise,
any Contribution intentionally submitted for inclusion in the Work
by You to the Licensor shall be under the terms and conditions of
this License, without any additional terms or conditions.
Notwithstanding the above, nothing herein shall supersede or modify
the terms of any separate license agreement you may have executed
with Licensor regarding such Contributions.
6. Trademarks. This License does not grant permission to use the trade
names, trademarks, service marks, or product names of the Licensor,
except as required for reasonable and customary use in describing the
origin of the Work and reproducing the content of the NOTICE file.
7. Disclaimer of Warranty. Unless required by applicable law or
agreed to in writing, Licensor provides the Work (and each
Contributor provides its Contributions) on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
implied, including, without limitation, any warranties or conditions
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
PARTICULAR PURPOSE. You are solely responsible for determining the
appropriateness of using or redistributing the Work and assume any
risks associated with Your exercise of permissions under this License.
8. Limitation of Liability. In no event and under no legal theory,
whether in tort (including negligence), contract, or otherwise,
unless required by applicable law (such as deliberate and grossly
negligent acts) or agreed to in writing, shall any Contributor be
liable to You for damages, including any direct, indirect, special,
incidental, or consequential damages of any character arising as a
result of this License or out of the use or inability to use the
Work (including but not limited to damages for loss of goodwill,
work stoppage, computer failure or malfunction, or any and all
other commercial damages or losses), even if such Contributor
has been advised of the possibility of such damages.
9. Accepting Warranty or Additional Liability. While redistributing
the Work or Derivative Works thereof, You may choose to offer,
and charge a fee for, acceptance of support, warranty, indemnity,
or other liability obligations and/or rights consistent with this
License. However, in accepting such obligations, You may act only
on Your own behalf and on Your sole responsibility, not on behalf
of any other Contributor, and only if You agree to indemnify,
defend, and hold each Contributor harmless for any liability
incurred by, or claims asserted against, such Contributor by reason
of your accepting any such warranty or additional liability.
END OF TERMS AND CONDITIONS
APPENDIX: How to apply the Apache License to your work.
To apply the Apache License to your work, attach the following
boilerplate notice, with the fields enclosed by brackets "[]"
replaced with your own identifying information. (Don't include
the brackets!) The text should be enclosed in the appropriate
comment syntax for the file format. We also recommend that a
file or class name and description of purpose be included on the
same "printed page" as the copyright notice for easier
identification within third-party archives.
Copyright [yyyy] [name of copyright owner]
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.

View File

@@ -1,21 +0,0 @@
========================================================================
== NOTICE file corresponding to section 4 d of the Apache License, ==
== Version 2.0, in this case for the Spring Integration distribution. ==
========================================================================
This product includes software developed by
the Apache Software Foundation (http://www.apache.org).
The end-user documentation included with a redistribution, if any,
must include the following acknowledgement:
"This product includes software developed by the Spring Framework
Project (http://www.springframework.org)."
Alternatively, this acknowledgement may appear in the software itself,
if and wherever such third-party acknowledgements normally appear.
The names "Spring", "Spring Framework", and "Spring Integration" must
not be used to endorse or promote products derived from this software
without prior written permission. For written permission, please contact
enquiries@springsource.com.

View File

@@ -1,4 +0,0 @@
/**
* Root package of the IP Extensions.
*/
package org.springframework.integration.x.ip;

View File

@@ -1,194 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.serializer;
import java.io.IOException;
import java.io.InputStream;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer;
/**
* Base class for (de)Serializers that start with an HTTP-like protocol then
* switch to some other protocol.
*
* @author Gary Russell
* @since 3.0
*
*/
public abstract class AbstractHttpSwitchingDeserializer implements StatefulDeserializer<DataFrame> {
protected final Log logger = LogFactory.getLog(this.getClass());
protected volatile int maxMessageSize = 2048;
private final Map<InputStream, BasicState> streamState = new ConcurrentHashMap<InputStream, BasicState>();
protected final ByteArrayCrLfSerializer crlfDeserializer = new ByteArrayCrLfSerializer();
private static final Pattern requestLinePattern = Pattern.compile("GET *([^ ]+) *HTTP/");
public void setMaxMessageSize(int maxMessageSize) {
this.maxMessageSize = maxMessageSize;
}
protected ByteArrayCrLfSerializer getCrlfDeserializer() {
return crlfDeserializer;
}
@Override
public abstract DataFrame deserialize(InputStream inputStream) throws IOException;
protected BasicState getStreamState(InputStream inputStream) {
return streamState.get(inputStream);
}
/**
* Returns null if we've switched from HTTP-like protocol; headers otherwise.
* @param inputStream
* @return null or list of DataFrame, where the first frame contains the headers.
* Implementations may add additional frames.
* @throws IOException
*/
protected List<DataFrame> checkStreaming(InputStream inputStream) throws IOException {
BasicState isStreaming = this.streamState.get(inputStream);
if (isStreaming == null) { //Consume the headers - TODO - check status
StringBuilder headersBuilder = new StringBuilder();
byte[] headers = new byte[this.maxMessageSize];
String path = null;
String queryString = null;
int headersLength;
do {
headersLength = this.crlfDeserializer.fillToCrLf(inputStream, headers);
String header = new String(headers, 0, headersLength, "UTF-8");
if (path == null) {
if (header.startsWith("GET")) {
Matcher requestLineMatcher = requestLinePattern.matcher(header);
if (requestLineMatcher.find()) {
path = requestLineMatcher.group(1);
if (path.contains("?")) {
int queryStarts = path.indexOf("?");
queryString = path.substring(queryStarts + 1);
path = path.substring(0, queryStarts);
}
}
}
}
headersBuilder.append(header).append("\r\n");
}
while (headersLength > 0);
BasicState basicState = createState();
basicState.setPath(path);
basicState.setQueryString(queryString);
List<DataFrame> dataList = new ArrayList<DataFrame>();
List<DataFrame> decodedHeaders = decodeHeaders(headersBuilder.toString(), basicState, dataList);
this.streamState.put(inputStream, basicState);
return decodedHeaders;
}
return null;
}
protected BasicState createState() {
return new BasicState();
}
protected void checkClosure(int bite) throws IOException {
if (bite < 0) {
logger.debug("Socket closed during message assembly");
throw new IOException("Socket closed during message assembly");
}
}
protected List<DataFrame> decodeHeaders(String frameData, BasicState state, List<DataFrame> dataList) {
// TODO: Full header separation - mvc utils?
if (logger.isDebugEnabled()) {
logger.debug("Received:Headers\r\n" + frameData);
}
dataList.add(createDataFrame(DataFrame.TYPE_HEADERS, frameData));
return dataList;
}
protected DataFrame createDataFrame(int type, String frameData) {
return new DataFrame(type, frameData);
}
@Override
public void removeState(Object key) {
this.streamState.remove(key);
}
public BasicState getState(Object key) {
return this.streamState.get(key);
}
public static class BasicState {
private volatile String path;
private volatile String queryString;
private volatile DataFrame pendingFrame;
private final List<byte[]> fragments = new ArrayList<byte[]>();
public DataFrame getPendingFrame() {
return pendingFrame;
}
public void setPendingFrame(DataFrame pendingFrame) {
this.pendingFrame = pendingFrame;
}
public List<byte[]> getFragments() {
return fragments;
}
public String getPath() {
return path;
}
private void setPath(String path) {
this.path = path;
}
public String getQueryString() {
return queryString;
}
private void setQueryString(String queryString) {
this.queryString = queryString;
}
@Override
public String toString() {
return "BasicState [" +
(this.path != null ? ("path=" + this.path) : "") +
(this.queryString != null ? (", queryString=" + this.queryString) : "") +
(this.pendingFrame != null ? (", pendingFrame=" + this.pendingFrame) : "") +
(this.fragments.size() > 0 ? (", fragments.size()=" + this.fragments.size()) : "") +
"]";
}
}
}

View File

@@ -1,65 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.serializer;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class DataFrame {
public static final int TYPE_INVALID = 0;
public static final int TYPE_HEADERS = 1;
public static final int TYPE_DATA = 4;
public static final int TYPE_DATA_BINARY = 260;
protected final int type;
protected final String payload;
protected final byte[] binary;
public DataFrame(int type, String payload) {
this(type, payload, null);
}
public DataFrame(int type, byte[] binary) {
this(type, null, binary);
}
public DataFrame(int type, String payload, byte[] binary) {
this.type = type;
this.payload = payload;
this.binary = binary;
}
public int getType() {
return this.type;
}
public String getPayload() {
return this.payload;
}
public byte[] getBinary() {
return binary;
}
}

View File

@@ -1,28 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.serializer;
import org.springframework.core.serializer.Deserializer;
/**
* @author Gary Russell
* @since 3.0
*
*/
public interface StatefulDeserializer<T> extends Deserializer<T> {
void removeState(Object key);
}

View File

@@ -1,60 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.integration.ip.tcp.connection.TcpConnectionEvent;
import org.springframework.integration.ip.tcp.connection.TcpConnectionSupport;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketEvent extends TcpConnectionEvent {
private static final long serialVersionUID = -6788341703196233248L;
public enum WebSocketEventType implements EventType {
HANDSHAKE_COMPLETE,
WEBSOCKET_CLOSED
}
private final String path;
private final String queryString;
public WebSocketEvent(TcpConnectionSupport connection, WebSocketEventType type, String path, String queryString) {
super(connection, type, (String) new DirectFieldAccessor(connection).getPropertyValue("connectionFactoryName"));
this.path = path;
this.queryString = queryString;
}
public String getPath() {
return path;
}
public String getQueryString() {
return queryString;
}
@Override
public String toString() {
return super.toString().replace("TcpConnectionEvent", "WebSocketEvent")
.replace("]", ", path=" + this.path + ", queryString=" + this.queryString + "]");
}
}

View File

@@ -1,98 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import org.springframework.integration.x.ip.serializer.DataFrame;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketFrame extends DataFrame {
public static final int TYPE_FRAGMENTED_CONTROL = 256;
public static final int TYPE_INVALID_UTF8 = 512;
public static final int TYPE_PING = 5;
public static final int TYPE_PONG = 6;
public static final int TYPE_OPEN = 7;
public static final int TYPE_CLOSE = 8;
private static final String[] typeToString = new String[] {
"Invalid", "Headers", "**", "**", "Data", "Ping", "Pong", "Open", "Close"
};
private volatile short status = -1;
private volatile int rsv;
public WebSocketFrame(int type, String payload) {
super(type, payload);
}
public WebSocketFrame(int type, byte[] binary) {
super(type, binary);
}
public WebSocketFrame(int type, String payload, byte[] binary) {
super(type, payload, binary);
}
public short getStatus() {
return status;
}
public void setStatus(short status) {
this.status = status;
}
public void setRsv(int rsv) {
this.rsv = rsv;
}
public int getRsv() {
return rsv;
}
@Override
public String toString() {
int len = 0;
boolean trunc = false;
if (this.payload!= null) {
len = Math.min(100, payload.length());
trunc = len < payload.length();
}
String typeAsString;
if ((type & 0xff) < typeToString.length) {
typeAsString = typeToString[type & 0xff];
}
else {
typeAsString = Integer.toString(type);
}
return "WebSocketFrame [type=" + typeAsString + "(" +
type + ")"+ (payload == null ? "" : ", payload=" + payload.substring(0, len) +
(trunc ? "..." : "")) +
", binary=" + binary +
(binary != null ? ", binary.length=" + binary.length : "") +
", status=" + status + ", rsv=" + rsv + "]";
}
}

View File

@@ -1,33 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketHeaders {
private WebSocketHeaders() {}
private final static String WS = "websocket_";
public final static String PATH = WS + "path";
public final static String QUERY_STRING = WS + "queryString";
}

View File

@@ -1,44 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import java.util.Map;
import org.springframework.integration.ip.tcp.connection.TcpConnection;
import org.springframework.integration.ip.tcp.connection.TcpMessageMapper;
import org.springframework.integration.x.ip.websocket.WebSocketTcpConnectionInterceptorFactory.WebSocketTcpConnectionInterceptor;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketMessageMapper extends TcpMessageMapper {
private final WebSocketTcpConnectionInterceptorFactory connectionInterceptorFactory;
public WebSocketMessageMapper(WebSocketTcpConnectionInterceptorFactory connectionInterceptorFactory) {
this.connectionInterceptorFactory = connectionInterceptorFactory;
}
@Override
protected Map<String, ?> supplyCustomHeaders(TcpConnection connection) {
WebSocketTcpConnectionInterceptor interceptor = this.connectionInterceptorFactory.locateInterceptor(connection);
return interceptor == null ? null : interceptor.getAdditionalHeaders();
}
}

View File

@@ -1,555 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.io.UnsupportedEncodingException;
import java.nio.ByteBuffer;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import org.apache.commons.codec.binary.Base64;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.core.serializer.Serializer;
import org.springframework.integration.MessagingException;
import org.springframework.integration.ip.tcp.serializer.SoftEndOfStreamException;
import org.springframework.integration.x.ip.serializer.AbstractHttpSwitchingDeserializer;
import org.springframework.integration.x.ip.serializer.DataFrame;
import org.springframework.util.Assert;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketSerializer extends AbstractHttpSwitchingDeserializer implements Serializer<Object> {
private static final String HTTP_1_1_101_WEB_SOCKET_PROTOCOL_HANDSHAKE_SPRING_INTEGRATION =
"HTTP/1.1 101 Web Socket Protocol Handshake - Spring Integration\r\n";
private static final Set<Short> INVALID_STATUS = new HashSet<Short>(
Arrays.asList((short) 1004, (short) 1005, (short) 1006, (short) 1012, (short) 1013, (short) 1014, (short) 1015));
private volatile boolean server;
private boolean validateUtf8;
private volatile Boolean streamChecked;
private volatile boolean nio;
private volatile DirectFieldAccessor streamAccessor;
public void setServer(boolean server) {
this.server = server;
}
/**
* Validate UTF-8 (required for Autobahn tests).
* @param validateUtf8
*/
public void setValidateUtf8(boolean validateUtf8) {
this.validateUtf8 = validateUtf8;
}
@Override
protected DataFrame createDataFrame(int type, String frameData) {
return new WebSocketFrame(type, frameData);
}
@Override
protected BasicState createState() {
return new WebSocketState();
}
@Override
public void serialize(final Object frame, OutputStream outputStream)
throws IOException {
String data = "";
WebSocketFrame theFrame = null;
if (frame instanceof String) {
data = (String) frame;
theFrame = new WebSocketFrame(WebSocketFrame.TYPE_DATA, data);
}
else if (frame instanceof WebSocketFrame) {
theFrame = (WebSocketFrame) frame;
data = theFrame.getPayload();
}
if (data != null && data.startsWith("HTTP/1.1")) {
outputStream.write(data.getBytes());
return;
}
int lenBytes;
int payloadLen = this.server ? 0 : 0x80; //masked
boolean close = theFrame.getType() == WebSocketFrame.TYPE_CLOSE;
boolean ping = theFrame.getType() == WebSocketFrame.TYPE_PING;
boolean pong = theFrame.getType() == WebSocketFrame.TYPE_PONG;
byte[] bytes = theFrame.getBinary() != null ? theFrame.getBinary() : data.getBytes("UTF-8");
int length = bytes.length;
if (close) {
length += 2;
}
if (length >= Math.pow(2, 16)) {
lenBytes = 8;
payloadLen |= 127;
}
else if (length > 125) {
lenBytes = 2;
payloadLen |= 126;
}
else {
lenBytes = 0;
payloadLen |= length;
}
int mask = (int) System.currentTimeMillis();
ByteBuffer buffer = ByteBuffer.allocate(length + 6 + lenBytes);
if (ping) {
buffer.put((byte) 0x89);
}
else if (pong) {
buffer.put((byte) 0x8a);
}
else if (close) {
buffer.put((byte) 0x88);
}
else if (theFrame.getType() == WebSocketFrame.TYPE_DATA_BINARY) {
buffer.put((byte) 0x82);
}
else {
// Final fragment; text
buffer.put((byte) 0x81);
}
buffer.put((byte) payloadLen);
if (lenBytes == 2) {
buffer.putShort((short) length);
}
else if (lenBytes == 8) {
buffer.putLong(length);
}
byte[] maskBytes = new byte[4];
if (!server) {
buffer.putInt(mask);
buffer.position(buffer.position() - 4);
buffer.get(maskBytes);
}
if (close) {
buffer.putShort(theFrame.getStatus());
// TODO: mask status when client
}
for (int i = 0; i < bytes.length; i++) {
if (server) {
buffer.put(bytes[i]);
}
else {
buffer.put((byte) (bytes[i] ^ maskBytes[i % 4]));
}
}
outputStream.write(buffer.array(), 0, buffer.position());
}
@Override
public DataFrame deserialize(InputStream inputStream) throws IOException {
if (this.streamChecked == null) {
this.nio = inputStream.getClass().getName().endsWith("TcpNioConnection$ChannelInputStream");
this.streamAccessor = new DirectFieldAccessor(inputStream);
this.streamChecked = Boolean.TRUE;
}
DataFrame frame = null;
BasicState state = this.getState(inputStream);
if (state != null) {
frame = state.getPendingFrame();
}
while (frame == null || (frame.getPayload() == null && frame.getBinary() == null)) {
frame = doDeserialize(inputStream, frame);
if (frame.getPayload() == null && frame.getBinary() == null) {
state.setPendingFrame(frame);
}
}
return frame;
}
private DataFrame doDeserialize(InputStream inputStream, DataFrame protoFrame) throws IOException {
List<DataFrame> headers = checkStreaming(inputStream);
if (headers != null) {
return headers.get(0);
}
int bite;
if (logger.isDebugEnabled()) {
logger.debug("Available to read:" + inputStream.available());
}
boolean done = false;
int len = 0;
int n = 0;
int dataInx = 0;
byte[] buffer = null;
boolean fin = false;
boolean ping = false;
boolean pong = false;
boolean close = false;
boolean binary = false;
boolean invalid = false;
String invalidText = null;
boolean fragmentedControl = false;
int lenBytes = 0;
byte[] mask = new byte[4];
int maskInx = 0;
int rsv = 0;
while (!done ) {
bite = inputStream.read();
// logger.debug("Read:" + Integer.toHexString(bite));
if (this.nio) {
bite = checkclosed(bite, inputStream);
}
if (bite < 0 && n == 0) {
throw new SoftEndOfStreamException("Stream closed between payloads");
}
checkClosure(bite);
switch (n++) {
case 0:
fin = (bite & 0x80) > 0;
rsv = (bite & 0x70) >> 4;
bite &= 0x0f;
switch (bite) {
case 0x00:
logger.debug("Continuation, fin=" + fin);
if (protoFrame == null) {
invalid = true;
invalidText = "Unexpected continuation frame";
}
else {
binary = protoFrame.getType() == WebSocketFrame.TYPE_DATA_BINARY;
}
this.getState(inputStream).setPendingFrame(null);
break;
case 0x01:
logger.debug("Text, fin=" + fin);
if (protoFrame != null) {
invalid = true;
invalidText = "Expected continuation frame";
}
break;
case 0x02:
logger.debug("Binary, fin=" + fin);
if (protoFrame != null) {
invalid = true;
invalidText = "Expected continuation frame";
}
binary = true;
break;
case 0x08:
logger.debug("Close, fin=" + fin);
fragmentedControl = !fin;
close = true;
break;
case 0x09:
ping = true;
binary = true;
fragmentedControl = !fin;
logger.debug("Ping, fin=" + fin);
break;
case 0x0a:
pong = true;
fragmentedControl = !fin;
logger.debug("Pong, fin=" + fin);
break;
case 0x03:
case 0x04:
case 0x05:
case 0x06:
case 0x07:
case 0x0b:
case 0x0c:
case 0x0d:
case 0x0e:
case 0x0f:
invalid = true;
invalidText = "Reserved opcode " + Integer.toHexString(bite);
break;
default:
throw new IOException("Unexpected opcode " + Integer.toHexString(bite));
}
break;
case 1:
if (this.server) {
if ((bite & 0x80) == 0) {
throw new IOException("Illegal: Expected masked data from client");
}
bite &= 0x7f;
}
if ((bite & 0x80) > 0) {
throw new IOException("Illegal: Received masked data from server");
}
if (bite < 126) {
len = bite;
buffer = new byte[len];
}
else if (bite == 126) {
lenBytes = 2;
}
else {
lenBytes = 8;
}
break;
case 2:
case 3:
case 4:
case 5:
if (lenBytes > 4 && bite != 0) {
throw new IOException("Max supported length exceeded");
}
case 6:
if (lenBytes > 3 && (bite & 0x80) > 0) {
throw new IOException("Max supported length exceeded");
}
case 7:
case 8:
case 9:
if (lenBytes-- > 0) {
len = len << 8 | (bite & 0xff);
if (lenBytes == 0) {
buffer = new byte[len];
}
break;
}
default:
if (this.server && maskInx < 4) {
mask[maskInx++] = (byte) bite;
}
else {
if (this.server) {
bite ^= mask[dataInx % 4];
}
buffer[dataInx++] = (byte) bite;
}
done = (server ? maskInx == 4 : true) && dataInx >= len;
}
};
WebSocketFrame frame;
if (fragmentedControl) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_FRAGMENTED_CONTROL, "Fragmented control frame", buffer);
}
else if (invalid) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_INVALID, invalidText, buffer);
}
else if (!fin) {
List<byte[]> fragments = this.getState(inputStream).getFragments();
fragments.add(buffer);
logger.debug("Fragment");
return new WebSocketFrame(binary ? WebSocketFrame.TYPE_DATA_BINARY : WebSocketFrame.TYPE_DATA, (String) null);
}
else if (ping) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_PING, buffer);
}
else if (pong) {
String data = new String(buffer, "UTF-8");
frame = new WebSocketFrame(WebSocketFrame.TYPE_PONG, data);
}
else if (close) {
String data = new String(buffer, "UTF-8");
if (data.length() >= 2) {
data = data.substring(2);
}
WebSocketFrame closeFrame = new WebSocketFrame(WebSocketFrame.TYPE_CLOSE, data);
short status = 1000;
if (buffer.length >= 2) {
status = (short) ((buffer[0] << 8) | (buffer[1] & 0xff));
closeFrame.setStatus(status);
}
if (buffer.length == 1 || buffer.length > 125 ||
(buffer.length > 2 && !validateUtf8IfNecessary(buffer, 2, data)) ||
status < 1000 || INVALID_STATUS.contains(status) || (status >= 1016 && status < 3000) || status >= 5000) {
// Simply close in this case; no close reply
((WebSocketState) this.getState(inputStream)).setCloseInitiated(true);
}
frame = closeFrame;
}
else {
List<byte[]> fragments = this.getState(inputStream).getFragments();
if (fragments.size() == 0) {
if (binary) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_DATA_BINARY, buffer);
}
else {
String data = new String(buffer, "UTF-8");
if (!validateUtf8IfNecessary(buffer, 0, data)) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_INVALID_UTF8, "Invalid UTF-8", buffer);
}
else {
frame = new WebSocketFrame(WebSocketFrame.TYPE_DATA, data);
}
}
}
else {
fragments.add(buffer);
int utf8Len = 0;
for (byte[] fragment : fragments) {
utf8Len += fragment.length;
}
byte[] reconstructed = new byte[utf8Len];
int utf8Pos = 0;
for (byte[] fragment : fragments) {
System.arraycopy(fragment, 0, reconstructed, utf8Pos, fragment.length);
utf8Pos += fragment.length;
}
fragments.clear();
if (binary) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_DATA_BINARY, reconstructed);
}
else {
String data = new String(reconstructed, "UTF-8");
if (!validateUtf8IfNecessary(reconstructed, 0, data)) {
frame = new WebSocketFrame(WebSocketFrame.TYPE_INVALID_UTF8, "Invalid UTF-8", reconstructed);
}
else {
frame = new WebSocketFrame(WebSocketFrame.TYPE_DATA, data);
}
}
}
}
if (rsv > 0) {
frame.setRsv(rsv);
}
return frame;
}
/**
* TODO: workaround for INT-2936
*/
private int checkclosed(int bite, InputStream inputStream) {
if (bite < 0) { // possibly a closed stream
try {
if ((Boolean) streamAccessor.getPropertyValue("isClosed") &&
inputStream.available() == 0) {
return -1;
}
else {
return bite & 0xff;
}
}
catch (Exception e) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to check closed", e);
}
return bite;
}
}
else {
return bite;
}
}
private boolean validateUtf8IfNecessary(byte[] buffer, int offset, String data) {
if (this.validateUtf8) {
try {
byte[] bytes = data.getBytes("UTF-8");
if (bytes.length != buffer.length - offset) {
return false;
}
for (int i = 0; i < bytes.length; i++) {
if (buffer[i + offset] != bytes[i]) {
return false;
}
}
}
catch (UnsupportedEncodingException e) {
throw new MessagingException("UTF-8 Conversion error");
}
}
return true;
}
@Override
protected void checkClosure(int bite) throws IOException {
if (bite < 0) {
logger.debug("Socket closed during message assembly");
throw new IOException("Socket closed during message assembly");
}
}
@Override
public void removeState(Object inputStream) {
super.removeState(inputStream);
}
public WebSocketFrame generateHandshake(WebSocketFrame frame) throws Exception {
Assert.isTrue(frame.getType() == WebSocketFrame.TYPE_HEADERS, "Expected headers:" + frame);
String[] headers = frame.getPayload().split("\\r\\n");
String key = null;
String version = null;
for (String header : headers) {
if (header.toLowerCase().startsWith("sec-websocket-key")) {
key = header.split(":")[1].trim();
}
else if (header.toLowerCase().startsWith("sec-websocket-version")) {
version = header.split(":")[1].trim();
}
}
if (key == null) {
throw new WebSocketUpgradeException("400 Bad Request: No sec-websocket-key header detected");
}
else if (!"13".equals(version)) {
throw new WebSocketUpgradeException("426 Upgrade Required", "sec-websocket-version: 13\r\n");
}
String handshake = HTTP_1_1_101_WEB_SOCKET_PROTOCOL_HANDSHAKE_SPRING_INTEGRATION +
"Upgrade: WebSocket\r\n" +
"Connection: Upgrade\r\n" +
"Sec-WebSocket-Accept: " + this.generateWebSocketAccept(key) + "\r\n\r\n";
return new WebSocketFrame(WebSocketFrame.TYPE_DATA, handshake);
}
private String generateWebSocketAccept(String key) throws NoSuchAlgorithmException {
MessageDigest md = MessageDigest.getInstance("SHA-1");
String toDigest = key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";
byte[] acceptStringBytes = md.digest(toDigest.getBytes());
acceptStringBytes = Base64.encodeBase64(acceptStringBytes);
String acceptString = new String(acceptStringBytes);
return acceptString;
}
public static class WebSocketState extends BasicState {
private volatile boolean closeInitiated;
private volatile boolean expectingPong;
public boolean isCloseInitiated() {
return this.closeInitiated;
}
public void setCloseInitiated(boolean closeInitiated) {
this.closeInitiated = closeInitiated;
}
public boolean isExpectingPong() {
return this.expectingPong;
}
public void setExpectingPong(boolean expectingPong) {
this.expectingPong = expectingPong;
}
}
}

View File

@@ -1,369 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentHashMap;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.core.serializer.Deserializer;
import org.springframework.integration.Message;
import org.springframework.integration.MessageHandlingException;
import org.springframework.integration.MessageHeaders;
import org.springframework.integration.MessagingException;
import org.springframework.integration.aggregator.ResequencingMessageGroupProcessor;
import org.springframework.integration.aggregator.ResequencingMessageHandler;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.core.MessageHandler;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.ip.tcp.connection.TcpConnection;
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactory;
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorSupport;
import org.springframework.integration.ip.tcp.connection.TcpConnectionSupport;
import org.springframework.integration.ip.tcp.connection.TcpNioConnection;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.x.ip.websocket.WebSocketEvent.WebSocketEventType;
import org.springframework.integration.x.ip.websocket.WebSocketSerializer.WebSocketState;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.util.Assert;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketTcpConnectionInterceptorFactory extends IntegrationObjectSupport
implements TcpConnectionInterceptorFactory {
private static final Message<WebSocketFrame> PING = MessageBuilder.withPayload(
new WebSocketFrame(WebSocketFrame.TYPE_PING, "Ping from SI")).build();
private static final Log logger = LogFactory.getLog(WebSocketTcpConnectionInterceptor.class);
private final Map<TcpConnection, WebSocketTcpConnectionInterceptor> connections =
new ConcurrentHashMap<TcpConnection, WebSocketTcpConnectionInterceptor>();
private volatile TaskScheduler taskScheduler;
private volatile long pingInterval = 25000;
private final Runnable pinger = new Runnable() {
@Override
public void run() {
// Add 100ms to allow for heuristics
long pingFilter = System.currentTimeMillis() - pingInterval + 100;
for (Entry<TcpConnection, WebSocketTcpConnectionInterceptor> entry : connections.entrySet()) {
TcpConnection connection = entry.getKey();
String connectionId = connection.getConnectionId();
if (entry.getValue().getLastReceiveTime() <= pingFilter) {
try {
if (logger.isDebugEnabled()) {
logger.debug("Sending Ping to " + connectionId);
}
connection.send(PING);
}
catch (Exception e) {
logger.error("Failed to send Ping to " + connectionId, e);
connection.close();
}
}
else {
if (logger.isTraceEnabled()) {
logger.trace("Skipping PING for " + connectionId + " due to recent send");
}
}
}
if (pingInterval > 0) {
taskScheduler.schedule(pinger, new Date(System.currentTimeMillis() + pingInterval));
}
}
};
@Override
public void setTaskScheduler(TaskScheduler taskScheduler) {
this.taskScheduler = taskScheduler;
}
/**
* The time between PINGs sent for idle connections. Must be less than half the
* socket timeout (if any).
* @param pingInterval
*/
public void setPingInterval(long pingInterval) {
this.pingInterval = pingInterval;
}
@Override
protected void onInit() throws Exception {
super.onInit();
if (this.pingInterval > 0) {
if (this.taskScheduler == null) {
this.taskScheduler = this.getTaskScheduler();
}
this.taskScheduler.schedule(this.pinger, new Date(System.currentTimeMillis() + this.pingInterval));
}
}
@Override
public TcpConnectionInterceptorSupport getInterceptor() {
return new WebSocketTcpConnectionInterceptor();
}
public WebSocketTcpConnectionInterceptor locateInterceptor(TcpConnection connection) {
return this.connections.get(connection);
}
public class WebSocketTcpConnectionInterceptor extends TcpConnectionInterceptorSupport {
private volatile boolean shook;
private final DirectChannel resequenceChannel = new DirectChannel();
private final EventDrivenConsumer resequencer;
private long lastReceiveTime;
public WebSocketTcpConnectionInterceptor() {
super();
ResequencingMessageHandler handler = new ResequencingMessageHandler(new ResequencingMessageGroupProcessor());
handler.setReleasePartialSequences(true);
DirectChannel resequenced = new DirectChannel();
resequenced.setBeanName("resequencedWSFrames");
handler.setOutputChannel(resequenced);
this.resequencer = new EventDrivenConsumer(this.resequenceChannel, handler);
resequenced.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
doOnMessage(message);
}
});
this.resequencer.afterPropertiesSet();
this.resequencer.start();
}
public long getLastReceiveTime() {
return lastReceiveTime;
}
/**
* When using NIO, we have to resequence the messages because frames may
* arrive out of order. This is particularly an issue for some of the
* Autobahn tests where, for example, many pings are sent and the test
* expects the pongs to come back in the same order.
*/
@Override
public boolean onMessage(Message<?> message) {
this.lastReceiveTime = System.currentTimeMillis();
if (this.getTheConnection() instanceof TcpNioConnection && message.getHeaders().getCorrelationId() != null) {
resequenceChannel.send(message);
return true;
}
else {
return this.doOnMessage(message);
}
}
public boolean doOnMessage(Message<?> message) {
Assert.isInstanceOf(WebSocketFrame.class, message.getPayload());
WebSocketFrame payload = (WebSocketFrame) message.getPayload();
WebSocketState state = getState(message);
if (logger.isTraceEnabled()) {
logger.trace(state);
}
if (payload.getRsv() > 0) {
if (logger.isDebugEnabled()) {
logger.debug("Reserved bits:" + payload.getRsv());
}
this.protocolViolation(message);
}
else if (payload.getType() == WebSocketFrame.TYPE_CLOSE) {
try {
if (logger.isDebugEnabled()) {
logger.debug("Close, status:" + payload.getStatus());
}
// If we initiated the close, just close.
if (!state.isCloseInitiated()) {
if (payload.getStatus() < 0) {
payload.setStatus((short) 1000);
}
this.send(message);
}
WebSocketEvent event = new WebSocketEvent(this.getTheConnection(),
WebSocketEventType.WEBSOCKET_CLOSED, state.getPath(), state.getQueryString());
this.getTheConnection().publishEvent(event);
this.close();
}
catch (Exception e) {
throw new MessageHandlingException(message, "Send failed", e);
}
}
else if (state == null || state.isCloseInitiated()) {
if (logger.isWarnEnabled()) {
logger.warn("Message dropped - close initiated:" + message);
}
}
else if ((payload.getType() & 0xff) == WebSocketFrame.TYPE_INVALID) {
if (logger.isDebugEnabled()) {
logger.debug("Invalid:" + payload.getPayload());
}
this.protocolViolation(message);
}
else if (payload.getType() == WebSocketFrame.TYPE_FRAGMENTED_CONTROL) {
if (logger.isDebugEnabled()) {
logger.debug("Fragmented Control Op");
}
this.protocolViolation(message);
}
else if (payload.getType() == WebSocketFrame.TYPE_PING) {
try {
if (logger.isDebugEnabled()) {
logger.debug("Ping received on " + this.getConnectionId() + ":"
+ new String(payload.getBinary(), "UTF-8"));
}
if (payload.getBinary().length > 125) {
this.protocolViolation(message);
}
else {
WebSocketFrame pong = new WebSocketFrame(WebSocketFrame.TYPE_PONG, payload.getBinary());
this.send(MessageBuilder.withPayload(pong)
.copyHeaders(message.getHeaders())
.build());
}
}
catch (Exception e) {
throw new MessageHandlingException(message, "Send failed", e);
}
}
else if (payload.getType() == WebSocketFrame.TYPE_PONG) {
if (logger.isDebugEnabled()) {
logger.debug("Pong received on " + this.getConnectionId());
}
}
else if (this.shook) {
return super.onMessage(message);
}
else {
try {
doHandshake(payload, message.getHeaders());
this.shook = true;
WebSocketEvent event = new WebSocketEvent(this.getTheConnection(),
WebSocketEventType.HANDSHAKE_COMPLETE, state.getPath(), state.getQueryString());
this.getTheConnection().publishEvent(event);
}
catch (Exception e) {
throw new MessageHandlingException(message, "Handshake failed", e);
}
}
return true;
}
private WebSocketState getState(Object object) {
Object stateKey = null;
stateKey = this.getTheConnection().getDeserializerStateKey();
Assert.notNull(stateKey, "StateKey must not be null:" + object);
WebSocketState state = (WebSocketState) this.getRequiredDeserializer().getState(stateKey);
Assert.notNull(state, "State must not be null:" + object);
return state;
}
private void protocolViolation(Message<?> message) {
if (logger.isDebugEnabled()) {
logger.debug("Protocol violation - closing; " + message);
}
WebSocketFrame frame = (WebSocketFrame) message.getPayload();
String error = "Protocol Error" + frame.getPayload() == null ? "" : (":" + frame.getPayload());
WebSocketFrame close = new WebSocketFrame(WebSocketFrame.TYPE_CLOSE, error);
close.setStatus(frame.getType() == WebSocketFrame.TYPE_INVALID_UTF8 ? (short) 1007 : (short) 1002);
try {
Object stateKey = this.getTheConnection().getDeserializerStateKey();
if (stateKey != null) {
WebSocketState webSocketState = (WebSocketState) this.getRequiredDeserializer().getState(stateKey);
if (webSocketState != null) {
webSocketState.setCloseInitiated(true);
}
this.send(MessageBuilder.withPayload(close)
.copyHeaders(message.getHeaders())
.build());
}
}
catch (Exception e) {
throw new MessageHandlingException(message, "Send failed", e);
}
}
@Override
public void close() {
connections.remove(this.getTheConnection());
Object stateKey = this.getTheConnection().getDeserializerStateKey();
if (stateKey != null) {
this.getRequiredDeserializer().removeState(stateKey);
}
super.close();
}
private void doHandshake(WebSocketFrame frame, MessageHeaders messageHeaders) throws Exception {
try {
WebSocketFrame handshake = this.getRequiredDeserializer().generateHandshake(frame);
this.send(MessageBuilder.withPayload(handshake)
.copyHeaders(messageHeaders)
.build());
}
catch (WebSocketUpgradeException e) {
this.send(MessageBuilder
.withPayload(
new WebSocketFrame(WebSocketFrame.TYPE_DATA, "HTTP/1.1 " +
e.getMessage() + e.getHeaders()))
.copyHeaders(messageHeaders)
.build());
this.close();
}
}
private WebSocketSerializer getRequiredDeserializer() {
Deserializer<?> deserializer = this.getDeserializer();
Assert.state(deserializer instanceof WebSocketSerializer,
"Deserializer must be a WebSocketSerializer");
return (WebSocketSerializer) deserializer;
}
public Map<String, String> getAdditionalHeaders() {
Map<String, String> headers = new HashMap<String, String>();
WebSocketState state = this.getState(this.getConnectionId());
if (state.getPath() != null) {
headers.put(WebSocketHeaders.PATH, state.getPath());
}
if (state.getQueryString() != null) {
headers.put(WebSocketHeaders.QUERY_STRING, state.getQueryString());
}
return headers;
}
@Override
public void setTheConnection(TcpConnectionSupport theConnection) {
connections.put(theConnection, this);
super.setTheConnection(theConnection);
}
}
}

View File

@@ -1,42 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
/**
* @author Gary Russell
* @since 3.0
*
*/
@SuppressWarnings("serial")
public class WebSocketUpgradeException extends RuntimeException {
private final String headers;
public WebSocketUpgradeException(String message) {
super(message + "\r\n");
this.headers = "\r\n";
}
public WebSocketUpgradeException(String message, String headers) {
super(message + "\r\n");
this.headers = headers + "\r\n";
}
protected String getHeaders() {
return headers;
}
}

View File

@@ -1,88 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.serializer;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.InputStream;
import org.junit.Test;
import org.springframework.integration.x.ip.serializer.AbstractHttpSwitchingDeserializer.BasicState;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class AbstractHttpSwitchingDeserializerTests {
private final AbstractHttpSwitchingDeserializer deserializer = new AbstractHttpSwitchingDeserializer() {
@Override
public DataFrame deserialize(InputStream inputStream) throws IOException {
return checkStreaming(inputStream).get(0);
}
};
@Test
public void testPathNoQuery() throws Exception {
String simplePath = "GET /foo HTTP/1.1\r\n\r\n";
InputStream stream = new ByteArrayInputStream(simplePath.getBytes());
DataFrame frame = deserializer.deserialize(stream);
assertEquals(DataFrame.TYPE_HEADERS, frame.getType());
BasicState state = deserializer.getState(stream);
assertNotNull(state);
assertEquals("/foo", state.getPath());
assertNull(state.getQueryString());
}
@Test
public void testPathAndQuery() throws Exception {
String simplePath = "GET /foo?bar HTTP/1.1\r\n\r\n";
InputStream stream = new ByteArrayInputStream(simplePath.getBytes());
DataFrame frame = deserializer.deserialize(stream);
assertEquals(DataFrame.TYPE_HEADERS, frame.getType());
BasicState state = deserializer.getState(stream);
assertNotNull(state);
assertEquals("/foo", state.getPath());
assertEquals("bar", state.getQueryString());
}
@Test
public void testPathEmptyQuery() throws Exception {
String simplePath = "GET /foo? HTTP/1.1\r\n\r\n";
InputStream stream = new ByteArrayInputStream(simplePath.getBytes());
DataFrame frame = deserializer.deserialize(stream);
assertEquals(DataFrame.TYPE_HEADERS, frame.getType());
BasicState state = deserializer.getState(stream);
assertNotNull(state);
assertEquals("/foo", state.getPath());
assertEquals("", state.getQueryString());
}
}

View File

@@ -1,101 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:int-ip="http://www.springframework.org/schema/integration/ip"
xmlns:int-stream="http://www.springframework.org/schema/integration/stream"
xmlns:int="http://www.springframework.org/schema/integration"
xsi:schemaLocation="http://www.springframework.org/schema/integration/ip http://www.springframework.org/schema/integration/ip/spring-integration-ip.xsd
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
http://www.springframework.org/schema/integration/stream http://www.springframework.org/schema/integration/stream/spring-integration-stream.xsd
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
<!-- ws://localhost:18080 -->
<int-ip:tcp-connection-factory id="ws"
type="server" port="18080"
using-nio="false"
so-timeout="600000"
interceptor-factory-chain="interceptors"
serializer="wsSerializer"
deserializer="wsSerializer" />
<int-ip:tcp-inbound-channel-adapter connection-factory="ws" channel="echoChannel" />
<!-- wss://localhost:28080 -->
<!-- We MUST use apply-sequence with NIO so we can resequence frames -->
<int-ip:tcp-connection-factory id="wss"
type="server" port="28080"
using-nio="true" apply-sequence="true"
so-timeout="600000"
ssl-context-support="sslCtxSup"
interceptor-factory-chain="interceptors"
serializer="wsSerializer"
deserializer="wsSerializer" />
<bean id="sslCtxSup" class="org.springframework.integration.ip.tcp.connection.DefaultTcpSSLContextSupport">
<constructor-arg value="key.store"/>
<constructor-arg value="trust.store"/>
<constructor-arg value="secret"/>
<constructor-arg value="secret"/>
</bean>
<int-ip:tcp-inbound-channel-adapter connection-factory="wss" channel="addSSLHeader" />
<int:header-enricher input-channel="addSSLHeader" output-channel="echoChannel">
<int:header name="SSL" value="true"/>
</int:header-enricher>
<!-- End SSL config -->
<bean id="wsSerializer" class="org.springframework.integration.x.ip.websocket.WebSocketSerializer">
<property name="server" value="true" />
<property name="validateUtf8" value="true" />
</bean>
<bean id="interceptors" class="org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactoryChain">
<property name="interceptors">
<bean class="org.springframework.integration.x.ip.websocket.WebSocketTcpConnectionInterceptorFactory" />
</property>
</bean>
<bean id="service" class="org.springframework.integration.x.ip.websocket.WebSocketServerTests$DemoService" />
<int:chain input-channel="echoChannel" output-channel="toChooseBrowser">
<!-- payload is a DataFrame -->
<int:service-activator expression="payload" />
</int:chain>
<int:header-value-router input-channel="toChooseBrowser"
header-name="SSL" default-output-channel="toBrowser">
<int:mapping value="false" channel="toBrowser"/>
<int:mapping value="true" channel="toBrowserSSL"/>
</int:header-value-router>
<!-- outbound send to websocket (ws:) -->
<int:channel id="toBrowser" />
<int-ip:tcp-outbound-channel-adapter connection-factory="ws" channel="toBrowser" />
<!-- outbound send to websocket (wss:) -->
<int:channel id="toBrowserSSL" />
<int-ip:tcp-outbound-channel-adapter connection-factory="wss" channel="toBrowserSSL" />
<int:channel id="stdOut" />
<int-stream:stdout-channel-adapter append-newline="true" channel="stdOut" />
<!-- Remove state after error on write -->
<int:channel id="removeSocket" />
<int:chain input-channel="removeSocket">
<int:transformer expression="payload.failedMessage.headers['ip_connectionId']" />
<int:service-activator ref="service" method="remove" />
</int:chain>
</beans>

View File

@@ -1,34 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import org.springframework.context.support.ClassPathXmlApplicationContext;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class AutobahnTests {
public static void main(String[] args) throws Exception {
new ClassPathXmlApplicationContext("Autobahn-context.xml", AutobahnTests.class);
System.out.println("Hit Enter To Terminate...");
System.in.read();
System.exit(0);
}
}

View File

@@ -1,71 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:int-ip="http://www.springframework.org/schema/integration/ip"
xmlns:int-stream="http://www.springframework.org/schema/integration/stream"
xmlns:int="http://www.springframework.org/schema/integration"
xsi:schemaLocation="http://www.springframework.org/schema/integration/ip http://www.springframework.org/schema/integration/ip/spring-integration-ip.xsd
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
http://www.springframework.org/schema/integration/stream http://www.springframework.org/schema/integration/stream/spring-integration-stream.xsd
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
<int-ip:tcp-connection-factory id="ws"
type="server" port="18080"
using-nio="true"
so-timeout="60000"
interceptor-factory-chain="interceptors"
serializer="wsSerializer"
deserializer="wsSerializer"
mapper="mapper" />
<bean id="mapper" class="org.springframework.integration.x.ip.websocket.WebSocketMessageMapper">
<constructor-arg ref="webSocketInterceptorFactory"/>
</bean>
<bean id="wsSerializer" class="org.springframework.integration.x.ip.websocket.WebSocketSerializer">
<property name="server" value="true" />
</bean>
<bean id="interceptors" class="org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactoryChain">
<property name="interceptors" ref="webSocketInterceptorFactory" />
</bean>
<bean id="webSocketInterceptorFactory"
class="org.springframework.integration.x.ip.websocket.WebSocketTcpConnectionInterceptorFactory" />
<int-ip:tcp-inbound-channel-adapter connection-factory="ws" channel="startStopChannel" />
<bean id="service" class="org.springframework.integration.x.ip.websocket.WebSocketServerTests$DemoService" />
<int:chain input-channel="startStopChannel">
<!-- payload is a SockJsFrame with the data in its payload property -->
<int:transformer expression="payload.payload" />
<int:service-activator method="startStop" ref="service" />
</int:chain>
<!-- outbound send to websocket -->
<int:inbound-channel-adapter ref="service" method="getNext" channel="out">
<int:poller fixed-delay="1000" error-channel="removeSocket" />
</int:inbound-channel-adapter>
<int:splitter input-channel="out" output-channel="toBrowser" />
<int:channel id="toBrowser" />
<int-ip:tcp-outbound-channel-adapter connection-factory="ws" channel="toBrowser" />
<int:channel id="stdOut" />
<int-stream:stdout-channel-adapter append-newline="true" channel="stdOut" />
<!-- Remove state after error on write -->
<int:channel id="removeSocket" />
<int:chain input-channel="removeSocket">
<int:transformer expression="payload.failedMessage.headers['ip_connectionId']" />
<int:service-activator ref="service" method="remove" />
</int:chain>
</beans>

View File

@@ -1,139 +0,0 @@
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.x.ip.websocket;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.BeansException;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
import org.springframework.context.ApplicationListener;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.Message;
import org.springframework.integration.annotation.Header;
import org.springframework.integration.annotation.Headers;
import org.springframework.integration.ip.IpHeaders;
import org.springframework.integration.ip.tcp.connection.AbstractConnectionFactory;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.x.ip.websocket.WebSocketEvent.WebSocketEventType;
/**
* @author Gary Russell
* @since 3.0
*
*/
public class WebSocketServerTests{
public static void main(String[] args) throws Exception {
new ClassPathXmlApplicationContext(WebSocketServerTests.class.getSimpleName() + "-context.xml", WebSocketServerTests.class);
System.out.println("Hit Enter To Terminate...");
System.in.read();
System.exit(0);
}
public static class DemoService implements ApplicationListener<WebSocketEvent>, ApplicationContextAware {
private static final Log logger = LogFactory.getLog(DemoService.class);
private final Map<String, AtomicInteger> clients = new HashMap<String, AtomicInteger>();
private final Map<String, AtomicInteger> paused = new HashMap<String, AtomicInteger>();
private volatile ApplicationContext applicationContext;
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
this.applicationContext = applicationContext;
}
public void startStop(String command, @Header(IpHeaders.CONNECTION_ID) String connectionId,
@Headers Map<String, ?> headers) {
if (headers != null) {
logger.info("Received '"
+ command
+ "' from '"
+ connectionId
+ "' path:"
+ headers.get(WebSocketHeaders.PATH) + " query-string:" + headers.get(WebSocketHeaders.QUERY_STRING));
}
if ("stop".equalsIgnoreCase(command)) {
AtomicInteger clientInt = clients.remove(connectionId);
if (clientInt != null) {
paused.put(connectionId, clientInt);
}
logger.info("Connection " + connectionId + " stopped");
}
else if ("start".equalsIgnoreCase(command)) {
AtomicInteger clientInt = paused.remove(connectionId);
clientInt = clientInt == null ? new AtomicInteger() : clientInt;
clients.put(connectionId, clientInt);
logger.info("Connection " + connectionId + " (re)started");
}
else {
logger.info("Unexpected command: " + command);
}
}
public List<Message<?>> getNext() {
List<Message<?>> messages = new ArrayList<Message<?>>();
for (Entry<String, AtomicInteger> entry : clients.entrySet()) {
Message<String> message = MessageBuilder.withPayload(Integer.toString(entry.getValue().incrementAndGet()))
.setHeader(IpHeaders.CONNECTION_ID, entry.getKey())
.build();
messages.add(message);
logger.warn("Sending " + message.getPayload() + " to connection " + entry.getKey());
}
if (messages.size() == 0) {
return null;
}
else {
return messages;
}
}
public void remove(String connetionId) {
logger.warn("Error on write; removing " + connetionId);
clients.remove(connetionId);
}
@Override
public void onApplicationEvent(WebSocketEvent event) {
logger.info(event);
if (WebSocketEventType.HANDSHAKE_COMPLETE.equals(event.getType())) {
startStop("start", event.getConnectionId(), null);
try {
logger.info("Handshake complete for new connection on port "
+ this.applicationContext.getBean(event.getConnectionFactoryName(),
AbstractConnectionFactory.class).getPort());
}
catch (Exception e) {
logger.error("Failed to get port", e);
}
}
else if (WebSocketEventType.WEBSOCKET_CLOSED.equals(event.getType())) {
clients.remove(event.getConnectionId());
}
}
}
}

View File

@@ -1,83 +0,0 @@
<!--
/*
* Copyright 2002-2013 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-->
<!DOCTYPE HTML PUBLIC "-//W3C//DTD HTML 4.01 Transitional//EN"
"http://www.w3.org/TR/html4/loose.dtd">
<html>
<head><title>Spring Integration Web Socket Test</title></head>
<body>
<script type="text/javascript">
var webSocket;
if (window.WebSocket) {
webSocket = new WebSocket("ws://localhost:18080/myapp?foo");
webSocket.onmessage = function(event) {
document.getElementById("theField").value = event.data;
}
webSocket.onopen = function(event) {
document.getElementById("status").value = "WebSocket opened";
};
webSocket.onclose = function(event) {
document.getElementById("status").value = "WebSocket closed";
alert("WebSocket closed.");
};
} else {
alert("Your browser does not support Websockets. (Use Chrome)");
}
function sendToSI(message) {
if (!window.WebSocket) {
return;
}
if (webSocket.readyState == WebSocket.OPEN) {
webSocket.send(message);
} else {
alert("The socket is not open.");
}
document.getElementById("status").value = "sent " + message;
if (message == "start") {
document.getElementById("message").value = "stop";
}
else if (message == "stop") {
document.getElementById("message").value = "start";
}
}
function doClose(message) {
if (!window.WebSocket) {
return;
}
webSocket.close();
}
</script>
<center>
<form onsubmit="return false;" id="form">
<br/><br/>
<input type="text" id="message" name="message" value="stop" />
<br/><br/>
<input type="button" value="Send To Spring Integration" onclick="sendToSI(this.form.message.value)" align="center"/>
<br/><br/>
<input type="button" value="Close WebSocket" onclick="doClose()" align="center"/>
</form>
<br/>
<textarea rows="1" cols="5" id="theField" readonly="readonly"></textarea>
<br/><br/>
<textarea rows="1" id="status" readonly="readonly"></textarea>
</center>
</body>
</html>

View File

@@ -1,40 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE log4j:configuration SYSTEM "log4j.dtd">
<log4j:configuration xmlns:log4j="http://jakarta.apache.org/log4j/">
<!-- Appenders -->
<appender name="console" class="org.apache.log4j.ConsoleAppender">
<param name="Target" value="System.out" />
<layout class="org.apache.log4j.PatternLayout">
<param name="ConversionPattern" value="%d{HH:mm:ss.SSS} %-5p [%t][%c] %m%n" />
</layout>
</appender>
<!-- Loggers -->
<logger name="org.springframework">
<level value="warn" />
</logger>
<logger name="org.springframework.integration">
<level value="info" />
</logger>
<logger name="org.springframework.integration.x">
<level value="info" />
</logger>
<logger name="org.springframework.integration.endpoint.SourcePollingChannelAdapter">
<level value="info" />
</logger>
<logger name="org.springframework.integration.samples">
<level value="debug" />
</logger>
<!-- Root Logger -->
<root>
<priority value="warn" />
<appender-ref ref="console" />
</root>
</log4j:configuration>