| 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(); |