// snaprun.par: several programs run at once, each answered with its exit status. // // The arguments arrive as one string of fields with a NUL after each but the // last (`par.plan` in src/par.bend): how many may run at once, how many gigabytes // each was given, the directory they write into, the deadline, and then each // job as its arguments, each with a `+` in front, and an empty field that ends // it. A marked argument is never empty, so the empty field cannot be mistaken // for one whatever the arguments hold, and NUL is the one char no argument // can hold. The empty job, an empty field alone, is a job `par` refused: it // answers 127 and runs nothing. No shell sees any of it. // // Every job's file is this run's: a job that is not run, because it was // refused or the deadline had passed, still has its file emptied, so nothing // a run before this one left there is read back as its output. // // A child's output goes to `/` rather than down a pipe. One pipe a // child would mean polling every one of them at once or deadlocking on // whichever filled its buffer first, and a file is what the caller reads back // anyway. // // The deadline is an absolute second, or 0 for none, and it is enforced here // because here is where the queueing happens. A job bounded by what was left // when the batch started is not bounded at all: with a width of four and // thirty jobs, the last one starts long after that and still gets the whole // window. So a job is skipped outright once the deadline has gone by, and one // that does start carries an alarm set to what is actually left. // // A width of 0 means the caller has no opinion and this works one out. Cores // are the obvious ceiling and the wrong one on their own: every one of these // runs under a memory cap, so what the machine can hold divides the answer as // surely as what it can compute. What it can hold is the memory going spare // rather than the memory it has — the cap bounds what one job may take, so a // width read off the total would happily start seven jobs entitled to more // than the machine has left and leave the kernel to sort it out. // // Spare means MemAvailable, which is the kernel's own estimate of what a new // job could get, page cache it would reclaim included. `sysinfo` has no field // for it: its freeram plus bufferram omit Cached, so the width would be too // small. sysinfo is still the fallback for a system with no /proc. #include #include #include #include #include #include #include #include #include // the gigabytes a new job could have, by the kernel's own reckoning, or 0 when // this system does not say static double snaprun_par_spare(void) { FILE* f = fopen("/proc/meminfo", "r"); if (f) { char name[64]; unsigned long kb = 0; while (fscanf(f, "%63s %lu kB\n", name, &kb) == 2) { if (strcmp(name, "MemAvailable:") == 0) { fclose(f); return (double)kb / (1024.0 * 1024.0); } } fclose(f); } struct sysinfo si; if (sysinfo(&si) == 0) { double unit = si.mem_unit ? (double)si.mem_unit : 1.0; return (((double)si.freeram + (double)si.bufferram) * unit) / (1024.0 * 1024.0 * 1024.0); } return 0.0; } // how many of these may run at once, when the caller did not say static int snaprun_par_width(int cap_gb) { long cores = sysconf(_SC_NPROCESSORS_ONLN); int n = cores > 0 ? (int)cores : 1; double gb = snaprun_par_spare(); if (cap_gb > 0 && gb > 0.0) { int fits = (int)(gb / (double)cap_gb); if (fits < n) { n = fits; } } return n < 1 ? 1 : n; } // the job a finished child belonged to, so its status lands in the right place static int snaprun_par_slot(pid_t* pids, int n, pid_t pid) { for (int i = 0; i < n; i++) { if (pids[i] == pid) { return i; } } return -1; } Term snaprun_par_run(Env e, Term* f, IoWork* w) { uint64_t n = 0; char* cmd = io_cstr(e, f[0], &n); // the NULs between fields are already terminators, so each field is its // own C string; count them size_t lines = 1; for (size_t i = 0; i < (size_t)n; i++) { if (cmd[i] == '\0') { lines++; } } char** line = malloc(lines * sizeof(char*)); size_t at = 0; line[at++] = cmd; // a field starts after every NUL, the last one included: the empty field // that ends the last job is the empty string at the end, which io_cstr // terminates for (size_t i = 0; i < (size_t)n; i++) { if (cmd[i] == '\0') { line[at++] = cmd + i + 1; } } int want = lines > 0 ? atoi(line[0]) : 0; int cap_gb = lines > 1 ? atoi(line[1]) : 0; const char* dir = lines > 2 ? line[2] : "."; long by = lines > 3 ? atol(line[3]) : 0; int width = want > 0 ? want : snaprun_par_width(cap_gb); // the jobs, each a vector into the fields already split: the marked fields // up to the empty one that ends the job, each with its mark stepped over size_t jobs = 0; for (size_t i = 4; i < at; i++) { if (line[i][0] == '\0') { jobs++; } } char*** argvs = malloc((jobs ? jobs : 1) * sizeof(char**)); size_t j = 0; size_t from = 4; for (size_t i = 4; i < at && j < jobs; i++) { if (line[i][0] == '\0') { size_t argc = i - from; char** argv = malloc((argc + 1) * sizeof(char*)); for (size_t k = 0; k < argc; k++) { argv[k] = line[from + k] + 1; } argv[argc] = NULL; argvs[j++] = argv; from = i + 1; } } pid_t* pids = malloc((jobs ? jobs : 1) * sizeof(pid_t)); int* codes = malloc((jobs ? jobs : 1) * sizeof(int)); for (size_t i = 0; i < jobs; i++) { pids[i] = -1; codes[i] = 127; } size_t next = 0; int live = 0; while (next < jobs || live > 0) { while (next < jobs && live < width) { // what is left of the run when this job starts, which is not what was // left when the one before it did long spare = by > 0 ? by - (long)time(NULL) : 0; char path[4096]; snprintf(path, sizeof(path), "%s/%zu", dir, next); // a refused job, and one the deadline has passed, run nothing, and // their files are emptied so no earlier run's output is read as theirs int refused = argvs[next][0] == NULL; if (refused || (by > 0 && spare <= 0)) { int gone = open(path, O_WRONLY | O_CREAT | O_TRUNC, 0644); if (gone >= 0) { close(gone); } codes[next] = refused ? 127 : 124; next++; continue; } pid_t pid = fork(); if (pid == 0) { // a child that cannot have its own file would otherwise inherit this // program's stdout and print a job's output into the run's verdict int out = open(path, O_WRONLY | O_CREAT | O_TRUNC, 0644); if (out < 0) { _exit(127); } int nul = open("/dev/null", O_RDONLY); dup2(nul, 0); dup2(out, 1); dup2(out, 2); // the alarm outlives the exec, and SIGALRM ends a process by default, // so this is the deadline the job cannot run past if (by > 0) { alarm((unsigned int)spare); } execvp(argvs[next][0], argvs[next]); _exit(127); } pids[next] = pid; if (pid > 0) { live++; } next++; } if (live > 0) { int status = 0; pid_t done = wait(&status); if (done < 0) { // a signal interrupting the wait is not the jobs finishing if (errno == EINTR) { continue; } break; } // only a child of this batch's counts against the width. `wait` reaps // whatever finished, and giving up a slot to someone else's child would // let this return while jobs of its own were still writing their files. int slot = snaprun_par_slot(pids, (int)jobs, done); if (slot >= 0) { // the alarm is the deadline, so a job it ended reads as a job that // timed out rather than as one killed by a signal nobody sent codes[slot] = WIFEXITED(status) ? WEXITSTATUS(status) : WTERMSIG(status) == SIGALRM ? 124 : 128 + WTERMSIG(status); // the slot stops answering to this pid, so that a pid the kernel // handed out again could not overwrite a job already scored pids[slot] = -1; live--; } } } // the statuses in the order the jobs were given, one to a line size_t cap = jobs * 8 + 2; char* out = malloc(cap); size_t len = 0; for (size_t i = 0; i < jobs; i++) { len += (size_t)snprintf(out + len, cap - len, i + 1 < jobs ? "%d\n" : "%d", codes[i]); } for (size_t i = 0; i < jobs; i++) { free(argvs[i]); } free(codes); free(pids); free(argvs); free(line); free(cmd); Term s = io_str(e, out, len); free(out); return s; } static void __attribute__((constructor)) snaprun_par_use(void) { io_eff(CID_SNAPRUN_PAR, snaprun_par_run, 0); }