Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 9 additions & 7 deletions src/net/capture_engine.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ void CaptureEngine::processPackets(int worker_id, mqueue<Packet>* work_queue) {
MemcacheCommand mc = parse(packet);
if (mc.isResponse()) {
enqueue(mc);
resCount += 1;
resCount += mc.getObjectNumber();
#ifdef _DEBUG
logger->trace(CONTEXT,
"worker %d, packet %ld, key %s", worker_id, packet.id(),
Expand Down Expand Up @@ -172,14 +172,16 @@ MemcacheCommand CaptureEngine::parse(const Packet& packet) const
*/
void CaptureEngine::enqueue(const MemcacheCommand& mc)
{
Elem e(mc.getObjectKey(), mc.getObjectSize());
for (uint32_t idx = 0; idx < mc.getObjectNumber(); ++idx) {
Elem e(mc.getObjectKey(idx), mc.getObjectSize(idx));
#ifdef _DEBUG
logger->trace(CONTEXT,
"Produced stat: %s, %d", e.first.c_str(), e.second);
logger->trace(CONTEXT,
"Produced stat: %s, %d", e.first.c_str(), e.second);
#endif
barrier_lock.lock();
barrier->produce(e);
barrier_lock.unlock();
barrier_lock.lock();
barrier->produce(e);
barrier_lock.unlock();
}
}

} // end namespace
145 changes: 112 additions & 33 deletions src/net/memcache_command.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ extern "C" {
#include <netinet/ip.h>
#include <netinet/tcp.h>
#include <netinet/if_ether.h>
#include "protocol_binary.h"
}

static inline std::string ipv4addressToString(const void * src) {
Expand All @@ -27,14 +28,6 @@ using namespace std;
MemcacheCommand MemcacheCommand::create(const Packet& pkt,
const bpf_u_int32 captureAddress)
{
static ssize_t ether_header_sz = sizeof(struct ether_header);
static ssize_t ip_sz = sizeof(struct ip);
static ssize_t tcphdr_sz = sizeof(struct tcphdr);

const struct ether_header* ethernetHeader;
const struct ip* ipHeader;
const struct tcphdr* tcpHeader;

const Packet::Header* pkthdr = &pkt.getHeader();
const Packet::Data* packet = pkt.getData();

Expand All @@ -46,14 +39,16 @@ MemcacheCommand MemcacheCommand::create(const Packet& pkt,

// must be an IP packet
// TODO add support for dumping localhost
ethernetHeader = (struct ether_header*)packet;
const struct ether_header* ethernetHeader = (struct ether_header*)packet;
ssize_t ethernetHeaderSize = sizeof(struct ether_header); //14 bytes
auto etype = ntohs(ethernetHeader->ether_type);
if (etype != ETHERTYPE_IP) {
return MemcacheCommand();
}

// must be TCP - TODO add support for UDP
ipHeader = (struct ip*)(packet + ether_header_sz);
const struct ip* ipHeader = (struct ip*)(packet + ethernetHeaderSize);
ssize_t ipHeaderSize = ipHeader->ip_hl * 4;
auto itype = ipHeader->ip_p;
if (itype != IPPROTO_TCP) {
return MemcacheCommand();
Expand All @@ -69,39 +64,56 @@ MemcacheCommand MemcacheCommand::create(const Packet& pkt,
// FIXME will remove once we add back the direction parsing
(void)possible_request;

tcpHeader = (struct tcphdr*)(packet + ether_header_sz + ip_sz);
const struct tcphdr* tcpHeader = (struct tcphdr*)(packet + ethernetHeaderSize
+ ipHeaderSize);
ssize_t tcpHeaderSize = tcpHeader->doff * 4;
(void)tcpHeader;
data = (u_char*)(packet + ether_header_sz + ip_sz + tcphdr_sz);
dataLength = pkthdr->len - (ether_header_sz + ip_sz + tcphdr_sz);
data = (u_char*)(packet + ethernetHeaderSize + ipHeaderSize + tcpHeaderSize);
dataLength = pkthdr->len - (ethernetHeaderSize + ipHeaderSize + tcpHeaderSize);
if (dataLength > pkthdr->caplen) {
dataLength = pkthdr->caplen;
}

// TODO revert to detecting request/response and doing the right thing
if (dataLength <= 0) {
return MemcacheCommand();
}
return MemcacheCommand::makeResponse(data, dataLength, sourceAddress);
}

// protected default constructor
MemcacheCommand::MemcacheCommand()
: cmdType_(MC_UNKNOWN),
sourceAddress_(),
commandName_(),
objectKey_(),
objectSize_(0)
commandName_()
{}

// protected constructor
MemcacheCommand::MemcacheCommand(const memcache_command_t cmdType,
const string sourceAddress,
const string commandName)
: cmdType_(cmdType),
sourceAddress_(sourceAddress),
commandName_(commandName)
{}

MemcacheCommand::MemcacheCommand(const memcache_command_t cmdType,
const string sourceAddress,
const string commandName,
const string objectKey,
uint32_t objectSize)
: cmdType_(cmdType),
sourceAddress_(sourceAddress),
commandName_(commandName),
objectKey_(objectKey),
objectSize_(objectSize)
{}
commandName_(commandName)
{
pushObject(objectKey, objectSize);
}

void MemcacheCommand::pushObject(const std::string objectKey, uint32_t objectSize)
{
objectKeyList_.push_back(objectKey);
objectSizeList_.push_back(objectSize);
}

// static protected
MemcacheCommand MemcacheCommand::makeRequest(u_char*, int, string)
Expand All @@ -110,30 +122,97 @@ MemcacheCommand MemcacheCommand::makeRequest(u_char*, int, string)
return MemcacheCommand();
}

// static protected
MemcacheCommand MemcacheCommand::makeResponse(u_char *data, int length,
string sourceAddress)
MemcacheCommand MemcacheCommand::parseAsciiResponse(u_char *data, int length,
string sourceAddress)
{
static pcrecpp::RE re("VALUE (\\S+) \\d+ (\\d+)",
static pcrecpp::RE re("(VALUE (\\S+) \\d+ (\\d+))",
pcrecpp::RE_Options(PCRE_MULTILINE));
static int minimum_length = 11; // 'VALUE a 0 1'

string whole;
string key;
int size = -1;
string input = "";
for (int i = 0; i < length; i++) {
int cid = (int)data[i];
if (isprint(cid) || cid == 10 || cid == 13) {
input += (char)data[i];

MemcacheCommand mc(MC_RESPONSE, sourceAddress, "");
int offset = 0;
while (length - offset >= minimum_length) {
if (!re.PartialMatch(data + offset, &whole, &key, &size)) {
break;
}
if (size >= 0) {
mc.pushObject(key, size);
offset += whole.length() + 2 + size + 2; // 2 for '\r\n', 2 for '\r\n'
} else {
break;
}
}
if (input.length() < 11) {

if (mc.getObjectNumber() > 0) {
return mc;
} else {
return MemcacheCommand();
}
re.PartialMatch(input, &key, &size);
if (size >= 0) {
return MemcacheCommand(MC_RESPONSE, sourceAddress, "", key, size);
}

MemcacheCommand MemcacheCommand::parseBinaryResponse(u_char *data, int length,
string sourceAddress)
{
static const int HEADER_LENGTH = sizeof(protocol_binary_response_header);
static const int MINIMUM_LENGTH = HEADER_LENGTH;

string key;
int valuelen = 0;
MemcacheCommand mc(MC_RESPONSE, sourceAddress, "");
int offset = 0;
while (length - offset >= MINIMUM_LENGTH) {
protocol_binary_response_header header;
memcpy((char*)&header, data + offset, HEADER_LENGTH);
header.response.status = ntohs(header.response.status);
header.response.keylen = ntohs(header.response.keylen);
header.response.bodylen = ntohl(header.response.bodylen);
valuelen = header.response.bodylen - header.response.extlen
- header.response.keylen;

if (header.response.status == PROTOCOL_BINARY_RESPONSE_SUCCESS ||
header.response.status == PROTOCOL_BINARY_RESPONSE_AUTH_CONTINUE) {
switch (header.response.opcode) {
case PROTOCOL_BINARY_CMD_GET:
case PROTOCOL_BINARY_CMD_GETQ:
Logger::getLogger("command")->debug(CONTEXT, "get reponse without key");
break;
case PROTOCOL_BINARY_CMD_GETK:
case PROTOCOL_BINARY_CMD_GETKQ:
key.assign((char*)data + offset + HEADER_LENGTH + header.response.extlen,
header.response.keylen);
mc.pushObject(key, valuelen);
break;
default:
break;
}
}

offset += HEADER_LENGTH + header.response.bodylen;
}

if (mc.getObjectNumber() > 0) {
return mc;
} else {
return MemcacheCommand();
}
}

// static protected
MemcacheCommand MemcacheCommand::makeResponse(u_char *data, int length,
string sourceAddress)
{
if (length < 1) {
return MemcacheCommand();
}
if (data[0] == PROTOCOL_BINARY_RES) {
return parseBinaryResponse(data, length, sourceAddress);
} else {
return parseAsciiResponse(data, length, sourceAddress);
}
}

} // end namespace
24 changes: 20 additions & 4 deletions src/net/memcache_command.h
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,17 @@ class MemcacheCommand
// only when isRequest is true
std::string getCommandName() const { return commandName_; }

ssize_t getObjectNumber() const { return objectKeyList_.size(); }

// sometimes when isResponse is true, sometimes when isRequest is true
std::string getObjectKey() const { return objectKey_; }
std::string getObjectKey(uint32_t idx = 0) const {
return (idx < objectKeyList_.size() ? objectKeyList_[idx] : "");
}

// only when isResponse is true
uint32_t getObjectSize() const { return objectSize_; }
uint32_t getObjectSize(uint32_t idx = 0) const {
return (idx < objectSizeList_.size() ? objectSizeList_[idx] : 0);
}

// source address for request
std::string getSourceAddress() const { return sourceAddress_; }
Expand All @@ -47,24 +53,34 @@ class MemcacheCommand
// Default constructor is protected
MemcacheCommand();

MemcacheCommand(const memcache_command_t cmdType,
const std::string sourceAddress,
const std::string commandName);

MemcacheCommand(const memcache_command_t cmdType,
const std::string sourceAddress,
const std::string commandName,
const std::string objectKey,
uint32_t objectSize);

void pushObject(const std::string objectKey, uint32_t objectSize);

static MemcacheCommand makeRequest(u_char *data,
int dataLength,
std::string sourceAddress);
static MemcacheCommand makeResponse(u_char *data,
int dataLength,
std::string sourceAddress);
static MemcacheCommand parseAsciiResponse(u_char *data, int length,
std::string sourceAddress);
static MemcacheCommand parseBinaryResponse(u_char *data, int length,
std::string sourceAddress);

const memcache_command_t cmdType_;
const std::string sourceAddress_;
const std::string commandName_;
const std::string objectKey_;
const uint32_t objectSize_;
std::vector<std::string> objectKeyList_;
std::vector<uint32_t> objectSizeList_;
};

} // end namespace
Expand Down
Loading