mooncached.lua

1
local server = require "net.server";
2
local log = require "util.logger".init("memcached");
3
local memcache = require "util.memcache";
4
 
5
local cache = memcache.new();
6
 
7
local memcached_listener = {};
8
local command_handlers = {};
9
 
10
--- Network handlers
11
 
12
function memcached_listener.onconnect(conn)
13
end
14
 
15
function memcached_listener.onincoming(conn, line)
16
	local command, params_pos = line:match("^(%S+) ?()");
17
	local command_handler = command_handlers[command];
18
	if command_handler then
19
		local ok, err = command_handler(conn, line:sub(params_pos));
20
		if ok == false then
21
			conn:write("CLIENT_ERROR "..err.."\r\n");
22
		end
23
	elseif command then
24
		log("warn", "Client sent unknown command: %s", command);
25
		conn:write("ERROR\r\n");
26
	end
27
end
28
 
29
function memcached_listener.ondisconnect(conn, err)
30
end
31
 
32
--- Command handlers
33
 
34
local function generic_store_command(store_method, conn, params)
35
	local key, flags, exptime, bytes, reply = params:match("(%S+) (%d+) (%d+) (%d+) ?(.*)$");
36
	flags, exptime, bytes, reply = tonumber(flags), tonumber(exptime), tonumber(bytes), reply ~= "noreply";
37
	if not (flags and exptime and bytes) then
38
		if reply then
39
			return false, "Invalid parameter(s)";
40
		else
41
			return nil;
42
		end
43
	end
44
	conn:set_mode("*a");
45
	local received_count, received_buffer = 0, {};
46
	local function handle_data(conn, data)
47
		log("debug", "Received data of length "..#data.." out of "..bytes);
48
		received_count = received_count + #data;
49
		received_buffer[#received_buffer+1] = data;
50
		if received_count >= bytes then
51
			received_buffer = table.concat(received_buffer);
52
			local ok, err = store_method(cache, key, flags, exptime, received_buffer:sub(1,bytes));
53
			if reply then
54
				if ok then
55
					if err == true then
56
						conn:send("STORED\r\n");
57
					else
58
						conn:send("NOT_STORED\r\n");
59
					end
60
				else
61
					conn:send("SERVER_ERROR "..(err or "Unknown error").."\r\n");
62
				end
63
			end
64
			conn:setlistener(memcached_listener);
65
			conn:set_mode("*l");
66
			if received_count > bytes then
67
				log("debug", "Re-handling %d extra bytes", received_count-bytes);
68
				memcached_listener.onincoming(conn, received_buffer:sub(bytes+1));
69
			end
70
		end
71
	end
72
	conn:setlistener({
73
		onincoming = handle_data;
74
		ondisconnect = memcached_listener.ondisconnect;
75
	});
76
	log("debug", "Waiting for "..bytes.." bytes from client");
77
	return true;
78
end
79
 
80
function command_handlers.set(conn, params)
81
	return generic_store_command(cache.set, conn, params);
82
end
83
 
84
function command_handlers.add(conn, params)
85
	return generic_store_command(cache.add, conn, params);
86
end
87
 
88
function command_handlers.replace(conn, params)
89
	return generic_store_command(cache.replace, conn, params);
90
end
91
 
92
local function generic_increment_decrement_command(method, conn, params)
93
	local key, amount = params:match("^(%S+) (%d+)");
94
	local reply = params:match(" (noreply)$") ~= "noreply";
95
	amount = tonumber(amount);
96
	local ok, err;
97
	if not (key and amount) then
98
		if reply then
99
			return false, "Invalid parameter(s)";
100
		else
101
			return nil;
102
		end
103
	end
104
	local ok, new_value = method(cache, key, amount);
105
	if ok and reply then
106
		conn:write(new_value.."\r\n");
107
	end
108
	return true;
109
end
110
 
111
function command_handlers.incr(conn, params)
112
	return generic_increment_decrement_command(cache.incr, conn, params);
113
end
114
 
115
function command_handlers.decr(conn, params)
116
	return generic_increment_decrement_command(cache.decr, conn, params);
117
end
118
 
119
function command_handlers.get(conn, keys)
120
	for key in keys:gmatch("%S+") do
121
		local flags, data = cache:get(key);
122
		if data then
123
			conn:write("VALUE "..key.." "..flags.." "..#data.."\r\n"..data.."\r\n");
124
		end
125
	end
126
	conn:write("END\r\n");
127
	return true;
128
end
129
 
130
function command_handlers.delete(conn, params)
131
	local key, keyend = params:match("^(%S+)()");
132
	local time, reply = params:match(" (%d+)", keyend), params:match(" (noreply)$", keyend);
133
	time, reply = tonumber(time), reply ~= "noreply";
134
	local ok, err;
135
	if not key then
136
		ok, err = false, "Unable to determine key from request";
137
	else
138
		ok, err = cache:delete(key, time);
139
		if ok then
140
			if err then
141
				conn:write("DELETED\r\n");
142
			else
143
				conn:write("NOT_FOUND\r\n");
144
			end
145
		end
146
	end
147
	if not reply then
148
		return nil;
149
	end
150
	return ok, err;
151
end
152
 
153
function command_handlers.version(conn)
154
	conn:write("VERSION Mooncached 0.1\r\n");
155
	return true;
156
end
157
 
158
function command_handlers.quit(conn)
159
	conn:close();
160
	return true;
161
end
162
 
163
logger.setwriter(function (name, level, format, ...) return print(name, level, format:format(...)); end);
164
 
165
server.addserver("*", 11211, memcached_listener, "*l");
166
 
167
server.loop();