🐛 fix: Fix AvgSpeed concurrent modification exception, Optimize the number of requests for getIdleChatFiles
This commit is contained in:
parent
c92a780245
commit
c018f98aa1
3 changed files with 34 additions and 18 deletions
|
|
@ -1,5 +1,6 @@
|
||||||
package telegram.files;
|
package telegram.files;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.TreeMap;
|
import java.util.TreeMap;
|
||||||
|
|
@ -172,7 +173,7 @@ public class AvgSpeed {
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
List<Long> speeds = speedPoints.values().stream()
|
List<Long> speeds = new ArrayList<>(speedPoints.values()).stream()
|
||||||
.map(point -> point.speed)
|
.map(point -> point.speed)
|
||||||
.filter(speed -> speed > 0)
|
.filter(speed -> speed > 0)
|
||||||
.sorted()
|
.sorted()
|
||||||
|
|
|
||||||
|
|
@ -260,7 +260,7 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
searchChatMessages.filter = TdApiHelp.getSearchMessagesFilter(filter.get("type"));
|
searchChatMessages.filter = TdApiHelp.getSearchMessagesFilter(filter.get("type"));
|
||||||
|
|
||||||
return (Objects.equals(filter.get("status"), FileRecord.DownloadStatus.idle.name()) ?
|
return (Objects.equals(filter.get("status"), FileRecord.DownloadStatus.idle.name()) ?
|
||||||
this.getIdleChatFiles(searchChatMessages) :
|
this.getIdleChatFiles(searchChatMessages, 0) :
|
||||||
this.execute(searchChatMessages))
|
this.execute(searchChatMessages))
|
||||||
.compose(foundChatMessages ->
|
.compose(foundChatMessages ->
|
||||||
DataVerticle.fileRepository.getFilesByUniqueId(TdApiHelp.getFileUniqueIds(Arrays.asList(foundChatMessages.messages)))
|
DataVerticle.fileRepository.getFilesByUniqueId(TdApiHelp.getFileUniqueIds(Arrays.asList(foundChatMessages.messages)))
|
||||||
|
|
@ -269,7 +269,11 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private Future<TdApi.FoundChatMessages> getIdleChatFiles(TdApi.SearchChatMessages searchChatMessages) {
|
private Future<TdApi.FoundChatMessages> getIdleChatFiles(TdApi.SearchChatMessages searchChatMessages, int seq) {
|
||||||
|
if (seq != 0) {
|
||||||
|
// Increase the limit and reduce the number of requests
|
||||||
|
searchChatMessages.limit = 100;
|
||||||
|
}
|
||||||
return this.execute(searchChatMessages)
|
return this.execute(searchChatMessages)
|
||||||
.compose(foundChatMessages -> {
|
.compose(foundChatMessages -> {
|
||||||
TdApi.Message[] messages = Stream.of(foundChatMessages.messages)
|
TdApi.Message[] messages = Stream.of(foundChatMessages.messages)
|
||||||
|
|
@ -284,9 +288,9 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
.orElse(false)
|
.orElse(false)
|
||||||
)
|
)
|
||||||
.toArray(TdApi.Message[]::new);
|
.toArray(TdApi.Message[]::new);
|
||||||
if (ArrayUtil.isEmpty(messages)) {
|
if (ArrayUtil.isEmpty(messages) && foundChatMessages.nextFromMessageId != 0) {
|
||||||
searchChatMessages.fromMessageId = foundChatMessages.nextFromMessageId;
|
searchChatMessages.fromMessageId = foundChatMessages.nextFromMessageId;
|
||||||
return getIdleChatFiles(searchChatMessages);
|
return getIdleChatFiles(searchChatMessages, seq + 1);
|
||||||
} else {
|
} else {
|
||||||
foundChatMessages.messages = messages;
|
foundChatMessages.messages = messages;
|
||||||
return Future.succeededFuture(foundChatMessages);
|
return Future.succeededFuture(foundChatMessages);
|
||||||
|
|
@ -461,7 +465,14 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
file.local.path,
|
file.local.path,
|
||||||
FileRecord.DownloadStatus.completed,
|
FileRecord.DownloadStatus.completed,
|
||||||
System.currentTimeMillis()
|
System.currentTimeMillis()
|
||||||
).compose(r -> Future.failedFuture("File is already downloaded successfully"));
|
).compose(r -> {
|
||||||
|
sendFileStatusHttpEvent(file, r);
|
||||||
|
if (r == null || r.isEmpty()) {
|
||||||
|
return Future.failedFuture("File is downloaded completed, but update status failed");
|
||||||
|
} else {
|
||||||
|
return Future.failedFuture("File is already downloaded successfully");
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
if (isPaused && !file.local.isDownloadingActive) {
|
if (isPaused && !file.local.isDownloadingActive) {
|
||||||
return Future.failedFuture("File is not downloading");
|
return Future.failedFuture("File is not downloading");
|
||||||
|
|
@ -674,6 +685,17 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
JsonObject.of("telegramId", this.getId(), "payload", JsonObject.mapFrom(payload)));
|
JsonObject.of("telegramId", this.getId(), "payload", JsonObject.mapFrom(payload)));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void sendFileStatusHttpEvent(TdApi.File file, JsonObject fileUpdated) {
|
||||||
|
if (fileUpdated == null || fileUpdated.isEmpty()) return;
|
||||||
|
sendHttpEvent(EventPayload.build(EventPayload.TYPE_FILE_STATUS, new JsonObject()
|
||||||
|
.put("fileId", file.id)
|
||||||
|
.put("downloadStatus", fileUpdated.getString("downloadStatus"))
|
||||||
|
.put("localPath", fileUpdated.getString("localPath"))
|
||||||
|
.put("completionDate", fileUpdated.getLong("completionDate"))
|
||||||
|
.put("downloadedSize", file.local.downloadedSize)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
private void handleAuthorizationResult(TdApi.Object object) {
|
private void handleAuthorizationResult(TdApi.Object object) {
|
||||||
switch (object.getConstructor()) {
|
switch (object.getConstructor()) {
|
||||||
case TdApi.Error.CONSTRUCTOR:
|
case TdApi.Error.CONSTRUCTOR:
|
||||||
|
|
@ -828,16 +850,7 @@ public class TelegramVerticle extends AbstractVerticle {
|
||||||
localPath,
|
localPath,
|
||||||
downloadStatus,
|
downloadStatus,
|
||||||
completionDate)
|
completionDate)
|
||||||
.onSuccess(r -> {
|
.onSuccess(r -> sendFileStatusHttpEvent(file, r));
|
||||||
if (r == null || r.isEmpty()) return;
|
|
||||||
sendHttpEvent(EventPayload.build(EventPayload.TYPE_FILE_STATUS, new JsonObject()
|
|
||||||
.put("fileId", file.id)
|
|
||||||
.put("downloadStatus", r.getString("downloadStatus"))
|
|
||||||
.put("localPath", r.getString("localPath"))
|
|
||||||
.put("completionDate", r.getLong("completionDate"))
|
|
||||||
.put("downloadedSize", downloadedSize)
|
|
||||||
));
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
if (lastFileEventTime == 0 || System.currentTimeMillis() - lastFileEventTime > 1000) {
|
if (lastFileEventTime == 0 || System.currentTimeMillis() - lastFileEventTime > 1000) {
|
||||||
sendHttpEvent(EventPayload.build(EventPayload.TYPE_FILE, updateFile));
|
sendHttpEvent(EventPayload.build(EventPayload.TYPE_FILE, updateFile));
|
||||||
|
|
|
||||||
|
|
@ -100,17 +100,19 @@ public class DataVerticleTest {
|
||||||
);
|
);
|
||||||
String updateLocalPath = "local_path";
|
String updateLocalPath = "local_path";
|
||||||
Long completionDate = 1L;
|
Long completionDate = 1L;
|
||||||
|
int newFileId = 2;
|
||||||
DataVerticle.fileRepository.create(fileRecord)
|
DataVerticle.fileRepository.create(fileRecord)
|
||||||
.compose(r -> DataVerticle.fileRepository.updateStatus(r.id(), r.uniqueId(), updateLocalPath, FileRecord.DownloadStatus.downloading, completionDate))
|
.compose(r -> DataVerticle.fileRepository.updateStatus(newFileId, r.uniqueId(), updateLocalPath, FileRecord.DownloadStatus.downloading, completionDate))
|
||||||
.compose(r -> {
|
.compose(r -> {
|
||||||
testContext.verify(() -> {
|
testContext.verify(() -> {
|
||||||
Assertions.assertEquals(updateLocalPath, r.getString("localPath"));
|
Assertions.assertEquals(updateLocalPath, r.getString("localPath"));
|
||||||
Assertions.assertEquals(FileRecord.DownloadStatus.downloading.name(), r.getString("downloadStatus"));
|
Assertions.assertEquals(FileRecord.DownloadStatus.downloading.name(), r.getString("downloadStatus"));
|
||||||
Assertions.assertEquals(completionDate, r.getLong("completionDate"));
|
Assertions.assertEquals(completionDate, r.getLong("completionDate"));
|
||||||
});
|
});
|
||||||
return DataVerticle.fileRepository.getByPrimaryKey(fileRecord.id(), fileRecord.uniqueId());
|
return DataVerticle.fileRepository.getByPrimaryKey(newFileId, fileRecord.uniqueId());
|
||||||
})
|
})
|
||||||
.onComplete(testContext.succeeding(r -> testContext.verify(() -> {
|
.onComplete(testContext.succeeding(r -> testContext.verify(() -> {
|
||||||
|
Assertions.assertEquals(newFileId, r.id());
|
||||||
Assertions.assertEquals(FileRecord.DownloadStatus.downloading.name(), r.downloadStatus());
|
Assertions.assertEquals(FileRecord.DownloadStatus.downloading.name(), r.downloadStatus());
|
||||||
Assertions.assertEquals(updateLocalPath, r.localPath());
|
Assertions.assertEquals(updateLocalPath, r.localPath());
|
||||||
testContext.completeNow();
|
testContext.completeNow();
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue