mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-28 04:25:03 +02:00
UBERF-8508: Get rid of Mongo in storage adapter (#6989)
Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
+15
-67
@@ -16,7 +16,12 @@
|
||||
import { type Attachment } from '@hcengineering/attachment'
|
||||
import { type Blob, type MeasureContext, type Ref, type WorkspaceId, RateLimiter } from '@hcengineering/core'
|
||||
import { DOMAIN_ATTACHMENT } from '@hcengineering/model-attachment'
|
||||
import { type ListBlobResult, type StorageAdapter, type StorageAdapterEx } from '@hcengineering/server-core'
|
||||
import {
|
||||
type ListBlobResult,
|
||||
type StorageAdapter,
|
||||
type StorageAdapterEx,
|
||||
type UploadedObjectInfo
|
||||
} from '@hcengineering/server-core'
|
||||
import { type Db } from 'mongodb'
|
||||
import { PassThrough } from 'stream'
|
||||
|
||||
@@ -25,54 +30,6 @@ export interface MoveFilesParams {
|
||||
move: boolean
|
||||
}
|
||||
|
||||
export async function syncFiles (
|
||||
ctx: MeasureContext,
|
||||
workspaceId: WorkspaceId,
|
||||
exAdapter: StorageAdapterEx
|
||||
): Promise<void> {
|
||||
if (exAdapter.adapters === undefined) return
|
||||
|
||||
for (const [name, adapter] of [...exAdapter.adapters.entries()].reverse()) {
|
||||
await adapter.make(ctx, workspaceId)
|
||||
|
||||
await retryOnFailure(ctx, 5, async () => {
|
||||
let time = Date.now()
|
||||
let count = 0
|
||||
|
||||
const iterator = await adapter.listStream(ctx, workspaceId)
|
||||
try {
|
||||
while (true) {
|
||||
const dataBulk = await iterator.next()
|
||||
if (dataBulk.length === 0) break
|
||||
|
||||
for (const data of dataBulk) {
|
||||
const blob = await exAdapter.stat(ctx, workspaceId, data._id)
|
||||
if (blob !== undefined) {
|
||||
if (blob.provider !== name && name === exAdapter.defaultAdapter) {
|
||||
await exAdapter.syncBlobFromStorage(ctx, workspaceId, data._id, exAdapter.defaultAdapter)
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
await exAdapter.syncBlobFromStorage(ctx, workspaceId, data._id, name)
|
||||
|
||||
count += 1
|
||||
if (count % 100 === 0) {
|
||||
const duration = Date.now() - time
|
||||
time = Date.now()
|
||||
|
||||
console.log('...processed', count, Math.round(duration / 1000) + 's')
|
||||
}
|
||||
}
|
||||
}
|
||||
console.log('processed', count)
|
||||
} finally {
|
||||
await iterator.close()
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
export async function moveFiles (
|
||||
ctx: MeasureContext,
|
||||
workspaceId: WorkspaceId,
|
||||
@@ -81,15 +38,13 @@ export async function moveFiles (
|
||||
): Promise<void> {
|
||||
if (exAdapter.adapters === undefined) return
|
||||
|
||||
const target = exAdapter.adapters.get(exAdapter.defaultAdapter)
|
||||
const target = exAdapter.adapters[0].adapter
|
||||
if (target === undefined) return
|
||||
|
||||
// We assume that the adapter moves all new files to the default adapter
|
||||
await target.make(ctx, workspaceId)
|
||||
|
||||
for (const [name, adapter] of exAdapter.adapters.entries()) {
|
||||
if (name === exAdapter.defaultAdapter) continue
|
||||
|
||||
for (const { name, adapter } of exAdapter.adapters.slice(1).reverse()) {
|
||||
console.log('moving from', name, 'limit', 'concurrency', params.concurrency)
|
||||
|
||||
// we attempt retry the whole process in case of failure
|
||||
@@ -192,14 +147,9 @@ async function processAdapter (
|
||||
}
|
||||
|
||||
for (const data of dataBulk) {
|
||||
let targetBlob: Blob | ListBlobResult | undefined = targetBlobs.get(data._id)
|
||||
const targetBlob: Blob | ListBlobResult | undefined = targetBlobs.get(data._id)
|
||||
if (targetBlob !== undefined) {
|
||||
console.log('Target blob already exists', targetBlob._id)
|
||||
|
||||
const aggrBlob = await exAdapter.stat(ctx, workspaceId, data._id)
|
||||
if (aggrBlob === undefined || aggrBlob?.provider !== targetBlob.provider) {
|
||||
targetBlob = await exAdapter.syncBlobFromStorage(ctx, workspaceId, targetBlob._id, exAdapter.defaultAdapter)
|
||||
}
|
||||
// We could safely delete source blob
|
||||
toRemove.push(data._id)
|
||||
}
|
||||
@@ -211,15 +161,13 @@ async function processAdapter (
|
||||
console.error('blob not found', data._id)
|
||||
continue
|
||||
}
|
||||
targetBlob = await rateLimiter.exec(async () => {
|
||||
const info = await rateLimiter.exec(async () => {
|
||||
try {
|
||||
const result = await retryOnFailure(
|
||||
ctx,
|
||||
5,
|
||||
async () => {
|
||||
await processFile(ctx, source, target, workspaceId, sourceBlob)
|
||||
// We need to sync and update aggregator table for now.
|
||||
return await exAdapter.syncBlobFromStorage(ctx, workspaceId, sourceBlob._id, exAdapter.defaultAdapter)
|
||||
return await processFile(ctx, source, target, workspaceId, sourceBlob)
|
||||
},
|
||||
50
|
||||
)
|
||||
@@ -232,8 +180,8 @@ async function processAdapter (
|
||||
}
|
||||
})
|
||||
|
||||
if (targetBlob !== undefined) {
|
||||
// We could safely delete source blob
|
||||
// We could safely delete source blob
|
||||
if (info !== undefined) {
|
||||
toRemove.push(sourceBlob._id)
|
||||
}
|
||||
processedBytes += sourceBlob.size
|
||||
@@ -266,14 +214,14 @@ async function processFile (
|
||||
target: Pick<StorageAdapter, 'put'>,
|
||||
workspaceId: WorkspaceId,
|
||||
blob: Blob
|
||||
): Promise<void> {
|
||||
): Promise<UploadedObjectInfo> {
|
||||
const readable = await source.get(ctx, workspaceId, blob._id)
|
||||
try {
|
||||
readable.on('end', () => {
|
||||
readable.destroy()
|
||||
})
|
||||
const stream = readable.pipe(new PassThrough())
|
||||
await target.put(ctx, workspaceId, blob._id, stream, blob.contentType, blob.size)
|
||||
return await target.put(ctx, workspaceId, blob._id, stream, blob.contentType, blob.size)
|
||||
} finally {
|
||||
readable.destroy()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user