advance.js 52 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690
  1. var session = require('./session');
  2. var fs = require('fs');
  3. var Async = require('./async');
  4. var EventProxy = require('./event').EventProxy;
  5. var util = require('./util');
  6. // 文件分块上传全过程,暴露的分块上传接口
  7. function sliceUploadFile(params, callback) {
  8. var self = this;
  9. var ep = new EventProxy();
  10. var TaskId = params.TaskId;
  11. var Bucket = params.Bucket;
  12. var Region = params.Region;
  13. var Key = params.Key;
  14. var FilePath = params.FilePath;
  15. var ChunkSize = params.ChunkSize || params.SliceSize || self.options.ChunkSize;
  16. var AsyncLimit = params.AsyncLimit;
  17. var StorageClass = params.StorageClass;
  18. var ServerSideEncryption = params.ServerSideEncryption;
  19. var FileSize;
  20. var onProgress;
  21. var onHashProgress = params.onHashProgress;
  22. // 上传过程中出现错误,返回错误
  23. ep.on('error', function (err) {
  24. if (!self._isRunningTask(TaskId)) return;
  25. err.UploadId = params.UploadData.UploadId || '';
  26. return callback(err);
  27. });
  28. // 上传分块完成,开始 uploadSliceComplete 操作
  29. ep.on('upload_complete', function (UploadCompleteData) {
  30. var _UploadCompleteData = util.extend(
  31. {
  32. UploadId: params.UploadData.UploadId || '',
  33. },
  34. UploadCompleteData
  35. );
  36. callback(null, _UploadCompleteData);
  37. });
  38. // 上传分块完成,开始 uploadSliceComplete 操作
  39. ep.on('upload_slice_complete', function (UploadData) {
  40. var metaHeaders = {};
  41. util.each(params.Headers, function (val, k) {
  42. var shortKey = k.toLowerCase();
  43. if (shortKey.indexOf('x-cos-meta-') === 0 || shortKey === 'pic-operations') {
  44. metaHeaders[k] = val;
  45. }
  46. });
  47. uploadSliceComplete.call(
  48. self,
  49. {
  50. Bucket: Bucket,
  51. Region: Region,
  52. Key: Key,
  53. UploadId: UploadData.UploadId,
  54. SliceList: UploadData.SliceList,
  55. Headers: metaHeaders,
  56. },
  57. function (err, data) {
  58. if (!self._isRunningTask(TaskId)) return;
  59. session.removeUsing(UploadData.UploadId);
  60. if (err) {
  61. onProgress(null, true);
  62. return ep.emit('error', err);
  63. }
  64. session.removeUploadId.call(self, UploadData.UploadId);
  65. onProgress({ loaded: FileSize, total: FileSize }, true);
  66. ep.emit('upload_complete', data);
  67. }
  68. );
  69. });
  70. // 获取 UploadId 完成,开始上传每个分片
  71. ep.on('get_upload_data_finish', function (UploadData) {
  72. // 处理 UploadId 缓存
  73. var uuid = session.getFileId(params.FileStat, params.ChunkSize, Bucket, Key);
  74. uuid && session.saveUploadId.call(self, uuid, UploadData.UploadId, self.options.UploadIdCacheLimit); // 缓存 UploadId
  75. session.setUsing(UploadData.UploadId); // 标记 UploadId 为正在使用
  76. // 获取 UploadId
  77. onProgress(null, true); // 任务状态开始 uploading
  78. uploadSliceList.call(
  79. self,
  80. {
  81. TaskId: TaskId,
  82. Bucket: Bucket,
  83. Region: Region,
  84. Key: Key,
  85. FilePath: FilePath,
  86. FileSize: FileSize,
  87. SliceSize: ChunkSize,
  88. AsyncLimit: AsyncLimit,
  89. ServerSideEncryption: ServerSideEncryption,
  90. UploadData: UploadData,
  91. Headers: params.Headers,
  92. onProgress: onProgress,
  93. },
  94. function (err, data) {
  95. if (!self._isRunningTask(TaskId)) return;
  96. if (err) {
  97. onProgress(null, true);
  98. return ep.emit('error', err);
  99. }
  100. ep.emit('upload_slice_complete', data);
  101. }
  102. );
  103. });
  104. // 开始获取文件 UploadId,里面会视情况计算 ETag,并比对,保证文件一致性,也优化上传
  105. ep.on('get_file_size_finish', function () {
  106. onProgress = util.throttleOnProgress.call(self, FileSize, params.onProgress);
  107. if (params.UploadData.UploadId) {
  108. ep.emit('get_upload_data_finish', params.UploadData);
  109. } else {
  110. var _params = util.extend(
  111. {
  112. TaskId: TaskId,
  113. Bucket: Bucket,
  114. Region: Region,
  115. Key: Key,
  116. Headers: params.Headers,
  117. StorageClass: StorageClass,
  118. FilePath: FilePath,
  119. FileSize: FileSize,
  120. SliceSize: ChunkSize,
  121. onHashProgress: onHashProgress,
  122. },
  123. params
  124. );
  125. getUploadIdAndPartList.call(self, _params, function (err, UploadData) {
  126. if (!self._isRunningTask(TaskId)) return;
  127. if (err) return ep.emit('error', err);
  128. params.UploadData.UploadId = UploadData.UploadId;
  129. params.UploadData.PartList = UploadData.PartList;
  130. ep.emit('get_upload_data_finish', params.UploadData);
  131. });
  132. }
  133. });
  134. // 获取上传文件大小
  135. FileSize = params.ContentLength;
  136. delete params.ContentLength;
  137. !params.Headers && (params.Headers = {});
  138. util.each(params.Headers, function (item, key) {
  139. if (key.toLowerCase() === 'content-length') {
  140. delete params.Headers[key];
  141. }
  142. });
  143. // 控制分片大小
  144. (function () {
  145. var SIZE = [1, 2, 4, 8, 16, 32, 64, 128, 256, 512, 1024, 1024 * 2, 1024 * 4, 1024 * 5];
  146. var AutoChunkSize = 1024 * 1024;
  147. for (var i = 0; i < SIZE.length; i++) {
  148. AutoChunkSize = SIZE[i] * 1024 * 1024;
  149. if (FileSize / AutoChunkSize <= self.options.MaxPartNumber) break;
  150. }
  151. params.ChunkSize = params.SliceSize = ChunkSize = Math.max(ChunkSize, AutoChunkSize);
  152. })();
  153. // 开始上传
  154. if (FileSize === 0) {
  155. params.Body = '';
  156. params.ContentLength = 0;
  157. params.SkipTask = true;
  158. self.putObject(params, callback);
  159. } else {
  160. ep.emit('get_file_size_finish');
  161. }
  162. }
  163. // 获取上传任务的 UploadId
  164. function getUploadIdAndPartList(params, callback) {
  165. var TaskId = params.TaskId;
  166. var Bucket = params.Bucket;
  167. var Region = params.Region;
  168. var Key = params.Key;
  169. var StorageClass = params.StorageClass;
  170. var self = this;
  171. // 计算 ETag
  172. var ETagMap = {};
  173. var FileSize = params.FileSize;
  174. var SliceSize = params.SliceSize;
  175. var SliceCount = Math.ceil(FileSize / SliceSize);
  176. var FinishSliceCount = 0;
  177. var FinishSize = 0;
  178. var onHashProgress = util.throttleOnProgress.call(self, FileSize, params.onHashProgress);
  179. var getChunkETag = function (PartNumber, callback) {
  180. var start = SliceSize * (PartNumber - 1);
  181. var end = Math.min(start + SliceSize, FileSize);
  182. var ChunkSize = end - start;
  183. if (ETagMap[PartNumber]) {
  184. callback(null, {
  185. PartNumber: PartNumber,
  186. ETag: ETagMap[PartNumber],
  187. Size: ChunkSize,
  188. });
  189. } else {
  190. util.fileSlice(params.FilePath, start, end, function (chunkItem) {
  191. util.getFileMd5(chunkItem, function (err, md5) {
  192. if (err) return callback(util.error(err));
  193. var ETag = '"' + md5 + '"';
  194. ETagMap[PartNumber] = ETag;
  195. FinishSliceCount += 1;
  196. FinishSize += ChunkSize;
  197. onHashProgress({ loaded: FinishSize, total: FileSize });
  198. callback(null, {
  199. PartNumber: PartNumber,
  200. ETag: ETag,
  201. Size: ChunkSize,
  202. });
  203. });
  204. });
  205. }
  206. };
  207. // 通过和文件的 md5 对比,判断 UploadId 是否可用
  208. var isAvailableUploadList = function (PartList, callback) {
  209. var PartCount = PartList.length;
  210. // 如果没有分片,通过
  211. if (PartCount === 0) {
  212. return callback(null, true);
  213. }
  214. // 检查分片数量
  215. if (PartCount > SliceCount) {
  216. return callback(null, false);
  217. }
  218. // 检查分片大小
  219. if (PartCount > 1) {
  220. var PartSliceSize = Math.max(PartList[0].Size, PartList[1].Size);
  221. if (PartSliceSize !== SliceSize) {
  222. return callback(null, false);
  223. }
  224. }
  225. // 逐个分片计算并检查 ETag 是否一致
  226. var next = function (index) {
  227. if (index < PartCount) {
  228. var Part = PartList[index];
  229. getChunkETag(Part.PartNumber, function (err, chunk) {
  230. if (chunk && chunk.ETag === Part.ETag && chunk.Size === Part.Size) {
  231. next(index + 1);
  232. } else {
  233. callback(null, false);
  234. }
  235. });
  236. } else {
  237. callback(null, true);
  238. }
  239. };
  240. next(0);
  241. };
  242. var ep = new EventProxy();
  243. ep.on('error', function (errData) {
  244. if (!self._isRunningTask(TaskId)) return;
  245. return callback(errData);
  246. });
  247. // 存在 UploadId
  248. ep.on('upload_id_available', function (UploadData) {
  249. // 转换成 map
  250. var map = {};
  251. var list = [];
  252. util.each(UploadData.PartList, function (item) {
  253. map[item.PartNumber] = item;
  254. });
  255. for (var PartNumber = 1; PartNumber <= SliceCount; PartNumber++) {
  256. var item = map[PartNumber];
  257. if (item) {
  258. item.PartNumber = PartNumber;
  259. item.Uploaded = true;
  260. } else {
  261. item = {
  262. PartNumber: PartNumber,
  263. ETag: null,
  264. Uploaded: false,
  265. };
  266. }
  267. list.push(item);
  268. }
  269. UploadData.PartList = list;
  270. callback(null, UploadData);
  271. });
  272. // 不存在 UploadId, 初始化生成 UploadId
  273. ep.on('no_available_upload_id', function () {
  274. if (!self._isRunningTask(TaskId)) return;
  275. var _params = util.extend(
  276. {
  277. Bucket: Bucket,
  278. Region: Region,
  279. Key: Key,
  280. Headers: util.clone(params.Headers),
  281. Query: util.clone(params.Query),
  282. StorageClass: StorageClass,
  283. },
  284. params
  285. );
  286. self.multipartInit(_params, function (err, data) {
  287. if (!self._isRunningTask(TaskId)) return;
  288. if (err) return ep.emit('error', err);
  289. var UploadId = data.UploadId;
  290. if (!UploadId) {
  291. return callback(util.error(new Error('no such upload id')));
  292. }
  293. ep.emit('upload_id_available', { UploadId: UploadId, PartList: [] });
  294. });
  295. });
  296. // 如果已存在 UploadId,找一个可以用的 UploadId
  297. ep.on('has_and_check_upload_id', function (UploadIdList) {
  298. // 串行地,找一个内容一致的 UploadId
  299. UploadIdList = UploadIdList.reverse();
  300. Async.eachLimit(
  301. UploadIdList,
  302. 1,
  303. function (UploadId, asyncCallback) {
  304. if (!self._isRunningTask(TaskId)) return;
  305. // 如果正在上传,跳过
  306. if (session.using[UploadId]) {
  307. asyncCallback(); // 检查下一个 UploadId
  308. return;
  309. }
  310. // 判断 UploadId 是否可用
  311. wholeMultipartListPart.call(
  312. self,
  313. {
  314. Bucket: Bucket,
  315. Region: Region,
  316. Key: Key,
  317. UploadId: UploadId,
  318. },
  319. function (err, PartListData) {
  320. if (!self._isRunningTask(TaskId)) return;
  321. if (err) {
  322. session.removeUsing(UploadId);
  323. return ep.emit('error', err);
  324. }
  325. var PartList = PartListData.PartList;
  326. PartList.forEach(function (item) {
  327. item.PartNumber *= 1;
  328. item.Size *= 1;
  329. item.ETag = item.ETag || '';
  330. });
  331. isAvailableUploadList(PartList, function (err, isAvailable) {
  332. if (!self._isRunningTask(TaskId)) return;
  333. if (err) return ep.emit('error', err);
  334. if (isAvailable) {
  335. asyncCallback({
  336. UploadId: UploadId,
  337. PartList: PartList,
  338. }); // 马上结束
  339. } else {
  340. asyncCallback(); // 检查下一个 UploadId
  341. }
  342. });
  343. }
  344. );
  345. },
  346. function (AvailableUploadData) {
  347. if (!self._isRunningTask(TaskId)) return;
  348. onHashProgress(null, true);
  349. if (AvailableUploadData && AvailableUploadData.UploadId) {
  350. ep.emit('upload_id_available', AvailableUploadData);
  351. } else {
  352. ep.emit('no_available_upload_id');
  353. }
  354. }
  355. );
  356. });
  357. // 在本地缓存找可用的 UploadId
  358. ep.on('seek_local_avail_upload_id', function (RemoteUploadIdList) {
  359. // 在本地找可用的 UploadId
  360. var uuid = session.getFileId(params.FileStat, params.ChunkSize, Bucket, Key);
  361. var LocalUploadIdList = session.getUploadIdList.call(self, uuid);
  362. if (!uuid || !LocalUploadIdList) {
  363. ep.emit('has_and_check_upload_id', RemoteUploadIdList);
  364. return;
  365. }
  366. var next = function (index) {
  367. // 如果本地找不到可用 UploadId,再一个个遍历校验远端
  368. if (index >= LocalUploadIdList.length) {
  369. ep.emit('has_and_check_upload_id', RemoteUploadIdList);
  370. return;
  371. }
  372. var UploadId = LocalUploadIdList[index];
  373. // 如果不在远端 UploadId 列表里,跳过并删除
  374. if (!util.isInArray(RemoteUploadIdList, UploadId)) {
  375. session.removeUploadId.call(self, UploadId);
  376. next(index + 1);
  377. return;
  378. }
  379. // 如果正在上传,跳过
  380. if (session.using[UploadId]) {
  381. next(index + 1);
  382. return;
  383. }
  384. // 判断 UploadId 是否存在线上
  385. wholeMultipartListPart.call(
  386. self,
  387. {
  388. Bucket: Bucket,
  389. Region: Region,
  390. Key: Key,
  391. UploadId: UploadId,
  392. },
  393. function (err, PartListData) {
  394. if (!self._isRunningTask(TaskId)) return;
  395. if (err) {
  396. // 如果 UploadId 获取会出错,跳过并删除
  397. session.removeUploadId.call(self, UploadId);
  398. next(index + 1);
  399. } else {
  400. // 找到可用 UploadId
  401. ep.emit('upload_id_available', {
  402. UploadId: UploadId,
  403. PartList: PartListData.PartList,
  404. });
  405. }
  406. }
  407. );
  408. };
  409. next(0);
  410. });
  411. // 获取线上 UploadId 列表
  412. ep.on('get_remote_upload_id_list', function () {
  413. // 获取符合条件的 UploadId 列表,因为同一个文件可以有多个上传任务。
  414. wholeMultipartList.call(
  415. self,
  416. {
  417. Bucket: Bucket,
  418. Region: Region,
  419. Key: Key,
  420. },
  421. function (err, data) {
  422. if (!self._isRunningTask(TaskId)) return;
  423. if (err) return ep.emit('error', err);
  424. // 整理远端 UploadId 列表
  425. var RemoteUploadIdList = util
  426. .filter(data.UploadList, function (item) {
  427. return (
  428. item.Key === Key && (!StorageClass || item.StorageClass.toUpperCase() === StorageClass.toUpperCase())
  429. );
  430. })
  431. .reverse()
  432. .map(function (item) {
  433. return item.UploadId || item.UploadID;
  434. });
  435. if (RemoteUploadIdList.length) {
  436. ep.emit('seek_local_avail_upload_id', RemoteUploadIdList);
  437. } else {
  438. // 远端没有 UploadId,清理缓存的 UploadId
  439. var uuid = session.getFileId(params.FileStat, params.ChunkSize, Bucket, Key),
  440. LocalUploadIdList;
  441. if (uuid && (LocalUploadIdList = session.getUploadIdList.call(self, uuid))) {
  442. util.each(LocalUploadIdList, function (UploadId) {
  443. session.removeUploadId.call(self, UploadId);
  444. });
  445. }
  446. ep.emit('no_available_upload_id');
  447. }
  448. }
  449. );
  450. });
  451. // 开始找可用 UploadId
  452. ep.emit('get_remote_upload_id_list');
  453. }
  454. // 获取符合条件的全部上传任务 (条件包括 Bucket, Region, Prefix)
  455. function wholeMultipartList(params, callback) {
  456. var self = this;
  457. var UploadList = [];
  458. var sendParams = {
  459. Bucket: params.Bucket,
  460. Region: params.Region,
  461. Prefix: params.Key,
  462. };
  463. var next = function () {
  464. self.multipartList(sendParams, function (err, data) {
  465. if (err) return callback(err);
  466. UploadList.push.apply(UploadList, data.Upload || []);
  467. if (data.IsTruncated === 'true') {
  468. // 列表不完整
  469. sendParams.KeyMarker = data.NextKeyMarker;
  470. sendParams.UploadIdMarker = data.NextUploadIdMarker;
  471. next();
  472. } else {
  473. callback(null, { UploadList: UploadList });
  474. }
  475. });
  476. };
  477. next();
  478. }
  479. // 获取指定上传任务的分块列表
  480. function wholeMultipartListPart(params, callback) {
  481. var self = this;
  482. var PartList = [];
  483. var sendParams = {
  484. Bucket: params.Bucket,
  485. Region: params.Region,
  486. Key: params.Key,
  487. UploadId: params.UploadId,
  488. };
  489. var next = function () {
  490. self.multipartListPart(sendParams, function (err, data) {
  491. if (err) return callback(err);
  492. PartList.push.apply(PartList, data.Part || []);
  493. if (data.IsTruncated === 'true') {
  494. // 列表不完整
  495. sendParams.PartNumberMarker = data.NextPartNumberMarker;
  496. next();
  497. } else {
  498. callback(null, { PartList: PartList });
  499. }
  500. });
  501. };
  502. next();
  503. }
  504. // 上传文件分块,包括
  505. /*
  506. UploadId (上传任务编号)
  507. AsyncLimit (并发量),
  508. SliceList (上传的分块数组),
  509. FilePath (本地文件的位置),
  510. SliceSize (文件分块大小)
  511. FileSize (文件大小)
  512. onProgress (上传成功之后的回调函数)
  513. */
  514. function uploadSliceList(params, cb) {
  515. var self = this;
  516. var TaskId = params.TaskId;
  517. var Bucket = params.Bucket;
  518. var Region = params.Region;
  519. var Key = params.Key;
  520. var UploadData = params.UploadData;
  521. var FileSize = params.FileSize;
  522. var SliceSize = params.SliceSize;
  523. var ChunkParallel = Math.min(params.AsyncLimit || self.options.ChunkParallelLimit || 1, 256);
  524. var FilePath = params.FilePath;
  525. var SliceCount = Math.ceil(FileSize / SliceSize);
  526. var FinishSize = 0;
  527. var ServerSideEncryption = params.ServerSideEncryption;
  528. var needUploadSlices = util.filter(UploadData.PartList, function (SliceItem) {
  529. if (SliceItem['Uploaded']) {
  530. FinishSize += SliceItem['PartNumber'] >= SliceCount ? FileSize % SliceSize || SliceSize : SliceSize;
  531. }
  532. return !SliceItem['Uploaded'];
  533. });
  534. var onProgress = params.onProgress;
  535. Async.eachLimit(
  536. needUploadSlices,
  537. ChunkParallel,
  538. function (SliceItem, asyncCallback) {
  539. if (!self._isRunningTask(TaskId)) return;
  540. var PartNumber = SliceItem['PartNumber'];
  541. var currentSize =
  542. Math.min(FileSize, SliceItem['PartNumber'] * SliceSize) - (SliceItem['PartNumber'] - 1) * SliceSize;
  543. var preAddSize = 0;
  544. uploadSliceItem.call(
  545. self,
  546. {
  547. TaskId: TaskId,
  548. Bucket: Bucket,
  549. Region: Region,
  550. Key: Key,
  551. SliceSize: SliceSize,
  552. FileSize: FileSize,
  553. PartNumber: PartNumber,
  554. ServerSideEncryption: ServerSideEncryption,
  555. FilePath: FilePath,
  556. UploadData: UploadData,
  557. Headers: params.Headers,
  558. onProgress: function (data) {
  559. FinishSize += data.loaded - preAddSize;
  560. preAddSize = data.loaded;
  561. onProgress({ loaded: FinishSize, total: FileSize });
  562. },
  563. },
  564. function (err, data) {
  565. if (!self._isRunningTask(TaskId)) return;
  566. if (err) {
  567. FinishSize -= preAddSize;
  568. } else {
  569. FinishSize += currentSize - preAddSize;
  570. SliceItem.ETag = data.ETag;
  571. }
  572. onProgress({ loaded: FinishSize, total: FileSize });
  573. asyncCallback(err || null, data);
  574. }
  575. );
  576. },
  577. function (err) {
  578. if (!self._isRunningTask(TaskId)) return;
  579. if (err) return cb(err);
  580. cb(null, {
  581. UploadId: UploadData.UploadId,
  582. SliceList: UploadData.PartList,
  583. });
  584. }
  585. );
  586. }
  587. // 上传指定分片
  588. function uploadSliceItem(params, callback) {
  589. var self = this;
  590. var TaskId = params.TaskId;
  591. var Bucket = params.Bucket;
  592. var Region = params.Region;
  593. var Key = params.Key;
  594. var FileSize = params.FileSize;
  595. var FilePath = params.FilePath;
  596. var PartNumber = params.PartNumber * 1;
  597. var SliceSize = params.SliceSize;
  598. var ServerSideEncryption = params.ServerSideEncryption;
  599. var UploadData = params.UploadData;
  600. var ChunkRetryTimes = self.options.ChunkRetryTimes + 1;
  601. var Headers = params.Headers || {};
  602. var start = SliceSize * (PartNumber - 1);
  603. var ContentLength = SliceSize;
  604. var end = start + SliceSize;
  605. if (end > FileSize) {
  606. end = FileSize;
  607. ContentLength = end - start;
  608. }
  609. var headersWhiteList = ['x-cos-traffic-limit', 'x-cos-mime-limit'];
  610. var headers = {};
  611. util.each(Headers, function (v, k) {
  612. if (headersWhiteList.indexOf(k) > -1) {
  613. headers[k] = v;
  614. }
  615. });
  616. util.fileSlice(FilePath, start, end, function (md5Body) {
  617. util.getFileMd5(md5Body, function (err, md5) {
  618. var contentMd5 = md5 ? util.binaryBase64(md5) : '';
  619. var PartItem = UploadData.PartList[PartNumber - 1];
  620. var switchHost = false;
  621. Async.retry(
  622. ChunkRetryTimes,
  623. function (tryCallback) {
  624. if (!self._isRunningTask(TaskId)) return;
  625. util.fileSlice(FilePath, start, end, function (Body) {
  626. self.multipartUpload(
  627. {
  628. TaskId: TaskId,
  629. Bucket: Bucket,
  630. Region: Region,
  631. Key: Key,
  632. ContentLength: ContentLength,
  633. PartNumber: PartNumber,
  634. UploadId: UploadData.UploadId,
  635. ServerSideEncryption: ServerSideEncryption,
  636. Body: Body,
  637. Headers: headers,
  638. onProgress: params.onProgress,
  639. ContentMD5: contentMd5,
  640. SwitchHost: switchHost,
  641. },
  642. function (err, data) {
  643. if (!self._isRunningTask(TaskId)) return;
  644. if (err) {
  645. switchHost = err.switchHost || false;
  646. }
  647. if (err) return tryCallback(err);
  648. PartItem.Uploaded = true;
  649. return tryCallback(null, data);
  650. }
  651. );
  652. });
  653. },
  654. function (err, data) {
  655. if (err) {
  656. delete err.switchHost;
  657. }
  658. if (!self._isRunningTask(TaskId)) return;
  659. return callback(err, data);
  660. }
  661. );
  662. });
  663. });
  664. }
  665. // 完成分块上传
  666. function uploadSliceComplete(params, callback) {
  667. var Bucket = params.Bucket;
  668. var Region = params.Region;
  669. var Key = params.Key;
  670. var UploadId = params.UploadId;
  671. var SliceList = params.SliceList;
  672. var self = this;
  673. var ChunkRetryTimes = this.options.ChunkRetryTimes + 1;
  674. var Headers = params.Headers;
  675. var Parts = SliceList.map(function (item) {
  676. return {
  677. PartNumber: item.PartNumber,
  678. ETag: item.ETag,
  679. };
  680. });
  681. // 完成上传的请求也做重试
  682. Async.retry(
  683. ChunkRetryTimes,
  684. function (tryCallback) {
  685. self.multipartComplete(
  686. {
  687. Bucket: Bucket,
  688. Region: Region,
  689. Key: Key,
  690. UploadId: UploadId,
  691. Parts: Parts,
  692. Headers: Headers,
  693. },
  694. tryCallback
  695. );
  696. },
  697. function (err, data) {
  698. callback(err, data);
  699. }
  700. );
  701. }
  702. // 抛弃分块上传任务
  703. /*
  704. AsyncLimit (抛弃上传任务的并发量),
  705. UploadId (上传任务的编号,当 Level 为 task 时候需要)
  706. Level (抛弃分块上传任务的级别,task : 抛弃指定的上传任务,file : 抛弃指定的文件对应的上传任务,其他值 :抛弃指定Bucket 的全部上传任务)
  707. */
  708. function abortUploadTask(params, callback) {
  709. var Bucket = params.Bucket;
  710. var Region = params.Region;
  711. var Key = params.Key;
  712. var UploadId = params.UploadId;
  713. var Level = params.Level || 'task';
  714. var AsyncLimit = params.AsyncLimit;
  715. var self = this;
  716. var ep = new EventProxy();
  717. ep.on('error', function (errData) {
  718. return callback(errData);
  719. });
  720. // 已经获取到需要抛弃的任务列表
  721. ep.on('get_abort_array', function (AbortArray) {
  722. abortUploadTaskArray.call(
  723. self,
  724. {
  725. Bucket: Bucket,
  726. Region: Region,
  727. Key: Key,
  728. Headers: params.Headers,
  729. AsyncLimit: AsyncLimit,
  730. AbortArray: AbortArray,
  731. },
  732. callback
  733. );
  734. });
  735. if (Level === 'bucket') {
  736. // Bucket 级别的任务抛弃,抛弃该 Bucket 下的全部上传任务
  737. wholeMultipartList.call(
  738. self,
  739. {
  740. Bucket: Bucket,
  741. Region: Region,
  742. },
  743. function (err, data) {
  744. if (err) return callback(err);
  745. ep.emit('get_abort_array', data.UploadList || []);
  746. }
  747. );
  748. } else if (Level === 'file') {
  749. // 文件级别的任务抛弃,抛弃该文件的全部上传任务
  750. if (!Key) return callback(util.error(new Error('abort_upload_task_no_key')));
  751. wholeMultipartList.call(
  752. self,
  753. {
  754. Bucket: Bucket,
  755. Region: Region,
  756. Key: Key,
  757. },
  758. function (err, data) {
  759. if (err) return callback(err);
  760. ep.emit('get_abort_array', data.UploadList || []);
  761. }
  762. );
  763. } else if (Level === 'task') {
  764. // 单个任务级别的任务抛弃,抛弃指定 UploadId 的上传任务
  765. if (!UploadId) return callback(util.error(new Error('abort_upload_task_no_id')));
  766. if (!Key) return callback(util.error(new Error('abort_upload_task_no_key')));
  767. ep.emit('get_abort_array', [
  768. {
  769. Key: Key,
  770. UploadId: UploadId,
  771. },
  772. ]);
  773. } else {
  774. return callback(util.error(new Error('abort_unknown_level')));
  775. }
  776. }
  777. // 批量抛弃分块上传任务
  778. function abortUploadTaskArray(params, callback) {
  779. var Bucket = params.Bucket;
  780. var Region = params.Region;
  781. var Key = params.Key;
  782. var AbortArray = params.AbortArray;
  783. var AsyncLimit = params.AsyncLimit || 1;
  784. var self = this;
  785. var index = 0;
  786. var resultList = new Array(AbortArray.length);
  787. Async.eachLimit(
  788. AbortArray,
  789. AsyncLimit,
  790. function (AbortItem, nextItem) {
  791. var eachIndex = index;
  792. if (Key && Key !== AbortItem.Key) {
  793. resultList[eachIndex] = { error: { KeyNotMatch: true } };
  794. nextItem(null);
  795. return;
  796. }
  797. var UploadId = AbortItem.UploadId || AbortItem.UploadID;
  798. self.multipartAbort(
  799. {
  800. Bucket: Bucket,
  801. Region: Region,
  802. Key: AbortItem.Key,
  803. Headers: params.Headers,
  804. UploadId: UploadId,
  805. },
  806. function (err) {
  807. var task = {
  808. Bucket: Bucket,
  809. Region: Region,
  810. Key: AbortItem.Key,
  811. UploadId: UploadId,
  812. };
  813. resultList[eachIndex] = { error: err, task: task };
  814. nextItem(null);
  815. }
  816. );
  817. index++;
  818. },
  819. function (err) {
  820. if (err) return callback(err);
  821. var successList = [];
  822. var errorList = [];
  823. for (var i = 0, len = resultList.length; i < len; i++) {
  824. var item = resultList[i];
  825. if (item['task']) {
  826. if (item['error']) {
  827. errorList.push(item['task']);
  828. } else {
  829. successList.push(item['task']);
  830. }
  831. }
  832. }
  833. return callback(null, {
  834. successList: successList,
  835. errorList: errorList,
  836. });
  837. }
  838. );
  839. }
  840. // 高级上传
  841. function uploadFile(params, callback) {
  842. var self = this;
  843. // 判断多大的文件使用分片上传
  844. var SliceSize = params.SliceSize === undefined ? self.options.SliceSize : params.SliceSize;
  845. // 开始处理每个文件
  846. var taskList = [];
  847. fs.stat(params.FilePath, function (err, stat) {
  848. if (err) {
  849. return callback(err);
  850. }
  851. var isDir = stat.isDirectory();
  852. var FileSize = (params.ContentLength = stat.size || 0);
  853. var fileInfo = { TaskId: '' };
  854. // 整理 option,用于返回给回调
  855. util.each(params, function (v, k) {
  856. if (typeof v !== 'object' && typeof v !== 'function') {
  857. fileInfo[k] = v;
  858. }
  859. });
  860. // 处理文件 TaskReady
  861. var _onTaskReady = params.onTaskReady;
  862. var onTaskReady = function (tid) {
  863. fileInfo.TaskId = tid;
  864. _onTaskReady && _onTaskReady(tid);
  865. };
  866. params.onTaskReady = onTaskReady;
  867. // 处理文件完成
  868. var _onFileFinish = params.onFileFinish;
  869. var onFileFinish = function (err, data) {
  870. _onFileFinish && _onFileFinish(err, data, fileInfo);
  871. callback && callback(err, data);
  872. };
  873. // 添加上传任务
  874. var api = FileSize <= SliceSize || isDir ? 'putObject' : 'sliceUploadFile';
  875. if (api === 'putObject') {
  876. params.Body = isDir ? '' : fs.createReadStream(params.FilePath);
  877. params.Body.isSdkCreated = true;
  878. }
  879. taskList.push({
  880. api: api,
  881. params: params,
  882. callback: onFileFinish,
  883. });
  884. self._addTasks(taskList);
  885. });
  886. }
  887. // 批量上传文件
  888. function uploadFiles(params, callback) {
  889. var self = this;
  890. // 判断多大的文件使用分片上传
  891. var SliceSize = params.SliceSize === undefined ? self.options.SliceSize : params.SliceSize;
  892. // 汇总返回进度
  893. var TotalSize = 0;
  894. var TotalFinish = 0;
  895. var onTotalProgress = util.throttleOnProgress.call(self, TotalFinish, params.onProgress);
  896. // 汇总返回回调
  897. var unFinishCount = params.files.length;
  898. var _onTotalFileFinish = params.onFileFinish;
  899. var resultList = Array(unFinishCount);
  900. var onTotalFileFinish = function (err, data, options) {
  901. onTotalProgress(null, true);
  902. _onTotalFileFinish && _onTotalFileFinish(err, data, options);
  903. resultList[options.Index] = {
  904. options: options,
  905. error: err,
  906. data: data,
  907. };
  908. if (--unFinishCount <= 0 && callback) {
  909. callback(null, { files: resultList });
  910. }
  911. };
  912. // 开始处理每个文件
  913. var taskList = [];
  914. var count = params.files.length;
  915. util.each(params.files, function (fileParams, index) {
  916. var isDir = false;
  917. var FileSize = 0;
  918. if (fileParams.Body) {
  919. util.getFileSize('putObject', fileParams, function (err, size) {
  920. FileSize = fileParams.ContentLengt = size;
  921. });
  922. } else if (fileParams.FilePath) {
  923. var stat;
  924. try {
  925. stat = fs.statSync(fileParams.FilePath);
  926. } catch (e) {}
  927. isDir = stat ? stat.isDirectory() : false;
  928. FileSize = fileParams.ContentLength = stat ? stat.size : 0;
  929. }
  930. var fileInfo = { Index: index, TaskId: '' };
  931. // 更新文件总大小
  932. TotalSize += FileSize;
  933. // 整理 option,用于返回给回调
  934. util.each(fileParams, function (v, k) {
  935. if (typeof v !== 'object' && typeof v !== 'function') {
  936. fileInfo[k] = v;
  937. }
  938. });
  939. // 处理单个文件 TaskReady
  940. var _onTaskReady = fileParams.onTaskReady;
  941. var onTaskReady = function (tid) {
  942. fileInfo.TaskId = tid;
  943. _onTaskReady && _onTaskReady(tid);
  944. };
  945. fileParams.onTaskReady = onTaskReady;
  946. // 处理单个文件进度
  947. var PreAddSize = 0;
  948. var _onProgress = fileParams.onProgress;
  949. var onProgress = function (info) {
  950. TotalFinish = TotalFinish - PreAddSize + info.loaded;
  951. PreAddSize = info.loaded;
  952. _onProgress && _onProgress(info);
  953. onTotalProgress({ loaded: TotalFinish, total: TotalSize });
  954. };
  955. fileParams.onProgress = onProgress;
  956. // 处理单个文件完成
  957. var _onFileFinish = fileParams.onFileFinish;
  958. var onFileFinish = function (err, data) {
  959. _onFileFinish && _onFileFinish(err, data);
  960. onTotalFileFinish && onTotalFileFinish(err, data, fileInfo);
  961. };
  962. // 添加上传任务,传入 Body 则只支持简单上传
  963. var api = FileSize <= SliceSize || isDir || fileParams.Body ? 'putObject' : 'sliceUploadFile';
  964. if (api === 'putObject' && fileParams.FilePath && !fileParams.Body) {
  965. fileParams.Body = isDir ? '' : fs.createReadStream(fileParams.FilePath);
  966. fileParams.Body.isSdkCreated = true;
  967. }
  968. taskList.push({
  969. api: api,
  970. params: fileParams,
  971. callback: onFileFinish,
  972. });
  973. --count === 0 && self._addTasks(taskList);
  974. });
  975. }
  976. // 分片复制文件
  977. function sliceCopyFile(params, callback) {
  978. var ep = new EventProxy();
  979. var self = this;
  980. var Bucket = params.Bucket;
  981. var Region = params.Region;
  982. var Key = params.Key;
  983. var CopySource = params.CopySource;
  984. var m = util.getSourceParams.call(this, CopySource);
  985. if (!m) {
  986. callback(util.error(new Error('CopySource format error')));
  987. return;
  988. }
  989. var SourceBucket = m.Bucket;
  990. var SourceRegion = m.Region;
  991. var SourceKey = decodeURIComponent(m.Key);
  992. var CopySliceSize = params.CopySliceSize === undefined ? self.options.CopySliceSize : params.CopySliceSize;
  993. CopySliceSize = Math.max(0, CopySliceSize);
  994. var ChunkSize = params.CopyChunkSize || this.options.CopyChunkSize;
  995. var ChunkParallel = this.options.CopyChunkParallelLimit;
  996. var ChunkRetryTimes = this.options.ChunkRetryTimes + 1;
  997. var ChunkCount = 0;
  998. var FinishSize = 0;
  999. var FileSize;
  1000. var onProgress;
  1001. var SourceResHeaders = {};
  1002. var SourceHeaders = {};
  1003. var TargetHeader = {};
  1004. // 分片复制完成,开始 multipartComplete 操作
  1005. ep.on('copy_slice_complete', function (UploadData) {
  1006. var metaHeaders = {};
  1007. util.each(params.Headers, function (val, k) {
  1008. if (k.toLowerCase().indexOf('x-cos-meta-') === 0) metaHeaders[k] = val;
  1009. });
  1010. var Parts = util.map(UploadData.PartList, function (item) {
  1011. return {
  1012. PartNumber: item.PartNumber,
  1013. ETag: item.ETag,
  1014. };
  1015. });
  1016. // 完成上传的请求也做重试
  1017. Async.retry(
  1018. ChunkRetryTimes,
  1019. function (tryCallback) {
  1020. self.multipartComplete(
  1021. {
  1022. Bucket: Bucket,
  1023. Region: Region,
  1024. Key: Key,
  1025. UploadId: UploadData.UploadId,
  1026. Parts: Parts,
  1027. },
  1028. tryCallback
  1029. );
  1030. },
  1031. function (err, data) {
  1032. session.removeUsing(UploadData.UploadId); // 标记 UploadId 没被使用了,因为复制没提供重试,所以只要出错,就是 UploadId 停用了。
  1033. if (err) {
  1034. onProgress(null, true);
  1035. return callback(err);
  1036. }
  1037. session.removeUploadId.call(self, UploadData.UploadId);
  1038. onProgress({ loaded: FileSize, total: FileSize }, true);
  1039. callback(null, data);
  1040. }
  1041. );
  1042. });
  1043. ep.on('get_copy_data_finish', function (UploadData) {
  1044. // 处理 UploadId 缓存
  1045. var uuid = session.getCopyFileId(CopySource, SourceResHeaders, ChunkSize, Bucket, Key);
  1046. uuid && session.saveUploadId.call(self, uuid, UploadData.UploadId, self.options.UploadIdCacheLimit); // 缓存 UploadId
  1047. session.setUsing(UploadData.UploadId); // 标记 UploadId 为正在使用
  1048. var needCopySlices = util.filter(UploadData.PartList, function (SliceItem) {
  1049. if (SliceItem['Uploaded']) {
  1050. FinishSize += SliceItem['PartNumber'] >= ChunkCount ? FileSize % ChunkSize || ChunkSize : ChunkSize;
  1051. }
  1052. return !SliceItem['Uploaded'];
  1053. });
  1054. Async.eachLimit(
  1055. needCopySlices,
  1056. ChunkParallel,
  1057. function (SliceItem, asyncCallback) {
  1058. var PartNumber = SliceItem.PartNumber;
  1059. var CopySourceRange = SliceItem.CopySourceRange;
  1060. var currentSize = SliceItem.end - SliceItem.start;
  1061. Async.retry(
  1062. ChunkRetryTimes,
  1063. function (tryCallback) {
  1064. copySliceItem.call(
  1065. self,
  1066. {
  1067. Bucket: Bucket,
  1068. Region: Region,
  1069. Key: Key,
  1070. CopySource: CopySource,
  1071. UploadId: UploadData.UploadId,
  1072. PartNumber: PartNumber,
  1073. CopySourceRange: CopySourceRange,
  1074. },
  1075. tryCallback
  1076. );
  1077. },
  1078. function (err, data) {
  1079. if (err) return asyncCallback(err);
  1080. FinishSize += currentSize;
  1081. onProgress({ loaded: FinishSize, total: FileSize });
  1082. SliceItem.ETag = data.ETag;
  1083. asyncCallback(err || null, data);
  1084. }
  1085. );
  1086. },
  1087. function (err) {
  1088. if (err) {
  1089. session.removeUsing(UploadData.UploadId); // 标记 UploadId 没被使用了,因为复制没提供重试,所以只要出错,就是 UploadId 停用了。
  1090. onProgress(null, true);
  1091. return callback(err);
  1092. }
  1093. ep.emit('copy_slice_complete', UploadData);
  1094. }
  1095. );
  1096. });
  1097. ep.on('get_chunk_size_finish', function () {
  1098. var createNewUploadId = function () {
  1099. self.multipartInit(
  1100. {
  1101. Bucket: Bucket,
  1102. Region: Region,
  1103. Key: Key,
  1104. Headers: TargetHeader,
  1105. },
  1106. function (err, data) {
  1107. if (err) return callback(err);
  1108. params.UploadId = data.UploadId;
  1109. ep.emit('get_copy_data_finish', { UploadId: params.UploadId, PartList: params.PartList });
  1110. }
  1111. );
  1112. };
  1113. // 在本地找可用的 UploadId
  1114. var uuid = session.getCopyFileId(CopySource, SourceResHeaders, ChunkSize, Bucket, Key);
  1115. var LocalUploadIdList = session.getUploadIdList.call(self, uuid);
  1116. if (!uuid || !LocalUploadIdList) return createNewUploadId();
  1117. var next = function (index) {
  1118. // 如果本地找不到可用 UploadId,再一个个遍历校验远端
  1119. if (index >= LocalUploadIdList.length) return createNewUploadId();
  1120. var UploadId = LocalUploadIdList[index];
  1121. // 如果正在被使用,跳过
  1122. if (session.using[UploadId]) return next(index + 1);
  1123. // 判断 UploadId 是否存在线上
  1124. wholeMultipartListPart.call(
  1125. self,
  1126. {
  1127. Bucket: Bucket,
  1128. Region: Region,
  1129. Key: Key,
  1130. UploadId: UploadId,
  1131. },
  1132. function (err, PartListData) {
  1133. if (err) {
  1134. // 如果 UploadId 获取会出错,跳过并删除
  1135. session.removeUploadId.call(self, UploadId);
  1136. next(index + 1);
  1137. } else {
  1138. // 如果异步回来 UploadId 已经被用了,也跳过
  1139. if (session.using[UploadId]) return next(index + 1);
  1140. // 找到可用 UploadId
  1141. var finishETagMap = {};
  1142. var offset = 0;
  1143. util.each(PartListData.PartList, function (PartItem) {
  1144. var size = parseInt(PartItem.Size);
  1145. var end = offset + size - 1;
  1146. finishETagMap[PartItem.PartNumber + '|' + offset + '|' + end] = PartItem.ETag;
  1147. offset += size;
  1148. });
  1149. util.each(params.PartList, function (PartItem) {
  1150. var ETag = finishETagMap[PartItem.PartNumber + '|' + PartItem.start + '|' + PartItem.end];
  1151. if (ETag) {
  1152. PartItem.ETag = ETag;
  1153. PartItem.Uploaded = true;
  1154. }
  1155. });
  1156. ep.emit('get_copy_data_finish', { UploadId: UploadId, PartList: params.PartList });
  1157. }
  1158. }
  1159. );
  1160. };
  1161. next(0);
  1162. });
  1163. ep.on('get_file_size_finish', function () {
  1164. // 控制分片大小
  1165. (function () {
  1166. var SIZE = [1, 2, 4, 8, 16, 32, 64, 128, 256, 512, 1024, 1024 * 2, 1024 * 4, 1024 * 5];
  1167. var AutoChunkSize = 1024 * 1024;
  1168. for (var i = 0; i < SIZE.length; i++) {
  1169. AutoChunkSize = SIZE[i] * 1024 * 1024;
  1170. if (FileSize / AutoChunkSize <= self.options.MaxPartNumber) break;
  1171. }
  1172. params.ChunkSize = ChunkSize = Math.max(ChunkSize, AutoChunkSize);
  1173. ChunkCount = Math.ceil(FileSize / ChunkSize);
  1174. var list = [];
  1175. for (var partNumber = 1; partNumber <= ChunkCount; partNumber++) {
  1176. var start = (partNumber - 1) * ChunkSize;
  1177. var end = partNumber * ChunkSize < FileSize ? partNumber * ChunkSize - 1 : FileSize - 1;
  1178. var item = {
  1179. PartNumber: partNumber,
  1180. start: start,
  1181. end: end,
  1182. CopySourceRange: 'bytes=' + start + '-' + end,
  1183. };
  1184. list.push(item);
  1185. }
  1186. params.PartList = list;
  1187. })();
  1188. if (params.Headers['x-cos-metadata-directive'] === 'Replaced') {
  1189. TargetHeader = params.Headers;
  1190. } else {
  1191. TargetHeader = SourceHeaders;
  1192. }
  1193. TargetHeader['x-cos-storage-class'] = params.Headers['x-cos-storage-class'] || SourceHeaders['x-cos-storage-class'];
  1194. TargetHeader = util.clearKey(TargetHeader);
  1195. /**
  1196. * 对于归档存储的对象,如果未恢复副本,则不允许 Copy
  1197. */
  1198. if (SourceHeaders['x-cos-storage-class'] === 'ARCHIVE' || SourceHeaders['x-cos-storage-class'] === 'DEEP_ARCHIVE') {
  1199. var restoreHeader = SourceHeaders['x-cos-restore'];
  1200. if (!restoreHeader || restoreHeader === 'ongoing-request="true"') {
  1201. callback(util.error(new Error('Unrestored archive object is not allowed to be copied')));
  1202. return;
  1203. }
  1204. }
  1205. /**
  1206. * 去除一些无用的头部,规避 multipartInit 出错
  1207. * 这些头部通常是在 putObjectCopy 时才使用
  1208. */
  1209. delete TargetHeader['x-cos-copy-source'];
  1210. delete TargetHeader['x-cos-metadata-directive'];
  1211. delete TargetHeader['x-cos-copy-source-If-Modified-Since'];
  1212. delete TargetHeader['x-cos-copy-source-If-Unmodified-Since'];
  1213. delete TargetHeader['x-cos-copy-source-If-Match'];
  1214. delete TargetHeader['x-cos-copy-source-If-None-Match'];
  1215. ep.emit('get_chunk_size_finish');
  1216. });
  1217. // 获取远端复制源文件的大小
  1218. self.headObject(
  1219. {
  1220. Bucket: SourceBucket,
  1221. Region: SourceRegion,
  1222. Key: SourceKey,
  1223. },
  1224. function (err, data) {
  1225. if (err) {
  1226. if (err.statusCode && err.statusCode === 404) {
  1227. callback(util.error(err, { ErrorStatus: SourceKey + ' Not Exist' }));
  1228. } else {
  1229. callback(err);
  1230. }
  1231. return;
  1232. }
  1233. FileSize = params.FileSize = data.headers['content-length'];
  1234. if (FileSize === undefined || !FileSize) {
  1235. callback(
  1236. util.error(
  1237. new Error(
  1238. 'get Content-Length error, please add "Content-Length" to CORS ExposeHeader setting.( 获取Content-Length失败,请在CORS ExposeHeader设置中添加Content-Length,请参考文档:https://cloud.tencent.com/document/product/436/13318 )'
  1239. )
  1240. )
  1241. );
  1242. return;
  1243. }
  1244. onProgress = util.throttleOnProgress.call(self, FileSize, params.onProgress);
  1245. // 开始上传
  1246. if (FileSize <= CopySliceSize) {
  1247. if (!params.Headers['x-cos-metadata-directive']) {
  1248. params.Headers['x-cos-metadata-directive'] = 'Copy';
  1249. }
  1250. self.putObjectCopy(params, function (err, data) {
  1251. if (err) {
  1252. onProgress(null, true);
  1253. return callback(err);
  1254. }
  1255. onProgress({ loaded: FileSize, total: FileSize }, true);
  1256. callback(err, data);
  1257. });
  1258. } else {
  1259. var resHeaders = data.headers;
  1260. SourceResHeaders = resHeaders;
  1261. SourceHeaders = {
  1262. 'Cache-Control': resHeaders['cache-control'],
  1263. 'Content-Disposition': resHeaders['content-disposition'],
  1264. 'Content-Encoding': resHeaders['content-encoding'],
  1265. 'Content-Type': resHeaders['content-type'],
  1266. Expires: resHeaders['expires'],
  1267. 'x-cos-storage-class': resHeaders['x-cos-storage-class'],
  1268. };
  1269. util.each(resHeaders, function (v, k) {
  1270. var metaPrefix = 'x-cos-meta-';
  1271. if (k.indexOf(metaPrefix) === 0 && k.length > metaPrefix.length) {
  1272. SourceHeaders[k] = v;
  1273. }
  1274. });
  1275. ep.emit('get_file_size_finish');
  1276. }
  1277. }
  1278. );
  1279. }
  1280. // 复制指定分片
  1281. function copySliceItem(params, callback) {
  1282. var TaskId = params.TaskId;
  1283. var Bucket = params.Bucket;
  1284. var Region = params.Region;
  1285. var Key = params.Key;
  1286. var CopySource = params.CopySource;
  1287. var UploadId = params.UploadId;
  1288. var PartNumber = params.PartNumber * 1;
  1289. var CopySourceRange = params.CopySourceRange;
  1290. var ChunkRetryTimes = this.options.ChunkRetryTimes + 1;
  1291. var self = this;
  1292. Async.retry(
  1293. ChunkRetryTimes,
  1294. function (tryCallback) {
  1295. self.uploadPartCopy(
  1296. {
  1297. TaskId: TaskId,
  1298. Bucket: Bucket,
  1299. Region: Region,
  1300. Key: Key,
  1301. CopySource: CopySource,
  1302. UploadId: UploadId,
  1303. PartNumber: PartNumber,
  1304. CopySourceRange: CopySourceRange,
  1305. },
  1306. function (err, data) {
  1307. tryCallback(err || null, data);
  1308. }
  1309. );
  1310. },
  1311. function (err, data) {
  1312. return callback(err, data);
  1313. }
  1314. );
  1315. }
  1316. // 分片下载文件
  1317. function downloadFile(params, callback) {
  1318. if (this.options.ObjectKeySimplifyCheck) {
  1319. // getObject 的 Key 需要格式化,避免调用成 getBucket
  1320. var formatKey = util.simplifyPath(params.Key);
  1321. if (formatKey === '/') {
  1322. callback(util.error(new Error('The Getobject Key is illegal')));
  1323. return;
  1324. }
  1325. }
  1326. var self = this;
  1327. var TaskId = params.TaskId || util.uuid();
  1328. var Bucket = params.Bucket;
  1329. var Region = params.Region;
  1330. var Key = params.Key;
  1331. var FilePath = params.FilePath;
  1332. var FileSize;
  1333. var FinishSize = 0;
  1334. var onProgress;
  1335. var ChunkSize = params.ChunkSize || 1024 * 1024;
  1336. var ParallelLimit = params.ParallelLimit || 5;
  1337. var RetryTimes = params.RetryTimes || 3;
  1338. var ep = new EventProxy();
  1339. var PartList;
  1340. var aborted = false;
  1341. var head = {};
  1342. ep.on('error', function (err) {
  1343. callback(err);
  1344. });
  1345. ep.on('get_file_info', function () {
  1346. // 获取远端复制源文件的大小
  1347. self.headObject(
  1348. {
  1349. Bucket: Bucket,
  1350. Region: Region,
  1351. Key: Key,
  1352. },
  1353. function (err, data) {
  1354. if (err) return ep.emit('error', err);
  1355. // 获取文件大小
  1356. FileSize = params.FileSize = data.headers['content-length']
  1357. ? parseInt(data.headers['content-length'])
  1358. : undefined;
  1359. if (FileSize === undefined) {
  1360. callback(
  1361. util.error(
  1362. new Error(
  1363. 'get Content-Length error, please add "Content-Length" to CORS ExposeHeader setting.( 获取Content-Length失败,请在CORS ExposeHeader设置中添加Content-Length,请参考文档:https://cloud.tencent.com/document/product/436/13318 )'
  1364. )
  1365. )
  1366. );
  1367. return;
  1368. }
  1369. // 归档文件不支持下载
  1370. const resHeaders = data.headers;
  1371. const storageClass = resHeaders['x-cos-storage-class'] || '';
  1372. const restoreStatus = resHeaders['x-cos-restore'] || '';
  1373. if (
  1374. ['DEEP_ARCHIVE', 'ARCHIVE'].includes(storageClass) &&
  1375. (!restoreStatus || restoreStatus === 'ongoing-request="true"')
  1376. ) {
  1377. // 自定义返回的错误码 与cos api无关
  1378. return callback({
  1379. statusCode: 403,
  1380. header: resHeaders,
  1381. code: 'CannotDownload',
  1382. message: 'Archive object can not download, please restore to Standard storage class.',
  1383. });
  1384. }
  1385. // 整理文件信息
  1386. head = {
  1387. ETag: data.ETag,
  1388. size: FileSize,
  1389. mtime: resHeaders['last-modified'],
  1390. crc64ecma: resHeaders['x-cos-hash-crc64ecma'],
  1391. };
  1392. // 处理进度反馈
  1393. onProgress = util.throttleOnProgress.call(self, FileSize, function (info) {
  1394. if (aborted) return;
  1395. params.onProgress(info);
  1396. });
  1397. if (FileSize <= ChunkSize) {
  1398. // 小文件直接单请求下载
  1399. self.getObject(
  1400. {
  1401. TaskId: TaskId,
  1402. Bucket: Bucket,
  1403. Region: Region,
  1404. Key: Key,
  1405. onProgress: onProgress,
  1406. Output: fs.createWriteStream(FilePath),
  1407. },
  1408. function (err, data) {
  1409. if (err) {
  1410. onProgress(null, true);
  1411. return callback(err);
  1412. }
  1413. onProgress({ loaded: FileSize, total: FileSize }, true);
  1414. callback(err, data);
  1415. }
  1416. );
  1417. } else {
  1418. // 大文件分片下载
  1419. ep.emit('calc_suitable_chunk_size');
  1420. }
  1421. }
  1422. );
  1423. });
  1424. // 计算合适的分片大小
  1425. ep.on('calc_suitable_chunk_size', function () {
  1426. // 控制分片大小
  1427. var SIZE = [1, 2, 4, 8, 16, 32, 64, 128, 256, 512, 1024, 1024 * 2, 1024 * 4, 1024 * 5];
  1428. var AutoChunkSize = 1024 * 1024;
  1429. for (var i = 0; i < SIZE.length; i++) {
  1430. AutoChunkSize = SIZE[i] * 1024 * 1024;
  1431. if (FileSize / AutoChunkSize <= self.options.MaxPartNumber) break;
  1432. }
  1433. params.ChunkSize = ChunkSize = Math.max(ChunkSize, AutoChunkSize);
  1434. var ChunkCount = Math.ceil(FileSize / ChunkSize);
  1435. var list = [];
  1436. for (var partNumber = 1; partNumber <= ChunkCount; partNumber++) {
  1437. var start = (partNumber - 1) * ChunkSize;
  1438. var end = partNumber * ChunkSize < FileSize ? partNumber * ChunkSize - 1 : FileSize - 1;
  1439. var item = {
  1440. PartNumber: partNumber,
  1441. start: start,
  1442. end: end,
  1443. };
  1444. list.push(item);
  1445. }
  1446. PartList = list;
  1447. ep.emit('prepare_file');
  1448. });
  1449. // 准备要下载的空文件
  1450. ep.on('prepare_file', function () {
  1451. fs.writeFile(FilePath, '', (err) => {
  1452. if (err) {
  1453. ep.emit('error', err.code === 'EISDIR' ? { code: 'exist_same_dir', message: FilePath } : err);
  1454. } else {
  1455. ep.emit('start_download_chunks');
  1456. }
  1457. });
  1458. });
  1459. // 计算合适的分片大小
  1460. var result;
  1461. ep.on('start_download_chunks', function () {
  1462. onProgress({ loaded: 0, total: FileSize }, true);
  1463. var maxPartNumber = PartList.length;
  1464. Async.eachLimit(
  1465. PartList,
  1466. ParallelLimit,
  1467. function (part, nextChunk) {
  1468. if (aborted) return;
  1469. Async.retry(
  1470. RetryTimes,
  1471. function (tryCallback) {
  1472. if (aborted) return;
  1473. // FinishSize
  1474. var Headers = util.clone(params.Headers);
  1475. Headers.Range = 'bytes=' + part.start + '-' + part.end;
  1476. const writeStream = fs.createWriteStream(FilePath, {
  1477. start: part.start,
  1478. flags: 'r+',
  1479. });
  1480. var preAddSize = 0;
  1481. var chunkReadSize = part.end - part.start;
  1482. self.getObject(
  1483. {
  1484. TaskId: TaskId,
  1485. Bucket: params.Bucket,
  1486. Region: params.Region,
  1487. Key: params.Key,
  1488. Query: params.Query,
  1489. Headers: Headers,
  1490. onProgress: function (data) {
  1491. if (aborted) return;
  1492. FinishSize += data.loaded - preAddSize;
  1493. preAddSize = data.loaded;
  1494. onProgress({ loaded: FinishSize, total: FileSize });
  1495. },
  1496. Output: writeStream,
  1497. },
  1498. function (err, data) {
  1499. if (aborted) return;
  1500. // 处理错误和进度
  1501. if (err) {
  1502. FinishSize -= preAddSize;
  1503. return tryCallback(err);
  1504. }
  1505. // 处理返回值
  1506. if (part.PartNumber === maxPartNumber) result = data;
  1507. var chunkHeaders = data.headers || {};
  1508. var contentRanges = chunkHeaders['content-range'] || ''; // content-range 格式:"bytes 3145728-4194303/68577051"
  1509. var totalSize = parseInt(contentRanges.split('/')[1] || 0);
  1510. // 只校验文件大小和 crc64 是否有变更
  1511. var changed;
  1512. if (chunkHeaders['x-cos-hash-crc64ecma'] !== head.crc64ecma)
  1513. changed = 'download error, x-cos-hash-crc64ecma has changed.';
  1514. else if (totalSize !== head.size) changed = 'download error, Last-Modified has changed.';
  1515. // else if (data.ETag !== head.ETag) error = 'download error, ETag has changed.';
  1516. // else if (chunkHeaders['last-modified'] !== head.mtime) error = 'download error, Last-Modified has changed.';
  1517. // 如果
  1518. if (changed) {
  1519. FinishSize -= preAddSize;
  1520. onProgress({ loaded: FinishSize, total: FileSize });
  1521. ep.emit('error', {
  1522. code: 'ObjectHasChanged',
  1523. message: changed,
  1524. statusCode: data.statusCode,
  1525. header: chunkHeaders,
  1526. });
  1527. self.emit('inner-kill-task', { TaskId: TaskId });
  1528. } else {
  1529. FinishSize += chunkReadSize - preAddSize;
  1530. part.loaded = true;
  1531. onProgress({ loaded: FinishSize, total: FileSize });
  1532. tryCallback(err, data);
  1533. }
  1534. }
  1535. );
  1536. },
  1537. function (err, data) {
  1538. if (aborted) return;
  1539. nextChunk(err, data);
  1540. }
  1541. );
  1542. },
  1543. function (err, data) {
  1544. if (aborted) return;
  1545. onProgress({ loaded: FileSize, total: FileSize }, true);
  1546. if (err) return ep.emit('error', err);
  1547. ep.emit('download_chunks_complete');
  1548. }
  1549. );
  1550. });
  1551. // 下载已完成
  1552. ep.on('download_chunks_complete', function () {
  1553. callback(null, result);
  1554. });
  1555. // 监听 取消任务
  1556. var killTask = function (info) {
  1557. var killingTaskId = info.TaskId || '';
  1558. if (killingTaskId === TaskId) {
  1559. aborted = true;
  1560. }
  1561. };
  1562. TaskId && self.on('inner-kill-task', killTask);
  1563. ep.emit('get_file_info');
  1564. }
  1565. var API_MAP = {
  1566. sliceUploadFile: sliceUploadFile,
  1567. abortUploadTask: abortUploadTask,
  1568. uploadFile: uploadFile,
  1569. uploadFiles: uploadFiles,
  1570. sliceCopyFile: sliceCopyFile,
  1571. downloadFile: downloadFile,
  1572. };
  1573. module.exports.init = function (COS, task) {
  1574. task.transferToTaskMethod(API_MAP, 'sliceUploadFile');
  1575. util.each(API_MAP, function (fn, apiName) {
  1576. COS.prototype[apiName] = util.apiWrapper(apiName, fn);
  1577. });
  1578. };