Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 0 additions & 5 deletions .vscode/settings.json

This file was deleted.

7 changes: 5 additions & 2 deletions dev.hxml
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
tests.hxml
-lib travix
-lib tink_io
-lib hxnodejs
-js bin/node/tests.js

-neko bin/neko/tests.n

# -lib hxnodejs
# -js bin/node/tests.js
1 change: 1 addition & 0 deletions extraParams.hxml
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
--macro tink.io.Boot.boot()
3 changes: 2 additions & 1 deletion haxe_libraries/tink_io.hxml
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,5 @@
-cp src

-lib tink_chunk
-lib tink_streams
-lib tink_streams
extraParams.hxml
46 changes: 46 additions & 0 deletions hxformat.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
{
"indentation": {
"character": " ",
"tabWidth": 2,
"trailingWhitespace": false
},
"sameLine": {
"functionBody": "keep",
"anonFunctionBody": "keep",
"ifBody": "keep",
"elseBody": "keep",
"forBody": "keep",
"whileBody": "keep",
"doWhileBody": "keep",
"tryBody": "keep",
"catchBody": "keep",
"caseBody": "keep",
"expressionCase": "keep",
"returnBody": "keep",
"returnBodySingleLine": "keep"
},
"wrapping": {
"maxLineLength": 160,
"callParameter": { "defaultWrap": "keep" },
"arrayWrap": { "defaultWrap": "keep" },
"objectLiteral": { "defaultWrap": "keep" },
"methodChain": { "defaultWrap": "keep" },
"opBoolChain": { "defaultWrap": "keep" },
"opAddSubChain": { "defaultWrap": "keep" }
},
"emptyLines": {
"finalNewline": false,
"maxAnywhereInFile": 2
},
"whitespace": {
"ifPolicy": "after",
"forPolicy": "none",
"whilePolicy": "none",
"catchPolicy": "after",
"doPolicy": "after",
"objectFieldColonPolicy": "after",
"colonPolicy": "none",
"arrowFunctionsPolicy": "around",
"binopPolicy": "around"
}
}
21 changes: 21 additions & 0 deletions src/tink/io/Boot.hx
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
package tink.io;

#if macro
// this whole file existed because of https://github.com/HaxeFoundation/haxe/issues/12985
class Boot {
static function boot() {
tink.SyntaxHub.transformMain.whenever(function(e) {
if (haxe.macro.Context.defined("java")) {
return macro {
@:pos(e.pos) tink.io.java.OnMainThread.init();
$e;
};
} else {
return e;
}
});
}
}
#else
#error
#end
7 changes: 7 additions & 0 deletions src/tink/io/Sink.hx
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,13 @@ abstract SinkYielding<FailingWith, Result>(SinkObject<FailingWith, Result>)
case { worker: w }: w;
});

#if (sys && target.threaded)
static public function ofSocket(name:String, socket:sys.net.Socket, pool:tink.io.std.SelectPool, ?options:{ ?worker:Worker }):RealSink
return tink.io.std.SocketSink.wrap(name, socket, pool, switch options {
case null | { worker: null }: Worker.get();
case { worker: w }: w;
});
#end

}

Expand Down
11 changes: 11 additions & 0 deletions src/tink/io/Source.hx
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,17 @@ abstract Source<E>(SourceObject<E>) from SourceObject<E> to SourceObject<E> to S
case v: v;
}), 0);
}

#if (sys && target.threaded)
@:noUsing static public inline function ofSocket(name:String, socket:sys.net.Socket, pool:tink.io.std.SelectPool, ?options:{ ?chunkSize: Int, ?worker:Worker }):RealSource {
if (options == null)
options = {};
return tink.io.std.SocketSource.wrap(name, socket, pool, options.worker.ensure(), switch options.chunkSize {
case null: 0x10000;
case v: v;
});
}
#end

public function chunked():Stream<Chunk, E>
return this;
Expand Down
2 changes: 1 addition & 1 deletion src/tink/io/java/JavaFileSource.hx
Original file line number Diff line number Diff line change
Expand Up @@ -66,4 +66,4 @@ private class ReadHandler implements CompletionHandler<Integer, ByteBuffer> {
cb.invoke(Fail(Error.withData('Read failed for "${parent.name}", reason: ' + exc.getMessage(), exc)));
});
}
}
}
6 changes: 5 additions & 1 deletion src/tink/io/java/JavaSocketSink.hx
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,11 @@ class JavaSocketSink extends SinkBase<Error, Noise> {
});

if (options.end)
ret.handle(function (end) channel.shutdownOutput());
ret.handle(function (end) {
try channel.shutdownOutput()
catch (e:java.nio.channels.ClosedChannelException) {}
catch (e:java.io.IOException) {}
});

return ret.map(function (c) return c.toResult(Noise));
}
Expand Down
28 changes: 14 additions & 14 deletions src/tink/io/java/OnMainThread.hx
Original file line number Diff line number Diff line change
@@ -1,17 +1,17 @@
package tink.io.java;

/**
Java NIO async channel callbacks run on a JDK thread pool without a Haxe event loop.
Re-dispatch stream completions so downstream code (e.g. haxe.Timer) runs on a safe thread.
**/
class OnMainThread {
public static function run(fn:Void->Void):Void {
#if tink_runloop
tink.RunLoop.current.work(fn);
#elseif java
haxe.MainLoop.add(fn);
#else
fn();
#end
}
}
public static function init() {
run(noop);
}

public static function run(fn:Void->Void):Void {
#if haxe5
sys.thread.Thread.main().events.run(fn);
#else
haxe.EntryPoint.runInMainThread(fn);
#end
}

static function noop() {}
}
Loading
Loading