diff --git a/services/datalake/pod-datalake/src/datalake/datalake.ts b/services/datalake/pod-datalake/src/datalake/datalake.ts index b74f1ee6c3..fc0955365d 100644 --- a/services/datalake/pod-datalake/src/datalake/datalake.ts +++ b/services/datalake/pod-datalake/src/datalake/datalake.ts @@ -19,7 +19,7 @@ import { Readable } from 'stream' import { type BlobDB } from './db' import { digestToUUID, stringToUUID } from './encodings' import { type BlobHead, type BlobBody, type BlobList, type BlobStorage, type Datalake, type Location } from './types' - +import { requestHLS } from '../handlers/video' import { type S3Bucket } from '../s3' export class DatalakeImpl implements Datalake { @@ -144,6 +144,9 @@ export class DatalakeImpl implements Datalake { } await bucket.put(ctx, filename, body, putOptions) await this.db.createBlobData(ctx, { workspace, name, hash, location, filename, size, type: contentType }) + if (contentType.startsWith('video/')) { + void requestHLS(ctx, workspace, name) + } return { name, size, contentType, lastModified, etag: hash } } } @@ -168,6 +171,10 @@ export class DatalakeImpl implements Datalake { await this.db.createBlobData(ctx, { workspace, name, hash, location, filename, size, type: contentType }) } + if (contentType.startsWith('video/')) { + void requestHLS(ctx, workspace, name) + } + return { name, size, contentType, lastModified, etag: hash } } diff --git a/services/datalake/pod-datalake/src/handlers/blob.ts b/services/datalake/pod-datalake/src/handlers/blob.ts index 15b7f75323..b1c4842116 100644 --- a/services/datalake/pod-datalake/src/handlers/blob.ts +++ b/services/datalake/pod-datalake/src/handlers/blob.ts @@ -18,7 +18,6 @@ import { type Request, type Response } from 'express' import { UploadedFile } from 'express-fileupload' import fs from 'fs' -import { requestHLS } from './video' import { cacheControl } from '../const' import { type Datalake } from '../datalake' import { getBufferSha256, getStreamSha256 } from '../hash' @@ -230,11 +229,6 @@ export async function handleUploadFormData ( ctx.info('uploaded', { workspace, name, etag: metadata.etag, type: contentType }) - if (contentType.startsWith('video/')) { - ctx.info('transcode', { workspace, name }) - void requestHLS(ctx, workspace, name) - } - return { key, metadata } } catch (err: any) { const error = err instanceof Error ? err.message : String(err) diff --git a/services/datalake/pod-datalake/src/handlers/multipart.ts b/services/datalake/pod-datalake/src/handlers/multipart.ts index 876a1a7bbf..5bb64c3f21 100644 --- a/services/datalake/pod-datalake/src/handlers/multipart.ts +++ b/services/datalake/pod-datalake/src/handlers/multipart.ts @@ -17,7 +17,6 @@ import { MeasureContext } from '@hcengineering/core' import { type Request, type Response } from 'express' import { cacheControl } from '../const' import { Datalake } from '../datalake' -import { requestHLS } from './video' export interface MultipartUpload { key: string @@ -119,12 +118,6 @@ export async function handleMultipartUploadComplete ( await bucket.completeMultipartUpload(ctx, uuid, { uploadId }, parts) const metadata = await datalake.create(ctx, workspace, name, uuid) - const contentType = req.headers['content-type'] ?? 'application/octet-stream' - if (contentType.startsWith('video/')) { - ctx.info('transcode', { workspace, name }) - void requestHLS(ctx, workspace, name) - } - ctx.info('multipart-complete', { workspace, name, uuid, uploadId }) res.status(200).json(metadata) diff --git a/services/datalake/pod-datalake/src/handlers/video.ts b/services/datalake/pod-datalake/src/handlers/video.ts index 5b9d35c4c8..c382d68971 100644 --- a/services/datalake/pod-datalake/src/handlers/video.ts +++ b/services/datalake/pod-datalake/src/handlers/video.ts @@ -26,11 +26,20 @@ interface StreamRequest { } export async function requestHLS (ctx: MeasureContext, workspace: string, name: string): Promise { + try { + ctx.info('request for hls', { workspace, name }) + await postTranscodingTask(ctx, workspace, name) + } catch (err) { + ctx.error('can not schedule a task', { err }) + } +} + +async function postTranscodingTask (ctx: MeasureContext, workspace: string, name: string): Promise { if (config.StreamUrl === undefined) { return } const streamReq: StreamRequest = { format: 'hls', source: name, workspace } - const token = generateToken(systemAccountUuid) + const token = generateToken(systemAccountUuid, undefined, { iss: 'datalake', aud: 'stream' }) const request = new Request(config.StreamUrl, { method: 'POST',