Merge pull request #1035 from Owain94/http-api-rx-calls

http-api: rx calls
This commit is contained in:
Ganom
2019-07-18 23:23:02 -04:00
committed by GitHub
10 changed files with 262 additions and 230 deletions

View File

@@ -26,6 +26,7 @@ package net.runelite.http.api.item;
import com.google.gson.JsonParseException;
import com.google.gson.reflect.TypeToken;
import io.reactivex.Observable;
import java.awt.image.BufferedImage;
import java.io.IOException;
import java.io.InputStream;
@@ -112,7 +113,7 @@ public class ItemClient
}
}
public BufferedImage getIcon(int itemId) throws IOException
public Observable<BufferedImage> getIcon(int itemId)
{
HttpUrl url = RuneLiteAPI.getApiBase().newBuilder()
.addPathSegment("item")
@@ -126,23 +127,26 @@ public class ItemClient
.url(url)
.build();
return Observable.defer(() ->
{
try (Response response = RuneLiteAPI.CLIENT.newCall(request).execute())
{
if (!response.isSuccessful())
{
logger.debug("Error grabbing icon {}: {}", itemId, response);
return null;
return Observable.just(null);
}
InputStream in = response.body().byteStream();
synchronized (ImageIO.class)
{
return ImageIO.read(in);
return Observable.just(ImageIO.read(in));
}
}
});
}
public SearchResult search(String itemName) throws IOException
public Observable<SearchResult> search(String itemName)
{
HttpUrl url = RuneLiteAPI.getApiBase().newBuilder()
.addPathSegment("item")
@@ -152,6 +156,8 @@ public class ItemClient
logger.debug("Built URI: {}", url);
return Observable.defer(() ->
{
Request request = new Request.Builder()
.url(url)
.build();
@@ -161,19 +167,20 @@ public class ItemClient
if (!response.isSuccessful())
{
logger.debug("Error looking up item {}: {}", itemName, response);
return null;
return Observable.just(null);
}
InputStream in = response.body().byteStream();
return RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), SearchResult.class);
return Observable.just(RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), SearchResult.class));
}
catch (JsonParseException ex)
{
throw new IOException(ex);
return Observable.error(ex);
}
});
}
public ItemPrice[] getPrices() throws IOException
public Observable<ItemPrice[]> getPrices()
{
HttpUrl.Builder urlBuilder = RuneLiteAPI.getApiBase().newBuilder()
.addPathSegment("item")
@@ -183,6 +190,9 @@ public class ItemClient
logger.debug("Built URI: {}", url);
return Observable.defer(() ->
{
Request request = new Request.Builder()
.url(url)
.build();
@@ -192,19 +202,20 @@ public class ItemClient
if (!response.isSuccessful())
{
logger.warn("Error looking up prices: {}", response);
return null;
return Observable.just(null);
}
InputStream in = response.body().byteStream();
return RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), ItemPrice[].class);
return Observable.just(RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), ItemPrice[].class));
}
catch (JsonParseException ex)
{
throw new IOException(ex);
return Observable.error(ex);
}
});
}
public Map<Integer, ItemStats> getStats() throws IOException
public Observable<Map<Integer, ItemStats>> getStats()
{
HttpUrl.Builder urlBuilder = RuneLiteAPI.getStaticBase().newBuilder()
.addPathSegment("item")
@@ -215,6 +226,9 @@ public class ItemClient
logger.debug("Built URI: {}", url);
return Observable.defer(() ->
{
Request request = new Request.Builder()
.url(url)
.build();
@@ -224,18 +238,19 @@ public class ItemClient
if (!response.isSuccessful())
{
logger.warn("Error looking up item stats: {}", response);
return null;
return Observable.just(null);
}
InputStream in = response.body().byteStream();
final Type typeToken = new TypeToken<Map<Integer, ItemStats>>()
{
}.getType();
return RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), typeToken);
return Observable.just(RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), typeToken));
}
catch (JsonParseException ex)
{
throw new IOException(ex);
}
return Observable.error(ex);
}
});
}
}

View File

@@ -58,7 +58,7 @@ public class OSBGrandExchangeClient
{
if (!response.isSuccessful())
{
throw new IOException("Error looking up item id: " + response);
return Observable.error(new IOException("Error looking up item id: " + response));
}
final InputStream in = response.body().byteStream();

View File

@@ -26,6 +26,9 @@
package net.runelite.http.api.worlds;
import com.google.gson.JsonParseException;
import io.reactivex.Observable;
import java.io.InputStream;
import java.io.InputStreamReader;
import net.runelite.http.api.RuneLiteAPI;
import okhttp3.HttpUrl;
import okhttp3.Request;
@@ -33,15 +36,11 @@ import okhttp3.Response;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
public class WorldClient
{
private static final Logger logger = LoggerFactory.getLogger(WorldClient.class);
public WorldResult lookupWorlds() throws IOException
public Observable<WorldResult> lookupWorlds()
{
HttpUrl url = RuneLiteAPI.getApiBase().newBuilder()
.addPathSegment("worlds.js")
@@ -49,6 +48,8 @@ public class WorldClient
logger.debug("Built URI: {}", url);
return Observable.defer(() ->
{
Request request = new Request.Builder()
.url(url)
.build();
@@ -58,15 +59,16 @@ public class WorldClient
if (!response.isSuccessful())
{
logger.debug("Error looking up worlds: {}", response);
return null;
return Observable.just(null);
}
InputStream in = response.body().byteStream();
return RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), WorldResult.class);
return Observable.just(RuneLiteAPI.GSON.fromJson(new InputStreamReader(in), WorldResult.class));
}
catch (JsonParseException ex)
{
throw new IOException(ex);
}
return Observable.error(ex);
}
});
}
}

View File

@@ -27,14 +27,16 @@ package net.runelite.client.callback;
import com.google.inject.Inject;
import java.util.Iterator;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.Executor;
import java.util.function.BooleanSupplier;
import javax.inject.Singleton;
import lombok.extern.slf4j.Slf4j;
import net.runelite.api.Client;
import org.jetbrains.annotations.NotNull;
@Singleton
@Slf4j
public class ClientThread
public class ClientThread implements Executor
{
private final ConcurrentLinkedQueue<BooleanSupplier> invokes = new ConcurrentLinkedQueue<>();
@@ -112,4 +114,14 @@ public class ClientThread
}
}
}
@Override
public void execute(@NotNull Runnable r)
{
invoke(() ->
{
r.run();
return true;
});
}
}

View File

@@ -30,9 +30,9 @@ import com.google.common.cache.LoadingCache;
import com.google.common.collect.ImmutableMap;
import com.google.gson.Gson;
import com.google.gson.reflect.TypeToken;
import io.reactivex.schedulers.Schedulers;
import java.awt.Color;
import java.awt.image.BufferedImage;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.lang.reflect.Type;
@@ -142,6 +142,7 @@ public class ItemManager
private final LoadingCache<OutlineKey, BufferedImage> itemOutlines;
private Map<Integer, ItemPrice> itemPrices = Collections.emptyMap();
private Map<Integer, ItemStats> itemStats = Collections.emptyMap();
@Inject
public ItemManager(
Client client,
@@ -209,9 +210,11 @@ public class ItemManager
private void loadPrices()
{
try
itemClient.getPrices()
.subscribeOn(Schedulers.io())
.subscribe(
(prices) ->
{
ItemPrice[] prices = itemClient.getPrices();
if (prices != null)
{
ImmutableMap.Builder<Integer, ItemPrice> map = ImmutableMap.builderWithExpectedSize(prices.length);
@@ -223,29 +226,27 @@ public class ItemManager
}
log.debug("Loaded {} prices", itemPrices.size());
}
catch (IOException e)
{
log.warn("error loading prices!", e);
}
},
(e) -> log.warn("error loading prices!", e)
);
}
private void loadStats()
{
try
itemClient.getStats()
.subscribeOn(Schedulers.io())
.subscribe(
(stats) ->
{
final Map<Integer, ItemStats> stats = itemClient.getStats();
if (stats != null)
{
itemStats = ImmutableMap.copyOf(stats);
}
log.debug("Loaded {} stats", itemStats.size());
}
catch (IOException e)
{
log.warn("error loading stats!", e);
}
},
(e) -> log.warn("error loading stats!", e)
);
}
private void onGameStateChanged(final GameStateChanged event)

View File

@@ -827,10 +827,10 @@ public class ChatCommandsPlugin extends Plugin
{
ItemPrice item = retrieveFromList(results, search);
CLIENT.lookupItem(item.getId())
.subscribeOn(Schedulers.single())
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.from(clientThread))
.subscribe(
(osbresult) ->
clientThread.invoke(() ->
{
int itemId = item.getId();
int itemPrice = itemManager.getItemPrice(itemId);
@@ -865,7 +865,7 @@ public class ChatCommandsPlugin extends Plugin
messageNode.setRuneLiteFormatMessage(response);
chatMessageManager.update(messageNode);
client.refreshChat();
})
}
);
}
}

View File

@@ -25,13 +25,14 @@
package net.runelite.client.plugins.defaultworld;
import com.google.inject.Provides;
import java.io.IOException;
import io.reactivex.schedulers.Schedulers;
import javax.inject.Inject;
import javax.inject.Singleton;
import lombok.extern.slf4j.Slf4j;
import net.runelite.api.Client;
import net.runelite.api.GameState;
import net.runelite.api.events.GameStateChanged;
import net.runelite.client.callback.ClientThread;
import net.runelite.client.config.ConfigManager;
import net.runelite.client.eventbus.EventBus;
import net.runelite.client.events.SessionOpen;
@@ -40,7 +41,6 @@ import net.runelite.client.plugins.PluginDescriptor;
import net.runelite.client.util.WorldUtil;
import net.runelite.http.api.worlds.World;
import net.runelite.http.api.worlds.WorldClient;
import net.runelite.http.api.worlds.WorldResult;
@PluginDescriptor(
name = "Default World",
@@ -60,6 +60,9 @@ public class DefaultWorldPlugin extends Plugin
@Inject
private EventBus eventBus;
@Inject
private ClientThread clientThread;
private final WorldClient worldClient = new WorldClient();
private int worldCache;
private boolean worldChangeRequired;
@@ -122,10 +125,12 @@ public class DefaultWorldPlugin extends Plugin
return;
}
try
worldClient.lookupWorlds()
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.from(clientThread))
.subscribe(
(worldResult) ->
{
final WorldResult worldResult = worldClient.lookupWorlds();
if (worldResult == null)
{
return;
@@ -150,11 +155,9 @@ public class DefaultWorldPlugin extends Plugin
{
log.warn("World {} not found.", correctedWorld);
}
}
catch (IOException e)
{
log.warn("Error looking up world {}. Error: {}", correctedWorld, e);
}
},
(e) -> log.warn("Error looking up world {}. Error: {}", correctedWorld, e)
);
}
private void applyWorld()

View File

@@ -370,10 +370,10 @@ public class ExaminePlugin extends Plugin
{
int finalQuantity = quantity;
CLIENT.lookupItem(id)
.subscribeOn(Schedulers.single())
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.from(clientThread))
.subscribe(
(osbresult) ->
clientThread.invoke(() ->
{
message
.append(ChatColorType.NORMAL)
@@ -400,7 +400,7 @@ public class ExaminePlugin extends Plugin
.append(ChatColorType.NORMAL)
.append("ea)");
}
}),
},
(e) -> log.error(e.toString())
);
}

View File

@@ -551,14 +551,14 @@ public class GrandExchangePlugin extends Plugin
}
CLIENT.lookupItem(itemId)
.subscribeOn(Schedulers.single())
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.from(clientThread))
.subscribe(
(osbresult) ->
clientThread.invoke(() ->
{
final String text = geText.getText() + OSB_GE_TEXT + StackFormatter.formatNumber(osbresult.getOverall_average());
geText.setText(text);
}),
},
(e) -> log.debug("Error getting price of item {}", itemId, e)
);
});

View File

@@ -29,8 +29,8 @@ import com.google.common.base.Stopwatch;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ObjectArrays;
import com.google.inject.Provides;
import io.reactivex.schedulers.Schedulers;
import java.awt.image.BufferedImage;
import java.io.IOException;
import java.time.Duration;
import java.time.Instant;
import java.util.Comparator;
@@ -502,10 +502,11 @@ public class WorldHopperPlugin extends Plugin
{
log.debug("Fetching worlds");
try
new WorldClient().lookupWorlds()
.subscribeOn(Schedulers.io())
.subscribe(
(worldResult) ->
{
WorldResult worldResult = new WorldClient().lookupWorlds();
if (worldResult != null)
{
worldResult.getWorlds().sort(Comparator.comparingInt(World::getId));
@@ -513,11 +514,9 @@ public class WorldHopperPlugin extends Plugin
this.lastFetch = Instant.now();
updateList();
}
}
catch (IOException ex)
{
log.warn("Error looking up worlds", ex);
}
},
(ex) -> log.warn("Error looking up worlds", ex)
);
}
/**