0001 #if os(Linux) 0002 import Glibc 0003 private let system_fork = Glibc.fork 0004 #else 0005 import Darwin.C 0006 @_silgen_name("fork") private func system_fork() -> Int32 0007 #endif 0008 0009 import Nest 0010 0011 0012 enum Address
Arbiter.swift:238 let pid = system_fork(): CustomStringConvertible { 0013 case IP
Arbiter.swift:44 let addresses: [Address]Arbiter.swift:50 init(application: RequestType -> ResponseType, workers: Int, addresses: [Address], timeout: Int) {Curassow.swift:12 extension Address : ArgumentConvertible {Curassow.swift:35 Option("bind", Address.IP(hostname: "0.0.0.0", port: 8000), description: "The address to bind sockets."),(hostname: String, port: UInt16) 0014 0015 func socket
Arbiter.swift:17 case let IP(hostname, port):Arbiter.swift:28 case let IP(hostname, port):(backlog: Int32) throws -> Socket { 0016 switch self { 0017 case let IP(hostname, port): 0018 let socket = try Socket() 0019 try socket.bind(hostname, port: port) 0020 try socket.listen(backlog) 0021 // TODO: Set socket non blocking 0022 return socket 0023 } 0024 } 0025 0026 var description: String { 0027 switch self { 0028 case let IP(hostname, port): 0029 return "\(hostname):\(port)" 0030 } 0031 } 0032 } 0033 0034 0035 /// Arbiter maintains the worker processes 0036 class Arbiter<Worker
Arbiter.swift:59 listeners.append(try address.socket(backlog)): WorkerType> { 0037 let logger
Arbiter.swift:39 var workers: [pid_t: Worker] = [:]Arbiter.swift:236 let worker = Worker(logger: logger, listeners: listeners, timeout: timeout / 2, application: application)= Logger() 0038 var listeners
Arbiter.swift:60 logger.info("Listening at http://\(address) (\(getpid()))")Arbiter.swift:112 logger.info("Shutting down")Arbiter.swift:205 logger.critical("Worker timeout (pid: \(pid))")Arbiter.swift:246 logger.info("Worker exiting (pid: \(workerPid))"): [Socket] = [] 0039 var workers
Arbiter.swift:59 listeners.append(try address.socket(backlog))Arbiter.swift:99 listeners.forEach { $0.close() }: [pid_t: Worker] = [:] 0040 let timeout
Arbiter.swift:164 workers.removeValueForKey(pid)Arbiter.swift:180 let neededWorkers = numberOfWorkers - workers.countArbiter.swift:193 for (pid, worker) in workers {Arbiter.swift:199 workers.removeValueForKey(pid)Arbiter.swift:208 workers.removeValueForKey(pid)Arbiter.swift:217 let killCount = workers.count - numberOfWorkersArbiter.swift:220 if let (pid, _) = workers.popFirst() {Arbiter.swift:229 for pid in workers.keys {: Int 0041 let backlog
Arbiter.swift:118 let timeout = timeval(tv_sec: self.timeout, tv_usec: 0)Arbiter.swift:196 if lastUpdate >= timeout {Arbiter.swift:236 let worker = Worker(logger: logger, listeners: listeners, timeout: timeout / 2, application: application): Int32 = 2048 0042 0043 var numberOfWorkers
Arbiter.swift:59 listeners.append(try address.socket(backlog)): Int 0044 let addresses
Arbiter.swift:144 numberOfWorkers += 1Arbiter.swift:150 if numberOfWorkers > 1 {Arbiter.swift:151 numberOfWorkers -= 1Arbiter.swift:180 let neededWorkers = numberOfWorkers - workers.countArbiter.swift:217 let killCount = workers.count - numberOfWorkers: [Address] 0045 0046 let application: RequestType -> ResponseType 0047 0048 var signalHandler
Arbiter.swift:58 for address in addresses {: SignalHandler! 0049 0050 init(application: RequestType -> ResponseType, workers: Int, addresses: [Address], timeout: Int) { 0051 self.application = application 0052 self.numberOfWorkers = workers 0053 self.addresses = addresses 0054 self.timeout = timeout 0055 } 0056 0057 func createSockets
Arbiter.swift:65 signalHandler = try SignalHandler()Arbiter.swift:66 signalHandler.register(.Interrupt, handleINT)Arbiter.swift:67 signalHandler.register(.Quit, handleQUIT)Arbiter.swift:68 signalHandler.register(.Terminate, handleTerminate)Arbiter.swift:69 signalHandler.register(.TTIN, handleTTIN)Arbiter.swift:70 signalHandler.register(.TTOU, handleTTOU)Arbiter.swift:71 signalHandler.register(.Child, handleChild)Arbiter.swift:72 sharedHandler = signalHandlerArbiter.swift:88 if !signalHandler.process() {Arbiter.swift:119 let (read, _, _) = select([signalHandler.pipe[0]], [], [], timeout: timeout)Arbiter.swift:123 while try signalHandler.pipe[0].read(1).count > 0 {}() throws { 0058 for address in addresses { 0059 listeners.append(try address.socket(backlog)) 0060 logger.info("Listening at http://\(address) (\(getpid()))") 0061 } 0062 } 0063 0064 func registerSignals
Arbiter.swift:83 try createSockets()() throws { 0065 signalHandler = try SignalHandler() 0066 signalHandler.register(.Interrupt, handleINT) 0067 signalHandler.register(.Quit, handleQUIT) 0068 signalHandler.register(.Terminate, handleTerminate) 0069 signalHandler.register(.TTIN, handleTTIN) 0070 signalHandler.register(.TTOU, handleTTOU) 0071 signalHandler.register(.Child, handleChild) 0072 sharedHandler = signalHandler 0073 SignalHandler.registerSignals() 0074 } 0075 0076 var running
Arbiter.swift:82 try registerSignals()= false 0077 0078 // Main run loop for the master process 0079 func run() throws { 0080 running = true 0081 0082 try registerSignals() 0083 try createSockets() 0084 0085 manageWorkers() 0086 0087 while running { 0088 if !signalHandler.process() { 0089 sleep() 0090 murderWorkers() 0091 manageWorkers() 0092 } 0093 } 0094 0095 halt() 0096 } 0097 0098 func stop
Arbiter.swift:80 running = trueArbiter.swift:87 while running {Arbiter.swift:107 running = falseArbiter.swift:139 running = false(graceful: Bool = true) { 0099 listeners.forEach { $0.close() } 0100 0101 if graceful { 0102 killWorkers(SIGTERM) 0103 } else { 0104 killWorkers(SIGQUIT) 0105 } 0106 0107 running = false 0108 } 0109 0110 func halt
Arbiter.swift:111 stop()Arbiter.swift:131 stop(false)Arbiter.swift:135 stop(false)(exitStatus: Int32 = 0) { 0111 stop() 0112 logger.info("Shutting down") 0113 exit(exitStatus) 0114 } 0115 0116 /// Sleep, waiting for stuff to happen on our signal pipe 0117 func sleep
Arbiter.swift:95 halt()() { 0118 let timeout = timeval(tv_sec: self.timeout, tv_usec: 0) 0119 let (read, _, _) = select([signalHandler.pipe[0]], [], [], timeout: timeout) 0120 0121 if !read.isEmpty { 0122 do { 0123 while try signalHandler.pipe[0].read(1).count > 0 {} 0124 } catch {} 0125 } 0126 } 0127 0128 // MARK: Handle Signals 0129 0130 func handleINT
Arbiter.swift:89 sleep()() { 0131 stop(false) 0132 } 0133 0134 func handleQUIT
Arbiter.swift:66 signalHandler.register(.Interrupt, handleINT)() { 0135 stop(false) 0136 } 0137 0138 func handleTerminate
Arbiter.swift:67 signalHandler.register(.Quit, handleQUIT)() { 0139 running = false 0140 } 0141 0142 /// Increases the amount of workers by one 0143 func handleTTIN
Arbiter.swift:68 signalHandler.register(.Terminate, handleTerminate)() { 0144 numberOfWorkers += 1 0145 manageWorkers() 0146 } 0147 0148 /// Decreases the amount of workers by one 0149 func handleTTOU
Arbiter.swift:69 signalHandler.register(.TTIN, handleTTIN)() { 0150 if numberOfWorkers > 1 { 0151 numberOfWorkers -= 1 0152 manageWorkers() 0153 } 0154 } 0155 0156 func handleChild
Arbiter.swift:70 signalHandler.register(.TTOU, handleTTOU)() { 0157 while true { 0158 var stat: Int32 = 0 0159 let pid = waitpid(-1, &stat, WNOHANG) 0160 if pid == -1 { 0161 break 0162 } 0163 0164 workers.removeValueForKey(pid) 0165 } 0166 0167 manageWorkers() 0168 } 0169 0170 // MARK: Worker 0171 0172 // Maintain number of workers by spawning or killing as required. 0173 func manageWorkers
Arbiter.swift:71 signalHandler.register(.Child, handleChild)() { 0174 spawnWorkers() 0175 murderExcessWorkers() 0176 } 0177 0178 // Spawn workers until we have enough 0179 func spawnWorkers
Arbiter.swift:85 manageWorkers()Arbiter.swift:91 manageWorkers()Arbiter.swift:145 manageWorkers()Arbiter.swift:152 manageWorkers()Arbiter.swift:167 manageWorkers()() { 0180 let neededWorkers = numberOfWorkers - workers.count 0181 if neededWorkers > 0 { 0182 for _ in 0..<neededWorkers { 0183 spawnWorker() 0184 } 0185 } 0186 } 0187 0188 // Murder workers that have timed out 0189 func murderWorkers
Arbiter.swift:174 spawnWorkers()() { 0190 var currentTime = timeval() 0191 gettimeofday(¤tTime, nil) 0192 0193 for (pid, worker) in workers { 0194 let lastUpdate = currentTime.tv_sec - worker.temp.lastUpdate.tv_sec 0195 0196 if lastUpdate >= timeout { 0197 if worker.aborted { 0198 if kill(pid, SIGKILL) == ESRCH { 0199 workers.removeValueForKey(pid) 0200 } 0201 } else { 0202 var worker = worker 0203 worker.aborted = true 0204 0205 logger.critical("Worker timeout (pid: \(pid))") 0206 0207 if kill(pid, SIGABRT) == ESRCH { 0208 workers.removeValueForKey(pid) 0209 } 0210 } 0211 } 0212 } 0213 } 0214 0215 // Murder unused workers, oldest first 0216 func murderExcessWorkers
Arbiter.swift:90 murderWorkers()() { 0217 let killCount = workers.count - numberOfWorkers 0218 if killCount > 0 { 0219 for _ in 0..<killCount { 0220 if let (pid, _) = workers.popFirst() { 0221 kill(pid, SIGKILL) 0222 } 0223 } 0224 } 0225 } 0226 0227 // Kill all workers with given signal 0228 func killWorkers
Arbiter.swift:175 murderExcessWorkers()(signal: Int32) { 0229 for pid in workers.keys { 0230 kill(pid, signal) 0231 } 0232 } 0233 0234 // Spawns a new worker process 0235 func spawnWorker
Arbiter.swift:102 killWorkers(SIGTERM)Arbiter.swift:104 killWorkers(SIGQUIT)() { 0236 let worker = Worker(logger: logger, listeners: listeners, timeout: timeout / 2, application: application) 0237 0238 let pid = system_fork() 0239 if pid != 0 { 0240 workers[pid] = worker 0241 return 0242 } 0243 0244 let workerPid = getpid() 0245 worker.run() 0246 logger.info("Worker exiting (pid: \(workerPid))") 0247 exit(0) 0248 } 0249 } 0250
Arbiter.swift:183 spawnWorker()