/
PushService.java
150 lines (128 loc) · 5.66 KB
/
PushService.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
/*
* SymmetricDS is an open source database synchronization solution.
*
* Copyright (C) Chris Henson <chenson42@users.sourceforge.net>
*
* This library is free software; you can redistribute it and/or
* modify it under the terms of the GNU Lesser General Public
* License as published by the Free Software Foundation; either
* version 3 of the License, or (at your option) any later version.
*
* This library is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
* Lesser General Public License for more details.
*
* You should have received a copy of the GNU Lesser General Public
* License along with this library; if not, see
* <http://www.gnu.org/licenses/>.
*/
package org.jumpmind.symmetric.service.impl;
import java.io.BufferedReader;
import java.net.ConnectException;
import java.net.SocketException;
import java.util.List;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.jumpmind.symmetric.common.ErrorConstants;
import org.jumpmind.symmetric.model.BatchInfo;
import org.jumpmind.symmetric.model.Node;
import org.jumpmind.symmetric.model.NodeSecurity;
import org.jumpmind.symmetric.service.IAcknowledgeService;
import org.jumpmind.symmetric.service.IDataExtractorService;
import org.jumpmind.symmetric.service.IDataService;
import org.jumpmind.symmetric.service.INodeService;
import org.jumpmind.symmetric.service.IPushService;
import org.jumpmind.symmetric.transport.AuthenticationException;
import org.jumpmind.symmetric.transport.ConnectionRejectedException;
import org.jumpmind.symmetric.transport.IOutgoingWithResponseTransport;
import org.jumpmind.symmetric.transport.ITransportManager;
import org.jumpmind.symmetric.transport.TransportException;
public class PushService extends AbstractService implements IPushService {
private static final Log logger = LogFactory.getLog(PushService.class);
private IDataExtractorService extractor;
private IAcknowledgeService ackService;
private ITransportManager transportManager;
private INodeService nodeService;
private IDataService dataService;
public void pushData() {
List<Node> nodes = nodeService.findNodesToPushTo();
if (nodes != null && nodes.size() > 0) {
for (Node node : nodes) {
logger.info("Push requested for " + node);
if (pushToNode(node)) {
logger.info("Push completed for " + node);
} else {
logger.info("Push unsuccessful for " + node);
}
}
}
}
private boolean pushToNode(Node remote) {
IOutgoingWithResponseTransport transport = null;
boolean success = false;
try {
NodeSecurity nodeSecurity = nodeService.findNodeSecurity(remote.getNodeId());
if (nodeSecurity != null) {
if (nodeSecurity.isInitialLoadEnabled()) {
dataService.insertReloadEvent(remote);
}
}
transport = transportManager.getPushTransport(remote, nodeService.findIdentity());
if (extractor.extract(remote, transport)) {
logger.info("Push data sent to " + remote);
BufferedReader reader = transport.readResponse();
String ackString = reader.readLine();
String ackExtendedString = reader.readLine();
if (logger.isDebugEnabled()) {
logger.debug("Reading ack: " + ackString);
logger.debug("Reading extended ack: " + ackExtendedString);
}
List<BatchInfo> batches = transportManager.readAcknowledgement(ackString, ackExtendedString);
for (BatchInfo batchInfo : batches) {
if (logger.isDebugEnabled()) {
logger.debug("Saving ack: " + batchInfo.getBatchId() + ", "
+ (batchInfo.isOk() ? "OK" : "error"));
}
ackService.ack(batchInfo);
}
}
success = true;
} catch (ConnectException ex) {
logger.warn(ErrorConstants.COULD_NOT_CONNECT_TO_TRANSPORT + " url=" + remote.getSyncURL());
} catch (ConnectionRejectedException ex) {
logger.warn(ErrorConstants.TRANSPORT_REJECTED_CONNECTION);
} catch (SocketException ex) {
logger.warn(ex.getMessage());
} catch (TransportException ex) {
logger.warn(ex.getMessage());
} catch (AuthenticationException ex) {
logger.warn(ErrorConstants.NOT_AUTHENTICATED);
} catch (Exception e) {
// just report the error because we want to push to other nodes
// in our list
logger.error(e, e);
} finally {
try {
transport.close();
} catch (Exception e) {
}
}
return success;
}
public void setExtractor(IDataExtractorService extractor) {
this.extractor = extractor;
}
public void setTransportManager(ITransportManager tm) {
this.transportManager = tm;
}
public void setNodeService(INodeService nodeService) {
this.nodeService = nodeService;
}
public void setAckService(IAcknowledgeService ackService) {
this.ackService = ackService;
}
public void setDataService(IDataService dataService) {
this.dataService = dataService;
}
}