bug6059: implemented getResourceData(...); rudimentary search for DAWs

This commit is contained in:
Axel Uhl
2024-12-20 17:41:17 +01:00
parent b155970ee3
commit 5415f52ae6
4 changed files with 64 additions and 8 deletions
@@ -5,6 +5,7 @@ import java.nio.channels.ServerSocketChannel;
import com.igtimi.IgtimiData.DataMsg;
import com.igtimi.IgtimiData.DataPoint;
import com.igtimi.IgtimiStream.Msg;
import com.sap.sailing.domain.igtimiadapter.BulkFixReceiver;
import com.sap.sailing.domain.igtimiadapter.DataAccessWindow;
import com.sap.sailing.domain.igtimiadapter.Device;
@@ -17,6 +18,7 @@ import com.sap.sailing.domain.igtimiadapter.server.replication.ReplicableRiotSer
import com.sap.sailing.domain.igtimiadapter.server.replication.RiotReplicationOperation;
import com.sap.sailing.domain.igtimiadapter.server.riot.impl.RiotServerImpl;
import com.sap.sailing.domain.igtimiadapter.shared.IgtimiWindReceiver;
import com.sap.sse.common.TimeRange;
import com.sap.sse.replication.Replicable;
/**
@@ -96,7 +98,12 @@ public interface RiotServer extends Replicable<ReplicableRiotServer, RiotReplica
DataAccessWindow getDataAccessWindowById(long id);
Iterable<DataAccessWindow> getDataAccessWindows(Iterable<String> serialNumbers, TimeRange timeRange);
void addDataAccessWindow(DataAccessWindow daw);
void removeDataAccessWindow(long dawId);
// TODO clarify what this means in the presence of replication, when run on a Replica
Iterable<Msg> getMessages(String deviceSerialNumber, TimeRange timeRange);
}
@@ -14,6 +14,7 @@ import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.nio.channels.SocketChannel;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
@@ -34,6 +35,8 @@ import com.sap.sailing.domain.igtimiadapter.server.riot.RiotConnection;
import com.sap.sailing.domain.igtimiadapter.server.riot.RiotMessageListener;
import com.sap.sailing.domain.igtimiadapter.server.riot.RiotServer;
import com.sap.sse.common.Duration;
import com.sap.sse.common.TimeRange;
import com.sap.sse.common.Util;
import com.sap.sse.replication.interfaces.impl.AbstractReplicableWithObjectInputStream;
import com.sap.sse.shared.util.Wait;
import com.sap.sse.util.ObjectInputStreamResolvingAgainstCache;
@@ -50,6 +53,7 @@ public class RiotServerImpl extends AbstractReplicableWithObjectInputStream<Repl
private final Thread communicatorThread;
private boolean running;
private final MongoObjectFactory mongoObjectFactory;
private final DomainObjectFactory domainObjectFactory;
/**
* The active connections managed by this server. Heartbeats will be sent into the channels found in the key set on
@@ -81,6 +85,7 @@ public class RiotServerImpl extends AbstractReplicableWithObjectInputStream<Repl
this.dataAccessWindows = new ConcurrentHashMap<>();
this.devices = new ConcurrentHashMap<>();
this.connections = new ConcurrentHashMap<>();
this.domainObjectFactory = domainObjectFactory;
this.mongoObjectFactory = mongoObjectFactory;
for (final Device device : domainObjectFactory.getDevices()) {
devices.put(device.getId(), device);
@@ -263,6 +268,20 @@ public class RiotServerImpl extends AbstractReplicableWithObjectInputStream<Repl
return dataAccessWindows.get(id);
}
@Override
public Iterable<DataAccessWindow> getDataAccessWindows(Iterable<String> deviceSerialNumbers, TimeRange timeRange) {
final Set<DataAccessWindow> result = new HashSet<>();
final Set<String> deviceSerialNumbersAsSet = new HashSet<>();
Util.addAll(deviceSerialNumbers, deviceSerialNumbersAsSet);
// TODO provide a more efficient implementation if this turns out to become a performance bottleneck, e.g., by keeping the DataAccessWindows in a time-sorted TreeSet
for (final DataAccessWindow daw : getDataAccessWindows()) {
if (deviceSerialNumbersAsSet.contains(daw.getDeviceSerialNumber()) && timeRange.intersects(daw.getTimeRange())) {
result.add(daw);
}
}
return result;
}
@Override
public void addDataAccessWindow(DataAccessWindow daw) {
apply(s->s.internalAddDataAccessWindow(daw));
@@ -323,4 +342,9 @@ public class RiotServerImpl extends AbstractReplicableWithObjectInputStream<Repl
objectOutputStream.writeObject(resources);
objectOutputStream.writeObject(dataAccessWindows);
}
@Override
public Iterable<Msg> getMessages(String deviceSerialNumber, TimeRange timeRange) {
return domainObjectFactory.getMessages(deviceSerialNumber, timeRange);
}
}