feat: implement async file I/O with open/close/readBytes and scheduler error handling
Add foreign class File with static open, read, size methods and instance methods for close, size, readBytes. Implement async file operations using libuv callbacks that resume fibers on completion. Introduce schedulerResumeError and resumeError_ to propagate I/O errors as fiber errors. Add setExitCode to CLI vm for runtime error exit handling. Bump MAX_METHODS_PER_CLASS from 2 to 9 to accommodate new File methods.
This commit is contained in:
+155
-16
@@ -9,32 +9,171 @@
|
||||
|
||||
#include <stdio.h>
|
||||
|
||||
// Called by libuv when the stat call for size completes.
|
||||
static void sizeCallback(uv_fs_t* request)
|
||||
void fileAllocate(WrenVM* vm)
|
||||
{
|
||||
// Store the file descriptor in the foreign data, so that we can get to it
|
||||
// in the finalizer.
|
||||
int* fd = (int*)wrenAllocateForeign(vm, sizeof(int));
|
||||
*fd = (int)wrenGetArgumentDouble(vm, 1);
|
||||
}
|
||||
|
||||
void fileFinalize(WrenVM* vm)
|
||||
{
|
||||
int fd = *(int*)wrenGetArgumentForeign(vm, 0);
|
||||
|
||||
// Already closed.
|
||||
if (fd == -1) return;
|
||||
|
||||
uv_fs_t request;
|
||||
uv_fs_close(getLoop(), &request, fd, NULL);
|
||||
uv_fs_req_cleanup(&request);
|
||||
}
|
||||
|
||||
// If [request] failed with an error, sends the runtime error to the VM and
|
||||
// frees the request.
|
||||
//
|
||||
// Returns true if an error was reported.
|
||||
static bool handleRequestError(uv_fs_t* request)
|
||||
{
|
||||
if (request->result >= 0) return false;
|
||||
|
||||
WrenValue* fiber = (WrenValue*)request->data;
|
||||
schedulerResumeError(fiber, uv_strerror((int)request->result));
|
||||
uv_fs_req_cleanup(request);
|
||||
free(request);
|
||||
return true;
|
||||
}
|
||||
|
||||
// Allocates a new request that resumes [fiber] when it completes.
|
||||
uv_fs_t* createRequest(WrenValue* fiber)
|
||||
{
|
||||
uv_fs_t* request = (uv_fs_t*)malloc(sizeof(uv_fs_t));
|
||||
request->data = fiber;
|
||||
return request;
|
||||
}
|
||||
|
||||
// Releases resources used by [request].
|
||||
//
|
||||
// Returns the fiber that should be resumed after [request] completes.
|
||||
WrenValue* freeRequest(uv_fs_t* request)
|
||||
{
|
||||
WrenValue* fiber = (WrenValue*)request->data;
|
||||
|
||||
if (request->result != 0)
|
||||
{
|
||||
schedulerResumeString(fiber, uv_strerror((int)request->result));
|
||||
uv_fs_req_cleanup(request);
|
||||
return;
|
||||
}
|
||||
uv_fs_req_cleanup(request);
|
||||
free(request);
|
||||
|
||||
return fiber;
|
||||
}
|
||||
|
||||
static void openCallback(uv_fs_t* request)
|
||||
{
|
||||
if (handleRequestError(request)) return;
|
||||
|
||||
double fd = (double)request->result;
|
||||
WrenValue* fiber = freeRequest(request);
|
||||
|
||||
schedulerResumeDouble(fiber, fd);
|
||||
}
|
||||
|
||||
void fileOpen(WrenVM* vm)
|
||||
{
|
||||
const char* path = wrenGetArgumentString(vm, 1);
|
||||
uv_fs_t* request = createRequest(wrenGetArgumentValue(vm, 2));
|
||||
|
||||
// TODO: Allow controlling flags and modes.
|
||||
uv_fs_open(getLoop(), request, path, O_RDONLY, S_IRUSR, openCallback);
|
||||
}
|
||||
|
||||
// Called by libuv when the stat call for size completes.
|
||||
static void sizeCallback(uv_fs_t* request)
|
||||
{
|
||||
if (handleRequestError(request)) return;
|
||||
|
||||
double size = (double)request->statbuf.st_size;
|
||||
uv_fs_req_cleanup(request);
|
||||
WrenValue* fiber = freeRequest(request);
|
||||
|
||||
schedulerResumeDouble(fiber, size);
|
||||
}
|
||||
|
||||
void fileStartSize(WrenVM* vm)
|
||||
void fileSizePath(WrenVM* vm)
|
||||
{
|
||||
const char* path = wrenGetArgumentString(vm, 1);
|
||||
WrenValue* fiber = wrenGetArgumentValue(vm, 2);
|
||||
|
||||
// Store the fiber to resume when the request completes.
|
||||
uv_fs_t* request = (uv_fs_t*)malloc(sizeof(uv_fs_t));
|
||||
request->data = fiber;
|
||||
|
||||
uv_fs_t* request = createRequest(wrenGetArgumentValue(vm, 2));
|
||||
uv_fs_stat(getLoop(), request, path, sizeCallback);
|
||||
}
|
||||
|
||||
static void closeCallback(uv_fs_t* request)
|
||||
{
|
||||
if (handleRequestError(request)) return;
|
||||
|
||||
WrenValue* fiber = freeRequest(request);
|
||||
schedulerResume(fiber);
|
||||
}
|
||||
|
||||
void fileClose(WrenVM* vm)
|
||||
{
|
||||
int* foreign = (int*)wrenGetArgumentForeign(vm, 0);
|
||||
int fd = *foreign;
|
||||
|
||||
// If it's already closed, we're done.
|
||||
if (fd == -1)
|
||||
{
|
||||
wrenReturnBool(vm, true);
|
||||
return;
|
||||
}
|
||||
|
||||
// Mark it closed immediately.
|
||||
*foreign = -1;
|
||||
|
||||
uv_fs_t* request = createRequest(wrenGetArgumentValue(vm, 1));
|
||||
uv_fs_close(getLoop(), request, fd, closeCallback);
|
||||
wrenReturnBool(vm, false);
|
||||
}
|
||||
|
||||
void fileDescriptor(WrenVM* vm)
|
||||
{
|
||||
int* foreign = (int*)wrenGetArgumentForeign(vm, 0);
|
||||
int fd = *foreign;
|
||||
wrenReturnDouble(vm, fd);
|
||||
}
|
||||
|
||||
static void readBytesCallback(uv_fs_t* request)
|
||||
{
|
||||
if (handleRequestError(request)) return;
|
||||
|
||||
uv_buf_t buffer = request->bufs[0];
|
||||
WrenValue* fiber = freeRequest(request);
|
||||
|
||||
// TODO: Having to copy the bytes here is a drag. It would be good if Wren's
|
||||
// embedding API supported a way to *give* it bytes that were previously
|
||||
// allocated using Wren's own allocator.
|
||||
schedulerResumeBytes(fiber, buffer.base, buffer.len);
|
||||
|
||||
// TODO: Likewise, freeing this after we resume is lame.
|
||||
free(buffer.base);
|
||||
}
|
||||
|
||||
void fileReadBytes(WrenVM* vm)
|
||||
{
|
||||
uv_fs_t* request = createRequest(wrenGetArgumentValue(vm, 2));
|
||||
|
||||
int fd = *(int*)wrenGetArgumentForeign(vm, 0);
|
||||
// TODO: Assert fd != -1.
|
||||
|
||||
uv_buf_t buffer;
|
||||
buffer.len = (size_t)wrenGetArgumentDouble(vm, 1);
|
||||
buffer.base = (char*)malloc(buffer.len);
|
||||
|
||||
// TODO: Allow passing in offset.
|
||||
uv_fs_read(getLoop(), request, fd, &buffer, 1, 0, readBytesCallback);
|
||||
}
|
||||
|
||||
void fileSize(WrenVM* vm)
|
||||
{
|
||||
uv_fs_t* request = createRequest(wrenGetArgumentValue(vm, 1));
|
||||
|
||||
int fd = *(int*)wrenGetArgumentForeign(vm, 0);
|
||||
// TODO: Assert fd != -1.
|
||||
|
||||
uv_fs_fstat(getLoop(), request, fd, sizeCallback);
|
||||
}
|
||||
|
||||
+60
-6
@@ -1,15 +1,69 @@
|
||||
import "scheduler" for Scheduler
|
||||
|
||||
class File {
|
||||
static size(path) {
|
||||
foreign class File {
|
||||
static open(path) {
|
||||
if (!(path is String)) Fiber.abort("Path must be a string.")
|
||||
|
||||
startSize_(path, Fiber.current)
|
||||
var result = Scheduler.runNextScheduled_()
|
||||
if (result is String) Fiber.abort(result)
|
||||
open_(path, Fiber.current)
|
||||
var fd = Scheduler.runNextScheduled_()
|
||||
return new_(fd)
|
||||
}
|
||||
|
||||
static open(path, fn) {
|
||||
var file = open(path)
|
||||
var fiber = Fiber.new { fn.call(file) }
|
||||
|
||||
// Poor man's finally. Can we make this more elegant?
|
||||
var result = fiber.try()
|
||||
file.close()
|
||||
|
||||
// TODO: Want something like rethrow since now the callstack ends here. :(
|
||||
if (fiber.error != null) Fiber.abort(fiber.error)
|
||||
return result
|
||||
}
|
||||
|
||||
foreign static startSize_(path, fiber)
|
||||
static read(path) {
|
||||
return File.open(path) {|file| file.readBytes(file.size) }
|
||||
}
|
||||
|
||||
static size(path) {
|
||||
if (!(path is String)) Fiber.abort("Path must be a string.")
|
||||
|
||||
sizePath_(path, Fiber.current)
|
||||
return Scheduler.runNextScheduled_()
|
||||
}
|
||||
|
||||
construct new_(fd) {}
|
||||
|
||||
close() {
|
||||
if (close_(Fiber.current)) return
|
||||
Scheduler.runNextScheduled_()
|
||||
}
|
||||
|
||||
isOpen { descriptor != -1 }
|
||||
|
||||
size {
|
||||
if (!isOpen) Fiber.abort("File is not open.")
|
||||
|
||||
size_(Fiber.current)
|
||||
return Scheduler.runNextScheduled_()
|
||||
}
|
||||
|
||||
readBytes(count) {
|
||||
if (!isOpen) Fiber.abort("File is not open.")
|
||||
if (!(count is Num)) Fiber.abort("Count must be an integer.")
|
||||
if (!count.isInteger) Fiber.abort("Count must be an integer.")
|
||||
if (count < 0) Fiber.abort("Count cannot be negative.")
|
||||
|
||||
readBytes_(count, Fiber.current)
|
||||
return Scheduler.runNextScheduled_()
|
||||
}
|
||||
|
||||
foreign static open_(path, fiber)
|
||||
foreign static sizePath_(path, fiber)
|
||||
|
||||
foreign close_(fiber)
|
||||
foreign descriptor
|
||||
foreign readBytes_(count, fiber)
|
||||
foreign size_(fiber)
|
||||
}
|
||||
|
||||
+60
-6
@@ -2,16 +2,70 @@
|
||||
static const char* ioModuleSource =
|
||||
"import \"scheduler\" for Scheduler\n"
|
||||
"\n"
|
||||
"class File {\n"
|
||||
" static size(path) {\n"
|
||||
"foreign class File {\n"
|
||||
" static open(path) {\n"
|
||||
" if (!(path is String)) Fiber.abort(\"Path must be a string.\")\n"
|
||||
"\n"
|
||||
" startSize_(path, Fiber.current)\n"
|
||||
" var result = Scheduler.runNextScheduled_()\n"
|
||||
" if (result is String) Fiber.abort(result)\n"
|
||||
" open_(path, Fiber.current)\n"
|
||||
" var fd = Scheduler.runNextScheduled_()\n"
|
||||
" return new_(fd)\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" static open(path, fn) {\n"
|
||||
" var file = open(path)\n"
|
||||
" var fiber = Fiber.new { fn.call(file) }\n"
|
||||
"\n"
|
||||
" // Poor man's finally. Can we make this more elegant?\n"
|
||||
" var result = fiber.try()\n"
|
||||
" file.close()\n"
|
||||
"\n"
|
||||
" // TODO: Want something like rethrow since now the callstack ends here. :(\n"
|
||||
" if (fiber.error != null) Fiber.abort(fiber.error)\n"
|
||||
" return result\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" foreign static startSize_(path, fiber)\n"
|
||||
" static read(path) {\n"
|
||||
" return File.open(path) {|file| file.readBytes(file.size) }\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" static size(path) {\n"
|
||||
" if (!(path is String)) Fiber.abort(\"Path must be a string.\")\n"
|
||||
"\n"
|
||||
" sizePath_(path, Fiber.current)\n"
|
||||
" return Scheduler.runNextScheduled_()\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" construct new_(fd) {}\n"
|
||||
"\n"
|
||||
" close() {\n"
|
||||
" if (close_(Fiber.current)) return\n"
|
||||
" Scheduler.runNextScheduled_()\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" isOpen { descriptor != -1 }\n"
|
||||
"\n"
|
||||
" size {\n"
|
||||
" if (!isOpen) Fiber.abort(\"File is not open.\")\n"
|
||||
"\n"
|
||||
" size_(Fiber.current)\n"
|
||||
" return Scheduler.runNextScheduled_()\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" readBytes(count) {\n"
|
||||
" if (!isOpen) Fiber.abort(\"File is not open.\")\n"
|
||||
" if (!(count is Num)) Fiber.abort(\"Count must be an integer.\")\n"
|
||||
" if (!count.isInteger) Fiber.abort(\"Count must be an integer.\")\n"
|
||||
" if (count < 0) Fiber.abort(\"Count cannot be negative.\")\n"
|
||||
"\n"
|
||||
" readBytes_(count, Fiber.current)\n"
|
||||
" return Scheduler.runNextScheduled_()\n"
|
||||
" }\n"
|
||||
"\n"
|
||||
" foreign static open_(path, fiber)\n"
|
||||
" foreign static sizePath_(path, fiber)\n"
|
||||
"\n"
|
||||
" foreign close_(fiber)\n"
|
||||
" foreign descriptor\n"
|
||||
" foreign readBytes_(count, fiber)\n"
|
||||
" foreign size_(fiber)\n"
|
||||
"}\n";
|
||||
|
||||
+36
-6
@@ -12,33 +12,63 @@
|
||||
// one.
|
||||
static WrenValue* resume;
|
||||
static WrenValue* resumeWithArg;
|
||||
static WrenValue* resumeError;
|
||||
|
||||
void schedulerCaptureMethods(WrenVM* vm)
|
||||
{
|
||||
resume = wrenGetMethod(vm, "scheduler", "Scheduler", "resume_(_)");
|
||||
resumeWithArg = wrenGetMethod(vm, "scheduler", "Scheduler", "resume_(_,_)");
|
||||
resumeError = wrenGetMethod(vm, "scheduler", "Scheduler", "resumeError_(_,_)");
|
||||
}
|
||||
|
||||
static void callResume(WrenValue* resumeMethod, WrenValue* fiber,
|
||||
const char* argTypes, ...)
|
||||
{
|
||||
va_list args;
|
||||
va_start(args, argTypes);
|
||||
WrenInterpretResult result = wrenCallVarArgs(getVM(), resumeMethod, NULL,
|
||||
argTypes, args);
|
||||
va_end(args);
|
||||
|
||||
wrenReleaseValue(getVM(), fiber);
|
||||
|
||||
// If a runtime error occurs in response to an async operation and nothing
|
||||
// catches the error in the fiber, then exit the CLI.
|
||||
if (result == WREN_RESULT_RUNTIME_ERROR)
|
||||
{
|
||||
uv_stop(getLoop());
|
||||
setExitCode(70); // EX_SOFTWARE.
|
||||
}
|
||||
}
|
||||
|
||||
void schedulerResume(WrenValue* fiber)
|
||||
{
|
||||
wrenCall(getVM(), resume, NULL, "v", fiber);
|
||||
wrenReleaseValue(getVM(), fiber);
|
||||
callResume(resume, fiber, "v", fiber);
|
||||
}
|
||||
|
||||
void schedulerResumeBytes(WrenValue* fiber, const char* bytes, size_t length)
|
||||
{
|
||||
callResume(resumeWithArg, fiber, "va", fiber, bytes, length);
|
||||
}
|
||||
|
||||
void schedulerResumeDouble(WrenValue* fiber, double value)
|
||||
{
|
||||
wrenCall(getVM(), resumeWithArg, NULL, "vd", fiber, value);
|
||||
wrenReleaseValue(getVM(), fiber);
|
||||
callResume(resumeWithArg, fiber, "vd", fiber, value);
|
||||
}
|
||||
|
||||
void schedulerResumeString(WrenValue* fiber, const char* text)
|
||||
{
|
||||
wrenCall(getVM(), resumeWithArg, NULL, "vs", fiber, text);
|
||||
wrenReleaseValue(getVM(), fiber);
|
||||
callResume(resumeWithArg, fiber, "vs", fiber, text);
|
||||
}
|
||||
|
||||
void schedulerResumeError(WrenValue* fiber, const char* error)
|
||||
{
|
||||
callResume(resumeError, fiber, "vs", fiber, error);
|
||||
}
|
||||
|
||||
void schedulerReleaseMethods()
|
||||
{
|
||||
if (resume != NULL) wrenReleaseValue(getVM(), resume);
|
||||
if (resumeWithArg != NULL) wrenReleaseValue(getVM(), resumeWithArg);
|
||||
if (resumeError != NULL) wrenReleaseValue(getVM(), resumeError);
|
||||
}
|
||||
|
||||
@@ -4,8 +4,10 @@
|
||||
#include "wren.h"
|
||||
|
||||
void schedulerResume(WrenValue* fiber);
|
||||
void schedulerResumeBytes(WrenValue* fiber, const char* bytes, size_t length);
|
||||
void schedulerResumeDouble(WrenValue* fiber, double value);
|
||||
void schedulerResumeString(WrenValue* fiber, const char* text);
|
||||
void schedulerResumeError(WrenValue* fiber, const char* error);
|
||||
|
||||
void schedulerReleaseMethods();
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ class Scheduler {
|
||||
// Called by native code.
|
||||
static resume_(fiber) { fiber.transfer() }
|
||||
static resume_(fiber, arg) { fiber.transfer(arg) }
|
||||
static resumeError_(fiber, error) { fiber.transferError(error) }
|
||||
|
||||
static runNextScheduled_() {
|
||||
if (__scheduled == null || __scheduled.isEmpty) {
|
||||
|
||||
@@ -13,6 +13,7 @@ static const char* schedulerModuleSource =
|
||||
" // Called by native code.\n"
|
||||
" static resume_(fiber) { fiber.transfer() }\n"
|
||||
" static resume_(fiber, arg) { fiber.transfer(arg) }\n"
|
||||
" static resumeError_(fiber, error) { fiber.transferError(error) }\n"
|
||||
"\n"
|
||||
" static runNextScheduled_() {\n"
|
||||
" if (__scheduled == null || __scheduled.isEmpty) {\n"
|
||||
|
||||
Reference in New Issue
Block a user