Files
Companion/src/main/java/io/dietz/ed/companion/component/ZeroMqHandler.java
T
2021-05-07 22:28:29 +00:00

108 lines
3.8 KiB
Java

package io.dietz.ed.companion.component;
import java.io.IOException;
import java.time.LocalDateTime;
import java.util.zip.DataFormatException;
import java.util.zip.Inflater;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.zeromq.SocketType;
import org.zeromq.ZContext;
import org.zeromq.ZMQ;
import io.dietz.ed.companion.domain.MarketWithRoot;
import io.dietz.ed.companion.models.eddb.Market;
import io.dietz.ed.companion.repositories.MarketRepository;
import io.dietz.ed.companion.repositories.StationRepository;
public class ZeroMqHandler implements Runnable {
private ZContext context;
@Value("${eddn.filter:https://eddn.edcd.io/schemas/commodity/3}")
private String filterSchema;
@Value("${eddn.host:eddn.edcd.io}")
private String host;
@Value("${eddn.port:9500}")
private String port;
@Autowired
private ObjectMapper objectMapper;
@Autowired
private MarketRepository marketRepository;
@Autowired
private StationRepository stationRepository;
public void msg(String msg) {
System.out.println(msg);
}
@Override
public void run() {
pump();
context.close();
}
public synchronized void pump() {
Inflater inflater = new Inflater();
context = new ZContext();
ZMQ.Socket socket = context.createSocket(SocketType.SUB);
socket.subscribe("");
socket.setReceiveTimeOut(30000);
socket.connect("tcp://" + host + ":" + port);
msg("EDDN Relay connected");
ZMQ.Poller poller = context.createPoller(2);
poller.register(socket, ZMQ.Poller.POLLIN);
byte[] output = new byte[256 * 1024];
while (true) {
int poll = poller.poll(-1);
if (poll == ZMQ.Poller.POLLIN) {
if (poller.pollin(0)) {
byte[] recv = socket.recv(ZMQ.NOBLOCK);
if (recv.length > 0) {
// decompress
inflater.reset();
inflater.setInput(recv);
try {
int outlen = inflater.inflate(output);
String outputString = new String(output, 0, outlen, "UTF-8");
// outputString contains a json message
if (outputString.contains(filterSchema)) {
//msg(outputString);
MarketWithRoot marketContainer = objectMapper.readValue(outputString, MarketWithRoot.class);
msg(marketContainer.getMarket().toString());
var test = marketRepository.findById(marketContainer.getMarket().marketId).orElseGet(() -> {
msg("Market created!");
Market market = new Market(marketContainer.getMarket().marketId);
marketRepository.saveAndFlush(market);
return market;
});
if(test.getStation() != null) {
test.getStation().setUpdatedAt(LocalDateTime.now());
stationRepository.saveAndFlush(test.getStation());
msg(test.getStation().getName() + " updated!");
}
}
} catch (DataFormatException | IOException e) {
e.printStackTrace();
}
}
}
}
}
}
}