Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion src/http/error-handler.connection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ describe.each(['HTTP', 'S3', 'TUS'] as const)('%s error connections', (protocol)
},
onResponseError,
})
await app.register(authenticatedRoutes, { tusServer: tus })
await app.register(authenticatedRoutes, { tusServer: tus, signed: false })
} else {
await app.register(closeConnectionOnError)
if (protocol === 'S3') {
Expand Down
1 change: 1 addition & 0 deletions src/http/plugins/sync-hooks.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,7 @@ describe('sync request lifecycle hooks', () => {
register: (app: FastifyInstance) =>
app.register(publicRoutes, {
tusServer: { handle: vi.fn() } as unknown as Server,
signed: false,
}),
hooks: ['preHandler'],
},
Expand Down
6 changes: 6 additions & 0 deletions src/http/routes/tus/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,7 @@ describe('public tus route request context', () => {
rawRes.end()
}),
} as unknown as Server,
signed: false,
})
})

Expand All @@ -142,4 +143,9 @@ describe('public tus route request context', () => {
sbReqId: 'sb-req-123',
})
})

it('disposes the db when the response closes', async () => {
await app.inject({ method: 'OPTIONS', url: '/public/object' })
expect(observedUpload?.db.dispose).toHaveBeenCalled()
})
})
43 changes: 28 additions & 15 deletions src/http/routes/tus/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import {
onResponseError,
onUploadFinish,
SIGNED_URL_SUFFIX,
verifySignedUploadRequest,
} from './lifecycle'

const {
Expand Down Expand Up @@ -162,6 +163,7 @@ function createTusServer(
}

const resourceId = UploadId.fromString(uploadId)
await verifySignedUploadRequest(req, resourceId)

const bucket = await req.upload.storage
.asSuperUser()
Expand Down Expand Up @@ -210,6 +212,7 @@ export default async function routes(fastify: FastifyInstance) {

fastify.register(authenticatedRoutes, {
tusServer,
signed: false,
})
})

Expand All @@ -221,7 +224,7 @@ export default async function routes(fastify: FastifyInstance) {

fastify.register(authenticatedRoutes, {
tusServer,
operation: '_signed',
signed: true,
})
},
{ prefix: SIGNED_URL_SUFFIX }
Expand All @@ -231,6 +234,7 @@ export default async function routes(fastify: FastifyInstance) {
fastify.register(async (fastify) => {
fastify.register(publicRoutes, {
tusServer,
signed: false,
})
})

Expand All @@ -242,7 +246,7 @@ export default async function routes(fastify: FastifyInstance) {

fastify.register(publicRoutes, {
tusServer,
operation: '_signed',
signed: true,
})
},
{ prefix: SIGNED_URL_SUFFIX }
Expand All @@ -252,7 +256,8 @@ export default async function routes(fastify: FastifyInstance) {
function setTusRequestContext(
req: FastifyRequest,
reply: FastifyReply,
done: HookHandlerDoneFunction
done: HookHandlerDoneFunction,
isSigned: boolean
) {
// TUS protocol rejections write directly and skip Fastify's onSend hook.
const writeHead = reply.raw.writeHead
Expand All @@ -262,6 +267,7 @@ function setTusRequestContext(
}
return Reflect.apply(writeHead, reply.raw, args)
}
reply.raw.once('close', () => req.db?.dispose())

;(req.raw as MultiPartRequest).log = req.log
;(req.raw as MultiPartRequest).upload = {
Expand All @@ -270,14 +276,16 @@ function setTusRequestContext(
owner: req.owner,
db: req.db,
isUpsert: req.headers['x-upsert'] === 'true',
isSigned,
reqId: req.id,
sbReqId: req.sbReqId,
}
done()
}

export const authenticatedRoutes = fastifyPlugin(
async (fastify: FastifyInstance, options: { tusServer: Server; operation?: string }) => {
async (fastify: FastifyInstance, options: { tusServer: Server; signed: boolean }) => {
const operationSuffix = options.signed ? '_signed' : ''
fastify.register(async function authorizationContext(fastify) {
fastify.addContentTypeParser('application/offset+octet-stream', (request, payload, done) =>
done(null)
Expand All @@ -289,14 +297,16 @@ export const authenticatedRoutes = fastifyPlugin(
})
})

fastify.addHook('preHandler', setTusRequestContext)
fastify.addHook('preHandler', (req, res, done) =>
setTusRequestContext(req, res, done, options.signed)
)

fastify.post(
'/',
{
schema: { summary: 'Handle POST request for TUS Resumable uploads', tags: ['resumable'] },
config: {
operation: `${ROUTE_OPERATIONS.TUS_CREATE_UPLOAD}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_CREATE_UPLOAD}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -309,7 +319,7 @@ export const authenticatedRoutes = fastifyPlugin(
{
schema: { summary: 'Handle POST request for TUS Resumable uploads', tags: ['resumable'] },
config: {
operation: `${ROUTE_OPERATIONS.TUS_CREATE_UPLOAD}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_CREATE_UPLOAD}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -322,7 +332,7 @@ export const authenticatedRoutes = fastifyPlugin(
{
schema: { summary: 'Handle PUT request for TUS Resumable uploads', tags: ['resumable'] },
config: {
operation: `${ROUTE_OPERATIONS.TUS_UPLOAD_PART}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_UPLOAD_PART}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -337,7 +347,7 @@ export const authenticatedRoutes = fastifyPlugin(
tags: ['resumable'],
},
config: {
operation: `${ROUTE_OPERATIONS.TUS_UPLOAD_PART}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_UPLOAD_PART}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -349,7 +359,7 @@ export const authenticatedRoutes = fastifyPlugin(
{
schema: { summary: 'Handle HEAD request for TUS Resumable uploads', tags: ['resumable'] },
config: {
operation: `${ROUTE_OPERATIONS.TUS_GET_UPLOAD}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_GET_UPLOAD}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -364,7 +374,7 @@ export const authenticatedRoutes = fastifyPlugin(
tags: ['resumable'],
},
config: {
operation: `${ROUTE_OPERATIONS.TUS_DELETE_UPLOAD}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_DELETE_UPLOAD}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -376,13 +386,16 @@ export const authenticatedRoutes = fastifyPlugin(
)

export const publicRoutes = fastifyPlugin(
async (fastify: FastifyInstance, options: { tusServer: Server; operation?: string }) => {
async (fastify: FastifyInstance, options: { tusServer: Server; signed: boolean }) => {
const operationSuffix = options.signed ? '_signed' : ''
fastify.register(async (fastify) => {
fastify.addContentTypeParser('application/offset+octet-stream', (request, payload, done) =>
done(null)
)

fastify.addHook('preHandler', setTusRequestContext)
fastify.addHook('preHandler', (req, res, done) =>
setTusRequestContext(req, res, done, options.signed)
)

fastify.options(
'/',
Expand All @@ -393,7 +406,7 @@ export const publicRoutes = fastifyPlugin(
description: 'Handle OPTIONS request for TUS Resumable uploads',
},
config: {
operation: `${ROUTE_OPERATIONS.TUS_OPTIONS}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_OPTIONS}${operationSuffix}`,
},
},
async (req, res) => {
Expand All @@ -410,7 +423,7 @@ export const publicRoutes = fastifyPlugin(
description: 'Handle OPTIONS request for TUS Resumable uploads',
},
config: {
operation: `${ROUTE_OPERATIONS.TUS_OPTIONS}${options.operation || ''}`,
operation: `${ROUTE_OPERATIONS.TUS_OPTIONS}${operationSuffix}`,
},
},
async (req, res) => {
Expand Down
47 changes: 0 additions & 47 deletions src/http/routes/tus/lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,3 @@
import { EventEmitter } from 'node:events'
import type { ServerResponse } from 'node:http'
import { logSchema } from '@internal/monitoring'
import { Uploader } from '@storage/uploader'
import type { DataStore } from '@tus/server'
Expand All @@ -16,21 +14,16 @@ function createRawTusRequest({
method?: string
sbReqId?: string
} = {}) {
const response = new EventEmitter()
const reqLog = {
error: vi.fn(),
warn: vi.fn(),
}
const dispose = vi.fn()

const request = {
headers,
log: reqLog,
method,
upload: {
db: {
dispose,
},
isUpsert: false,
owner: 'owner-123',
storage: {
Expand All @@ -46,60 +39,20 @@ function createRawTusRequest({
} as unknown as MultiPartRequest

return {
dispose,
rawReq: {
method,
runtime: {
name: 'node',
node: {
req: request,
res: response as unknown as ServerResponse,
},
},
} as unknown as Parameters<typeof onIncomingRequest>[0],
reqLog,
response,
}
}

describe('tus lifecycle logging', () => {
it('disposes the db synchronously when the response finishes', async () => {
const { dispose, rawReq, response } = createRawTusRequest({
method: 'HEAD',
})

await onIncomingRequest(rawReq, uploadId, {} as DataStore)

response.emit('finish')

expect(dispose).toHaveBeenCalledOnce()
})

it('disposes the db when the response closes without finishing', async () => {
const { dispose, rawReq, response } = createRawTusRequest({
method: 'HEAD',
})

await onIncomingRequest(rawReq, uploadId, {} as DataStore)

response.emit('close')

expect(dispose).toHaveBeenCalledOnce()
})

it('disposes the db only once when a finished response subsequently closes', async () => {
const { dispose, rawReq, response } = createRawTusRequest({
method: 'HEAD',
})

await onIncomingRequest(rawReq, uploadId, {} as DataStore)

response.emit('finish')
response.emit('close')

expect(dispose).toHaveBeenCalledOnce()
})

it('logs upload metadata parse failures with sbReqId through logSchema', async () => {
const warningSpy = vi.spyOn(logSchema, 'warning').mockImplementation(() => undefined)
const { rawReq, reqLog } = createRawTusRequest({
Expand Down
Loading
Loading