Compare commits

...
34 Commits
Author SHA1 Message Date
roberto@sirius 79fa25c4e0 latest (i hope) pip hack 2010-05-01 08:10:36 +02:00
roberto@sirius 02050675a8 pip hack [5] 2010-05-01 08:06:44 +02:00
roberto@sirius 0d5857a1bc pip hack [4] 2010-05-01 08:04:05 +02:00
roberto@sirius e3a51be1f9 pip hack [3] 2010-05-01 07:58:53 +02:00
roberto@sirius 57cb098a6e move to setuptools 2010-05-01 07:39:56 +02:00
roberto@sirius 5f285eef56 pip hack [2] 2010-05-01 07:29:07 +02:00
roberto@sirius 7d5cf6f318 pip hack 2010-05-01 07:14:14 +02:00
roberto@silente 99d305c7a1 snmp big endian fixes 2010-04-30 16:41:57 +02:00
roberto@sirius b9cc3b6a7d snmp fix 2010-04-30 15:51:32 +02:00
roberto@sirius e1f8c69d3d minor fixes 2010-04-30 15:16:18 +02:00
roberto@sirius 67b6fbd15b admin modifier fix 2010-04-30 14:46:27 +02:00
roberto@sirius 50f407eb9e no_server fix 2010-04-30 14:41:30 +02:00
roberto@freebsd b127699c0a FreeBSD/Dragonfly little fix 2010-04-30 14:12:51 +02:00
roberto@netbsd 4d74fa6759 small netbsd include fix 2010-04-30 16:04:02 +02:00
roberto@sirius 2b9d6cb7c2 uwsgi_log everywhere 2010-04-30 12:28:01 +02:00
roberto@solaris e0b67bba58 uwsgi_log instead of fprintf [1] 2010-04-30 12:22:07 +02:00
roberto@silente 3cdf17df75 strlcopy insted of strcpy 2010-04-30 09:54:25 +02:00
roberto@silente 43c086494b be more verbose on extreme cases 2010-04-30 09:18:53 +02:00
roberto@silente ae26e90b3a OpenBSD ha no RLIMIT_AS and swapcontext(), added an harakiri message to avoid user confusion 2010-04-30 09:01:26 +02:00
roberto@sirius d35a7d8725 updated urack.rb for unbit 2010-04-27 12:47:29 +02:00
roberto@sirius 043e7ddde4 uWSGI 0.9.5-rc2 2010-04-20 10:51:54 +02:00
roberto@sirius 3930960159 psycogreen example 2010-04-20 10:30:57 +02:00
roberto@sirius be12175c8a psycopg2 async example 2010-04-20 10:19:00 +02:00
roberto@sirius 02a1847539 async memory fixes 2010-04-19 16:01:13 +02:00
roberto@solaris 0235535658 /dev/poll complete support 2010-04-19 13:43:36 +02:00
roberto@sirius 3af75ca64c default uwsgiconfig.py 2010-04-19 10:14:21 +02:00
roberto@sirius efbe28502d new wsgi_request struct fixes 2010-04-19 10:13:38 +02:00
roberto@solaris 47847ee471 Solaris/Sparc fixes 2010-04-19 08:34:46 +02:00
roberto@sirius 688ffe32a9 urack refactoring 2010-04-15 17:02:29 +02:00
roberto@sirius c6bbe01f3e fix for Plack handler 2010-04-15 12:59:24 +02:00
roberto@sirius 002af6128c UWSGI_FD for plack handler 2010-04-15 11:47:39 +02:00
roberto@sirius 5b5c92a9d7 fix for ruby 1.8.6 2010-04-15 11:01:18 +02:00
roberto@sirius 7d31d73fc5 UWSGI_FD for rack handler 2010-04-15 10:56:25 +02:00
roberto@sirius c5c2c83c65 uWSGI 0.9.5-rc1 2010-04-14 07:26:32 +02:00
33 changed files with 1105 additions and 905 deletions
+2
View File
@@ -1,2 +1,4 @@
d3850b52334a91d005b0fe45e1e38ba0997a1fb8 no_server mode
2d8a617299153ed2efa6ecee7920e9d36c29e0f5 0.9.5beta1
04722a2d0121080d14a64df097c278860db5ce8a 0.9.5rc1
820cc64ab296dc4589ef72a33a6b804c22a0430c 0.9.5rc2
+34
View File
@@ -1,3 +1,37 @@
*** may 2010 ***
* 0.9.5 [20100501] *
- hook based request/modifier management
- plugin support via dlopen
- on-the-fly management flag
- integrated proxy
- logging via udp
- improved spooler for cron-like apps
- async support
- green thread platform (uGreen) on top of teh async mode
- transparent Erlang integration
- embedded snmp agent
- nagios mode
- improved python 3.x support
- improved xml configuration
- new build system
- address space usage limiting
- lot of portability fixes
- lot of optimizations and code refactoring
*** april 2010 ***
* Fourth Maintenance release [20100427] *
- backported non-yielding app optimization from 0.9.5
- backported udp logging from 0.9.5
- backported uwsgi_error() from 0.9.5 (no more dumb perror())
- fix a rare (but possible) segmentation fault
- added --version
- updated apache2 module
*** march 2010 ***
* Third Maintenance release [20100310] *
+31 -26
View File
@@ -14,7 +14,7 @@ int async_queue_init(int serverfd) {
epfd = epoll_create(256);
if (epfd < 0) {
perror("epoll_create()");
uwsgi_error("epoll_create()");
return -1 ;
}
@@ -23,7 +23,7 @@ int async_queue_init(int serverfd) {
ee.data.fd = serverfd;
if (epoll_ctl(epfd, EPOLL_CTL_ADD, serverfd, &ee)) {
perror("epoll_ctl()");
uwsgi_error("epoll_ctl()");
close(epfd);
return -1;
}
@@ -42,10 +42,10 @@ int async_wait(int queuefd, void *events, int nevents, int block, int timeout) {
timeout = timeout*1000;
}
//fprintf(stderr,"waiting with timeout %d nevents %d\n", timeout, nevents);
//uwsgi_log("waiting with timeout %d nevents %d\n", timeout, nevents);
ret = epoll_wait(queuefd, (struct epoll_event *) events, nevents, timeout);
if (ret < 0) {
perror("epoll_wait()");
uwsgi_error("epoll_wait()");
}
return ret ;
}
@@ -58,7 +58,7 @@ int async_add(int queuefd, int fd, int etype) {
ee.data.fd = fd;
if (epoll_ctl(queuefd, EPOLL_CTL_ADD, fd, &ee)) {
perror("epoll_ctl()");
uwsgi_error("epoll_ctl()");
return -1;
}
@@ -73,7 +73,7 @@ int async_mod(int queuefd, int fd, int etype) {
ee.data.fd = fd;
if (epoll_ctl(queuefd, EPOLL_CTL_MOD, fd, &ee)) {
perror("epoll_ctl()");
uwsgi_error("epoll_ctl()");
return -1;
}
@@ -88,7 +88,7 @@ int async_del(int queuefd, int fd, int etype) {
ee.data.fd = fd;
if (epoll_ctl(queuefd, EPOLL_CTL_DEL, fd, &ee)) {
perror("epoll_ctl()");
uwsgi_error("epoll_ctl()");
return -1;
}
@@ -101,21 +101,25 @@ int async_queue_init(int serverfd) {
int dpfd ;
struct pollfd dpev;
dpfd = open("/dev/poll", O_RDWR);
if (dpfd < 0) {
perror("open()");
uwsgi_error("open()");
return -1 ;
}
dpev.fd = serverfd;
dpev.events = POLLIN ;
dpev.revents = 0;
if (write(dpfd, &dpev, sizeof(struct pollfd)) < 0) {
perror("write()");
uwsgi_error("write()");
return -1;
}
return dpfd;
}
@@ -135,10 +139,10 @@ int async_wait(int queuefd, void *events, int nevents, int block, int timeout) {
dv.dp_nfds = nevents;
dv.dp_timeout = timeout;
//fprintf(stderr,"waiting with timeout %d nevents %d\n", timeout, nevents);
//uwsgi_log("waiting with timeout %d nevents %d\n", timeout, nevents);
ret = ioctl(queuefd, DP_POLL, &dv);
if (ret < 0) {
perror("ioctl()");
uwsgi_error("ioctl()");
}
return ret ;
}
@@ -148,9 +152,10 @@ int async_add(int queuefd, int fd, int etype) {
pl.fd = fd ;
pl.events = etype ;
pl.revents = 0 ;
if (write(queuefd, &pl, sizeof(struct pollfd))) {
perror("write()");
if (write(queuefd, &pl, sizeof(struct pollfd)) < 0) {
uwsgi_error("write()");
return -1;
}
@@ -163,8 +168,8 @@ int async_mod(int queuefd, int fd, int etype) {
}
int async_del(int queuefd, int fd, int etype) {
// to remove an fd from /dev/poll you have to simply close it
return 0;
// use POLLREMOVE to remove an fd
return async_add(queuefd, fd, POLLREMOVE);
}
#else
@@ -175,13 +180,13 @@ int async_queue_init(int serverfd) {
kfd = kqueue();
if (kfd < 0) {
perror("kqueue()");
uwsgi_error("kqueue()");
return -1 ;
}
EV_SET(&kev, serverfd, EVFILT_READ, EV_ADD, 0, 0, 0);
if (kevent(kfd, &kev, 1, NULL, 0, NULL) < 0) {
perror("kevent()");
uwsgi_error("kevent()");
return -1;
}
@@ -211,7 +216,7 @@ int async_wait(int queuefd, void *events, int nevents, int block, int timeout) {
}
if (ret < 0) {
perror("kevent()");
uwsgi_error("kevent()");
}
return ret;
@@ -223,7 +228,7 @@ int async_add(int queuefd, int fd, int etype) {
EV_SET(&kev, fd, etype, EV_ADD, 0, 0, 0);
if (kevent(queuefd, &kev, 1, NULL, 0, NULL) < 0) {
perror("kevent()");
uwsgi_error("kevent()");
return -1;
}
return 0;
@@ -234,13 +239,13 @@ int async_mod(int queuefd, int fd, int etype) {
EV_SET(&kev, fd, ASYNC_OUT, EV_DISABLE, 0, 0, 0);
if (kevent(queuefd, &kev, 1, NULL, 0, NULL) < 0) {
perror("kevent()");
uwsgi_error("kevent()");
return -1;
}
EV_SET(&kev, fd, etype, EV_ADD, 0, 0, 0);
if (kevent(queuefd, &kev, 1, NULL, 0, NULL) < 0) {
perror("kevent()");
uwsgi_error("kevent()");
return -1;
}
return 0;
@@ -251,7 +256,7 @@ int async_del(int queuefd, int fd, int etype) {
EV_SET(&kev, fd, etype, EV_DELETE, 0, 0, 0);
if (kevent(queuefd, &kev, 1, NULL, 0, NULL) < 0) {
perror("kevent()");
uwsgi_error("kevent()");
return -1;
}
@@ -260,11 +265,11 @@ int async_del(int queuefd, int fd, int etype) {
#endif
struct wsgi_request *next_wsgi_req(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req) {
inline struct wsgi_request *next_wsgi_req(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req) {
uint8_t *ptr = (uint8_t *) wsgi_req ;
ptr += sizeof(struct wsgi_request)+(uwsgi->buffer_size-1) ;
ptr += sizeof(struct wsgi_request) ;
return (struct wsgi_request *) ptr ;
}
@@ -340,7 +345,7 @@ struct wsgi_request *find_wsgi_req_by_id(struct uwsgi_server *uwsgi, int async_i
uint8_t *ptr = (uint8_t *) uwsgi->wsgi_requests ;
ptr += (sizeof(struct wsgi_request)+(uwsgi->buffer_size-1)) * async_id ;
ptr += sizeof(struct wsgi_request) * async_id ;
return (struct wsgi_request *) ptr ;
}
@@ -389,7 +394,7 @@ void async_write_all(struct uwsgi_server *uwsgi, char *data, size_t len) {
if (wsgi_req->async_status == UWSGI_PAUSED) {
rlen = write(wsgi_req->poll.fd, data, len);
if (rlen < 0) {
perror("write()");
uwsgi_error("write()");
}
else {
wsgi_req->response_size += rlen ;
+8 -1
View File
@@ -17,7 +17,14 @@ sub new {
sub run {
my ($self, $app) = @_;
my $server = IO::Socket::INET->new(LocalPort => $self->{port}, LocalAddr => $self->{host}, Listen => 100, ReuseAddr => 1);
my $server ;
if (exists($ENV{'UWSGI_FD'})) {
$server = IO::Socket::UNIX->new_from_fd($ENV{'UWSGI_FD'}, '+<');
}
else {
$server = IO::Socket::INET->new(LocalPort => $self->{port}, LocalAddr => $self->{host}, Listen => 100, ReuseAddr => 1);
}
while ( my $client = $server->accept ) {
+276 -372
View File
@@ -1,403 +1,307 @@
# Copyright 2009 Unbit S.a.s. <info@unbit.it>
# see the COPYRIGHT file
require 'rubygems'
require 'socket'
require 'rack'
require 'rack/content_length'
require 'rack/rewindable_input'
require 'stringio'
$stdout.sync = true
options = {}
require 'optparse'
if Process.uid == 0
puts "Never run uRack as root !!!"
exit
options = {
:processes => 1,
:master => false,
:logging => true,
:static => true,
:config => nil,
:cachehack => false
}
opts = OptionParser.new do |opts|
opts.on("-p", "--processes=nproc", Integer, '') { |p| options[:processes] = p }
opts.on("-M", "--master", '') { options[:master] = true }
opts.on("-L", "--no-logging", '') { options[:logging] = false }
opts.on("-S", "--no-static", '') { options[:static] = false }
opts.on("-C", "--cache-hack", '') { options[:cachehack] = true }
opts.parse! ARGV
end
$stderr.puts "[#{Time.new}] starting uRack"
if ARGV[0]
options[:config] = ARGV[0]
end
def human_round(float)
(float * (10 ** 2)).round / (10 ** 2).to_f
end
unless ENV.has_key?('UNBIT_RACK_PATH')
if File.exists?("#{ARGV.last}/config/environment.rb")
RAILS_ROOT = String.new(ARGV.last)
else
RAILS_ROOT = String.new(Dir.getwd)
end
options[:sockname] = '/tmp/uwsgi.sock'
options[:sock_chmod] = 0
options[:stderr_logfile] = nil
else
$domain = ARGV.last.gsub('[','').gsub(']','')
$unbit_log_base = "[uRack/Unbit on #{$domain}]"
RAILS_ROOT = String.new(ENV['UNBIT_RACK_PATH'])
options[:unbit] = true
$limit_as = human_round(Process.getrlimit(Process::RLIMIT_AS)[1].to_f/1024/1024)
puts "#{$unbit_log_base} process address space limit is #{$limit_as} MB"
$uidsec_size = syscall(357,0,0)
puts "#{$unbit_log_base} need #{$uidsec_size} bytes to store uidsec_struct..."
$uidsec = '1' * $uidsec_size
puts "#{$unbit_log_base} uidsec_struct allocated."
end
options[:environment] = (ENV['RAILS_ENV'] || "development").dup
options[:processes] = 1
options[:serve_file] = nil
options[:max_input_size] = 8
options[:unbit_debug] = false
options[:master] = true
options[:gc_freq] = nil
require 'optparse'
ARGV.clone.options do |opts|
opts.on('-s', '--socket=socket', String, 'Unix socket path.', 'Default: /tmp/uwsgi.sock') { |v|
options[:sockname] = v
}
opts.on('-C','--chmod','chmod to 666 the unix socket.') {
options[:sock_chmod] = true
}
opts.on('-F','--serve-file','serve static file.') {
options[:serve_file] = true
}
opts.on('-p', '--processes=processes', Integer, 'Number of processes to spawn.', 'Default: 1') {|v|
options[:processes] = v
}
opts.on('-g', '--gc-freq=requests', Integer, 'Number of requests between GC.', 'Default: 1') {|v|
options[:gc_freq] = v
}
opts.on('-d', '--daemon=logfile', 'Put processes in background.', 'Default: all processes stay in foreground') {|v|
options[:stderr_logfile] = v
}
opts.on('-i', '--max-input-size=size', 'Max POST data size (in Kbyte). Bigger data goes into a temporary file.', 'Default: 8') {|v|
options[:max_input_size] = v
}
opts.on('-D', '--unbit-debug', 'Enable debug-level logging.', 'Default: disabled') {|v|
options[:unbit_debug] = true
}
opts.on('-M', '--master-process', 'Enable the master process manager.', 'Default: disabled') {|v|
options[:master] = true
}
opts.on("-e", '--environment=name', String, 'Specifies the environment to run this server under (test/development/production).','Default: development') { |v|
options[:environment] = v
}
opts.separator ""
opts.parse!
end
puts "[#{Time.new}] uRack: loading app [#{RAILS_ROOT}]..."
starttime = Time.now
require RAILS_ROOT + "/config/environment"
options[:app_name] = RAILS_ROOT
require 'socket'
require 'rubygems'
require 'rack'
require 'rack/utils'
require 'rack/content_length'
$requests = 0
$workers = Array.new
module Rack
module Handler
class Unbit
class Tempfile < ::Tempfile
def _close
@tmpfile.close if @tmpfile
@data[1] = nil if @data
@tmpfile = nil
end
end
def self.run(app, options={})
# parse the socket options
# note: on unbit we use the stdin as communication socket
server = nil
master_pid = Process.pid
options[:processes] ||= 1
class RailsCachingHack
def initialize app
@app = app
end
unless options[:unbit]
begin
::File.delete(options[:sockname])
rescue
end
def stripuri(uri)
uri = uri.split('/')
uri.slice!(1)
val = uri.join('/')
end
def call env
env['PATH_INFO'] = stripuri(env['PATH_INFO']) if env['PATH_INFO']
env['REQUEST_URI'] = stripuri(env['REQUEST_URI']) if env['REQUEST_URI']
env['REQUEST_URI'] = '/' if env['REQUEST_URI'] == ''
@app.call env
end
end
server = UNIXServer.new(options[:sockname])
if options[:sock_chmod]
::File.chmod(0666, options[:sockname])
end
module Handler
class Unbit
if options[:stderr_logfile]
cwd = Dir.getwd
Process.daemon
Dir.chdir(cwd)
$stdout.reopen(options[:stderr_logfile],'a')
$stderr.reopen(options[:stderr_logfile],'a')
# log files need to be unbuffered !
$stdout.sync = true
$stderr.sync = true
# pid is changed after the .daemon call
master_pid = Process.pid
end
else
server = UNIXServer.for_fd($stdin.fileno)
end
def self.run(app, options={})
server = UNIXServer.for_fd(0)
$workers = Array.new
can_spawn = true
if options[:master]
$0 = 'uRack master'
$workers[0] = Process.pid
ulog("master process enabled (pid: #{$workers[0]})")
$workers[1] = Process.fork
if $workers[1].to_i > 0
ulog("spawned worker 1 (pid: #{$workers[1]})")
elsif $workers[1].to_i == 0
$0 = "uRack worker 1"
can_spawn = false
end
else
$0 = 'uRack worker 1'
$workers[0] = nil
$workers[1] = Process.pid
ulog("spawned worker 1 (pid: #{$workers[1]})")
end
if options[:master]
$0 = "uRack #{ARGV.last}"
pid = Process.fork
if pid.to_i > 0
puts "[#{Time.new}] uRack: spawned rack worker 1 (pid: #{pid})"
$workers[1] = pid ;
end
end
# the check on master pid is necessary
# without it the first worker will execute this part
if options[:processes] > 1 and Process.pid == master_pid
for p in 2..options[:processes]
pid = Process.fork
break if pid.to_i == 0
puts "[#{Time.new}] uRack: spawned rack worker #{p} (pid: #{pid})"
$workers[p] = pid ;
end
end
# am i the master ?
if options[:master] and Process.pid == master_pid
Signal.trap('TERM') do
puts "[#{Time.new}] uRack: gracefully killing uRack"
for wid in 1..options[:processes]
Process.kill('TERM', $workers[wid]);
end
exit 0;
end
Signal.trap('QUIT') do
puts "[#{Time.new}] uRack: brutally killing uRack"
for wid in 1..options[:processes]
Process.kill('INT', $workers[wid]);
end
exit 0;
end
puts "[#{Time.new}] uRack: the master process manager is alive."
while 1
pid = Process.waitpid
puts "[#{Time.new}] uRack: worker died ! (pid: #{pid})"
newpid = Process.fork
break if newpid.to_i == 0
puts "[#{Time.new}] uRack: respawned rack worker (pid: #{newpid})"
for wid in 1..options[:processes]
if $workers[wid] == pid
$workers[wid] = newpid
break
end
end
end
end
if Process.pid != master_pid
$0 = $0.gsub(/^uRack/,'urack')
Signal.trap("TERM") do
puts "[#{Time.new}] uRack: grecefully killing process #{Process.pid}..."
if $in_request == 0
puts "[#{Time.new}] uRack: goodbye to process #{Process.pid}"
exit 0
end
$manage_next_request = nil
end
Signal.trap('INT') do
puts "[#{Time.new}] uRack: brutally killing process #{Process.pid}..."
exit 0
end
Signal.trap('PIPE') do
puts "[#{Time.new}] uRack: the webserver (or the client) has closed the connection with process #{Process.pid} !!!"
end
end
max_input_size = options[:max_input_size].to_i * 1024
requests = 0
null_post = StringIO.new('')
options[:app_name] ||= Dir.getwd
app = Rack::ContentLength.new(app)
$manage_next_request = true
if options[:gc_freq].to_i > 1
GC.disable
end
while $manage_next_request
$in_request = 0
client = server.accept
$in_request = 1
timed_out = false
ver, size, arg1 = client.recvfrom(4)[0].unpack('CvC')
vars = client.recvfrom(size)[0]
env = Hash.new
i = 0
while i < size
kl = vars[i, 2].unpack('v')[0]
i = i + 2
key = vars[i, kl]
i = i + kl
vl = vars[i, 2].unpack('v')[0]
i = i + 2
value = vars[i, vl]
i = i + vl
env[key] = value
end
rack_input = null_post
env.delete 'CONTENT_LENGTH' if env['CONTENT_LENGTH'] == ''
if env['CONTENT_LENGTH'].to_i > 0
cl = env['CONTENT_LENGTH'].to_i
if cl <= max_input_size
poststr = ''
while poststr.size < cl
poststr = poststr + client.recvfrom(cl-poststr.size)[0] ;
end
rack_input = StringIO.new(poststr)
else
rack_input = Tempfile.new('UnbitRack')
rack_input.chmod(0000)
bytes_read = 0
while bytes_read < cl
buf = client.recvfrom(4096)[0]
rack_input.write(buf)
bytes_read += buf.size
end
rack_input.rewind
end
end
env["SCRIPT_NAME"] = ""
env.update({"rack.version" => [0,1],
"rack.input" => rack_input,
"rack.errors" => $stderr,
"rack.multithread" => false,
"rack.multiprocess" => true,
"rack.run_once" => false,
"rack.url_scheme" => "http"
})
env["QUERY_STRING"] ||= ""
env["HTTP_VERSION"] ||= env["SERVER_PROTOCOL"]
env["REQUEST_PATH"] ||= "/"
begin
#puts env.inspect
$start_of_request = Time.new
status, headers, body = app.call(env)
$speed = Time.new - $start_of_request
client.print "#{env["HTTP_VERSION"]} #{status} #{Rack::Utils::HTTP_STATUS_CODES[status]}\r\n"
headers.each {|k,vs|
vs.split("\n").each { |v|
client.print "#{k}: #{v}\r\n"
}
}
client.print "\r\n"
client.flush
body.each {|part|
client.print part
client.flush
}
rescue Errno::EPIPE
puts "[#{Time.new}] uRack: the webserver (or the client) has closed the connection with process #{Process.pid} !!!"
timed_out = true
ensure
if rack_input.respond_to? :unlink
rack_input.unlink
end
unless rack_input.closed?
rack_input.close
end
client.close
requests = requests+1
# logging
begin
# the syscall 356 is only available on unbit kernels
# other systems need to use the proc file way
# stat = ::File.open('/proc/self/stat', 'r')
# procline = stat.readline
# statd = procline.split /\s+/
# stat.close
$stderr.puts "[#{Time.new}] uRack: [#{options[:app_name]}] req: #{requests} ip: #{env['REMOTE_ADDR']} pid: #{Process.pid} as: #{human_round(syscall(356).to_f/1024/1024)} MB => #{env['REQUEST_METHOD']} #{env['REQUEST_URI']} in #{$speed} secs [#{status}]#{' TIMED OUT !!!' if timed_out}"
rescue
$stderr.puts "[#{Time.new}] uRack: [#{options[:app_name]}] req: #{requests} unable to get /proc/self/stat or syscall 356"
end
if syscall(357, $uidsec, 0) == $uidsec_size
if $uidsec[120..123].unpack('i')[0] > 0
$stderr.puts "[#{Time.new}] uRack: found a memory allocation error for request #{requests} (pid: #{Process.pid}). Better to kill myself..."
$manage_next_request = nil
end
end
$uidsec = '1' * $uidsec_size
if options[:unbit_debug]
current_as = human_round( (syscall(356).to_f/1024/1024) - $last_as )
$stderr.puts "#{$unbit_log_base} resource status after request #{requests}: AS for this request: #{current_as} MB | OBJ for this request: #{ObjectSpace.each_object {} - $last_obj} | OBJ total: #{ObjectSpace.each_object {}}"
$last_as = human_round(syscall(356).to_f/1024/1024) ;
$last_obj = ObjectSpace.each_object {}
end
if options[:gc_freq].to_i > 1
if requests % options[:gc_freq].to_i == 0
$stderr.puts "[#{Time.new}] uRack: [#{options[:app_name]}] calling GC for pid #{Process.pid} after #{requests} requests."
GC.enable ; GC.start ; GC.disable
end
end
end
end
puts "[#{Time.new}] uRack: goodbye to process #{Process.pid}"
if options[:processes] > 1 and can_spawn
for p in 2..options[:processes]
$workers[p] = Process.fork
if $workers[p].to_i == 0
$0 = "uRack worker #{p}"
break
elsif $workers[p].to_i > 0
ulog("spawned worker #{p} (pid: #{$workers[p]})")
end
end
end
# am i the master ?
if options[:master] and Process.pid == $workers[0]
while 1
i_am_a_child = false
pid = Process.waitpid
for wid in 1..options[:processes]
if $workers[wid] == pid
ulog("worker #{wid} died ! (pid: #{pid})")
$workers[wid] = Process.fork
if $workers[wid].to_i == 0
$0 = "uRack worker #{wid}"
i_am_a_child = true
break
elsif $workers[wid].to_i > 0
ulog("respawned worker #{wid} (pid: #{$workers[wid]})")
end
end
end
break if i_am_a_child
end
end
while client = server.accept
serve client, app, options
end
end
def self.serve(client, app, options)
speed = Time.now
head, sender = client.recvfrom(4)
unless head
client.close
return
end
mod1, size, mod2 = head.unpack('CvC')
if size == 0 or size.nil?
client.close
return
end
vars, sender = client.recvfrom(size)
if vars.length != size
client.close
return
end
env = Hash.new
i = 0
while i < size
kl = vars[i, 2].unpack('v')[0]
i = i + 2
key = vars[i, kl]
i = i + kl
vl = vars[i, 2].unpack('v')[0]
i = i + 2
value = vars[i, vl]
i = i + vl
env[key] = value
end
env.delete "HTTP_CONTENT_LENGTH"
env.delete "HTTP_CONTENT_TYPE"
env["SCRIPT_NAME"] = "" if env["SCRIPT_NAME"] == "/"
env["QUERY_STRING"] ||= ""
env["HTTP_VERSION"] ||= env["SERVER_PROTOCOL"]
env["REQUEST_PATH"] ||= "/"
env.delete "PATH_INFO" if env["PATH_INFO"] == ""
env.delete "CONTENT_TYPE" if env["CONTENT_TYPE"] == ""
env.delete "CONTENT_LENGTH" if env["CONTENT_LENGTH"] == ""
if env["CONTENT_LENGTH"].to_i > 4096
rack_input = Rack::RewindableInput::Tempfile.new('Rack_unbit_Input')
rack_input.chmod(0000)
rack_input.set_encoding(Encoding::BINARY) if rack_input.respond_to?(:set_encoding)
rack_input.binmode
rack_input.unlink
remains = env["CONTENT_LENGTH"].to_i
while remains > 0
if remains >= 4096
buf, sender = client.recvfrom(4096)
else
buf, sender = client.recvfrom(remains)
end
rack_input.write( buf )
remains -= buf.length
end
elsif env["CONTENT_LENGTH"].to_i > 0
rack_input = StringIO.new(client.recvfrom(env["CONTENT_LENGTH"].to_i)[0])
else
rack_input = StringIO.new('')
end
rack_input.rewind
env.update({"rack.version" => [1,1],
"rack.input" => rack_input,
"rack.errors" => $stderr,
"rack.multithread" => false,
"rack.multiprocess" => true,
"rack.run_once" => false,
"rack.url_scheme" => ["yes", "on", "1"].include?(env["HTTPS"]) ? "https" : "http"
})
app = Rack::ContentLength.new(app)
disconnected = false
begin
status, headers, body = app.call(env)
begin
send_headers client, env["HTTP_VERSION"] ,status, headers
send_body client, body
ensure
body.close if body.respond_to? :close
end
rescue Errno::EPIPE, Errno::ECONNRESET
disconnected = true
ensure
rack_input.close
client.close
end
$requests = $requests+1
ulog("req: #{$requests} ip: #{env['REMOTE_ADDR']} pid: #{Process.pid} as: #{human_round(syscall(356).to_f/1024/1024)} MB => #{env['REQUEST_METHOD']} #{env['REQUEST_URI']} in #{Time.now-speed} secs [#{status}]#{' DISCONNECTED !!!' if disconnected}") if options[:logging]
end
def self.send_headers(client, protocol, status, headers)
client.print "#{protocol} #{status} #{Rack::Utils::HTTP_STATUS_CODES[status]}\r\n"
headers.each { |k, vs|
vs.split("\n").each { |v|
client.print "#{k}: #{v}\r\n"
}
}
client.print "\r\n"
client.flush
end
def self.send_body(client, body)
body.each { |part|
client.print part
client.flush
}
end
end
end
end
def ulog(message)
$stderr.puts "[#{Time.new}] uRack: #{message}"
end
ulog("starting at #{Dir.getwd}")
limit_as = human_round(Process.getrlimit(Process::RLIMIT_AS)[1].to_f/1024/1024)
ulog("your process address space limit is #{limit_as} MB")
server = Rack::Handler.get('Unbit')
app = Rack::Builder.new {
if options[:serve_file]
puts "[#{Time.new}] uRack: file serving enabled."
use Rails::Rack::Static
starttime = Time.now
app = nil
if options[:config]
cfgfile = File.read(options[:config])
if cfgfile[/^#\\(.*)/]
opts.parse! $1.split(/\s+/)
end
if ActionController.const_defined?(:Dispatcher) && (ActionController::Dispatcher.instance_methods.include?(:call) || ActionController::Dispatcher.instance_methods.include?("call"))
run ActionController::Dispatcher.new
app = eval "Rack::Builder.new {( " + cfgfile + "\n )}.to_app", nil, options[:config]
else
# falling back to rails
if Dir.getwd =~ /\/public\/?$/
require '../config/environment'
else
require 'thin'
run Rack::Adapter::Rails.new(:environment => ENV['RAILS_ENV'])
require 'config/environment'
end
}.to_app
if options[:unbit]
$after_spawn_used_as = human_round(syscall(356).to_f/1024/1024)
$stderr.puts "#{$unbit_log_base} now you have #{$limit_as-$after_spawn_used_as} MB of address space available (used #{$after_spawn_used_as}MB after app startup)"
app = Rack::Builder.new {
if options[:cachehack]
use Rack::RailsCachingHack
end
if options[:static]
use Rails::Rack::Static
end
if ActionController.const_defined?(:Dispatcher) && (ActionController::Dispatcher.instance_methods.include?(:call) || ActionController::Dispatcher.instance_methods.include?("call"))
run ActionController::Dispatcher.new
else
require 'thin'
run Rack::Adapter::Rails.new(:environment => ENV['RAILS_ENV'])
end
}
end
secs = Time.now-starttime
puts "[#{Time.new}] uRack: ready to serve requests after #{secs.to_i} secs (pid: #{Process.pid})"
if options[:unbit_debug]
$last_as = $after_spawn_used_as ;
$last_obj = ObjectSpace.each_object {}
end
ulog("your app is ready (in #{Time.now-starttime} seconds).")
server.run(app, options)
+6 -1
View File
@@ -1,13 +1,18 @@
require 'socket'
require 'rack/content_length'
require 'rack/rewindable_input'
require 'stringio'
module Rack
module Handler
class Uwsgi
def self.run(app, options={})
server = TCPServer.new(options[:Host], options[:Port])
if ENV['UWSGI_FD']
server = UNIXServer.for_fd(ENV['UWSGI_FD'].to_i)
else
server = TCPServer.new(options[:Host], options[:Port])
end
while client = server.accept
serve client, app
end
+3 -3
View File
@@ -1,6 +1,6 @@
Package: uwsgi
Version: 0.9.2
Version: 0.9.5
Maintainer: Unbit
Description: Fast, developer-friendly wsgi server
Architecture: i386
Depends: python2.5, libxml2
Architecture: all
Depends: python, libxml2
+23 -23
View File
@@ -47,7 +47,7 @@ PyObject *py_erlang_recv_message(PyObject * self, PyObject * args) {
if (timeout > 0) {
eret = poll(&erpoll, 1, timeout * 1000);
if (eret < 0) {
perror("poll()");
uwsgi_error("poll()");
goto clear;
}
else if (eret == 0) {
@@ -318,7 +318,7 @@ int init_erlang(char *nodename, char *cookie) {
ip = strchr(nodename, '@');
if (ip == NULL) {
fprintf(stderr, "*** invalid erlang node name ***\n");
uwsgi_log( "*** invalid erlang node name ***\n");
return -1;
}
@@ -326,12 +326,12 @@ int init_erlang(char *nodename, char *cookie) {
// get the cookie from the home
cookiehome = getenv("HOME");
if (!cookiehome) {
fprintf(stderr, "unable to get erlang cookie from your home.\n");
uwsgi_log( "unable to get erlang cookie from your home.\n");
return -1;
}
cookiefile = malloc(strlen(cookiehome) + 1 + strlen(".erlang.cookie") + 1);
if (!cookiefile) {
perror("malloc()");
uwsgi_error("malloc()");
}
cookiefile[0] = 0;
strcat(cookiefile, cookiehome);
@@ -339,14 +339,14 @@ int init_erlang(char *nodename, char *cookie) {
cookiefd = open(cookiefile, O_RDONLY);
if (cookiefd < 0) {
perror("open()");
uwsgi_error("open()");
free(cookiefile);
return -1;
}
memset(cookievalue, 0, 128);
if (read(cookiefd, cookievalue, 127) < 1) {
fprintf(stderr, "invalid cookie found in %s\n", cookiefile);
uwsgi_log( "invalid cookie found in %s\n", cookiefile);
close(cookiefd);
free(cookiefile);
return -1;
@@ -358,7 +358,7 @@ int init_erlang(char *nodename, char *cookie) {
node = malloc((ip - nodename) + 1);
if (node == NULL) {
perror("malloc()");
uwsgi_error("malloc()");
return -1;
}
memset(node, 0, (ip - nodename) + 1);
@@ -367,13 +367,13 @@ int init_erlang(char *nodename, char *cookie) {
erl_init(NULL, 0);
if (erl_connect_xinit(ip + 1, node, nodename, NULL, cookie, 0) == -1) {
fprintf(stderr, "*** unable to initialize erlang c-node ***\n");
uwsgi_log( "*** unable to initialize erlang c-node ***\n");
return -1;
}
efd = socket(AF_INET, SOCK_STREAM, 0);
if (efd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1;
}
@@ -384,37 +384,37 @@ int init_erlang(char *nodename, char *cookie) {
rlen = 1;
if (setsockopt(efd, SOL_SOCKET, SO_REUSEADDR, &rlen, sizeof(rlen))) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
close(efd);
return -1;
}
if (bind(efd, (struct sockaddr *) &e_addr, sizeof(struct sockaddr_in)) < 0) {
perror("bind()");
uwsgi_error("bind()");
close(efd);
return -1;
}
rlen = sizeof(struct sockaddr_in);
if (getsockname(efd, (struct sockaddr *) &e_addr, (socklen_t *) & rlen)) {
perror("getsockname()");
uwsgi_error("getsockname()");
close(efd);
return -1;
}
if (listen(efd, uwsgi.listen_queue)) {
perror("listen()");
uwsgi_error("listen()");
close(efd);
return -1;
}
if (erl_publish(ntohs(e_addr.sin_port)) < 0) {
fprintf(stderr, "*** unable to subscribe with EPMD ***\n");
uwsgi_log( "*** unable to subscribe with EPMD ***\n");
close(efd);
return -1;
}
fprintf(stderr, "Erlang C-Node initialized on port %d you can access it with name %s\n", ntohs(e_addr.sin_port), nodename);
uwsgi_log( "Erlang C-Node initialized on port %d you can access it with name %s\n", ntohs(e_addr.sin_port), nodename);
for (uwsgi_function = uwsgi_erlang_methods; uwsgi_function->ml_name != NULL; uwsgi_function++) {
PyObject *func = PyCFunction_New(uwsgi_function, NULL);
@@ -512,7 +512,7 @@ ETERM *py_to_eterm(PyObject * pobj) {
free(eobj3);
}
else {
fprintf(stderr, "UNMANAGED PYTHON TYPE: %s\n", pobj->ob_type->tp_name);
uwsgi_log( "UNMANAGED PYTHON TYPE: %s\n", pobj->ob_type->tp_name);
}
clear:
@@ -560,7 +560,7 @@ PyObject *eterm_to_py(ETERM * obj) {
eobj = PyInt_FromLong(ERL_INT_VALUE(obj));
break;
case ERL_BINARY:
fprintf(stderr, "FOUND A BINARY %.*s\n", ERL_BIN_SIZE(obj), ERL_BIN_PTR(obj));
uwsgi_log( "FOUND A BINARY %.*s\n", ERL_BIN_SIZE(obj), ERL_BIN_PTR(obj));
break;
case ERL_PID:
eobj = PyDict_New();
@@ -581,7 +581,7 @@ PyObject *eterm_to_py(ETERM * obj) {
break;
}
default:
fprintf(stderr, "UNMANAGED ETERM TYPE: %d\n", ERL_TYPE(obj));
uwsgi_log( "UNMANAGED ETERM TYPE: %d\n", ERL_TYPE(obj));
break;
}
@@ -603,13 +603,13 @@ void erlang_loop(struct wsgi_request *wsgi_req) {
PyObject *callable = PyDict_GetItemString(uwsgi.embedded_dict, "erlang_func");
if (!callable) {
PyErr_Print();
fprintf(stderr, "- you have not defined a uwsgi.erlang_func callable, Erlang message manager will be disabled until you define it -\n");
uwsgi_log( "- you have not defined a uwsgi.erlang_func callable, Erlang message manager will be disabled until you define it -\n");
}
PyObject *pargs = PyTuple_New(1);
if (!pargs) {
PyErr_Print();
fprintf(stderr, "- error preparing arg tuple for uwsgi.erlang_func callable, Erlang message manager will be disabled -\n");
uwsgi_log( "- error preparing arg tuple for uwsgi.erlang_func callable, Erlang message manager will be disabled -\n");
}
while (uwsgi.workers[uwsgi.mywid].manage_next_request) {
@@ -623,7 +623,7 @@ void erlang_loop(struct wsgi_request *wsgi_req) {
UWSGI_SET_ERLANGING;
for (;;) {
if (erl_receive_msg(wsgi_req->poll.fd, (unsigned char *) &wsgi_req->buffer, uwsgi.buffer_size, &em) == ERL_MSG) {
if (erl_receive_msg(wsgi_req->poll.fd, (unsigned char *) wsgi_req->buffer, uwsgi.buffer_size, &em) == ERL_MSG) {
if (em.type == ERL_TICK)
continue;
@@ -632,7 +632,7 @@ void erlang_loop(struct wsgi_request *wsgi_req) {
}
if (!callable) {
fprintf(stderr, "- you still have not defined a uwsgi.erlang_func callable, Erlang message rejected -\n");
uwsgi_log( "- you still have not defined a uwsgi.erlang_func callable, Erlang message rejected -\n");
}
PyObject *zero = eterm_to_py(em.msg);
@@ -694,7 +694,7 @@ static void erlang_log() {
uwsgi.workers[uwsgi.mywid].rss_size = 0;
uwsgi.workers[uwsgi.mywid].vsz_size = 0;
}
fprintf(stderr, "[Erlang worker %d pid %d] request %llu done {rss: %llu vsz: %llu}\n", uwsgi.mywid, uwsgi.mypid, uwsgi.workers[uwsgi.mywid].requests, uwsgi.workers[uwsgi.mywid].rss_size, uwsgi.workers[uwsgi.mywid].vsz_size);
uwsgi_log( "[Erlang worker %d pid %d] request %llu done {rss: %llu vsz: %llu}\n", uwsgi.mywid, uwsgi.mypid, uwsgi.workers[uwsgi.mywid].requests, uwsgi.workers[uwsgi.mywid].rss_size, uwsgi.workers[uwsgi.mywid].vsz_size);
}
#else
+60 -18
View File
@@ -1,23 +1,39 @@
#include "uwsgi.h"
#if defined(__FreeBSD__) || defined(__NetBSD__) || defined(__DragonFly__) || defined(__sun__) || defined(__OpenBSD__)
#if defined(__FreeBSD__) || defined(__NetBSD__) || defined(__DragonFly__) || defined(__OpenBSD__)
#include <kvm.h>
#include <sys/sysctl.h>
#include <sys/user.h>
#elif defined(__sun__)
/* Terrible Hack !!! */
#ifndef _LP64
#undef _FILE_OFFSET_BITS
#endif
#include <procfs.h>
#define _FILE_OFFSET_BITS 64
#endif
#if defined(__NetBSD__) || defined(__FreeBSD__) || defined(__DragonFly__)
#include <sys/sysctl.h>
#endif
extern struct uwsgi_server uwsgi;
void log_request(struct wsgi_request *wsgi_req) {
// optimize this (please)
char *time_request;
time_t microseconds, microseconds2;
static char *empty = "";
char *first_part = empty;
int rlen;
int app_req = -1;
char *msg2 = " ";
char *via = msg2;
char mempkt[4096];
char logpkt[4096];
struct iovec logvec[2] ;
int logvecpos = 0 ;
#ifdef UWSGI_SENDFILE
char *msg1 = " via sendfile() ";
@@ -43,20 +59,23 @@ void log_request(struct wsgi_request *wsgi_req) {
if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1) {
#ifndef UNBIT
if (uwsgi.synclog) {
snprintf(uwsgi.sync_page, uwsgi.page_size, "{address space usage: %lld bytes/%lluMB} {rss usage: %llu bytes/%lluMB} ", uwsgi.workers[uwsgi.mywid].vsz_size, uwsgi.workers[uwsgi.mywid].vsz_size / 1024 / 1024, uwsgi.workers[uwsgi.mywid].rss_size, uwsgi.workers[uwsgi.mywid].rss_size / 1024 / 1024);
first_part = uwsgi.sync_page;
}
else {
fprintf(stderr, "{address space usage: %lld bytes/%lluMB} {rss usage: %llu bytes/%lluMB} ", uwsgi.workers[uwsgi.mywid].vsz_size, uwsgi.workers[uwsgi.mywid].vsz_size / 1024 / 1024, uwsgi.workers[uwsgi.mywid].rss_size, uwsgi.workers[uwsgi.mywid].rss_size / 1024 / 1024);
}
rlen = snprintf(mempkt, 4096, "{address space usage: %lld bytes/%lluMB} {rss usage: %llu bytes/%lluMB} ",
uwsgi.workers[uwsgi.mywid].vsz_size, uwsgi.workers[uwsgi.mywid].vsz_size / 1024 / 1024,
uwsgi.workers[uwsgi.mywid].rss_size, uwsgi.workers[uwsgi.mywid].rss_size / 1024 / 1024);
#else
fprintf(stderr, "{address space usage: %lld bytes/%lluMB} ", uwsgi.workers[uwsgi.mywid].vsz_size, uwsgi.workers[uwsgi.mywid].vsz_size / 1024 / 1024);
rlen = snprintf(mempkt, 4096, "{address space usage: %lld bytes/%lluMB} ",
uwsgi.workers[uwsgi.mywid].vsz_size, uwsgi.workers[uwsgi.mywid].vsz_size / 1024 / 1024);
#endif
logvec[logvecpos].iov_base = mempkt ;
logvec[logvecpos].iov_len = rlen ;
logvecpos++;
}
fprintf(stderr, "%s[pid: %d|app: %d|req: %d/%llu] %.*s (%.*s) {%d vars in %d bytes} [%.*s] %.*s %.*s => generated %llu bytes in %ld msecs%s(%.*s %d) %d headers in %d bytes (%d async switches on async core %d)\n",
first_part,
rlen = snprintf(logpkt, 4096, "[pid: %d|app: %d|req: %d/%llu] %.*s (%.*s) {%d vars in %d bytes} [%.*s] %.*s %.*s => generated %llu bytes in %ld msecs%s(%.*s %d) %d headers in %d bytes (%d async switches on async core %d)\n",
uwsgi.mypid,
wsgi_req->app_id,
app_req,
@@ -77,6 +96,11 @@ void log_request(struct wsgi_request *wsgi_req) {
wsgi_req->headers_size,
wsgi_req->async_switches, wsgi_req->async_id);
logvec[logvecpos].iov_base = logpkt ;
logvec[logvecpos].iov_len = rlen ;
// do not check for errors
rlen = writev(2, logvec, logvecpos+1);
}
@@ -91,11 +115,24 @@ void get_memusage() {
if (procfile) {
i = fscanf(procfile, "%*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %llu %lld", &uwsgi.workers[uwsgi.mywid].vsz_size, &uwsgi.workers[uwsgi.mywid].rss_size);
if (i != 2) {
fprintf(stderr, "warning: invalid record in /proc/self/stat\n");
uwsgi_log( "warning: invalid record in /proc/self/stat\n");
}
fclose(procfile);
}
uwsgi.workers[uwsgi.mywid].rss_size = uwsgi.workers[uwsgi.mywid].rss_size * uwsgi.page_size;
#elif defined (__sun__)
psinfo_t info;
int procfd ;
procfd = open("/proc/self/psinfo", O_RDONLY);
if (procfd >= 0) {
if ( read(procfd, (char *) &info, sizeof(info)) > 0) {
uwsgi.workers[uwsgi.mywid].rss_size = (uint64_t) info.pr_rssize * 1024 ;
uwsgi.workers[uwsgi.mywid].vsz_size = (uint64_t) info.pr_size * 1024 ;
}
close(procfd);
}
#elif defined( __APPLE__)
/* darwin documentation says that the value are in pages, but they are bytes !!! */
struct task_basic_info t_info;
@@ -105,25 +142,30 @@ void get_memusage() {
uwsgi.workers[uwsgi.mywid].rss_size = t_info.resident_size;
uwsgi.workers[uwsgi.mywid].vsz_size = t_info.virtual_size;
}
#elif defined(__FreeBSD__) || defined(__NetBSD__) || defined(__DragonFly__) || defined(__sun__) || defined(__OpenBSD__)
#elif defined(__FreeBSD__) || defined(__NetBSD__) || defined(__DragonFly__) || defined(__OpenBSD__)
kvm_t *kv;
int cnt;
kv = kvm_open(NULL, NULL, NULL, O_RDONLY, NULL);
if (kv) {
#if defined(__FreeBSD__) || defined(__DragonFly__)
struct kinfo_proc *kproc;
kproc = kvm_getprocs(kv, KERN_PROC_PID, uwsgi.mypid, &cnt);
if (kproc && cnt > 0) {
uwsgi.workers[uwsgi.mywid].vsz_size = kproc->ki_size;
uwsgi.workers[uwsgi.mywid].rss_size = kproc->ki_rssize * uwsgi.page_size;
}
#elif defined(__NetBSD__)
#elif defined(__NetBSD__) || defined(__OpenBSD__)
struct kinfo_proc2 *kproc2;
kproc2 = kvm_getproc2(kv, KERN_PROC_PID, uwsgi.mypid, sizeof(struct kinfo_proc2), &cnt);
if (kproc2 && cnt > 0) {
#ifdef __OpenBSD__
uwsgi.workers[uwsgi.mywid].vsz_size = (kproc2->p_vm_dsize + kproc2->p_vm_ssize + kproc2->p_vm_tsize) * uwsgi.page_size;
#else
uwsgi.workers[uwsgi.mywid].vsz_size = kproc2->p_vm_msize * uwsgi.page_size;
#endif
uwsgi.workers[uwsgi.mywid].rss_size = kproc2->p_vm_rssize * uwsgi.page_size;
}
#endif
+3 -3
View File
@@ -25,18 +25,18 @@ void nagios(struct uwsgi_server *uwsgi) {
uwsgi->wsgi_req->uh.pktsize = 0;
uwsgi->wsgi_req->uh.modifier2 = 0;
if (write(nagios_poll.fd, uwsgi->wsgi_req, 4) != 4) {
perror("write()");
uwsgi_error("write()");
fprintf(stdout, "UWSGI CRITICAL: could not send ping packet to workers\n");
exit(2);
}
nagios_poll.events = POLLIN;
if (!uwsgi_parse_response(&nagios_poll, uwsgi->shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) uwsgi->wsgi_req, &uwsgi->wsgi_req->buffer)) {
if (!uwsgi_parse_response(&nagios_poll, uwsgi->shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) uwsgi->wsgi_req, uwsgi->wsgi_req->buffer)) {
fprintf(stdout, "UWSGI CRITICAL: timed out waiting for response\n");
exit(2);
}
else {
if (uwsgi->wsgi_req->uh.pktsize > 0) {
fprintf(stdout, "UWSGI WARNING: %.*s\n", uwsgi->wsgi_req->uh.pktsize, &uwsgi->wsgi_req->buffer);
fprintf(stdout, "UWSGI WARNING: %.*s\n", uwsgi->wsgi_req->uh.pktsize, uwsgi->wsgi_req->buffer);
exit(1);
}
else {
+31 -30
View File
@@ -23,7 +23,7 @@ ssize_t send_udp_message(uint8_t modifier1, char *host, char *message, uint16_t
fd = socket(AF_INET, SOCK_DGRAM, 0);
if (fd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1 ;
}
@@ -49,7 +49,7 @@ ssize_t send_udp_message(uint8_t modifier1, char *host, char *message, uint16_t
ret = sendto(fd, udpbuff, message_size+4, 0, (struct sockaddr *) &udp_addr, sizeof(udp_addr));
if (ret < 0) {
perror("sendto()");
uwsgi_error("sendto()");
}
close(fd);
@@ -69,13 +69,13 @@ int uwsgi_enqueue_message(char *host, int port, uint8_t modifier1, uint8_t modif
timeout = 1;
if (size > 0xFFFF) {
fprintf(stderr, "invalid object (marshalled) size\n");
uwsgi_log( "invalid object (marshalled) size\n");
return -1;
}
uwsgi_poll.fd = socket(AF_INET, SOCK_STREAM, 0);
if (uwsgi_poll.fd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1;
}
@@ -87,7 +87,7 @@ int uwsgi_enqueue_message(char *host, int port, uint8_t modifier1, uint8_t modif
uwsgi_poll.events = POLLIN;
if (timed_connect(&uwsgi_poll, (const struct sockaddr *) &uws_addr, sizeof(struct sockaddr_in), timeout)) {
perror("connect()");
uwsgi_error("connect()");
close(uwsgi_poll.fd);
return -1;
}
@@ -98,14 +98,14 @@ int uwsgi_enqueue_message(char *host, int port, uint8_t modifier1, uint8_t modif
cnt = write(uwsgi_poll.fd, &uh, 4);
if (cnt != 4) {
perror("write()");
uwsgi_error("write()");
close(uwsgi_poll.fd);
return -1;
}
cnt = write(uwsgi_poll.fd, message, size);
if (cnt != size) {
perror("write()");
uwsgi_error("write()");
close(uwsgi_poll.fd);
return -1;
}
@@ -127,7 +127,7 @@ PyObject *uwsgi_send_message(const char *host, int port, uint8_t modifier1, uint
timeout = 1;
if (size > 0xFFFF) {
fprintf(stderr, "invalid object (marshalled) size\n");
uwsgi_log( "invalid object (marshalled) size\n");
Py_INCREF(Py_None);
return Py_None;
}
@@ -136,7 +136,7 @@ PyObject *uwsgi_send_message(const char *host, int port, uint8_t modifier1, uint
uwsgi_mpoll.fd = socket(AF_INET, SOCK_STREAM, 0);
if (uwsgi_mpoll.fd < 0) {
perror("socket()");
uwsgi_error("socket()");
Py_INCREF(Py_None);
return Py_None;
}
@@ -149,7 +149,7 @@ PyObject *uwsgi_send_message(const char *host, int port, uint8_t modifier1, uint
UWSGI_SET_BLOCKING;
if (timed_connect(&uwsgi_mpoll, (const struct sockaddr *) &uws_addr, sizeof(struct sockaddr_in), timeout)) {
perror("connect()");
uwsgi_error("connect()");
close(uwsgi_mpoll.fd);
Py_INCREF(Py_None);
return Py_None;
@@ -161,7 +161,7 @@ PyObject *uwsgi_send_message(const char *host, int port, uint8_t modifier1, uint
cnt = write(uwsgi_mpoll.fd, &uh, 4);
if (cnt != 4) {
perror("write()");
uwsgi_error("write()");
close(uwsgi_mpoll.fd);
Py_INCREF(Py_None);
return Py_None;
@@ -169,7 +169,7 @@ PyObject *uwsgi_send_message(const char *host, int port, uint8_t modifier1, uint
cnt = write(uwsgi_mpoll.fd, message, size);
if (cnt != size) {
perror("write()");
uwsgi_error("write()");
close(uwsgi_mpoll.fd);
Py_INCREF(Py_None);
return Py_None;
@@ -208,11 +208,11 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header
/* first 4 byte header */
rlen = poll(upoll, 1, timeout * 1000);
if (rlen < 0) {
perror("poll()");
uwsgi_error("poll()");
exit(1);
}
else if (rlen == 0) {
fprintf(stderr, "timeout. skip request\n");
uwsgi_log( "timeout. skip request\n");
close(upoll->fd);
return 0;
}
@@ -222,17 +222,17 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header
while (i < 4) {
rlen = poll(upoll, 1, timeout * 1000);
if (rlen < 0) {
perror("poll()");
uwsgi_error("poll()");
exit(1);
}
else if (rlen == 0) {
fprintf(stderr, "timeout waiting for header. skip request.\n");
uwsgi_log( "timeout waiting for header. skip request.\n");
close(upoll->fd);
break;
}
rlen = read(upoll->fd, (char *) (uh) + i, 4 - i);
if (rlen <= 0) {
fprintf(stderr, "broken header. skip request.\n");
uwsgi_log( "broken header. skip request.\n");
close(upoll->fd);
break;
}
@@ -243,7 +243,7 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header
}
}
else if (rlen <= 0) {
fprintf(stderr, "invalid request header size: %d...skip\n", rlen);
uwsgi_log( "invalid request header size: %d...skip\n", rlen);
close(upoll->fd);
return 0;
}
@@ -254,28 +254,28 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header
/* check for max buffer size */
if (uh->pktsize > uwsgi.buffer_size) {
fprintf(stderr, "invalid request block size: %d...skip\n", uh->pktsize);
uwsgi_log( "invalid request block size: %d...skip\n", uh->pktsize);
close(upoll->fd);
return 0;
}
//fprintf(stderr,"ready for reading %d bytes\n", wsgi_req.size);
//uwsgi_log("ready for reading %d bytes\n", wsgi_req.size);
i = 0;
while (i < uh->pktsize) {
rlen = poll(upoll, 1, timeout * 1000);
if (rlen < 0) {
perror("poll()");
uwsgi_error("poll()");
exit(1);
}
else if (rlen == 0) {
fprintf(stderr, "timeout. skip request. (expecting %d bytes, got %d)\n", uh->pktsize, i);
uwsgi_log( "timeout. skip request. (expecting %d bytes, got %d)\n", uh->pktsize, i);
close(upoll->fd);
break;
}
rlen = read(upoll->fd, buffer + i, uh->pktsize - i);
if (rlen <= 0) {
fprintf(stderr, "broken vars. skip request.\n");
uwsgi_log( "broken vars. skip request.\n");
close(upoll->fd);
break;
}
@@ -292,7 +292,7 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header
int uwsgi_parse_vars(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req) {
char *buffer = &wsgi_req->buffer;
char *buffer = wsgi_req->buffer;
char *ptrbuf, *bufferend;
@@ -312,7 +312,7 @@ int uwsgi_parse_vars(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req)
#endif
/* key cannot be null */
if (!strsize) {
fprintf(stderr, "uwsgi key cannot be null. skip this request.\n");
uwsgi_log( "uwsgi key cannot be null. skip this request.\n");
return -1;
}
@@ -386,17 +386,18 @@ int uwsgi_parse_vars(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req)
wsgi_req->var_cnt++;
}
else {
fprintf(stderr, "max vec size reached. skip this header.\n");
uwsgi_log( "max vec size reached. skip this header.\n");
return -1;
}
// var value
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptrbuf;
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = strsize;
//uwsgi_log("%.*s = %.*s\n", wsgi_req->hvec[wsgi_req->var_cnt-1].iov_len, wsgi_req->hvec[wsgi_req->var_cnt-1].iov_base, wsgi_req->hvec[wsgi_req->var_cnt].iov_len, wsgi_req->hvec[wsgi_req->var_cnt].iov_base);
if (wsgi_req->var_cnt < uwsgi->vec_size - (4 + 1)) {
wsgi_req->var_cnt++;
}
else {
fprintf(stderr, "max vec size reached. skip this header.\n");
uwsgi_log( "max vec size reached. skip this header.\n");
return -1;
}
ptrbuf += strsize;
@@ -436,7 +437,7 @@ int uwsgi_ping_node(int node, struct wsgi_request *wsgi_req) {
uwsgi_poll.fd = socket(AF_INET, SOCK_STREAM, 0);
if (uwsgi_poll.fd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1;
}
@@ -449,12 +450,12 @@ int uwsgi_ping_node(int node, struct wsgi_request *wsgi_req) {
wsgi_req->uh.pktsize = 0;
wsgi_req->uh.modifier2 = 0;
if (write(uwsgi_poll.fd, wsgi_req, 4) != 4) {
perror("write()");
uwsgi_error("write()");
return -1;
}
uwsgi_poll.events = POLLIN;
if (!uwsgi_parse_response(&uwsgi_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, &wsgi_req->buffer)) {
if (!uwsgi_parse_response(&uwsgi_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer)) {
return -1;
}
+25 -25
View File
@@ -139,14 +139,14 @@ void uwsgi_proxy(int proxyfd) {
int next_node = -1;
fprintf(stderr, "spawned uWSGI proxy (pid: %d)\n", getpid());
uwsgi_log( "spawned uWSGI proxy (pid: %d)\n", getpid());
fprintf(stderr, "allocating space for %d concurrent proxy connections\n", max_connections);
uwsgi_log( "allocating space for %d concurrent proxy connections\n", max_connections);
// allocate memory for connections
upcs = malloc(sizeof(struct uwsgi_proxy_connection) * max_connections);
if (!upcs) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
memset(upcs, 0, sizeof(struct uwsgi_proxy_connection) * max_connections);
@@ -169,7 +169,7 @@ void uwsgi_proxy(int proxyfd) {
#endif
if (!eevents) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
@@ -182,7 +182,7 @@ void uwsgi_proxy(int proxyfd) {
nevents = async_wait(efd, eevents, max_events, -1, 0);
if (nevents < 0) {
perror("epoll_wait()");
uwsgi_error("epoll_wait()");
continue;
}
@@ -195,7 +195,7 @@ void uwsgi_proxy(int proxyfd) {
// new connection, accept it
ev.ASYNC_FD = accept(proxyfd, (struct sockaddr *) &upc_addr, &upc_len);
if (ev.ASYNC_FD < 0) {
perror("accept()");
uwsgi_error("accept()");
continue;
}
upcs[ev.ASYNC_FD].node = -1;
@@ -204,7 +204,7 @@ void uwsgi_proxy(int proxyfd) {
upcs[ev.ASYNC_FD].dest_fd = socket(AF_INET, SOCK_STREAM, 0);
if (upcs[ev.ASYNC_FD].dest_fd < 0) {
perror("socket()");
uwsgi_error("socket()");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
@@ -212,7 +212,7 @@ void uwsgi_proxy(int proxyfd) {
// set nonblocking
if (ioctl(upcs[ev.ASYNC_FD].dest_fd, FIONBIO, &nonblocking)) {
perror("ioctl()");
uwsgi_error("ioctl()");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
@@ -221,7 +221,7 @@ void uwsgi_proxy(int proxyfd) {
upcs[ev.ASYNC_FD].retry = 0;
next_node = uwsgi_proxy_find_next_node(next_node);
if (next_node == -1) {
fprintf(stderr, "unable to find an available worker in the cluster !\n");
uwsgi_log( "unable to find an available worker in the cluster !\n");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
@@ -250,7 +250,7 @@ void uwsgi_proxy(int proxyfd) {
// re-set blocking
if (ioctl(upcs[upcs[ev.ASYNC_FD].dest_fd].dest_fd, FIONBIO, &blocking)) {
perror("ioctl()");
uwsgi_error("ioctl()");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
@@ -271,7 +271,7 @@ void uwsgi_proxy(int proxyfd) {
}
else {
// connection failed, retry with the next node ?
perror("connect()");
uwsgi_error("connect()");
// close only when all node are tried
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
@@ -280,7 +280,7 @@ void uwsgi_proxy(int proxyfd) {
}
else {
fprintf(stderr, "!!! something horrible happened to the uWSGI proxy, reloading it !!!\n");
uwsgi_log( "!!! something horrible happened to the uWSGI proxy, reloading it !!!\n");
exit(1);
}
}
@@ -289,14 +289,14 @@ void uwsgi_proxy(int proxyfd) {
if (eevents[i].ASYNC_IS_IN) {
// is this a connected client/worker ?
//fprintf(stderr,"ready %d\n", upcs[eevents[i].data.fd].status);
//uwsgi_log("ready %d\n", upcs[eevents[i].data.fd].status);
if (!upcs[eevents[i].ASYNC_FD].status) {
if (upcs[eevents[i].ASYNC_FD].dest_fd >= 0) {
rlen = read(eevents[i].ASYNC_FD, buffer, 4096);
if (rlen < 0) {
perror("read()");
uwsgi_error("read()");
uwsgi_proxy_close(upcs, eevents[i].ASYNC_FD);
continue;
}
@@ -307,7 +307,7 @@ void uwsgi_proxy(int proxyfd) {
else {
wlen = write(upcs[eevents[i].ASYNC_FD].dest_fd, buffer, rlen);
if (wlen != rlen) {
perror("write()");
uwsgi_error("write()");
uwsgi_proxy_close(upcs, eevents[i].ASYNC_FD);
continue;
}
@@ -323,7 +323,7 @@ void uwsgi_proxy(int proxyfd) {
continue;
}
else {
fprintf(stderr, "UNKNOWN STATUS %d\n", upcs[eevents[i].ASYNC_FD].status);
uwsgi_log( "UNKNOWN STATUS %d\n", upcs[eevents[i].ASYNC_FD].status);
continue;
}
}
@@ -333,15 +333,15 @@ void uwsgi_proxy(int proxyfd) {
#ifdef UWSGI_PROXY_USE_KQUEUE
if (getsockopt(eevents[i].ASYNC_FD, SOL_SOCKET, SO_ERROR, (void *) (&soopt), &solen) < 0) {
perror("getsockopt()");
uwsgi_error("getsockopt()");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
/* is something bad ? */
if (soopt) {
fprintf(stderr, "connect() %s\n", strerror(soopt));
uwsgi_log( "connect() %s\n", strerror(soopt));
// increase errors on node
fprintf(stderr, "*** marking cluster node %d/%s as failed ***\n", upcs[eevents[i].ASYNC_FD].node, uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].name);
uwsgi_log( "*** marking cluster node %d/%s as failed ***\n", upcs[eevents[i].ASYNC_FD].node, uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].name);
uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].errors++;
uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].status = UWSGI_NODE_FAILED;
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
@@ -366,32 +366,32 @@ void uwsgi_proxy(int proxyfd) {
}
// re-set blocking
if (ioctl(ev.ASYNC_FD, FIONBIO, &blocking)) {
perror("ioctl()");
uwsgi_error("ioctl()");
uwsgi_proxy_close(upcs, ev.ASYNC_FD);
continue;
}
}
else {
fprintf(stderr, "strange event for %d\n", (int) eevents[i].ASYNC_FD);
uwsgi_log( "strange event for %d\n", (int) eevents[i].ASYNC_FD);
}
}
else {
if (upcs[eevents[i].ASYNC_FD].status == UWSGI_PROXY_CONNECTING) {
if (getsockopt(eevents[i].ASYNC_FD, SOL_SOCKET, SO_ERROR, (void *) (&soopt), &solen) < 0) {
perror("getsockopt()");
uwsgi_error("getsockopt()");
}
/* is something bad ? */
if (soopt) {
fprintf(stderr, "connect() %s\n", strerror(soopt));
uwsgi_log( "connect() %s\n", strerror(soopt));
}
// increase errors on node
fprintf(stderr, "*** marking cluster node %d/%s as failed ***\n", upcs[eevents[i].ASYNC_FD].node, uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].name);
uwsgi_log( "*** marking cluster node %d/%s as failed ***\n", upcs[eevents[i].ASYNC_FD].node, uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].name);
uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].errors++;
uwsgi.shared->nodes[upcs[eevents[i].ASYNC_FD].node].status = UWSGI_NODE_FAILED;
}
else {
fprintf(stderr, "STRANGE EVENT !!! %d %d %d\n", (int) eevents[i].ASYNC_FD, (int) eevents[i].ASYNC_EV, upcs[eevents[i].ASYNC_FD].status);
uwsgi_log( "STRANGE EVENT !!! %d %d %d\n", (int) eevents[i].ASYNC_FD, (int) eevents[i].ASYNC_EV, upcs[eevents[i].ASYNC_FD].status);
}
uwsgi_proxy_close(upcs, eevents[i].ASYNC_FD);
continue;
+6 -7
View File
@@ -8,11 +8,10 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
ssize_t sf_len = 0 ;
#endif
//fprintf(stderr,"managing request for %d %p\n", wsgi_req->async_id, wsgi_req);
// return or yield ?
if (PyString_Check((PyObject *)wsgi_req->async_result)) {
if ((wsize = write(wsgi_req->poll.fd, PyString_AsString(wsgi_req->async_result), PyString_Size(wsgi_req->async_result))) < 0) {
perror("write()");
uwsgi_error("write()");
goto clear;
}
wsgi_req->response_size += wsize;
@@ -22,8 +21,7 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
#ifdef PYTHREE
if (PyBytes_Check((PyObject *)wsgi_req->async_result)) {
if ((wsize = write(wsgi_req->poll.fd, PyBytes_AsString(wsgi_req->async_result), PyBytes_Size(wsgi_req->async_result))) < 0) {
perror("write()");
Py_DECREF(wsgi_req->async_result);
uwsgi_error("write()");
goto clear;
}
wsgi_req->response_size += wsize;
@@ -47,13 +45,13 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
}
#endif
// ok its a yield
if (!wsgi_req->async_placeholder) {
wsgi_req->async_placeholder = PyObject_GetIter(wsgi_req->async_result);
if (!wsgi_req->async_placeholder) {
goto clear2;
}
Py_DECREF((PyObject *)wsgi_req->async_result);
#ifdef UWSGI_ASYNC
if (uwsgi->async > 1) {
return UWSGI_AGAIN;
@@ -63,6 +61,7 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
pychunk = PyIter_Next(wsgi_req->async_placeholder) ;
if (!pychunk) {
@@ -74,7 +73,7 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
if (PyString_Check(pychunk)) {
if ((wsize = write(wsgi_req->poll.fd, PyString_AsString(pychunk), PyString_Size(pychunk))) < 0) {
perror("write()");
uwsgi_error("write()");
Py_DECREF(pychunk);
goto clear;
}
@@ -84,7 +83,7 @@ int manage_python_response(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
#ifdef PYTHREE
if (PyBytes_Check(pychunk)) {
if ((wsize = write(wsgi_req->poll.fd, PyBytes_AsString(pychunk), PyBytes_Size(pychunk))) < 0) {
perror("write()");
uwsgi_error("write()");
Py_DECREF(pychunk);
goto clear;
}
+8 -8
View File
@@ -35,7 +35,7 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
if (!wsgi_req->sendfile_fd_size) {
if (fstat(fd, &stat_buf)) {
perror("fstat()");
uwsgi_error("fstat()");
return 0;
}
else {
@@ -61,7 +61,7 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
}
if (sf_ret) {
perror("sendfile()");
uwsgi_error("sendfile()");
return 0;
}
@@ -79,7 +79,7 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
}
if (sf_ret) {
perror("sendfile()");
uwsgi_error("sendfile()");
return 0;
}
@@ -94,7 +94,7 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
}
if (sf_ret < 0) {
perror("sendfile()");
uwsgi_error("sendfile()");
return 0;
}
@@ -119,12 +119,12 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
if (uwsgi->async > 1) {
jlen = read(fd, nosf_buf, wsgi_req->sendfile_fd_chunk);
if (jlen <= 0) {
perror("read()");
uwsgi_error("read()");
return 0;
}
jlen = write(sockfd, nosf_buf, jlen);
if (jlen <= 0) {
perror("write()");
uwsgi_error("write()");
return 0;
}
return jlen ;
@@ -133,13 +133,13 @@ ssize_t uwsgi_sendfile(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
while (i < wsgi_req->sendfile_fd_size) {
jlen = read(fd, nosf_buf, wsgi_req->sendfile_fd_chunk);
if (jlen <= 0) {
perror("read()");
uwsgi_error("read()");
break;
}
i += jlen;
jlen = write(sockfd, nosf_buf, jlen);
if (jlen <= 0) {
perror("write()");
uwsgi_error("write()");
break;
}
rlen += jlen ;
+7 -3
View File
@@ -3,9 +3,10 @@ import sys
import uwsgiconfig as uc
import shutil
from distutils.core import setup, Distribution
from distutils.command.install import install
from distutils.command.build_ext import build_ext
from setuptools import setup
from setuptools.dist import Distribution
from setuptools.command.install import install
from setuptools.command.build_ext import build_ext
class uWSGIBuilder(build_ext):
@@ -17,6 +18,9 @@ class uWSGIBuilder(build_ext):
class uWSGIInstall(install):
def run(self):
# hack, hack and still hack. We need to find a solution for 0.9.6
if self.record:
record_file = open(self.record,'w')
uc.parse_vars()
uc.build_uwsgi(sys.prefix + '/bin/' + uc.UWSGI_BIN_NAME)
+4
View File
@@ -2,3 +2,7 @@ def application(env, start_response):
#print env
start_response('200 Ok', [('Content-type', 'text/plain')])
yield "hello world"
yield "hello world2"
for i in xrange(1,1000):
yield str(i)
+13 -8
View File
@@ -43,17 +43,21 @@ void manage_snmp(int fd, uint8_t * buffer, int size, struct sockaddr_in *client_
uint64_t request_id;
uint64_t version;
// KISS for memory management
if (size > SNMP_WATERMARK)
return;
ptr++;
// check total sequence size
if (*ptr > SNMP_WATERMARK || *ptr < 13)
return;
ptr++;
// check snmp version
if (*ptr != SNMP_INTEGER)
return;
@@ -66,6 +70,7 @@ void manage_snmp(int fd, uint8_t * buffer, int size, struct sockaddr_in *client_
// check for community string (this must be set from the python vm using uwsgi.snmp_community or with --snmp-community arg)
if (*ptr != SNMP_STRING)
return;
@@ -206,7 +211,7 @@ void manage_snmp(int fd, uint8_t * buffer, int size, struct sockaddr_in *client_
if (size > 0) {
if (sendto(fd, buffer, size, 0, (struct sockaddr *) client_addr, sizeof(struct sockaddr_in)) < 0) {
perror("sendto()");
uwsgi_error("sendto()");
}
}
@@ -266,7 +271,7 @@ static uint8_t snmp_int_to_snmp(uint64_t snmp_val, uint8_t oid_type, uint8_t * b
uint8_t tlen;
int i, j;
uint8_t *ptr = (uint8_t *) & snmp_val;
uint8_t *ptr = (uint8_t *) &snmp_val;
// check for counter, counter64 or gauge
@@ -334,7 +339,7 @@ PyObject *py_snmp_counter32(PyObject * self, PyObject * args) {
uint8_t oid_num;
uint32_t oid_val = 0;
if (!PyArg_ParseTuple(args, "bi:snmp_set_counter32", &oid_num, &oid_val)) {
if (!PyArg_ParseTuple(args, "bI:snmp_set_counter32", &oid_num, &oid_val)) {
return NULL;
}
@@ -358,7 +363,7 @@ PyObject *py_snmp_counter64(PyObject * self, PyObject * args) {
uint8_t oid_num;
uint64_t oid_val = 0;
if (!PyArg_ParseTuple(args, "bl:snmp_set_counter64", &oid_num, &oid_val)) {
if (!PyArg_ParseTuple(args, "bK:snmp_set_counter64", &oid_num, &oid_val)) {
return NULL;
}
@@ -382,7 +387,7 @@ PyObject *py_snmp_gauge(PyObject * self, PyObject * args) {
uint8_t oid_num;
uint32_t oid_val = 0;
if (!PyArg_ParseTuple(args, "bi:snmp_set_gauge", &oid_num, &oid_val)) {
if (!PyArg_ParseTuple(args, "bI:snmp_set_gauge", &oid_num, &oid_val)) {
return NULL;
}
@@ -410,11 +415,11 @@ PyObject *py_snmp_community(PyObject * self, PyObject * args) {
}
if (strlen(snmp_community) > 72) {
fprintf(stderr, "*** warning the supplied SNMP community string will be truncated to 72 chars ***\n");
uwsgi_log( "*** warning the supplied SNMP community string will be truncated to 72 chars ***\n");
memcpy(uwsgi.shared->snmp_community, snmp_community, 72);
}
else {
strcpy(uwsgi.shared->snmp_community, snmp_community);
strlcpy(uwsgi.shared->snmp_community, snmp_community, 73);
}
Py_INCREF(Py_True);
@@ -440,7 +445,7 @@ void snmp_init() {
Py_DECREF(func);
}
fprintf(stderr, "SNMP python functions initialized.\n");
uwsgi_log( "SNMP python functions initialized.\n");
}
+35 -35
View File
@@ -7,54 +7,54 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst
int serverfd;
struct sockaddr_un *uws_addr;
fprintf(stderr, "binding on UNIX socket: %s\n", socket_name);
uwsgi_log( "binding on UNIX socket: %s\n", socket_name);
// leave 1 byte for abstract namespace (108 linux -> 104 bsd/mac)
if (strlen(socket_name) > 102) {
fprintf(stderr, "invalid socket name\n");
uwsgi_log( "invalid socket name\n");
exit(1);
}
uws_addr = malloc(sizeof(struct sockaddr_un));
if (uws_addr == NULL) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
memset(uws_addr, 0, sizeof(struct sockaddr_un));
serverfd = socket(AF_UNIX, SOCK_STREAM, 0);
if (serverfd < 0) {
perror("socket()");
uwsgi_error("socket()");
exit(1);
}
if (abstract_socket == 0) {
if (unlink(socket_name) != 0 && errno != ENOENT) {
perror("unlink()");
uwsgi_error("unlink()");
}
}
if (abstract_socket == 1) {
fprintf(stderr, "setting abstract socket mode (warning: only Linux supports this)\n");
uwsgi_log( "setting abstract socket mode (warning: only Linux supports this)\n");
}
uws_addr->sun_family = AF_UNIX;
strcpy(uws_addr->sun_path + abstract_socket, socket_name);
strlcpy(uws_addr->sun_path + abstract_socket, socket_name, 102+1);
if (bind(serverfd, (struct sockaddr *) uws_addr, strlen(socket_name) + abstract_socket + ((void *) uws_addr->sun_path - (void *) uws_addr)) != 0) {
perror("bind()");
uwsgi_error("bind()");
exit(1);
}
if (listen(serverfd, listen_queue) != 0) {
perror("listen()");
uwsgi_error("listen()");
exit(1);
}
// chmod unix socket for lazy users
if (chmod_socket == 1 && abstract_socket == 0) {
fprintf(stderr, "chmod() socket to 666 for lazy and brave users\n");
uwsgi_log( "chmod() socket to 666 for lazy and brave users\n");
if (chmod(socket_name, S_IRUSR | S_IWUSR | S_IRGRP | S_IWGRP | S_IROTH | S_IWOTH) != 0) {
perror("chmod()");
uwsgi_error("chmod()");
}
}
@@ -91,15 +91,15 @@ int bind_to_sctp(char *socket_name, int listen_queue, char *sctp_port) {
serverfd = socket(AF_INET, SOCK_STREAM, IPPROTO_SCTP);
if (serverfd < 0) {
perror("socket()");
uwsgi_error("socket()");
exit(1);
}
fprintf(stderr, "binding on %d SCTP interfaces on port: %d\n", num_ip, ntohs(uws_addr[0].sin_port));
uwsgi_log( "binding on %d SCTP interfaces on port: %d\n", num_ip, ntohs(uws_addr[0].sin_port));
if (sctp_bindx(serverfd, (struct sockaddr *) uws_addr, num_ip, SCTP_BINDX_ADD_ADDR) != 0) {
perror("sctp_bindx()");
uwsgi_error("sctp_bindx()");
exit(1);
}
@@ -107,11 +107,11 @@ int bind_to_sctp(char *socket_name, int listen_queue, char *sctp_port) {
sctp_im.sinit_num_ostreams = 0xFFFF;
if (setsockopt(serverfd, IPPROTO_SCTP, SCTP_INITMSG, &sctp_im, sizeof(sctp_im))) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
}
if (listen(serverfd, listen_queue) != 0) {
perror("listen()");
uwsgi_error("listen()");
exit(1);
}
@@ -150,7 +150,7 @@ int bind_to_udp(char *socket_name) {
serverfd = socket(AF_INET, SOCK_DGRAM, 0);
if (serverfd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1;
}
@@ -168,19 +168,19 @@ int bind_to_udp(char *socket_name) {
}
#endif
fprintf(stderr, "binding on UDP port: %d\n", ntohs(uws_addr.sin_port));
uwsgi_log( "binding on UDP port: %d\n", ntohs(uws_addr.sin_port));
if (bind(serverfd, (struct sockaddr *) &uws_addr, sizeof(uws_addr)) != 0) {
perror("bind()");
uwsgi_error("bind()");
close(serverfd);
return -1;
}
#ifdef UWSGI_MULTICAST
if (uwsgi.multicast_group) {
fprintf(stderr, "joining uWSGI multicast group: %s:%d\n", uwsgi.multicast_group, ntohs(uws_addr.sin_port));
uwsgi_log( "joining uWSGI multicast group: %s:%d\n", uwsgi.multicast_group, ntohs(uws_addr.sin_port));
if (setsockopt(serverfd, IPPROTO_IP, IP_ADD_MEMBERSHIP, &mc, sizeof(mc))) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
}
}
#endif
@@ -209,14 +209,14 @@ int connect_to_tcp(char *socket_name, int port, int timeout) {
uwsgi_poll.fd = socket(AF_INET, SOCK_STREAM, 0);
if (uwsgi_poll.fd < 0) {
perror("socket()");
uwsgi_error("socket()");
return -1;
}
uwsgi_poll.events = POLLIN;
if (timed_connect(&uwsgi_poll, (const struct sockaddr *) &uws_addr, sizeof(struct sockaddr_in), timeout)) {
perror("connect()");
uwsgi_error("connect()");
close(uwsgi_poll.fd);
return -1;
}
@@ -247,12 +247,12 @@ int bind_to_tcp(char *socket_name, int listen_queue, char *tcp_port) {
serverfd = socket(AF_INET, SOCK_STREAM, 0);
if (serverfd < 0) {
perror("socket()");
uwsgi_error("socket()");
exit(1);
}
if (setsockopt(serverfd, SOL_SOCKET, SO_REUSEADDR, (const void *) &reuse, sizeof(int)) < 0) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
exit(1);
}
@@ -260,29 +260,29 @@ int bind_to_tcp(char *socket_name, int listen_queue, char *tcp_port) {
#ifdef __linux__
if (setsockopt(serverfd, IPPROTO_TCP, TCP_DEFER_ACCEPT, &uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], sizeof(int))) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
}
#elif defined(__apple__) || defined(__freebsd__)
struct accept_filter_arg afa;
strcpy(afa.af_name, "dataready");
afa.af_arg[0] = 0;
if (setsockopt(serverfd, SOL_SOCKET, SO_ACCEPTFILTER, &afa, sizeof(struct accept_filter_arg))) {
perror("setsockopt()");
uwsgi_error("setsockopt()");
}
#endif
}
fprintf(stderr, "binding on TCP port: %d\n", ntohs(uws_addr.sin_port));
uwsgi_log( "binding on TCP port: %d\n", ntohs(uws_addr.sin_port));
if (bind(serverfd, (struct sockaddr *) &uws_addr, sizeof(uws_addr)) != 0) {
perror("bind()");
uwsgi_error("bind()");
exit(1);
}
if (listen(serverfd, listen_queue) != 0) {
perror("listen()");
uwsgi_error("listen()");
exit(1);
}
@@ -300,12 +300,12 @@ int timed_connect(struct pollfd *fdpoll, const struct sockaddr *addr, int addr_s
arg = fcntl(fdpoll->fd, F_GETFL, NULL);
if (arg < 0) {
perror("fcntl()");
uwsgi_error("fcntl()");
return -1;
}
arg |= O_NONBLOCK;
if (fcntl(fdpoll->fd, F_SETFL, arg) < 0) {
perror("fcntl()");
uwsgi_error("fcntl()");
return -1;
}
@@ -321,13 +321,13 @@ int timed_connect(struct pollfd *fdpoll, const struct sockaddr *addr, int addr_s
cnt = poll(fdpoll, 1, timeout * 1000);
/* check for errors */
if (cnt < 0 && errno != EINTR) {
perror("poll()");
uwsgi_error("poll()");
return -1;
}
/* something hapened on the socket ... */
else if (cnt > 0) {
if (getsockopt(fdpoll->fd, SOL_SOCKET, SO_ERROR, (void *) (&soopt), &solen) < 0) {
perror("getsockopt()");
uwsgi_error("getsockopt()");
return -1;
}
/* is something bad ? */
@@ -348,7 +348,7 @@ int timed_connect(struct pollfd *fdpoll, const struct sockaddr *addr, int addr_s
/* re-set blocking socket */
arg &= (~O_NONBLOCK);
if (fcntl(fdpoll->fd, F_SETFL, arg) < 0) {
perror("fcntl()");
uwsgi_error("fcntl()");
return -1;
}
+36 -36
View File
@@ -13,7 +13,7 @@ int spool_request(struct uwsgi_server *uwsgi, char *filename, int rn, char *buff
struct uwsgi_header uh;
if (gethostname(hostname, 256)) {
perror("gethostname()");
uwsgi_error("gethostname()");
return 0;
}
@@ -27,16 +27,16 @@ int spool_request(struct uwsgi_server *uwsgi, char *filename, int rn, char *buff
fd = open(filename, O_CREAT | O_EXCL | O_WRONLY, S_IRUSR | S_IWUSR);
if (fd < 0) {
perror("open()");
uwsgi_error("open()");
return 0;
}
#ifdef __sun__
if (lockf(fd, F_LOCK, 0)) {
perror("lockf()");
uwsgi_error("lockf()");
#else
if (flock(fd, LOCK_EX)) {
perror("flock()");
uwsgi_error("flock()");
#endif
close(fd);
return 0;
@@ -59,14 +59,14 @@ int spool_request(struct uwsgi_server *uwsgi, char *filename, int rn, char *buff
close(fd);
fprintf(stderr, "written %d bytes to spool file %s.\n", size + 4, filename);
uwsgi_log( "written %d bytes to spool file %s.\n", size + 4, filename);
return 1;
clear:
perror("write()");
uwsgi_error("write()");
unlink(filename);
close(fd);
return 0;
@@ -92,14 +92,14 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
spool_tuple = PyTuple_New(1);
if (!spool_tuple) {
fprintf(stderr, "could not create spooler tuple.\n");
uwsgi_log( "could not create spooler tuple.\n");
exit(1);
}
spool_env = PyDict_New();
if (!spool_env) {
fprintf(stderr, "could not create spooler env.\n");
uwsgi_log( "could not create spooler env.\n");
exit(1);
}
@@ -109,12 +109,12 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
}
if (chdir(uwsgi->spool_dir)) {
perror("chdir()");
uwsgi_error("chdir()");
exit(1);
}
// asked by Marco Beri
fprintf(stderr, "lowering spooler priority to %d\n", PRIO_MAX);
uwsgi_log( "lowering spooler priority to %d\n", PRIO_MAX);
setpriority(PRIO_PROCESS, getpid(), PRIO_MAX);
for (;;) {
@@ -133,33 +133,33 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
continue;
}
if (!access(dp->d_name, R_OK | W_OK)) {
fprintf(stderr, "managing spool request %s...\n", dp->d_name);
uwsgi_log( "managing spool request %s...\n", dp->d_name);
spooler_callable = PyDict_GetItemString(uwsgi_module_dict, "spooler");
if (!spooler_callable) {
fprintf(stderr, "you have to define uwsgi.spooler to use the spooler !!!\n");
uwsgi_log( "you have to define uwsgi.spooler to use the spooler !!!\n");
continue;
}
spool_fd = open(dp->d_name, O_RDONLY);
if (spool_fd < 0) {
perror("open()");
uwsgi_error("open()");
continue;
}
#ifdef __sun__
if (lockf(spool_fd, F_LOCK, 0)) {
perror("lockf()");
uwsgi_error("lockf()");
#else
if (flock(spool_fd, LOCK_EX)) {
perror("flock()");
uwsgi_error("flock()");
#endif
close(spool_fd);
continue;
}
if (read(spool_fd, &uh, 4) != 4) {
perror("read()");
uwsgi_error("read()");
close(spool_fd);
continue;
}
@@ -173,7 +173,7 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
while (datasize < uh.pktsize) {
rlen = read(spool_fd, &uwstrlen, 2);
if (rlen != 2) {
perror("read()");
uwsgi_error("read()");
goto next_spool;
}
datasize += rlen;
@@ -182,12 +182,12 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
if (uwstrlen > 0) {
key = malloc(uwstrlen + 1);
if (!key) {
perror("malloc()");
uwsgi_error("malloc()");
goto retry_later;
}
rlen = read(spool_fd, key, uwstrlen);
if (rlen != uwstrlen) {
perror("read()");
uwsgi_error("read()");
free(key);
goto next_spool;
}
@@ -197,7 +197,7 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
rlen = read(spool_fd, &uwstrlen, 2);
if (rlen != 2) {
perror("read()");
uwsgi_error("read()");
free(key);
goto next_spool;
}
@@ -207,13 +207,13 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
val = malloc(uwstrlen + 1);
if (!val) {
free(key);
perror("malloc()");
uwsgi_error("malloc()");
goto retry_later;
}
rlen = read(spool_fd, val, uwstrlen);
if (rlen != uwstrlen) {
perror("read()");
uwsgi_error("read()");
free(key);
goto next_spool;
}
@@ -241,25 +241,25 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
spool_result = python_call(spooler_callable, spool_tuple);
if (!spool_result) {
PyErr_Print();
fprintf(stderr, "error detected. spool request canceled.\n");
uwsgi_log( "error detected. spool request canceled.\n");
goto next_spool;
}
if (PyInt_Check(spool_result)) {
if (PyInt_AsLong(spool_result) == 17) {
Py_DECREF(spool_result);
fprintf(stderr, "retry this task later...\n");
uwsgi_log( "retry this task later...\n");
goto retry_later;
}
}
Py_DECREF(spool_result);
fprintf(stderr, "done with task/spool %s\n", dp->d_name);
uwsgi_log( "done with task/spool %s\n", dp->d_name);
next_spool:
if (unlink(dp->d_name)) {
perror("unlink");
fprintf(stderr, "something horrible happened to the spooler. Better to kill it.\n");
uwsgi_error("unlink");
uwsgi_log( "something horrible happened to the spooler. Better to kill it.\n");
exit(1);
}
retry_later:
@@ -271,7 +271,7 @@ void spooler(struct uwsgi_server *uwsgi, PyObject * uwsgi_module_dict) {
closedir(sdir);
}
else {
perror("opendir()");
uwsgi_error("opendir()");
}
}
@@ -283,29 +283,29 @@ int uwsgi_request_spooler(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_
char spool_filename[1024];
if (uwsgi->spool_dir == NULL) {
fprintf(stderr, "the spooler is inactive !!!...skip\n");
uwsgi_log( "the spooler is inactive !!!...skip\n");
wsgi_req->uh.modifier1 = 255;
wsgi_req->uh.pktsize = 0;
wsgi_req->uh.modifier2 = 0;
i = write(wsgi_req->poll.fd, wsgi_req, 4);
if (i != 4) {
perror("write()");
uwsgi_error("write()");
}
return -1;
}
fprintf(stderr, "managing spool request...\n");
i = spool_request(uwsgi, spool_filename, uwsgi->workers[0].requests + 1, &wsgi_req->buffer, wsgi_req->uh.pktsize);
uwsgi_log( "managing spool request...\n");
i = spool_request(uwsgi, spool_filename, uwsgi->workers[0].requests + 1, wsgi_req->buffer, wsgi_req->uh.pktsize);
wsgi_req->uh.modifier1 = 255;
wsgi_req->uh.pktsize = 0;
if (i > 0) {
wsgi_req->uh.modifier2 = 1;
if (write(wsgi_req->poll.fd, wsgi_req, 4) != 4) {
fprintf(stderr, "disconnected client, remove spool file.\n");
uwsgi_log( "disconnected client, remove spool file.\n");
/* client disconnect, remove spool file */
if (unlink(spool_filename)) {
perror("unlink()");
fprintf(stderr, "something horrible happened !!! check your spooler ASAP !!!\n");
uwsgi_error("unlink()");
uwsgi_log( "something horrible happened !!! check your spooler ASAP !!!\n");
goodbye_cruel_world();
}
}
@@ -316,7 +316,7 @@ int uwsgi_request_spooler(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_
wsgi_req->uh.modifier2 = 0;
i = write(wsgi_req->poll.fd, wsgi_req, 4);
if (i != 4) {
perror("write()");
uwsgi_error("write()");
}
}
+9 -9
View File
@@ -25,7 +25,7 @@ PyObject *py_uwsgi_stackless_worker(PyObject * self, PyObject * args) {
int async_id = wsgi_req->async_id;
//fprintf(stderr,"i am the tasklet worker\n");
//uwsgi_log("i am the tasklet worker\n");
for(;;) {
@@ -62,19 +62,19 @@ void stackless_init(struct uwsgi_server *uwsgi) {
uwsgi->stackless_table = malloc( sizeof(struct stackless_req*) * uwsgi->async);
if (!uwsgi->stackless_table) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
for(i=0;i<uwsgi->async;i++) {
uwsgi->stackless_table[i] = malloc(sizeof(struct stackless_req));
if (!uwsgi->stackless_table[i]) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
memset(uwsgi->stackless_table[i], 0, sizeof(struct stackless_req));
}
fprintf(stderr,"initializing %d tasklet...", uwsgi->async);
uwsgi_log("initializing %d tasklet...", uwsgi->async);
// creating uwsgi->async tasklets
for(i=0;i<uwsgi->async;i++) {
@@ -90,7 +90,7 @@ void stackless_init(struct uwsgi_server *uwsgi) {
wsgi_req = next_wsgi_req(uwsgi, wsgi_req) ;
}
fprintf(stderr,"done\n");
uwsgi_log("done\n");
}
void stackless_loop(struct uwsgi_server *uwsgi) {
@@ -114,7 +114,7 @@ void stackless_loop(struct uwsgi_server *uwsgi) {
if (uwsgi->async_events[i].ASYNC_FD == uwsgi->serverfd) {
//pass the connection to the first available tasklet
fprintf(stderr,"sending new connection...\n");
uwsgi_log("sending new connection...\n");
PyChannel_Send(uwsgi->workers_channel, Py_True);
}
@@ -130,11 +130,11 @@ void stackless_loop(struct uwsgi_server *uwsgi) {
//int_tasklet = (PyTaskletObject *) PyStackless_RunWatchdog( 1000 );
/*
fprintf(stderr,"done watchdog %p\n", int_tasklet);
uwsgi_log("done watchdog %p\n", int_tasklet);
if (!PyTasklet_IsCurrent(int_tasklet)) {
fprintf(stderr,"re-insert: %d\n", 1);// PyTasklet_Insert(int_tasklet));
uwsgi_log("re-insert: %d\n", 1);// PyTasklet_Insert(int_tasklet));
}
fprintf(stderr,"recycle\n");
uwsgi_log("recycle\n");
*/
}
+1 -1
View File
@@ -2,5 +2,5 @@ import time
def application(env, start_response):
start_response( '200 OK', [ ('Content-Type','text/html') ])
for i in range(1,100000):
for i in range(1,1000):
yield "<h1>%s at %s</h1>" % (i, str(time.time()))
+45
View File
@@ -0,0 +1,45 @@
import uwsgi
import psycopg2
def ugreen_wait_callback(conn, timeout=-1):
"""A wait callback useful to allow uWSGI/uGreen to work with Psycopg."""
while True:
state = conn.poll()
if state == psycopg2.extensions.POLL_OK:
break
elif state == psycopg2.extensions.POLL_READ:
uwsgi.green_wait_fdread(conn.fileno())
elif state == psycopg2.extensions.POLL_WRITE:
uwsgi.green_wait_fdwrite(conn.fileno())
else:
raise Exception("Unexpected result from poll: %r", state)
# set the wait callback
psycopg2.extensions.set_wait_callback(ugreen_wait_callback)
def application(env, start_response):
start_response('200 Ok', [('Content-type', 'text/html')])
# connect
conn = psycopg2.connect("dbname=prova user=postgres")
# get cursor
curs = conn.cursor()
yield "<table>"
# run query
curs.execute("SELECT * FROM tests")
while True:
row = curs.fetchone()
if not row: break
yield "<tr><td>%s</td></tr>" % str(row)
yield "</table>"
conn.close()
+50
View File
@@ -0,0 +1,50 @@
import uwsgi
import psycopg2
def async_wait(conn):
# conn can be a connection or a cursor
if not hasattr(conn, 'poll'):
conn = conn.connection
# interesting part: suspend until ready
while True:
state = conn.poll()
if state == psycopg2.extensions.POLL_OK:
break
elif state == psycopg2.extensions.POLL_READ:
uwsgi.green_wait_fdread(conn.fileno())
elif state == psycopg2.extensions.POLL_WRITE:
uwsgi.green_wait_fdwrite(conn.fileno())
else:
raise Exception("Unexpected result from poll: %r", state)
def application(env, start_response):
start_response('200 Ok', [('Content-type', 'text/html')])
conn = psycopg2.connect("dbname=prova user=postgres", async=True)
# suspend until connection
async_wait(conn)
curs = conn.cursor()
yield "<table>"
curs.execute("SELECT * FROM tests")
# suspend until result
async_wait(curs)
while True:
row = curs.fetchone()
if not row: break
yield "<tr><td>%s</td></tr>" % str(row)
yield "</table>"
conn.close()
+16 -14
View File
@@ -18,7 +18,7 @@ void u_green_write_all(struct uwsgi_server *uwsgi, char *data, size_t len) {
if (wsgi_req->async_status == UWSGI_PAUSED) {
rlen = write(wsgi_req->poll.fd, data, len);
if (rlen < 0) {
perror("write()");
uwsgi_error("write()");
// mark core as plagued
wsgi_req->async_plagued = 1 ;
}
@@ -262,6 +262,7 @@ static void u_green_request(struct uwsgi_server *uwsgi, struct wsgi_request *wsg
}
wsgi_req->async_status = UWSGI_OK;
u_green_schedule_to_main(uwsgi, async_id);
if (wsgi_req_recv(wsgi_req)) {
@@ -295,12 +296,12 @@ void u_green_init(struct uwsgi_server *uwsgi) {
u_stack_size = uwsgi->ugreen_stackpages * uwsgi->page_size ;
}
fprintf(stderr,"initializing %d uGreen threads with stack size of %lu (%lu KB)\n", uwsgi->async, (unsigned long) u_stack_size, (unsigned long) u_stack_size/1024);
uwsgi_log("initializing %d uGreen threads with stack size of %lu (%lu KB)\n", uwsgi->async, (unsigned long) u_stack_size, (unsigned long) u_stack_size/1024);
uwsgi->ugreen_contexts = malloc( sizeof(ucontext_t*) * uwsgi->async);
if (!uwsgi->ugreen_contexts) {
perror("malloc()\n");
uwsgi_error("malloc()\n");
exit(1);
}
@@ -308,24 +309,26 @@ void u_green_init(struct uwsgi_server *uwsgi) {
for(i=0;i<uwsgi->async;i++) {
uwsgi->ugreen_contexts[i] = malloc( sizeof(ucontext_t) );
if (!uwsgi->ugreen_contexts[i]) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
getcontext(uwsgi->ugreen_contexts[i]);
uwsgi->ugreen_contexts[i]->uc_stack.ss_sp = mmap(NULL, u_stack_size + (uwsgi->page_size*2) , PROT_READ | PROT_WRITE | PROT_EXEC, MAP_ANON | MAP_PRIVATE, -1, 0) + uwsgi->page_size;
if (!uwsgi->ugreen_contexts[i]->uc_stack.ss_sp) {
perror("mmap()");
uwsgi_error("mmap()");
exit(1);
}
// set guard pages for stack
if (mprotect(uwsgi->ugreen_contexts[i]->uc_stack.ss_sp - uwsgi->page_size, uwsgi->page_size, PROT_NONE)) {
perror("mprotect()");
uwsgi_error("mprotect()");
exit(1);
}
if (mprotect(uwsgi->ugreen_contexts[i]->uc_stack.ss_sp + u_stack_size, uwsgi->page_size, PROT_NONE)) {
perror("mprotect()");
uwsgi_error("mprotect()");
exit(1);
}
uwsgi->ugreen_contexts[i]->uc_stack.ss_size = u_stack_size ;
uwsgi->ugreen_contexts[i]->uc_link = NULL;
makecontext(uwsgi->ugreen_contexts[i], (void (*) (void)) &u_green_request, 3, uwsgi, wsgi_req, i);
@@ -397,25 +400,23 @@ void u_green_loop(struct uwsgi_server *uwsgi) {
while(uwsgi->workers[uwsgi->mywid].manage_next_request) {
uwsgi->async_running = u_green_blocking(uwsgi) ;
timeout = u_green_get_timeout(uwsgi);
uwsgi->async_nevents = async_wait(uwsgi->async_queue, uwsgi->async_events, uwsgi->async, uwsgi->async_running, timeout);
u_green_expire_timeouts(uwsgi);
if (uwsgi->async_nevents < 0) {
continue;
}
u_green_expire_timeouts(uwsgi);
if (uwsgi->async_nevents > 0) {
wsgi_req = find_first_accepting_wsgi_req(uwsgi);
if (!wsgi_req) goto cycle;
}
for(i=0; i<uwsgi->async_nevents;i++) {
if (uwsgi->async_events[i].ASYNC_FD == uwsgi->serverfd) {
wsgi_req = find_first_accepting_wsgi_req(uwsgi);
if (!wsgi_req) goto cycle;
u_green_schedule_to_req(uwsgi, wsgi_req);
}
else {
@@ -431,6 +432,7 @@ void u_green_loop(struct uwsgi_server *uwsgi) {
}
cycle:
wsgi_req = find_wsgi_req_by_id(uwsgi, current) ;
if (wsgi_req->async_status != UWSGI_ACCEPTING && wsgi_req->async_status != UWSGI_PAUSED && wsgi_req->async_waiting_fd == -1 && !wsgi_req->async_timeout) {
u_green_schedule_to_req(uwsgi, wsgi_req);
+70 -35
View File
@@ -7,10 +7,28 @@ extern struct uwsgi_server uwsgi;
uint16_t uwsgi_swap16(uint16_t x) {
return (uint16_t) ((x & 0xff) << 8 | (x & 0xff00) >> 8);
}
uint32_t uwsgi_swap32(uint32_t x) {
x = ( (x<<8) & 0xFF00FF00 ) | ( (x>>8) & 0x00FF00FF );
return (x>>16) | (x<<16);
}
// thanks to ffmpeg project for this idea :P
uint64_t uwsgi_swap64(uint64_t x) {
union {
uint64_t ll;
uint32_t l[2];
} w, r;
w.ll = x;
r.l[0] = uwsgi_swap32(w.l[1]);
r.l[1] = uwsgi_swap32(w.l[0]);
return r.ll;
}
#endif
void set_harakiri(int sec) {
if (uwsgi.workers) {
if (uwsgi.master_process) {
if (sec == 0) {
uwsgi.workers[uwsgi.mywid].harakiri = 0;
}
@@ -38,7 +56,7 @@ void daemonize(char *logfile) {
pid = fork();
if (pid < 0) {
perror("fork()");
uwsgi_error("fork()");
exit(1);
}
if (pid != 0) {
@@ -46,7 +64,7 @@ void daemonize(char *logfile) {
}
if (setsid() < 0) {
perror("setsid()");
uwsgi_error("setsid()");
exit(1);
}
@@ -54,7 +72,7 @@ void daemonize(char *logfile) {
/* refork... */
pid = fork();
if (pid < 0) {
perror("fork()");
uwsgi_error("fork()");
exit(1);
}
if (pid != 0) {
@@ -65,14 +83,14 @@ void daemonize(char *logfile) {
/*if (chdir("/") != 0) {
perror("chdir()");
uwsgi_error("chdir()");
exit(1);
} */
fdin = open("/dev/null", O_RDWR);
if (fdin < 0) {
perror("open()");
uwsgi_error("open()");
exit(1);
}
@@ -82,13 +100,13 @@ void daemonize(char *logfile) {
if (udp_port) {
udp_port[0] = 0 ;
if ( !udp_port[1] || !logfile[0] ) {
fprintf(stderr,"invalid udp address\n");
uwsgi_log("invalid udp address\n");
exit(1);
}
fd = socket(AF_INET, SOCK_DGRAM, 0);
if (fd < 0) {
perror("socket()");
uwsgi_error("socket()");
exit(1);
}
@@ -99,7 +117,7 @@ void daemonize(char *logfile) {
udp_addr.sin_addr.s_addr = inet_addr(logfile);
if (connect(fd, (const struct sockaddr *) &udp_addr, sizeof(struct sockaddr_in)) < 0) {
perror("connect()");
uwsgi_error("connect()");
exit(1);
}
}
@@ -107,7 +125,7 @@ void daemonize(char *logfile) {
#endif
fd = open(logfile, O_RDWR | O_CREAT | O_APPEND, S_IRUSR | S_IWUSR | S_IRGRP);
if (fd < 0) {
perror("open()");
uwsgi_error("open()");
exit(1);
}
#ifdef UWSGI_UDP
@@ -116,20 +134,20 @@ void daemonize(char *logfile) {
/* stdin */
if (dup2(fdin, 0) < 0) {
perror("dup2()");
uwsgi_error("dup2()");
exit(1);
}
/* stdout */
if (dup2(fd, 1) < 0) {
perror("dup2()");
uwsgi_error("dup2()");
exit(1);
}
/* stderr */
if (dup2(fd, 2) < 0) {
perror("dup2()");
uwsgi_error("dup2()");
exit(1);
}
@@ -148,21 +166,21 @@ char *uwsgi_get_cwd() {
cwd = malloc(newsize);
if (cwd == NULL) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
if (getcwd(cwd, newsize) == NULL) {
newsize = errno;
fprintf(stderr, "need a bigger buffer (%d bytes) for getcwd(). doing reallocation.\n", newsize);
uwsgi_log("need a bigger buffer (%d bytes) for getcwd(). doing reallocation.\n", newsize);
free(cwd);
cwd = malloc(newsize);
if (cwd == NULL) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
if (getcwd(cwd, newsize) == NULL) {
perror("getcwd()");
uwsgi_error("getcwd()");
exit(1);
}
}
@@ -190,49 +208,49 @@ void internal_server_error(int fd, char *message) {
void uwsgi_as_root() {
if (!getuid()) {
fprintf(stderr, "uWSGI running as root, you can use --uid/--gid/--chroot options\n");
uwsgi_log("uWSGI running as root, you can use --uid/--gid/--chroot options\n");
if (uwsgi.chroot) {
fprintf(stderr, "chroot() to %s\n", uwsgi.chroot);
uwsgi_log("chroot() to %s\n", uwsgi.chroot);
if (chroot(uwsgi.chroot)) {
perror("chroot()");
uwsgi_error("chroot()");
exit(1);
}
#ifdef __linux__
if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG]) {
fprintf(stderr, "*** Warning, on linux system you have to bind-mount the /proc fs in your chroot to get memory debug/report.\n");
uwsgi_log("*** Warning, on linux system you have to bind-mount the /proc fs in your chroot to get memory debug/report.\n");
}
#endif
}
if (uwsgi.gid) {
fprintf(stderr, "setgid() to %d\n", uwsgi.gid);
uwsgi_log("setgid() to %d\n", uwsgi.gid);
if (setgid(uwsgi.gid)) {
perror("setgid()");
uwsgi_error("setgid()");
exit(1);
}
}
if (uwsgi.uid) {
fprintf(stderr, "setuid() to %d\n", uwsgi.uid);
uwsgi_log("setuid() to %d\n", uwsgi.uid);
if (setuid(uwsgi.uid)) {
perror("setuid()");
uwsgi_error("setuid()");
exit(1);
}
}
if (!getuid()) {
fprintf(stderr, " *** WARNING: you are running uWSGI as root !!! (use the --uid flag) *** \n");
uwsgi_log(" *** WARNING: you are running uWSGI as root !!! (use the --uid flag) *** \n");
}
}
else {
if (uwsgi.chroot) {
fprintf(stderr, "cannot chroot() as non-root user\n");
uwsgi_log("cannot chroot() as non-root user\n");
exit(1);
}
if (uwsgi.gid) {
fprintf(stderr, "cannot setgid() as non-root user\n");
uwsgi_log("cannot setgid() as non-root user\n");
exit(1);
}
if (uwsgi.uid) {
fprintf(stderr, "cannot setuid() as non-root user\n");
uwsgi_log("cannot setuid() as non-root user\n");
exit(1);
}
}
@@ -250,6 +268,7 @@ void uwsgi_close_request(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_r
if (uwsgi->shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1)
get_memusage();
// close the connection with the webserver
if (!wsgi_req->fd_closed) {
close(wsgi_req->poll.fd);
@@ -257,10 +276,10 @@ void uwsgi_close_request(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_r
uwsgi->workers[0].requests++;
uwsgi->workers[uwsgi->mywid].requests++;
// after_request hook
(*uwsgi->shared->after_hooks[wsgi_req->uh.modifier1]) (uwsgi, wsgi_req);
// leave harakiri mode
if (uwsgi->workers[uwsgi->mywid].harakiri > 0) {
set_harakiri(0);
@@ -277,6 +296,7 @@ void uwsgi_close_request(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_r
goodbye_cruel_world();
}
}
void wsgi_req_setup(struct wsgi_request *wsgi_req, int async_id) {
@@ -292,6 +312,7 @@ void wsgi_req_setup(struct wsgi_request *wsgi_req, int async_id) {
wsgi_req->async_waiting_fd = -1;
#endif
wsgi_req->hvec = &uwsgi.async_hvec[wsgi_req->async_id];
wsgi_req->buffer = uwsgi.async_buf[wsgi_req->async_id];
}
@@ -301,7 +322,7 @@ int wsgi_req_recv(struct wsgi_request *wsgi_req) {
if (uwsgi.shared->options[UWSGI_OPTION_LOGGING]) gettimeofday(&wsgi_req->start_of_request, NULL);
if (!uwsgi_parse_response(&wsgi_req->poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, &wsgi_req->buffer)) {
if (!uwsgi_parse_response(&wsgi_req->poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer)) {
return -1;
}
@@ -320,7 +341,7 @@ int wsgi_req_accept(int fd, struct wsgi_request *wsgi_req) {
wsgi_req->poll.fd = accept(fd, (struct sockaddr *) &wsgi_req->c_addr, (socklen_t *) &wsgi_req->c_len);
if (wsgi_req->poll.fd < 0) {
perror("accept()");
uwsgi_error("accept()");
return -1;
}
@@ -346,7 +367,7 @@ void sanitize_args(struct uwsgi_server *uwsgi) {
#ifdef UWSGI_PROFILER
if (uwsgi->enable_profiler) {
fprintf(stderr,"*** Profiler enabled, do not use it in production environment !!! ***\n");
uwsgi_log("*** Profiler enabled, do not use it in production environment !!! ***\n");
uwsgi->async = 1;
}
#endif
@@ -355,7 +376,7 @@ void sanitize_args(struct uwsgi_server *uwsgi) {
#ifdef UWSGI_THREADING
if (uwsgi->ugreen) {
if (uwsgi->has_threads) {
fprintf(stderr,"--- python threads will be disabled in uGreen mode ---\n");
uwsgi_log("--- python threads will be disabled in uGreen mode ---\n");
uwsgi->has_threads = 0;
}
}
@@ -387,7 +408,7 @@ void parse_sys_envs(char **envs, struct option *long_options) {
if (!strncmp(*uenvs, "UWSGI_", 6)) {
earg = malloc(strlen(*uenvs+6)+1);
if (!earg) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
env_to_arg(*uenvs+6, earg);
@@ -423,3 +444,17 @@ void parse_sys_envs(char **envs, struct option *long_options) {
}
}
//use this instead of fprintf to avoid buffering mess with udp logging
void uwsgi_log(const char *fmt, ...) {
va_list ap;
char logpkt[4096];
int rlen;
va_start (ap, fmt);
rlen = vsnprintf(logpkt, 4096, fmt, ap );
va_end(ap);
// do not check for errors
rlen = write(2, logpkt, rlen);
}
+206 -177
View File
File diff suppressed because it is too large Load Diff
+22 -7
View File
@@ -4,13 +4,22 @@
#define UWSGI_VERSION "0.9.5-dev"
#define uwsgi_error(x) uwsgi_log("%s: %s [%s line %d]\n", x, strerror(errno), __FILE__, __LINE__);
#include <stdio.h>
#include <stdlib.h>
#include <signal.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/types.h>
#include <netinet/in.h>
#include <netinet/tcp.h>
#include <stdarg.h>
// linux has not strlcpy
#ifdef __linux
#define strlcpy(x, y, z) strcpy(x, y)
#endif
#ifdef UWSGI_SCTP
#include <netinet/sctp.h>
@@ -58,8 +67,9 @@
#undef _POSIX_C_SOURCE
#endif
#ifdef __sun__
#undef _FILE_OFFSET_BITS
#define WAIT_ANY (-1)
#include <sys/filio.h>
#define PRIO_MAX 20
#endif
#define MAX_PYARGV 10
@@ -160,6 +170,10 @@ PyAPI_FUNC(PyObject *) PyMarshal_ReadObjectFromString(char *, Py_ssize_t);
#ifdef __linux__
#include <endian.h>
#elif __sun__
#include <sys/byteorder.h>
#ifdef _BIG_ENDIAN
#define __BIG_ENDIAN__ 1
#endif
#elif __apple__
#include <libkern/OSByteOrder.h>
#else
@@ -337,8 +351,7 @@ struct wsgi_request {
PyTaskletObject* tasklet;
#endif
// buffer MUST BE THE LAST VAR !!!
char buffer;
char *buffer;
};
struct uwsgi_server {
@@ -369,6 +382,7 @@ struct uwsgi_server {
#endif
struct iovec *async_hvec;
char **async_buf;
struct rlimit rl;
@@ -416,9 +430,6 @@ struct uwsgi_server {
int page_size;
char *sync_page;
int synclog;
char *test_module;
char *pidfile;
@@ -656,6 +667,8 @@ void set_harakiri(int);
#ifdef __BIG_ENDIAN__
uint16_t uwsgi_swap16(uint16_t);
uint32_t uwsgi_swap32(uint32_t);
uint64_t uwsgi_swap64(uint64_t);
#endif
int init_uwsgi_app(PyObject *, PyObject *);
@@ -717,7 +730,7 @@ struct wsgi_request *find_first_available_wsgi_req(struct uwsgi_server *);
struct wsgi_request *find_wsgi_req_by_fd(struct uwsgi_server *, int, int);
struct wsgi_request *find_wsgi_req_by_id(struct uwsgi_server *, int);
struct wsgi_request *next_wsgi_req(struct uwsgi_server *, struct wsgi_request *);
inline struct wsgi_request *next_wsgi_req(struct uwsgi_server *, struct wsgi_request *);
int async_add(int, int , int) ;
int async_mod(int, int , int) ;
@@ -812,3 +825,5 @@ void sanitize_args(struct uwsgi_server *);
void env_to_arg(char *, char *);
void parse_sys_envs(char **, struct option *);
void uwsgi_log(const char *, ...);
+20 -14
View File
@@ -4,27 +4,27 @@
int uwsgi_request_ping(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req) {
char len;
fprintf(stderr, "PING\n");
uwsgi_log( "PING\n");
wsgi_req->uh.modifier2 = 1;
wsgi_req->uh.pktsize = 0;
len = strlen(uwsgi->shared->warning_message);
if (len > 0) {
// endianess check is not needed as the warning message can be max 80 chars
// TODO: check endianess ?
wsgi_req->uh.pktsize = len;
}
if (write(wsgi_req->poll.fd, wsgi_req, 4) != 4) {
perror("write()");
uwsgi_error("write()");
}
if (len > 0) {
if (write(wsgi_req->poll.fd, uwsgi->shared->warning_message, len)
!= len) {
perror("write()");
uwsgi_error("write()");
}
}
return 0;
return UWSGI_OK;
}
/* uwsgi ADMIN|10 */
@@ -33,19 +33,25 @@ int uwsgi_request_admin(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_re
int i;
if (wsgi_req->uh.pktsize >= 4) {
memcpy(&opt_value, &wsgi_req->buffer, 4);
// TODO: check endianess
memcpy(&opt_value, wsgi_req->buffer, 4);
// TODO: check endianess ?
}
fprintf(stderr, "setting internal option %d to %d\n", wsgi_req->uh.modifier2, opt_value);
uwsgi_log( "setting internal option %d to %d\n", wsgi_req->uh.modifier2, opt_value);
uwsgi->shared->options[wsgi_req->uh.modifier2] = opt_value;
// ACK
wsgi_req->uh.modifier1 = 255;
wsgi_req->uh.pktsize = 0;
wsgi_req->uh.modifier2 = 1;
i = write(wsgi_req->poll.fd, wsgi_req, 4);
if (i != 4) {
perror("write()");
uwsgi_error("write()");
}
return UWSGI_OK;
}
@@ -62,7 +68,7 @@ int uwsgi_request_fastfunc(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi
ffunc = PyList_GetItem(uwsgi->fastfuncslist, wsgi_req->uh.modifier2);
if (ffunc) {
fprintf(stderr, "managing fastfunc %d\n", wsgi_req->uh.modifier2);
uwsgi_log( "managing fastfunc %d\n", wsgi_req->uh.modifier2);
return uwsgi_python_call(uwsgi, wsgi_req, ffunc, NULL);
}
@@ -76,7 +82,7 @@ int uwsgi_request_marshal(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_
PyObject *umm = PyDict_GetItemString(uwsgi->embedded_dict,
"message_manager_marshal");
if (umm) {
PyObject *ummo = PyMarshal_ReadObjectFromString(&wsgi_req->buffer,
PyObject *ummo = PyMarshal_ReadObjectFromString(wsgi_req->buffer,
wsgi_req->uh.pktsize);
if (ummo) {
if (!PyTuple_SetItem(uwsgi->embedded_args, 0, ummo)) {
@@ -96,15 +102,15 @@ int uwsgi_request_marshal(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_
PyString_Size(marshalled);
if (write(wsgi_req->poll.fd, wsgi_req, 4) == 4) {
if (write(wsgi_req->poll.fd, PyString_AsString(marshalled), wsgi_req->uh.pktsize) != wsgi_req->uh.pktsize) {
perror("write()");
uwsgi_error("write()");
}
}
else {
perror("write()");
uwsgi_error("write()");
}
}
else {
fprintf(stderr, "marshalled object is too big. skip\n");
uwsgi_log( "marshalled object is too big. skip\n");
}
Py_DECREF(marshalled);
}
+20 -20
View File
@@ -9,12 +9,12 @@ extern struct uwsgi_server uwsgi;
#ifdef __APPLE__
#define UWSGI_LOCK OSSpinLockLock((OSSpinLock *) uwsgi.sharedareamutex);
#define UWSGI_UNLOCK OSSpinLockUnlock((OSSpinLock *) uwsgi.sharedareamutex);
#elif defined(__linux__)
#elif defined(__linux__) || defined(__sun__) || defined(__FreeBSD__)
#define UWSGI_LOCK pthread_mutex_lock((pthread_mutex_t *) uwsgi.sharedareamutex + sizeof(pthread_mutexattr_t));
#define UWSGI_UNLOCK pthread_mutex_unlock((pthread_mutex_t *) uwsgi.sharedareamutex + sizeof(pthread_mutexattr_t));
#else
#define UWSGI_LOCK if (flock(uwsgi.serverfd, LOCK_EX)) { perror("flock()"); }
#define UWSGI_UNLOCK if (flock(uwsgi.serverfd, LOCK_UN)) { perror("flock()"); }
#define UWSGI_LOCK if (flock(uwsgi.serverfd, LOCK_EX)) { uwsgi_error("flock()"); }
#define UWSGI_UNLOCK if (flock(uwsgi.serverfd, LOCK_UN)) { uwsgi_error("flock()"); }
#endif
#define UWSGI_LOGBASE "[- uWSGI -"
@@ -51,7 +51,7 @@ PyObject *py_uwsgi_warning(PyObject * self, PyObject * args) {
len = strlen(message);
if (len > 80) {
fprintf(stderr, "- warning message must be max 80 chars, it will be truncated -");
uwsgi_log( "- warning message must be max 80 chars, it will be truncated -");
memcpy(uwsgi.shared->warning_message, message, 80);
uwsgi.shared->warning_message[80] = 0;
}
@@ -74,10 +74,10 @@ PyObject *py_uwsgi_log(PyObject * self, PyObject * args) {
tt = time(NULL);
if (logline[strlen(logline)] != '\n') {
fprintf(stderr, UWSGI_LOGBASE " %.*s] %s\n", 24, ctime(&tt), logline);
uwsgi_log( UWSGI_LOGBASE " %.*s] %s\n", 24, ctime(&tt), logline);
}
else {
fprintf(stderr, UWSGI_LOGBASE " %.*s] %s", 24, ctime(&tt), logline);
uwsgi_log( UWSGI_LOGBASE " %.*s] %s", 24, ctime(&tt), logline);
}
Py_INCREF(Py_True);
@@ -430,7 +430,7 @@ PyObject *py_uwsgi_send_multi_message(PyObject * self, PyObject * args) {
clen = PyTuple_Size(arg_cluster);
multipoll = malloc(clen * sizeof(struct pollfd));
if (!multipoll) {
perror("malloc");
uwsgi_error("malloc");
Py_INCREF(Py_None);
return Py_None;
}
@@ -438,7 +438,7 @@ PyObject *py_uwsgi_send_multi_message(PyObject * self, PyObject * args) {
buffer = malloc(uwsgi.buffer_size * clen);
if (!buffer) {
perror("malloc");
uwsgi_error("malloc");
free(multipoll);
Py_INCREF(Py_None);
return Py_None;
@@ -493,11 +493,11 @@ PyObject *py_uwsgi_send_multi_message(PyObject * self, PyObject * args) {
while (managed < clen) {
pret = poll(multipoll, clen, PyInt_AsLong(arg_timeout) * 1000);
if (pret < 0) {
perror("poll()");
uwsgi_error("poll()");
goto megamulticlear;
}
else if (pret == 0) {
fprintf(stderr, "timeout on multiple send !\n");
uwsgi_log( "timeout on multiple send !\n");
goto megamulticlear;
}
else {
@@ -579,15 +579,15 @@ PyObject *py_uwsgi_load_plugin(PyObject * self, PyObject * args) {
plugin_handle = dlopen(plugin_name, RTLD_NOW | RTLD_GLOBAL);
if (!plugin_handle) {
fprintf(stderr, "%s\n", dlerror());
uwsgi_log( "%s\n", dlerror());
}
else {
plugin_init = dlsym(plugin_handle, "uwsgi_init");
if (plugin_init) {
if ((*plugin_init) (&uwsgi, pargs)) {
fprintf(stderr, "plugin initialization returned error\n");
uwsgi_log( "plugin initialization returned error\n");
if (dlclose(plugin_handle)) {
fprintf(stderr, "unable to unload plugin\n");
uwsgi_log( "unable to unload plugin\n");
}
Py_INCREF(Py_None);
@@ -607,7 +607,7 @@ PyObject *py_uwsgi_load_plugin(PyObject * self, PyObject * args) {
}
else {
fprintf(stderr, "%s\n", dlerror());
uwsgi_log( "%s\n", dlerror());
}
}
@@ -797,7 +797,7 @@ PyObject *py_uwsgi_workers(PyObject * self, PyObject * args) {
PyObject *py_uwsgi_reload(PyObject * self, PyObject * args) {
if (kill(uwsgi.workers[0].pid, SIGHUP)) {
perror("kill()");
uwsgi_error("kill()");
Py_INCREF(Py_None);
return Py_None;
}
@@ -830,7 +830,7 @@ PyObject *py_uwsgi_worker_id(PyObject * self, PyObject * args) {
}
PyObject *py_uwsgi_disconnect(PyObject * self, PyObject * args) {
fprintf(stderr, "detaching uWSGI from current connection...\n");
uwsgi_log( "detaching uWSGI from current connection...\n");
struct wsgi_request *wsgi_req = current_wsgi_req(&uwsgi);
@@ -900,13 +900,13 @@ void init_uwsgi_module_spooler(PyObject * current_uwsgi_module) {
uwsgi_module_dict = PyModule_GetDict(current_uwsgi_module);
if (!uwsgi_module_dict) {
fprintf(stderr, "could not get uwsgi module __dict__\n");
uwsgi_log( "could not get uwsgi module __dict__\n");
exit(1);
}
spool_buffer = malloc(uwsgi.buffer_size);
if (!spool_buffer) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
@@ -925,7 +925,7 @@ void init_uwsgi_module_advanced(PyObject * current_uwsgi_module) {
uwsgi_module_dict = PyModule_GetDict(current_uwsgi_module);
if (!uwsgi_module_dict) {
fprintf(stderr, "could not get uwsgi module __dict__\n");
uwsgi_log( "could not get uwsgi module __dict__\n");
exit(1);
}
@@ -942,7 +942,7 @@ void init_uwsgi_module_sharedarea(PyObject * current_uwsgi_module) {
uwsgi_module_dict = PyModule_GetDict(current_uwsgi_module);
if (!uwsgi_module_dict) {
fprintf(stderr, "could not get uwsgi module __dict__\n");
uwsgi_log( "could not get uwsgi module __dict__\n");
exit(1);
}
+7 -1
View File
@@ -136,6 +136,8 @@ def unbit_setup():
def parse_vars():
global UGREEN
version = sys.version_info
uver = "%d.%d" % (version[0], version[1])
@@ -144,14 +146,18 @@ def parse_vars():
if str(PYLIB_PATH) != '':
ldflags.insert(0,'-L' + PYLIB_PATH)
kvm_list = ['SunOS', 'FreeBSD', 'OpenBSD', 'NetBSD', 'DragonFly']
kvm_list = ['FreeBSD', 'OpenBSD', 'NetBSD', 'DragonFly']
if uwsgi_os == 'SunOS':
ldflags.append('-lsendfile')
ldflags.remove('-rdynamic')
if uwsgi_os in kvm_list:
ldflags.append('-lkvm')
if uwsgi_os == 'OpenBSD':
UGREEN = False
if EMBEDDED:
cflags.append('-DUWSGI_EMBEDDED')
gcc_list.append('uwsgi_pymodule')
+5 -5
View File
@@ -98,12 +98,12 @@ int uwsgi_request_wsgi(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
/* Standard WSGI request */
if (!wsgi_req->uh.pktsize) {
fprintf(stderr, "Invalid WSGI request. skip.\n");
uwsgi_log( "Invalid WSGI request. skip.\n");
return -1;
}
if (uwsgi_parse_vars(uwsgi, wsgi_req)) {
fprintf(stderr,"Invalid WSGI request. skip.\n");
uwsgi_log("Invalid WSGI request. skip.\n");
return -1;
}
@@ -158,13 +158,13 @@ int uwsgi_request_wsgi(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
if (wsgi_req->protocol_len < 5) {
fprintf(stderr, "INVALID PROTOCOL: %.*s", wsgi_req->protocol_len, wsgi_req->protocol);
uwsgi_log( "INVALID PROTOCOL: %.*s", wsgi_req->protocol_len, wsgi_req->protocol);
internal_server_error(wsgi_req->poll.fd, "invalid HTTP protocol !!!");
goto clear;
}
if (strncmp(wsgi_req->protocol, "HTTP/", 5)) {
fprintf(stderr, "INVALID PROTOCOL: %.*s", wsgi_req->protocol_len, wsgi_req->protocol);
uwsgi_log( "INVALID PROTOCOL: %.*s", wsgi_req->protocol_len, wsgi_req->protocol);
internal_server_error(wsgi_req->poll.fd, "invalid HTTP protocol !!!");
goto clear;
}
@@ -183,7 +183,7 @@ int uwsgi_request_wsgi(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req
for (i = 0; i < wsgi_req->var_cnt; i += 2) {
//fprintf(stderr,"%.*s: %.*s\n", wsgi_req->hvec[i].iov_len, wsgi_req->hvec[i].iov_base, wsgi_req->hvec[i+1].iov_len, wsgi_req->hvec[i+1].iov_base);
//uwsgi_log("%.*s: %.*s\n", wsgi_req->hvec[i].iov_len, wsgi_req->hvec[i].iov_base, wsgi_req->hvec[i+1].iov_len, wsgi_req->hvec[i+1].iov_base);
pydictkey = PyString_FromStringAndSize(wsgi_req->hvec[i].iov_base, wsgi_req->hvec[i].iov_len);
pydictvalue = PyString_FromStringAndSize(wsgi_req->hvec[i + 1].iov_base, wsgi_req->hvec[i + 1].iov_len);
PyDict_SetItem(wsgi_req->async_environ, pydictkey, pydictvalue);
+5 -5
View File
@@ -27,7 +27,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) {
}
if (!PyString_Check(head)) {
fprintf(stderr, "http status must be a string !\n");
uwsgi_log( "http status must be a string !\n");
goto clear;
}
@@ -94,7 +94,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) {
goto clear;
}
if (!PyList_Check(headers)) {
fprintf(stderr, "http headers must be in a python list\n");
uwsgi_log( "http headers must be in a python list\n");
goto clear;
}
wsgi_req->header_cnt = PyList_Size(headers);
@@ -111,7 +111,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) {
goto clear;
}
if (!PyTuple_Check(head)) {
fprintf(stderr, "http header must be defined in a tuple !\n");
uwsgi_log( "http header must be defined in a tuple !\n");
goto clear;
}
h_key = PyTuple_GetItem(head, 0);
@@ -140,7 +140,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) {
#endif
wsgi_req->hvec[j + 3].iov_base = nl;
wsgi_req->hvec[j + 3].iov_len = NL_SIZE;
//fprintf(stderr, "%.*s: %.*s\n", wsgi_req->hvec[j].iov_len, (char *)wsgi_req->hvec[j].iov_base, wsgi_req->hvec[j+2].iov_len, (char *) wsgi_req->hvec[j+2].iov_base);
//uwsgi_log( "%.*s: %.*s\n", wsgi_req->hvec[j].iov_len, (char *)wsgi_req->hvec[j].iov_base, wsgi_req->hvec[j+2].iov_len, (char *) wsgi_req->hvec[j+2].iov_base);
}
@@ -151,7 +151,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) {
wsgi_req->headers_size = writev(wsgi_req->poll.fd, wsgi_req->hvec, j + 1);
if (wsgi_req->headers_size < 0) {
perror("writev()");
uwsgi_error("writev()");
}
Py_INCREF(uwsgi.wsgi_writeout);
+18 -18
View File
@@ -23,21 +23,21 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
doc = xmlReadFile(uwsgi.xml_config, NULL, 0);
if (doc == NULL) {
fprintf(stderr, "[uWSGI] could not parse file %s.\n", uwsgi.xml_config);
uwsgi_log( "[uWSGI] could not parse file %s.\n", uwsgi.xml_config);
exit(1);
}
if (long_options) {
fprintf(stderr, "[uWSGI] parsing config file %s\n", uwsgi.xml_config);
uwsgi_log( "[uWSGI] parsing config file %s\n", uwsgi.xml_config);
}
element = xmlDocGetRootElement(doc);
if (element == NULL) {
fprintf(stderr, "[uWSGI] invalid xml config file.\n");
uwsgi_log( "[uWSGI] invalid xml config file.\n");
exit(1);
}
if (strcmp((char *) element->name, "uwsgi")) {
fprintf(stderr, "[uWSGI] invalid xml root element, <uwsgi> expected.\n");
uwsgi_log( "[uWSGI] invalid xml root element, <uwsgi> expected.\n");
exit(1);
}
@@ -52,13 +52,13 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
break;
if (!strcmp((char *) node->name, aopt->name)) {
if (!node->children && aopt->has_arg) {
fprintf(stderr, "[uWSGI] %s option need a value. skip.\n", aopt->name);
uwsgi_log( "[uWSGI] %s option need a value. skip.\n", aopt->name);
exit(1);
}
if (node->children) {
if (!node->children->content && aopt->has_arg) {
fprintf(stderr, "[uWSGI] %s option need a value. skip.\n", aopt->name);
uwsgi_log( "[uWSGI] %s option need a value. skip.\n", aopt->name);
exit(1);
}
}
@@ -89,7 +89,7 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
if (!strcmp((char *) node->name, "app")) {
xml_uwsgi_mountpoint = xmlGetProp(node, (const xmlChar *) "mountpoint");
if (xml_uwsgi_mountpoint == NULL) {
fprintf(stderr, "no mountpoint defined for app. skip.\n");
uwsgi_log( "no mountpoint defined for app. skip.\n");
continue;
}
wsgi_req->script_name = (char *) xml_uwsgi_mountpoint;
@@ -99,12 +99,12 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
if (node2->type == XML_ELEMENT_NODE) {
if (!strcmp((char *) node2->name, "script")) {
if (!node2->children) {
fprintf(stderr, "no wsgi script defined. skip.\n");
uwsgi_log( "no wsgi script defined. skip.\n");
continue;
}
xml_uwsgi_script = node2->children->content;
if (xml_uwsgi_script == NULL) {
fprintf(stderr, "no wsgi script defined. skip.\n");
uwsgi_log( "no wsgi script defined. skip.\n");
continue;
}
wsgi_req->wsgi_script = (char *) xml_uwsgi_script;
@@ -145,7 +145,7 @@ void uwsgi_endElement(void *userData, const char *name) {
}
else if (current_xmlnode_has_arg) {
if (!current_xmlnode_text_len) {
fprintf(stderr,"option %s requires an argument\n", name);
uwsgi_log("option %s requires an argument\n", name);
exit(1);
}
// HACK: use the first char of closing tag for nulling string
@@ -188,9 +188,9 @@ void uwsgi_startApp(void *userData, const char *name, const char **attrs) {
if (!strcmp(name, "app")) {
current_xmlnode = 0 ;
fprintf(stderr,"%s = %s\n", attrs[0], attrs[1]);
uwsgi_log("%s = %s\n", attrs[0], attrs[1]);
if (strcmp(attrs[0], "mountpoint")) {
fprintf(stderr,"invalid attribute for app tag. must be 'mountpoint'\n");
uwsgi_log("invalid attribute for app tag. must be 'mountpoint'\n");
exit(1);
}
if (attrs[1]) {
@@ -204,7 +204,7 @@ void uwsgi_startApp(void *userData, const char *name, const char **attrs) {
}
else if (!strcmp(name, "script")) {
if (!wsgi_req->script_name_len) {
fprintf(stderr,"you have not specified an app mountpoint.\n");
uwsgi_log("you have not specified an app mountpoint.\n");
exit(1);
}
current_xmlnode = 1 ;
@@ -249,24 +249,24 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
xmlfd = open(uwsgi.xml_config, O_RDONLY);
if (xmlfd < 0) {
perror("open()");
uwsgi_error("open()");
exit(1);
}
if (fstat(xmlfd, &stat_buf)) {
perror("fstat()");
uwsgi_error("fstat()");
exit(1);
}
xmlbuf = malloc(stat_buf.st_size);
if (!xmlbuf) {
perror("malloc()");
uwsgi_error("malloc()");
exit(1);
}
rlen = read(xmlfd, xmlbuf, stat_buf.st_size);
if (rlen != stat_buf.st_size) {
perror("read()");
uwsgi_error("read()");
exit(1);
}
close(xmlfd);
@@ -283,7 +283,7 @@ void uwsgi_xml_config(struct wsgi_request *wsgi_req, struct option *long_options
}
if (!XML_Parse(parser, xmlbuf, stat_buf.st_size, 1)) {
fprintf(stderr, "%s at line %d\n", XML_ErrorString(XML_GetErrorCode(parser)), (int) XML_GetCurrentLineNumber(parser));
uwsgi_log( "%s at line %d\n", XML_ErrorString(XML_GetErrorCode(parser)), (int) XML_GetCurrentLineNumber(parser));
exit(1);
}