From cec309d3c88a11b22ed7aa622700c791fa5ae04f Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 12:34:21 +0500 Subject: [PATCH 01/14] [Colossus] increase default timeout value in extrinsic wrapper --- storage-node/src/services/runtime/extrinsics.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/storage-node/src/services/runtime/extrinsics.ts b/storage-node/src/services/runtime/extrinsics.ts index 74f80f48b8..6ac4d52e79 100644 --- a/storage-node/src/services/runtime/extrinsics.ts +++ b/storage-node/src/services/runtime/extrinsics.ts @@ -263,7 +263,9 @@ export async function inviteStorageBucketOperator( async function extrinsicWrapper( extrinsic: () => Promise, throwErr = false, - timeoutMs = 25000 // 25s - default extrinsic timeout + // 5 mins - based on the default transactions validity of Substrate based chains with + // 6s block time: https://polkadot.js.org/docs/api/FAQ/#how-long-do-transactions-live + timeoutMs = 300_000 ): Promise { try { await timeout(extrinsic(), timeoutMs) From 701aff3c4caca87f5f2c80c9f2a150f4817f8a1e Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 12:36:13 +0500 Subject: [PATCH 02/14] [Colossus] update 'target' option in tsconfig.json --- storage-node/tsconfig.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/storage-node/tsconfig.json b/storage-node/tsconfig.json index d0cdf037fa..319873e6dd 100644 --- a/storage-node/tsconfig.json +++ b/storage-node/tsconfig.json @@ -7,7 +7,7 @@ "rootDir": "src", "strict": true, "strictNullChecks": true, - "target": "es2017", + "target": "es2020", "skipLibCheck": true, "baseUrl": ".", "esModuleInterop": true, From 6b82071fa6386aa5cd41750ece60633ae41e07a1 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 12:39:56 +0500 Subject: [PATCH 03/14] [Colossus] update graphql queries --- storage-node/src/services/queryNode/api.ts | 66 +++++++++++++------ .../queryNode/queries/queries.graphql | 16 +++++ 2 files changed, 62 insertions(+), 20 deletions(-) diff --git a/storage-node/src/services/queryNode/api.ts b/storage-node/src/services/queryNode/api.ts index 099322577d..42edf1fd24 100644 --- a/storage-node/src/services/queryNode/api.ts +++ b/storage-node/src/services/queryNode/api.ts @@ -1,47 +1,51 @@ import { ApolloClient, - NormalizedCacheObject, + DocumentNode, HttpLink, - defaultDataIdFromObject, InMemoryCache, - DocumentNode, - split, + NormalizedCacheObject, + defaultDataIdFromObject, from, + split, } from '@apollo/client' import { onError } from '@apollo/client/link/error' +import { WebSocketLink } from '@apollo/client/link/ws' +import { getMainDefinition } from '@apollo/client/utilities' import fetch from 'cross-fetch' +import stringify from 'fast-safe-stringify' +import ws from 'ws' +import logger from '../logger' import { + DataObjectDetailsFragment, + DataObjectsWithBagAndBucketsFragment, GetBagConnection, GetBagConnectionQuery, GetBagConnectionQueryVariables, + GetDataObjectConnection, + GetDataObjectConnectionQuery, + GetDataObjectConnectionQueryVariables, + GetDataObjectsByIds, + GetDataObjectsByIdsQuery, + GetDataObjectsByIdsQueryVariables, GetStorageBucketDetails, - GetStorageBucketDetailsQuery, + GetStorageBucketDetailsByWorkerId, GetStorageBucketDetailsByWorkerIdQuery, GetStorageBucketDetailsByWorkerIdQueryVariables, + GetStorageBucketDetailsQuery, GetStorageBucketDetailsQueryVariables, - StorageBucketDetailsFragment, - StorageBagDetailsFragment, - DataObjectDetailsFragment, - GetDataObjectConnectionQuery, - GetDataObjectConnectionQueryVariables, - GetDataObjectConnection, - StorageBucketIdsFragment, GetStorageBucketsConnection, GetStorageBucketsConnectionQuery, GetStorageBucketsConnectionQueryVariables, - GetStorageBucketDetailsByWorkerId, + QueryNodeState, + QueryNodeStateFields, QueryNodeStateFieldsFragment, QueryNodeStateSubscription, QueryNodeStateSubscriptionVariables, - QueryNodeState, - QueryNodeStateFields, + StorageBagDetailsFragment, + StorageBucketDetailsFragment, + StorageBucketIdsFragment, } from './generated/queries' import { Maybe, StorageBagWhereInput } from './generated/schema' -import { WebSocketLink } from '@apollo/client/link/ws' -import { getMainDefinition } from '@apollo/client/utilities' -import ws from 'ws' -import logger from '../logger' -import stringify from 'fast-safe-stringify' /** * Defines query paging limits. @@ -314,6 +318,28 @@ export class QueryNodeApi { return fullResult } + /** + * Returns data objects info by IDs. + * + * @param bagIds - query filter: data object IDs + */ + public async getDataObjectsByIds(ids: string[]): Promise> { + const allIds = [...ids] // Copy to avoid modifying the original array + const fullResult: DataObjectsWithBagAndBucketsFragment[] = [] + while (allIds.length) { + const idsBatch = allIds.splice(0, 1000) + fullResult.push( + ...((await this.multipleEntitiesQuery( + GetDataObjectsByIds, + { limit: MAX_RESULTS_PER_QUERY, ids: idsBatch }, + 'storageDataObjects' + )) || []) + ) + } + + return fullResult + } + /** * Returns storage bucket IDs. * diff --git a/storage-node/src/services/queryNode/queries/queries.graphql b/storage-node/src/services/queryNode/queries/queries.graphql index 0c1d9024b1..f4acd709f3 100644 --- a/storage-node/src/services/queryNode/queries/queries.graphql +++ b/storage-node/src/services/queryNode/queries/queries.graphql @@ -115,6 +115,22 @@ query getDataObjectConnection($bagIds: StorageBagWhereInput, $limit: Int, $curso } } +fragment DataObjectsWithBagAndBuckets on StorageDataObject { + id + storageBag { + id + storageBuckets { + id + } + } +} + +query getDataObjectsByIds($ids: [ID!], $limit: Int) { + storageDataObjects(where: { id_in: $ids }, limit: $limit) { + ...DataObjectsWithBagAndBuckets + } +} + fragment QueryNodeStateFields on ProcessorState { chainHead lastCompleteBlock From abf8d850185755fd28bed68f654119fe25d9d5b1 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 16:07:52 +0500 Subject: [PATCH 04/14] register new uloads in acceptPendingObjectsService --- .../src/services/sync/synchronizer.ts | 17 +++++++++++------ .../src/services/webApi/controllers/common.ts | 19 +++++++++++++++---- .../services/webApi/controllers/filesApi.ts | 18 +++++++----------- 3 files changed, 33 insertions(+), 21 deletions(-) diff --git a/storage-node/src/services/sync/synchronizer.ts b/storage-node/src/services/sync/synchronizer.ts index 30f041e217..1318a15b81 100644 --- a/storage-node/src/services/sync/synchronizer.ts +++ b/storage-node/src/services/sync/synchronizer.ts @@ -1,16 +1,21 @@ -import { getStorageObligationsFromRuntime, DataObligations } from './storageObligations' -import logger from '../../services/logger' -import { getDataObjectIDs } from '../../services/caching/localDataObjects' -import { SyncTask, DownloadFileTask, PrepareDownloadFileTask } from './tasks' -import { WorkingStack, TaskProcessorSpawner, TaskSink } from './workingProcess' -import _ from 'lodash' import { ApiPromise } from '@polkadot/api' +import _ from 'lodash' +import { getDataObjectIDs } from '../../services/caching/localDataObjects' +import logger from '../../services/logger' +import { DataObligations, getStorageObligationsFromRuntime } from './storageObligations' +import { DownloadFileTask, PrepareDownloadFileTask, SyncTask } from './tasks' +import { TaskProcessorSpawner, TaskSink, WorkingStack } from './workingProcess' /** * Temporary directory name for data uploading. */ export const TempDirName = 'temp' +/** + * Temporary Directory name for data objects not yet accepted (pending) in runtime. + */ +export const PendingDirName = 'pending' + /** * Runs the data synchronization workflow. It compares the current node's * storage obligations with the local storage and fixes the difference. diff --git a/storage-node/src/services/webApi/controllers/common.ts b/storage-node/src/services/webApi/controllers/common.ts index 3057c56768..a4aa63b575 100644 --- a/storage-node/src/services/webApi/controllers/common.ts +++ b/storage-node/src/services/webApi/controllers/common.ts @@ -1,10 +1,11 @@ -import * as express from 'express' -import { ExtrinsicFailedError } from '../../runtime/api' -import { BagIdValidationError } from '../../helpers/bagTypes' import { ApiPromise } from '@polkadot/api' import { KeyringPair } from '@polkadot/keyring/types' -import { ErrorResponse } from '../types' +import * as express from 'express' import { QueryNodeApi } from '../../../services/queryNode/api' +import { BagIdValidationError } from '../../helpers/bagTypes' +import { ExtrinsicFailedError } from '../../runtime/api' +import { AcceptPendingObjectsService } from '../../sync/acceptPendingObjects' +import { ErrorResponse } from '../types' /** * Dedicated error for the web api requests. @@ -126,6 +127,16 @@ export type AppConfig = { */ tempFileUploadingDir: string + /** + * Temporary directory for data objects in pending acceptance state + */ + pendingDataObjectsDir: string + + /** + * Service to periodically check for pending data objects, and send `accept_pending_data_objects` batch TXs + */ + acceptPendingObjectsService: AcceptPendingObjectsService + /** * Environment configuration */ diff --git a/storage-node/src/services/webApi/controllers/filesApi.ts b/storage-node/src/services/webApi/controllers/filesApi.ts index b5d44bc9ab..5c944771fd 100644 --- a/storage-node/src/services/webApi/controllers/filesApi.ts +++ b/storage-node/src/services/webApi/controllers/filesApi.ts @@ -109,8 +109,6 @@ export async function uploadFile( const fileObj = getFileObject(req) cleanupFileName = fileObj.path - const workerId = res.locals.workerId - const api = res.locals.api const bagId = parseBagId(uploadRequest.bagId) @@ -120,19 +118,17 @@ export async function uploadFile( // Prepare new file name const dataObjectId = uploadRequest.dataObjectId - const uploadsDir = res.locals.uploadsDir - const newPath = path.join(uploadsDir, dataObjectId) + const pendingObjectsDir = res.locals.pendingDataObjectsDir + const newPath = path.join(pendingObjectsDir, dataObjectId) - registerNewDataObjectId(dataObjectId) - await addDataObjectIdToCache(dataObjectId) - - // Overwrites existing file. + // Move file to pending objects Dir. await fsPromises.rename(fileObj.path, newPath) - cleanupFileName = newPath - await acceptPendingDataObjects(api, bagId, bucketKeyPair, workerId, new BN(uploadRequest.storageBucketId), [ + res.locals.acceptPendingObjectsService.push( new BN(uploadRequest.dataObjectId), - ]) + new BN(uploadRequest.storageBucketId), + uploadRequest.bagId + ) res.status(201).json({ id: hash, From 7d6772d3ad149e1ef1de915a46a4ef7f4d05cff6 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 16:22:57 +0500 Subject: [PATCH 05/14] update '/state/data' endpoint response to return pending data objects status too --- storage-node/src/api-spec/openapi.yaml | 6 ++++ .../services/webApi/controllers/stateApi.ts | 30 ++++++++++++++----- 2 files changed, 28 insertions(+), 8 deletions(-) diff --git a/storage-node/src/api-spec/openapi.yaml b/storage-node/src/api-spec/openapi.yaml index e2d3f50ff7..37d4624d71 100644 --- a/storage-node/src/api-spec/openapi.yaml +++ b/storage-node/src/api-spec/openapi.yaml @@ -281,6 +281,12 @@ components: tempDownloads: type: integer format: int64 + pendingDirSize: + type: integer + format: int64 + pendingObjects: + type: integer + format: int64 VersionResponse: type: object required: diff --git a/storage-node/src/services/webApi/controllers/stateApi.ts b/storage-node/src/services/webApi/controllers/stateApi.ts index 418fba329d..f9b3707974 100644 --- a/storage-node/src/services/webApi/controllers/stateApi.ts +++ b/storage-node/src/services/webApi/controllers/stateApi.ts @@ -1,12 +1,13 @@ -import { getDataObjectIDs } from '../../../services/caching/localDataObjects' import * as express from 'express' -import _ from 'lodash' -import { getDataObjectIDsByBagId } from '../../sync/storageObligations' -import { sendResponseWithError, AppConfig } from './common' import fastFolderSize from 'fast-folder-size' -import { promisify } from 'util' import fs from 'fs' +import _ from 'lodash' import NodeCache from 'node-cache' +import { promisify } from 'util' +import { getDataObjectIDs } from '../../../services/caching/localDataObjects' +import logger from '../../logger' +import { QueryNodeApi } from '../../queryNode/api' +import { getDataObjectIDsByBagId } from '../../sync/storageObligations' import { DataObjectResponse, DataStatsResponse, @@ -14,8 +15,7 @@ import { StatusResponse, VersionResponse, } from '../types' -import { QueryNodeApi } from '../../queryNode/api' -import logger from '../../logger' +import { AppConfig, sendResponseWithError } from './common' const fsPromises = fs.promises // Expiration period in seconds for the local cache. @@ -57,6 +57,7 @@ export async function getLocalDataStats( try { const uploadsDir = res.locals.uploadsDir const tempFileDir = res.locals.tempFileUploadingDir + const pendingObjectsDir = res.locals.pendingDataObjectsDir const fastFolderSizeAsync = promisify(fastFolderSize) const tempFolderExists = fs.existsSync(tempFileDir) @@ -68,6 +69,8 @@ export async function getLocalDataStats( let objectNumber = stats.length let tempDownloads = 0 let tempDirSize = 0 + let pendingObjects = 0 + let pendingDirSize = 0 if (tempFolderExists) { if (objectNumber > 0) { objectNumber-- @@ -75,11 +78,20 @@ export async function getLocalDataStats( const tempDirStatsPromise = fsPromises.readdir(tempFileDir) const tempDirSizePromise = fastFolderSizeAsync(tempFileDir) + const pendingDirStatsPromise = fsPromises.readdir(pendingObjectsDir) + const pendingDirSizePromise = fastFolderSizeAsync(pendingObjectsDir) - const [tempDirStats, tempSize] = await Promise.all([tempDirStatsPromise, tempDirSizePromise]) + const [tempDirStats, tempSize, pendingDirStats, pendingSize] = await Promise.all([ + tempDirStatsPromise, + tempDirSizePromise, + pendingDirStatsPromise, + pendingDirSizePromise, + ]) tempDirSize = tempSize ?? 0 tempDownloads = tempDirStats.length + pendingDirSize = pendingSize ?? 0 + pendingObjects = pendingDirStats.length } res.status(200).json({ @@ -87,6 +99,8 @@ export async function getLocalDataStats( totalSize: totalSize ?? 0, tempDownloads, tempDirSize, + pendingObjects, + pendingDirSize, }) } catch (err) { sendResponseWithError(res, next, err, 'local_data_stats') From 63012b70bddd97f10622f35678b0c582f7307d0c Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 16:56:25 +0500 Subject: [PATCH 06/14] update 'GET /files/{id}' endpoint to return the object status --- .../services/webApi/controllers/filesApi.ts | 41 ++++++++++--------- 1 file changed, 21 insertions(+), 20 deletions(-) diff --git a/storage-node/src/services/webApi/controllers/filesApi.ts b/storage-node/src/services/webApi/controllers/filesApi.ts index 5c944771fd..f9f7c82222 100644 --- a/storage-node/src/services/webApi/controllers/filesApi.ts +++ b/storage-node/src/services/webApi/controllers/filesApi.ts @@ -1,31 +1,28 @@ -import { acceptPendingDataObjects } from '../../runtime/extrinsics' -import { createUploadToken, verifyTokenSignature } from '../../helpers/auth' -import { hashFile } from '../../helpers/hashing' -import { registerNewDataObjectId } from '../../caching/newUploads' -import { addDataObjectIdToCache } from '../../caching/localDataObjects' -import { createNonce, getTokenExpirationTime } from '../../caching/tokenNonceKeeper' -import { getFileInfo } from '../../helpers/fileInfo' -import logger from '../../logger' import { ApiPromise } from '@polkadot/api' +import { PalletStorageBagIdType as BagId, PalletMembershipMembershipObject as Membership } from '@polkadot/types/lookup' +import { hexToString } from '@polkadot/util' +import BN from 'bn.js' import * as express from 'express' import fs from 'fs' import path from 'path' +import { timeout } from 'promise-timeout' import send from 'send' -import { hexToString } from '@polkadot/util' +import { QueryNodeApi } from '../../../services/queryNode/api' +import { createNonce, getTokenExpirationTime } from '../../caching/tokenNonceKeeper' +import { createUploadToken, verifyTokenSignature } from '../../helpers/auth' import { parseBagId } from '../../helpers/bagTypes' -import { timeout } from 'promise-timeout' -import { WebApiError, sendResponseWithError, getHttpStatusCodeByError, AppConfig } from './common' +import { getFileInfo } from '../../helpers/fileInfo' +import { hashFile } from '../../helpers/hashing' +import logger from '../../logger' import { getStorageBucketIdsByWorkerId } from '../../sync/storageObligations' -import { PalletMembershipMembershipObject as Membership, PalletStorageBagIdType as BagId } from '@polkadot/types/lookup' -import BN from 'bn.js' import { + GetFileHeadersRequestParams, + GetFileRequestParams, UploadFileQueryParams, - UploadTokenRequest, UploadTokenBody, - GetFileRequestParams, - GetFileHeadersRequestParams, + UploadTokenRequest, } from '../types' -import { QueryNodeApi } from '../../../services/queryNode/api' +import { AppConfig, WebApiError, getHttpStatusCodeByError, sendResponseWithError } from './common' const fsPromises = fs.promises /** @@ -39,7 +36,10 @@ export async function getFile( try { const dataObjectId = new BN(req.params.id) const uploadsDir = res.locals.uploadsDir - const fullPath = path.resolve(uploadsDir, dataObjectId.toString()) + const pendingObjectsDir = res.locals.pendingDataObjectsDir + const pending = res.locals.acceptPendingObjectsService.getPendingDataObject(dataObjectId.toString()) + + const fullPath = path.resolve(pending ? pendingObjectsDir : uploadsDir, dataObjectId.toString()) const fileInfo = await getFileInfo(fullPath) const fileStats = await fsPromises.stat(fullPath) @@ -51,6 +51,7 @@ export async function getFile( res.setHeader('Content-Disposition', 'inline') res.setHeader('Content-Type', fileInfo.mimeType) res.setHeader('Content-Length', fileStats.size) + res.setHeader('X-Status', pending ? 'pending' : 'accepted') }) stream.on('error', (err) => { @@ -125,8 +126,8 @@ export async function uploadFile( await fsPromises.rename(fileObj.path, newPath) res.locals.acceptPendingObjectsService.push( - new BN(uploadRequest.dataObjectId), - new BN(uploadRequest.storageBucketId), + uploadRequest.dataObjectId, + uploadRequest.storageBucketId, uploadRequest.bagId ) From 51901e12660b6e7cd5589cd29443da7a6db17bd5 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Tue, 21 Nov 2023 17:00:51 +0500 Subject: [PATCH 07/14] [Colossus] add 'AcceptPendingObjectsService' to periodically send accept pending data objects txs --- storage-node/src/commands/server.ts | 50 +++-- .../src/services/runtime/extrinsics.ts | 60 +++++ .../src/services/sync/acceptPendingObjects.ts | 212 ++++++++++++++++++ 3 files changed, 308 insertions(+), 14 deletions(-) create mode 100644 storage-node/src/services/sync/acceptPendingObjects.ts diff --git a/storage-node/src/commands/server.ts b/storage-node/src/commands/server.ts index 531d2b4b08..0fb0706a54 100644 --- a/storage-node/src/commands/server.ts +++ b/storage-node/src/commands/server.ts @@ -1,23 +1,24 @@ import { flags } from '@oclif/command' -import { createApp } from '../services/webApi/app' -import ApiCommandBase from '../command-base/ApiCommandBase' -import logger, { initNewLogger, DatePatternByFrequency, Frequency } from '../services/logger' -import { loadDataObjectIdCache } from '../services/caching/localDataObjects' import { ApiPromise } from '@polkadot/api' -import { performSync, TempDirName } from '../services/sync/synchronizer' -import sleep from 'sleep-promise' -import rimraf from 'rimraf' +import { KeyringPair } from '@polkadot/keyring/types' +import { Option } from '@polkadot/types-codec' +import { PalletStorageStorageBucketRecord } from '@polkadot/types/lookup' +import fs from 'fs' import _ from 'lodash' import path from 'path' +import rimraf from 'rimraf' +import sleep from 'sleep-promise' import { promisify } from 'util' -import ExitCodes from './../command-base/ExitCodes' -import fs from 'fs' -import { getStorageBucketIdsByWorkerId } from '../services/sync/storageObligations' -import { PalletStorageStorageBucketRecord } from '@polkadot/types/lookup' -import { Option } from '@polkadot/types-codec' -import { QueryNodeApi } from '../services/queryNode/api' -import { KeyringPair } from '@polkadot/keyring/types' +import ApiCommandBase from '../command-base/ApiCommandBase' import { customFlags } from '../command-base/CustomFlags' +import { loadDataObjectIdCache } from '../services/caching/localDataObjects' +import logger, { DatePatternByFrequency, Frequency, initNewLogger } from '../services/logger' +import { QueryNodeApi } from '../services/queryNode/api' +import { AcceptPendingObjectsService } from '../services/sync/acceptPendingObjects' +import { getStorageBucketIdsByWorkerId } from '../services/sync/storageObligations' +import { PendingDirName, TempDirName, performSync } from '../services/sync/synchronizer' +import { createApp } from '../services/webApi/app' +import ExitCodes from './../command-base/ExitCodes' const fsPromises = fs.promises /** @@ -127,6 +128,11 @@ Supported values: warn, error, debug, info. Default:debug`, default: 'daily', required: false, }), + maxBatchTxSize: flags.integer({ + description: 'Maximum number of `accept_pending_data_objects` in a batch transactions.', + default: 10, + required: false, + }), ...ApiCommandBase.flags, } @@ -205,6 +211,20 @@ Supported values: warn, error, debug, info. Default:debug`, await this.ensureDevelopmentChain() } + const pendingDataObjectsDir = path.join(flags.uploads, PendingDirName) + + const acceptPendingObjectsService = await AcceptPendingObjectsService.init( + api, + qnApi, + workerId, + flags.uploads, + pendingDataObjectsDir, + bucketKeyPairs, + writableBuckets, + flags.maxBatchTxSize, + 6000 // Every block + ) + // Don't run sync job if no buckets selected, to prevent purging // any assets. if (flags.sync && selectedBuckets.length) { @@ -244,6 +264,8 @@ Supported values: warn, error, debug, info. Default:debug`, maxFileSize, uploadsDir: flags.uploads, tempFileUploadingDir, + pendingDataObjectsDir, + acceptPendingObjectsService, process: this.config, enableUploadingAuth, downloadBuckets: selectedBuckets, diff --git a/storage-node/src/services/runtime/extrinsics.ts b/storage-node/src/services/runtime/extrinsics.ts index 6ac4d52e79..0b003a2715 100644 --- a/storage-node/src/services/runtime/extrinsics.ts +++ b/storage-node/src/services/runtime/extrinsics.ts @@ -5,6 +5,8 @@ import { PalletStorageBagIdType as BagId, PalletStorageDynamicBagType as Dynamic import BN from 'bn.js' import { timeout } from 'promise-timeout' import logger from '../../services/logger' +import { parseBagId } from '../helpers/bagTypes' +import { AcceptPendingDataObjectsParams } from '../sync/acceptPendingObjects' import { formatDispatchError, getEvent, getEvents, sendAndFollowNamedTx } from './api' /** @@ -129,6 +131,64 @@ export async function updateStorageBucketsForBags( return [success, failedCalls] } +/** + * Accepts pending data objects by storage provider in batch transaction. + * + * @remarks + * It sends an batch extrinsic to the runtime. + * + * @param api - runtime API promise + * @param workerId - runtime storage provider ID (worker ID) + * @param acceptPendingDataObjectsParams - acceptPendingDataObject extrinsic parameters + * @returns promise with a list of failedCalls. + */ +export async function acceptPendingDataObjectsBatch( + api: ApiPromise, + workerId: number, + acceptPendingDataObjectsParams: AcceptPendingDataObjectsParams[] +): Promise { + // a list of failed data objects + const failedDataObjects: string[] = [] + + const txsByTransactorAccount = acceptPendingDataObjectsParams.map(({ account, storageBucket }) => { + const txs = storageBucket.bags.map((bag) => + api.tx.storage.acceptPendingDataObjects( + workerId, + storageBucket.id, + parseBagId(bag.id), + api.createType('BTreeSet', bag.dataObjects) + ) + ) + + return [account, txs, storageBucket.bags] as const + }) + + for (const [account, txs, bags] of txsByTransactorAccount) { + const txBatch = api.tx.utility.forceBatch(txs) + + const success = await extrinsicWrapper(async () => { + await sendAndFollowNamedTx(api, account, txBatch, (result) => { + // Process individual ItemFailed events + const events = getEvents(result, 'utility', ['ItemCompleted', 'ItemFailed']) + events.forEach((e, i) => { + if (e.method === 'ItemFailed') { + failedDataObjects.push(...bags[i].dataObjects.toString()) + } + }) + }) + }) + + if (!success) { + // If the batch transaction failed, push all data objects to failed list + bags.forEach((bag) => { + failedDataObjects.push(...bag.dataObjects.toString()) + }) + } + } + + return failedDataObjects +} + /** * Accepts pending data objects by storage provider. * diff --git a/storage-node/src/services/sync/acceptPendingObjects.ts b/storage-node/src/services/sync/acceptPendingObjects.ts new file mode 100644 index 0000000000..7d5c6fbe69 --- /dev/null +++ b/storage-node/src/services/sync/acceptPendingObjects.ts @@ -0,0 +1,212 @@ +import { ApiPromise } from '@polkadot/api' +import { KeyringPair } from '@polkadot/keyring/types' +import fs from 'fs' +import path from 'path' +import { addDataObjectIdToCache } from '../caching/localDataObjects' +import { registerNewDataObjectId } from '../caching/newUploads' +import logger from '../logger' +import { QueryNodeApi } from '../queryNode/api' +import { acceptPendingDataObjectsBatch } from '../runtime/extrinsics' +const fsPromises = fs.promises + +export type AcceptPendingDataObjectsParams = { + account: KeyringPair + storageBucket: { + id: string + bags: { + id: string + dataObjects: string[] + }[] + } +} + +export class AcceptPendingObjectsService { + private pendingDataObjects: Map // dataObjectId -> [storageBucketId, bagId] + + private constructor( + private qnApi: QueryNodeApi, + private pendingDataObjectsDir: string, + private bucketKeyPairs: Map, + private uploadBuckets: string[] + ) { + this.pendingDataObjects = new Map() + } + + static async init( + api: ApiPromise, + qnApi: QueryNodeApi, + workerId: number, + uploadsDir: string, + pendingDataObjectsDir: string, + bucketKeyPairs: Map, + uploadBuckets: string[], + maxTxBatchSize: number, + intervalMs: number + ): Promise { + const acceptPendingObjectsService = new AcceptPendingObjectsService( + qnApi, + pendingDataObjectsDir, + bucketKeyPairs, + uploadBuckets + ) + + // Load pending data objects from the pending directory + await acceptPendingObjectsService.loadPendingDataObjects() + + const runWithInterval = () => { + acceptPendingObjectsService + .acceptPendingDataObjects(api, workerId, uploadsDir, maxTxBatchSize) + .catch((err) => logger.error(`Failed to register pending data objects as accepted in runtime: ${err}`)) + + setTimeout(runWithInterval, intervalMs) + } + runWithInterval() + + return acceptPendingObjectsService + } + + getPendingDataObject(objectId: string): [string, string] | undefined { + return this.pendingDataObjects.get(objectId) + } + + private async loadPendingDataObjectsFromIDs(pendingIds: string[]): Promise { + const pendingDataObjects = await this.qnApi.getDataObjectsByIds(pendingIds) + + pendingDataObjects.forEach((dataObject) => { + const storageBucket = dataObject.storageBag.storageBuckets.find(({ id }) => this.uploadBuckets.includes(id)) + if (storageBucket) { + this.push(dataObject.id, storageBucket.id, dataObject.storageBag.id) + } else { + logger.warn( + `Data object ${dataObject.id} in pending directory is not assigned to any of the upload buckets: ${this.uploadBuckets}.` + ) + } + }) + } + + private async loadPendingDataObjects(): Promise { + const pendingIds = await this.getPendingObjectsFromLocalDir(this.pendingDataObjectsDir) + + await this.loadPendingDataObjectsFromIDs(pendingIds) + + logger.debug(`Pending data objects ID cache loaded.`) + } + + public push(dataObjectId: string, storageBucketId: string, bagId: string): void { + this.pendingDataObjects.set(dataObjectId, [storageBucketId, bagId]) + } + + private async popN(n: number): Promise { + const params: AcceptPendingDataObjectsParams[] = [] + let count = 0 + + while (count < n && this.pendingDataObjects.size > 0) { + // Extract the first element from the map + const [dataObjectId, [storageBucketId, bagId]]: [string, [string, string]] = this.pendingDataObjects + .entries() + .next().value + this.pendingDataObjects.delete(dataObjectId) + + // Find or create the storage bucket in the params array + let storageBucket = params.find((p) => p.storageBucket.id === storageBucketId) + if (!storageBucket) { + const account = this.bucketKeyPairs.get(storageBucketId.toString()) + if (!account) { + logger.error(`No key pair found for storage bucket ${storageBucketId}.`) + continue + } + + storageBucket = { + account, + storageBucket: { + id: storageBucketId, + bags: [], + }, + } + params.push(storageBucket) + } + + // Find or create the bag in the storage bucket in the params array + let bag = storageBucket.storageBucket.bags.find((b) => b.id === bagId) + if (!bag) { + bag = { + id: bagId, + dataObjects: [], + } + storageBucket.storageBucket.bags.push(bag) + } + + // Add the data object to the bag in the params array + bag.dataObjects.push(dataObjectId) + + count++ + } + + return params + } + + /** + * Returns file names from the pending objects directory. + * + * @param directory - local directory to get file names from + */ + private async getPendingObjectsFromLocalDir(directory: string): Promise { + // Check if directory exists and if not, create it + if (!fs.existsSync(directory)) { + await fsPromises.mkdir(directory) + } + + // Read the directory contents + return await fsPromises.readdir(directory) + } + + private async acceptPendingDataObjects( + api: ApiPromise, + workerId: number, + uploadsDir: string, + maxTxBatchSize: number + ): Promise { + const params = await this.popN(maxTxBatchSize) + + let failedObjectsIds: string[] + try { + failedObjectsIds = await acceptPendingDataObjectsBatch(api, workerId, params) + } catch (err) { + // Re-push all params as failed since the entire batch failed + params.forEach((param) => { + param.storageBucket.bags.forEach((bag) => { + bag.dataObjects.forEach((dataObjectId) => { + this.push(dataObjectId, param.storageBucket.id, bag.id) + }) + }) + }) + + throw err + } + + // Handle failed calls + await this.loadPendingDataObjectsFromIDs(failedObjectsIds) + + const successfulObjectsIds = params.flatMap((param) => + param.storageBucket.bags.flatMap((bag) => + bag.dataObjects.filter( + (dataObjectId) => !failedObjectsIds.some((failedId) => failedId === dataObjectId.toString()) + ) + ) + ) + + // Handle successful calls + for (const dataObjectId of successfulObjectsIds) { + const dataObjectIdStr = dataObjectId.toString() + const currentPath = path.join(this.pendingDataObjectsDir, dataObjectIdStr) + const newPath = path.join(uploadsDir, dataObjectIdStr) + registerNewDataObjectId(dataObjectIdStr) + await addDataObjectIdToCache(dataObjectIdStr) + await fsPromises.rename(currentPath, newPath) + } + + if (successfulObjectsIds.length > 0) { + logger.info(`Successfully registered ${successfulObjectsIds.length} pending data objects as accepted in runtime.`) + } + } +} From b9686a74010fec363690eb4d751d6e6256afc2bb Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Wed, 29 Nov 2023 19:09:22 +0500 Subject: [PATCH 08/14] dont accept re-upload of already uploaded objects --- storage-node/src/services/webApi/app.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/storage-node/src/services/webApi/app.ts b/storage-node/src/services/webApi/app.ts index f82ac67cea..7576a137bd 100644 --- a/storage-node/src/services/webApi/app.ts +++ b/storage-node/src/services/webApi/app.ts @@ -245,4 +245,9 @@ async function validateUploadFileParams(req: express.Request, res: express.Respo if (dataObject.accepted.valueOf()) { throw new WebApiError(`Data object ${dataObjectId} has already been accepted by storage node`, 400) } + + const isObjectPending = res.locals.acceptPendingObjectsService.getPendingDataObject(dataObjectId.toString()) + if (isObjectPending) { + throw new WebApiError(`Data object ${dataObjectId} already exists`, 400) + } } From cdc4ea7c88f49fe68950b23bbb4862dc20b04504 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Wed, 29 Nov 2023 19:12:42 +0500 Subject: [PATCH 09/14] ignore pdending dir from uploads dir when loading data objects --- storage-node/src/commands/server.ts | 2 +- .../src/services/caching/localDataObjects.ts | 13 ++++++++++--- 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/storage-node/src/commands/server.ts b/storage-node/src/commands/server.ts index 0fb0706a54..48caabb677 100644 --- a/storage-node/src/commands/server.ts +++ b/storage-node/src/commands/server.ts @@ -204,7 +204,7 @@ Supported values: warn, error, debug, info. Default:debug`, await recreateTempDirectory(flags.uploads, TempDirName) if (fs.existsSync(flags.uploads)) { - await loadDataObjectIdCache(flags.uploads, TempDirName) + await loadDataObjectIdCache(flags.uploads, TempDirName, PendingDirName) } if (flags.dev) { diff --git a/storage-node/src/services/caching/localDataObjects.ts b/storage-node/src/services/caching/localDataObjects.ts index 604a62c8b0..8cd83ffea8 100644 --- a/storage-node/src/services/caching/localDataObjects.ts +++ b/storage-node/src/services/caching/localDataObjects.ts @@ -31,13 +31,20 @@ export async function getDataObjectIDs(): Promise { * @param uploadDir - uploading directory * @param tempDirName - temp directory name */ -export async function loadDataObjectIdCache(uploadDir: string, tempDirName: string): Promise { +export async function loadDataObjectIdCache( + uploadDir: string, + tempDirName: string, + pendingDirName: string +): Promise { await lock.acquireAsync() const localIds = await getLocalFileNames(uploadDir) - // Filter temporary directory name. + // Filter temporary & pending directory name. const tempDirectoryName = path.parse(tempDirName).name - const ids = localIds.filter((dataObjectId) => dataObjectId !== tempDirectoryName) + const pendingDirectoryName = path.parse(pendingDirName).name + const ids = localIds.filter( + (dataObjectId) => dataObjectId !== tempDirectoryName && dataObjectId !== pendingDirectoryName + ) idCache = new Set(ids) From 205073aa7a047f223141113486d7129a954ea2ee Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Wed, 29 Nov 2023 19:20:48 +0500 Subject: [PATCH 10/14] make AcceptPendingObjectsService a singleton class --- storage-node/src/commands/server.ts | 2 +- .../src/services/sync/acceptPendingObjects.ts | 41 +++++++++++-------- 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/storage-node/src/commands/server.ts b/storage-node/src/commands/server.ts index 48caabb677..9e87685274 100644 --- a/storage-node/src/commands/server.ts +++ b/storage-node/src/commands/server.ts @@ -213,7 +213,7 @@ Supported values: warn, error, debug, info. Default:debug`, const pendingDataObjectsDir = path.join(flags.uploads, PendingDirName) - const acceptPendingObjectsService = await AcceptPendingObjectsService.init( + const acceptPendingObjectsService = await AcceptPendingObjectsService.create( api, qnApi, workerId, diff --git a/storage-node/src/services/sync/acceptPendingObjects.ts b/storage-node/src/services/sync/acceptPendingObjects.ts index 7d5c6fbe69..bbc97220f5 100644 --- a/storage-node/src/services/sync/acceptPendingObjects.ts +++ b/storage-node/src/services/sync/acceptPendingObjects.ts @@ -21,18 +21,20 @@ export type AcceptPendingDataObjectsParams = { } export class AcceptPendingObjectsService { + private static instance: AcceptPendingObjectsService | null = null private pendingDataObjects: Map // dataObjectId -> [storageBucketId, bagId] private constructor( private qnApi: QueryNodeApi, private pendingDataObjectsDir: string, + private uploadsDir: string, private bucketKeyPairs: Map, private uploadBuckets: string[] ) { this.pendingDataObjects = new Map() } - static async init( + public static async create( api: ApiPromise, qnApi: QueryNodeApi, workerId: number, @@ -43,26 +45,29 @@ export class AcceptPendingObjectsService { maxTxBatchSize: number, intervalMs: number ): Promise { - const acceptPendingObjectsService = new AcceptPendingObjectsService( - qnApi, - pendingDataObjectsDir, - bucketKeyPairs, - uploadBuckets - ) - - // Load pending data objects from the pending directory - await acceptPendingObjectsService.loadPendingDataObjects() + if (this.instance === null) { + this.instance = new AcceptPendingObjectsService( + qnApi, + pendingDataObjectsDir, + uploadsDir, + bucketKeyPairs, + uploadBuckets + ) + await this.instance.loadPendingDataObjects() + this.instance.runWithInterval(api, workerId, maxTxBatchSize, intervalMs) + } + return this.instance + } - const runWithInterval = () => { - acceptPendingObjectsService - .acceptPendingDataObjects(api, workerId, uploadsDir, maxTxBatchSize) - .catch((err) => logger.error(`Failed to register pending data objects as accepted in runtime: ${err}`)) + private runWithInterval(api: ApiPromise, workerId: number, maxTxBatchSize: number, intervalMs: number) { + const run = () => { + this.acceptPendingDataObjects(api, workerId, maxTxBatchSize).catch((err) => + logger.error(`Failed to register pending data objects as accepted in runtime: ${err}`) + ) - setTimeout(runWithInterval, intervalMs) + setTimeout(run, intervalMs) } - runWithInterval() - - return acceptPendingObjectsService + run() } getPendingDataObject(objectId: string): [string, string] | undefined { From c1c0a598f5abc828837558476f160e25fbb92544 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Wed, 29 Nov 2023 19:22:33 +0500 Subject: [PATCH 11/14] rename push & popN methods --- .../src/services/sync/acceptPendingObjects.ts | 41 ++++++++++--------- .../services/webApi/controllers/filesApi.ts | 2 +- 2 files changed, 23 insertions(+), 20 deletions(-) diff --git a/storage-node/src/services/sync/acceptPendingObjects.ts b/storage-node/src/services/sync/acceptPendingObjects.ts index bbc97220f5..d3b38b4b59 100644 --- a/storage-node/src/services/sync/acceptPendingObjects.ts +++ b/storage-node/src/services/sync/acceptPendingObjects.ts @@ -1,6 +1,7 @@ import { ApiPromise } from '@polkadot/api' import { KeyringPair } from '@polkadot/keyring/types' import fs from 'fs' +import _ from 'lodash' import path from 'path' import { addDataObjectIdToCache } from '../caching/localDataObjects' import { registerNewDataObjectId } from '../caching/newUploads' @@ -101,21 +102,30 @@ export class AcceptPendingObjectsService { this.pendingDataObjects.set(dataObjectId, [storageBucketId, bagId]) } - private async popN(n: number): Promise { - const params: AcceptPendingDataObjectsParams[] = [] + private take(n: number): Array<[string, [string, string]]> { + const pending: Array<[string, [string, string]]> = [] let count = 0 - while (count < n && this.pendingDataObjects.size > 0) { - // Extract the first element from the map - const [dataObjectId, [storageBucketId, bagId]]: [string, [string, string]] = this.pendingDataObjects - .entries() - .next().value + for (const [dataObjectId, bucketAndBag] of this.pendingDataObjects.entries()) { + if (count >= n) break + pending.push([dataObjectId, bucketAndBag]) this.pendingDataObjects.delete(dataObjectId) + count++ + } + + return pending + } + + private async createAcceptPendingObjectsParams( + pendingObjects: Array<[string, [string, string]]> + ): Promise { + const params: AcceptPendingDataObjectsParams[] = [] - // Find or create the storage bucket in the params array + // Find or create the storage bucket in the params array + for (const [dataObjectId, [storageBucketId, bagId]] of pendingObjects) { let storageBucket = params.find((p) => p.storageBucket.id === storageBucketId) if (!storageBucket) { - const account = this.bucketKeyPairs.get(storageBucketId.toString()) + const account = this.bucketKeyPairs.get(storageBucketId) if (!account) { logger.error(`No key pair found for storage bucket ${storageBucketId}.`) continue @@ -143,8 +153,6 @@ export class AcceptPendingObjectsService { // Add the data object to the bag in the params array bag.dataObjects.push(dataObjectId) - - count++ } return params @@ -165,13 +173,8 @@ export class AcceptPendingObjectsService { return await fsPromises.readdir(directory) } - private async acceptPendingDataObjects( - api: ApiPromise, - workerId: number, - uploadsDir: string, - maxTxBatchSize: number - ): Promise { - const params = await this.popN(maxTxBatchSize) + private async acceptPendingDataObjects(api: ApiPromise, workerId: number, maxTxBatchSize: number): Promise { + const params = await this.createAcceptPendingObjectsParams(this.take(maxTxBatchSize)) let failedObjectsIds: string[] try { @@ -181,7 +184,7 @@ export class AcceptPendingObjectsService { params.forEach((param) => { param.storageBucket.bags.forEach((bag) => { bag.dataObjects.forEach((dataObjectId) => { - this.push(dataObjectId, param.storageBucket.id, bag.id) + this.add(dataObjectId, param.storageBucket.id, bag.id) }) }) }) diff --git a/storage-node/src/services/webApi/controllers/filesApi.ts b/storage-node/src/services/webApi/controllers/filesApi.ts index f9f7c82222..fb5e37162f 100644 --- a/storage-node/src/services/webApi/controllers/filesApi.ts +++ b/storage-node/src/services/webApi/controllers/filesApi.ts @@ -125,7 +125,7 @@ export async function uploadFile( // Move file to pending objects Dir. await fsPromises.rename(fileObj.path, newPath) - res.locals.acceptPendingObjectsService.push( + res.locals.acceptPendingObjectsService.add( uploadRequest.dataObjectId, uploadRequest.storageBucketId, uploadRequest.bagId From b9a65a4826b7208edff70c85268a8724ea5836fe Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Wed, 29 Nov 2023 19:26:45 +0500 Subject: [PATCH 12/14] remove stale dataobjects from pending dir & move accepted dataobjects to uploads dir --- .../queryNode/queries/queries.graphql | 1 + .../src/services/sync/acceptPendingObjects.ts | 65 ++++++++++++++----- 2 files changed, 49 insertions(+), 17 deletions(-) diff --git a/storage-node/src/services/queryNode/queries/queries.graphql b/storage-node/src/services/queryNode/queries/queries.graphql index f4acd709f3..199028a7f4 100644 --- a/storage-node/src/services/queryNode/queries/queries.graphql +++ b/storage-node/src/services/queryNode/queries/queries.graphql @@ -117,6 +117,7 @@ query getDataObjectConnection($bagIds: StorageBagWhereInput, $limit: Int, $curso fragment DataObjectsWithBagAndBuckets on StorageDataObject { id + isAccepted storageBag { id storageBuckets { diff --git a/storage-node/src/services/sync/acceptPendingObjects.ts b/storage-node/src/services/sync/acceptPendingObjects.ts index d3b38b4b59..b3a50d9642 100644 --- a/storage-node/src/services/sync/acceptPendingObjects.ts +++ b/storage-node/src/services/sync/acceptPendingObjects.ts @@ -78,16 +78,31 @@ export class AcceptPendingObjectsService { private async loadPendingDataObjectsFromIDs(pendingIds: string[]): Promise { const pendingDataObjects = await this.qnApi.getDataObjectsByIds(pendingIds) - pendingDataObjects.forEach((dataObject) => { - const storageBucket = dataObject.storageBag.storageBuckets.find(({ id }) => this.uploadBuckets.includes(id)) - if (storageBucket) { - this.push(dataObject.id, storageBucket.id, dataObject.storageBag.id) - } else { - logger.warn( - `Data object ${dataObject.id} in pending directory is not assigned to any of the upload buckets: ${this.uploadBuckets}.` - ) - } - }) + const deletedObjects = _.differenceWith(pendingIds, pendingDataObjects, (id, dataObject) => dataObject.id === id) + + // Remove stale objects from the pending directory + if (deletedObjects.length) { + logger.warn(`Found data objects deleted from runtime in Pending directory: ${deletedObjects.length}`) + await Promise.all(deletedObjects.map((id) => fsPromises.unlink(path.join(this.pendingDataObjectsDir, id)))) + } + + await Promise.all( + pendingDataObjects.map(async (dataObject) => { + const storageBucket = dataObject.storageBag.storageBuckets.find(({ id }) => this.uploadBuckets.includes(id)) + + if (storageBucket) { + if (dataObject.isAccepted) { + await this.movePendingDataObjectToUploadsDir(dataObject.id) + } else { + this.add(dataObject.id, storageBucket.id, dataObject.storageBag.id) + } + } else { + logger.warn( + `Data object ${dataObject.id} in pending directory is no longer assigned to any of the upload buckets: ${this.uploadBuckets}.` + ) + } + }) + ) } private async loadPendingDataObjects(): Promise { @@ -98,7 +113,28 @@ export class AcceptPendingObjectsService { logger.debug(`Pending data objects ID cache loaded.`) } - public push(dataObjectId: string, storageBucketId: string, bagId: string): void { + private async movePendingDataObjectToUploadsDir(dataObjectId: string): Promise { + const currentPath = path.join(this.pendingDataObjectsDir, dataObjectId) + const newPath = path.join(this.uploadsDir, dataObjectId) + + try { + // Check if the file already exists in the uploads directory (i.e. synced from other operators) + try { + await fsPromises.access(newPath, fs.constants.F_OK) + logger.warn(`File ${dataObjectId} already exists in uploads directory. Deleting current file.`) + await fsPromises.unlink(currentPath) + } catch { + // If the file does not exist in the uploads directory, proceed with the rename + registerNewDataObjectId(dataObjectId) + await addDataObjectIdToCache(dataObjectId) + await fsPromises.rename(currentPath, newPath) + } + } catch (err) { + logger.error(`Error handling data object ${dataObjectId}: ${err}`) + } + } + + public add(dataObjectId: string, storageBucketId: string, bagId: string): void { this.pendingDataObjects.set(dataObjectId, [storageBucketId, bagId]) } @@ -205,12 +241,7 @@ export class AcceptPendingObjectsService { // Handle successful calls for (const dataObjectId of successfulObjectsIds) { - const dataObjectIdStr = dataObjectId.toString() - const currentPath = path.join(this.pendingDataObjectsDir, dataObjectIdStr) - const newPath = path.join(uploadsDir, dataObjectIdStr) - registerNewDataObjectId(dataObjectIdStr) - await addDataObjectIdToCache(dataObjectIdStr) - await fsPromises.rename(currentPath, newPath) + await this.movePendingDataObjectToUploadsDir(dataObjectId.toString()) } if (successfulObjectsIds.length > 0) { From 84421316f732645bb15e99ea61223a170b45f257 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Thu, 30 Nov 2023 15:52:52 +0500 Subject: [PATCH 13/14] address CR --- storage-node/src/services/sync/acceptPendingObjects.ts | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/storage-node/src/services/sync/acceptPendingObjects.ts b/storage-node/src/services/sync/acceptPendingObjects.ts index b3a50d9642..60a702d822 100644 --- a/storage-node/src/services/sync/acceptPendingObjects.ts +++ b/storage-node/src/services/sync/acceptPendingObjects.ts @@ -62,11 +62,9 @@ export class AcceptPendingObjectsService { private runWithInterval(api: ApiPromise, workerId: number, maxTxBatchSize: number, intervalMs: number) { const run = () => { - this.acceptPendingDataObjects(api, workerId, maxTxBatchSize).catch((err) => - logger.error(`Failed to register pending data objects as accepted in runtime: ${err}`) - ) - - setTimeout(run, intervalMs) + this.acceptPendingDataObjects(api, workerId, maxTxBatchSize) + .catch((err) => logger.error(`Failed to register pending data objects as accepted in runtime: ${err}`)) + .finally(() => setTimeout(run, intervalMs)) } run() } From 426a2354f89b521cea04b6a204bfd70d4ce22422 Mon Sep 17 00:00:00 2001 From: Zeeshan Akram <97m.zeeshan@gmail.com> Date: Thu, 30 Nov 2023 16:10:12 +0500 Subject: [PATCH 14/14] bumped package version & updated changelog --- storage-node/CHANGELOG.md | 6 ++++++ storage-node/package.json | 2 +- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/storage-node/CHANGELOG.md b/storage-node/CHANGELOG.md index b8f61390d5..7f4f723b91 100644 --- a/storage-node/CHANGELOG.md +++ b/storage-node/CHANGELOG.md @@ -1,3 +1,9 @@ +### 3.9.0 + +- Added new `AcceptPendingObjectsService` that is responsible for periodically sending batch `accept_pending_data_objects` for all the pending data objects. The `POST /files` endpoint now no longer calls the `accept_pending_data_objects` extrinsic for individual uploads, instead, it registers all the pending objects with `AcceptPendingObjectsService` +- Updated `/state/data` endpoint response headers to return data objects status too i.e. (`pending` or `accepted`) +- **FIX**: Increase the default timeout value in the `extrinsicWrapper` function to match the transaction validity in the transaction pool + ### 3.8.1 - Hotfix: Fix call stack size exceeded when handling large number of initial object to sync. diff --git a/storage-node/package.json b/storage-node/package.json index 72e8bfdb47..2b955dd8f6 100644 --- a/storage-node/package.json +++ b/storage-node/package.json @@ -1,7 +1,7 @@ { "name": "storage-node", "description": "Joystream storage subsystem.", - "version": "3.8.1", + "version": "3.9.0", "author": "Joystream contributors", "bin": { "storage-node": "./bin/run"