diff --git a/src/main/java/io/github/protocol/bookkeeper/BookieInfo.java b/src/main/java/io/github/protocol/bookkeeper/BookieInfo.java new file mode 100644 index 0000000..0877f05 --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/BookieInfo.java @@ -0,0 +1,23 @@ +package io.github.protocol.bookkeeper; + +public class BookieInfo { + private long freeSpace; + + private long totalSpace; + + public long getFreeSpace() { + return freeSpace; + } + + public void setFreeSpace(long freeSpace) { + this.freeSpace = freeSpace; + } + + public long getTotalSpace() { + return totalSpace; + } + + public void setTotalSpace(long totalSpace) { + this.totalSpace = totalSpace; + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/BookieStatus.java b/src/main/java/io/github/protocol/bookkeeper/BookieStatus.java new file mode 100644 index 0000000..e3a4262 --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/BookieStatus.java @@ -0,0 +1,43 @@ +package io.github.protocol.bookkeeper; + +public class BookieStatus { + private boolean running; + + private boolean readOnly; + + private boolean shuttingDown; + + private boolean availableForHighPriorityWrites; + + public boolean isRunning() { + return running; + } + + public void setRunning(boolean running) { + this.running = running; + } + + public boolean isReadOnly() { + return readOnly; + } + + public void setReadOnly(boolean readOnly) { + this.readOnly = readOnly; + } + + public boolean isShuttingDown() { + return shuttingDown; + } + + public void setShuttingDown(boolean shuttingDown) { + this.shuttingDown = shuttingDown; + } + + public boolean isAvailableForHighPriorityWrites() { + return availableForHighPriorityWrites; + } + + public void setAvailableForHighPriorityWrites(boolean availableForHighPriorityWrites) { + this.availableForHighPriorityWrites = availableForHighPriorityWrites; + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/Bookies.java b/src/main/java/io/github/protocol/bookkeeper/Bookies.java new file mode 100644 index 0000000..4a29a59 --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/Bookies.java @@ -0,0 +1,40 @@ +package io.github.protocol.bookkeeper; + +import java.util.List; +import java.util.Map; + +public interface Bookies { + + Map bookieList() throws BookkeeperAdminException; + + Map listBookieInfo() throws BookkeeperAdminException; + + LastLogMark lastLogMark() throws BookkeeperAdminException; + + DiskFile listDiskFile() throws BookkeeperAdminException; + + void expandStorage() throws BookkeeperAdminException; + + void forceGc(boolean forceMajor, boolean forceMinor) throws BookkeeperAdminException; + + boolean isInForceGc() throws BookkeeperAdminException; + + void suspendGc(boolean major, boolean minor) throws BookkeeperAdminException; + + GcSuspendStatus gcSuspendStatus() throws BookkeeperAdminException; + + void resumeGc(boolean major, boolean minor) throws BookkeeperAdminException; + + List gcStatusList() throws BookkeeperAdminException; + + BookieStatus status() throws BookkeeperAdminException; + + void setReadOnly(boolean readOnly) throws BookkeeperAdminException; + + boolean isReadOnly() throws BookkeeperAdminException; + + boolean isReady() throws BookkeeperAdminException; + + BookieInfo bookieInfo() throws BookkeeperAdminException; + +} diff --git a/src/main/java/io/github/protocol/bookkeeper/BookiesImpl.java b/src/main/java/io/github/protocol/bookkeeper/BookiesImpl.java new file mode 100644 index 0000000..a2432b2 --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/BookiesImpl.java @@ -0,0 +1,228 @@ +package io.github.protocol.bookkeeper; + +import com.fasterxml.jackson.core.type.TypeReference; + +import java.io.IOException; +import java.net.http.HttpResponse; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +public class BookiesImpl implements Bookies { + + private final InnerHttpClient innerHttpClient; + + public BookiesImpl(InnerHttpClient innerHttpClient) { + this.innerHttpClient = innerHttpClient; + } + + @Override + public Map bookieList() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_LIST); + return JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public Map listBookieInfo() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_LIST_INFO); + return JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + + } + + @Override + public LastLogMark lastLogMark() throws BookkeeperAdminException { + Map map; + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_LAST_LOG_MARK); + map = JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + LastLogMark lastLogMark = new LastLogMark(); + HashMap lastLogMarks = new HashMap<>(); + for (Map.Entry entry : map.entrySet()) { + lastLogMarks.put(lastLogMark.extractValueFromLastLogMarkKey(entry.getKey()), + lastLogMark.extractValueFromLastLogMarkValue(entry.getValue())); + } + lastLogMark.setLogFileIdTxnMap(lastLogMarks); + return lastLogMark; + } + + @Override + public DiskFile listDiskFile() throws BookkeeperAdminException { + Map map; + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_LIST_DISK_FILE); + map = JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + + DiskFile diskFile = new DiskFile(); + diskFile.setIndexFiles(map.get("index files").split("\t")); + diskFile.setJournalFiles(map.get("journal files").split("\t")); + diskFile.setEntryLogFiles(map.get("entrylog files").split("\t")); + return diskFile; + } + + @Override + public void expandStorage() throws BookkeeperAdminException { + try { + innerHttpClient.put(UrlConst.BOOKIE_EXPAND_STORAGE); + } catch (IOException | InterruptedException e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public void forceGc(boolean forceMajor, boolean forceMinor) throws BookkeeperAdminException { + HashMap requestBody = new HashMap<>(); + requestBody.put("forceMajor", forceMajor); + requestBody.put("forceMinor", forceMinor); + try { + innerHttpClient.put(UrlConst.BOOKIE_GC, JacksonService.toJson(requestBody)); + } catch (IOException | InterruptedException e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public boolean isInForceGc() throws BookkeeperAdminException { + Map map; + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_GC); + map = JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + return Boolean.parseBoolean(map.get("is_in_force_gc")); + } + + @Override + public void suspendGc(boolean major, boolean minor) throws BookkeeperAdminException { + HashMap requestBody = new HashMap<>(); + requestBody.put("suspendMajor", major); + requestBody.put("suspendMinor", minor); + try { + innerHttpClient.put(UrlConst.BOOKIE_GC_SUSPEND_COMPACTION, JacksonService.toJson(requestBody)); + } catch (IOException | InterruptedException e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public GcSuspendStatus gcSuspendStatus() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_GC_SUSPEND_COMPACTION); + Map map = JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + boolean isMajorGcSuspended = map.get("isMajorGcSuspended"); + boolean isMinorGcSuspended = map.get("isMinorGcSuspended"); + return new GcSuspendStatus(isMajorGcSuspended, isMinorGcSuspended); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public void resumeGc(boolean major, boolean minor) throws BookkeeperAdminException { + HashMap requestBody = new HashMap<>(); + requestBody.put("resumeMajor", major); + requestBody.put("resumeMinor", minor); + try { + innerHttpClient.put(UrlConst.BOOKIE_GC_RESUME_COMPACTION, JacksonService.toJson(requestBody)); + } catch (IOException | InterruptedException e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public List gcStatusList() throws BookkeeperAdminException { + List> list; + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_GC_DETAILS); + list = JacksonService.toRefer(resp.body(), new TypeReference>>() { + }); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + return list.stream().map(map -> { + GarbageCollectionStatus garbageCollectionStatus = new GarbageCollectionStatus(); + garbageCollectionStatus.setForceCompacting((boolean) map.get("forceCompacting")); + garbageCollectionStatus.setMajorCompacting((boolean) map.get("majorCompacting")); + garbageCollectionStatus.setMinorCompacting((boolean) map.get("minorCompacting")); + garbageCollectionStatus.setLastMajorCompactionTime((long) map.get("lastMajorCompactionTime")); + garbageCollectionStatus.setLastMinorCompactionTime((long) map.get("lastMinorCompactionTime")); + garbageCollectionStatus.setMajorCompactionCounter((int) map.get("majorCompactionCounter")); + garbageCollectionStatus.setMinorCompactionCounter((int) map.get("minorCompactionCounter")); + return garbageCollectionStatus; + }).collect(Collectors.toList()); + } + + @Override + public BookieStatus status() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_STATE); + return JacksonService.toObject(resp.body(), BookieStatus.class); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public void setReadOnly(boolean readOnly) throws BookkeeperAdminException { + HashMap requestBody = new HashMap<>(); + requestBody.put("readOnly", readOnly); + try { + innerHttpClient.put(UrlConst.BOOKIE_STATE_READ_ONLY, JacksonService.toJson(requestBody)); + } catch (IOException | InterruptedException e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public boolean isReadOnly() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_STATE_READ_ONLY); + Map map = JacksonService.toRefer(resp.body(), new TypeReference>() { + }); + return map.get("readOnly"); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public boolean isReady() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_READY); + return "OK".equals(resp.body()); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } + + @Override + public BookieInfo bookieInfo() throws BookkeeperAdminException { + try { + HttpResponse resp = innerHttpClient.get(UrlConst.BOOKIE_INFO); + return JacksonService.toObject(resp.body(), BookieInfo.class); + } catch (Exception e) { + throw new BookkeeperAdminException(e); + } + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/DiskFile.java b/src/main/java/io/github/protocol/bookkeeper/DiskFile.java new file mode 100644 index 0000000..2e9123d --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/DiskFile.java @@ -0,0 +1,31 @@ +package io.github.protocol.bookkeeper; + +public class DiskFile { + private String[] indexFiles; + private String[] journalFiles; + private String[] entryLogFiles; + + public String[] getIndexFiles() { + return indexFiles; + } + + public void setIndexFiles(String[] indexFiles) { + this.indexFiles = indexFiles; + } + + public String[] getJournalFiles() { + return journalFiles; + } + + public void setJournalFiles(String[] journalFiles) { + this.journalFiles = journalFiles; + } + + public String[] getEntryLogFiles() { + return entryLogFiles; + } + + public void setEntryLogFiles(String[] entryLogFiles) { + this.entryLogFiles = entryLogFiles; + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/GarbageCollectionStatus.java b/src/main/java/io/github/protocol/bookkeeper/GarbageCollectionStatus.java new file mode 100644 index 0000000..081c01c --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/GarbageCollectionStatus.java @@ -0,0 +1,74 @@ +package io.github.protocol.bookkeeper; + +public class GarbageCollectionStatus { + + private boolean forceCompacting; + + private boolean majorCompacting; + + private boolean minorCompacting; + + private long lastMajorCompactionTime; + + private long lastMinorCompactionTime; + + private int majorCompactionCounter; + + private int minorCompactionCounter; + + public boolean isForceCompacting() { + return forceCompacting; + } + + public void setForceCompacting(boolean forceCompacting) { + this.forceCompacting = forceCompacting; + } + + public boolean isMajorCompacting() { + return majorCompacting; + } + + public void setMajorCompacting(boolean majorCompacting) { + this.majorCompacting = majorCompacting; + } + + public boolean isMinorCompacting() { + return minorCompacting; + } + + public void setMinorCompacting(boolean minorCompacting) { + this.minorCompacting = minorCompacting; + } + + public long getLastMajorCompactionTime() { + return lastMajorCompactionTime; + } + + public void setLastMajorCompactionTime(long lastMajorCompactionTime) { + this.lastMajorCompactionTime = lastMajorCompactionTime; + } + + public long getLastMinorCompactionTime() { + return lastMinorCompactionTime; + } + + public void setLastMinorCompactionTime(long lastMinorCompactionTime) { + this.lastMinorCompactionTime = lastMinorCompactionTime; + } + + public int getMajorCompactionCounter() { + return majorCompactionCounter; + } + + public void setMajorCompactionCounter(int majorCompactionCounter) { + this.majorCompactionCounter = majorCompactionCounter; + } + + public int getMinorCompactionCounter() { + return minorCompactionCounter; + } + + public void setMinorCompactionCounter(int minorCompactionCounter) { + this.minorCompactionCounter = minorCompactionCounter; + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/GcSuspendStatus.java b/src/main/java/io/github/protocol/bookkeeper/GcSuspendStatus.java new file mode 100644 index 0000000..54916db --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/GcSuspendStatus.java @@ -0,0 +1,30 @@ +package io.github.protocol.bookkeeper; + +public class GcSuspendStatus { + private boolean isMajorGcSuspended; + + private boolean isMinorGcSuspended; + + public boolean isMajorGcSuspended() { + return isMajorGcSuspended; + } + + public void setMajorGcSuspended(boolean majorGcSuspended) { + isMajorGcSuspended = majorGcSuspended; + } + + public boolean isMinorGcSuspended() { + return isMinorGcSuspended; + } + + public void setMinorGcSuspended(boolean minorGcSuspended) { + isMinorGcSuspended = minorGcSuspended; + } + + public GcSuspendStatus(){} + + public GcSuspendStatus(boolean isMajorGcSuspended, boolean isMinorGcSuspended) { + this.isMajorGcSuspended = isMajorGcSuspended; + this.isMinorGcSuspended = isMinorGcSuspended; + } +} diff --git a/src/main/java/io/github/protocol/bookkeeper/LastLogMark.java b/src/main/java/io/github/protocol/bookkeeper/LastLogMark.java new file mode 100644 index 0000000..b23eaaf --- /dev/null +++ b/src/main/java/io/github/protocol/bookkeeper/LastLogMark.java @@ -0,0 +1,30 @@ +package io.github.protocol.bookkeeper; + +import java.util.Arrays; +import java.util.Map; + +public class LastLogMark { + + private Map logFileIdTxnMap; + + public Map getLogFileIdTxnMap() { + return logFileIdTxnMap; + } + + public void setLogFileIdTxnMap(Map logFileIdTxnMap) { + this.logFileIdTxnMap = logFileIdTxnMap; + } + public int extractValueFromLastLogMarkKey(String data) { + String result = Arrays.stream(data.replaceAll("LastLogMark: Journal Id -", "") + .replaceAll(" ", "") + .split("\\(")).findFirst().get(); + return Integer.parseInt(result); + } + + public int extractValueFromLastLogMarkValue(String data) { + String result = data.replace("Pos - ", "") + .replaceAll(" ", ""); + return Integer.parseInt(result); + } + +} diff --git a/src/test/java/io/github/protocol/bookkeeper/BookiesImplTest.java b/src/test/java/io/github/protocol/bookkeeper/BookiesImplTest.java new file mode 100644 index 0000000..7a15b98 --- /dev/null +++ b/src/test/java/io/github/protocol/bookkeeper/BookiesImplTest.java @@ -0,0 +1,102 @@ +package io.github.protocol.bookkeeper; + +import io.github.embedded.bookkeeper.core.EmbeddedBookkeeperServer; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; + +public class BookiesImplTest { + + private static final EmbeddedBookkeeperServer SERVER = new EmbeddedBookkeeperServer(); + + private static Bookies bookiesImpl = null; + + @BeforeAll + public static void setup() throws Exception { + SERVER.start(); + Configuration conf = new Configuration(); + conf.setHost("localhost"); + conf.setPort(SERVER.getBkWebPort()); + bookiesImpl = new BookiesImpl(new InnerHttpClient(conf)); + } + + @Test + public void testBookiesList() throws BookkeeperAdminException { + Map map = bookiesImpl.bookieList(); + Assertions.assertEquals(1, map.size()); + String key = ""; + for (Map.Entry entry : map.entrySet()) { + key = entry.getKey(); + } + Assertions.assertTrue(key.contains(String.valueOf(SERVER.getBkPort()))); + } + + @Test + public void testListBookieInfo() throws BookkeeperAdminException { + Map info = bookiesImpl.listBookieInfo(); + Assertions.assertEquals(2, info.size()); + } + + @Test + public void testLastLogMark() throws BookkeeperAdminException { + LastLogMark lastLogMark = bookiesImpl.lastLogMark(); + Assertions.assertEquals(1, lastLogMark.getLogFileIdTxnMap().size()); + } + + @Test + public void testListDiskFile() throws BookkeeperAdminException { + bookiesImpl.listDiskFile(); + } + + @Test + public void testIsInForceGc() throws BookkeeperAdminException { + boolean inForceGc = bookiesImpl.isInForceGc(); + Assertions.assertFalse(inForceGc); + } + + @Test + public void testGcStatus() throws BookkeeperAdminException { + List statuses = bookiesImpl.gcStatusList(); + Assertions.assertEquals(1, statuses.size()); + GarbageCollectionStatus garbageCollectionStatus = statuses.get(0); + Assertions.assertFalse(garbageCollectionStatus.isMajorCompacting()); + + } + + @Test + public void testBookiesStatus() throws BookkeeperAdminException { + BookieStatus status = bookiesImpl.status(); + Assertions.assertTrue(status.isRunning()); + Assertions.assertFalse(status.isReadOnly()); + + } + + @Test + public void testIsReadOnly() throws BookkeeperAdminException { + boolean readOnly = bookiesImpl.isReadOnly(); + Assertions.assertFalse(readOnly); + } + + @Test + public void testBookiesIsReady() throws BookkeeperAdminException { + boolean ready = bookiesImpl.isReady(); + Assertions.assertTrue(ready); + } + + @Test + public void testBookiesInfo() throws BookkeeperAdminException { + BookieInfo bookieInfo = bookiesImpl.bookieInfo(); + Assertions.assertTrue(bookieInfo.getFreeSpace() > 0); + Assertions.assertTrue(bookieInfo.getTotalSpace() > 0); + } + + @AfterAll + public static void teardown() throws Exception { + SERVER.close(); + } + +}