Skip to content
Merged
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
16 changes: 1 addition & 15 deletions .github/tsan_suppressions.txt
Original file line number Diff line number Diff line change
Expand Up @@ -10,24 +10,10 @@
# hloop_process_events read — upstream idiom, worst case is one extra loop iteration.
race:hloop_stop

# Deliberate unsynchronized liveness peek: das_writer_is_connected polls hio_is_opened
# from the tick thread; a momentarily stale answer is the API's documented contract.
race:hio_is_opened

# Cross-thread WebSocket send: das_wss_send / das_wsc_send call hio_write from das threads;
# libhv mutex-guards the write queue but hio_write4 reads io state outside it (#3428).
race:hio_write4

# libhv's lazy default logger calls localtime while its date-header timer calls gmtime: one static buffer.
race:logger_create

# WebSocket client open/close callback bookkeeping inside hv::WebSocketClient (#3428).
race:hv::WebSocketClient::close

# Keep-alive handler reuse vs dasHV deferred (HTTP_STATUS_UNFINISHED) responses: the loop
# thread may parse the next pipelined message while the tick thread still fills the
# previous response object (#3428).
race:HttpHandler::onMessageComplete
race:hio_is_opened

# libdbus's own global locks, taken in one order by dbus_bus_get_private (the Linux tray's
# create) and the reverse by dbus_connection_close / unref (its destroy), both on the main
Expand Down
6 changes: 3 additions & 3 deletions modules/dasHV/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -176,9 +176,9 @@ IF ((NOT DAS_HV_INCLUDED) AND ((NOT ${DAS_HV_DISABLED}) OR (NOT DEFINED DAS_HV_D
ENDIF()
ExternalProject_Add(
LIBHV
URL https://github.com/ithewei/libhv/archive/5da5fc22cdfd1b59ba2b0520abb34f95d050975c.tar.gz
URL_HASH SHA256=b4f774936ccd6981df74b23f9e033e0b3ab584f6682d889490377dc9450b5f54
DOWNLOAD_EXTRACT_TIMESTAMP TRUE
URL https://github.com/ithewei/libhv/archive/303f50c7b23f250c5b620d92c04b41060e1701c2.tar.gz
URL_HASH SHA256=43f425dae61d060913c5ce280c9776dfb98c61bb965a38c5a38d4788a7295201
DOWNLOAD_EXTRACT_TIMESTAMP FALSE
PATCH_COMMAND ${CMAKE_COMMAND} -DLIBHV_SRC_DIR=<SOURCE_DIR> -P ${DAS_HV_DIR}/patch_libhv.cmake
PREFIX ${CMAKE_CURRENT_BINARY_DIR}/libhv
CMAKE_ARGS ${HV_CMAKE_FLAGS}
Expand Down
49 changes: 34 additions & 15 deletions modules/dasHV/src/dasHV.cpp
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
#include "daScript/misc/platform.h"

#include <future>
#include <atomic>
#include <cstdlib>

#include "../../../src/builtin/module_builtin_rtti.h"
Expand Down Expand Up @@ -167,12 +168,14 @@ class WebSocketClient_Adapter : public hv::WebSocketClient, public HvWebSocketCl
WebSocketClient_Adapter ( char * pClass, const StructInfo * info, Context * ctx )
: HvWebSocketClient_Adapter(info), classPtr(pClass), context(ctx) {
onopen = [=]() {
connected.store(true);
lock_guard<mutex> guard(lock);
que.emplace_back([=](){
onOpen();
});
};
onclose = [=]() {
connected.store(false);
lock_guard<mutex> guard(lock);
que.emplace_back([=](){
onClose();
Expand Down Expand Up @@ -223,6 +226,9 @@ class WebSocketClient_Adapter : public hv::WebSocketClient, public HvWebSocketCl
Context * context;
mutex lock;
vector<function<void()>> que;
atomic<bool> connected{false};
public:
bool isConnected() { return connected.load(); }
};

Handle<hv::WebSocketClient> makeWebSocketClient ( const void * pClass, const StructInfo * info, Context * context ) {
Expand Down Expand Up @@ -254,7 +260,11 @@ int das_wsc_send_buf ( Handle<hv::WebSocketClient> h, const char* msg, int32_t l
int das_wsc_close ( Handle<hv::WebSocketClient> h ) {
auto p = HandleRegistry<hv::WebSocketClient>::instance().lookup(h);
if ( !p ) return -1;
return p->close();
const auto & loop = p->loop();
if ( !loop || !loop->isRunning() ) return p->close();
auto client = p.get();
loop->runInLoop([client](){ client->close(); });
return 0;
}

void das_hv_set_log_file ( const char * path ) {
Expand All @@ -264,7 +274,7 @@ void das_hv_set_log_file ( const char * path ) {
bool das_wsc_is_connected ( Handle<hv::WebSocketClient> h ) {
auto p = HandleRegistry<hv::WebSocketClient>::instance().lookup(h);
if ( !p ) return false;
return p->isConnected();
return ((WebSocketClient_Adapter *) p.get())->isConnected();
}

void das_wsc_tick ( Handle<hv::WebSocketClient> h ) {
Expand Down Expand Up @@ -433,14 +443,16 @@ class WebServer_Adapter : public hv::WebSocketServer, public HvWebServer_Adapter
// worker_threads defaults to 1, so this is the loop the connection
// lives on — the send must run here, not on the tick thread.
auto connLoop = this->loop();
auto resp = make_shared<HttpResponse>(*ctx->response);
lock_guard<mutex> guard(lock);
que.emplace_back([context,at,lmb,ctx,connLoop](){
que.emplace_back([context,at,lmb,ctx,connLoop,resp](){
int st = das_invoke_lambda<int>::invoke<HttpRequest*,HttpResponse*>(
context, at, lmb, ctx->request.get(), ctx->response.get());
ctx->response->status_code = (http_status) st;
context, at, lmb, ctx->request.get(), resp.get());
resp->status_code = (http_status) st;
if ( connLoop ) {
connLoop->runInLoop([ctx](){ ctx->send(); });
connLoop->runInLoop([ctx,resp](){ *ctx->response = *resp; ctx->send(); });
} else {
*ctx->response = *resp;
ctx->send();
}
});
Expand Down Expand Up @@ -522,10 +534,15 @@ class WebServer_Adapter : public hv::WebSocketServer, public HvWebServer_Adapter
lock_guard<mutex> guard(writer_lock);
active_writers.erase(w);
}
bool is_writer_open ( hv::HttpResponseWriter * w ) {
lock_guard<mutex> guard(writer_lock);
auto it = active_writers.find(w);
return it != active_writers.end() && it->second.second->load();
}
HttpResponseWriterPtr find_writer ( hv::HttpResponseWriter * w ) {
lock_guard<mutex> guard(writer_lock);
auto it = active_writers.find(w);
return it != active_writers.end() ? it->second : HttpResponseWriterPtr();
return it != active_writers.end() ? it->second.first : HttpResponseWriterPtr();
}
// Streaming route: an async (writer) handler. libhv hands us a live HttpResponseWriter and keeps
// the connection open until we End() it. The das handler runs on the tick thread (its context +
Expand All @@ -535,9 +552,16 @@ class WebServer_Adapter : public hv::WebSocketServer, public HvWebServer_Adapter
void STREAM ( const char * relative_path, Lambda lmb, Context * context, LineInfoArg * at ) {
lock_guard<mutex> guard(lock);
router.Any(relative_path, [this,context,at,lmb](const HttpRequestPtr & req, const HttpResponseWriterPtr & writer) {
auto open = make_shared<atomic<bool>>(true);
auto mark_close = [writer,open](){
if ( !writer->isOpened() ) open->store(false);
else writer->onclose = [open](){ open->store(false); };
};
if ( auto wloop = this->loop(0) ) wloop->runInLoop(mark_close);
else mark_close();
{
lock_guard<mutex> wguard(writer_lock);
active_writers[writer.get()] = writer;
active_writers[writer.get()] = make_pair(writer, open);
}
{
lock_guard<mutex> qguard(lock);
Expand All @@ -556,7 +580,7 @@ class WebServer_Adapter : public hv::WebSocketServer, public HvWebServer_Adapter
mutex lock;
vector<function<void()>> que;
mutex writer_lock;
map<hv::HttpResponseWriter*, HttpResponseWriterPtr> active_writers;
map<hv::HttpResponseWriter*, pair<HttpResponseWriterPtr, shared_ptr<atomic<bool>>>> active_writers;
mutex channel_lock;
map<hv::WebSocketChannel*, Handle<hv::WebSocketChannel>> channel_handles;
};
Expand Down Expand Up @@ -861,15 +885,10 @@ void das_writer_release ( Handle<hv::WebSocketServer> h, hv::HttpResponseWriter
if ( auto adapter = lookup_server(h) ) adapter->release_writer(w);
}

// True while the writer's connection is still open (isOpened: io alive + not disconnected). The
// async write ops never report a dead peer, so a server polls this to evict abandoned streams.
// False for null / unknown / already-released writers.
bool das_writer_is_connected ( Handle<hv::WebSocketServer> h, hv::HttpResponseWriter * w ) {
auto adapter = lookup_server(h);
if ( !adapter || !w ) return false;
auto sp = adapter->find_writer(w);
if ( !sp ) return false;
return sp->isOpened();
return adapter->is_writer_open(w);
}

void das_wss_set_document_root ( Handle<hv::WebSocketServer> h, const char * dir ) {
Expand Down
17 changes: 11 additions & 6 deletions tests/stbimage/test_apng.das
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,11 @@ struct ChunkSummary {
bad : bool // signature mismatch or truncated stream
}

def private temp_output(ext : string) : string {
var error : string
return create_temp_file("stbimage_test", ext, error)
}

def be_u32_at(b0, b1, b2, b3 : uint8) : uint {
return ((uint(b0) << 24u) | (uint(b1) << 16u) | (uint(b2) << 8u) | uint(b3))
}
Expand Down Expand Up @@ -99,7 +104,7 @@ def make_gradient_frame(var pixels : array<uint8>; w : int; h : int; frame_idx :
[test]
def test_apng_basic(t : T?) {
t |> run("APNG: encode 4 frames, parse chunks, decode frame 0") <| @(tt : T?) {
let fname = "{get_das_root()}/tests/stbimage/_test_output.apng"
let fname = temp_output(".apng")
let W = 32; let H = 32
var writer = stbi_apng_begin(fname, W, H, 4)
tt |> success(writer != null, "stbi_apng_begin succeeded")
Expand Down Expand Up @@ -145,7 +150,7 @@ def test_apng_basic(t : T?) {
[test]
def test_apng_rgb_channels(t : T?) {
t |> run("APNG: 3-channel RGB encode works") <| @(tt : T?) {
let fname = "{get_das_root()}/tests/stbimage/_test_output_rgb.apng"
let fname = temp_output(".apng")
let W = 8; let H = 8
var writer = stbi_apng_begin(fname, W, H, 3)
tt |> success(writer != null, "begin RGB succeeded")
Expand Down Expand Up @@ -181,7 +186,7 @@ def test_apng_begin_bad_filename(t : T?) {
[test]
def test_apng_drop_accounting(t : T?) {
t |> run("APNG: submitted == emitted + dropped (worker race-tolerant invariant)") <| @(tt : T?) {
let fname = "{get_das_root()}/tests/stbimage/_test_output_overflow.apng"
let fname = temp_output(".apng")
// Big-ish frame so deflate takes noticeable time vs. enqueue speed.
let W = 512; let H = 512
var writer = stbi_apng_begin(fname, W, H, 4)
Expand Down Expand Up @@ -225,8 +230,8 @@ def test_apng_drop_accounting(t : T?) {
[test]
def test_apng_double_writers(t : T?) {
t |> run("APNG: two parallel writers — no shared state") <| @(tt : T?) {
let path1 = "{get_das_root()}/tests/stbimage/_test_output_p1.apng"
let path2 = "{get_das_root()}/tests/stbimage/_test_output_p2.apng"
let path1 = temp_output(".apng")
let path2 = temp_output(".apng")
var w1 = stbi_apng_begin(path1, 32, 32, 4)
var w2 = stbi_apng_begin(path2, 64, 64, 4)
tt |> success(w1 != null && w2 != null, "both writers opened")
Expand Down Expand Up @@ -270,7 +275,7 @@ def test_apng_double_writers(t : T?) {
[test]
def test_apng_dropped_zero_when_keeping_up(t : T?) {
t |> run("APNG: drop counter stays at 0 for small, slow stream") <| @(tt : T?) {
let fname = "{get_das_root()}/tests/stbimage/_test_output_slow.apng"
let fname = temp_output(".apng")
let W = 16; let H = 16
var writer = stbi_apng_begin(fname, W, H, 4)
if (writer == null) {
Expand Down
13 changes: 9 additions & 4 deletions tests/stbimage/test_write.das
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,11 @@ require daslib/safe_addr
require daslib/fio
require dastest/testing_boost public

def private temp_output(ext : string) : string {
var error : string
return create_temp_file("stbimage_test", ext, error)
}

def make_rgba_pixels(var pixels : uint8[16]) {
for (i in range(16)) {
pixels[i] = uint8(i * 16)
Expand All @@ -27,7 +32,7 @@ def test_write_png_roundtrip(t : T?) {
t |> run("write PNG to file and read back") <| @(tt : T?) {
var pixels : uint8[16]
make_rgba_pixels(pixels)
let fname = "{get_das_root()}/tests/stbimage/_test_output.png"
let fname = temp_output(".png")
let ok = stbi_write_png(fname, 2, 2, 4, safe_addr(pixels[0]), 8)
tt |> success(ok != 0, "write_png succeeded")
// read back
Expand All @@ -46,7 +51,7 @@ def test_write_bmp_roundtrip(t : T?) {
t |> run("write BMP to file and read back") <| @(tt : T?) {
var pixels : uint8[12]
make_rgb_pixels(pixels)
let fname = "{get_das_root()}/tests/stbimage/_test_output.bmp"
let fname = temp_output(".bmp")
let ok = stbi_write_bmp(fname, 2, 2, 3, safe_addr(pixels[0]))
tt |> success(ok != 0, "write_bmp succeeded")
var x, y, comp : int
Expand All @@ -64,7 +69,7 @@ def test_write_tga_roundtrip(t : T?) {
t |> run("write TGA to file and read back") <| @(tt : T?) {
var pixels : uint8[16]
make_rgba_pixels(pixels)
let fname = "{get_das_root()}/tests/stbimage/_test_output.tga"
let fname = temp_output(".tga")
let ok = stbi_write_tga(fname, 2, 2, 4, safe_addr(pixels[0]))
tt |> success(ok != 0, "write_tga succeeded")
var x, y, comp : int
Expand All @@ -82,7 +87,7 @@ def test_write_jpg_roundtrip(t : T?) {
t |> run("write JPG to file and read back") <| @(tt : T?) {
var pixels : uint8[12]
make_rgb_pixels(pixels)
let fname = "{get_das_root()}/tests/stbimage/_test_output.jpg"
let fname = temp_output(".jpg")
let ok = stbi_write_jpg(fname, 2, 2, 3, safe_addr(pixels[0]), 90)
tt |> success(ok != 0, "write_jpg succeeded")
var x, y, comp : int
Expand Down
Loading