Revert "+ upgrade to Rjc 0.6.4 (snapshot for now)"

This reverts commit db0522429a.

+ Downgrading RJC dependency to 0.6.3 since 0.6.4 is not yet released
This commit is contained in:
Costin Leau
2011-04-05 21:51:13 +03:00
parent f98805a6ce
commit 6cf1a58606
5 changed files with 45 additions and 141 deletions

View File

@@ -314,8 +314,8 @@
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>1.5</source>
<target>1.5</target>
<source>1.6</source>
<target>1.6</target>
<compilerArgument>-Xlint:all</compilerArgument>
<showWarnings>true</showWarnings>
<showDeprecation>false</showDeprecation>

View File

@@ -16,10 +16,10 @@
<spring.range>"[3.0.0, 4.0.0)"</spring.range>
<jredis.ver>03122010</jredis.ver>
<jedis.ver>1.5.2</jedis.ver>
<rjc.ver>0.6.4-SNAPSHOT</rjc.ver>
<rjc.ver>0.6.3</rjc.ver>
<jedis.range>"[1.0.0,2.0.0)"</jedis.range>
<jackson.range>"[1.6, 2.0.0)"</jackson.range>
<rjc.range>"[0.6.4, 0.6.4]"</rjc.range>
<rjc.range>"[0.6.3, 0.6.3]"</rjc.range>
</properties>
<dependencies>

View File

@@ -1,125 +0,0 @@
/*
* Copyright 2011 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.keyvalue.redis.connection.rjc;
import java.io.IOException;
import java.net.UnknownHostException;
import java.util.List;
import org.idevlab.rjc.ds.RedisConnection;
import org.idevlab.rjc.message.RedisNodeSubscriber;
import org.idevlab.rjc.protocol.Protocol.Command;
/**
* Basic decorator suppressing close() calls to the underlying connection.
* Used for reusing arbitrary connections with {@link RedisNodeSubscriber} without
* resorting to connection pooling.
*
* @author Costin Leau
*/
class CloseSuppressingRjcConnection implements RedisConnection {
private final RedisConnection delegate;
/**
* Constructs a new <code>CloseSuppressingRjcConnection</code> instance.
*
* @param delegate
*/
CloseSuppressingRjcConnection(RedisConnection delegate) {
this.delegate = delegate;
}
public void close() {
// no-op
}
public void connect() throws UnknownHostException, IOException {
delegate.connect();
}
public List<Object> getAll() {
return delegate.getAll();
}
public byte[] getBinaryBulkReply() {
return delegate.getBinaryBulkReply();
}
public List<byte[]> getBinaryMultiBulkReply() {
return delegate.getBinaryMultiBulkReply();
}
public String getBulkReply() {
return delegate.getBulkReply();
}
public String getHost() {
return delegate.getHost();
}
public Long getIntegerReply() {
return delegate.getIntegerReply();
}
public List<String> getMultiBulkReply() {
return delegate.getMultiBulkReply();
}
public List<Object> getObjectMultiBulkReply() {
return delegate.getObjectMultiBulkReply();
}
public Object getOne() {
return delegate.getOne();
}
public int getPort() {
return delegate.getPort();
}
public String getStatusCodeReply() {
return delegate.getStatusCodeReply();
}
public int getTimeout() {
return delegate.getTimeout();
}
public boolean isConnected() {
return delegate.isConnected();
}
public void rollbackTimeout() {
delegate.rollbackTimeout();
}
public void sendCommand(Command arg0, byte[]... arg1) {
delegate.sendCommand(arg0, arg1);
}
public void sendCommand(Command arg0, String... arg1) {
delegate.sendCommand(arg0, arg1);
}
public void sendCommand(Command arg0) {
delegate.sendCommand(arg0);
}
public void setTimeoutInfinite() {
delegate.setTimeoutInfinite();
}
}

View File

@@ -58,7 +58,7 @@ public class RjcConnection implements RedisConnection {
SingleDataSource connectionDataSource = new SingleDataSource(connection);
session = new SessionFactoryImpl(connectionDataSource).create();
subscriber = new RedisNodeSubscriber();
subscriber.setDataSource(new SingleDataSource(new CloseSuppressingRjcConnection(connection)));
subscriber.setDataSource(connectionDataSource);
client = new Client(connection);
this.dbIndex = dbIndex;
@@ -79,9 +79,13 @@ public class RjcConnection implements RedisConnection {
@Override
public void close() throws DataAccessException {
isClosed = true;
try {
subscriber.close();
} catch (Exception ex) {
// ignore
}
try {
session.close();
} catch (Exception ex) {
throw convertRjcAccessException(ex);
@@ -2024,9 +2028,8 @@ public class RjcConnection implements RedisConnection {
throw new UnsupportedOperationException();
}
subscription = new RjcSubscription(listener, subscriber);
subscription = new RjcSubscription(listener, subscriber, client);
subscription.pSubscribe(patterns);
subscriber.runSubscription();
} catch (Exception ex) {
throw convertRjcAccessException(ex);
@@ -2048,9 +2051,8 @@ public class RjcConnection implements RedisConnection {
throw new UnsupportedOperationException();
}
subscription = new RjcSubscription(listener, subscriber);
subscription = new RjcSubscription(listener, subscriber, client);
subscription.subscribe(channels);
subscriber.runSubscription();
} catch (Exception ex) {
throw convertRjcAccessException(ex);

View File

@@ -15,6 +15,7 @@
*/
package org.springframework.data.keyvalue.redis.connection.rjc;
import org.idevlab.rjc.Client;
import org.idevlab.rjc.message.RedisNodeSubscriber;
import org.springframework.data.keyvalue.redis.connection.MessageListener;
import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscription;
@@ -27,36 +28,62 @@ import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscript
class RjcSubscription extends AbstractSubscription {
private final RedisNodeSubscriber subscriber;
private final Client client;
// rjc does not support subscription while listening
// so we have to handle this ourselves through the client
private volatile boolean subscribed = false;
RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber) {
RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber, Client client) {
super(listener);
this.subscriber = subscriber;
subscriber.setMessageListener(new RjcMessageListener(listener));
subscriber.setPMessageListener(new RjcMessageListener(listener));
this.client = client;
}
@Override
protected void doClose() {
subscriber.close();
subscribed = false;
client.unsubscribe();
client.punsubscribe();
client.rollbackTimeout();
}
@Override
protected void doPsubscribe(byte[]... patterns) {
subscriber.psubscribe(RjcUtils.decodeMultiple(patterns));
String[] pats = RjcUtils.decodeMultiple(patterns);
if (subscribed) {
client.psubscribe(pats);
}
else {
subscriber.setPatterns(RjcUtils.addArray(subscriber.getPatterns(), pats));
subscribed = true;
subscriber.subscribe();
}
}
@Override
protected void doPUnsubscribe(boolean all, byte[]... patterns) {
subscriber.punsubscribe(RjcUtils.decodeMultiple(patterns));
client.punsubscribe(RjcUtils.decodeMultiple(patterns));
}
@Override
protected void doSubscribe(byte[]... channels) {
subscriber.subscribe(RjcUtils.decodeMultiple(channels));
String[] chs = RjcUtils.decodeMultiple(channels);
if (subscribed) {
client.subscribe(chs);
}
else {
subscriber.setPatterns(RjcUtils.addArray(subscriber.getChannels(), chs));
subscribed = true;
subscriber.subscribe();
}
}
@Override
protected void doUnsubscribe(boolean all, byte[]... channels) {
subscriber.punsubscribe(RjcUtils.decodeMultiple(channels));
client.punsubscribe(RjcUtils.decodeMultiple(channels));
}
}