Skip to content

Commit e8452fc

Browse files
committed
with test
1 parent 0e8a760 commit e8452fc

4 files changed

Lines changed: 45 additions & 14 deletions

File tree

r/adbcdrivermanager/R/async.R

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,25 @@ adbc_async_task_status <- function(task) {
2727
.Call(RAdbcAsyncTaskWaitFor, task, 0)
2828
}
2929

30+
adbc_async_task_set_callback <- function(task, callback, loop = later::current_loop()) {
31+
# If the task is completed, run the callback (or else the callback
32+
# will not run)
33+
if (adbc_async_task_status(task) == "ready") {
34+
result <- adbc_async_task_result(task)
35+
callback(result)
36+
} else {
37+
.Call(RAdbcAsyncTaskSetCallback, task, callback, loop$id)
38+
}
39+
40+
invisible(task)
41+
}
42+
43+
adbc_async_task_run_callback <- function(task) {
44+
callback <- task$callback
45+
result <- adbc_async_task_result(task)
46+
callback(result)
47+
}
48+
3049
adbc_async_task_wait_non_cancellable <- function(task, resolution = 0.05) {
3150
.Call(RAdbcAsyncTaskWaitFor, task, round(resolution * 1000))
3251
}

r/adbcdrivermanager/src/async.cc

Lines changed: 6 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -45,30 +45,23 @@ static void later_task_callback_wrapper(void* data);
4545
enum class RAdbcAsyncTaskStatus { NOT_STARTED, STARTED, READY };
4646

4747
struct RAdbcAsyncTask {
48-
RAdbcAsyncTask() : callback_sexp(R_NilValue), callback_data_sexp(R_NilValue) {}
48+
RAdbcAsyncTask() : callback_data_sexp(R_NilValue) {}
4949

50-
void SetCallback(SEXP callback, SEXP data, int loop_id) {
51-
if (callback_sexp != R_NilValue) {
52-
return;
53-
}
54-
55-
callback_sexp = callback;
50+
void SetCallback(SEXP data, int loop_id) {
5651
callback_data_sexp = data;
5752
later_loop_id = loop_id;
5853
later_ensure_initialized();
5954
}
6055

6156
void ScheduleCallbackIfSet() {
62-
if (callback_sexp != R_NilValue) {
57+
if (callback_data_sexp != R_NilValue) {
6358
later_execLaterNative2(&later_task_callback_wrapper, this, 0, later_loop_id);
64-
callback_sexp = R_NilValue;
6559
}
6660
}
6761

6862
AdbcError* return_error{nullptr};
6963
int* return_code{nullptr};
7064

71-
SEXP callback_sexp;
7265
SEXP callback_data_sexp;
7366
int later_loop_id{-1};
7467

@@ -79,9 +72,8 @@ struct RAdbcAsyncTask {
7972
static void later_task_callback_wrapper(void* data) {
8073
auto task = reinterpret_cast<RAdbcAsyncTask*>(data);
8174

82-
SEXP func_sym = PROTECT(Rf_install("adbc_async_run_callback"));
83-
SEXP func_call =
84-
PROTECT(Rf_lang3(func_sym, task->callback_sexp, task->callback_data_sexp));
75+
SEXP func_sym = PROTECT(Rf_install("adbc_async_task_run_callback"));
76+
SEXP func_call = PROTECT(Rf_lang2(func_sym, task->callback_data_sexp));
8577
SEXP pkg_chr = PROTECT(Rf_mkString("adbcdrivermanager"));
8678
SEXP pkg_ns = PROTECT(R_FindNamespace(pkg_chr));
8779
Rf_eval(func_call, pkg_ns);
@@ -140,7 +132,7 @@ extern "C" SEXP RAdbcAsyncTaskSetCallback(SEXP task_xptr, SEXP callback_sexp,
140132
int loop_id = adbc_as_int(loop_id_sexp);
141133

142134
SET_VECTOR_ELT(task_prot, 3, callback_sexp);
143-
task->SetCallback(callback_sexp, task_xptr, loop_id);
135+
task->SetCallback(task_xptr, loop_id);
144136
return R_NilValue;
145137
}
146138

r/adbcdrivermanager/src/init.c

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
/* generated by tools/make-callentries.R */
2323
SEXP RAdbcAsyncTaskNew(SEXP error_xptr);
24+
SEXP RAdbcAsyncTaskSetCallback(SEXP task_xptr, SEXP callback_sexp, SEXP loop_id_sexp);
2425
SEXP RAdbcAsyncTaskData(SEXP task_xptr);
2526
SEXP RAdbcAsyncTaskWaitFor(SEXP task_xptr, SEXP duration_ms_sexp);
2627
SEXP RAdbcAsyncTaskLaunchSleep(SEXP task_xptr, SEXP duration_ms_sexp);
@@ -110,6 +111,7 @@ SEXP RAdbcXptrSetProtected(SEXP xptr, SEXP prot);
110111

111112
static const R_CallMethodDef CallEntries[] = {
112113
{"RAdbcAsyncTaskNew", (DL_FUNC)&RAdbcAsyncTaskNew, 1},
114+
{"RAdbcAsyncTaskSetCallback", (DL_FUNC)&RAdbcAsyncTaskSetCallback, 3},
113115
{"RAdbcAsyncTaskData", (DL_FUNC)&RAdbcAsyncTaskData, 1},
114116
{"RAdbcAsyncTaskWaitFor", (DL_FUNC)&RAdbcAsyncTaskWaitFor, 2},
115117
{"RAdbcAsyncTaskLaunchSleep", (DL_FUNC)&RAdbcAsyncTaskLaunchSleep, 2},

r/adbcdrivermanager/tests/testthat/test-async.R

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,24 @@ test_that("async task waiter works", {
7474
)
7575
})
7676

77+
test_that("async tasks can set an R callback", {
78+
skip_if_not_installed("later")
79+
80+
async_called <- FALSE
81+
sleep_task <- adbc_async_sleep(200)
82+
adbc_async_task_set_callback(sleep_task, function(x) { async_called <<- TRUE })
83+
Sys.sleep(0.4)
84+
later::run_now()
85+
expect_true(async_called)
86+
87+
# Ensure the callback runs even if the task is already finished
88+
async_called <- FALSE
89+
sleep_task <- adbc_async_sleep(0)
90+
adbc_async_task_set_callback(sleep_task, function(x) { async_called <<- TRUE })
91+
Sys.sleep(0.1)
92+
expect_true(async_called)
93+
})
94+
7795
test_that("async task can be converted to a promise", {
7896
skip_if_not_installed("promises")
7997

0 commit comments

Comments
 (0)