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 */ 017package org.apache.activemq.transport.amqp.protocol; 018 019import java.util.ArrayList; 020import java.util.List; 021 022import org.apache.activemq.command.Command; 023import org.apache.qpid.proton.amqp.transport.ErrorCondition; 024import org.apache.qpid.proton.engine.Link; 025import org.apache.qpid.proton.engine.Sender; 026 027/** 028 * Abstract AmqpLink implementation that provide basic Link services. 029 */ 030public abstract class AmqpAbstractLink<LINK_TYPE extends Link> implements AmqpLink { 031 032 protected final AmqpSession session; 033 protected final LINK_TYPE endpoint; 034 035 protected boolean closed; 036 protected boolean opened; 037 protected List<Runnable> closeActions = new ArrayList<Runnable>(); 038 039 /** 040 * Creates a new AmqpLink type. 041 * 042 * @param session 043 * the AmqpSession that servers as the parent of this Link. 044 * @param endpoint 045 * the link endpoint this object represents. 046 */ 047 public AmqpAbstractLink(AmqpSession session, LINK_TYPE endpoint) { 048 this.session = session; 049 this.endpoint = endpoint; 050 } 051 052 @Override 053 public void open() { 054 if (!opened) { 055 getEndpoint().setContext(this); 056 getEndpoint().open(); 057 058 opened = true; 059 } 060 } 061 062 @Override 063 public void detach() { 064 if (!closed) { 065 if (getEndpoint() != null) { 066 getEndpoint().setContext(null); 067 getEndpoint().detach(); 068 getEndpoint().free(); 069 } 070 } 071 } 072 073 @Override 074 public void close(ErrorCondition error) { 075 if (!closed) { 076 077 if (getEndpoint() != null) { 078 if (getEndpoint() instanceof Sender) { 079 getEndpoint().setSource(null); 080 } else { 081 getEndpoint().setTarget(null); 082 } 083 getEndpoint().setCondition(error); 084 } 085 086 close(); 087 } 088 } 089 090 @Override 091 public void close() { 092 if (!closed) { 093 094 if (getEndpoint() != null) { 095 getEndpoint().setContext(null); 096 getEndpoint().close(); 097 getEndpoint().free(); 098 } 099 100 for (Runnable action : closeActions) { 101 action.run(); 102 } 103 104 closeActions.clear(); 105 opened = false; 106 closed = true; 107 } 108 } 109 110 /** 111 * @return true if this link has already been opened. 112 */ 113 public boolean isOpened() { 114 return opened; 115 } 116 117 /** 118 * @return true if this link has already been closed. 119 */ 120 public boolean isClosed() { 121 return closed; 122 } 123 124 /** 125 * @return the Proton Link type this link represents. 126 */ 127 public LINK_TYPE getEndpoint() { 128 return endpoint; 129 } 130 131 /** 132 * @return the parent AmqpSession for this Link instance. 133 */ 134 public AmqpSession getSession() { 135 return session; 136 } 137 138 @Override 139 public void addCloseAction(Runnable action) { 140 closeActions.add(action); 141 } 142 143 /** 144 * Shortcut method to hand off an ActiveMQ Command to the broker and assign 145 * a ResponseHandler to deal with any reply from the broker. 146 * 147 * @param command 148 * the Command object to send to the Broker. 149 */ 150 protected void sendToActiveMQ(Command command) { 151 session.getConnection().sendToActiveMQ(command, null); 152 } 153 154 /** 155 * Shortcut method to hand off an ActiveMQ Command to the broker and assign 156 * a ResponseHandler to deal with any reply from the broker. 157 * 158 * @param command 159 * the Command object to send to the Broker. 160 * @param handler 161 * the ResponseHandler that will handle the Broker's response. 162 */ 163 protected void sendToActiveMQ(Command command, ResponseHandler handler) { 164 session.getConnection().sendToActiveMQ(command, handler); 165 } 166}