[77] | 1 | package com.framsticks.communication; |
---|
| 2 | |
---|
| 3 | import com.framsticks.communication.queries.ApplicationRequest; |
---|
| 4 | import com.framsticks.communication.queries.RegistrationRequest; |
---|
| 5 | import com.framsticks.communication.queries.UseRequest; |
---|
| 6 | import com.framsticks.communication.queries.VersionRequest; |
---|
| 7 | import com.framsticks.communication.util.LoggingStateCallback; |
---|
| 8 | import com.framsticks.params.ListSource; |
---|
| 9 | import com.framsticks.util.*; |
---|
[84] | 10 | import com.framsticks.util.dispatching.AtOnceDispatcher; |
---|
| 11 | import com.framsticks.util.dispatching.Dispatcher; |
---|
| 12 | import com.framsticks.util.dispatching.Dispatching; |
---|
| 13 | import com.framsticks.util.lang.Pair; |
---|
| 14 | import com.framsticks.util.lang.Strings; |
---|
[77] | 15 | import org.apache.log4j.Logger; |
---|
| 16 | |
---|
| 17 | import java.io.IOException; |
---|
| 18 | import java.net.Socket; |
---|
| 19 | import java.net.SocketException; |
---|
| 20 | import java.util.*; |
---|
| 21 | import java.util.regex.Matcher; |
---|
| 22 | import java.util.regex.Pattern; |
---|
[85] | 23 | import com.framsticks.util.dispatching.RunAt; |
---|
[77] | 24 | |
---|
| 25 | /** |
---|
| 26 | * @author Piotr Sniegowski |
---|
| 27 | */ |
---|
| 28 | public class ClientConnection extends Connection { |
---|
| 29 | |
---|
[84] | 30 | private final static Logger log = Logger.getLogger(ClientConnection.class); |
---|
[77] | 31 | |
---|
[85] | 32 | protected final Map<String, Subscription<?>> subscriptions = new HashMap<>(); |
---|
[77] | 33 | |
---|
[88] | 34 | /** |
---|
| 35 | * @param connectedFunctor the connectedFunctor to set |
---|
| 36 | */ |
---|
| 37 | public void setConnectedFunctor(StateFunctor connectedFunctor) { |
---|
| 38 | this.connectedFunctor = connectedFunctor; |
---|
[84] | 39 | } |
---|
[77] | 40 | |
---|
[88] | 41 | |
---|
| 42 | protected StateFunctor connectedFunctor; |
---|
| 43 | |
---|
| 44 | @Override |
---|
| 45 | protected void joinableStart() { |
---|
[84] | 46 | try { |
---|
[90] | 47 | log.debug("connecting to " + address); |
---|
[77] | 48 | |
---|
[84] | 49 | socket = new Socket(hostName, port); |
---|
[77] | 50 | |
---|
[84] | 51 | socket.setSoTimeout(500); |
---|
[77] | 52 | |
---|
[90] | 53 | log.debug("connected to " + hostName + ":" + port); |
---|
[84] | 54 | connected = true; |
---|
| 55 | runThreads(); |
---|
[77] | 56 | |
---|
[84] | 57 | connectedFunctor.call(null); |
---|
| 58 | } catch (SocketException e) { |
---|
| 59 | log.error("failed to connect: " + e); |
---|
| 60 | connectedFunctor.call(e); |
---|
| 61 | } catch (IOException e) { |
---|
| 62 | log.error("buffer creation failure"); |
---|
| 63 | connectedFunctor.call(e); |
---|
[88] | 64 | // close(); |
---|
[84] | 65 | } |
---|
| 66 | } |
---|
[77] | 67 | |
---|
[88] | 68 | |
---|
[84] | 69 | private static abstract class InboundMessage { |
---|
| 70 | protected String currentFilePath; |
---|
| 71 | protected List<String> currentFileContent; |
---|
| 72 | protected final List<File> files = new ArrayList<File>(); |
---|
[77] | 73 | |
---|
[84] | 74 | public abstract void eof(); |
---|
[77] | 75 | |
---|
[84] | 76 | protected void initCurrentFile(String path) { |
---|
| 77 | currentFileContent = new LinkedList<String>(); |
---|
| 78 | currentFilePath = path; |
---|
| 79 | } |
---|
| 80 | protected void finishCurrentFile() { |
---|
| 81 | if (currentFileContent == null) { |
---|
| 82 | return; |
---|
| 83 | } |
---|
| 84 | files.add(new File(currentFilePath, new ListSource(currentFileContent))); |
---|
| 85 | currentFilePath = null; |
---|
| 86 | currentFileContent= null; |
---|
| 87 | } |
---|
[77] | 88 | |
---|
[84] | 89 | public abstract void startFile(String path); |
---|
[77] | 90 | |
---|
[84] | 91 | public void addLine(String line) { |
---|
| 92 | assert line != null; |
---|
| 93 | assert currentFileContent != null; |
---|
| 94 | currentFileContent.add(line.substring(0, line.length() - 1)); |
---|
| 95 | } |
---|
[77] | 96 | |
---|
[84] | 97 | public List<File> getFiles() { |
---|
| 98 | return files; |
---|
| 99 | } |
---|
| 100 | } |
---|
[77] | 101 | |
---|
[84] | 102 | private static class EventFire extends InboundMessage { |
---|
[85] | 103 | public final Subscription<?> subscription; |
---|
[77] | 104 | |
---|
[85] | 105 | private EventFire(Subscription<?> subscription) { |
---|
[84] | 106 | this.subscription = subscription; |
---|
| 107 | } |
---|
[77] | 108 | |
---|
[84] | 109 | public void startFile(String path) { |
---|
| 110 | assert path == null; |
---|
| 111 | initCurrentFile(null); |
---|
| 112 | } |
---|
[77] | 113 | |
---|
[84] | 114 | @Override |
---|
| 115 | public void eof() { |
---|
| 116 | finishCurrentFile(); |
---|
[85] | 117 | |
---|
| 118 | subscription.dispatchCall(getFiles()); |
---|
[84] | 119 | } |
---|
| 120 | } |
---|
[77] | 121 | |
---|
[85] | 122 | private static class SentQuery<C> extends InboundMessage { |
---|
[84] | 123 | Request request; |
---|
[85] | 124 | ResponseCallback<? extends C> callback; |
---|
| 125 | Dispatcher<C> dispatcher; |
---|
[77] | 126 | |
---|
[84] | 127 | public void startFile(String path) { |
---|
| 128 | finishCurrentFile(); |
---|
| 129 | if (path == null) { |
---|
| 130 | assert request instanceof ApplicationRequest; |
---|
| 131 | path = ((ApplicationRequest)request).getPath(); |
---|
| 132 | } |
---|
| 133 | initCurrentFile(path); |
---|
| 134 | } |
---|
[77] | 135 | |
---|
[84] | 136 | public void eof() { |
---|
| 137 | finishCurrentFile(); |
---|
| 138 | //no-operation |
---|
| 139 | } |
---|
[77] | 140 | |
---|
[84] | 141 | @Override |
---|
| 142 | public String toString() { |
---|
| 143 | return request.toString(); |
---|
| 144 | } |
---|
[85] | 145 | |
---|
| 146 | public void dispatchResponseProcess(final Response response) { |
---|
[90] | 147 | Dispatching.dispatchIfNotActive(dispatcher, new RunAt<C>() { |
---|
[85] | 148 | @Override |
---|
| 149 | public void run() { |
---|
| 150 | callback.process(response); |
---|
| 151 | } |
---|
| 152 | }); |
---|
| 153 | } |
---|
[84] | 154 | } |
---|
[85] | 155 | private Map<Integer, SentQuery<?>> queryMap = new HashMap<>(); |
---|
[77] | 156 | |
---|
| 157 | |
---|
[84] | 158 | protected final String hostName; |
---|
| 159 | protected final int port; |
---|
[77] | 160 | |
---|
[84] | 161 | private static Pattern addressPattern = Pattern.compile("^([^:]*)(:([0-9]+))?$"); |
---|
[77] | 162 | |
---|
[84] | 163 | public ClientConnection(String address) { |
---|
[88] | 164 | super(address); |
---|
[84] | 165 | Matcher matcher = addressPattern.matcher(address); |
---|
| 166 | if (!matcher.matches()) { |
---|
| 167 | log.fatal("invalid address: " + address); |
---|
| 168 | hostName = null; |
---|
| 169 | port = 0; |
---|
| 170 | return; |
---|
| 171 | } |
---|
| 172 | hostName = matcher.group(1); |
---|
| 173 | port = matcher.group(3) != null ? Integer.parseInt(matcher.group(3)) : 9009; |
---|
| 174 | } |
---|
[77] | 175 | |
---|
[85] | 176 | private SentQuery<?> currentlySentQuery; |
---|
[77] | 177 | |
---|
[85] | 178 | |
---|
| 179 | public <C extends Connection> void send(Request request, ResponseCallback<C> callback) { |
---|
| 180 | //TODO RunAt |
---|
| 181 | send(request, AtOnceDispatcher.getInstance(), callback); |
---|
[84] | 182 | } |
---|
[77] | 183 | |
---|
[85] | 184 | public <C> void send(Request request, Dispatcher<C> dispatcher, ResponseCallback<? extends C> callback) { |
---|
[77] | 185 | |
---|
[84] | 186 | if (!isConnected()) { |
---|
| 187 | log.fatal("not connected"); |
---|
| 188 | return; |
---|
| 189 | } |
---|
[85] | 190 | final SentQuery<C> sentQuery = new SentQuery<C>(); |
---|
[84] | 191 | sentQuery.request = request; |
---|
| 192 | sentQuery.callback = callback; |
---|
| 193 | sentQuery.dispatcher = dispatcher; |
---|
[77] | 194 | |
---|
[90] | 195 | senderThread.dispatch(new RunAt<Connection>(){ |
---|
[84] | 196 | @Override |
---|
| 197 | public void run() { |
---|
[85] | 198 | Integer id; |
---|
[84] | 199 | synchronized (ClientConnection.this) { |
---|
[77] | 200 | |
---|
[84] | 201 | while (!(requestIdEnabled || currentlySentQuery == null)) { |
---|
| 202 | try { |
---|
| 203 | ClientConnection.this.wait(); |
---|
| 204 | } catch (InterruptedException ignored) { |
---|
| 205 | break; |
---|
| 206 | } |
---|
| 207 | } |
---|
[85] | 208 | if (requestIdEnabled) { |
---|
| 209 | queryMap.put(nextQueryId, sentQuery); |
---|
| 210 | id = nextQueryId++; |
---|
| 211 | } else { |
---|
| 212 | currentlySentQuery = sentQuery; |
---|
| 213 | id = null; |
---|
| 214 | } |
---|
[84] | 215 | } |
---|
| 216 | String command = sentQuery.request.getCommand(); |
---|
| 217 | StringBuilder message = new StringBuilder(); |
---|
| 218 | message.append(command); |
---|
| 219 | if (id != null) { |
---|
| 220 | message.append(" ").append(id); |
---|
| 221 | } |
---|
| 222 | sentQuery.request.construct(message); |
---|
| 223 | String out = message.toString(); |
---|
[77] | 224 | |
---|
[84] | 225 | output.println(out); |
---|
| 226 | log.debug("sending query: " + out); |
---|
[77] | 227 | |
---|
[84] | 228 | } |
---|
| 229 | }); |
---|
| 230 | /* |
---|
| 231 | synchronized (this) { |
---|
| 232 | log.debug("queueing query: " + query); |
---|
| 233 | queryQueue.offer(sentQuery); |
---|
| 234 | notifyAll(); |
---|
| 235 | } |
---|
| 236 | */ |
---|
| 237 | } |
---|
[77] | 238 | |
---|
| 239 | |
---|
[84] | 240 | @Override |
---|
| 241 | public String toString() { |
---|
[88] | 242 | return "client connection " + address; |
---|
[84] | 243 | } |
---|
[77] | 244 | |
---|
[85] | 245 | public <C> void subscribe(final String path, final Dispatcher<C> dispatcher, final SubscriptionCallback<? extends C> callback) { |
---|
| 246 | send(new RegistrationRequest().path(path), new ResponseCallback<Connection>() { |
---|
[84] | 247 | @Override |
---|
| 248 | public void process(Response response) { |
---|
| 249 | if (!response.getOk()) { |
---|
| 250 | log.error("failed to register on event: " + path); |
---|
| 251 | callback.subscribed(null); |
---|
| 252 | return; |
---|
| 253 | } |
---|
| 254 | assert response.getFiles().isEmpty(); |
---|
[85] | 255 | Subscription<C> subscription = new Subscription<C>(ClientConnection.this, path, response.getComment(), dispatcher); |
---|
[84] | 256 | log.debug("registered on event: " + subscription); |
---|
| 257 | synchronized (subscriptions) { |
---|
| 258 | subscriptions.put(subscription.getRegisteredPath(), subscription); |
---|
| 259 | } |
---|
| 260 | subscription.setEventCallback(callback.subscribed(subscription)); |
---|
| 261 | if (subscription.getEventCallback() == null) { |
---|
| 262 | log.info("subscription for " + path + " aborted"); |
---|
[85] | 263 | subscription.unsubscribe(new LoggingStateCallback<C>(log, "abort subscription")); |
---|
[84] | 264 | } |
---|
| 265 | } |
---|
| 266 | }); |
---|
| 267 | } |
---|
[77] | 268 | |
---|
[84] | 269 | public void negotiateProtocolVersion(StateFunctor stateFunctor) { |
---|
| 270 | protocolVersion = -1; |
---|
| 271 | sendQueryVersion(1, stateFunctor); |
---|
| 272 | } |
---|
[77] | 273 | |
---|
[84] | 274 | public void sendQueryVersion(final int version, final StateFunctor stateFunctor) { |
---|
[85] | 275 | send(new VersionRequest().version(version), new StateCallback<Connection>() { |
---|
[84] | 276 | @Override |
---|
| 277 | public void call(Exception e) { |
---|
| 278 | if (e != null) { |
---|
| 279 | log.fatal("failed to upgrade protocol to version: " + version); |
---|
| 280 | return; |
---|
| 281 | } |
---|
| 282 | protocolVersion = version; |
---|
| 283 | if (version < 4) { |
---|
| 284 | /** it is an implicit loop here*/ |
---|
| 285 | sendQueryVersion(version + 1, stateFunctor); |
---|
| 286 | return; |
---|
| 287 | } |
---|
[85] | 288 | send(new UseRequest().feature("request_id"), new StateCallback<Connection>() { |
---|
[84] | 289 | @Override |
---|
| 290 | public void call(Exception e) { |
---|
| 291 | requestIdEnabled = e == null; |
---|
| 292 | /* |
---|
| 293 | synchronized (ClientConnection.this) { |
---|
| 294 | ClientConnection.this.notifyAll(); |
---|
| 295 | } |
---|
| 296 | */ |
---|
| 297 | if (!requestIdEnabled) { |
---|
| 298 | log.fatal("protocol negotiation failed"); |
---|
| 299 | stateFunctor.call(new Exception("protocol negotiation failed", e)); |
---|
| 300 | return; |
---|
| 301 | } |
---|
| 302 | stateFunctor.call(null); |
---|
| 303 | } |
---|
| 304 | }); |
---|
[77] | 305 | |
---|
[84] | 306 | } |
---|
| 307 | }); |
---|
| 308 | } |
---|
[77] | 309 | |
---|
| 310 | |
---|
[85] | 311 | private synchronized SentQuery<?> fetchQuery(Integer id, boolean remove) { |
---|
[84] | 312 | if (id == null) { |
---|
| 313 | if (requestIdEnabled) { |
---|
| 314 | return null; |
---|
| 315 | } |
---|
[85] | 316 | SentQuery<?> result = currentlySentQuery; |
---|
[84] | 317 | if (remove) { |
---|
| 318 | currentlySentQuery = null; |
---|
| 319 | notifyAll(); |
---|
| 320 | } |
---|
| 321 | return result; |
---|
| 322 | } |
---|
| 323 | if (queryMap.containsKey(id)) { |
---|
[85] | 324 | SentQuery<?> result = queryMap.get(id); |
---|
[84] | 325 | if (remove) { |
---|
| 326 | queryMap.remove(id); |
---|
| 327 | } |
---|
| 328 | return result; |
---|
| 329 | } |
---|
| 330 | return null; |
---|
| 331 | } |
---|
[77] | 332 | |
---|
[84] | 333 | private int nextQueryId = 0; |
---|
[77] | 334 | |
---|
[88] | 335 | protected void processMessage(InboundMessage inboundMessage) { |
---|
[84] | 336 | if (inboundMessage == null) { |
---|
| 337 | log.error("failed to use any inbound message"); |
---|
| 338 | return; |
---|
| 339 | } |
---|
[77] | 340 | |
---|
[84] | 341 | String line; |
---|
| 342 | while (!(line = getLine()).startsWith("eof")) { |
---|
| 343 | // log.debug("line: " + line); |
---|
| 344 | inboundMessage.addLine(line); |
---|
| 345 | } |
---|
| 346 | inboundMessage.eof(); |
---|
| 347 | } |
---|
[77] | 348 | |
---|
[88] | 349 | protected void processEvent(String rest) { |
---|
[84] | 350 | Matcher matcher = eventPattern.matcher(rest); |
---|
| 351 | if (!matcher.matches()) { |
---|
| 352 | log.error("invalid event line: " + rest); |
---|
| 353 | return; |
---|
| 354 | } |
---|
[85] | 355 | Subscription<?> subscription = subscriptions.get(matcher.group(1)); |
---|
[84] | 356 | if (subscription == null) { |
---|
| 357 | log.error("non subscribed event: " + matcher.group(1)); |
---|
| 358 | return; |
---|
| 359 | } |
---|
| 360 | EventFire event = new EventFire(subscription); |
---|
| 361 | event.startFile(null); |
---|
| 362 | processMessage(event); |
---|
| 363 | } |
---|
[77] | 364 | |
---|
| 365 | |
---|
[88] | 366 | protected void processMessageStartingWith(String line) { |
---|
[84] | 367 | Pair<String, String> command = Strings.splitIntoPair(line, ' ', "\n"); |
---|
| 368 | if (command.first.equals("event")) { |
---|
| 369 | processEvent(command.second); |
---|
| 370 | return; |
---|
| 371 | } |
---|
| 372 | Pair<Integer, String> rest = parseRest(command.second); |
---|
[77] | 373 | |
---|
[84] | 374 | if (command.first.equals("file")) { |
---|
[85] | 375 | SentQuery<?> sentQuery = fetchQuery(rest.first, false); |
---|
[84] | 376 | sentQuery.startFile(rest.second); |
---|
| 377 | processMessage(sentQuery); |
---|
| 378 | return; |
---|
| 379 | } |
---|
[77] | 380 | |
---|
[85] | 381 | SentQuery<?> sentQuery = fetchQuery(rest.first, true); |
---|
[84] | 382 | if (sentQuery == null) { |
---|
| 383 | return; |
---|
| 384 | } |
---|
| 385 | log.debug("parsing response for request " + sentQuery); |
---|
[77] | 386 | |
---|
[85] | 387 | sentQuery.dispatchResponseProcess(new Response(command.first.equals("ok"), rest.second, sentQuery.getFiles())); |
---|
[84] | 388 | } |
---|
[77] | 389 | |
---|
[84] | 390 | @Override |
---|
[88] | 391 | protected void receiverThreadRoutine() { |
---|
[84] | 392 | while (connected) { |
---|
[88] | 393 | try { |
---|
| 394 | processMessageStartingWith(getLine()); |
---|
| 395 | } catch (Exception e) { |
---|
| 396 | break; |
---|
| 397 | } |
---|
[84] | 398 | } |
---|
| 399 | } |
---|
[77] | 400 | |
---|
[88] | 401 | |
---|
[77] | 402 | } |
---|