loader 重构

This commit is contained in:
laoyuyu
2020-01-01 14:43:52 +08:00
170 changed files with 2862 additions and 21661 deletions

View File

@@ -16,14 +16,11 @@
package com.arialyy.aria.m3u8;
import android.text.TextUtils;
import com.arialyy.aria.core.common.RecordHandler;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.download.DownloadEntity;
import com.arialyy.aria.core.download.M3U8Entity;
import com.arialyy.aria.core.loader.IRecordHandler;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.loader.AbsLoader;
import com.arialyy.aria.core.wrapper.AbsTaskWrapper;
import com.arialyy.aria.core.loader.AbsNormalLoader;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.FileUtil;
import java.io.BufferedReader;
@@ -35,11 +32,11 @@ import java.io.IOException;
import java.io.InputStreamReader;
import java.nio.charset.Charset;
public abstract class BaseM3U8Loader extends AbsLoader {
public abstract class BaseM3U8Loader extends AbsNormalLoader {
protected M3U8TaskOption mM3U8Option;
public BaseM3U8Loader(IEventListener listener, DTaskWrapper wrapper) {
super(listener, wrapper);
public BaseM3U8Loader(DTaskWrapper wrapper, IEventListener listener) {
super(wrapper, listener);
mM3U8Option = (M3U8TaskOption) wrapper.getM3u8Option();
mTempFile = new File(wrapper.getEntity().getFilePath());
}
@@ -58,7 +55,7 @@ public abstract class BaseM3U8Loader extends AbsLoader {
return String.format("%s/%s.ts", dirCache, threadId);
}
protected String getCacheDir() {
public String getCacheDir() {
String cacheDir = mM3U8Option.getCacheDir();
if (TextUtils.isEmpty(cacheDir)) {
cacheDir = FileUtil.getTsCacheDir(getEntity().getFilePath(), mM3U8Option.getBandWidth());
@@ -74,7 +71,7 @@ public abstract class BaseM3U8Loader extends AbsLoader {
*/
public boolean generateIndexFile(boolean isLive) {
File tempFile =
new File(String.format(M3U8InfoThread.M3U8_INDEX_FORMAT, getEntity().getFilePath()));
new File(String.format(M3U8InfoTask.M3U8_INDEX_FORMAT, getEntity().getFilePath()));
if (!tempFile.exists()) {
ALog.e(TAG, "源索引文件不存在");
return false;
@@ -136,17 +133,10 @@ public abstract class BaseM3U8Loader extends AbsLoader {
return false;
}
@Override public long getCurrentLocation() {
@Override public long getCurrentProgress() {
return isRunning() ? getStateManager().getCurrentProgress() : getEntity().getCurrentProgress();
}
@Override protected IRecordHandler getRecordHandler(AbsTaskWrapper wrapper) {
RecordHandler handler = new RecordHandler(wrapper);
M3U8RecordHandler adapter = new M3U8RecordHandler((DTaskWrapper) wrapper);
handler.setAdapter(adapter);
return handler;
}
protected DownloadEntity getEntity() {
return (DownloadEntity) mTaskWrapper.getEntity();
}

View File

@@ -1,6 +1,6 @@
package com.arialyy.aria.m3u8;
public class IdGenerator {
public final class IdGenerator {
/**
* SnowFlake算法 64位Long类型生成唯一ID 第一位0表明正数 2-4241位表示毫秒时间戳差值起始值自定义
* 43-5210位机器编号5位数据中心编号5位进程编号 53-6412位毫秒内计数器 本机内存生成,性能高

View File

@@ -24,7 +24,8 @@ import com.arialyy.aria.core.common.CompleteInfo;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.download.DownloadEntity;
import com.arialyy.aria.core.download.M3U8Entity;
import com.arialyy.aria.core.inf.OnFileInfoCallback;
import com.arialyy.aria.core.loader.IInfoTask;
import com.arialyy.aria.core.loader.ILoaderVisitor;
import com.arialyy.aria.core.processor.IBandWidthUrlConverter;
import com.arialyy.aria.core.processor.IKeyUrlConverter;
import com.arialyy.aria.core.wrapper.AbsTaskWrapper;
@@ -59,30 +60,33 @@ import java.util.regex.Pattern;
* https://www.cnblogs.com/renhui/p/10351870.html
* https://blog.csdn.net/Guofengpu/article/details/54922865
*/
final public class M3U8InfoThread implements Runnable {
final public class M3U8InfoTask implements IInfoTask {
public static final String M3U8_INDEX_FORMAT = "%s.index";
private final String TAG = "M3U8InfoThread";
private DownloadEntity mEntity;
private DTaskWrapper mTaskWrapper;
private int mConnectTimeOut;
private OnFileInfoCallback onFileInfoCallback;
private OnGetLivePeerCallback onGetPeerCallback;
private HttpTaskOption mHttpOption;
private M3U8TaskOption mM3U8Option;
private Callback mCallback;
/**
* 是否停止获取切片信息{@code true}停止获取切片信息
*/
private boolean isStop = false;
@Override public void accept(ILoaderVisitor visitor) {
visitor.addComponent(this);
}
public interface OnGetLivePeerCallback {
void onGetPeer(String url, String extInf);
}
public M3U8InfoThread(DTaskWrapper taskWrapper, OnFileInfoCallback callback) {
public M3U8InfoTask(DTaskWrapper taskWrapper) {
this.mTaskWrapper = taskWrapper;
mEntity = taskWrapper.getEntity();
mConnectTimeOut = AriaConfig.getInstance().getDConfig().getConnectTimeOut();
onFileInfoCallback = callback;
mHttpOption = (HttpTaskOption) taskWrapper.getTaskOption();
mM3U8Option = (M3U8TaskOption) taskWrapper.getM3u8Option();
mEntity.getM3U8Entity().setLive(mTaskWrapper.getRequestType() == AbsTaskWrapper.M3U8_LIVE);
@@ -108,6 +112,10 @@ final public class M3U8InfoThread implements Runnable {
}
}
@Override public void setCallback(Callback callback) {
mCallback = callback;
}
private void handleConnect(HttpURLConnection conn) throws IOException {
int code = conn.getResponseCode();
if (code == HttpURLConnection.HTTP_OK) {
@@ -191,7 +199,7 @@ final public class M3U8InfoThread implements Runnable {
CompleteInfo info = new CompleteInfo();
info.obj = extInf;
onFileInfoCallback.onComplete(mEntity.getKey(), info);
mCallback.onSucceed(mEntity.getKey(), info);
if (fos != null) {
fos.close();
}
@@ -255,9 +263,9 @@ final public class M3U8InfoThread implements Runnable {
m3U8Entity.keyUrl) + ".key";
} else if (param.startsWith("IV")) {
m3U8Entity.iv = param.split("=")[1];
}else if (param.startsWith("KEYFORMAT")){
} else if (param.startsWith("KEYFORMAT")) {
m3U8Entity.keyFormat = param.split("=")[1];
}else if (param.startsWith("KEYFORMATVERSIONS")){
} else if (param.startsWith("KEYFORMATVERSIONS")) {
m3U8Entity.keyFormatVersion = param.split("=")[1];
}
}
@@ -282,8 +290,8 @@ final public class M3U8InfoThread implements Runnable {
private void handleUrlReTurn(HttpURLConnection conn, String newUrl) throws IOException {
ALog.d(TAG, "30x跳转新url为【" + newUrl + "");
if (TextUtils.isEmpty(newUrl) || newUrl.equalsIgnoreCase("null")) {
if (onFileInfoCallback != null) {
onFileInfoCallback.onFail(mEntity, new TaskException(TAG, "获取重定向链接失败"), false);
if (mCallback != null) {
mCallback.onFail(mEntity, new TaskException(TAG, "获取重定向链接失败"), false);
}
return;
}
@@ -341,7 +349,7 @@ final public class M3U8InfoThread implements Runnable {
}
private void failDownload(String errorInfo, boolean needRetry) {
onFileInfoCallback.onFail(mEntity, new M3U8Exception(TAG, errorInfo), needRetry);
mCallback.onFail(mEntity, new M3U8Exception(TAG, errorInfo), needRetry);
}
/**
@@ -364,7 +372,7 @@ final public class M3U8InfoThread implements Runnable {
if (keyUrlConverter != null) {
keyUrl = keyUrlConverter.convert(keyUrl);
}
if (TextUtils.isEmpty(keyUrl)){
if (TextUtils.isEmpty(keyUrl)) {
ALog.e(TAG, "m3u8密钥key url 为空");
return;
}

View File

@@ -29,7 +29,7 @@ import com.arialyy.aria.util.CommonUtil;
/**
* 下载监听类
*/
public class M3U8Listener extends BaseDListener implements IDLoadListener {
public final class M3U8Listener extends BaseDListener implements IDLoadListener {
public M3U8Listener(AbsTask<DTaskWrapper> task, Handler outHandler) {
super(task, outHandler);

View File

@@ -33,10 +33,10 @@ import java.util.ArrayList;
* @Author lyy
* @Date 2019-09-24
*/
public class M3U8RecordHandler extends RecordHandler {
public final class M3U8RecordHandler extends RecordHandler {
private M3U8TaskOption mOption;
M3U8RecordHandler(DTaskWrapper wrapper) {
public M3U8RecordHandler(DTaskWrapper wrapper) {
super(wrapper);
mOption = (M3U8TaskOption) wrapper.getM3u8Option();
}
@@ -64,7 +64,7 @@ public class M3U8RecordHandler extends RecordHandler {
// 重新下载所有切片
boolean reDownload =
(m3U8Entity.getPeerNum() <= 0 || (mOption.isGenerateIndexFile() && !new File(
String.format(M3U8InfoThread.M3U8_INDEX_FORMAT, getEntity().getFilePath())).exists()));
String.format(M3U8InfoTask.M3U8_INDEX_FORMAT, getEntity().getFilePath())).exists()));
for (ThreadRecord record : mTaskRecord.threadRecords) {
File temp = new File(BaseM3U8Loader.getTsFilePath(cacheDir, record.threadId));

View File

@@ -27,7 +27,7 @@ import java.util.List;
/**
* m3u8任务配信息
*/
public class M3U8TaskOption implements ITaskOption {
public final class M3U8TaskOption implements ITaskOption {
/**
* 所有ts文件的下载地址

View File

@@ -43,7 +43,7 @@ import java.util.Set;
/**
* Created by lyy on 2017/1/18. 下载线程
*/
public class M3U8ThreadTaskAdapter extends AbsThreadTaskAdapter {
public final class M3U8ThreadTaskAdapter extends AbsThreadTaskAdapter {
private final String TAG = "M3U8ThreadTask";
private HttpTaskOption mHttpTaskOption;

View File

@@ -0,0 +1,158 @@
/*
* Copyright (C) 2016 AriaLyy(https://github.com/AriaLyy/Aria)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.arialyy.aria.m3u8.live;
import android.os.Handler;
import android.os.Looper;
import android.os.Message;
import com.arialyy.aria.core.TaskRecord;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.inf.IThreadStateManager;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.listener.ISchedulers;
import com.arialyy.aria.core.loader.ILoaderVisitor;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
import java.nio.charset.Charset;
import static com.arialyy.aria.m3u8.M3U8InfoTask.M3U8_INDEX_FORMAT;
final class LiveStateManager implements IThreadStateManager {
private final String TAG = CommonUtil.getClassName(getClass());
private M3U8Listener mListener;
private long mProgress; //当前总进度
private Looper mLooper;
private DTaskWrapper mTaskWrapper;
private M3U8TaskOption mM3U8Option;
private FileOutputStream mIndexFos;
private M3U8LiveLoader mLoader;
/**
* @param listener 任务事件
*/
LiveStateManager(DTaskWrapper wrapper, IEventListener listener) {
mTaskWrapper = wrapper;
mListener = (M3U8Listener) listener;
mM3U8Option = (M3U8TaskOption) mTaskWrapper.getM3u8Option();
}
private Handler.Callback mCallback = new Handler.Callback() {
@Override public boolean handleMessage(Message msg) {
int peerIndex = msg.getData().getInt(ISchedulers.DATA_M3U8_PEER_INDEX);
switch (msg.what) {
case STATE_STOP:
if (mLoader.isBreak()) {
ALog.d(TAG, "任务停止");
quitLooper();
}
break;
case STATE_CANCEL:
if (mLoader.isBreak()) {
ALog.d(TAG, "任务取消");
quitLooper();
}
break;
case STATE_COMPLETE:
mLoader.notifyLock(true, peerIndex);
if (mM3U8Option.isGenerateIndexFile() && !mLoader.isBreak()) {
addExtInf(mLoader.getCurExtInfo().url, mLoader.getCurExtInfo().extInf);
}
mListener.onPeerComplete(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
break;
case STATE_RUNNING:
mProgress += (long) msg.obj;
break;
case STATE_FAIL:
mLoader.notifyLock(false, peerIndex);
mListener.onPeerFail(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
break;
}
return true;
}
};
void setLoader(M3U8LiveLoader loader) {
mLoader = loader;
}
/**
* 退出looper循环
*/
private void quitLooper() {
ALog.d(TAG, "quitLooper");
mLooper.quit();
if (mIndexFos != null) {
try {
mIndexFos.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
/**
* 给索引文件添加extInfo信息
*/
private void addExtInf(String url, String extInf) {
File indexFile =
new File(String.format(M3U8_INDEX_FORMAT, mTaskWrapper.getEntity().getFilePath()));
if (!indexFile.exists()) {
ALog.e(TAG, String.format("索引文件【%s】不存在添加peer的extInf失败", indexFile.getPath()));
return;
}
try {
if (mIndexFos == null) {
mIndexFos = new FileOutputStream(indexFile, true);
}
mIndexFos.write(extInf.concat("\r\n").getBytes(Charset.forName("UTF-8")));
mIndexFos.write(url.concat("\r\n").getBytes(Charset.forName("UTF-8")));
} catch (IOException e) {
e.printStackTrace();
}
}
@Override public boolean isFail() {
return false;
}
@Override public boolean isComplete() {
return false;
}
@Override public long getCurrentProgress() {
return mProgress;
}
@Override public void setLooper(TaskRecord taskRecord, Looper looper) {
mLooper = looper;
}
@Override public Handler.Callback getHandlerCallback() {
return mCallback;
}
@Override public void accept(ILoaderVisitor visitor) {
visitor.addComponent(this);
}
}

View File

@@ -17,42 +17,46 @@ package com.arialyy.aria.m3u8.live;
import android.os.Handler;
import android.os.Looper;
import android.os.Message;
import com.arialyy.aria.core.TaskRecord;
import android.text.TextUtils;
import com.arialyy.aria.core.ThreadRecord;
import com.arialyy.aria.core.common.AbsEntity;
import com.arialyy.aria.core.common.CompleteInfo;
import com.arialyy.aria.core.common.SubThreadConfig;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.inf.IThreadState;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.listener.ISchedulers;
import com.arialyy.aria.core.inf.IThreadStateManager;
import com.arialyy.aria.core.loader.IInfoTask;
import com.arialyy.aria.core.loader.IRecordHandler;
import com.arialyy.aria.core.loader.IThreadTaskBuilder;
import com.arialyy.aria.core.manager.ThreadTaskManager;
import com.arialyy.aria.core.processor.ILiveTsUrlConverter;
import com.arialyy.aria.core.processor.ITsMergeHandler;
import com.arialyy.aria.core.task.ThreadTask;
import com.arialyy.aria.exception.BaseException;
import com.arialyy.aria.exception.M3U8Exception;
import com.arialyy.aria.exception.TaskException;
import com.arialyy.aria.m3u8.BaseM3U8Loader;
import com.arialyy.aria.m3u8.IdGenerator;
import com.arialyy.aria.m3u8.M3U8InfoTask;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import com.arialyy.aria.m3u8.M3U8ThreadTaskAdapter;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.FileUtil;
import java.io.File;
import java.io.FileOutputStream;
import java.io.FilenameFilter;
import java.io.IOException;
import java.nio.charset.Charset;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
import static com.arialyy.aria.m3u8.M3U8InfoThread.M3U8_INDEX_FORMAT;
/**
* M3U8点播文件下载器
*/
public class M3U8LiveLoader extends BaseM3U8Loader {
final class M3U8LiveLoader extends BaseM3U8Loader {
/**
* 最大执行数
*/
@@ -63,28 +67,34 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
private Condition mCondition = LOCK.newCondition();
private LinkedBlockingQueue<ExtInfo> mPeerQueue = new LinkedBlockingQueue<>();
private ExtInfo mCurExtInfo;
private FileOutputStream mIndexFos;
private LiveStateManager mManager;
private M3U8InfoTask mInfoTask;
private ScheduledThreadPoolExecutor mTimer;
private List<String> mPeerUrls = new ArrayList<>();
private Looper mLooper;
M3U8LiveLoader(M3U8Listener listener, DTaskWrapper wrapper) {
super(listener, wrapper);
M3U8LiveLoader(DTaskWrapper wrapper, M3U8Listener listener) {
super(wrapper, listener);
if (((M3U8TaskOption) wrapper.getM3u8Option()).isGenerateIndexFile()) {
ALog.i(TAG, "直播文件下载,创建索引文件的操作将导致只能同时下载一个切片");
EXEC_MAX_NUM = 1;
}
}
@Override protected IThreadState createStateManager(Looper looper) {
LiveStateManager manager = new LiveStateManager(looper, mListener);
mStateHandler = new Handler(looper, manager);
return manager;
ExtInfo getCurExtInfo() {
return mCurExtInfo;
}
void offerPeer(ExtInfo extInfo) {
private void offerPeer(ExtInfo extInfo) {
mPeerQueue.offer(extInfo);
}
@Override protected void handleTask() {
@Override protected void handleTask(Looper looper) {
if (isBreak()) {
return;
}
mLooper = looper;
startLoaderLiveInfo();
new Thread(new Runnable() {
@Override public void run() {
String cacheDir = getCacheDir();
@@ -99,7 +109,7 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
}
mCurExtInfo = extInfo;
ThreadTask task = createThreadTask(cacheDir, index, extInfo.url);
getTaskList().put(index, task);
getTaskList().add(task);
mFlagQueue.offer(startThreadTask(task, task.getConfig().peerIndex));
index++;
}
@@ -120,7 +130,7 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
return mTempFile.length();
}
private void notifyLock(boolean success, int peerId) {
void notifyLock(boolean success, int peerId) {
try {
LOCK.lock();
long id = mFlagQueue.take();
@@ -137,17 +147,6 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
}
}
@Override protected void onPostStop() {
super.onPostStop();
if (mIndexFos != null) {
try {
mIndexFos.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
/**
* 启动线程任务
*
@@ -155,7 +154,7 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
*/
private long startThreadTask(ThreadTask task, int indexId) {
ThreadTaskManager.getInstance().startThread(mTaskWrapper.getKey(), task);
((M3U8Listener) mListener).onPeerStart(mTaskWrapper.getKey(),
((M3U8Listener) getListener()).onPeerStart(mTaskWrapper.getKey(),
task.getConfig().tempFile.getPath(),
indexId);
return IdGenerator.getInstance().nextId();
@@ -196,7 +195,7 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
*
* @return {@code true} 合并成功,{@code false}合并失败
*/
boolean mergeFile() {
private boolean mergeFile() {
ITsMergeHandler mergeHandler = mM3U8Option.getMergeHandler();
String cacheDir = getCacheDir();
List<String> partPath = new ArrayList<>();
@@ -234,101 +233,117 @@ public class M3U8LiveLoader extends BaseM3U8Loader {
}
}
/**
* M3U8线程状态管理直播不处理停止状态、删除、失败的状态
*/
private class LiveStateManager implements IThreadState {
private final String TAG = "M3U8ThreadStateManager";
@Override public void addComponent(IRecordHandler recordHandler) {
mRecordHandler = recordHandler;
}
/**
* 任务状态回调
*/
private M3U8Listener mListener;
private long mProgress; //当前总进度
private Looper mLooper;
/**
* @param listener 任务事件
*/
LiveStateManager(Looper looper, IEventListener listener) {
mLooper = looper;
mListener = (M3U8Listener) listener;
}
/**
* 退出looper循环
*/
private void quitLooper() {
ALog.d(TAG, "quitLooper");
mLooper.quit();
}
@Override public boolean handleMessage(Message msg) {
int peerIndex = msg.getData().getInt(ISchedulers.DATA_M3U8_PEER_INDEX);
switch (msg.what) {
case STATE_STOP:
if (isBreak()) {
ALog.d(TAG, "任务停止");
quitLooper();
}
break;
case STATE_CANCEL:
if (isBreak()) {
ALog.d(TAG, "任务取消");
quitLooper();
}
break;
case STATE_COMPLETE:
notifyLock(true, peerIndex);
if (mM3U8Option.isGenerateIndexFile() && !isBreak()) {
addExtInf(mCurExtInfo.url, mCurExtInfo.extInf);
}
mListener.onPeerComplete(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
break;
case STATE_RUNNING:
mProgress += (long) msg.obj;
break;
case STATE_FAIL:
notifyLock(false, peerIndex);
mListener.onPeerFail(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
break;
@Override public void addComponent(IInfoTask infoTask) {
mInfoTask = (M3U8InfoTask) infoTask;
mInfoTask.setCallback(new IInfoTask.Callback() {
@Override public void onSucceed(String key, CompleteInfo info) {
ALog.d(TAG, "更新直播的m3u8文件");
}
return false;
}
/**
* 给索引文件添加extInfo信息
*/
private void addExtInf(String url, String extInf) {
File indexFile =
new File(String.format(M3U8_INDEX_FORMAT, getEntity().getFilePath()));
if (!indexFile.exists()) {
ALog.e(TAG, String.format("索引文件【%s】不存在添加peer的extInf失败", indexFile.getPath()));
return;
@Override public void onFail(AbsEntity entity, BaseException e, boolean needRetry) {
}
try {
if (mIndexFos == null) {
mIndexFos = new FileOutputStream(indexFile, true);
});
mInfoTask.setOnGetPeerCallback(new M3U8InfoTask.OnGetLivePeerCallback() {
@Override public void onGetPeer(String url, String extInf) {
if (mPeerUrls.contains(url)) {
return;
}
mIndexFos.write(extInf.concat("\r\n").getBytes(Charset.forName("UTF-8")));
mIndexFos.write(url.concat("\r\n").getBytes(Charset.forName("UTF-8")));
} catch (IOException e) {
e.printStackTrace();
mPeerUrls.add(url);
ILiveTsUrlConverter converter = mM3U8Option.getLiveTsUrlConverter();
if (converter != null) {
if (TextUtils.isEmpty(mM3U8Option.getBandWidthUrl())) {
url = converter.convert(getEntity().getUrl(), url);
} else {
url = converter.convert(mM3U8Option.getBandWidthUrl(), url);
}
}
if (TextUtils.isEmpty(url) || !url.startsWith("http")) {
fail(new M3U8Exception(TAG, String.format("ts地址错误url%s", url)), false);
return;
}
offerPeer(new M3U8LiveLoader.ExtInfo(url, extInf));
}
});
}
private void fail(BaseException e, boolean needRetry) {
getListener().onFail(needRetry, e);
handleComplete();
}
private void handleComplete() {
if (mInfoTask != null) {
mInfoTask.setStop(true);
closeInfoTimer();
if (mM3U8Option.isGenerateIndexFile()) {
if (generateIndexFile(true)) {
getListener().onComplete();
} else {
getListener().onFail(false, new TaskException(TAG, "创建索引文件失败"));
}
} else if (mM3U8Option.isMergeFile()) {
if (mergeFile()) {
getListener().onComplete();
} else {
getListener().onFail(false, new M3U8Exception(TAG, "合并文件失败"));
}
} else {
getListener().onComplete();
}
}
}
@Override public boolean isFail() {
return false;
/**
* 开始循环加载m3u8信息
*/
private void startLoaderLiveInfo() {
mTimer = new ScheduledThreadPoolExecutor(1);
mTimer.scheduleWithFixedDelay(new Runnable() {
@Override public void run() {
mInfoTask.run();
}
}, 0, mM3U8Option.getLiveUpdateInterval(), TimeUnit.MILLISECONDS);
}
private void closeInfoTimer() {
if (mTimer != null && !mTimer.isShutdown()) {
mTimer.shutdown();
}
}
@Override public boolean isComplete() {
return false;
/**
* 需要在{@link #addComponent(IRecordHandler)} 后调用
*/
@Override public void addComponent(IThreadStateManager threadState) {
mManager = (LiveStateManager) threadState;
mManager.setLooper(mRecordHandler.getRecord(0), mLooper);
mManager.setLoader(this);
mStateHandler = new Handler(mLooper, mManager.getHandlerCallback());
}
/**
* @deprecated m3u8 不需要实现这个
*/
@Deprecated
@Override public void addComponent(IThreadTaskBuilder builder) {
}
@Override protected void checkComponent() {
if (mRecordHandler == null) {
throw new NullPointerException("任务记录组件为空");
}
@Override public long getCurrentProgress() {
return mProgress;
if (mInfoTask == null) {
throw new NullPointerException(("文件信息组件为空"));
}
if (mManager == null) {
throw new NullPointerException("任务状态管理组件为空");
}
}

View File

@@ -15,30 +15,17 @@
*/
package com.arialyy.aria.m3u8.live;
import android.text.TextUtils;
import com.arialyy.aria.core.common.AbsEntity;
import com.arialyy.aria.core.common.CompleteInfo;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.inf.OnFileInfoCallback;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.loader.AbsLoader;
import com.arialyy.aria.core.loader.AbsNormalLoaderUtil;
import com.arialyy.aria.core.processor.ILiveTsUrlConverter;
import com.arialyy.aria.core.loader.LoaderStructure;
import com.arialyy.aria.core.wrapper.AbsTaskWrapper;
import com.arialyy.aria.exception.BaseException;
import com.arialyy.aria.exception.M3U8Exception;
import com.arialyy.aria.exception.TaskException;
import com.arialyy.aria.http.HttpTaskOption;
import com.arialyy.aria.m3u8.M3U8InfoThread;
import com.arialyy.aria.m3u8.M3U8InfoTask;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8RecordHandler;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import com.arialyy.aria.util.ALog;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import com.arialyy.aria.util.CommonUtil;
/**
* M3U8直播文件下载工具对于直播来说需要定时更新m3u8文件
@@ -50,12 +37,7 @@ import java.util.concurrent.TimeUnit;
* 5、不处理直播切片下载失败的状态
*/
public class M3U8LiveUtil extends AbsNormalLoaderUtil {
private final String TAG = "M3U8LiveDownloadUtil";
private M3U8InfoThread mInfoThread;
private ScheduledThreadPoolExecutor mTimer;
private ExecutorService mInfoPool = Executors.newCachedThreadPool();
private List<String> mPeerUrls = new ArrayList<>();
private M3U8TaskOption mM3U8Option;
private final String TAG = CommonUtil.getClassName(getClass());
public M3U8LiveUtil(AbsTaskWrapper wrapper, IEventListener listener) {
super(wrapper, listener);
@@ -65,117 +47,20 @@ public class M3U8LiveUtil extends AbsNormalLoaderUtil {
return (DTaskWrapper) super.getTaskWrapper();
}
@Override protected AbsLoader createLoader() {
@Override public M3U8LiveLoader getLoader() {
getTaskWrapper().generateM3u8Option(M3U8TaskOption.class);
getTaskWrapper().generateTaskOption(HttpTaskOption.class);
mM3U8Option = (M3U8TaskOption) getTaskWrapper().getM3u8Option();
return new M3U8LiveLoader((M3U8Listener) getListener(), getTaskWrapper());
return
mLoader == null ? new M3U8LiveLoader(getTaskWrapper(), (M3U8Listener) getListener())
: (M3U8LiveLoader) mLoader;
}
@Override protected Runnable createInfoThread() {
return null;
}
private Runnable createLiveInfoThread() {
M3U8InfoThread infoThread =
new M3U8InfoThread(getTaskWrapper(), new OnFileInfoCallback() {
@Override public void onComplete(String key, CompleteInfo info) {
ALog.d(TAG, "更新直播的m3u8文件");
}
@Override public void onFail(AbsEntity entity, BaseException e, boolean needRetry) {
fail(e, needRetry);
}
});
infoThread.setOnGetPeerCallback(new M3U8InfoThread.OnGetLivePeerCallback() {
@Override public void onGetPeer(String url, String extInf) {
if (mPeerUrls.contains(url)) {
return;
}
mPeerUrls.add(url);
ILiveTsUrlConverter converter = mM3U8Option.getLiveTsUrlConverter();
if (converter != null) {
if (TextUtils.isEmpty(mM3U8Option.getBandWidthUrl())) {
url = converter.convert(getTaskWrapper().getEntity().getUrl(), url);
} else {
url = converter.convert(mM3U8Option.getBandWidthUrl(), url);
}
}
if (TextUtils.isEmpty(url) || !url.startsWith("http")) {
fail(new M3U8Exception(TAG, String.format("ts地址错误url%s", url)), false);
return;
}
getLoader().offerPeer(new M3U8LiveLoader.ExtInfo(url, extInf));
}
});
return infoThread;
}
@Override protected void onCancel() {
super.onCancel();
if (mInfoThread != null) {
mInfoThread.setStop(true);
}
}
/**
* 对于直播来说是没有停止的,停止就代表完成
*/
@Override protected void onStop() {
super.onStop();
handleComplete();
}
private void handleComplete() {
if (mInfoThread != null) {
mInfoThread.setStop(true);
closeTimer();
if (((M3U8TaskOption) getTaskWrapper().getM3u8Option()).isGenerateIndexFile()) {
if (getLoader().generateIndexFile(true)) {
getListener().onComplete();
} else {
getListener().onFail(false, new TaskException(TAG, "创建索引文件失败"));
}
} else if (mM3U8Option.isMergeFile()) {
if (getLoader().mergeFile()) {
getListener().onComplete();
} else {
getListener().onFail(false, new M3U8Exception(TAG, "合并文件失败"));
}
} else {
getListener().onComplete();
}
}
}
@Override protected void onStart() {
super.onStart();
startTimer();
}
private void startTimer() {
mTimer = new ScheduledThreadPoolExecutor(1);
mTimer.scheduleWithFixedDelay(new Runnable() {
@Override public void run() {
mInfoThread = (M3U8InfoThread) createLiveInfoThread();
mInfoPool.execute(mInfoThread);
}
}, 0, mM3U8Option.getLiveUpdateInterval(), TimeUnit.MILLISECONDS);
getLoader().start();
}
private void closeTimer() {
if (mTimer != null && !mTimer.isShutdown()) {
mTimer.shutdown();
}
}
@Override protected void fail(BaseException e, boolean needRetry) {
super.fail(e, needRetry);
handleComplete();
}
@Override public M3U8LiveLoader getLoader() {
return (M3U8LiveLoader) super.getLoader();
@Override public LoaderStructure BuildLoaderStructure() {
LoaderStructure structure = new LoaderStructure();
structure.addComponent(new M3U8RecordHandler(getTaskWrapper()))
.addComponent(new M3U8InfoTask(getTaskWrapper()))
.addComponent(new LiveStateManager(getTaskWrapper(), getListener()));
structure.accept(getLoader());
return structure;
}
}

View File

@@ -15,35 +15,36 @@
*/
package com.arialyy.aria.m3u8.vod;
import android.os.Bundle;
import android.os.Handler;
import android.os.Looper;
import android.os.Message;
import android.text.TextUtils;
import android.util.SparseArray;
import com.arialyy.aria.core.TaskRecord;
import com.arialyy.aria.core.ThreadRecord;
import com.arialyy.aria.core.common.AbsEntity;
import com.arialyy.aria.core.common.CompleteInfo;
import com.arialyy.aria.core.common.SubThreadConfig;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.event.Event;
import com.arialyy.aria.core.event.EventMsgUtil;
import com.arialyy.aria.core.event.PeerIndexEvent;
import com.arialyy.aria.core.inf.IThreadState;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.listener.ISchedulers;
import com.arialyy.aria.core.inf.IThreadStateManager;
import com.arialyy.aria.core.loader.IInfoTask;
import com.arialyy.aria.core.loader.IRecordHandler;
import com.arialyy.aria.core.loader.IThreadTaskBuilder;
import com.arialyy.aria.core.manager.ThreadTaskManager;
import com.arialyy.aria.core.processor.ITsMergeHandler;
import com.arialyy.aria.core.processor.IVodTsUrlConverter;
import com.arialyy.aria.core.task.ThreadTask;
import com.arialyy.aria.exception.BaseException;
import com.arialyy.aria.exception.TaskException;
import com.arialyy.aria.exception.M3U8Exception;
import com.arialyy.aria.m3u8.BaseM3U8Loader;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import com.arialyy.aria.m3u8.M3U8ThreadTaskAdapter;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import com.arialyy.aria.util.FileUtil;
import java.io.File;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ExecutorService;
@@ -55,7 +56,7 @@ import java.util.concurrent.locks.ReentrantLock;
/**
* M3U8点播文件下载器
*/
public class M3U8VodLoader extends BaseM3U8Loader {
final class M3U8VodLoader extends BaseM3U8Loader {
/**
* 最大执行数
*/
@@ -81,9 +82,10 @@ public class M3U8VodLoader extends BaseM3U8Loader {
private ExecutorService mJumpThreadPool;
private Thread jumpThread = null;
private M3U8TaskOption mM3U8Option;
private Looper mLooper;
M3U8VodLoader(M3U8Listener listener, DTaskWrapper wrapper) {
super(listener, wrapper);
M3U8VodLoader(DTaskWrapper wrapper, M3U8Listener listener) {
super(wrapper, listener);
mM3U8Option = (M3U8TaskOption) wrapper.getM3u8Option();
mFlagQueue = new ArrayBlockingQueue<>(mM3U8Option.getMaxTsQueueNum());
EXEC_MAX_NUM = mM3U8Option.getMaxTsQueueNum();
@@ -91,10 +93,37 @@ public class M3U8VodLoader extends BaseM3U8Loader {
EventMsgUtil.getDefault().register(this);
}
@Override protected IThreadState createStateManager(Looper looper) {
mManager = new VodStateManager(looper, mRecord, mListener);
mStateHandler = new Handler(looper, mManager);
return mManager;
@Override protected M3U8Listener getListener() {
return (M3U8Listener) super.getListener();
}
SparseArray<ThreadRecord> getBeforePeer() {
return mBeforePeer;
}
int getCompleteNum() {
return mCompleteNum;
}
void setCompleteNum(int mCompleteNum) {
this.mCompleteNum = mCompleteNum;
}
int getCurrentFlagSize() {
mCurrentFlagSize = mFlagQueue.size();
return mCurrentFlagSize;
}
void setCurrentFlagSize(int currentFlagSize) {
mCurrentFlagSize = currentFlagSize;
}
boolean isJump() {
return isJump;
}
File getTempFile() {
return mTempFile;
}
@Override public void onDestroy() {
@@ -115,7 +144,15 @@ public class M3U8VodLoader extends BaseM3U8Loader {
return super.isBreak() || isDestroy;
}
@Override protected void handleTask() {
@Override protected void handleTask(Looper looper) {
if (isBreak()) {
return;
}
mLooper = looper;
mInfoTask.run();
}
private void startThreadTask() {
Thread th = new Thread(new Runnable() {
@Override public void run() {
while (!isBreak()) {
@@ -193,7 +230,7 @@ public class M3U8VodLoader extends BaseM3U8Loader {
*/
private void addTaskToQueue(ThreadRecord tr) throws InterruptedException {
ThreadTask task = createThreadTask(mCacheDir, tr, tr.threadId);
getTaskList().put(tr.threadId, task);
getTaskList().add(task);
getEntity().getM3U8Entity().setPeerIndex(tr.threadId);
TempFlag flag = startThreadTask(task, tr.threadId);
if (flag != null) {
@@ -222,10 +259,10 @@ public class M3U8VodLoader extends BaseM3U8Loader {
}
mManager.updateStateCount();
if (mCompleteNum <= 0) {
mListener.onStart(0);
getListener().onStart(0);
} else {
int percent = mCompleteNum * 100 / mRecord.threadRecords.size();
mListener.onResume(percent);
getListener().onResume(percent);
}
}
@@ -344,7 +381,7 @@ public class M3U8VodLoader extends BaseM3U8Loader {
/**
* 从指定位置恢复任务
*/
private synchronized void resumeTask() {
synchronized void resumeTask() {
if (isBreak()) {
ALog.e(TAG, "任务已停止,恢复任务失败");
return;
@@ -388,11 +425,7 @@ public class M3U8VodLoader extends BaseM3U8Loader {
}
}
private M3U8Listener getListener() {
return (M3U8Listener) mListener;
}
private void notifyWaitLock(boolean isComplete) {
void notifyWaitLock(boolean isComplete) {
try {
LOCK.lock();
if (isComplete) {
@@ -449,236 +482,87 @@ public class M3U8VodLoader extends BaseM3U8Loader {
return threadTask;
}
@Override public void addComponent(IRecordHandler recordHandler) {
mRecordHandler = recordHandler;
mRecord = mRecordHandler.getRecord(0);
}
@Override public void addComponent(IInfoTask infoTask) {
mInfoTask = infoTask;
final List<String> urls = new ArrayList<>();
mInfoTask.setCallback(new IInfoTask.Callback() {
@Override public void onSucceed(String key, CompleteInfo info) {
IVodTsUrlConverter converter = mM3U8Option.getVodUrlConverter();
if (converter != null) {
if (TextUtils.isEmpty(mM3U8Option.getBandWidthUrl())) {
urls.addAll(
converter.convert(getEntity().getUrl(), (List<String>) info.obj));
} else {
urls.addAll(
converter.convert(mM3U8Option.getBandWidthUrl(), (List<String>) info.obj));
}
} else {
urls.addAll((Collection<? extends String>) info.obj);
}
if (urls.isEmpty()) {
fail(new M3U8Exception(TAG, "获取地址失败"), false);
return;
} else if (!urls.get(0).startsWith("http")) {
fail(new M3U8Exception(TAG, "地址错误请使用IVodTsUrlConverter处理你的url信息"), false);
return;
}
mM3U8Option.setUrls(urls);
if (isStop) {
getListener().onStop(getEntity().getCurrentProgress());
} else if (isCancel) {
getListener().onCancel();
} else {
startThreadTask();
}
}
@Override public void onFail(AbsEntity entity, BaseException e, boolean needRetry) {
fail(e, needRetry);
}
});
}
protected void fail(BaseException e, boolean needRetry) {
if (isBreak()) {
return;
}
getListener().onFail(needRetry, e);
onDestroy();
}
/**
* M3U8线程状态管理
* 需要在 {@link #addComponent(IRecordHandler)}后调用
*/
private class VodStateManager implements IThreadState {
private final String TAG = CommonUtil.getClassName(VodStateManager.class);
@Override public void addComponent(IThreadStateManager threadState) {
mManager = (VodStateManager) threadState;
mStateHandler = new Handler(mLooper, mManager.getHandlerCallback());
mManager.setVodLoader(this);
mManager.setLooper(mRecord, mLooper);
}
/**
* 任务状态回调
*/
private IEventListener listener;
private int startThreadNum; // 启动的线程总数
private int cancelNum = 0; // 已经取消的线程的数
private int stopNum = 0; // 已经停止的线程数
private int failNum = 0; // 失败的线程数
private long percent; //当前总进度,百分比进度
private long progress;
private TaskRecord taskRecord; // 任务记录
private Looper looper;
/**
* m3u8 不需要实现这个
*/
@Deprecated
@Override public void addComponent(IThreadTaskBuilder builder) {
/**
* @param taskRecord 任务记录
* @param listener 任务事件
*/
VodStateManager(Looper looper, TaskRecord taskRecord, IEventListener listener) {
this.looper = looper;
this.taskRecord = taskRecord;
for (ThreadRecord record : taskRecord.threadRecords) {
if (!record.isComplete) {
startThreadNum++;
}
}
this.listener = listener;
progress = getEntity().getCurrentProgress();
}
@Override protected void checkComponent() {
if (mRecordHandler == null) {
throw new NullPointerException("任务记录组件为空");
}
private void updateStateCount() {
cancelNum = 0;
stopNum = 0;
failNum = 0;
if (mInfoTask == null) {
throw new NullPointerException(("文件信息组件为空"));
}
/**
* 退出looper循环
*/
private void quitLooper() {
ALog.d(TAG, "quitLooper");
looper.quit();
}
@Override public boolean handleMessage(Message msg) {
int peerIndex = msg.getData().getInt(ISchedulers.DATA_M3U8_PEER_INDEX);
switch (msg.what) {
case STATE_STOP:
stopNum++;
removeSignThread((ThreadTask) msg.obj);
// 处理跳转位置后,恢复任务
if (isJump && (stopNum == mCurrentFlagSize || mCurrentFlagSize == 0) && !isBreak()) {
resumeTask();
return true;
}
if (isBreak()) {
ALog.d(TAG, String.format("vod任务【%s】停止", mTempFile.getName()));
quitLooper();
}
break;
case STATE_CANCEL:
cancelNum++;
removeSignThread((ThreadTask) msg.obj);
if (isBreak()) {
ALog.d(TAG, String.format("vod任务【%s】取消", mTempFile.getName()));
quitLooper();
}
break;
case STATE_FAIL:
failNum++;
for (ThreadRecord tr : mRecord.threadRecords) {
if (tr.threadId == peerIndex) {
mBeforePeer.put(peerIndex, tr);
break;
}
}
getListener().onPeerFail(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
if (isFail()) {
ALog.d(TAG, String.format("vod任务【%s】失败", mTempFile.getName()));
Bundle b = msg.getData();
listener.onFail(b.getBoolean(KEY_RETRY, true),
(BaseException) b.getSerializable(KEY_ERROR_INFO));
quitLooper();
}
break;
case STATE_COMPLETE:
if (isBreak()) {
quitLooper();
}
mCompleteNum++;
// 正在切换位置时,切片完成,队列减小
if (isJump) {
mCurrentFlagSize--;
if (mCurrentFlagSize < 0) {
mCurrentFlagSize = 0;
}
}
removeSignThread((ThreadTask) msg.obj);
getListener().onPeerComplete(mTaskWrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
handlerPercent();
if (!isJump) {
notifyWaitLock(true);
}
if (isComplete()) {
ALog.d(TAG, String.format(
"startThreadNum = %s, stopNum = %s, cancelNum = %s, failNum = %s, completeNum = %s, flagQueueSize = %s",
startThreadNum, stopNum, cancelNum, failNum, mCompleteNum, mFlagQueue.size()));
ALog.d(TAG, String.format("vod任务【%s】完成", mTempFile.getName()));
if (mM3U8Option.isGenerateIndexFile()) {
if (generateIndexFile(false)) {
listener.onComplete();
} else {
listener.onFail(false, new TaskException(TAG, "创建索引文件失败"));
}
} else if (mM3U8Option.isMergeFile()) {
if (mergeFile()) {
listener.onComplete();
} else {
listener.onFail(false, null);
}
} else {
listener.onComplete();
}
quitLooper();
}
break;
case STATE_RUNNING:
progress += (long) msg.obj;
break;
}
return true;
}
private void removeSignThread(ThreadTask threadTask) {
int index = getTaskList().indexOfValue(threadTask);
if (index != -1) {
getTaskList().removeAt(index);
}
ThreadTaskManager.getInstance().removeSingleTaskThread(mTaskWrapper.getKey(), threadTask);
}
/**
* 设置进度
*/
private void handlerPercent() {
int completeNum = mM3U8Option.getCompleteNum();
completeNum++;
mM3U8Option.setCompleteNum(completeNum);
int percent = completeNum * 100 / taskRecord.threadRecords.size();
getEntity().setPercent(percent);
getEntity().update();
this.percent = percent;
}
@Override public boolean isFail() {
printInfo("isFail");
return failNum != 0 && failNum == mFlagQueue.size() && !isJump;
}
@Override public boolean isComplete() {
if (mM3U8Option.isIgnoreFailureTs()) {
return mCompleteNum + failNum >= taskRecord.threadRecords.size() && !isJump;
} else {
return mCompleteNum == taskRecord.threadRecords.size() && !isJump;
}
}
@Override public long getCurrentProgress() {
return progress;
}
private void printInfo(String tag) {
if (false) {
ALog.d(tag, String.format(
"startThreadNum = %s, stopNum = %s, cancelNum = %s, failNum = %s, completeNum = %s, flagQueueSize = %s",
startThreadNum, stopNum, cancelNum, failNum, mCompleteNum, mFlagQueue.size()));
}
}
/**
* 合并文件
*
* @return {@code true} 合并成功,{@code false}合并失败
*/
private boolean mergeFile() {
ITsMergeHandler mergeHandler = mM3U8Option.getMergeHandler();
String cacheDir = getCacheDir();
List<String> partPath = new ArrayList<>();
for (ThreadRecord tr : taskRecord.threadRecords) {
partPath.add(BaseM3U8Loader.getTsFilePath(cacheDir, tr.threadId));
}
boolean isSuccess;
if (mergeHandler != null) {
isSuccess = mergeHandler.merge(getEntity().getM3U8Entity(), partPath);
if (mergeHandler.getClass().isAnonymousClass()) {
mM3U8Option.setMergeHandler(null);
}
} else {
isSuccess = FileUtil.mergeFile(taskRecord.filePath, partPath);
}
if (isSuccess) {
// 合并成功,删除缓存文件
File[] files = new File(cacheDir).listFiles();
for (File f : files) {
if (f.exists()) {
f.delete();
}
}
File cDir = new File(cacheDir);
if (cDir.exists()) {
cDir.delete();
}
return true;
} else {
ALog.e(TAG, "合并失败");
return false;
}
if (mManager == null) {
throw new NullPointerException("任务状态管理组件为空");
}
}

View File

@@ -15,25 +15,17 @@
*/
package com.arialyy.aria.m3u8.vod;
import android.text.TextUtils;
import com.arialyy.aria.core.common.AbsEntity;
import com.arialyy.aria.core.common.CompleteInfo;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.inf.OnFileInfoCallback;
import com.arialyy.aria.core.listener.IEventListener;
import com.arialyy.aria.core.loader.AbsLoader;
import com.arialyy.aria.core.loader.AbsNormalLoader;
import com.arialyy.aria.core.loader.AbsNormalLoaderUtil;
import com.arialyy.aria.core.processor.IVodTsUrlConverter;
import com.arialyy.aria.core.loader.LoaderStructure;
import com.arialyy.aria.core.wrapper.AbsTaskWrapper;
import com.arialyy.aria.exception.BaseException;
import com.arialyy.aria.exception.M3U8Exception;
import com.arialyy.aria.http.HttpTaskOption;
import com.arialyy.aria.m3u8.M3U8InfoThread;
import com.arialyy.aria.m3u8.M3U8InfoTask;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8RecordHandler;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
/**
* M3U8点播文件下载工具
@@ -43,10 +35,7 @@ import java.util.List;
* 3、完成所有分片下载后合并ts文件
* 4、删除该隐藏文件夹
*/
public class M3U8VodUtil extends AbsNormalLoaderUtil {
private List<String> mUrls = new ArrayList<>();
private M3U8TaskOption mM3U8Option;
public final class M3U8VodUtil extends AbsNormalLoaderUtil {
public M3U8VodUtil(AbsTaskWrapper wrapper, IEventListener listener) {
super(wrapper, listener);
@@ -56,48 +45,19 @@ public class M3U8VodUtil extends AbsNormalLoaderUtil {
return (DTaskWrapper) super.getTaskWrapper();
}
@Override protected AbsLoader createLoader() {
@Override public AbsNormalLoader getLoader() {
getTaskWrapper().generateM3u8Option(M3U8TaskOption.class);
getTaskWrapper().generateTaskOption(HttpTaskOption.class);
mM3U8Option = (M3U8TaskOption) getTaskWrapper().getM3u8Option();
return new M3U8VodLoader((M3U8Listener) getListener(), getTaskWrapper());
return mLoader == null ? new M3U8VodLoader(getTaskWrapper(), (M3U8Listener) getListener())
: mLoader;
}
@Override protected Runnable createInfoThread() {
return new M3U8InfoThread(getTaskWrapper(), new OnFileInfoCallback() {
@Override public void onComplete(String key, CompleteInfo info) {
IVodTsUrlConverter converter = mM3U8Option.getVodUrlConverter();
if (converter != null) {
if (TextUtils.isEmpty(mM3U8Option.getBandWidthUrl())) {
mUrls.addAll(
converter.convert(getTaskWrapper().getEntity().getUrl(), (List<String>) info.obj));
} else {
mUrls.addAll(
converter.convert(mM3U8Option.getBandWidthUrl(), (List<String>) info.obj));
}
} else {
mUrls.addAll((Collection<? extends String>) info.obj);
}
if (mUrls.isEmpty()) {
fail(new M3U8Exception(TAG, "获取地址失败"), false);
return;
} else if (!mUrls.get(0).startsWith("http")) {
fail(new M3U8Exception(TAG, "地址错误请使用IVodTsUrlConverter处理你的url信息"), false);
return;
}
mM3U8Option.setUrls(mUrls);
if (isStop()) {
getListener().onStop(getTaskWrapper().getEntity().getCurrentProgress());
} else if (isCancel()) {
getListener().onCancel();
} else {
getLoader().start();
}
}
@Override public void onFail(AbsEntity entity, BaseException e, boolean needRetry) {
fail(e, needRetry);
}
});
@Override public LoaderStructure BuildLoaderStructure() {
LoaderStructure structure = new LoaderStructure();
structure.addComponent(new M3U8RecordHandler(getTaskWrapper()))
.addComponent(new M3U8InfoTask(getTaskWrapper()))
.addComponent(new VodStateManager(getTaskWrapper(), (M3U8Listener) getListener()));
structure.accept(getLoader());
return structure;
}
}

View File

@@ -0,0 +1,303 @@
/*
* Copyright (C) 2016 AriaLyy(https://github.com/AriaLyy/Aria)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.arialyy.aria.m3u8.vod;
import android.os.Bundle;
import android.os.Handler;
import android.os.Looper;
import android.os.Message;
import com.arialyy.aria.core.TaskRecord;
import com.arialyy.aria.core.ThreadRecord;
import com.arialyy.aria.core.download.DTaskWrapper;
import com.arialyy.aria.core.download.DownloadEntity;
import com.arialyy.aria.core.inf.IThreadStateManager;
import com.arialyy.aria.core.listener.ISchedulers;
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.BaseException;
import com.arialyy.aria.exception.TaskException;
import com.arialyy.aria.m3u8.BaseM3U8Loader;
import com.arialyy.aria.m3u8.M3U8Listener;
import com.arialyy.aria.m3u8.M3U8TaskOption;
import com.arialyy.aria.util.ALog;
import com.arialyy.aria.util.CommonUtil;
import com.arialyy.aria.util.FileUtil;
import java.io.File;
import java.util.ArrayList;
import java.util.List;
/**
* m3u8 点播下载状态管理器
*/
public final class VodStateManager implements IThreadStateManager {
private final String TAG = CommonUtil.getClassName(getClass());
private M3U8Listener listener;
private int startThreadNum; // 启动的线程总数
private int cancelNum = 0; // 已经取消的线程的数
private int stopNum = 0; // 已经停止的线程数
private int failNum = 0; // 失败的线程数
private long percent; //当前总进度,百分比进度
private long progress;
private TaskRecord taskRecord; // 任务记录
private Looper looper;
private DTaskWrapper wrapper;
private M3U8TaskOption m3U8Option;
private M3U8VodLoader loader;
/**
* @param listener 任务事件
*/
VodStateManager(DTaskWrapper wrapper, M3U8Listener listener) {
this.wrapper = wrapper;
this.listener = listener;
m3U8Option = (M3U8TaskOption) wrapper.getM3u8Option();
progress = wrapper.getEntity().getCurrentProgress();
}
private Handler.Callback callback = new Handler.Callback() {
@Override public boolean handleMessage(Message msg) {
int peerIndex = msg.getData().getInt(ISchedulers.DATA_M3U8_PEER_INDEX);
switch (msg.what) {
case STATE_STOP:
stopNum++;
removeSignThread((ThreadTask) msg.obj);
// 处理跳转位置后,恢复任务
if (loader.isJump()
&& (stopNum == loader.getCurrentFlagSize() || loader.getCurrentFlagSize() == 0)
&& !loader.isBreak()) {
loader.resumeTask();
return true;
}
if (loader.isBreak()) {
ALog.d(TAG, String.format("vod任务【%s】停止", loader.getTempFile().getName()));
quitLooper();
}
break;
case STATE_CANCEL:
cancelNum++;
removeSignThread((ThreadTask) msg.obj);
if (loader.isBreak()) {
ALog.d(TAG, String.format("vod任务【%s】取消", loader.getTempFile().getName()));
quitLooper();
}
break;
case STATE_FAIL:
failNum++;
for (ThreadRecord tr : taskRecord.threadRecords) {
if (tr.threadId == peerIndex) {
loader.getBeforePeer().put(peerIndex, tr);
break;
}
}
getListener().onPeerFail(wrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
if (isFail()) {
ALog.d(TAG, String.format("vod任务【%s】失败", loader.getTempFile().getName()));
Bundle b = msg.getData();
listener.onFail(b.getBoolean(KEY_RETRY, true),
(BaseException) b.getSerializable(KEY_ERROR_INFO));
quitLooper();
}
break;
case STATE_COMPLETE:
if (loader.isBreak()) {
quitLooper();
}
loader.setCompleteNum(loader.getCompleteNum() + 1);
// 正在切换位置时,切片完成,队列减小
if (loader.isJump()) {
loader.setCurrentFlagSize(loader.getCurrentFlagSize() - 1);
if (loader.getCurrentFlagSize() < 0) {
loader.setCurrentFlagSize(0);
}
}
removeSignThread((ThreadTask) msg.obj);
getListener().onPeerComplete(wrapper.getKey(),
msg.getData().getString(ISchedulers.DATA_M3U8_PEER_PATH), peerIndex);
handlerPercent();
if (!loader.isJump()) {
loader.notifyWaitLock(true);
}
if (isComplete()) {
ALog.d(TAG, String.format(
"startThreadNum = %s, stopNum = %s, cancelNum = %s, failNum = %s, completeNum = %s, flagQueueSize = %s",
startThreadNum, stopNum, cancelNum, failNum, loader.getCompleteNum(),
loader.getCurrentFlagSize()));
ALog.d(TAG, String.format("vod任务【%s】完成", loader.getTempFile().getName()));
if (m3U8Option.isGenerateIndexFile()) {
if (loader.generateIndexFile(false)) {
listener.onComplete();
} else {
listener.onFail(false, new TaskException(TAG, "创建索引文件失败"));
}
} else if (m3U8Option.isMergeFile()) {
if (mergeFile()) {
listener.onComplete();
} else {
listener.onFail(false, null);
}
} else {
listener.onComplete();
}
quitLooper();
}
break;
case STATE_RUNNING:
progress += (long) msg.obj;
break;
}
return true;
}
};
void updateStateCount() {
cancelNum = 0;
stopNum = 0;
failNum = 0;
}
@Override public void setLooper(TaskRecord taskRecord, Looper looper) {
this.looper = looper;
this.taskRecord = taskRecord;
for (ThreadRecord record : taskRecord.threadRecords) {
if (!record.isComplete) {
startThreadNum++;
}
}
}
@Override public Handler.Callback getHandlerCallback() {
return callback;
}
private DownloadEntity getEntity() {
return wrapper.getEntity();
}
private M3U8Listener getListener() {
return listener;
}
void setVodLoader(M3U8VodLoader loader) {
this.loader = loader;
}
/**
* 退出looper循环
*/
private void quitLooper() {
ALog.d(TAG, "quitLooper");
looper.quit();
}
private void removeSignThread(ThreadTask threadTask) {
loader.getTaskList().remove(threadTask);
ThreadTaskManager.getInstance().removeSingleTaskThread(wrapper.getKey(), threadTask);
}
/**
* 设置进度
*/
private void handlerPercent() {
int completeNum = m3U8Option.getCompleteNum();
completeNum++;
m3U8Option.setCompleteNum(completeNum);
int percent = completeNum * 100 / taskRecord.threadRecords.size();
getEntity().setPercent(percent);
getEntity().update();
this.percent = percent;
}
@Override public boolean isFail() {
printInfo("isFail");
return failNum != 0 && failNum == loader.getCurrentFlagSize() && !loader.isJump();
}
@Override public boolean isComplete() {
if (m3U8Option.isIgnoreFailureTs()) {
return loader.getCompleteNum() + failNum >= taskRecord.threadRecords.size()
&& !loader.isJump();
} else {
return loader.getCompleteNum() == taskRecord.threadRecords.size() && !loader.isJump();
}
}
@Override public long getCurrentProgress() {
return progress;
}
private void printInfo(String tag) {
if (false) {
ALog.d(tag, String.format(
"startThreadNum = %s, stopNum = %s, cancelNum = %s, failNum = %s, completeNum = %s, flagQueueSize = %s",
startThreadNum, stopNum, cancelNum, failNum, loader.getCompleteNum(),
loader.getCurrentFlagSize()));
}
}
/**
* 合并文件
*
* @return {@code true} 合并成功,{@code false}合并失败
*/
private boolean mergeFile() {
ITsMergeHandler mergeHandler = m3U8Option.getMergeHandler();
String cacheDir = loader.getCacheDir();
List<String> partPath = new ArrayList<>();
for (ThreadRecord tr : taskRecord.threadRecords) {
partPath.add(BaseM3U8Loader.getTsFilePath(cacheDir, tr.threadId));
}
boolean isSuccess;
if (mergeHandler != null) {
isSuccess = mergeHandler.merge(getEntity().getM3U8Entity(), partPath);
if (mergeHandler.getClass().isAnonymousClass()) {
m3U8Option.setMergeHandler(null);
}
} else {
isSuccess = FileUtil.mergeFile(taskRecord.filePath, partPath);
}
if (isSuccess) {
// 合并成功,删除缓存文件
File[] files = new File(cacheDir).listFiles();
for (File f : files) {
if (f.exists()) {
f.delete();
}
}
File cDir = new File(cacheDir);
if (cDir.exists()) {
cDir.delete();
}
return true;
} else {
ALog.e(TAG, "合并失败");
return false;
}
}
@Override public void accept(ILoaderVisitor visitor) {
visitor.addComponent(this);
}
}