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
2 changes: 1 addition & 1 deletion DESCRIPTION
Original file line number Diff line number Diff line change
@@ -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:
Expand Down
9 changes: 9 additions & 0 deletions NEWS.md
Original file line number Diff line number Diff line change
@@ -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`
Expand Down
29 changes: 15 additions & 14 deletions src/parse.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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);
Expand All @@ -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);
}
Expand All @@ -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 */
Expand All @@ -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 */
Expand All @@ -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;
Expand Down
63 changes: 63 additions & 0 deletions tests/testthat/concurrent_parser.c
Original file line number Diff line number Diff line change
@@ -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 <R.h>
#include <Rinternals.h>
#include <pthread.h>
#include <string.h>
#include <dparser.h>

#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);
}
41 changes: 41 additions & 0 deletions tests/testthat/test-concurrent-dparse.R
Original file line number Diff line number Diff line change
@@ -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)
})
Loading