| 1 | |
package com.ctrip.framework.apollo.biz.message; |
| 2 | |
|
| 3 | |
import com.ctrip.framework.apollo.biz.entity.ReleaseMessage; |
| 4 | |
import com.ctrip.framework.apollo.biz.repository.ReleaseMessageRepository; |
| 5 | |
import com.dianping.cat.Cat; |
| 6 | |
import com.dianping.cat.message.Message; |
| 7 | |
import com.dianping.cat.message.Transaction; |
| 8 | |
|
| 9 | |
import org.slf4j.Logger; |
| 10 | |
import org.slf4j.LoggerFactory; |
| 11 | |
import org.springframework.beans.factory.annotation.Autowired; |
| 12 | |
import org.springframework.stereotype.Component; |
| 13 | |
|
| 14 | |
import java.util.Objects; |
| 15 | |
|
| 16 | |
|
| 17 | |
|
| 18 | |
|
| 19 | |
@Component |
| 20 | 3 | public class DatabaseMessageSender implements MessageSender { |
| 21 | 1 | private static final Logger logger = LoggerFactory.getLogger(DatabaseMessageSender.class); |
| 22 | |
|
| 23 | |
@Autowired |
| 24 | |
private ReleaseMessageRepository releaseMessageRepository; |
| 25 | |
|
| 26 | |
@Override |
| 27 | |
public void sendMessage(String message, String channel) { |
| 28 | 2 | logger.info("Sending message {} to channel {}", message, channel); |
| 29 | 2 | if (!Objects.equals(channel, Topics.APOLLO_RELEASE_TOPIC)) { |
| 30 | 1 | logger.warn("Channel {} not supported by DatabaseMessageSender!"); |
| 31 | 1 | return; |
| 32 | |
} |
| 33 | |
|
| 34 | 1 | Cat.logEvent("Apollo.AdminService.ReleaseMessage", message); |
| 35 | 1 | Transaction transaction = Cat.newTransaction("Apollo.AdminService", "sendMessage"); |
| 36 | |
try { |
| 37 | 1 | releaseMessageRepository.save(new ReleaseMessage(message)); |
| 38 | 1 | transaction.setStatus(Message.SUCCESS); |
| 39 | 0 | } catch (Throwable ex) { |
| 40 | 0 | logger.error("Sending message to database failed", ex); |
| 41 | 0 | transaction.setStatus(ex); |
| 42 | |
} finally { |
| 43 | 1 | transaction.complete(); |
| 44 | 1 | } |
| 45 | 1 | } |
| 46 | |
} |