| 31 | import static io.questdb.std.IOUringAccessor.*; |
| 32 | |
| 33 | public class IOURingImpl implements IOURing { |
| 34 | |
| 35 | // Holds <id, res> tuples for recently consumed cqes. |
| 36 | private final long[] cachedCqes; |
| 37 | private final long cqKheadAddr; |
| 38 | private final int cqKringMask; |
| 39 | private final long cqKtailAddr; |
| 40 | private final long cqesAddr; |
| 41 | private final IOURingFacade facade; |
| 42 | private final long ringAddr; |
| 43 | private final long ringFd; |
| 44 | private final long sqKheadAddr; |
| 45 | private final int sqKringEntries; |
| 46 | private final int sqKringMask; |
| 47 | private final long sqesAddr; |
| 48 | // Index of cached cqe tuple. |
| 49 | private int cachedIndex; |
| 50 | // Count of cached cqe tuples. |
| 51 | private int cachedSize; |
| 52 | private boolean closed = false; |
| 53 | private long idSeq; |
| 54 | |
| 55 | public IOURingImpl(IOURingFacade facade, int capacity) { |
| 56 | assert Numbers.isPow2(capacity); |
| 57 | this.facade = facade; |
| 58 | final long res = facade.create(capacity); |
| 59 | if (res < 0) { |
| 60 | throw CairoException.critical((int) -res).put("Cannot create io_uring instance"); |
| 61 | } |
| 62 | this.ringAddr = res; |
| 63 | |
| 64 | int ringFd = Unsafe.getInt(ringAddr + RING_FD_OFFSET); |
| 65 | |
| 66 | this.sqesAddr = Unsafe.getLong(ringAddr + SQ_SQES_OFFSET); |
| 67 | this.sqKheadAddr = Unsafe.getLong(ringAddr + SQ_KHEAD_OFFSET); |
| 68 | final long sqMaskAddr = Unsafe.getLong(ringAddr + SQ_KRING_MASK_OFFSET); |
| 69 | this.sqKringMask = Unsafe.getInt(sqMaskAddr); |
| 70 | final long sqEntriesAddr = Unsafe.getLong(ringAddr + SQ_KRING_ENTRIES_OFFSET); |
| 71 | this.sqKringEntries = Unsafe.getInt(sqEntriesAddr); |
| 72 | |
| 73 | this.cqesAddr = Unsafe.getLong(ringAddr + CQ_CQES_OFFSET); |
| 74 | this.cqKheadAddr = Unsafe.getLong(ringAddr + CQ_KHEAD_OFFSET); |
| 75 | this.cqKtailAddr = Unsafe.getLong(ringAddr + CQ_KTAIL_OFFSET); |
| 76 | final long cqMaskAddr = Unsafe.getLong(ringAddr + CQ_KRING_MASK_OFFSET); |
| 77 | this.cqKringMask = Unsafe.getInt(cqMaskAddr); |
| 78 | final long cqEntriesAddr = Unsafe.getLong(ringAddr + CQ_KRING_ENTRIES_OFFSET); |
| 79 | int cqKringEntries = Unsafe.getInt(cqEntriesAddr); |
| 80 | cachedCqes = new long[2 * cqKringEntries]; |
| 81 | |
| 82 | this.ringFd = Files.createUniqueFd(ringFd); |
| 83 | } |
| 84 | |
| 85 | @Override |
| 86 | public void close() { |
| 87 | if (closed) { |
| 88 | return; |
| 89 | } |
| 90 | Files.detach(ringFd); |
nothing calls this directly
no outgoing calls
no test coverage detected
searching dependent graphs…