-
Notifications
You must be signed in to change notification settings - Fork 868
/
OTransactionPhase1Task.java
238 lines (213 loc) · 9.05 KB
/
OTransactionPhase1Task.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
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
package com.orientechnologies.orient.server.distributed.impl.task;
import com.orientechnologies.common.concur.lock.OLockException;
import com.orientechnologies.orient.client.remote.message.OMessageHelper;
import com.orientechnologies.orient.client.remote.message.tx.ORecordOperationRequest;
import com.orientechnologies.orient.core.Orient;
import com.orientechnologies.orient.core.command.OCommandDistributedReplicateRequest;
import com.orientechnologies.orient.core.db.ODatabaseDocumentInternal;
import com.orientechnologies.orient.core.db.record.ORecordOperation;
import com.orientechnologies.orient.core.exception.OConcurrentModificationException;
import com.orientechnologies.orient.core.id.ORecordId;
import com.orientechnologies.orient.core.record.ORecord;
import com.orientechnologies.orient.core.record.ORecordInternal;
import com.orientechnologies.orient.core.serialization.serializer.record.binary.ORecordSerializerNetworkV37;
import com.orientechnologies.orient.core.storage.ORecordDuplicatedException;
import com.orientechnologies.orient.core.storage.impl.local.paginated.wal.OLogSequenceNumber;
import com.orientechnologies.orient.core.tx.OTransactionIndexChanges;
import com.orientechnologies.orient.core.tx.OTransactionInternal;
import com.orientechnologies.orient.server.OServer;
import com.orientechnologies.orient.server.distributed.ODistributedRequestId;
import com.orientechnologies.orient.server.distributed.ODistributedServerManager;
import com.orientechnologies.orient.server.distributed.ORemoteTaskFactory;
import com.orientechnologies.orient.server.distributed.impl.ODatabaseDocumentDistributed;
import com.orientechnologies.orient.server.distributed.impl.OTransactionOptimisticDistributed;
import com.orientechnologies.orient.server.distributed.impl.task.transaction.*;
import com.orientechnologies.orient.server.distributed.task.OAbstractReplicatedTask;
import com.orientechnologies.orient.server.distributed.task.ODistributedLockException;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
/**
* @author Luigi Dell'Aquila (l.dellaquila - at - orientdb.com)
*/
public class OTransactionPhase1Task extends OAbstractReplicatedTask {
public static final int FACTORYID = 43;
private volatile boolean hasResponse;
private OLogSequenceNumber lastLSN;
private List<ORecordOperation> ops;
private List<ORecordOperationRequest> operations;
private OCommandDistributedReplicateRequest.QUORUM_TYPE quorumType = OCommandDistributedReplicateRequest.QUORUM_TYPE.WRITE;
private transient int retryCount = 0;
public OTransactionPhase1Task() {
ops = new ArrayList<>();
operations = new ArrayList<>();
}
public OTransactionPhase1Task(List<ORecordOperation> ops) {
this.ops = ops;
operations = new ArrayList<>();
genOps(ops);
}
public void genOps(List<ORecordOperation> ops) {
for (ORecordOperation txEntry : ops) {
if (txEntry.type == ORecordOperation.LOADED)
continue;
ORecordOperationRequest request = new ORecordOperationRequest();
request.setType(txEntry.type);
request.setVersion(txEntry.getRecord().getVersion());
request.setId(txEntry.getRecord().getIdentity());
request.setRecordType(ORecordInternal.getRecordType(txEntry.getRecord()));
switch (txEntry.type) {
case ORecordOperation.CREATED:{
byte[] newRec = ORecordSerializerNetworkV37.INSTANCE.toStream(txEntry.getRecord(), false);
request.setRecord(newRec);
request.setContentChanged(ORecordInternal.isContentChanged(txEntry.getRecord()));
}
break;
case ORecordOperation.UPDATED:
byte[] deltaRec = ORecordSerializerNetworkV37.INSTANCE.toStream(txEntry.getRecord(), true);
byte[] newRec = ORecordSerializerNetworkV37.INSTANCE.toStream(txEntry.getRecord(), false);
request.setRecord(newRec);
request.setContentChanged(ORecordInternal.isContentChanged(txEntry.getRecord()));
break;
case ORecordOperation.DELETED:
break;
}
operations.add(request);
}
}
@Override
public String getName() {
return "TxPhase1";
}
@Override
public OCommandDistributedReplicateRequest.QUORUM_TYPE getQuorumType() {
return quorumType;
}
@Override
public Object execute(ODistributedRequestId requestId, OServer iServer, ODistributedServerManager iManager,
ODatabaseDocumentInternal database) throws Exception {
convert(database);
OTransactionOptimisticDistributed tx = new OTransactionOptimisticDistributed(database, ops);
OTransactionResultPayload res1 = executeTransaction(requestId, (ODatabaseDocumentDistributed) database, tx, false, retryCount);
if (res1 == null) {
retryCount++;
((ODatabaseDocumentDistributed) database).getStorageDistributed().getLocalDistributedDatabase()
.reEnqueue(requestId.getNodeId(), requestId.getMessageId(), database.getName(), this);
hasResponse = false;
return null;
}
hasResponse = true;
return new OTransactionPhase1TaskResult(res1);
}
@Override
public boolean hasResponse() {
return hasResponse;
}
public static OTransactionResultPayload executeTransaction(ODistributedRequestId requestId, ODatabaseDocumentDistributed database,
OTransactionInternal tx, boolean local, int retryCount) {
OTransactionResultPayload payload;
try {
if (database.beginDistributedTx(requestId, tx, local, retryCount)) {
payload = new OTxSuccess();
} else {
return null;
}
} catch (OConcurrentModificationException ex) {
payload = new OTxConcurrentModification((ORecordId) ex.getRid(), ex.getEnhancedDatabaseVersion());
} catch (ODistributedLockException | OLockException ex) {
payload = new OTxLockTimeout();
} catch (ORecordDuplicatedException ex) {
//TODO:Check if can get out the key
payload = new OTxUniqueIndex((ORecordId) ex.getRid(), ex.getIndexName(), ex.getKey());
} catch (RuntimeException ex) {
payload = new OTxException(ex);
}
return payload;
}
@Override
public void fromStream(DataInput in, ORemoteTaskFactory factory) throws IOException {
int size = in.readInt();
for (int i = 0; i < size; i++) {
ORecordOperationRequest req = OMessageHelper.readTransactionEntry(in);
operations.add(req);
}
lastLSN = new OLogSequenceNumber(in);
}
private void convert(ODatabaseDocumentInternal database) {
for (ORecordOperationRequest req : operations) {
byte type = req.getType();
if (type == ORecordOperation.LOADED) {
continue;
}
ORecord record = null;
switch (type) {
case ORecordOperation.CREATED:{
record = ORecordSerializerNetworkV37.INSTANCE.fromStream(req.getRecord(), null, null);
ORecordInternal.setRecordSerializer(record, database.getSerializer());
}
break;
case ORecordOperation.UPDATED: {
record = ORecordSerializerNetworkV37.INSTANCE.fromStream(req.getRecord(), null, null);
ORecordInternal.setRecordSerializer(record, database.getSerializer());
}
break;
case ORecordOperation.DELETED:
record = database.getRecord(req.getId());
if (record == null) {
record = Orient.instance().getRecordFactoryManager()
.newInstance(req.getRecordType(), req.getId().getClusterId(), database);
}
break;
}
ORecordInternal.setIdentity(record, (ORecordId) req.getId());
ORecordInternal.setVersion(record, req.getVersion());
ORecordOperation op = new ORecordOperation(record, type);
ops.add(op);
}
operations.clear();
}
@Override
public void toStream(DataOutput out) throws IOException {
out.writeInt(operations.size());
for (ORecordOperationRequest operation : operations) {
OMessageHelper.writeTransactionEntry(out, operation);
}
lastLSN.toStream(out);
}
@Override
public int getFactoryId() {
return FACTORYID;
}
public void init(OTransactionInternal operations) {
for (Map.Entry<String, OTransactionIndexChanges> indexOp : operations.getIndexOperations().entrySet()) {
if (indexOp.getValue().resolveAssociatedIndex(indexOp.getKey(), operations.getDatabase().getMetadata().getIndexManager())
.isUnique()) {
quorumType = OCommandDistributedReplicateRequest.QUORUM_TYPE.ALL;
break;
}
}
this.ops = new ArrayList<>(operations.getRecordOperations());
genOps(this.ops);
}
public void setLastLSN(OLogSequenceNumber lastLSN) {
this.lastLSN = lastLSN;
}
@Override
public OLogSequenceNumber getLastLSN() {
return lastLSN;
}
@Override
public boolean isIdempotent() {
return false;
}
@Override
public int[] getPartitionKey() {
if (operations.size() > 0)
return operations.stream().mapToInt((x) -> x.getId().getClusterId()).toArray();
else
return ops.stream().mapToInt((x) -> x.getRID().getClusterId()).toArray();
}
}