- 修复组任务,其中一个子任务在获取文件长度失败后,重新恢复组合任务,组合任务状态变为完成的问题 https://github.com/AriaLyy/Aria/issues/628

This commit is contained in:
laoyuyu
2020-03-03 21:03:37 +08:00
parent 54cbdb7ee0
commit 669ac6b09c
15 changed files with 131 additions and 99 deletions

View File

@@ -51,6 +51,7 @@ import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
@@ -75,10 +76,11 @@ final class M3U8VodLoader extends BaseM3U8Loader {
private SparseArray<ThreadRecord> mAfterPeer = new SparseArray<>();
private PeerIndexEvent mCurrentEvent;
private String mCacheDir;
private int aIndex = 0, bIndex = 0;
private int mCurrentFlagSize;
private AtomicInteger afterPeerIndex = new AtomicInteger();
private AtomicInteger beforePeerIndex = new AtomicInteger();
private AtomicInteger mCompleteNum = new AtomicInteger();
private AtomicInteger mCurrentFlagSize = new AtomicInteger();
private boolean isJump = false, isDestroy = false;
private int mCompleteNum = 0;
private ExecutorService mJumpThreadPool;
private Thread jumpThread = null;
private M3U8TaskOption mM3U8Option;
@@ -102,20 +104,20 @@ final class M3U8VodLoader extends BaseM3U8Loader {
}
int getCompleteNum() {
return mCompleteNum;
return mCompleteNum.get();
}
void setCompleteNum(int mCompleteNum) {
this.mCompleteNum = mCompleteNum;
void setCompleteNum(int completeNum) {
mCompleteNum.set(completeNum);
}
int getCurrentFlagSize() {
mCurrentFlagSize = mFlagQueue.size();
return mCurrentFlagSize;
mCurrentFlagSize.set(mFlagQueue.size());
return mCurrentFlagSize.get();
}
void setCurrentFlagSize(int currentFlagSize) {
mCurrentFlagSize = currentFlagSize;
mCurrentFlagSize.set(currentFlagSize);
}
boolean isJump() {
@@ -178,7 +180,7 @@ final class M3U8VodLoader extends BaseM3U8Loader {
try {
LOCK.lock();
while (mFlagQueue.size() < EXEC_MAX_NUM && !isBreak()) {
if (mCompleteNum == mRecord.threadRecords.size()) {
if (mCompleteNum.get() == mRecord.threadRecords.size()) {
break;
}
@@ -214,16 +216,18 @@ final class M3U8VodLoader extends BaseM3U8Loader {
ThreadRecord tr = null;
try {
// 优先下载peer指针之后的数据
if (bIndex == 0 && aIndex < mAfterPeer.size()) {
if (beforePeerIndex.get() == 0 && afterPeerIndex.get() < mAfterPeer.size()) {
//ALog.d(TAG, String.format("afterArray size:%s, index:%s", mAfterPeer.size(), aIndex));
tr = mAfterPeer.valueAt(aIndex);
aIndex++;
tr = mAfterPeer.valueAt(afterPeerIndex.get());
afterPeerIndex.getAndIncrement();
}
// 如果指针之后的数组没有切片了,则重新初始化指针位置,并获取指针之前的数组获取切片进行下载
if (mBeforePeer.size() > 0 && (tr == null || bIndex != 0) && bIndex < mBeforePeer.size()) {
tr = mBeforePeer.valueAt(bIndex);
bIndex++;
if (mBeforePeer.size() > 0
&& (tr == null || beforePeerIndex.get() != 0)
&& beforePeerIndex.get() < mBeforePeer.size()) {
tr = mBeforePeer.valueAt(beforePeerIndex.get());
beforePeerIndex.getAndIncrement();
}
} catch (Exception e) {
e.printStackTrace();
@@ -256,19 +260,19 @@ final class M3U8VodLoader extends BaseM3U8Loader {
return;
}
// 设置需要下载的切片
mCompleteNum = 0;
mCompleteNum.set(0);
for (ThreadRecord tr : mRecord.threadRecords) {
if (!tr.isComplete) {
mAfterPeer.put(tr.threadId, tr);
} else {
mCompleteNum++;
mCompleteNum.getAndIncrement();
}
}
getStateManager().updateStateCount();
if (mCompleteNum <= 0) {
if (mCompleteNum.get() <= 0) {
getListener().onStart(0);
} else {
int percent = mCompleteNum * 100 / mRecord.threadRecords.size();
int percent = mCompleteNum.get() * 100 / mRecord.threadRecords.size();
getListener().onResume(percent);
}
}
@@ -329,7 +333,7 @@ final class M3U8VodLoader extends BaseM3U8Loader {
isJump = true;
notifyWaitLock(false);
mCurrentFlagSize = mFlagQueue.size();
mCurrentFlagSize.set(mFlagQueue.size());
// 停止所有正在执行的线程任务
try {
TempFlag flag;
@@ -403,12 +407,12 @@ final class M3U8VodLoader extends BaseM3U8Loader {
mBeforePeer.clear();
mAfterPeer.clear();
mFlagQueue.clear();
aIndex = 0;
bIndex = 0;
mCompleteNum = 0;
afterPeerIndex.set(0);
beforePeerIndex.set(0);
mCompleteNum.set(0);
for (ThreadRecord tr : mRecord.threadRecords) {
if (tr.isComplete) {
mCompleteNum++;
mCompleteNum.getAndIncrement();
continue;
}
if (tr.threadId < mCurrentEvent.peerIndex) {

View File

@@ -29,8 +29,8 @@ import com.arialyy.aria.core.loader.ILoaderVisitor;
import com.arialyy.aria.core.manager.ThreadTaskManager;
import com.arialyy.aria.core.processor.ITsMergeHandler;
import com.arialyy.aria.core.task.ThreadTask;
import com.arialyy.aria.exception.AriaM3U8Exception;
import com.arialyy.aria.exception.AriaException;
import com.arialyy.aria.exception.AriaM3U8Exception;
import com.arialyy.aria.m3u8.BaseM3U8Loader;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8TaskOption;
@@ -40,6 +40,7 @@ import com.arialyy.aria.util.FileUtil;
import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
/**
* m3u8 点播下载状态管理器
@@ -49,9 +50,9 @@ public final class VodStateManager implements IThreadStateManager {
private M3U8Listener listener;
private int startThreadNum; // 启动的线程总数
private int cancelNum = 0; // 已经取消的线程的数
private int stopNum = 0; // 已经停止的线程数
private int failNum = 0; // 失败的线程数
private AtomicInteger cancelNum = new AtomicInteger(0); // 已经取消的线程的数
private AtomicInteger stopNum = new AtomicInteger(0); // 已经停止的线程数
private AtomicInteger failNum = new AtomicInteger(0); // 失败的线程数
private long progress;
private TaskRecord taskRecord; // 任务记录
private Looper looper;
@@ -74,11 +75,11 @@ public final class VodStateManager implements IThreadStateManager {
int peerIndex = msg.getData().getInt(ISchedulers.DATA_M3U8_PEER_INDEX);
switch (msg.what) {
case STATE_STOP:
stopNum++;
stopNum.getAndIncrement();
removeSignThread((ThreadTask) msg.obj);
// 处理跳转位置后,恢复任务
if (loader.isJump()
&& (stopNum == loader.getCurrentFlagSize() || loader.getCurrentFlagSize() == 0)
&& (stopNum.get() == loader.getCurrentFlagSize() || loader.getCurrentFlagSize() == 0)
&& !loader.isBreak()) {
loader.resumeTask();
return true;
@@ -90,7 +91,7 @@ public final class VodStateManager implements IThreadStateManager {
}
break;
case STATE_CANCEL:
cancelNum++;
cancelNum.getAndIncrement();
removeSignThread((ThreadTask) msg.obj);
if (loader.isBreak()) {
@@ -99,7 +100,7 @@ public final class VodStateManager implements IThreadStateManager {
}
break;
case STATE_FAIL:
failNum++;
failNum.getAndIncrement();
for (ThreadRecord tr : taskRecord.threadRecords) {
if (tr.threadId == peerIndex) {
loader.getBeforePeer().put(peerIndex, tr);
@@ -175,9 +176,9 @@ public final class VodStateManager implements IThreadStateManager {
};
void updateStateCount() {
cancelNum = 0;
stopNum = 0;
failNum = 0;
cancelNum.set(0);
stopNum.set(0);
failNum.set(0);
}
@Override public void setLooper(TaskRecord taskRecord, Looper looper) {
@@ -233,12 +234,12 @@ public final class VodStateManager implements IThreadStateManager {
@Override public boolean isFail() {
printInfo("isFail");
return failNum != 0 && failNum == loader.getCurrentFlagSize() && !loader.isJump();
return failNum.get() != 0 && failNum.get() == loader.getCurrentFlagSize() && !loader.isJump();
}
@Override public boolean isComplete() {
if (m3U8Option.isIgnoreFailureTs()) {
return loader.getCompleteNum() + failNum >= taskRecord.threadRecords.size()
return loader.getCompleteNum() + failNum.get() >= taskRecord.threadRecords.size()
&& !loader.isJump();
} else {
return loader.getCompleteNum() == taskRecord.threadRecords.size() && !loader.isJump();