refaturing message system, step 3: making use of queue
Signed-off-by: Stephan Richter <s.richter@srsoftware.de>
This commit is contained in:
@@ -1,15 +1,18 @@
|
||||
/* © SRSoftware 2025 */
|
||||
package de.srsoftware.umbrella.message;
|
||||
|
||||
import de.srsoftware.umbrella.core.model.Envelope;
|
||||
import de.srsoftware.umbrella.core.model.Message;
|
||||
import de.srsoftware.umbrella.core.model.UmbrellaUser;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
public interface MessageQueue<T> {
|
||||
public Stream<T> getEnvelopes();
|
||||
public Stream<T> getEnvelopesFor(UmbrellaUser user);
|
||||
public Optional<T> getEnvelope(int hash);
|
||||
public void markRead(int hash, UmbrellaUser user);
|
||||
public void push(T message);
|
||||
public interface MessageQueue<T extends Message<?>> {
|
||||
public Stream<Envelope<T>> getEnvelopes();
|
||||
public List<Envelope<T>> getEnvelopesFor(UmbrellaUser user);
|
||||
public Optional<Envelope<T>> getEnvelope(int hash);
|
||||
Stream<UmbrellaUser> getReceivers();
|
||||
public Optional<Envelope<T>> markRead(int hash, UmbrellaUser user);
|
||||
public void push(Envelope<T> message);
|
||||
}
|
||||
|
||||
@@ -40,7 +40,6 @@ import jakarta.mail.internet.MimeMultipart;
|
||||
import jakarta.mail.util.ByteArrayDataSource;
|
||||
import java.io.IOException;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import org.json.JSONArray;
|
||||
import org.json.JSONObject;
|
||||
|
||||
@@ -77,15 +76,18 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
}
|
||||
private final String from,host,user,pass;
|
||||
private final int port;
|
||||
private final SqliteMessageDb db;
|
||||
private final MessageDb db;
|
||||
private final MessageQueue<TranslatedMessage> queue;
|
||||
private Session session;
|
||||
private final List<Envelope<?>> queue = new CopyOnWriteArrayList<>();
|
||||
private String debugAddress;
|
||||
private final String debugAddress;
|
||||
private final HashMap<Receiver,List<Exception>> exceptions = new HashMap<>();
|
||||
|
||||
public MessageSystem(Configuration config) throws UmbrellaException {
|
||||
var dbFile = config.get(CONFIG_DB).orElseThrow(() -> missingConfig(CONFIG_DB));
|
||||
db = new SqliteMessageDb(connect(dbFile));
|
||||
var sqlite = new SqliteMessageDb(connect(dbFile));
|
||||
db = sqlite;
|
||||
queue = sqlite;
|
||||
|
||||
debugAddress = config.get(DEBUG_ADDREESS).map(Object::toString).orElse(null);
|
||||
port = config.get(CONFIG_SMTP_PORT,587);
|
||||
host = config.get(CONFIG_SMTP_HOST).map(Object::toString).orElseThrow(() -> new RuntimeException("umbrella.modules.message.smtp.host not configured!"));
|
||||
@@ -151,12 +153,7 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
}
|
||||
|
||||
private boolean getMessage(HttpExchange ex, UmbrellaUser user, int hash) throws IOException {
|
||||
var envel = db.getEnvelope(hash).filter(env -> env.isFor(user));
|
||||
if (envel.isPresent()) return sendMessage(ex, user, envel.get());
|
||||
var envelope = queue.stream()
|
||||
.filter(msg -> msg.isFor(user))
|
||||
.filter(msg -> msg.hashCode() == hash)
|
||||
.findFirst();
|
||||
var envelope = queue.getEnvelope(hash).filter(env -> env.isFor(user));
|
||||
if (envelope.isPresent()) return sendMessage(ex, user, envelope.get());
|
||||
return notFound(ex);
|
||||
}
|
||||
@@ -166,9 +163,7 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
}
|
||||
|
||||
private boolean listMessages(HttpExchange ex, UmbrellaUser user) throws IOException {
|
||||
var msgs = db.getEnvelopesFor(user).map(e -> summary(e, user.language())).toList();
|
||||
if (!msgs.isEmpty()) return sendContent(ex,msgs);
|
||||
var messages = queue.stream().filter(e -> e.isFor(user)).map(e -> summary(e, user.language())).toList();
|
||||
var messages = queue.getEnvelopesFor(user).stream().map(e -> summary(e, user.language()));
|
||||
return sendContent(ex,messages);
|
||||
}
|
||||
|
||||
@@ -184,7 +179,7 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
|
||||
private boolean patchSettings(HttpExchange ex, UmbrellaUser user) throws IOException {
|
||||
var json = json(ex);
|
||||
Settings settings = null;
|
||||
Settings settings;
|
||||
if (json.has(INSTANTLY) && json.get(INSTANTLY) instanceof Boolean b && b){
|
||||
settings = new Instantly();
|
||||
} else {
|
||||
@@ -199,49 +194,37 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
private boolean markRead(HttpExchange ex, UmbrellaUser user, String path) {
|
||||
try {
|
||||
var hash = Integer.parseInt(path);
|
||||
db.markRead(hash, user);
|
||||
var envelope = queue.stream().filter(env -> env.hashCode() == hash).findFirst().orElse(null);
|
||||
if (envelope != null){
|
||||
envelope.receivers().remove(user);
|
||||
return sendMessage(ex,user,envelope);
|
||||
}
|
||||
var envelope = queue.markRead(hash, user);
|
||||
if (envelope.isPresent()) return sendMessage(ex,user,envelope.get());
|
||||
return notFound(ex);
|
||||
} catch (NumberFormatException | IOException e) {
|
||||
throw invalidField(HASH,LONG);
|
||||
}
|
||||
}
|
||||
|
||||
private boolean sendAt(UmbrellaUser user, Integer scheduledHour){
|
||||
try {
|
||||
return db.getSettings(user).sendAt(scheduledHour);
|
||||
} catch (UmbrellaException ignored) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
private synchronized void processMessages(Integer scheduledHour) {
|
||||
LOG.log(INFO,"Running {0}…",scheduledHour == null ? "instantly" : "scheduled at "+scheduledHour);
|
||||
var queue = new ArrayList<>(this.queue);
|
||||
var dueRecipients = new ArrayList<User>();
|
||||
List<User> recipients = queue.stream().map(Envelope::receivers).flatMap(Set::stream).filter(Objects::nonNull).distinct().toList();
|
||||
|
||||
{ // for known users: get notification preferences, fallback to _immediately_ for unknown users
|
||||
for (User recv : recipients) {
|
||||
if (recv instanceof UmbrellaUser uu) {
|
||||
try {
|
||||
if (!db.getSettings(uu).sendAt(scheduledHour)) continue;
|
||||
} catch (UmbrellaException ignored) {}
|
||||
}
|
||||
dueRecipients.add(recv);
|
||||
}
|
||||
}
|
||||
var dueRecipients = queue.getReceivers().filter(uu -> sendAt(uu,scheduledHour)).toList();
|
||||
|
||||
var date = new Date();
|
||||
|
||||
for (var receiver : dueRecipients){
|
||||
var combined = new CombinedMessage(t("Collected messages"),receiver);
|
||||
var envelopes = queue.stream().filter(env -> env.isFor(receiver)).toList();
|
||||
for (var envelope : envelopes) combined.merge(envelope.message());
|
||||
|
||||
try {
|
||||
var envelopes = queue.getEnvelopesFor(receiver);
|
||||
envelopes.stream().map(Envelope::message).forEach(combined::merge);
|
||||
send(combined,date);
|
||||
for (var envelope : envelopes){
|
||||
var audience = envelope.receivers();
|
||||
audience.remove(receiver);
|
||||
if (audience.isEmpty()) queue.remove(envelope);
|
||||
}
|
||||
envelopes.forEach(env -> queue.markRead(env.hashCode(),receiver));
|
||||
} catch (Exception ex){
|
||||
LOG.log(WARNING,"Failed to deliver mail ({0}) to {1}.",combined.subject(),receiver,ex);
|
||||
for (var message : combined.messages()) exceptions.computeIfAbsent(new Receiver(receiver,message), k -> new ArrayList<>()).add(ex);
|
||||
@@ -262,17 +245,6 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
));
|
||||
}
|
||||
|
||||
private static JSONObject summary(Envelope<?> envelope, String lang) {
|
||||
var message = envelope.message();
|
||||
if (message instanceof TranslatableMessage tm) message = tm.translate(lang);
|
||||
|
||||
var sender = message.sender().name();
|
||||
var subject = message.subject();
|
||||
var time = envelope.time().format(TIME_FORMATTER);
|
||||
var hash = envelope.hashCode();
|
||||
return new JSONObject(Map.of(SENDER,sender,SUBJECT,subject,TIMESTAMP,time,HASH,hash));
|
||||
}
|
||||
|
||||
|
||||
private void send(CombinedMessage message, Date date) throws MessagingException {
|
||||
var receiver = message.receiver();
|
||||
@@ -311,9 +283,29 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public void send(Envelope<?> envelope) {
|
||||
queue.add(envelope);
|
||||
db.push(envelope);
|
||||
switch (envelope.message()){
|
||||
case TranslatedMessage ignored:
|
||||
queue.push((Envelope<TranslatedMessage>) envelope);
|
||||
break;
|
||||
case TranslatableMessage tm:
|
||||
Map<String,Envelope<TranslatedMessage>> map = new HashMap<>();
|
||||
for (var receiver : envelope.receivers()){
|
||||
var lang = receiver.language();
|
||||
var env = map.get(lang);
|
||||
if (env == null){
|
||||
TranslatedMessage translated = tm.translate(lang);
|
||||
env = new Envelope<>(translated,new HashSet<>());
|
||||
map.put(lang,env);
|
||||
}
|
||||
env.receivers().add(receiver);
|
||||
}
|
||||
map.values().forEach(queue::push);
|
||||
default:
|
||||
return;
|
||||
}
|
||||
|
||||
new Thread(() -> processMessages(null)).start();
|
||||
}
|
||||
|
||||
@@ -330,8 +322,15 @@ public class MessageSystem extends BaseHandler implements PostBox, EventListener
|
||||
return session;
|
||||
}
|
||||
|
||||
public void setDebugAddress(String newVal) {
|
||||
this.debugAddress = newVal;
|
||||
private static JSONObject summary(Envelope<?> envelope, String lang) {
|
||||
var message = envelope.message();
|
||||
if (message instanceof TranslatableMessage tm) message = tm.translate(lang);
|
||||
|
||||
var sender = message.sender().name();
|
||||
var subject = message.subject();
|
||||
var time = envelope.time().format(TIME_FORMATTER);
|
||||
var hash = envelope.hashCode();
|
||||
return new JSONObject(Map.of(SENDER,sender,SUBJECT,subject,TIMESTAMP,time,HASH,hash));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -28,7 +28,7 @@ import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
public class SqliteMessageDb implements MessageDb, MessageQueue<Envelope<TranslatedMessage>> {
|
||||
public class SqliteMessageDb implements MessageDb, MessageQueue<TranslatedMessage> {
|
||||
private static final System.Logger LOG = System.getLogger(SqliteMessageDb.class.getSimpleName());
|
||||
private final Connection db;
|
||||
private static final String DB_VERSION = "message_db_version";
|
||||
@@ -97,10 +97,15 @@ CREATE TABLE IF NOT EXISTS {0} ( {1} VARCHAR(255) PRIMARY KEY, {2} VARCHAR(255)
|
||||
}
|
||||
|
||||
@Override
|
||||
public Stream<Envelope<TranslatedMessage>> getEnvelopesFor(UmbrellaUser user) {
|
||||
public List<Envelope<TranslatedMessage>> getEnvelopesFor(UmbrellaUser user) {
|
||||
throw new UmbrellaException(HTTP_SERVER_ERROR,"{class}.getEnvelopesFor({user}) not implemented!","class",getClass().getSimpleName(),"user",user.name()); // TODO
|
||||
}
|
||||
|
||||
@Override
|
||||
public Stream<UmbrellaUser> getReceivers() {
|
||||
return Stream.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Settings getSettings(UmbrellaUser user) throws UmbrellaException {
|
||||
try {
|
||||
@@ -120,7 +125,7 @@ CREATE TABLE IF NOT EXISTS {0} ( {1} VARCHAR(255) PRIMARY KEY, {2} VARCHAR(255)
|
||||
}
|
||||
|
||||
@Override
|
||||
public void markRead(int hash, UmbrellaUser user) {
|
||||
public Optional<Envelope<TranslatedMessage>> markRead(int hash, UmbrellaUser user) {
|
||||
throw new UmbrellaException(HTTP_SERVER_ERROR,"{class}.markRead(hash, user) not implemented!","class",getClass().getSimpleName()); // TODO
|
||||
// TODO: throw exception if message not found!
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user