From 929944ef091aad3742e39bc966bb0f98c54639bd Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 13:13:20 -0400 Subject: [PATCH 1/7] use NonBlockingFileIO for Process.execute --- Sources/Core/Process+Execute.swift | 53 +++++++++++++++++------------- 1 file changed, 31 insertions(+), 22 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index 48bb6ff0..fe080761 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -94,26 +94,30 @@ extension Process { if program.hasPrefix("/") { let stdout = Pipe() let stderr = Pipe() - - // will be set to false when the program is done - var running = true - // readabilityHandler doesn't work on linux, so we are left with this hack - DispatchQueue.global().async { - while running { - let stdout = stdout.fileHandleForReading.availableData - if !stdout.isEmpty { - output(.stdout(stdout)) - } + let threadPool = BlockingIOThreadPool(numberOfThreads: 3) + threadPool.start() + let file = NonBlockingFileIO(threadPool: threadPool) + let allocator = ByteBufferAllocator() + + let outfile = FileHandle(descriptor: stdout.fileHandleForReading.fileDescriptor) + _ = file.readChunked(fileHandle: outfile, byteCount: .max, allocator: allocator, eventLoop: worker.eventLoop) { chunk in + if chunk.readableBytes > 0 { + var chunk = chunk + let data = chunk.readData(length: chunk.readableBytes)! + output(.stdout(data)) } + return worker.future() } - DispatchQueue.global().async { - while running { - let stderr = stderr.fileHandleForReading.availableData - if !stderr.isEmpty { - output(.stderr(stderr)) - } + + let errfile = FileHandle(descriptor: stderr.fileHandleForReading.fileDescriptor) + _ = file.readChunked(fileHandle: errfile, byteCount: .max, allocator: allocator, eventLoop: worker.eventLoop) { chunk in + if chunk.readableBytes > 0 { + var chunk = chunk + let data = chunk.readData(length: chunk.readableBytes)! + output(.stderr(data)) } + return worker.future() } // stdout.fileHandleForReading.readabilityHandler = { handle in @@ -130,15 +134,20 @@ extension Process { // } // output(.stderr(data)) // } - - let promise = worker.eventLoop.newPromise(Int32.self) - DispatchQueue.global().async { + + let res = threadPool.runIfActive(eventLoop: worker.eventLoop) { () -> Int32 in let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) process.waitUntilExit() - running = false - promise.succeed(result: process.terminationStatus) + return process.terminationStatus } - return promise.futureResult + + res.always { + try? errfile.close() + try? outfile.close() + threadPool.shutdownGracefully { _ in } + } + + return res } else { var resolvedPath: String? return asyncExecute("/bin/sh", ["-c", "which \(program)"], on: worker) { o in From 8d2803df792755ac43c63bc92afdde17d097de35 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 13:33:13 -0400 Subject: [PATCH 2/7] use DispatchIO for process execute --- Sources/Core/Process+Execute.swift | 61 +++++++++++++++--------------- 1 file changed, 31 insertions(+), 30 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index fe080761..c4e5d059 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -92,32 +92,28 @@ extension Process { /// - returns: A future containing the termination status of the process. public static func asyncExecute(_ program: String, _ arguments: [String], on worker: Worker, _ output: @escaping (ProcessOutput) -> ()) -> Future { if program.hasPrefix("/") { + // create process data pipes let stdout = Pipe() let stderr = Pipe() - let threadPool = BlockingIOThreadPool(numberOfThreads: 3) - threadPool.start() - let file = NonBlockingFileIO(threadPool: threadPool) - let allocator = ByteBufferAllocator() + // create dispatch sources for the pipes + let stdoutsource = DispatchSource.makeReadSource(fileDescriptor: stdout.fileHandleForReading.fileDescriptor) + let stderrsource = DispatchSource.makeReadSource(fileDescriptor: stderr.fileHandleForReading.fileDescriptor) - let outfile = FileHandle(descriptor: stdout.fileHandleForReading.fileDescriptor) - _ = file.readChunked(fileHandle: outfile, byteCount: .max, allocator: allocator, eventLoop: worker.eventLoop) { chunk in - if chunk.readableBytes > 0 { - var chunk = chunk - let data = chunk.readData(length: chunk.readableBytes)! - output(.stdout(data)) - } - return worker.future() + // setup read handlers for the output sources + stdoutsource.setEventHandler { + let data = stdout.fileHandleForReading.availableData + guard !data.isEmpty else { + return + } + output(.stdout(data)) } - - let errfile = FileHandle(descriptor: stderr.fileHandleForReading.fileDescriptor) - _ = file.readChunked(fileHandle: errfile, byteCount: .max, allocator: allocator, eventLoop: worker.eventLoop) { chunk in - if chunk.readableBytes > 0 { - var chunk = chunk - let data = chunk.readData(length: chunk.readableBytes)! - output(.stderr(data)) + stderrsource.setEventHandler { + let data = stderr.fileHandleForReading.availableData + guard !data.isEmpty else { + return } - return worker.future() + output(.stderr(data)) } // stdout.fileHandleForReading.readabilityHandler = { handle in @@ -135,19 +131,24 @@ extension Process { // output(.stderr(data)) // } - let res = threadPool.runIfActive(eventLoop: worker.eventLoop) { () -> Int32 in + let promise = worker.eventLoop.newPromise(Int32.self) + DispatchQueue.global().async { + // start the output sources + stdoutsource.resume() + stderrsource.resume() + + // launch and run the process let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) process.waitUntilExit() - return process.terminationStatus + + // cleanup output sources + stdoutsource.cancel() + stderrsource.cancel() + + // succeed with the termination status + promise.succeed(result: process.terminationStatus) } - - res.always { - try? errfile.close() - try? outfile.close() - threadPool.shutdownGracefully { _ in } - } - - return res + return promise.futureResult } else { var resolvedPath: String? return asyncExecute("/bin/sh", ["-c", "which \(program)"], on: worker) { o in From 509879a6c039149454c311a32995f68f3dbc30c9 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 13:38:43 -0400 Subject: [PATCH 3/7] use terminationHandler instead of blocking --- Sources/Core/Process+Execute.swift | 36 +++++++++--------------------- 1 file changed, 11 insertions(+), 25 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index c4e5d059..d7c07cef 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -115,39 +115,25 @@ extension Process { } output(.stderr(data)) } - - // stdout.fileHandleForReading.readabilityHandler = { handle in - // let data = handle.availableData - // guard !data.isEmpty else { - // return - // } - // output(.stdout(data)) - // } - // stderr.fileHandleForReading.readabilityHandler = { handle in - // let data = handle.availableData - // guard !data.isEmpty else { - // return - // } - // output(.stderr(data)) - // } + // start the output sources + stdoutsource.resume() + stderrsource.resume() + + // launch and run the process + let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) + + // succeed with the termination status let promise = worker.eventLoop.newPromise(Int32.self) - DispatchQueue.global().async { - // start the output sources - stdoutsource.resume() - stderrsource.resume() - - // launch and run the process - let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) - process.waitUntilExit() - + process.terminationHandler = { process in // cleanup output sources stdoutsource.cancel() stderrsource.cancel() - // succeed with the termination status + // complete the promise promise.succeed(result: process.terminationStatus) } + return promise.futureResult } else { var resolvedPath: String? From 84454478be9c1e56ff12e2b70406af670be25655 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 14:09:27 -0400 Subject: [PATCH 4/7] use DispatchIO stream --- Sources/Core/Process+Execute.swift | 71 +++++++++++++++++------------- 1 file changed, 41 insertions(+), 30 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index d7c07cef..d2bbe819 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -92,49 +92,54 @@ extension Process { /// - returns: A future containing the termination status of the process. public static func asyncExecute(_ program: String, _ arguments: [String], on worker: Worker, _ output: @escaping (ProcessOutput) -> ()) -> Future { if program.hasPrefix("/") { + // create queue for async work + // create process data pipes let stdout = Pipe() let stderr = Pipe() - - // create dispatch sources for the pipes - let stdoutsource = DispatchSource.makeReadSource(fileDescriptor: stdout.fileHandleForReading.fileDescriptor) - let stderrsource = DispatchSource.makeReadSource(fileDescriptor: stderr.fileHandleForReading.fileDescriptor) - - // setup read handlers for the output sources - stdoutsource.setEventHandler { - let data = stdout.fileHandleForReading.availableData - guard !data.isEmpty else { - return - } - output(.stdout(data)) + + // setup dispatch io for stdout + let stdoutPromise = worker.eventLoop.newPromise(Void.self) + let stdoutIO = DispatchIO(type: .stream, fileDescriptor: stdout.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in + close(stdout.fileHandleForReading.fileDescriptor) } - stderrsource.setEventHandler { - let data = stderr.fileHandleForReading.availableData - guard !data.isEmpty else { - return + stdoutIO.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in + if done { + stdoutPromise.succeed(result: ()) + } else if err != 0 { + stdoutPromise.fail(error: ProcessIOError(code: err)) + } else if let data = data, !data.isEmpty { + output(.stdout(Data(data))) } - output(.stderr(data)) } - - // start the output sources - stdoutsource.resume() - stderrsource.resume() + // setup dispatch io for stderr + let stderrPromise = worker.eventLoop.newPromise(Void.self) + let stderrIO = DispatchIO(type: .stream, fileDescriptor: stderr.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in + close(stderr.fileHandleForReading.fileDescriptor) + } + stderrIO.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in + if done { + stderrPromise.succeed(result: ()) + } else if err != 0 { + stderrPromise.fail(error: ProcessIOError(code: err)) + } else if let data = data, !data.isEmpty { + output(.stderr(Data(data))) + } + } + // launch and run the process let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) - - // succeed with the termination status - let promise = worker.eventLoop.newPromise(Int32.self) + + // create a new promise for the termination status and set callback + let processPromise = worker.eventLoop.newPromise(Int32.self) process.terminationHandler = { process in - // cleanup output sources - stdoutsource.cancel() - stderrsource.cancel() - // complete the promise - promise.succeed(result: process.terminationStatus) + processPromise.succeed(result: process.terminationStatus) } - return promise.futureResult + return stdoutPromise.futureResult.and(stderrPromise.futureResult) + .transform(to: processPromise.futureResult) } else { var resolvedPath: String? return asyncExecute("/bin/sh", ["-c", "which \(program)"], on: worker) { o in @@ -177,6 +182,10 @@ public struct ProcessExecuteError: Error { public var stdout: String } +private struct ProcessIOError: Error { + var code: Int32 +} + extension ProcessExecuteError: Debuggable { /// See `Debuggable.identifier`. public var identifier: String { @@ -189,4 +198,6 @@ extension ProcessExecuteError: Debuggable { } } +private let executeQueue = DispatchQueue(label: "codes.vapor.core.async.execute") + #endif From 90b2d45399f27918f63edc6ea11380fef2779ea8 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 14:17:49 -0400 Subject: [PATCH 5/7] add comments to dispatchio file code --- Sources/Core/Process+Execute.swift | 68 ++++++++++++++++-------------- 1 file changed, 37 insertions(+), 31 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index d2bbe819..de974cc4 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -92,41 +92,42 @@ extension Process { /// - returns: A future containing the termination status of the process. public static func asyncExecute(_ program: String, _ arguments: [String], on worker: Worker, _ output: @escaping (ProcessOutput) -> ()) -> Future { if program.hasPrefix("/") { - // create queue for async work + // generic dispatch io config code for stdout/stderr + func setupDispatchIO(for pipe: Pipe, onData: @escaping (Data) -> ()) -> (Future, DispatchIO) { + // create new promise to signal IO is done + let promise = worker.eventLoop.newPromise(Void.self) + + // create dispatch io stream on read handle of pipe + let io = DispatchIO(type: .stream, fileDescriptor: pipe.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in + // close pipe once stream is cancelled + close(pipe.fileHandleForReading.fileDescriptor) + } + + // start async read on stream + io.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in + if done { + // signal IO is done + promise.succeed(result: ()) + } else if err != 0 { + // signal IO failure + promise.fail(error: ProcessIOError(code: err)) + } else if let data = data, !data.isEmpty { + // convert DispatchData to Data and pass to callback + onData(Data(data)) + } + } + + // return created future and IO stream + return (promise.futureResult, io) + } // create process data pipes let stdout = Pipe() let stderr = Pipe() - // setup dispatch io for stdout - let stdoutPromise = worker.eventLoop.newPromise(Void.self) - let stdoutIO = DispatchIO(type: .stream, fileDescriptor: stdout.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in - close(stdout.fileHandleForReading.fileDescriptor) - } - stdoutIO.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in - if done { - stdoutPromise.succeed(result: ()) - } else if err != 0 { - stdoutPromise.fail(error: ProcessIOError(code: err)) - } else if let data = data, !data.isEmpty { - output(.stdout(Data(data))) - } - } - - // setup dispatch io for stderr - let stderrPromise = worker.eventLoop.newPromise(Void.self) - let stderrIO = DispatchIO(type: .stream, fileDescriptor: stderr.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in - close(stderr.fileHandleForReading.fileDescriptor) - } - stderrIO.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in - if done { - stderrPromise.succeed(result: ()) - } else if err != 0 { - stderrPromise.fail(error: ProcessIOError(code: err)) - } else if let data = data, !data.isEmpty { - output(.stderr(Data(data))) - } - } + // setup dispatch IO on each pipe + let (stdoutFuture, stdoutIO) = setupDispatchIO(for: stdout) { output(.stdout($0)) } + let (stderrFuture, stderrIO) = setupDispatchIO(for: stderr) { output(.stderr($0)) } // launch and run the process let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) @@ -134,11 +135,16 @@ extension Process { // create a new promise for the termination status and set callback let processPromise = worker.eventLoop.newPromise(Int32.self) process.terminationHandler = { process in + // close dispatch IO + stdoutIO.close() + stderrIO.close() + // complete the promise processPromise.succeed(result: process.terminationStatus) } - return stdoutPromise.futureResult.and(stderrPromise.futureResult) + // combine stdout/err and process result futures + return stdoutFuture.and(stderrFuture) .transform(to: processPromise.futureResult) } else { var resolvedPath: String? From b054783487e553e90ca7ec2ce7064610b4eb5a66 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 15:09:43 -0400 Subject: [PATCH 6/7] docker testing --- .dockerignore | 6 ++++++ docker-compose.yml | 7 +++++++ test.Dockerfile | 4 ++++ 3 files changed, 17 insertions(+) create mode 100644 .dockerignore create mode 100644 docker-compose.yml create mode 100644 test.Dockerfile diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 00000000..0973f967 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,6 @@ +.git +.build +DerivedData +Package.resolved +*.xcodeproj + diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 00000000..4c828fd5 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,7 @@ +version: '2' +services: + test: + build: + context: . + dockerfile: test.Dockerfile + diff --git a/test.Dockerfile b/test.Dockerfile new file mode 100644 index 00000000..11256561 --- /dev/null +++ b/test.Dockerfile @@ -0,0 +1,4 @@ +FROM swift:4.2 +COPY . . +ENTRYPOINT swift test + From 22e7a42dad6db87131a7427414db17898347ab47 Mon Sep 17 00:00:00 2001 From: tanner0101 Date: Wed, 17 Oct 2018 16:01:39 -0400 Subject: [PATCH 7/7] add debug prints --- Sources/Core/Process+Execute.swift | 26 ++++++++++++++++++-------- 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/Sources/Core/Process+Execute.swift b/Sources/Core/Process+Execute.swift index de974cc4..4751556a 100644 --- a/Sources/Core/Process+Execute.swift +++ b/Sources/Core/Process+Execute.swift @@ -99,12 +99,16 @@ extension Process { // create dispatch io stream on read handle of pipe let io = DispatchIO(type: .stream, fileDescriptor: pipe.fileHandleForReading.fileDescriptor, queue: executeQueue) { _ in + print("[EVENT] \(pipe.fileHandleForReading.fileDescriptor) cancel") // close pipe once stream is cancelled close(pipe.fileHandleForReading.fileDescriptor) } + io.setLimit(lowWater: 0) + // start async read on stream io.read(offset: 0, length: .max, queue: executeQueue) { done, data, err in + print("[EVENT] \(pipe.fileHandleForReading.fileDescriptor) \(done)") if done { // signal IO is done promise.succeed(result: ()) @@ -116,6 +120,7 @@ extension Process { onData(Data(data)) } } + io.resume() // return created future and IO stream return (promise.futureResult, io) @@ -124,14 +129,16 @@ extension Process { // create process data pipes let stdout = Pipe() let stderr = Pipe() - - // setup dispatch IO on each pipe - let (stdoutFuture, stdoutIO) = setupDispatchIO(for: stdout) { output(.stdout($0)) } - let (stderrFuture, stderrIO) = setupDispatchIO(for: stderr) { output(.stderr($0)) } - + print("stdout: \(stdout.fileHandleForReading.fileDescriptor)") + print("stderr: \(stderr.fileHandleForReading.fileDescriptor)") + // launch and run the process let process = launchProcess(path: program, arguments, stdout: stdout, stderr: stderr) - + + // setup dispatch IO on each pipe + let (stderrFuture, stderrIO) = setupDispatchIO(for: stderr) { output(.stderr($0)) } + let (stdoutFuture, stdoutIO) = setupDispatchIO(for: stdout) { output(.stdout($0)) } + // create a new promise for the termination status and set callback let processPromise = worker.eventLoop.newPromise(Int32.self) process.terminationHandler = { process in @@ -139,10 +146,13 @@ extension Process { stdoutIO.close() stderrIO.close() + print("process complete: \(process.terminationStatus)") + // complete the promise processPromise.succeed(result: process.terminationStatus) } - + process.launch() + // combine stdout/err and process result futures return stdoutFuture.and(stderrFuture) .transform(to: processPromise.futureResult) @@ -171,7 +181,6 @@ extension Process { process.arguments = arguments process.standardOutput = stdout process.standardError = stderr - process.launch() return process } } @@ -205,5 +214,6 @@ extension ProcessExecuteError: Debuggable { } private let executeQueue = DispatchQueue(label: "codes.vapor.core.async.execute") +private let executeQueue2 = DispatchQueue(label: "codes.vapor.core.async.execute2") #endif