/
dixsu
/
websocket-client
Обзор
Документация
Войти
/
dixsu
/
websocket-client
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
master
src/main/java/WebSocketInputStream.java
85 строк
3 KB
Сахно Роман Александрович
начало работы с POST как с потоками
19 фев 2025, 20:34
19 фев 2025, 20:34
1b600f0
Код
Авторство
О чём код?
import java.io.IOException; import java.io.InputStream; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; public class WebSocketInputStream extends InputStream { private final BlockingQueue<byte[]> dataQueue = new LinkedBlockingQueue<>(); private byte[] currentBuffer; private int currentBufferPosition; private boolean isClosed = false; @Override public int read() throws IOException { if (isClosed) { return -1; // Поток закрыт } // Если текущий буфер пуст или закончился, берем следующий if (currentBuffer == null || currentBufferPosition >= currentBuffer.length) { try { currentBuffer = dataQueue.take(); // Блокируемся, пока не появятся данные currentBufferPosition = 0; } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IOException("Thread interrupted while waiting for data", e); } } // Возвращаем следующий байт из текущего буфера return currentBuffer[currentBufferPosition++] & 0xFF; } @Override public int read(byte[] b, int off, int len) throws IOException { if (isClosed) { return -1; // Поток закрыт } int bytesRead = 0; while (bytesRead < len) { if (currentBuffer == null || currentBufferPosition >= currentBuffer.length) { try { currentBuffer = dataQueue.take(); // Блокируемся, пока не появятся данные currentBufferPosition = 0; } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IOException("Thread interrupted while waiting for data", e); } } int bytesToCopy = Math.min(len - bytesRead, currentBuffer.length - currentBufferPosition); System.arraycopy(currentBuffer, currentBufferPosition, b, off + bytesRead, bytesToCopy); currentBufferPosition += bytesToCopy; bytesRead += bytesToCopy; } return bytesRead; } @Override public void close() { isClosed = true; dataQueue.clear(); // Очищаем очередь при закрытии потока } /** * Добавляет данные в поток. * * @param data Данные для добавления. */ public void addData(byte[] data) { if (!isClosed) { dataQueue.add(data); } } /** * Помечает поток как завершенный. */ public void finish() { if (!isClosed) { dataQueue.add(new byte[0]); // Добавляем пустой массив, чтобы сигнализировать о завершении } } }