mirror of
https://github.com/MihailRis/voxelcore.git
synced 2026-10-04 18:41:51 +00:00
Merge pull request #853 from Onran0/io-streams-enhance
I/O Stream's enhance
This commit is contained in:
commit
b66b3881f9
13 changed files with 492 additions and 169 deletions
|
|
@ -92,6 +92,11 @@ socket:recv_async(
|
|||
[опционально] usetable: boolean=false
|
||||
) -> nil|table|Bytearray
|
||||
|
||||
-- Оборачивает сокет в io_stream (см. ../io_stream.md)
|
||||
socket:as_stream(
|
||||
[опционально] binary_mode: boolean=true
|
||||
) -> io_stream
|
||||
|
||||
-- Закрывает соединение
|
||||
socket:close()
|
||||
|
||||
|
|
|
|||
|
|
@ -18,7 +18,7 @@
|
|||
Поток имеет три различных вида режима:
|
||||
- Режим общего поведения (`general` / `mode`)
|
||||
|
||||
- Режим сброса (`flush` / `flushMode`)
|
||||
- Режим сброса (`flush` / `flush_mode`)
|
||||
|
||||
- Двоичный режим (`binary`)
|
||||
|
||||
|
|
@ -27,34 +27,34 @@
|
|||
Определяет, как поток обрабатывает чтение и запись,
|
||||
имеет три подрежима:
|
||||
|
||||
| Режим | Описание |
|
||||
|---------------|--------|
|
||||
| `"default"` | Прямой режим. `read` может вернуть меньше байт, чем запрошено. `write` сразу отправляет данные в низкоуровневый дескриптор. Нет буферизации. |
|
||||
| Режим | Описание |
|
||||
|---------------|-------------------------------------------------------------------------------------------------------------------------------------------------|
|
||||
| `"default"` | Прямой режим. `read` может вернуть меньше байт, чем запрошено. `write` сразу отправляет данные в низкоуровневый дескриптор. Нет буферизации. |
|
||||
| `"yield"` | Как `default`, но при нехватке данных в `read(n)` поток будет вызывать `coroutine.yield()`, пока не соберёт ровно `n` байт. Удобно для корутин. |
|
||||
| `"buffered"` | Включает внутренние буферы чтения и записи. `read` берёт данные из буфера, `write` — складывает в буфер. При превышении `maxBufferSize` — ошибка `buffer overflow`. |
|
||||
| `"buffered"` | Включает внутренние буферы чтения и записи. `read` берёт данные из буфера , `write` — складывает в буфер. |
|
||||
|
||||
### flush
|
||||
|
||||
Работает только в режиме `"buffered"`,
|
||||
имеет два подрежима:
|
||||
|
||||
| Режим | Что делает `flush()` |
|
||||
|------------------|----------------------|
|
||||
| `"all"` (по умолчанию) | Сначала сбрасывает буфер записи → низкоуровневый `write`, затем вызывает `ioLib.flush(descriptor)` |
|
||||
| `"buffer"` | Сбрасывает только буфер записи, без вызова системного `flush` |
|
||||
| Режим | Что делает `flush()` |
|
||||
|------------------------|-----------------------------------------------------------------------------------------------------|
|
||||
| `"all"` (по умолчанию) | Сначала сбрасывает буфер записи → низкоуровневый `write`, затем вызывает `io_lib.flush(descriptor)` |
|
||||
| `"buffer"` | Сбрасывает только буфер записи, без вызова системного `flush` |
|
||||
|
||||
|
||||
### binary
|
||||
**Независимый флаг** (включается через `set_binary_mode(true)` или при создании потока).
|
||||
Определяет, в каком виде методы `read` и `write` принимают и возвращают данные:
|
||||
|
||||
| binary = true | binary = false (по умолчанию) |
|
||||
|-----------------------------------------|----------------------------------------|
|
||||
| binary = true | binary = false (по умолчанию) |
|
||||
|------------------------------------------------------------|----------------------------------------|
|
||||
| Данные — это **байты** (`Bytearray`, таблица чисел 0..255) | Данные — это **текстовые строки** |
|
||||
| `read(n)` → `Bytearray` или `table<number>`| `read()` → одна строка |
|
||||
| `read("i4 f")` → распаковка через `byteutil.unpack` | `read(n)` → n строк в таблице |
|
||||
| `write(Bytearray)` → запись байтов | `write("hello")` → строка + `\n` |
|
||||
| `write("i4", 42)` → `byteutil.pack` | `write({"a","b"})` → две строки с `\n`|
|
||||
| `read(n)` → `Bytearray` или `table<int>` | `read()` → одна строка |
|
||||
| `read("i4 f")` → распаковка через `byteutil.unpack` | `read(n)` → n строк в таблице |
|
||||
| `write(Bytearray)` → запись байтов | `write("hello")` → строка + `\n` |
|
||||
| `write("i4", 42)` → `byteutil.pack` | `write({"a","b"})` → две строки с `\n` |
|
||||
|
||||
`read_line` / `write_line` — работают как в текстовом режиме
|
||||
|
||||
|
|
@ -90,25 +90,25 @@ io_stream:set_flush_mode(string)
|
|||
--[[
|
||||
Читает данные из потока
|
||||
В двоичном режиме:
|
||||
Если arg - number, то читает из потока arg байт и возвращает ввиде Bytearray или таблицы, если useTable = true
|
||||
Если arg - int, то читает из потока arg байт и возвращает ввиде Bytearray или таблицы, если use_table = true
|
||||
|
||||
Если arg - string, то функция интерпретирует arg как шаблон для byteutil. Прочитает кол-во байт, которое определено шаблоном, передаст их в byteutil.unpack и вернёт результат
|
||||
В текстовом режиме:
|
||||
Если arg - number, то читает нужное кол-во строк с окончанием CRLF/LF из arg и возвращает ввиде таблицы. Также, если trimEmptyLines = true, то удаляет пустые строки с начала и конца из итоговой таблицы
|
||||
Если arg - int, то читает нужное кол-во строк с окончанием CRLF/LF из arg и возвращает ввиде таблицы. Также, если trim_empty_lines = true, то удаляет пустые строки с начала и конца из итоговой таблицы
|
||||
|
||||
Если arg не определён, то читает одну строку с окончанием CRLF/LF и возвращает её.
|
||||
--]]
|
||||
io_stream:read(
|
||||
[опционально] arg: number | string,
|
||||
[опционально] useTable | trimEmptyLines: boolean
|
||||
) -> Bytearray | table<number> | string | table<string> | ...
|
||||
[опционально] arg: int | string,
|
||||
[опционально] use_table = false | trim_empty_lines = true: boolean
|
||||
) -> Bytearray | table<int> | string | table<string> | ...
|
||||
|
||||
--[[
|
||||
Записывает данные в поток
|
||||
В двоичном режиме:
|
||||
Если arg - string, то функция интерпретирует arg как шаблон для byteutil, передаст его и ... в byteutil.pack и результат запишет в поток
|
||||
|
||||
Если arg - Bytearray | table<number>, то записывает байты в поток
|
||||
Если arg - Bytearray | table<int>, то записывает байты в поток
|
||||
|
||||
В текстовом режиме:
|
||||
Если arg - string, то записывает строку в поток (вместе с окончанием LF)
|
||||
|
|
@ -116,7 +116,7 @@ io_stream:read(
|
|||
Если arg - table<string>, то записывает каждую строку из таблицы отдельно
|
||||
--]]
|
||||
io_stream:write(
|
||||
arg: Bytearray | table<number> | string | table<string>,
|
||||
arg: Bytearray | table<int> | string | table<string>,
|
||||
[опционально] ...
|
||||
)
|
||||
|
||||
|
|
@ -128,53 +128,59 @@ io_stream:write_line(string)
|
|||
|
||||
--[[
|
||||
В двоичном режиме:
|
||||
Читает все доступные байты из потока и возвращает ввиде Bytearray или table<number>, если useTable = true
|
||||
Читает все доступные байты из потока и возвращает ввиде Bytearray или table<int>, если use_table = true
|
||||
|
||||
В текстовом режиме:
|
||||
Читает все доступные строки из потока в table<string> если useTable = true, или в одну строку вместе с окончаниями, если нет
|
||||
Читает все доступные строки из потока в table<string> если use_table = true, или в одну строку вместе с окончаниями, если нет
|
||||
|
||||
--]]
|
||||
io_stream:read_fully(
|
||||
[опционально] useTable: boolean
|
||||
) -> Bytearray | table<number> | table<string> | string
|
||||
[опционально] use_table: boolean
|
||||
) -> Bytearray | table<int> | table<string> | string
|
||||
|
||||
--[[
|
||||
Устанавливает позицию в потоке
|
||||
Если length определён, то возвращает true, если length байт доступно к чтению. Иначе возвращает false.
|
||||
|
||||
Если не определён, то возвращает количество байт, которое можно прочитать.
|
||||
|
||||
В не буферизированном режиме потока может всегда возвращать 0 или false, если поток не поддерживает available.
|
||||
--]]
|
||||
io_stream:available(
|
||||
[опционально] length: int
|
||||
) -> int | boolean
|
||||
|
||||
--[[
|
||||
Устанавливает позицию в потоке (всегда в байтах)
|
||||
Режимы:
|
||||
b - Задаёт позицию относительно начало файла
|
||||
c - Задаёт позицию относительно текущей позиции
|
||||
e - Задаёт позицию относительно конца файла
|
||||
|
||||
Может бросать ошибку, если поток не поддерживает seek.
|
||||
--]]
|
||||
io_stream:seek(
|
||||
mode: string
|
||||
offset: number
|
||||
offset: int
|
||||
)
|
||||
|
||||
-- Возвращает текущую позицию в потоке от начала.
|
||||
-- Может бросать ошибку, если поток не поддерживает tell.
|
||||
io_stream:tell() -> int
|
||||
```
|
||||
|
||||
## Методы Buffered-режима
|
||||
|
||||
```lua
|
||||
--[[
|
||||
Если length определён, то возвращает true, если length байт доступно к чтению. Иначе возвращает false
|
||||
|
||||
Если не определён, то возвращает количество байт, которое можно прочитать
|
||||
|
||||
--]]
|
||||
io_stream:available(
|
||||
[опционально] length: number
|
||||
) -> number | boolean
|
||||
|
||||
-- Возвращает максимальный размер буферов
|
||||
io_stream:get_max_buffer_size() -> number
|
||||
io_stream:get_max_buffer_size() -> int
|
||||
|
||||
-- Задаёт новый максимальный размер буферов
|
||||
io_stream:set_max_buffer_size(max_size: number)
|
||||
io_stream:set_max_buffer_size(max_size: int)
|
||||
```
|
||||
|
||||
## Методы контроля состояния потока
|
||||
|
||||
```lua
|
||||
|
||||
-- Возвращает true, если поток открыт на данный момент
|
||||
io_stream:is_alive() -> bool
|
||||
|
||||
|
|
@ -185,16 +191,27 @@ io_stream:is_closed() -> bool
|
|||
io_stream:close()
|
||||
|
||||
-- Записывает все данные из write-буфера в поток в buffer/all flush-режимах
|
||||
-- Вызывает ioLib.flush() в all flush-режиме
|
||||
-- Вызывает ioLib.flush() в all flush-режиме, или ничего не делает, если
|
||||
-- ioLib не поддерживает flush.
|
||||
io_stream:flush()
|
||||
|
||||
-- Создаёт новый поток из Bytearray.
|
||||
-- Может использоваться одновременно как для чтения, так и для записи.
|
||||
-- Результат записи будет записан в тот же Bytearray, что был передан в функцию.
|
||||
io_stream.wrap_bytearray(
|
||||
buffer: Bytearray,
|
||||
|
||||
-- по-умолчанию равен true, поскольку функция из вводных аргументов
|
||||
-- будет использоваться преимущественно для работы с двоичными данными.
|
||||
[опционально] binary_mode: boolean = true
|
||||
) -> io_stream
|
||||
|
||||
-- Создаёт новый поток с переданным дескриптором и использующим переданную I/O библиотеку. (Более подробно в core:io_stream.lua)
|
||||
io_stream.new(
|
||||
descriptor: int,
|
||||
binaryMode: bool,
|
||||
ioLib: table,
|
||||
binary_mode: boolean,
|
||||
io_lib: table,
|
||||
[опционально] mode: string = "default",
|
||||
[опционально] flushMode: string = "all"
|
||||
[опционально] flush_mode: string = "all"
|
||||
) -> io_stream
|
||||
```
|
||||
131
res/modules/internal/stream_providers/bytearray.lua
Normal file
131
res/modules/internal/stream_providers/bytearray.lua
Normal file
|
|
@ -0,0 +1,131 @@
|
|||
local io_stream = require "core:io_stream"
|
||||
|
||||
local lib = { }
|
||||
|
||||
local buffers = { }
|
||||
local positions = { }
|
||||
|
||||
local next_descriptor = 0
|
||||
|
||||
local function open_descriptor(buffer)
|
||||
next_descriptor = next_descriptor + 1
|
||||
|
||||
buffers[next_descriptor] = buffer
|
||||
positions[next_descriptor] = 1
|
||||
|
||||
return next_descriptor
|
||||
end
|
||||
|
||||
local function require_descriptor(descriptor)
|
||||
if not buffers[descriptor] then
|
||||
error("unknown descriptor")
|
||||
end
|
||||
end
|
||||
|
||||
function lib.read(descriptor, length)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
local buf = buffers[descriptor]
|
||||
local buf_length = #buf
|
||||
local pos = positions[descriptor]
|
||||
|
||||
local to_read = math.min(buf_length - pos + 1, length)
|
||||
|
||||
if to_read <= 0 then
|
||||
return Bytearray()
|
||||
end
|
||||
|
||||
local segment = buf:slice(pos, to_read)
|
||||
|
||||
positions[descriptor] = pos + to_read
|
||||
|
||||
return segment
|
||||
end
|
||||
|
||||
function lib.write(descriptor, data)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
local buf = buffers[descriptor]
|
||||
local pos = positions[descriptor]
|
||||
|
||||
local buf_length = #buf
|
||||
local data_length = #data
|
||||
|
||||
local end_pos = pos + data_length - 1
|
||||
|
||||
-- size ensuring
|
||||
if end_pos > buf_length then
|
||||
for i = buf_length + 1, end_pos do
|
||||
buf[i] = 0
|
||||
end
|
||||
end
|
||||
|
||||
for i = 1, data_length do
|
||||
buf[i + pos - 1] = data[i]
|
||||
end
|
||||
|
||||
positions[descriptor] = pos + data_length
|
||||
end
|
||||
|
||||
function lib.seek(descriptor, mode, offset)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
local buf = buffers[descriptor]
|
||||
local buf_length = #buf
|
||||
|
||||
local base
|
||||
|
||||
if mode == 'b' then
|
||||
base = 1
|
||||
elseif mode == 'c' then
|
||||
base = positions[descriptor]
|
||||
elseif mode == 'e' then
|
||||
base = buf_length + 1
|
||||
else error('invalid seek mode') end
|
||||
|
||||
local new_pos = base + offset
|
||||
|
||||
if new_pos < 1 then
|
||||
error('failed to seek stream')
|
||||
end
|
||||
|
||||
positions[descriptor] = new_pos
|
||||
end
|
||||
|
||||
function lib.tell(descriptor)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
return positions[descriptor]
|
||||
end
|
||||
|
||||
function lib.available(descriptor)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
local buf = buffers[descriptor]
|
||||
local pos = positions[descriptor]
|
||||
|
||||
return math.max(#buf - pos + 1, 0)
|
||||
end
|
||||
|
||||
function lib.is_alive(descriptor)
|
||||
return buffers[descriptor] ~= nil
|
||||
end
|
||||
|
||||
function lib.close(descriptor)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
buffers[descriptor] = nil
|
||||
positions[descriptor] = nil
|
||||
end
|
||||
|
||||
return function(buffer, binary_mode)
|
||||
if binary_mode == nil then
|
||||
binary_mode = true
|
||||
end
|
||||
|
||||
return io_stream.new(
|
||||
open_descriptor(buffer),
|
||||
binary_mode,
|
||||
lib
|
||||
)
|
||||
end
|
||||
|
|
@ -6,6 +6,7 @@ local lib = {
|
|||
seek = file.__seek_descriptor,
|
||||
tell = file.__tell_descriptor,
|
||||
flush = file.__flush_descriptor,
|
||||
available = file.__available_descriptor,
|
||||
is_alive = file.__has_descriptor,
|
||||
close = file.__close_descriptor
|
||||
}
|
||||
|
|
@ -17,6 +18,7 @@ file.__write_descriptor = nil
|
|||
file.__seek_descriptor = nil
|
||||
file.__tell_descriptor = nil
|
||||
file.__flush_descriptor = nil
|
||||
file.__available_descriptor = nil
|
||||
file.__has_descriptor = nil
|
||||
file.__close_descriptor = nil
|
||||
|
||||
|
|
|
|||
|
|
@ -9,6 +9,7 @@ int close(int fd);
|
|||
ssize_t read(int fd, void *buf, size_t count);
|
||||
ssize_t write(int fd, const void *buf, size_t count);
|
||||
int fcntl(int fd, int cmd, ...);
|
||||
int ioctl(int fd, unsigned long request, ...);
|
||||
|
||||
const char *strerror(int errnum);
|
||||
]]
|
||||
|
|
@ -20,6 +21,7 @@ local O_WRONLY = 0x1
|
|||
local O_RDWR = 0x2
|
||||
local O_NONBLOCK = 0x800
|
||||
local F_GETFL = 3
|
||||
local FIONREAD = 0x541B
|
||||
|
||||
local function getError()
|
||||
local err = FFI.errno()
|
||||
|
|
@ -58,12 +60,16 @@ function lib.write(fd, bytearray)
|
|||
end
|
||||
end
|
||||
|
||||
function lib.seek(fd, mode, offset)
|
||||
error("cannot seek the named pipe")
|
||||
end
|
||||
function lib.available(fd)
|
||||
if fd == nil or fd < 0 then return 0 end
|
||||
|
||||
function lib.flush(fd)
|
||||
-- no flush on unix
|
||||
local bytes_ready = FFI.new("int[1]", 0)
|
||||
|
||||
if C.ioctl(fd, FIONREAD, bytes_ready) == -1 then
|
||||
return 0
|
||||
end
|
||||
|
||||
return tonumber(bytes_ready[0])
|
||||
end
|
||||
|
||||
function lib.is_alive(fd)
|
||||
|
|
|
|||
|
|
@ -97,14 +97,26 @@ function lib.write(handle, bytearray)
|
|||
end
|
||||
end
|
||||
|
||||
function lib.seek(handle, mode, offset)
|
||||
error("cannot seek the named pipe")
|
||||
end
|
||||
|
||||
function lib.flush(handle)
|
||||
C.FlushFileBuffers(handle)
|
||||
end
|
||||
|
||||
function lib.available(handle)
|
||||
if handle == nil or handle == INVALID_HANDLE_VALUE then
|
||||
return 0
|
||||
end
|
||||
|
||||
local bytes_available = FFI.new("DWORD[1]")
|
||||
|
||||
local success = C.PeekNamedPipe(handle, nil, 0, nil, bytes_available, nil)
|
||||
|
||||
if success == 0 then
|
||||
return 0
|
||||
end
|
||||
|
||||
return tonumber(bytes_available[0])
|
||||
end
|
||||
|
||||
function lib.is_alive(handle)
|
||||
if handle == nil or handle == INVALID_HANDLE_VALUE then
|
||||
return false
|
||||
|
|
|
|||
64
res/modules/internal/stream_providers/socket.lua
Normal file
64
res/modules/internal/stream_providers/socket.lua
Normal file
|
|
@ -0,0 +1,64 @@
|
|||
local io_stream = require "core:io_stream"
|
||||
|
||||
local lib = { }
|
||||
|
||||
local sockets = { }
|
||||
|
||||
local next_descriptor = 0
|
||||
|
||||
local function open_descriptor(socket)
|
||||
next_descriptor = next_descriptor + 1
|
||||
|
||||
sockets[next_descriptor] = socket
|
||||
|
||||
return next_descriptor
|
||||
end
|
||||
|
||||
local function require_descriptor(descriptor)
|
||||
if not sockets[descriptor] then
|
||||
error("unknown descriptor")
|
||||
end
|
||||
end
|
||||
|
||||
function lib.read(descriptor, length)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
return sockets[descriptor]:recv(length)
|
||||
end
|
||||
|
||||
function lib.write(descriptor, data)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
sockets[descriptor]:send(data)
|
||||
end
|
||||
|
||||
function lib.available(descriptor)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
return sockets[descriptor]:available()
|
||||
end
|
||||
|
||||
function lib.is_alive(descriptor)
|
||||
local socket = sockets[descriptor]
|
||||
|
||||
return socket ~= nil and socket:is_alive()
|
||||
end
|
||||
|
||||
function lib.close(descriptor)
|
||||
require_descriptor(descriptor)
|
||||
|
||||
sockets[descriptor]:close()
|
||||
sockets[descriptor] = nil
|
||||
end
|
||||
|
||||
return function(socket, binary_mode)
|
||||
if binary_mode == nil then
|
||||
binary_mode = true
|
||||
end
|
||||
|
||||
return io_stream.new(
|
||||
open_descriptor(socket),
|
||||
binary_mode,
|
||||
lib
|
||||
)
|
||||
end
|
||||
|
|
@ -25,13 +25,13 @@ local ALL_FLUSH_MODES = {
|
|||
local CR = string.byte('\r')
|
||||
local LF = string.byte('\n')
|
||||
|
||||
local function readFully(result, readFunc)
|
||||
local function read_fully(result, read_func)
|
||||
local isTable = type(result) == "table"
|
||||
|
||||
local buf
|
||||
|
||||
repeat
|
||||
buf = readFunc(MAX_BUFFER_SIZE)
|
||||
buf = read_func(MAX_BUFFER_SIZE)
|
||||
|
||||
if isTable then
|
||||
for i = 1, #buf do
|
||||
|
|
@ -42,41 +42,46 @@ local function readFully(result, readFunc)
|
|||
end
|
||||
|
||||
--[[
|
||||
|
||||
descriptor - descriptor of stream for provided I/O library
|
||||
binaryMode - if enabled, most methods will expect bytes instead of strings
|
||||
ioLib - I/O library. Should include the following functions:
|
||||
read(descriptor: int, length: int) -> Bytearray
|
||||
May return bytearray with a smaller size if bytes have not arrived yet or have run out
|
||||
May throw error if descriptor isn't readable
|
||||
write(descriptor: int, data: Bytearray)
|
||||
flush(descriptor: int)
|
||||
May throw error if descriptor isn't writeable
|
||||
[optional] seek(descriptor: int, mode: string, offset: int)
|
||||
Mode may be 'b' (relative begin), 'c' (relative current), 'e' (relative end + 1)
|
||||
[optional] tell(descriptor: int) -> int
|
||||
[optional] flush(descriptor: int)
|
||||
[optional] available(descriptor: int) -> int
|
||||
is_alive(descriptor: int) -> bool
|
||||
close(descriptor: int)
|
||||
--]]
|
||||
|
||||
function io_stream.new(descriptor, binaryMode, ioLib, mode, flushMode)
|
||||
function io_stream.new(descriptor, binary_mode, io_lib, mode, flush_mode)
|
||||
mode = mode or DEFAULT_MODE
|
||||
flushMode = flushMode or FLUSH_MODE_ALL
|
||||
flush_mode = flush_mode or FLUSH_MODE_ALL
|
||||
|
||||
local self = setmetatable({}, io_stream)
|
||||
|
||||
self.descriptor = descriptor
|
||||
self.binaryMode = binaryMode
|
||||
self.maxBufferSize = MAX_BUFFER_SIZE
|
||||
self.ioLib = ioLib
|
||||
self.binary_mode = binary_mode
|
||||
self.max_buffer_size = MAX_BUFFER_SIZE
|
||||
self.io_lib = io_lib
|
||||
|
||||
self:set_mode(mode)
|
||||
self:set_flush_mode(flushMode)
|
||||
self:set_flush_mode(flush_mode)
|
||||
|
||||
return self
|
||||
end
|
||||
|
||||
function io_stream:is_binary_mode()
|
||||
return self.binaryMode
|
||||
return self.binary_mode
|
||||
end
|
||||
|
||||
function io_stream:set_binary_mode(binaryMode)
|
||||
self.binaryMode = binaryMode ~= nil
|
||||
function io_stream:set_binary_mode(binary_mode)
|
||||
self.binary_mode = binary_mode ~= nil
|
||||
end
|
||||
|
||||
function io_stream:get_mode()
|
||||
|
|
@ -88,124 +93,141 @@ function io_stream:set_mode(mode)
|
|||
error("invalid stream mode: "..mode)
|
||||
end
|
||||
|
||||
if self.mode == BUFFERED_MODE then
|
||||
self.writeBuffer:clear()
|
||||
self.readBuffer:clear()
|
||||
if self.write_buffer then
|
||||
self.write_buffer:clear()
|
||||
self.write_buffer = nil
|
||||
end
|
||||
|
||||
if mode == BUFFERED_MODE and not self.writeBuffer then
|
||||
self.writeBuffer = Bytearray()
|
||||
self.readBuffer = Bytearray()
|
||||
if self.read_buffer then
|
||||
self.read_buffer:clear()
|
||||
self.read_buffer = nil
|
||||
end
|
||||
|
||||
if mode == BUFFERED_MODE then
|
||||
self.write_buffer = Bytearray(self.max_buffer_size)
|
||||
self.read_buffer = Bytearray(self.max_buffer_size)
|
||||
end
|
||||
|
||||
self.mode = mode
|
||||
end
|
||||
|
||||
function io_stream:get_flush_mode()
|
||||
return self.flushMode
|
||||
return self.flush_mode
|
||||
end
|
||||
|
||||
function io_stream:set_flush_mode(flushMode)
|
||||
if not table.has(ALL_FLUSH_MODES, flushMode) then
|
||||
error("invalid flush mode: "..flushMode)
|
||||
function io_stream:set_flush_mode(flush_mode)
|
||||
if not table.has(ALL_FLUSH_MODES, flush_mode) then
|
||||
error("invalid flush mode: " .. flush_mode)
|
||||
end
|
||||
|
||||
self.flushMode = flushMode
|
||||
self.flush_mode = flush_mode
|
||||
end
|
||||
|
||||
function io_stream:get_max_buffer_size()
|
||||
return self.maxBufferSize
|
||||
return self.max_buffer_size
|
||||
end
|
||||
|
||||
function io_stream:set_max_buffer_size(maxBufferSize)
|
||||
self.maxBufferSize = maxBufferSize
|
||||
function io_stream:set_max_buffer_size(max_buffer_size)
|
||||
self.max_buffer_size = max_buffer_size
|
||||
|
||||
self.write_buffer = Bytearray(self.max_buffer_size)
|
||||
self.read_buffer = Bytearray(self.max_buffer_size)
|
||||
end
|
||||
|
||||
function io_stream:available(length)
|
||||
local available = self.io_lib.available and self.io_lib.available(self.descriptor) or 0
|
||||
|
||||
if self.mode == BUFFERED_MODE then
|
||||
self:__update_read_buffer()
|
||||
available = available + #self.read_buffer
|
||||
end
|
||||
|
||||
if not length then
|
||||
return #self.readBuffer
|
||||
else
|
||||
return #self.readBuffer >= length
|
||||
end
|
||||
if not length then
|
||||
return available
|
||||
else
|
||||
return available >= length
|
||||
end
|
||||
end
|
||||
|
||||
function io_stream:__update_read_buffer()
|
||||
local readed = Bytearray()
|
||||
|
||||
readFully(readed, function(length) return self.ioLib.read(self.descriptor, length) end)
|
||||
|
||||
self.readBuffer:append(readed)
|
||||
|
||||
if #self.readBuffer > self.maxBufferSize then
|
||||
error "buffer overflow"
|
||||
end
|
||||
end
|
||||
|
||||
function io_stream:__read(length)
|
||||
function io_stream:__read(length, from_read_fully)
|
||||
if self.mode == YIELD_MODE then
|
||||
if from_read_fully then
|
||||
return self.io_lib.read(self.descriptor, length)
|
||||
end
|
||||
|
||||
local buffer = Bytearray()
|
||||
|
||||
while #buffer < length do
|
||||
buffer:append(self.ioLib.read(self.descriptor, length - #buffer))
|
||||
buffer:append(self.io_lib.read(self.descriptor, length - #buffer))
|
||||
|
||||
if #buffer < length then coroutine.yield() end
|
||||
end
|
||||
|
||||
return buffer
|
||||
elseif self.mode == BUFFERED_MODE then
|
||||
self:__update_read_buffer()
|
||||
local buf_len = #self.read_buffer
|
||||
|
||||
if #self.readBuffer < length then
|
||||
error "buffer underflow"
|
||||
if buf_len < length then
|
||||
self.read_buffer:append(
|
||||
self.io_lib.read(self.descriptor, self.max_buffer_size - buf_len)
|
||||
)
|
||||
end
|
||||
|
||||
local copy
|
||||
buf_len = #self.read_buffer
|
||||
|
||||
if #self.readBuffer == length then
|
||||
copy = Bytearray()
|
||||
|
||||
copy:append(self.readBuffer)
|
||||
length = math.min(buf_len, length)
|
||||
|
||||
self.readBuffer:clear()
|
||||
local copy = self.read_buffer:slice(1, length)
|
||||
|
||||
if buf_len == length then
|
||||
self.read_buffer:clear()
|
||||
else
|
||||
copy = Bytearray()
|
||||
|
||||
for i = 1, length do
|
||||
copy[i] = self.readBuffer[i]
|
||||
end
|
||||
|
||||
self.readBuffer:remove(1, length)
|
||||
self.read_buffer:remove(1, length)
|
||||
end
|
||||
|
||||
return copy
|
||||
elseif self.mode == DEFAULT_MODE then
|
||||
return self.ioLib.read(self.descriptor, length)
|
||||
return self.io_lib.read(self.descriptor, length)
|
||||
end
|
||||
end
|
||||
|
||||
function io_stream:__write(data)
|
||||
if self.mode == BUFFERED_MODE then
|
||||
self.writeBuffer:append(data)
|
||||
local data_length = #data
|
||||
|
||||
if #self.writeBuffer > self.maxBufferSize then
|
||||
error "buffer overflow"
|
||||
if #self.write_buffer + data_length > self.max_buffer_size then
|
||||
self:flush()
|
||||
end
|
||||
|
||||
if data_length > self.max_buffer_size then
|
||||
local to_write = math.floor(data_length / self.max_buffer_size) * self.max_buffer_size
|
||||
local to_save = data_length - to_write
|
||||
|
||||
self.io_lib.write(self.descriptor, data:slice(1, to_write))
|
||||
|
||||
self:flush()
|
||||
|
||||
self.write_buffer = data:slice(to_write + 1, to_save)
|
||||
else self.write_buffer:append(data) end
|
||||
elseif self.mode == DEFAULT_MODE or self.mode == YIELD_MODE then
|
||||
return self.ioLib.write(self.descriptor, data)
|
||||
return self.io_lib.write(self.descriptor, data)
|
||||
end
|
||||
end
|
||||
|
||||
function io_stream:read_fully(useTable)
|
||||
if self.binaryMode then
|
||||
local result = useTable and Bytearray() or { }
|
||||
function io_stream:read_fully(use_table)
|
||||
if self.binary_mode then
|
||||
local result = use_table and { } or Bytearray()
|
||||
|
||||
readFully(result, function() return self:__read(self.maxBufferSize) end)
|
||||
local avail = self:available()
|
||||
|
||||
if avail == 0 then
|
||||
avail = self.max_buffer_size
|
||||
end
|
||||
|
||||
read_fully(result, function() return self:__read(avail, true) end)
|
||||
|
||||
return result
|
||||
else
|
||||
if useTable then
|
||||
if use_table then
|
||||
local lines = { }
|
||||
|
||||
local line
|
||||
|
|
@ -220,7 +242,7 @@ function io_stream:read_fully(useTable)
|
|||
else
|
||||
local result = Bytearray()
|
||||
|
||||
readFully(result, function() return self:__read(self.maxBufferSize) end)
|
||||
read_fully(result, function() return self:__read(self.max_buffer_size) end)
|
||||
|
||||
return utf8.tostring(result)
|
||||
end
|
||||
|
|
@ -262,57 +284,61 @@ function io_stream:write_line(str)
|
|||
self:__write(utf8.tobytes(str .. "\n"))
|
||||
end
|
||||
|
||||
function io_stream:read(arg, useTable)
|
||||
local argType = type(arg)
|
||||
function io_stream:read(arg, use_table)
|
||||
local arg_type = type(arg)
|
||||
|
||||
if self.binaryMode then
|
||||
local byteArr
|
||||
if self.binary_mode then
|
||||
local byte_arr
|
||||
|
||||
if argType == "number" then
|
||||
if arg_type == "number" then
|
||||
-- using 'arg' as length
|
||||
|
||||
byteArr = self:__read(arg)
|
||||
byte_arr = self:__read(arg)
|
||||
|
||||
if useTable == true then
|
||||
if use_table == true then
|
||||
local t = { }
|
||||
|
||||
for i = 1, #byteArr do
|
||||
t[i] = byteArr[i]
|
||||
for i = 1, #byte_arr do
|
||||
t[i] = byte_arr[i]
|
||||
end
|
||||
|
||||
return t
|
||||
else
|
||||
return byteArr
|
||||
return byte_arr
|
||||
end
|
||||
elseif argType == "string" then
|
||||
elseif arg_type == "string" then
|
||||
return byteutil.unpack(
|
||||
arg,
|
||||
self:__read(byteutil.get_size(arg))
|
||||
)
|
||||
elseif argType == nil then
|
||||
elseif arg_type == "nil" then
|
||||
error(
|
||||
"in binary mode the first argument must be a string data format"..
|
||||
" for the library \"byteutil\" or the number of bytes to read"
|
||||
)
|
||||
else
|
||||
error("unknown argument type: "..argType)
|
||||
error("unknown argument type: "..arg_type)
|
||||
end
|
||||
else
|
||||
if not arg then
|
||||
return self:read_line()
|
||||
else
|
||||
local linesCount = arg
|
||||
local trimLastEmptyLines = useTable or true
|
||||
local lines_count = arg
|
||||
local trim_last_empty_lines = use_table
|
||||
|
||||
if linesCount < 0 then error "count of lines to read must be positive" end
|
||||
if use_table == nil then
|
||||
trim_last_empty_lines = true
|
||||
end
|
||||
|
||||
if lines_count < 0 then error "count of lines to read must be positive" end
|
||||
|
||||
local result = { }
|
||||
|
||||
for i = 1, linesCount do
|
||||
for i = 1, lines_count do
|
||||
result[i] = self:read_line()
|
||||
end
|
||||
|
||||
if trimLastEmptyLines then
|
||||
if trim_last_empty_lines then
|
||||
local i = #result
|
||||
|
||||
while i >= 0 do
|
||||
|
|
@ -340,45 +366,53 @@ function io_stream:read(arg, useTable)
|
|||
end
|
||||
|
||||
function io_stream:write(arg, ...)
|
||||
local argType = type(arg)
|
||||
local arg_type = type(arg)
|
||||
|
||||
if self.binaryMode then
|
||||
local byteArr
|
||||
if self.binary_mode then
|
||||
local byte_arr
|
||||
|
||||
if argType ~= "string" then
|
||||
if arg_type ~= "string" then
|
||||
-- using arg as bytes table/bytearray
|
||||
|
||||
if argType == "table" then
|
||||
byteArr = Bytearray(arg)
|
||||
if arg_type == "table" then
|
||||
byte_arr = Bytearray(arg)
|
||||
else
|
||||
byteArr = arg
|
||||
byte_arr = arg
|
||||
end
|
||||
else
|
||||
byteArr = byteutil.pack(arg, ...)
|
||||
byte_arr = byteutil.pack(arg, ...)
|
||||
end
|
||||
|
||||
self:__write(byteArr)
|
||||
self:__write(byte_arr)
|
||||
else
|
||||
if argType == "string" then
|
||||
if arg_type == "string" then
|
||||
self:write_line(arg)
|
||||
elseif argType == "table" then
|
||||
elseif arg_type == "table" then
|
||||
for i = 1, #arg do
|
||||
self:write_line(arg[i])
|
||||
end
|
||||
else error("unknown argument type: "..argType) end
|
||||
else error("unknown argument type: "..arg_type) end
|
||||
end
|
||||
end
|
||||
|
||||
function io_stream:seek(mode, offset)
|
||||
self.ioLib.seek(self.descriptor, mode, offset)
|
||||
if not self.io_lib.seek then
|
||||
error("cannot seek this stream")
|
||||
end
|
||||
|
||||
self.io_lib.seek(self.descriptor, mode, offset)
|
||||
end
|
||||
|
||||
function io_stream:tell()
|
||||
return self.ioLib.tell(self.descriptor)
|
||||
if not self.io_lib.tell then
|
||||
error("cannot tell this stream")
|
||||
end
|
||||
|
||||
return self.io_lib.tell(self.descriptor)
|
||||
end
|
||||
|
||||
function io_stream:is_alive()
|
||||
return self.ioLib.is_alive(self.descriptor)
|
||||
return self.io_lib.is_alive(self.descriptor)
|
||||
end
|
||||
|
||||
function io_stream:is_closed()
|
||||
|
|
@ -387,20 +421,24 @@ end
|
|||
|
||||
function io_stream:close()
|
||||
if self.mode == BUFFERED_MODE then
|
||||
self.readBuffer:clear()
|
||||
self.writeBuffer:clear()
|
||||
self.read_buffer:clear()
|
||||
self.write_buffer:clear()
|
||||
end
|
||||
|
||||
return self.ioLib.close(self.descriptor)
|
||||
return self.io_lib.close(self.descriptor)
|
||||
end
|
||||
|
||||
function io_stream:flush()
|
||||
if self.mode == BUFFERED_MODE and #self.writeBuffer > 0 then
|
||||
self.ioLib.write(self.descriptor, self.writeBuffer)
|
||||
self.writeBuffer:clear()
|
||||
if self.mode == BUFFERED_MODE and #self.write_buffer > 0 then
|
||||
self.io_lib.write(self.descriptor, self.write_buffer)
|
||||
self.write_buffer:clear()
|
||||
end
|
||||
|
||||
if self.flushMode ~= FLUSH_MODE_ONLY_BUFFER then self.ioLib.flush(self.descriptor) end
|
||||
if self.flush_mode ~= FLUSH_MODE_ONLY_BUFFER then
|
||||
if self.io_lib.flush then
|
||||
self.io_lib.flush(self.descriptor)
|
||||
elseif self:is_closed() then error("stream is closed") end
|
||||
end
|
||||
end
|
||||
|
||||
return io_stream
|
||||
return io_stream
|
||||
|
|
@ -49,6 +49,7 @@ local Socket = {__index={
|
|||
end
|
||||
return self:recv(length, usetable)
|
||||
end,
|
||||
as_stream=network.__as_stream,
|
||||
close=function(self) return network.__close(self.id) end,
|
||||
available=function(self) return network.__available(self.id) or 0 end,
|
||||
is_alive=function(self) return network.__is_alive(self.id) end,
|
||||
|
|
|
|||
|
|
@ -336,6 +336,10 @@ else
|
|||
os.pid = ffi.C.getpid()
|
||||
end
|
||||
|
||||
require("core:io_stream").wrap_bytearray = require "core:internal/stream_providers/bytearray"
|
||||
|
||||
network.__as_stream = require "core:internal/stream_providers/socket"
|
||||
|
||||
math.randomseed(time.uptime() * 1536227939)
|
||||
|
||||
rules = require "core:internal/rules"
|
||||
|
|
|
|||
|
|
@ -66,6 +66,33 @@ void io_descriptors::flush(int id) {
|
|||
}
|
||||
}
|
||||
|
||||
int io_descriptors::available(int id) {
|
||||
if (!is_readable(id))
|
||||
return 0;
|
||||
|
||||
auto* stream = descriptors[id]->in.get();
|
||||
|
||||
auto current_pos = stream->tellg();
|
||||
|
||||
if (current_pos == -1) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
stream->seekg(0, std::ios::end);
|
||||
|
||||
auto end_pos = stream->tellg();
|
||||
|
||||
stream->seekg(current_pos, std::ios::beg);
|
||||
|
||||
auto remaining = end_pos - current_pos;
|
||||
|
||||
if (remaining > std::numeric_limits<int>::max()) {
|
||||
return std::numeric_limits<int>::max();
|
||||
}
|
||||
|
||||
return static_cast<int>(remaining);
|
||||
}
|
||||
|
||||
bool io_descriptors::has_descriptor(int id) {
|
||||
return id >= 0 && id < static_cast<int>(::descriptors.size()) &&
|
||||
::descriptors[id].has_value();
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ namespace scripting::io_descriptors {
|
|||
std::ostream& require_output(int id);
|
||||
|
||||
void flush(int id);
|
||||
int available(int id);
|
||||
|
||||
bool has_descriptor(int id);
|
||||
|
||||
|
|
|
|||
|
|
@ -400,6 +400,20 @@ static int l_flush_descriptor(lua::State* L) {
|
|||
return 0;
|
||||
}
|
||||
|
||||
static int l_available_descriptor(lua::State* L) {
|
||||
int descriptor = lua::tointeger(L, 1);
|
||||
|
||||
if (!io_descriptors::has_descriptor(descriptor)) {
|
||||
throw std::runtime_error("unknown descriptor");
|
||||
}
|
||||
|
||||
if (!io_descriptors::is_readable(descriptor)) {
|
||||
throw std::runtime_error("descriptor is not readable");
|
||||
}
|
||||
|
||||
return lua::pushinteger(L, io_descriptors::available(descriptor));
|
||||
}
|
||||
|
||||
static int l_close_descriptor(lua::State* L) {
|
||||
int descriptor = lua::tointeger(L, 1);
|
||||
|
||||
|
|
@ -447,6 +461,7 @@ const luaL_Reg filelib[] = {
|
|||
{"__seek_descriptor", lua::wrap<l_seek_descriptor>},
|
||||
{"__tell_descriptor", lua::wrap<l_tell_descriptor>},
|
||||
{"__flush_descriptor", lua::wrap<l_flush_descriptor>},
|
||||
{"__available_descriptor", lua::wrap<l_available_descriptor>},
|
||||
{"__close_descriptor", lua::wrap<l_close_descriptor>},
|
||||
{"__close_all_descriptors", lua::wrap<l_close_all_descriptors>},
|
||||
{nullptr, nullptr}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue