001    /**
002     * Licensed to the Apache Software Foundation (ASF) under one or more
003     * contributor license agreements.  See the NOTICE file distributed with
004     * this work for additional information regarding copyright ownership.
005     * The ASF licenses this file to You under the Apache License, Version 2.0
006     * (the "License"); you may not use this file except in compliance with
007     * the License.  You may obtain a copy of the License at
008     *
009     *      http://www.apache.org/licenses/LICENSE-2.0
010     *
011     * Unless required by applicable law or agreed to in writing, software
012     * distributed under the License is distributed on an "AS IS" BASIS,
013     * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
014     * See the License for the specific language governing permissions and
015     * limitations under the License.
016     */
017    package org.apache.activemq.camel.component.broker;
018    
019    import org.apache.activemq.broker.ProducerBrokerExchange;
020    import org.apache.activemq.broker.inteceptor.MessageInterceptor;
021    import org.apache.activemq.command.Message;
022    import org.apache.camel.Endpoint;
023    import org.apache.camel.Exchange;
024    import org.apache.camel.ExchangePattern;
025    import org.apache.camel.Processor;
026    import org.apache.camel.component.jms.JmsBinding;
027    import org.apache.camel.impl.DefaultConsumer;
028    import org.slf4j.Logger;
029    import org.slf4j.LoggerFactory;
030    
031    public class BrokerConsumer extends DefaultConsumer implements MessageInterceptor {
032        protected final transient Logger logger = LoggerFactory.getLogger(BrokerConsumer.class);
033        private final JmsBinding jmsBinding = new JmsBinding();
034    
035        public BrokerConsumer(Endpoint endpoint, Processor processor) {
036            super(endpoint, processor);
037        }
038    
039        @Override
040        protected void doStart() throws Exception {
041            super.doStart();
042            ((BrokerEndpoint) getEndpoint()).addMessageInterceptor(this);
043        }
044    
045        @Override
046        protected void doStop() throws Exception {
047            ((BrokerEndpoint) getEndpoint()).removeMessageInterceptor(this);
048            super.doStop();
049        }
050    
051        @Override
052        public void intercept(ProducerBrokerExchange producerExchange, Message message) {
053            Exchange exchange = getEndpoint().createExchange(ExchangePattern.InOnly);
054    
055            exchange.setIn(new BrokerJmsMessage((javax.jms.Message) message, jmsBinding));
056            exchange.setProperty(Exchange.BINDING, jmsBinding);
057            exchange.setProperty(BrokerEndpoint.PRODUCER_BROKER_EXCHANGE, producerExchange);
058            try {
059                getProcessor().process(exchange);
060            } catch (Exception e) {
061                logger.error("Failed to process " + exchange, e);
062            }
063        }
064    }