001// --------------------------------------------------------------------------------
002// Copyright 2002-2026 Echo Three, LLC
003//
004// Licensed under the Apache License, Version 2.0 (the "License");
005// you may not use this file except in compliance with the License.
006// You may obtain a copy of the License at
007//
008//     http://www.apache.org/licenses/LICENSE-2.0
009//
010// Unless required by applicable law or agreed to in writing, software
011// distributed under the License is distributed on an "AS IS" BASIS,
012// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
013// See the License for the specific language governing permissions and
014// limitations under the License.
015// --------------------------------------------------------------------------------
016
017package com.echothree.control.user.core.server.command;
018
019import com.echothree.model.control.party.common.PartyTypes;
020import com.echothree.model.data.core.server.entity.EventSubscriber;
021import com.echothree.util.common.command.BaseResult;
022import com.echothree.util.server.control.BaseSimpleCommand;
023import com.echothree.util.server.control.CommandSecurityDefinition;
024import com.echothree.util.server.control.PartyTypeDefinition;
025import java.util.HashSet;
026import java.util.List;
027import java.util.Set;
028import javax.enterprise.context.Dependent;
029import javax.inject.Inject;
030
031@Dependent
032public class ProcessQueuedEventsCommand
033        extends BaseSimpleCommand {
034
035    private final static CommandSecurityDefinition COMMAND_SECURITY_DEFINITION;
036    
037    static {
038        COMMAND_SECURITY_DEFINITION = new CommandSecurityDefinition(List.of(
039                new PartyTypeDefinition(PartyTypes.UTILITY.name(), null)
040        ));
041    }
042
043    
044    /** Creates a new instance of ProcessQueuedEventsCommand */
045    public ProcessQueuedEventsCommand() {
046        super(COMMAND_SECURITY_DEFINITION, false);
047    }
048
049    @Override
050    protected BaseResult execute() {
051        var remainingTime = 50 * 1000; // 50 seconds
052        var queuedEvents = eventControl.getQueuedEventsForUpdate();
053
054        for(var queuedEvent : queuedEvents) {
055            var startTime = System.currentTimeMillis();
056            Set<EventSubscriber> eventSubscribers = new HashSet<>();
057            var event = queuedEvent.getEvent();
058
059            if(event != null) {
060                // TODO: this should not be necessary, bug 444
061                var eventType = event.getEventType();
062                var entityInstance = event.getEntityInstance();
063                var entityType = entityInstance.getEntityType();
064                var eventSubscriberEventTypes = eventControl.getEventSubscriberEventTypes(eventType);
065                var eventSubscriberEntityTypes = eventControl.getEventSubscriberEntityTypes(entityType, eventType);
066                var eventSubscriberEntityInstances = eventControl.getEventSubscriberEntityInstances(entityInstance, eventType);
067
068                eventSubscriberEventTypes.stream().map((eventSubscriberEventType) -> eventSubscriberEventType.getEventSubscriber()).filter((eventSubscriber) -> !eventSubscribers.contains(eventSubscriber)).map((eventSubscriber) -> {
069                    eventControl.createQueuedSubscriberEvent(eventSubscriber, event);
070                    return eventSubscriber;
071                }).forEach((eventSubscriber) -> {
072                    eventSubscribers.add(eventSubscriber);
073                });
074
075                eventSubscriberEntityTypes.stream().map((eventSubscriberEntityType) -> eventSubscriberEntityType.getEventSubscriber()).filter((eventSubscriber) -> !eventSubscribers.contains(eventSubscriber)).map((eventSubscriber) -> {
076                    eventControl.createQueuedSubscriberEvent(eventSubscriber, event);
077                    return eventSubscriber;
078                }).forEach((eventSubscriber) -> {
079                    eventSubscribers.add(eventSubscriber);
080                });
081                
082                eventSubscriberEntityInstances.stream().map((eventSubscriberEntityInstance) -> eventSubscriberEntityInstance.getEventSubscriber()).filter((eventSubscriber) -> !eventSubscribers.contains(eventSubscriber)).map((eventSubscriber) -> {
083                    eventControl.createQueuedSubscriberEvent(eventSubscriber, event);
084                    return eventSubscriber;
085                }).forEach((eventSubscriber) -> {
086                    eventSubscribers.add(eventSubscriber);
087                });
088
089                eventControl.removeQueuedEvent(queuedEvent);
090
091                remainingTime -= System.currentTimeMillis() - startTime;
092                if(remainingTime < 0) {
093                    break;
094                }
095            }
096        }
097
098        return null;
099    }
100}