diff --git a/DESCRIPTION b/DESCRIPTION index 816d486..4ca613c 100644 --- a/DESCRIPTION +++ b/DESCRIPTION @@ -1,6 +1,6 @@ Package: dparser Title: Port of 'Dparser' Package -Version: 1.3.1-13 +Version: 1.3.2 Authors@R: c(person("Matthew", "Fidler", email = "matthew.fidler@gmail.com", role = c("aut", "cre")), person("John", "Plevyak", role = c("aut", "cph"), email = "jplevyak@gmail.com")) Imports: diff --git a/NEWS.md b/NEWS.md index 6aae696..db2ba91 100644 --- a/NEWS.md +++ b/NEWS.md @@ -1,5 +1,14 @@ # dparser 1.3.2 +- `dparse()` can now be called from several threads at once, one parser + per thread. The first reduction path was a single process-wide static + vector, so concurrent parses shared and reallocated it, segfaulting + rxode2's multi-threaded `rxOptExpr()` (nlmixr2/rxode2#1427); it now lives + on the reducing call's stack. Callers must resolve the `dparser.h` entry + points on the main thread first, and syntax errors, ambiguity errors and + R-level callbacks still call the R API, so only error-free parses without + R callbacks are safe off the main thread. + - `buf_read()` (and therefore `sbuf_read()`) is hardened in three ways without changing its `int *len` ABI: * Files larger than `INT_MAX - 2` bytes are now rejected with `-1` diff --git a/src/parse.c b/src/parse.c index be16250..c637a0b 100644 --- a/src/parse.c +++ b/src/parse.c @@ -1319,14 +1319,14 @@ static void shift_all(Parser *p, char *pos) { } } -static VecZNode path1; /* static first path for speed */ - -static VecZNode *new_VecZNode(VecVecZNode *paths, int n, int parent) { +/* The first path is the caller's (reduce_one's) stack vector, not a static, so + concurrent dparse() calls, one parser per thread, share nothing here. */ +static VecZNode *new_VecZNode(VecVecZNode *paths, int n, int parent, VecZNode *path1) { int i; VecZNode *pv; if (!paths->n) - pv = &path1; + pv = path1; else pv = MALLOC(sizeof *pv); vec_clear(pv); @@ -1335,7 +1335,7 @@ static VecZNode *new_VecZNode(VecVecZNode *paths, int n, int parent) { return pv; } -static void build_paths_internal(ZNode *z, VecVecZNode *paths, int parent, int n, int n_to_go) { +static void build_paths_internal(ZNode *z, VecVecZNode *paths, int parent, int n, int n_to_go, VecZNode *path1) { uint j, k, l; vec_add(paths->v[parent], z); @@ -1344,26 +1344,26 @@ static void build_paths_internal(ZNode *z, VecVecZNode *paths, int parent, int n for (j = 0, l = 0; j < z->sns.v[k]->zns.n; j++) { if (z->sns.v[k]->zns.v[j]) { if (k + l) { - vec_add(paths, new_VecZNode(paths, n - (n_to_go - 1), parent)); + vec_add(paths, new_VecZNode(paths, n - (n_to_go - 1), parent, path1)); parent = paths->n - 1; } - build_paths_internal(z->sns.v[k]->zns.v[j], paths, parent, n, n_to_go - 1); + build_paths_internal(z->sns.v[k]->zns.v[j], paths, parent, n, n_to_go - 1, path1); l++; } } } -static void build_paths(ZNode *z, VecVecZNode *paths, int nchildren_to_go) { +static void build_paths(ZNode *z, VecVecZNode *paths, int nchildren_to_go, VecZNode *path1) { if (!nchildren_to_go) return; - vec_add(paths, new_VecZNode(paths, 0, -1)); - build_paths_internal(z, paths, 0, nchildren_to_go, nchildren_to_go); + vec_add(paths, new_VecZNode(paths, 0, -1, path1)); + build_paths_internal(z, paths, 0, nchildren_to_go, nchildren_to_go, path1); } -static void free_paths(VecVecZNode *paths) { +static void free_paths(VecVecZNode *paths, VecZNode *path1) { uint i; for (i = 0; i < paths->n; i++) { vec_free(paths->v[i]); - if (paths->v[i] != &path1) FREE(paths->v[i]); + if (paths->v[i] != path1) FREE(paths->v[i]); } vec_free(paths); } @@ -1374,6 +1374,7 @@ static void reduce_one(Parser *p, Reduction *r) { ZNode *first_z; uint i, j, n = r->reduction->nelements; VecVecZNode paths; + VecZNode path1; /* first path, see new_VecZNode() */ VecZNode *path; if (!r->znode) { /* epsilon reduction */ @@ -1382,7 +1383,7 @@ static void reduce_one(Parser *p, Reduction *r) { } else { DBG(Rprintf("reduce %d %p %d\n", (int)(r->snode->state - p->t->state), (void *)sn, n)); vec_clear(&paths); - build_paths(r->znode, &paths, n); + build_paths(r->znode, &paths, n, &path1); for (i = 0; i < paths.n; i++) { path = paths.v[i]; if (r->new_snode) { /* prune paths by new right epsilon node */ @@ -1399,7 +1400,7 @@ static void reduce_one(Parser *p, Reduction *r) { if (pn) for (j = 0; j < first_z->sns.n; j++) goto_PNode(p, &sn->loc, pn, first_z->sns.v[j]); } - free_paths(&paths); + free_paths(&paths, &path1); } unref_sn(p, sn); r->next = p->free_reductions; diff --git a/tests/testthat/concurrent_parser.c b/tests/testthat/concurrent_parser.c new file mode 100644 index 0000000..de80523 --- /dev/null +++ b/tests/testthat/concurrent_parser.c @@ -0,0 +1,63 @@ +/* + Runs dparse() from several threads at once, one parser per thread, the way + rxode2's rxOptExpr() does (nlmixr2/rxode2#1427). +*/ +#include +#include +#include +#include +#include + +#define MAX_THREADS 64 + +extern D_ParserTables parser_tables_gram; + +typedef struct { const char *txt; int reps; int bad; } cpJob; + +static int cpParseOnce(const char *txt) { + D_Parser *p = new_D_Parser(&parser_tables_gram, 100); + D_ParseNode *pn; + int ok; + if (p == NULL) return 0; + p->save_parse_tree = 1; + p->error_recovery = 0; + pn = dparse(p, (char *) txt, (int) strlen(txt)); + ok = pn != NULL && p->syntax_errors == 0; + if (pn != NULL) free_D_ParseNode(p, pn); + free_D_Parser(p); + return ok; +} + +static void *cpWorker(void *arg) { + cpJob *j = (cpJob *) arg; + int i; + for (i = 0; i < j->reps; i++) { + if (!cpParseOnce(j->txt)) j->bad++; + } + return NULL; +} + +SEXP concurrent_dparse(SEXP txtS, SEXP nthrS, SEXP repsS) { + const char *txt = CHAR(STRING_ELT(txtS, 0)); + int nthr = INTEGER(nthrS)[0], reps = INTEGER(repsS)[0], t, bad = 0; + pthread_t th[MAX_THREADS]; + cpJob jobs[MAX_THREADS]; + if (nthr < 1 || nthr > MAX_THREADS) Rf_error("nthr must be in 1..%d", MAX_THREADS); + /* serial parse first: also resolves dparser.h's R_GetCCallable() pointers on + the main thread, which worker threads must not do */ + if (!cpParseOnce(txt)) Rf_error("serial parse failed"); + for (t = 0; t < nthr; t++) { + jobs[t].txt = txt; jobs[t].reps = reps; jobs[t].bad = 0; + if (pthread_create(&th[t], NULL, cpWorker, &jobs[t]) != 0) { + int u; + /* jobs[] lives on this stack, which Rf_error() unwinds */ + for (u = 0; u < t; u++) pthread_join(th[u], NULL); + Rf_error("pthread_create failed"); + } + } + for (t = 0; t < nthr; t++) { + pthread_join(th[t], NULL); + bad += jobs[t].bad; + } + return Rf_ScalarInteger(bad); +} diff --git a/tests/testthat/test-concurrent-dparse.R b/tests/testthat/test-concurrent-dparse.R new file mode 100644 index 0000000..43d0b29 --- /dev/null +++ b/tests/testthat/test-concurrent-dparse.R @@ -0,0 +1,41 @@ +test_that("dparse() is safe to run from several threads at once (rxode2#1427)", { + skip_on_cran() + skip_on_os("windows") + wd <- tempfile("dparser-concurrent") + dir.create(wd) + on.exit(unlink(wd, recursive = TRUE), add = TRUE) + file.copy(test_path("concurrent_parser.c"), wd) + # long right-hand sides make each reduction grow the shared first path + writeLines(c( + "statements: statement*;", + "statement: identifier '=' expr ';';", + "expr: expr '+' expr $left 1", + " | expr '*' expr $left 2", + " | identifier '(' expr ',' expr ',' expr ')'", + " | '(' expr ')'", + " | identifier", + " | number;", + "identifier: \"[a-zA-Z_][a-zA-Z0-9_]*\";", + "number: \"[0-9]+\";" + ), file.path(wd, "concurrent.g")) + mkdparse(file.path(wd, "concurrent.g"), wd, use_r_header = TRUE, verbose = FALSE) + writeLines(c( + sprintf("PKG_CPPFLAGS=-I\"%s\"", dpIncludeDir()), + "PKG_CFLAGS=-pthread", + "PKG_LIBS=-pthread" + ), file.path(wd, "Makevars")) + so <- file.path(wd, paste0("concurrent_parser", .Platform$dynlib.ext)) + owd <- setwd(wd) + out <- system2(file.path(R.home("bin"), "R"), + c("CMD", "SHLIB", "-o", basename(so), "concurrent_parser.c", + "concurrent.g.d_parser.c"), + stdout = TRUE, stderr = TRUE) + setwd(owd) + if (!file.exists(so)) stop(paste(c("could not build the test parser:", out), collapse = "\n")) + dll <- dyn.load(so) + on.exit(dyn.unload(so), add = TRUE) + txt <- paste(sprintf("a%d = f(b + c * %d, g(x, y, z), (d + e) * h);", 1:200, 1:200), + collapse = "\n") + expect_equal(.Call(dll$concurrent_dparse, txt, 1L, 5L), 0L) + expect_equal(.Call(dll$concurrent_dparse, txt, 4L, 50L), 0L) +})