Files
huly-platform/foundations/server/packages/datalake/src/client.ts
T
+1 98652c6476 Include sub projects (#10201)
* Add bump-changes

* Add utility tests

* Add utility tests

* Bump to new version of esbuild and typescript

* v0.7.3

* use platform rig 0.7.10

* upgrade: memory engine optimized; change name  to  (was recommended by Copilot and Onnikov, TODO: CHANGE CLIENT TOO!!!)

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* Fix rate limits bug

* Bump versions

* Fix lock file

* Fix bug in queue cleanup

* Add more tests for queue

* Add api-test tests

* Initial commit

* Improve hierarchy + tests

Add tests for hierarchy and few performance/memory  optimizations.

* Add more hierarchy tests

* Move from Huly platform repository

* Add docker tests setup

* Fix test to be executed only once

* Add connection tests

* Fix package include source files

* More tests

* Create README.md

* Fix pnpm lock

* Fix packages publish

* Remove broken tests

* feat: adjust hulylake client for storage adapter

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Fix export

* Fix publish

* Fix message update (#114)

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Bump version

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Bump versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Fix lang store (#115)

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* update hulylake client

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump version

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Add hulylake storage adapter

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump version

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix validation issues

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix: do not fail on deseralization error and add logs

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* bump version -> 0.1.14

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix collaboration test

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Update prettier and new update-deps script

Prettier + svelte support

* fix unstable ydoc tests

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Add tx ordering middleware

* Fix ordering tests

* Fix Kafka close of admin

* Add tests for measurement and understand overhead

* Fix not updated lock file

* Fix update-deps

* Fix update-deps

* Use latest platform-rig

* Fix deps

* Add rush check to CI

* Use latest versions

* Bump versions

* Fix lock file

* validate json patch

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* bump version -> 0.1.15

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix merge unit tests

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Script to sync eslint deps

* Fix deps

* Fix tests

* Fix platform-rig detection

* Update to latest platform-rig

* Update to latest platform rig and core

* Bump typescript

* Bump typescript

* Rollback eslint plugins

* Fix lock file

* Bump platform-rig

* Update to latest platform-rig

* update to latest platform-rig

* Allow to compile svelte files

* Add ui-test component for checking compile

* Fix log levels rename compile ui -> compile ui-esbuild

* Fix build

* Bump esbuild svelte version

* Chore: use fixed versions in update-deps

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Chore: commit changes

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Update deps

* Add tests for session manager

* Fix txOrdering implementation

* Bump ordering

* Prevent metrics zero values in measure

+ Fix format svelte files

* Revert update-deps script logic

* v0.7.19

* update to latest platform-rig

* Update deps

* Fix pnpm

* Session counters

* Fix pnpm lock

* Add storage client

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump core

* Fix pnpm

* Get rid of communication dependency

* Add copilot memory file

* Use proper name for instructions file

* Fix instructions

* Use domain instead of test name in gauges

* Update instructions file

* fix front service upload

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* remove incorrect test

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Move packages to huly.core

* Move packages to core, since they are not utils

* Add global user profile

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Fix lock file

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Add support for memory limit check

* Bump version

* Fix pnpm

* report more accurate upload progress

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fic validation issues

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Fix deps

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Move LowLevelStorage to server

* Fix linting

* Revert "Fix linting"

This reverts commit 54631d353e.

* Revert "Move LowLevelStorage to server"

This reverts commit aafb8f6f12.

* feature: add regorus engine with permit file

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* Fix one second counters for memory usage

* Fix kafka test

* use fresh core

* Version bump

* fix: key parameter added

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* Fix readme and few author mistakes

* Export domain schemas

* Bump version

* Tests (#117)

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Bump version

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* feat: compact compact worker (#4)

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* bump version -> 0.1.16

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Add TypeIdentifier

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Add change logs

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* rename send -> try_send

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix pnpm lock

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Add identifier middleware, bump core

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Add subsciption methods to account client

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Fix lock file

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Fix reaction notification (#118)

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Bump version

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Improve find methods schemas to convert to valid types

Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com>

* Add change description

Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com>

* Do not transcode while recording

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Open telemetry support

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* use proper content type in multipart upload

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* fix build (#26)

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Fix peers (#120)

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Bump version

Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>

* Add ActivityCollaborativeChange

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Update pnpm

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Allow to suspend errors on with

* Fix pnpm cache

* update versions

* v0.7.17 for all

* v0.7.11

* v0.7.14

* remove arc from worker

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* bump version -> 0.1.17

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Rank for attributes

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Update pnpm

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix one second counters

* Fix withContext and allow pass options

* Fix formatting

* Use updated deps

* Bump versions

* Update deps

* Update deps to platform.core

* add support for textColor mark

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* add support for textStyle mark

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump versions again

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Fix

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Rework on second timers

* fix merge of large blobs feched from s3

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* bump version -> 0.1.18

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* New subscription methods in account-client

* Update lock file

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Send error on find for wrong domain

* Suspend connect custom errors events in traces

* Bump client

* Bump core

* update deps

* Fix lock file

* Sorting for TypeIdentifier

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Bump version

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* add workspace usage info

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Bump versions

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>

* Improve pg security perfomance

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix identifier middleware

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Update TxAccessLevel interface

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Allow guest to update its identities

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Add password login locked platform status

Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>

* Fix Uptrace normalizeMarkdown errors

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Add change log

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Add txMatch to permission

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* update pnpm lock

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Bump

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix permission middleware

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix enum sorting

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Enable formatting check

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>

* Enable formatting check

* Add change

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>

* Fix Uptrace NaN error

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* feature: removed actors, improved performance

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* feature: ping from server to clients added

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* feature: ping from server to clients added

Signed-off-by: Leonid Kaganov <lleo@lleo.me>

* Compress kafka messages and fix exception in findAll

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Bump versions

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Bump versions

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Rush change

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Fix compression param

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Trigger change

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Clean up

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Trigger change

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Bump markdown version

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Enable sub projects

* Fix wrong double symbol scripts

* Include foundation packages

* Add support for custom exclude filters

Add support for custom exclude filters - by Andrey Sobolev - haiodo@gmail.com

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>

* Bump

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>

* Fix Uptrace filter is not a function error

Signed-off-by: Artem Savchenko <armisav@gmail.com>

* Sync versions

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>

---------

Signed-off-by: Leonid Kaganov <lleo@lleo.me>
Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>
Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com>
Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com>
Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>
Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com>
Signed-off-by: Artem Savchenko <armisav@gmail.com>
Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
Co-authored-by: Leonid Kaganov <lleo@lleo.me>
Co-authored-by: Alexander Onnikov <Alexander.Onnikov@xored.com>
Co-authored-by: Alexander Onnikov <Alexander.Onnikov@gmail.com>
Co-authored-by: Kristina <kristin.fefelova@gmail.com>
Co-authored-by: Alexey Zinoviev <alexey.zinoviev@xored.com>
Co-authored-by: Denis Bykhov <bykhov.denis@gmail.com>
Co-authored-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com>
Co-authored-by: Alexander Onnikov <aonnikov@hardcoreeng.com>
Co-authored-by: Artem Savchenko <armisav@gmail.com>
2025-11-26 19:15:30 +05:00

546 lines
15 KiB
TypeScript

//
// Copyright © 2024 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
//
import { type MeasureContext, type WorkspaceUuid, concatLink } from '@hcengineering/core'
import { Readable } from 'stream'
import { DatalakeError, NetworkError, NotFoundError } from './error'
import { unwrapETag } from './utils'
/** @public */
export interface ObjectMetadata {
lastModified: number
name: string
contentType: string
etag: string
size?: number
}
/** @public */
export interface ListObjectOutput {
cursor: string | undefined
blobs: Omit<ObjectMetadata, 'lastModified'>[]
}
/** @public */
export interface StatObjectOutput {
lastModified: number
type: string
etag?: string
size?: number
}
/** @public */
export interface UploadObjectParams {
lastModified: number
type: string
size?: number
}
interface BlobUploadError {
key: string
error: string
}
interface BlobUploadSuccess {
key: string
id: string
metadata: ObjectMetadata
}
type BlobUploadResult = BlobUploadSuccess | BlobUploadError
interface MultipartUpload {
uploadId: string
}
interface MultipartUploadPart {
partNumber: number
etag: string
}
export interface R2UploadParams {
location: string
bucket: string
}
/** @public */
export interface WorkspaceStats {
count: number
size: number
}
/** @public */
export class DatalakeClient {
private readonly headers: Record<string, string>
constructor (
private readonly endpoint: string,
private readonly token: string
) {
this.headers = { Authorization: 'Bearer ' + token }
}
getObjectUrl (ctx: MeasureContext, workspace: WorkspaceUuid, objectName: string): string {
const path = `/blob/${workspace}/${encodeURIComponent(objectName)}`
return concatLink(this.endpoint, path)
}
async getWorkspaceStats (ctx: MeasureContext, workspace: WorkspaceUuid): Promise<WorkspaceStats> {
const path = `/stats/${workspace}`
const url = new URL(concatLink(this.endpoint, path))
const response = await fetchSafe(ctx, url, { headers: { ...this.headers } })
return (await response.json()) as WorkspaceStats
}
async listObjects (
ctx: MeasureContext,
workspace: WorkspaceUuid,
cursor: string | undefined,
limit: number = 100
): Promise<ListObjectOutput> {
const path = `/blob/${workspace}`
const url = new URL(concatLink(this.endpoint, path))
url.searchParams.append('limit', String(limit))
if (cursor !== undefined) {
url.searchParams.append('cursor', cursor)
}
const response = await fetchSafe(ctx, url, { headers: { ...this.headers } })
return (await response.json()) as ListObjectOutput
}
async getObject (ctx: MeasureContext, workspace: WorkspaceUuid, objectName: string): Promise<Readable | undefined> {
const url = this.getObjectUrl(ctx, workspace, objectName)
let response
try {
response = await fetchSafe(ctx, url, { headers: { ...this.headers } })
} catch (err: any) {
if (err.name === 'NotFoundError') {
return undefined
}
throw err
}
if (response.body == null) {
ctx.error('bad datalake response', { objectName })
throw new DatalakeError('Missing response body')
}
return Readable.fromWeb(response.body as any)
}
async getPartialObject (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
offset: number,
length?: number
): Promise<Readable | undefined> {
const url = this.getObjectUrl(ctx, workspace, objectName)
const headers = {
...this.headers,
Range: length !== undefined ? `bytes=${offset}-${offset + length - 1}` : `bytes=${offset}`
}
let response
try {
response = await fetchSafe(ctx, url, { headers })
} catch (err: any) {
if (err.name === 'NotFoundError') {
return undefined
}
throw err
}
if (response.body == null) {
ctx.error('bad datalake response', { objectName })
throw new DatalakeError('Missing response body')
}
return Readable.fromWeb(response.body as any)
}
async statObject (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string
): Promise<StatObjectOutput | undefined> {
const url = this.getObjectUrl(ctx, workspace, objectName)
let response: Response
try {
response = await fetchSafe(ctx, url, {
method: 'HEAD',
headers: { ...this.headers }
})
} catch (err: any) {
if (err.name === 'NotFoundError') {
return
}
throw err
}
const headers = response.headers
const lastModified = Date.parse(headers.get('Last-Modified') ?? '')
const size = parseInt(headers.get('Content-Length') ?? '0', 10)
return {
lastModified: isNaN(lastModified) ? 0 : lastModified,
size: isNaN(size) ? 0 : size,
type: headers.get('Content-Type') ?? '',
etag: unwrapETag(headers.get('ETag') ?? '')
}
}
async deleteObject (ctx: MeasureContext, workspace: WorkspaceUuid, objectName: string): Promise<void> {
const url = this.getObjectUrl(ctx, workspace, objectName)
try {
await fetchSafe(ctx, url, {
method: 'DELETE',
headers: { ...this.headers }
})
} catch (err: any) {
if (err.name !== 'NotFoundError') {
throw err
}
}
}
async putObject (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
stream: Readable | Buffer | string,
params: UploadObjectParams
): Promise<ObjectMetadata> {
let size = params.size
if (size === undefined) {
if (Buffer.isBuffer(stream)) {
size = stream.length
} else if (typeof stream === 'string') {
size = Buffer.byteLength(stream)
} else {
// TODO: Implement size calculation for Readable streams
ctx.warn('unknown object size', { workspace, objectName })
}
}
if (size === undefined || size < 64 * 1024 * 1024) {
return await ctx.with(
'direct-upload',
{},
(ctx) => this.uploadWithFormData(ctx, workspace, objectName, stream, { ...params, size }),
{ workspace, objectName }
)
} else {
return await ctx.with(
'multipart-upload',
{},
(ctx) => this.uploadWithMultipart(ctx, workspace, objectName, stream, { ...params, size }),
{ workspace, objectName }
)
}
}
async uploadWithFormData (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
stream: Readable | Buffer | string,
params: UploadObjectParams
): Promise<ObjectMetadata> {
const path = `/upload/form-data/${workspace}`
const url = concatLink(this.endpoint, path)
const buffer = await toBuffer(stream)
const file = new File([new Uint8Array(buffer)], objectName, {
type: params.type,
lastModified: params.lastModified
})
const form = new FormData()
form.append('file', file)
const response = await fetchSafe(ctx, url, {
method: 'POST',
body: form as unknown as BodyInit,
headers: {
...this.headers
}
})
const result = (await response.json()) as BlobUploadResult[]
if (result.length !== 1) {
throw new DatalakeError('Bad datalake response: ' + result.toString())
}
const uploadResult = result[0]
if ('error' in uploadResult) {
throw new DatalakeError('Upload failed: ' + uploadResult.error)
}
return uploadResult.metadata
}
async uploadWithMultipart (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
stream: Readable | Buffer | string,
params: UploadObjectParams
): Promise<ObjectMetadata> {
const chunkSize = 10 * 1024 * 1024
const multipart = await this.multipartUploadStart(ctx, workspace, objectName, params)
try {
const parts: MultipartUploadPart[] = []
let partNumber = 1
for await (const chunk of getChunks(stream, chunkSize)) {
const part = await this.multipartUploadPart(ctx, workspace, objectName, multipart, partNumber, chunk)
parts.push(part)
partNumber++
}
return await this.multipartUploadComplete(ctx, workspace, objectName, multipart, parts)
} catch (err: any) {
await this.multipartUploadAbort(ctx, workspace, objectName, multipart)
throw err
}
}
// S3
async uploadFromS3 (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
params: {
url: string
accessKeyId: string
secretAccessKey: string
}
): Promise<void> {
const path = `/upload/s3/${workspace}/${encodeURIComponent(objectName)}`
const url = concatLink(this.endpoint, path)
await fetchSafe(ctx, url, {
method: 'POST',
headers: {
...this.headers,
'Content-Type': 'application/json'
},
body: JSON.stringify(params)
})
}
async getS3UploadParams (ctx: MeasureContext, workspace: WorkspaceUuid): Promise<R2UploadParams> {
const path = `/upload/s3/${workspace}`
const url = concatLink(this.endpoint, path)
const response = await fetchSafe(ctx, url, { headers: { ...this.headers } })
const json = (await response.json()) as R2UploadParams
return json
}
async createFromS3 (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
params: {
filename: string
}
): Promise<void> {
const path = `/upload/s3/${workspace}/${encodeURIComponent(objectName)}`
const url = concatLink(this.endpoint, path)
await fetchSafe(ctx, url, {
method: 'POST',
headers: {
...this.headers,
'Content-Type': 'application/json'
},
body: JSON.stringify(params)
})
}
// Multipart
private async multipartUploadStart (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
params: UploadObjectParams
): Promise<MultipartUpload> {
const path = `/upload/multipart/${workspace}/${encodeURIComponent(objectName)}`
const url = concatLink(this.endpoint, path)
try {
const headers = {
...this.headers,
'Content-Type': params.type,
'Content-Length': params.size?.toString() ?? '0',
'Last-Modified': new Date(params.lastModified).toUTCString()
}
const response = await fetchSafe(ctx, url, { method: 'POST', headers })
return (await response.json()) as MultipartUpload
} catch (err: any) {
ctx.error('failed to start multipart upload', { workspace, objectName, err })
throw new DatalakeError('Failed to start multipart upload')
}
}
private async multipartUploadPart (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
multipart: MultipartUpload,
partNumber: number,
data: Readable | Buffer | string
): Promise<MultipartUploadPart> {
const path = `/upload/multipart/${workspace}/${encodeURIComponent(objectName)}/part`
const url = new URL(concatLink(this.endpoint, path))
url.searchParams.set('uploadId', multipart.uploadId)
url.searchParams.set('partNumber', partNumber.toString())
const body = data instanceof Readable ? (Readable.toWeb(data) as ReadableStream) : data
try {
const response = await fetchSafe(ctx, url, {
method: 'PUT',
body: body as BodyInit,
headers: { ...this.headers }
})
return (await response.json()) as MultipartUploadPart
} catch (err: any) {
ctx.error('failed to upload multipart part', { workspace, objectName, err })
throw new DatalakeError('Failed to upload multipart part')
}
}
private async multipartUploadComplete (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
multipart: MultipartUpload,
parts: MultipartUploadPart[]
): Promise<ObjectMetadata> {
const path = `/upload/multipart/${workspace}/${encodeURIComponent(objectName)}/complete`
const url = new URL(concatLink(this.endpoint, path))
url.searchParams.set('uploadId', multipart.uploadId)
try {
const res = await fetchSafe(ctx, url, {
method: 'POST',
body: JSON.stringify({ parts }),
headers: {
'Content-Type': 'application/json',
...this.headers
}
})
return (await res.json()) as ObjectMetadata
} catch (err: any) {
ctx.error('failed to complete multipart upload', { workspace, objectName, err })
throw new DatalakeError('Failed to complete multipart upload')
}
}
private async multipartUploadAbort (
ctx: MeasureContext,
workspace: WorkspaceUuid,
objectName: string,
multipart: MultipartUpload
): Promise<void> {
const path = `/upload/multipart/${workspace}/${encodeURIComponent(objectName)}/abort`
const url = new URL(concatLink(this.endpoint, path))
url.searchParams.set('uploadId', multipart.uploadId)
try {
await fetchSafe(ctx, url, { method: 'POST', headers: { ...this.headers } })
} catch (err: any) {
ctx.error('failed to abort multipart upload', { workspace, objectName, err })
throw new DatalakeError('Failed to abort multipart upload')
}
}
}
async function toBuffer (data: Buffer | string | Readable): Promise<Buffer> {
if (Buffer.isBuffer(data)) {
return data
} else if (typeof data === 'string') {
return Buffer.from(data)
} else if (data instanceof Readable) {
const chunks: Buffer[] = []
for await (const chunk of data) {
chunks.push(chunk)
}
return Buffer.concat(chunks as any)
} else {
throw new TypeError('Unsupported data type')
}
}
async function * getChunks (data: Buffer | string | Readable, chunkSize: number): AsyncGenerator<Buffer> {
if (Buffer.isBuffer(data)) {
let offset = 0
while (offset < data.length) {
yield data.subarray(offset, offset + chunkSize)
offset += chunkSize
}
} else if (typeof data === 'string') {
const buffer = Buffer.from(data)
yield * getChunks(buffer, chunkSize)
} else if (data instanceof Readable) {
let buffer = Buffer.alloc(0)
for await (const chunk of data) {
buffer = Buffer.concat([buffer, chunk])
while (buffer.length >= chunkSize) {
yield buffer.subarray(0, chunkSize)
buffer = buffer.subarray(chunkSize)
}
}
if (buffer.length > 0) {
yield buffer
}
}
}
async function fetchSafe (ctx: MeasureContext, url: string | URL, init?: RequestInit): Promise<Response> {
let response
try {
response = await ctx.with('fetch', {}, () => fetch(url, init), { url: url.toString() }, { span: 'disable' })
} catch (err: any) {
ctx.error('network error', { err })
throw new NetworkError(`Network error ${err}`)
}
if (!response.ok) {
const text = await response.text()
if (response.status === 404) {
throw new NotFoundError(text)
} else {
throw new DatalakeError(text)
}
}
return response
}