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}