001/*
002 * (C) Copyright 2017 Nuxeo (http://nuxeo.com/) and others.
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 * Contributors:
017 *     Florent Guillaume
018 */
019package org.nuxeo.runtime.pubsub;
020
021import java.util.List;
022import java.util.Map;
023import java.util.function.BiConsumer;
024
025import org.apache.commons.logging.Log;
026import org.apache.commons.logging.LogFactory;
027
028/**
029 * Abstract implementation of {@link PubSubProvider}.
030 * <p>
031 * This deals with subscribers registration and dispatch.
032 *
033 * @since 9.1
034 */
035public abstract class AbstractPubSubProvider implements PubSubProvider {
036
037    private final Log log = LogFactory.getLog(AbstractPubSubProvider.class);
038
039    protected String namespace;
040
041    /** List of subscribers for each topic. */
042    protected Map<String, List<BiConsumer<String, byte[]>>> subscribers;
043
044    @Override
045    public void initialize(Map<String, List<BiConsumer<String, byte[]>>> subscribers) {
046        this.subscribers = subscribers;
047    }
048
049    @Override
050    public void close() {
051        // DO NOT subscribers.clear(), we do not own this map
052    }
053
054    public void localPublish(String topic, byte[] message) {
055        if (subscribers == null) {
056            // not yet initialized
057            return;
058        }
059        List<BiConsumer<String, byte[]>> subs = subscribers.get(topic);
060        if (subs != null) {
061            for (BiConsumer<String, byte[]> subscriber : subs) {
062                try {
063                    subscriber.accept(topic, message);
064                } catch (RuntimeException e) {
065                    if (Thread.currentThread().isInterrupted()) {
066                        throw e;
067                    }
068                    // don't break everything if a subscriber is ill-behaved
069                    log.error("Exception in subscriber for topic: " + topic, e);
070                }
071            }
072        }
073    }
074
075}