2020-10-13 01:31:32 +02:00
|
|
|
package it.tdlight.common;
|
|
|
|
|
2021-01-24 18:11:25 +01:00
|
|
|
import it.tdlight.jni.TdApi;
|
2021-03-06 17:27:08 +01:00
|
|
|
import it.tdlight.jni.TdApi.Object;
|
|
|
|
import java.util.ArrayList;
|
|
|
|
import java.util.List;
|
|
|
|
import java.util.Map;
|
|
|
|
import java.util.Map.Entry;
|
|
|
|
import java.util.Objects;
|
2020-10-13 15:12:13 +02:00
|
|
|
import java.util.concurrent.ConcurrentHashMap;
|
2020-10-13 01:31:32 +02:00
|
|
|
import java.util.concurrent.atomic.AtomicLong;
|
|
|
|
import java.util.concurrent.atomic.AtomicReference;
|
2021-01-24 18:11:25 +01:00
|
|
|
import org.slf4j.Logger;
|
|
|
|
import org.slf4j.LoggerFactory;
|
2020-10-13 01:31:32 +02:00
|
|
|
|
|
|
|
public class InternalClientManager implements AutoCloseable {
|
|
|
|
|
2021-01-24 18:11:25 +01:00
|
|
|
private static final Logger logger = LoggerFactory.getLogger(InternalClientManager.class);
|
2020-10-13 01:31:32 +02:00
|
|
|
private static final AtomicReference<InternalClientManager> INSTANCE = new AtomicReference<>(null);
|
|
|
|
|
|
|
|
private final String implementationName;
|
|
|
|
private final ResponseReceiver responseReceiver = new ResponseReceiver(this::handleClientEvents);
|
2020-10-13 15:12:13 +02:00
|
|
|
private final ConcurrentHashMap<Integer, ClientEventsHandler> registeredClientEventHandlers = new ConcurrentHashMap<>();
|
2020-10-13 01:31:32 +02:00
|
|
|
private final AtomicLong currentQueryId = new AtomicLong();
|
|
|
|
|
|
|
|
private InternalClientManager(String implementationName) {
|
2020-10-13 02:02:24 +02:00
|
|
|
try {
|
|
|
|
Init.start();
|
|
|
|
} catch (Throwable ex) {
|
|
|
|
ex.printStackTrace();
|
|
|
|
System.exit(1);
|
|
|
|
}
|
2020-10-13 01:31:32 +02:00
|
|
|
this.implementationName = implementationName;
|
|
|
|
}
|
|
|
|
|
|
|
|
public static InternalClientManager get(String implementationName) {
|
|
|
|
return INSTANCE.updateAndGet(val -> val == null ? new InternalClientManager(implementationName) : val);
|
|
|
|
}
|
|
|
|
|
2021-01-24 18:11:25 +01:00
|
|
|
private void handleClientEvents(int clientId, boolean isClosed, long[] clientEventIds, TdApi.Object[] clientEvents) {
|
2020-10-13 15:12:13 +02:00
|
|
|
ClientEventsHandler handler = registeredClientEventHandlers.get(clientId);
|
|
|
|
|
2020-10-13 01:31:32 +02:00
|
|
|
if (handler != null) {
|
|
|
|
handler.handleEvents(isClosed, clientEventIds, clientEvents);
|
|
|
|
} else {
|
2021-03-06 17:27:08 +01:00
|
|
|
java.util.List<Entry<Long, TdApi.Object>> droppedEvents = getEffectivelyDroppedEvents(clientEventIds, clientEvents);
|
|
|
|
|
|
|
|
if (!droppedEvents.isEmpty()) {
|
|
|
|
logger.error("Unknown client id \"{}\"! {} events have been dropped!", clientId, droppedEvents.size());
|
|
|
|
for (Entry<Long, Object> droppedEvent : droppedEvents) {
|
|
|
|
logger.error("The following event, with id \"{}\", has been dropped: {}",
|
|
|
|
droppedEvent.getKey(),
|
|
|
|
droppedEvent.getValue());
|
|
|
|
}
|
2021-01-24 18:02:00 +01:00
|
|
|
}
|
2020-10-13 01:31:32 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
if (isClosed) {
|
2021-03-31 04:33:39 +02:00
|
|
|
logger.trace("Removing Client {} from event handlers", clientId);
|
2020-10-13 15:12:13 +02:00
|
|
|
registeredClientEventHandlers.remove(clientId);
|
2021-03-31 04:33:39 +02:00
|
|
|
logger.trace("Removed Client {} from event handlers", clientId);
|
2020-10-13 01:31:32 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-03-06 17:27:08 +01:00
|
|
|
/**
|
|
|
|
* Get only events that have been dropped, ignoring synthetic errors related to the closure of a client
|
|
|
|
*/
|
|
|
|
private List<Entry<Long, TdApi.Object>> getEffectivelyDroppedEvents(long[] clientEventIds, TdApi.Object[] clientEvents) {
|
|
|
|
java.util.List<Entry<Long, TdApi.Object>> droppedEvents = new ArrayList<>(clientEvents.length);
|
|
|
|
for (int i = 0; i < clientEvents.length; i++) {
|
|
|
|
long id = clientEventIds[i];
|
|
|
|
TdApi.Object event = clientEvents[i];
|
|
|
|
boolean mustPrintError = true;
|
|
|
|
if (event instanceof TdApi.Error) {
|
|
|
|
var errorEvent = (TdApi.Error) event;
|
|
|
|
if (Objects.equals("Request aborted", errorEvent.message)) {
|
|
|
|
mustPrintError = false;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if (mustPrintError) {
|
|
|
|
droppedEvents.add(Map.entry(id, event));
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return droppedEvents;
|
|
|
|
}
|
|
|
|
|
2021-02-13 17:41:54 +01:00
|
|
|
public void registerClient(int clientId, ClientEventsHandler internalClient) {
|
2020-10-13 18:33:06 +02:00
|
|
|
boolean replaced = registeredClientEventHandlers.put(clientId, internalClient) != null;
|
|
|
|
if (replaced) {
|
|
|
|
throw new IllegalStateException("Client " + clientId + " already registered");
|
|
|
|
}
|
2021-02-25 23:36:49 +01:00
|
|
|
responseReceiver.registerClient(clientId);
|
2020-10-13 01:31:32 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
public String getImplementationName() {
|
|
|
|
return implementationName;
|
|
|
|
}
|
|
|
|
|
|
|
|
public long getNextQueryId() {
|
2020-10-13 04:10:20 +02:00
|
|
|
return currentQueryId.updateAndGet(value -> (value >= Long.MAX_VALUE ? 0 : value) + 1);
|
2020-10-13 01:31:32 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
@Override
|
|
|
|
public void close() throws InterruptedException {
|
|
|
|
responseReceiver.close();
|
|
|
|
}
|
|
|
|
}
|