-
Notifications
You must be signed in to change notification settings - Fork 504
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
add list parse for rdb &pass gtid.lwm & add unknow command parse (#769)
* pass gtid.lwm & add unknow command parse * add unknown op lwm --------- Co-authored-by: hailu <[email protected]>
- Loading branch information
1 parent
b903aab
commit aa41c04
Showing
9 changed files
with
387 additions
and
172 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
125 changes: 125 additions & 0 deletions
125
...s/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbListParser.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,125 @@ | ||
package com.ctrip.xpipe.redis.core.redis.rdb.parser; | ||
|
||
import com.ctrip.xpipe.redis.core.redis.exception.RdbParseEmptyKeyException; | ||
import com.ctrip.xpipe.redis.core.redis.operation.RedisOpType; | ||
import com.ctrip.xpipe.redis.core.redis.operation.op.RedisOpSingleKey; | ||
import com.ctrip.xpipe.redis.core.redis.rdb.RdbLength; | ||
import com.ctrip.xpipe.redis.core.redis.rdb.RdbParseContext; | ||
import com.ctrip.xpipe.redis.core.redis.rdb.RdbParser; | ||
import io.netty.buffer.ByteBuf; | ||
import org.slf4j.Logger; | ||
import org.slf4j.LoggerFactory; | ||
|
||
/** | ||
* @author hailu | ||
* @date 2024/1/17 19:06 | ||
*/ | ||
public class RdbListParser extends AbstractRdbParser<Integer> implements RdbParser<Integer> { | ||
|
||
private RdbParseContext context; | ||
|
||
private RdbParser<byte[]> rdbStringParser; | ||
|
||
private RdbLength len; | ||
|
||
private int readCnt; | ||
|
||
private STATE state = STATE.READ_INIT; | ||
|
||
private static final Logger logger = LoggerFactory.getLogger(RdbListParser.class); | ||
|
||
enum STATE { | ||
READ_INIT, | ||
READ_LEN, | ||
READ_VALUE, | ||
READ_END | ||
} | ||
|
||
public RdbListParser(RdbParseContext parseContext) { | ||
this.context = parseContext; | ||
this.rdbStringParser = (RdbParser<byte[]>) context.getOrCreateParser(RdbParseContext.RdbType.STRING); | ||
} | ||
|
||
@Override | ||
public Integer read(ByteBuf byteBuf) { | ||
|
||
while (!isFinish() && byteBuf.readableBytes() > 0) { | ||
|
||
switch (state) { | ||
|
||
case READ_INIT: | ||
len = null; | ||
readCnt = 0; | ||
state = STATE.READ_LEN; | ||
break; | ||
|
||
case READ_LEN: | ||
len = parseRdbLength(byteBuf); | ||
if (null != len) { | ||
if (len.getLenValue() > 0) { | ||
state = STATE.READ_VALUE; | ||
} else { | ||
throw new RdbParseEmptyKeyException("set key " + context.getKey()); | ||
} | ||
} | ||
break; | ||
|
||
case READ_VALUE: | ||
byte[] value = rdbStringParser.read(byteBuf); | ||
if (null != value) { | ||
rdbStringParser.reset(); | ||
propagateCmdIfNeed(value); | ||
|
||
readCnt++; | ||
if (readCnt >= len.getLenValue()) { | ||
state = STATE.READ_END; | ||
} else { | ||
state = STATE.READ_VALUE; | ||
} | ||
} | ||
break; | ||
|
||
case READ_END: | ||
default: | ||
|
||
} | ||
|
||
if (isFinish()) { | ||
propagateExpireAtIfNeed(context.getKey(), context.getExpireMilli()); | ||
} | ||
} | ||
|
||
if (isFinish()) return len.getLenValue(); | ||
else return null; | ||
} | ||
|
||
private void propagateCmdIfNeed(byte[] value) { | ||
if (null == value || null == context.getKey()) { | ||
return; | ||
} | ||
|
||
notifyRedisOp(new RedisOpSingleKey( | ||
RedisOpType.RPUSH, | ||
new byte[][] {RedisOpType.RPUSH.name().getBytes(), context.getKey().get(), value}, | ||
context.getKey(), value)); | ||
} | ||
|
||
@Override | ||
public boolean isFinish() { | ||
return STATE.READ_END.equals(state); | ||
} | ||
|
||
@Override | ||
public void reset() { | ||
super.reset(); | ||
if (rdbStringParser != null) { | ||
rdbStringParser.reset(); | ||
} | ||
this.state = STATE.READ_INIT; | ||
} | ||
|
||
@Override | ||
protected Logger getLogger() { | ||
return logger; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.