package io.vertx.ext.web.handler.sockjs.impl;
import io.vertx.core.AsyncResult;
import io.vertx.core.Handler;
import io.vertx.core.Vertx;
import io.vertx.core.buffer.Buffer;
import io.vertx.core.http.HttpServerRequest;
import io.vertx.core.logging.Logger;
import io.vertx.core.logging.LoggerFactory;
import io.vertx.core.shareddata.LocalMap;
import io.vertx.ext.web.Router;
import io.vertx.ext.web.RoutingContext;
import io.vertx.ext.web.handler.sockjs.SockJSHandlerOptions;
import io.vertx.ext.web.handler.sockjs.SockJSSocket;
import static io.vertx.core.buffer.Buffer.buffer;
class EventSourceTransport extends BaseTransport {
private static final Logger log = LoggerFactory.getLogger(EventSourceTransport.class);
EventSourceTransport(Vertx vertx, Router router, LocalMap<String, SockJSSession> sessions, SockJSHandlerOptions options,
Handler<SockJSSocket> sockHandler) {
super(vertx, sessions, options);
String eventSourceRE = COMMON_PATH_ELEMENT_RE + "eventsource";
router.getWithRegex(eventSourceRE).handler(rc -> {
if (log.isTraceEnabled()) log.trace("EventSource transport, get: " + rc.request().uri());
String sessionID = rc.request().getParam("param0");
SockJSSession session = getSession(rc, options.getSessionTimeout(), options.getHeartbeatInterval(), sessionID, sockHandler);
HttpServerRequest req = rc.request();
session.register(req, new EventSourceListener(options.getMaxBytesStreaming(), rc, session));
});
}
private class EventSourceListener extends BaseListener {
final int maxBytesStreaming;
boolean headersWritten;
int bytesSent;
boolean closed;
EventSourceListener(int maxBytesStreaming, RoutingContext rc, SockJSSession session) {
super(rc, session);
this.maxBytesStreaming = maxBytesStreaming;
addCloseHandler(rc.response(), session);
}
@Override
public void sendFrame(String body, Handler<AsyncResult<Void>> handler) {
if (log.isTraceEnabled()) log.trace("EventSource, sending frame");
if (!headersWritten) {
rc.response().putHeader("Content-Type", "text/event-stream");
setNoCacheHeaders(rc);
setJSESSIONID(options, rc);
rc.response().setChunked(true).write("\r\n");
headersWritten = true;
}
String sb = "data: " +
body +
"\r\n\r\n";
Buffer buff = buffer(sb);
rc.response().write(buff, handler);
bytesSent += buff.length();
if (bytesSent >= maxBytesStreaming) {
if (log.isTraceEnabled()) log.trace("More than maxBytes sent so closing connection");
close();
}
}
public void close() {
if (!closed) {
try {
session.resetListener();
rc.response().end();
rc.response().close();
} catch (IllegalStateException e) {
}
closed = true;
}
}
}
}