privatefinal Log log = LogFactory.getLog(WsSession.class); // must not be static privatestaticfinal StringManager sm = StringManager.getManager(WsSession.class);
// An ellipsis is a single character that looks like three periods in a row // and is used to indicate a continuation. privatestaticfinalbyte[] ELLIPSIS_BYTES = "\u2026".getBytes(StandardCharsets.UTF_8); // An ellipsis is three bytes in UTF-8 privatestaticfinalint ELLIPSIS_BYTES_LEN = ELLIPSIS_BYTES.length;
static { // Use fake end point and path. They are never used, they just need to // be sufficient to pass the validation tests.
ServerEndpointConfig.Builder builder = ServerEndpointConfig.Builder.create(Object.class, "/");
ServerEndpointConfig sec = builder.build();
SEC_CONFIGURATOR_USES_IMPL_DEFAULT = sec.getConfigurator().getClass()
.equals(DefaultServerEndpointConfigurator.class);
}
// Message handlers that require decoders may map to text messages, // binary messages, both or neither.
// The frame processing code expects binary message handlers to // accept ByteBuffer
// Use the POJO message handler wrappers as they are designed to wrap // arbitrary objects with MessageHandlers and can wrap MessageHandlers // just as easily.
if (!removed) { // ISE for now. Could swallow this silently / log this if the ISE // becomes a problem thrownew IllegalStateException(sm.getString("wsSession.removeHandlerFailed", listener));
}
}
@Override public String getProtocolVersion() {
checkState(); return Constants.WS_VERSION_HEADER_VALUE;
}
@Override public String getNegotiatedSubprotocol() {
checkState(); return subProtocol;
}
@Override public List<Extension> getNegotiatedExtensions() {
checkState(); return negotiatedExtensions;
}
// Send the close message to the remote endpoint.
sendCloseMessage(closeReasonMessage);
fireEndpointOnClose(closeReasonLocal); if (!state.compareAndSet(State.OUTPUT_CLOSING, State.OUTPUT_CLOSED) || closeSocket) { /* *Aclosemessagewasreceivedinanotherthreadorthisishandlinganerrorcondition.Eitherway,no *furtherclosemessageisexpectedtobereceived.Markthesessionasfullyclosed...
*/
state.set(State.CLOSED); // ... and close the network connection.
wsRemoteEndpoint.close();
}
// Fail any uncompleted messages.
IOException ioe = new IOException(sm.getString("wsSession.messageFailed"));
SendResult sr = new SendResult(ioe); for (FutureToSendHandler f2sh : futures.keySet()) {
f2sh.onResult(sr);
}
}
/** *Calledwhenaclosemessageisreceived.Shouldonlyeverhappenonce.Alsocalledafteraprotocolerrorwhen *theProtocolHandlerneedstoforcetheclosingoftheconnection. * *@paramcloseReasonThereasoncontainedwithinthereceivedclosemessage.
*/ publicvoid onClose(CloseReason closeReason) { if (state.compareAndSet(State.OPEN, State.CLOSING)) { // Standard close.
// Fire the onError event Thread t = Thread.currentThread();
ClassLoader cl = t.getContextClassLoader();
t.setContextClassLoader(applicationClassLoader); try {
localEndpoint.onError(this, throwable);
} finally {
t.setContextClassLoader(cl);
}
}
privatevoid sendCloseMessage(CloseReason closeReason) { // 125 is maximum size for the payload of a control message
ByteBuffer msg = ByteBuffer.allocate(125);
CloseCode closeCode = closeReason.getCloseCode(); // CLOSED_ABNORMALLY should not be put on the wire if (closeCode == CloseCodes.CLOSED_ABNORMALLY) { // PROTOCOL_ERROR is probably better than GOING_AWAY here
msg.putShort((short) CloseCodes.PROTOCOL_ERROR.getCode());
} else {
msg.putShort((short) closeCode.getCode());
}
String reason = closeReason.getReasonPhrase(); if (reason != null && reason.length() > 0) {
appendCloseReasonWithTruncation(msg, reason);
}
msg.flip(); try {
wsRemoteEndpoint.sendMessageBlock(Constants.OPCODE_CLOSE, msg, true);
} catch (IOException | IllegalStateException e) { // Failed to send close message. Close the socket and let the caller // deal with the Exception if (log.isDebugEnabled()) {
log.debug(sm.getString("wsSession.sendCloseFail", id), e);
}
wsRemoteEndpoint.close(); // Failure to send a close message is not unexpected in the case of // an abnormal closure (usually triggered by a failure to read/write // from/to the client. In this case do not trigger the endpoint's // error handling if (closeCode != CloseCodes.CLOSED_ABNORMALLY) {
localEndpoint.onError(this, e);
}
} finally {
webSocketContainer.unregisterSession(getSessionMapKey(), this);
}
}
/** *Useprotectedsounittestscanaccessthismethoddirectly. * *@parammsgThemessage *@paramreasonThereason
*/ protectedstaticvoid appendCloseReasonWithTruncation(ByteBuffer msg, String reason) { // Once the close code has been added there are a maximum of 123 bytes // left for the reason phrase. If it is truncated then care needs to be // taken to ensure the bytes are not truncated in the middle of a // multi-byte UTF-8 character. byte[] reasonBytes = reason.getBytes(StandardCharsets.UTF_8);
if (reasonBytes.length <= 123) { // No need to truncate
msg.put(reasonBytes);
} else { // Need to truncate int remaining = 123 - ELLIPSIS_BYTES_LEN; int pos = 0; byte[] bytesNext = reason.substring(pos, pos + 1).getBytes(StandardCharsets.UTF_8); while (remaining >= bytesNext.length) {
msg.put(bytesNext);
remaining -= bytesNext.length;
pos++;
bytesNext = reason.substring(pos, pos + 1).getBytes(StandardCharsets.UTF_8);
}
msg.put(ELLIPSIS_BYTES);
}
}
/** *Makethesessionawareofa{@linkFutureToSendHandler}thatwillneedtobeforciblyclosedifthesession *closesbeforethe{@linkFutureToSendHandler}completes. * *@paramf2shThehandler
*/ protectedvoid registerFuture(FutureToSendHandler f2sh) { // Ideally, this code should sync on stateLock so that the correct // action is taken based on the current state of the connection. // However, a sync on stateLock can't be used here as it will create the // possibility of a dead-lock. See BZ 61183. // Therefore, a slightly less efficient approach is used.
// Always register the future.
futures.put(f2sh, f2sh);
if (isOpen()) { // The session is open. The future has been registered with the open // session. Normal processing continues. return;
}
// The session is closing / closed. The future may or may not have been registered // in time for it to be processed during session closure.
if (f2sh.isDone()) { // The future has completed. It is not known if the future was // completed normally by the I/O layer or in error by doClose(). It // doesn't matter which. There is nothing more to do here. return;
}
// The session is closing / closed. The Future had not completed when last checked. // There is a small timing window that means the Future may have been // completed since the last check. There is also the possibility that // the Future was not registered in time to be cleaned up during session // close. // Attempt to complete the Future with an error result as this ensures // that the Future completes and any client code waiting on it does not // hang. It is slightly inefficient since the Future may have been // completed in another thread or another thread may be about to // complete the Future but knowing if this is the case requires the sync // on stateLock (see above). // Note: If multiple attempts are made to complete the Future, the // second and subsequent attempts are ignored.
IOException ioe = new IOException(sm.getString("wsSession.messageFailed"));
SendResult sr = new SendResult(ioe);
f2sh.onResult(sr);
}
protectedvoid checkExpiration() { // Local copies to ensure consistent behaviour during method execution long timeout = maxIdleTimeout; long timeoutRead = getMaxIdleTimeoutRead(); long timeoutWrite = getMaxIdleTimeoutWrite();
long currentTime = System.currentTimeMillis();
String key = null;
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.