Skip to content
Open
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
7 changes: 7 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,7 @@ type StorageConfigType = {
tusUseFileVersionSeparator: boolean
tusAllowS3Tags: boolean
tusLockType: 'postgres' | 's3'
tusBodyIdleTimeoutMs: number
s3ProtocolEnabled: boolean
s3ProtocolPrefix: string
s3ProtocolAllowForwardedHeader: boolean
Expand Down Expand Up @@ -480,6 +481,12 @@ export function getConfig(options?: { reload?: boolean }): StorageConfigType {
getOptionalConfigFromEnv('TUS_USE_FILE_VERSION_SEPARATOR') === 'true',
tusAllowS3Tags: getOptionalConfigFromEnv('TUS_ALLOW_S3_TAGS') !== 'false',
tusLockType: getOptionalConfigFromEnv('TUS_LOCK_TYPE') || 'postgres',
tusBodyIdleTimeoutMs: envIntegerInRange(
getOptionalConfigFromEnv('TUS_BODY_IDLE_TIMEOUT_MS'),
1000 * 60,
0,
MAX_TIMER_DELAY_MS
),

// S3 Protocol
s3ProtocolEnabled: getOptionalConfigFromEnv('S3_PROTOCOL_ENABLED') !== 'false',
Expand Down
182 changes: 180 additions & 2 deletions src/http/routes/tus/index.test.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,14 @@
import { EventEmitter } from 'node:events'
import * as http from 'node:http'
import * as https from 'node:https'
import type { AddressInfo } from 'node:net'
import { NodeHttpHandler } from '@smithy/node-http-handler'
import { HttpRequest } from '@smithy/protocol-http'
import type { Server } from '@tus/server'
import Fastify, { FastifyInstance } from 'fastify'
import Fastify, { FastifyInstance, FastifyReply, FastifyRequest } from 'fastify'
import { getConfig } from '../../../config'
import { requestContext } from '../../plugins/request-context'
import { createTusLockS3Client, publicRoutes } from './index'
import { createTusLockS3Client, handleTusRequestWithIdleTimeout, publicRoutes } from './index'
import type { MultiPartRequest } from './lifecycle'

describe('TUS S3 clients', () => {
Expand Down Expand Up @@ -149,3 +150,180 @@ describe('public tus route request context', () => {
expect(observedUpload?.db.dispose).toHaveBeenCalled()
})
})

class FakeSocket {
bytesRead = 0
}

class FakeIncomingMessage extends EventEmitter {
socket: FakeSocket | undefined = new FakeSocket()
complete = false
readableEnded = false
readableLength = 0
destroyed = false
headers: Record<string, string> = {}
executionError?: Error
destroy = vi.fn((_err?: Error) => {
this.destroyed = true
})
}

describe('handleTusRequestWithIdleTimeout', () => {
const { tusBodyIdleTimeoutMs } = getConfig()

beforeEach(() => {
vi.useFakeTimers()
})

afterEach(() => {
vi.useRealTimers()
})

function createReqRes(headers: Record<string, string> = {}) {
const raw = new FakeIncomingMessage()
raw.headers = headers
const req = { raw } as unknown as FastifyRequest
const res = { raw: {} } as unknown as FastifyReply
return { req, res, raw }
}

// A `handle` that never resolves on its own, so the idle-check timers
// get a chance to fire before the wrapper's own `finally` disarms them.
function pendingHandle() {
return vi.fn(() => new Promise<void>(() => {}))
}

test('skips arming the idle timer when the request declares no body', async () => {
const setTimeoutSpy = vi.spyOn(global, 'setTimeout')
const handle = vi.fn().mockResolvedValue(undefined)
const { req, res, raw } = createReqRes()
const tusServer = { handle } as unknown as Server

await handleTusRequestWithIdleTimeout(tusServer, req, res)

expect(handle).toHaveBeenCalledWith(raw, res.raw)
expect(setTimeoutSpy).not.toHaveBeenCalled()
})

test('skips arming the idle timer when the body is already fully received', async () => {
const setTimeoutSpy = vi.spyOn(global, 'setTimeout')
const handle = vi.fn().mockResolvedValue(undefined)
const { req, res, raw } = createReqRes({ 'content-length': '10' })
raw.complete = true
const tusServer = { handle } as unknown as Server

await handleTusRequestWithIdleTimeout(tusServer, req, res)

expect(handle).toHaveBeenCalledWith(raw, res.raw)
expect(setTimeoutSpy).not.toHaveBeenCalled()
})

test('destroys the connection when no bytes arrive before the idle check fires', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)

expect(raw.destroy).toHaveBeenCalledTimes(1)
expect(raw.destroy.mock.calls[0][0]).toMatchObject({
message: 'TUS request body idle timeout - no bytes received',
})
expect(raw.executionError).toBe(raw.destroy.mock.calls[0][0])
})

test('reschedules instead of timing out while bytesRead keeps increasing, then times out once it stalls', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)

// Bytes arrive just before the first check - should reschedule, not time out.
raw.socket!.bytesRead = 5
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)
expect(raw.destroy).not.toHaveBeenCalled()

// No further bytes arrive during the second window - now it should time out.
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)
expect(raw.destroy).toHaveBeenCalledTimes(1)
})

test('reschedules while the destination write is draining a backpressure backlog, then times out once nothing moves at all', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)

// Client burst fills the buffer faster than the destination write can
// drain it: bytesRead and readableLength climb together.
raw.socket!.bytesRead = 100_000
raw.readableLength = 65_536
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)
expect(raw.destroy).not.toHaveBeenCalled()

// bytesRead stops growing, but the destination write is still actively draining the backlog
raw.readableLength = 32_768
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)
expect(raw.destroy).not.toHaveBeenCalled()

// No further progress is genuine idleness and must time out, even with a backlog
// sitting in the buffer, or a permanently stuck write would hold the lock forever.
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)
expect(raw.destroy).toHaveBeenCalledTimes(1)
})

test('skips arming the idle timer when the request is already destroyed', async () => {
const setTimeoutSpy = vi.spyOn(global, 'setTimeout')
const handle = vi.fn().mockResolvedValue(undefined)
const { req, res, raw } = createReqRes({ 'content-length': '10' })
raw.destroyed = true
const tusServer = { handle } as unknown as Server

await handleTusRequestWithIdleTimeout(tusServer, req, res)

expect(handle).toHaveBeenCalledWith(raw, res.raw)
expect(setTimeoutSpy).not.toHaveBeenCalled()
})

test('disarms without re-destroying if the connection is destroyed by something else before the idle check runs', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)
raw.destroyed = true
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)

expect(raw.destroy).not.toHaveBeenCalled()
})

test('disarms the idle timer on an abrupt disconnect, which closes without ever emitting end', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)
raw.emit('close')

await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs * 2)
expect(raw.destroy).not.toHaveBeenCalled()
})

test('disarms without destroying the connection if the idle check runs as the body finishes arriving', async () => {
const handle = pendingHandle()
const { req, res, raw } = createReqRes({ 'content-length': '10' })
const tusServer = { handle } as unknown as Server

void handleTusRequestWithIdleTimeout(tusServer, req, res)
// The body finishes arriving in the same window the idle check runs,
// racing the 'end' listener that would otherwise disarm it first.
raw.complete = true
await vi.advanceTimersByTimeAsync(tusBodyIdleTimeoutMs)

expect(raw.destroy).not.toHaveBeenCalled()
expect(raw.executionError).toBeUndefined()
})
})
72 changes: 68 additions & 4 deletions src/http/routes/tus/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ const {
tusMaxConcurrentUploads,
tusAllowS3Tags,
tusLockType,
tusBodyIdleTimeoutMs,
uploadFileSizeLimit,
storageBackendType,
storageFilePath,
Expand Down Expand Up @@ -283,6 +284,69 @@ function setTusRequestContext(
done()
}

export async function handleTusRequestWithIdleTimeout(
tusServer: Server,
req: FastifyRequest,
res: FastifyReply
) {
const socket = req.raw.socket
const isDone = () => req.raw.complete || req.raw.readableEnded || req.raw.destroyed
const hasDeclaredBody =
req.raw.headers['transfer-encoding'] !== undefined ||
Number(req.raw.headers['content-length']) > 0
if (!socket || tusBodyIdleTimeoutMs <= 0 || !hasDeclaredBody || isDone()) {
Comment thread
itslenny marked this conversation as resolved.
return tusServer.handle(req.raw, res.raw)
}

// We use our own timer instead of socket or stream events, so we don't interfere with how the body gets read
// bytesRead alone is not enough, because it also stalls when our own write to disk or S3 is slow
// bytesConsumed tracks the write side, so we can tell those two cases apart
// If either one is moving, the upload is alive
let lastBytesRead = socket.bytesRead
let lastBytesConsumed = lastBytesRead - req.raw.readableLength
let idleTimer: NodeJS.Timeout

const disarm = () => {
clearTimeout(idleTimer)
req.raw.removeListener('end', disarm)
req.raw.removeListener('close', disarm)
}

const checkIdle = () => {
if (isDone()) {
disarm()
return
}

const currentBytesRead = socket.bytesRead
const currentBytesConsumed = currentBytesRead - req.raw.readableLength
if (currentBytesRead > lastBytesRead || currentBytesConsumed > lastBytesConsumed) {
lastBytesRead = currentBytesRead
lastBytesConsumed = currentBytesConsumed
idleTimer = setTimeout(checkIdle, tusBodyIdleTimeoutMs)
return
}

const err = ERRORS.TusError('TUS request body idle timeout - no bytes received', 408)
req.raw.executionError = err
req.raw.destroy(err)
}

idleTimer = setTimeout(checkIdle, tusBodyIdleTimeoutMs)
// Stop tracking idle time once the client has sent the full body, so
// slow lock acquisition or upload finalization afterward can't trip a
// "no bytes received" timeout
req.raw.once('end', disarm)
Comment thread
itslenny marked this conversation as resolved.
// A disconnect never fires 'end', so also disarm on 'close'
req.raw.once('close', disarm)

try {
return await tusServer.handle(req.raw, res.raw)
} finally {
disarm()
}
}

export const authenticatedRoutes = fastifyPlugin(
async (fastify: FastifyInstance, options: { tusServer: Server; signed: boolean }) => {
const operationSuffix = options.signed ? '_signed' : ''
Expand Down Expand Up @@ -310,7 +374,7 @@ export const authenticatedRoutes = fastifyPlugin(
},
},
async (req, res) => {
await options.tusServer.handle(req.raw, res.raw)
await handleTusRequestWithIdleTimeout(options.tusServer, req, res)
}
)

Expand All @@ -323,7 +387,7 @@ export const authenticatedRoutes = fastifyPlugin(
},
},
async (req, res) => {
await options.tusServer.handle(req.raw, res.raw)
await handleTusRequestWithIdleTimeout(options.tusServer, req, res)
}
)

Expand All @@ -336,7 +400,7 @@ export const authenticatedRoutes = fastifyPlugin(
},
},
async (req, res) => {
await options.tusServer.handle(req.raw, res.raw)
await handleTusRequestWithIdleTimeout(options.tusServer, req, res)
}
)
fastify.patch(
Expand All @@ -351,7 +415,7 @@ export const authenticatedRoutes = fastifyPlugin(
},
},
async (req, res) => {
await options.tusServer.handle(req.raw, res.raw)
await handleTusRequestWithIdleTimeout(options.tusServer, req, res)
}
)
fastify.head(
Expand Down
Loading
Loading