fix: wasn't preventing 2 uploads of same file at once
Massimo Melina committed
Jul 8, 2024 at 12:37 UTC
0436923e1d2d38bddef2875eecd1a90f56f7829e
2 files changed
+138
-112
src/upload.ts
+127
-110
@@ -37,8 +37,10 @@ function setUploadMeta(path: string, ctx: Koa.Context) {
37
}
38
39
// stay sync because we use this function with formidable()
40
-const cache: any = {}
40
+const diskSpaceCache: any = {}
41
+const openFiles = new Set()
42
export function uploadWriter(base: VfsNode, path: string, ctx: Koa.Context) {
43
+ let fullPath = ''
44
if (dirTraversal(path))
45
return fail(HTTP_FOOL)
46
if (statusCodeForMissingPerm(base, 'can_upload', ctx)) {
@@ -50,7 +52,7 @@ export function uploadWriter(base: VfsNode, path: string, ctx: Koa.Context) {
52
return fail()
53
}
54
// enforce minAvailableMb
53
- const fullPath = join(base.source!, path)
55
+ fullPath = join(base.source!, path)
56
const dir = dirname(fullPath)
57
const min = minAvailableMb.get() * (1 << 20)
58
const reqSize = Number(ctx.headers["content-length"])
@@ -64,12 +66,12 @@ export function uploadWriter(base: VfsNode, path: string, ctx: Koa.Context) {
66
let closestVfsNode: typeof base | undefined = base
67
while (closestVfsNode && !closestVfsNode.original) closestVfsNode = closestVfsNode.parent
68
const statDir = closestVfsNode!.source!
67
- if (!Object.hasOwn(cache, statDir)) {
68
- const c = cache[statDir] = getDiskSpaceSync(statDir)
69
+ if (!Object.hasOwn(diskSpaceCache, statDir)) {
70
+ const c = diskSpaceCache[statDir] = getDiskSpaceSync(statDir)
71
if (!c) throw 'miss'
70
- setTimeout(() => delete cache[statDir], 3_000) // invalidate shortly
72
+ setTimeout(() => delete diskSpaceCache[statDir], 3_000) // invalidate shortly
73
}
72
- const { free } = cache[statDir]
74
+ const { free } = diskSpaceCache[statDir]
75
if (typeof free !== 'number' || isNaN(free))
76
throw ''
77
if (reqSize > free - (min || 0))
@@ -78,125 +80,135 @@ export function uploadWriter(base: VfsNode, path: string, ctx: Koa.Context) {
80
catch(e: any) { // warn, but let it through
81
console.warn("can't check disk size:", e.message || String(e))
82
}
83
+ if (openFiles.has(fullPath))
84
+ return fail(HTTP_CONFLICT, 'uploading')
85
// optionally 'skip'
86
if (ctx.query.existing === 'skip' && fs.existsSync(fullPath))
83
- return fail(HTTP_CONFLICT)
84
- // if upload creates a folder, then add meta to it too
85
- if (fs.mkdirSync(dir, { recursive: true }))
86
- setUploadMeta(dir, ctx)
87
- // use temporary name while uploading
88
- const keepName = basename(fullPath).slice(-200)
89
- let tempName = join(dir, 'hfs$upload-' + keepName)
90
- const resumable = fs.existsSync(tempName) && tempName
91
- if (resumable)
92
- tempName = join(dir, 'hfs$upload2-' + keepName)
93
- // checks for resume feature
94
- let resume = Number(ctx.query.resume)
95
- const size = resumable && try_(() => fs.statSync(resumable).size)
96
- if (size === undefined) // stat failed
97
- return fail(HTTP_SERVER_ERROR)
98
- if (resume > size)
99
- return fail(HTTP_RANGE_NOT_SATISFIABLE)
100
- // warn frontend about resume possibility
101
- if (!resume && resumable) {
102
- const timeout = 30
103
- notifyClient(ctx, 'upload.resumable', { [path]: size, expires: Date.now() + timeout * 1000 })
104
- delayedDelete(resumable, timeout, () =>
105
- fs.rename(tempName, resumable, err => {
106
- if (!err)
107
- tempName = resumable
108
- }) )
109
- }
110
- // append if resuming
111
- const resuming = resume && resumable
112
- if (!resuming)
113
- resume = 0
114
- const writeStream = createStreamLimiter(reqSize ?? Infinity)
115
- if (resuming) {
116
- fs.rm(tempName, () => {})
117
- tempName = resumable
118
- }
119
- cancelDeletion(tempName)
120
- ctx.state.uploadDestinationPath = tempName
121
- // allow plugins to mess with the write-stream, because the read-stream can be complicated in case of multipart
122
- const obj = { ctx, writeStream }
123
- const resEvent = events.emit('uploadStart', obj)
124
- if (resEvent?.isDefaultPrevented()) return
87
+ return fail(HTTP_CONFLICT, 'exists')
88
+ openFiles.add(fullPath)
89
+ try {
90
+ // if upload creates a folder, then add meta to it too
91
+ if (fs.mkdirSync(dir, { recursive: true }))
92
+ setUploadMeta(dir, ctx)
93
+ // use temporary name while uploading
94
+ const keepName = basename(fullPath).slice(-200)
95
+ let tempName = join(dir, 'hfs$upload-' + keepName)
96
+ const resumable = fs.existsSync(tempName) && !openFiles.has(tempName) && tempName
97
+ if (resumable)
98
+ tempName = join(dir, 'hfs$upload2-' + keepName)
99
+ // checks for resume feature
100
+ let resume = Number(ctx.query.resume)
101
+ const size = resumable && try_(() => fs.statSync(resumable).size)
102
+ if (size === undefined) // stat failed
103
+ return fail(HTTP_SERVER_ERROR)
104
+ if (resume > size)
105
+ return fail(HTTP_RANGE_NOT_SATISFIABLE)
106
+ // warn frontend about resume possibility
107
+ if (!resume && resumable) {
108
+ const timeout = 30
109
+ notifyClient(ctx, 'upload.resumable', { [path]: size, expires: Date.now() + timeout * 1000 })
110
+ delayedDelete(resumable, timeout, () =>
111
+ fs.rename(tempName, resumable, err => {
112
+ if (!err)
113
+ tempName = resumable
114
+ }) )
115
+ }
116
+ // append if resuming
117
+ const resuming = resume && resumable
118
+ if (!resuming)
119
+ resume = 0
120
+ const writeStream = createStreamLimiter(reqSize ?? Infinity)
121
+ if (resuming) {
122
+ fs.rm(tempName, () => {})
123
+ tempName = resumable
124
+ }
125
+ cancelDeletion(tempName)
126
+ ctx.state.uploadDestinationPath = tempName
127
+ // allow plugins to mess with the write-stream, because the read-stream can be complicated in case of multipart
128
+ const obj = { ctx, writeStream }
129
+ const resEvent = events.emit('uploadStart', obj)
130
+ if (resEvent?.isDefaultPrevented()) return
131
126
- const fileStream = resuming ? fs.createWriteStream(resumable, { flags: 'r+', start: resume })
127
- : fs.createWriteStream(tempName)
128
- writeStream.pipe(fileStream)
129
- Object.assign(obj, { fileStream })
130
- trackProgress()
132
+ const fileStream = resuming ? fs.createWriteStream(resumable, { flags: 'r+', start: resume })
133
+ : fs.createWriteStream(tempName)
134
+ writeStream.pipe(fileStream)
135
+ Object.assign(obj, { fileStream })
136
+ trackProgress()
137
132
- const lockMiddleware = pendingPromise() // outside we need to know when all operations stopped
133
- writeStream.once('close', async () => {
134
- try {
135
- if (ctx.req.aborted) {
136
- if (resumable) // we don't want to be left with 2 temp files
137
- return delayedDelete(tempName, 0)
138
- const sec = deleteUnfinishedUploadsAfter.get()
139
- return _.isNumber(sec) && delayedDelete(tempName, sec)
140
- }
141
- let dest = fullPath
142
- if (dontOverwriteUploading.get() && !await overwriteAnyway() && fs.existsSync(dest)) {
143
- const ext = extname(dest)
144
- const base = dest.slice(0, -ext.length || Infinity)
145
- let i = 1
146
- do dest = `${base} (${i++})${ext}`
147
- while (fs.existsSync(dest))
148
- }
138
+ const lockMiddleware = pendingPromise() // outside we need to know when all operations stopped
139
+ writeStream.once('close', async () => {
140
try {
150
- await rename(tempName, dest)
151
- ctx.state.uploadDestinationPath = dest
152
- setUploadMeta(dest, ctx)
153
- if (ctx.query.comment)
154
- void setCommentFor(dest, String(ctx.query.comment))
155
- if (resumable)
156
- delayedDelete(resumable, 0)
157
- events.emit('uploadFinished', obj)
158
- if (resEvent) for (const cb of resEvent)
159
- if (_.isFunction(cb))
160
- cb(obj)
141
+ if (ctx.req.aborted) {
142
+ if (resumable) // we don't want to be left with 2 temp files
143
+ return delayedDelete(tempName, 0)
144
+ const sec = deleteUnfinishedUploadsAfter.get()
145
+ return _.isNumber(sec) && delayedDelete(tempName, sec)
146
+ }
147
+ let dest = fullPath
148
+ if (dontOverwriteUploading.get() && !await overwriteAnyway() && fs.existsSync(dest)) {
149
+ const ext = extname(dest)
150
+ const base = dest.slice(0, -ext.length || Infinity)
151
+ let i = 1
152
+ do dest = `${base} (${i++})${ext}`
153
+ while (fs.existsSync(dest))
154
+ }
155
+ try {
156
+ await rename(tempName, dest)
157
+ ctx.state.uploadDestinationPath = dest
158
+ setUploadMeta(dest, ctx)
159
+ if (ctx.query.comment)
160
+ void setCommentFor(dest, String(ctx.query.comment))
161
+ if (resumable)
162
+ delayedDelete(resumable, 0)
163
+ events.emit('uploadFinished', obj)
164
+ if (resEvent) for (const cb of resEvent)
165
+ if (_.isFunction(cb))
166
+ cb(obj)
167
+ }
168
+ catch (err: any) {
169
+ setUploadMeta(tempName, ctx)
170
+ console.error("couldn't rename temp to", dest, String(err))
171
+ }
172
}
162
- catch (err: any) {
163
- setUploadMeta(tempName, ctx)
164
- console.error("couldn't rename temp to", dest, String(err))
173
+ finally {
174
+ releaseFile()
175
+ lockMiddleware.resolve()
176
}
177
+ })
178
+ return Object.assign(obj.writeStream, {
179
+ lockMiddleware
180
+ })
181
+
182
+ function trackProgress() {
183
+ let lastGot = 0
184
+ let lastGotTime = 0
185
+ const opTotal = reqSize + resume
186
+ Object.assign(ctx.state, { op: 'upload', opTotal, opOffset: resume / opTotal, opProgress: 0 })
187
+ const conn = updateConnectionForCtx(ctx)
188
+ if (!conn) return
189
+ const h = setInterval(() => {
190
+ const now = Date.now()
191
+ const got = fileStream.bytesWritten
192
+ const inSpeed = roundSpeed((got - lastGot) / (now - lastGotTime))
193
+ lastGot = got
194
+ lastGotTime = now
195
+ updateConnection(conn, { inSpeed, got }, { opProgress: (resume + got) / opTotal })
196
+ }, 1000)
197
+ writeStream.once('close', () => clearInterval(h) )
198
}
167
- finally {
168
- lockMiddleware.resolve()
169
- }
170
- })
171
- return Object.assign(obj.writeStream, {
172
- lockMiddleware
173
- })
199
+ }
200
+ catch (e: any) {
201
+ releaseFile()
202
+ throw e
203
+ }
204
205
async function overwriteAnyway() {
206
if (ctx.query.overwrite === undefined // legacy pre-0.52
177
- && ctx.query.existing !== 'overwrite') return
207
+ && ctx.query.existing !== 'overwrite') return
208
const n = await getNodeByName(path, base)
209
return n && hasPermission(n, 'can_delete', ctx)
210
}
211
182
- function trackProgress() {
183
- let lastGot = 0
184
- let lastGotTime = 0
185
- const opTotal = reqSize + resume
186
- Object.assign(ctx.state, { op: 'upload', opTotal, opOffset: resume / opTotal, opProgress: 0 })
187
- const conn = updateConnectionForCtx(ctx)
188
- if (!conn) return
189
- const h = setInterval(() => {
190
- const now = Date.now()
191
- const got = fileStream.bytesWritten
192
- const inSpeed = roundSpeed((got - lastGot) / (now - lastGotTime))
193
- lastGot = got
194
- lastGotTime = now
195
- updateConnection(conn, { inSpeed, got }, { opProgress: (resume + got) / opTotal })
196
- }, 1000)
197
- writeStream.once('close', () => clearInterval(h) )
198
- }
199
-
212
function delayedDelete(path: string, secs: number, cb?: Callback) {
213
clearTimeout(waitingToBeDeleted[path])
214
waitingToBeDeleted[path] = setTimeout(() => {
@@ -210,7 +222,12 @@ export function uploadWriter(base: VfsNode, path: string, ctx: Koa.Context) {
222
delete waitingToBeDeleted[path]
223
}
224
225
+ function releaseFile() {
226
+ openFiles.delete(fullPath)
227
+ }
228
+
229
function fail(status?: number, msg?: string) {
230
+ releaseFile()
231
if (status)
232
ctx.status = status
233
if (msg)
tests/test.ts
+11
-2
@@ -6,6 +6,7 @@ import { findDefined, randomId, tryJson, wait } from '../src/cross'
6
import { httpStream, stream2string, XRequestOptions } from '../src/util-http'
7
import { ThrottledStream, ThrottleGroup } from '../src/ThrottledStream'
8
import { rm, writeFile } from 'fs/promises'
9
+import { Readable } from 'stream'
10
/*
11
import { PORT, srv } from '../src'
12
@@ -137,6 +138,14 @@ describe('after-login', () => {
138
it('upload.never', reqUpload('/random', 403))
139
it('upload.ok', reqUpload(UPLOAD_DEST, 200))
140
it('upload.crossing', reqUpload(UPLOAD_DEST.replace('temp', '../..'), 418))
141
+ it('upload.overlap', async () => {
142
+ const seconds = .3
143
+ const throttled = Readable.from(BIG_CONTENT).pipe(new ThrottledStream(new ThrottleGroup(BIG_CONTENT.length / 1000 / seconds)))
144
+ const first = reqUpload(UPLOAD_DEST, 200, throttled, BIG_CONTENT.length)()
145
+ await wait(100)
146
+ await reqUpload(UPLOAD_DEST, 409)() // should conflict
147
+ await first
148
+ })
149
const renameTo = 'z'
150
it('rename.ok', reqApi('rename', { uri: UPLOAD_DEST, dest: renameTo }, 200))
151
it('delete.miss renamed', reqApi('delete', { uri: UPLOAD_DEST }, 404))
@@ -173,11 +182,11 @@ function login(usr: string, pwd=password) {
182
reqApi(cmd, params, (x,res)=> res.statusCode < 400)())
183
}
184
176
-function reqUpload(dest: string, tester: Tester, body?: string, size?: number) {
185
+function reqUpload(dest: string, tester: Tester, body?: string | Readable, size?: number) {
186
const fn = join(__dirname, 'page/gpl.png')
187
return req(dest, tester, {
188
method: 'PUT',
180
- headers: { 'content-length': size ?? body?.length ?? statSync(fn).size },
189
+ headers: { 'content-length': size ?? (body as any)?.length ?? statSync(fn).size }, // it's ok that Readable.length is undefined
190
body: body ?? createReadStream(fn)
191
})
192
}