Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / power / CMakeLists.txt
Author[]Andy Green <andy@warmcat.com> 2026-09-30 21:18 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-10-04 04:09 UTC
Tree3d8fdd3bba978ef4cad490fdaa0b56aaf7fca8a8   Raw Patch
 
pool: builders keep a repo's named pools synced through sai-server
pool: builders keep a repo's named pools synced through sai-server

A .sai.json configuration can now name a "pool": a set of files shared
by the repo's tasks, meant for state that builds up over many tasks and
builders, like fuzzing corpora.

The builder keeps a copy of each pool its tasks use under
$HOME/pools/, and tells build steps where it is in SAI_POOL_DIR,
SAI_POOL_KNOWN and SAI_POOL_FINDINGS.  It syncs the pool before a task's
first build step starts (unless it did within the last minute), every
minute while the task runs, and once more after it ends, so a task that
is stopped, eg, to make way for real work, loses nothing it wrote.
Tasks using the same pool on one builder share one copy.

 - SAI_POOL_DIR holds <sub>/<sha1> files, content addressed as
   libFuzzer names its corpus files, synced both ways.  Nothing is
   deleted by syncing, except when a task replaces a sub, eg, after
   minimizing it: it leaves a marker with the sequence number it
   started from, and the server removes what it had up to there that
   is no longer listed

 - SAI_POOL_KNOWN is a copy of what the server provides, eg, known
   reproducers

 - anything in SAI_POOL_FINDINGS is sent to the server, then deleted

Each sync is its own connection to /builder: the link key, then a JSON
hello naming a task and its artifact upload nonce, which decide the repo
and pool; then binary records both ways.  sai-server keeps each repo's
pool in its own sqlite db as an append-only log of entries, so builders
pull only what happened after the last sequence number they saw, and
offer what's new locally for the server to ask for.  Names are checked
before either side uses them as paths, content addressed files are
checked against their names on both ends, and sizes are capped.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
diff --git a/READMEs/README-idle.md b/READMEs/README-idle.md index 4e2fbc6..4dad018 100644 --- a/READMEs/README-idle.md +++ b/READMEs/README-idle.md @@ -102,10 +102,10 @@ task. The whole step, including any build, has to fit in the slice, so a script that divides its time between several things should divide `SAI_IDLE_SECS` between them. -Anything the task wants to keep between slices, like a fuzzing corpus, has to be -kept somewhere under `$HOME` on the builder by the task itself for now; the job -dir is reused by the lane's next slice on the same builder, but not preserved -across builders. +Anything the task wants to keep between slices, like a fuzzing corpus, should go +in a pool (see [README-pool.md](README-pool.md)): the builder keeps it synced +with every other builder working on the repo, even when it has to stop a slice +to make way for real work. ## In the web UI diff --git a/READMEs/README-pool.md b/READMEs/README-pool.md new file mode 100644 index 0000000..6b85f3e --- /dev/null +++ b/READMEs/README-pool.md @@ -0,0 +1,104 @@ +# Pools + +A pool is a named set of files belonging to a repo, that the builders running +its tasks keep synced through sai-server. It's meant for state that work +builds up over many tasks and many builders, like fuzzing corpora: every +builder working on the repo's fuzzing adds to one shared corpus, and starts +from what all the others found. + +## Asking for a pool in .sai.json + +A configuration names the pool its tasks use: + +``` + "fuzz": { + "cmake": "./fuzz/run.sh", + "pool": "fuzz", + "idle": 2 + } +``` + +The name is up to 32 lowercase letters, digits, `-` and `_`. Configurations +naming the same pool in the same repo share it. + +## What the task sees + +The builder keeps its copy of the pool under `$HOME/pools/`, and tells the +task's build steps where it is: + +|variable|what it is| +|---|---| +|`SAI_POOL_DIR`|the shared files, synced both ways| +|`SAI_POOL_KNOWN`|files the server provides, eg, known reproducers; only ever synced from the server| +|`SAI_POOL_FINDINGS`|anything the task leaves in here is sent to the server, then deleted| + +The builder syncs the pool before the task's first build step starts (unless it +did within the last minute), every minute while the task runs, and once more +after it ends. The builder does it rather than the task, so a task that's +stopped, eg, because real work needed the builder, loses nothing it wrote. +Several tasks on the builder using the same pool share one copy. + +### Files in SAI_POOL_DIR + +Only files at `<sub>/<name>` are synced, where + + - `<sub>` is a directory name of up to 32 letters, digits, `-`, `_` and `.`, + not starting with `.`, eg, `corpus-h2` + + - `<name>` is the 40 character lowercase hex SHA-1 of the file's content + +That's how libFuzzer names its corpus files already, so a libFuzzer corpus +dir per target, eg, `$SAI_POOL_DIR/corpus-h2/`, just works. Anything else in +there is left alone and stays local. Files are only ever added by syncing, +never deleted, except by a replace (below). + +Files arriving from the server are written somewhere else first and renamed +into place, so a fuzzer reading the directory never sees one half written. + +### Findings + +Files at `$SAI_POOL_FINDINGS/<sub>/<name>` are sent to the server and deleted +once it has them. `<name>` can be anything up to 64 letters, digits, `-`, `_` +and `.`, not starting with `.`, and the file can be up to 8MiB. Write a +finding under a name starting with `.` and rename it when it's complete, so it +isn't sent half written. + +The server keeps findings in the pool's db; what happens to them after that is +up to later work (see the idle fuzzing design). + +### Replacing a sub, eg, after minimizing a corpus + +Tasks can't delete shared files, since deleting one locally doesn't mean the +others should lose it. To replace everything in a sub instead, eg, with a +corpus libFuzzer's `-merge=1` minimized: + +1. read the number in `$SAI_POOL_DIR/.sai-pool-seq` before you start: it's how + far the builder had synced from the server +2. make `$SAI_POOL_DIR/<sub>/` contain just the files you want to keep +3. write the number from step 1 into `$SAI_POOL_DIR/.sai-replace-<sub>` + +At the next sync, the server removes every file in the sub it had up to that +number that isn't in the directory now, and the other builders remove them +when they next sync. Files anyone added after that number are kept, since the +task didn't know about them. + +## Sizes + +A synced file can be up to 1MiB, and a pool holds up to 2 million files. + +## On the server + +Each repo's pool is a sqlite db next to the others, +`<database>-pool-<repo>-<pool>.sqlite3`. Its entries are an append-only log: +each file added or removed gets the next sequence number, so a builder only +asks for what happened after the last one it saw. A removed file leaves a +row saying so. + +## Trust + +Only a builder that has the fleet link key can sync, and only the pool of a +task it was given: it has to present the task's upload nonce. But what goes +in a pool comes from the tasks, which run what's pushed to the repo, so a pool +is only as trustworthy as the people who can push to it. Content addressed +files are checked against their names on both ends, names are checked before +they're used as paths, and sizes are capped. diff --git a/READMEs/README-sai-json.md b/READMEs/README-sai-json.md index 323e9aa..e385e87 100644 --- a/READMEs/README-sai-json.md +++ b/READMEs/README-sai-json.md @@ -108,3 +108,14 @@ Besides the configuration's normal task, each event then has that many idle tasks ("lanes") for it on each platform. They do not count towards the event's result, and only run in time builders would otherwise spend idle, when the builder conf allows it. See [README-idle.md](README-idle.md). + +#### pool + +A configuration can name a pool, a set of files shared by the repo's tasks +that the builders keep synced through sai-server, eg, a fuzzing corpus: + +``` + "pool": "fuzz" +``` + +See [README-pool.md](README-pool.md). diff --git a/src/builder/CMakeLists.txt b/src/builder/CMakeLists.txt index 322c30a..e542037 100644 --- a/src/builder/CMakeLists.txt +++ b/src/builder/CMakeLists.txt @@ -9,6 +9,7 @@ set(SRCS b-nspawn.c b-task.c b-artifacts.c + b-pool.c b-logproxy.c b-refproxy.c b-load.c @@ -17,6 +18,7 @@ set(SRCS b-deletion.c b-power.c ../common/c-utils.c + ../common/c-pool.c ../common/struct-metadata.c ) diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 4210da5..842aa82 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -612,7 +612,7 @@ static const char * const runscript_win_first = "set SAI_LOGPROXY=%s\n" "set SAI_LOGPROXY_TTY0=%s\n" "set SAI_LOGPROXY_TTY1=%s\n" - "%s" + "%s%s" "set HOME=%s\n" "set CI=true\n" "set BUILDKIT_PROGRESS=plain\n" @@ -629,7 +629,7 @@ static const char * const runscript_win_next = "set SAI_LOGPROXY=%s\n" "set SAI_LOGPROXY_TTY0=%s\n" "set SAI_LOGPROXY_TTY1=%s\n" - "%s" + "%s%s" "set HOME=%s\n" "set CI=true\n" "set BUILDKIT_PROGRESS=plain\n" @@ -658,7 +658,7 @@ static const char * const runscript_first = "export SAI_LOGPROXY=%s\n" "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" - "%s" + "%s%s" "export CI=true\n" "export BUILDKIT_PROGRESS=plain\n" "set -e\n" @@ -687,7 +687,7 @@ static const char * const runscript_next = "export SAI_LOGPROXY=%s\n" "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" - "%s" + "%s%s" "export CI=true\n" "export BUILDKIT_PROGRESS=plain\n" "set -e\n" @@ -715,7 +715,7 @@ static const char * const runscript_build = "export SAI_LOGPROXY=%s\n" "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" - "%s" + "%s%s" "export CI=true\n" "export BUILDKIT_PROGRESS=plain\n" "set -e\n" @@ -747,8 +747,8 @@ saib_spawn_script(struct sai_nspawn *ns) NULL }; #endif - char one_step[4096], idle_env[64]; - char st[2048]; + char one_step[4096], idle_env[64], pool_env[1024]; + char st[8192]; unsigned int timeout_secs; int fd, n; #if defined(__linux__) @@ -803,6 +803,9 @@ saib_spawn_script(struct sai_nspawn *ns) * whatever it does in the time. We stop it anyway if it goes on much * longer than that. */ + /* where the task's pool is, if it has one, see b-pool.c */ + saib_pool_env(ns, pool_env, sizeof(pool_env)); + idle_env[0] = '\0'; timeout_secs = builder.build_timeout_secs; if (ns->task->idle) { @@ -826,7 +829,7 @@ saib_spawn_script(struct sai_nspawn *ns) ns->task->parallel ? ns->task->parallel : 1, respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, idle_env, - builder.home, + pool_env, builder.home, ns->inp, ns->task->build_step > 1 ? "\\src" : "", one_step); #else @@ -850,7 +853,7 @@ saib_spawn_script(struct sai_nspawn *ns) ns->task->parallel ? ns->task->parallel : 1, respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, idle_env, - builder.home, one_step); + pool_env, builder.home, one_step); #endif /* but from the script's pov, it's chrooted at /home/sai */ diff --git a/src/builder/b-pool.c b/src/builder/b-pool.c new file mode 100644 index 0000000..5b4e6b0 --- /dev/null +++ b/src/builder/b-pool.c @@ -0,0 +1,1369 @@ +/* + * Sai builder - ./src/builder/b-pool.c + * + * Copyright (C) 2019 - 2026 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * Pools, builder side: see READMEs/README-pool.md + * + * A task whose .sai.json configuration names a pool gets a local copy of it + * under $HOME/pools/, which we keep synced with sai-server: before its first + * build step starts, every minute while it runs, and once more after it ends. + * Doing it here rather than in the task means a task that's killed loses + * nothing it had written, and several tasks on the builder using the same + * pool share one copy. + * + * corpus/ both ways, content addressed <sub>/<sha1> files + * known/ a copy of what the server has, eg, known reproducers + * findings/ anything the task leaves in <sub>/ here is sent to the + * server, then deleted + * + * Each sync is its own connection to the server, a stream of records after + * the hello (see private.h): + * + * - PULL both content addressed namespaces from where we got to last time + * - OFFER the server corpus files that appeared locally since the last sync, + * and PUT the ones it wants + * - PUT everything in findings/ + * - REPLACE any corpus sub a task asked for, see README-pool.md + */ + +#include <libwebsockets.h> + +#include "b-private.h" + +#include <sys/types.h> +#include <sys/stat.h> +#include <stdlib.h> +#include <fcntl.h> +#include <errno.h> +#if !defined(WIN32) +#include <unistd.h> +#endif + +/* how often we sync a pool while tasks are using it */ +#define SAIB_POOL_SYNC_INTERVAL_US (60 * LWS_US_PER_SEC) +/* a task's first build step waits for a pull unless there was one this recently */ +#define SAIB_POOL_PULL_FRESH_US (60 * LWS_US_PER_SEC) +/* the longest a step waits for its pool to sync before starting anyway */ +#define SAIB_POOL_WAIT_MAX_US (180 * LWS_US_PER_SEC) +/* the longest one sync can take before we give up on it */ +#define SAIB_POOL_SESSION_MAX_US (15 * 60 * LWS_US_PER_SEC) +/* how often we offer the server every corpus file, not just new ones */ +#define SAIB_POOL_FULL_OFFER_SECS (24 * 3600) +/* files this much older than the last sync are offered again anyway */ +#define SAIB_POOL_OFFER_SLACK_SECS 60 +/* how many times the sync after the last task ends is tried */ +#define SAIB_POOL_FINAL_TRIES 3 + +typedef struct saib_pool { + lws_dll2_t list; /* builder.pool_owner */ + lws_dll2_owner_t waiters; /* nspawns waiting for a pull */ + + lws_sorted_usec_list_t sul_sync; /* next sync */ + lws_sorted_usec_list_t sul_wait; /* waiters stop waiting */ + lws_sorted_usec_list_t sul_session; /* sync taking too long */ + + struct sai_plat_server *spm; + struct lws_ss_handle *ss; /* the sync going on, if any */ + + char name[33]; + char key[160]; /* server-repo-pool, purified */ + char dir[300]; + + /* the task we sync on behalf of: the latest one to use the pool */ + char task_uuid[65]; + char nonce[33]; + + uint64_t cursor[2]; /* corpus, known */ + uint64_t offer_since; /* wall clock, see offer scan */ + uint64_t full_offer; /* wall clock of last full offer */ + lws_usec_t last_pull; + + int users; /* nspawns using it */ + int final_tries; + char resync; /* sync again when this one ends */ +} saib_pool_t; + +/* one record queued to go, or waiting to be acknowledged */ + +typedef struct saib_pool_rec { + lws_dll2_t list; + char path[384]; /* unlink on ACK, if any */ + char name[SAI_POOL_REC_NAME_MAX + 1]; + uint64_t replace_base; + size_t len; + size_t ofs; + char is_replace; + /* len bytes of record follow */ +} saib_pool_rec_t; + +/* a corpus file the server said it wants */ + +typedef struct saib_pool_put { + lws_dll2_t list; + char name[SAI_POOL_REC_NAME_MAX + 1]; +} saib_pool_put_t; + +typedef struct saib_pool_ss { + struct lws_ss_handle *ss; + void *opaque_data; + + saib_pool_t *pool; + + lws_dll2_owner_t txq; /* saib_pool_rec_t to send */ + lws_dll2_owner_t acks; /* saib_pool_rec_t awaiting ACK */ + lws_dll2_owner_t puts; /* saib_pool_put_t to load */ + struct lwsac *ac_puts; + + uint8_t *rxb; + size_t rxb_len; + size_t rxb_alloc; + + uint64_t scan_start; /* wall clock at the offer scan */ + + int pulls_left; + int offers_left; + + char sent_auth; + char sent_hello; + char pulled; + char full_offer; + char done; +} saib_pool_ss_t; + +static void +saib_pool_sync_start(saib_pool_t *pool); + +static const char * const ns_dir[] = { "corpus", "known", "findings" }; + +/* + * Local state + */ + +static const char * const state_paths[] = { + "corpus", "known", "offer_since", "full_offer", +}; + +static signed char +saib_pool_state_cb(struct lejp_ctx *ctx, char reason) +{ + saib_pool_t *pool = (saib_pool_t *)ctx->user; + uint64_t v; + + if (reason != LEJPCB_VAL_NUM_INT || !ctx->path_match) + return 0; + + v = (uint64_t)strtoull(ctx->buf, NULL, 10); + + switch (ctx->path_match - 1) { + case 0: + pool->cursor[SAI_POOL_NS_CORPUS] = v; + break; + case 1: + pool->cursor[SAI_POOL_NS_KNOWN] = v; + break; + case 2: + pool->offer_since = v; + break; + case 3: + pool->full_offer = v; + break; + } + + return 0; +} + +static void +saib_pool_state_load(saib_pool_t *pool) +{ + struct lejp_ctx ctx; + char path[384]; + uint8_t buf[256]; + int fd, n, m = 0; + + lws_snprintf(path, sizeof(path), "%s/.sai-pool-state", pool->dir); + fd = lws_open(path, O_RDONLY); + if (fd < 0) + return; + + lejp_construct(&ctx, saib_pool_state_cb, pool, state_paths, + LWS_ARRAY_SIZE(state_paths)); + do { + n = (int)read(fd, buf, sizeof(buf)); + if (n <= 0) + break; + m = lejp_parse(&ctx, buf, n); + } while (m == LEJP_CONTINUE); + lejp_destruct(&ctx); + close(fd); + + if (m < 0 && m != LEJP_CONTINUE) { + lwsl_warn("%s: %s unreadable, starting over\n", __func__, path); + pool->cursor[0] = pool->cursor[1] = 0; + pool->offer_since = pool->full_offer = 0; + } +} + +/* + * Write a file so it's never seen half-written: via tmp, which must be in a + * place nothing looks at, eg, not a corpus dir a fuzzer is reading + */ + +static int +saib_pool_write_file(const char *tmp, const char *path, const void *buf, + size_t len) +{ + int fd; + + fd = open(tmp, O_CREAT | O_TRUNC | O_WRONLY +#if defined(WIN32) + | _O_BINARY +#endif + , 0644); + if (fd < 0) + return -1; + if (len && (size_t)write(fd, buf, +#if defined(WIN32) + (unsigned int) +#endif + len) != len) { + close(fd); + unlink(tmp); + return -1; + } + close(fd); + +#if defined(WIN32) + unlink(path); +#endif + if (rename(tmp, path)) { + unlink(tmp); + return -1; + } + + return 0; +} + +static void +saib_pool_state_save(saib_pool_t *pool) +{ + char path[384], tmp[384], buf[256]; + int n; + + n = lws_snprintf(buf, sizeof(buf), "{\"corpus\":%llu,\"known\":%llu," + "\"offer_since\":%llu,\"full_offer\":%llu}", + (unsigned long long)pool->cursor[SAI_POOL_NS_CORPUS], + (unsigned long long)pool->cursor[SAI_POOL_NS_KNOWN], + (unsigned long long)pool->offer_since, + (unsigned long long)pool->full_offer); + lws_snprintf(path, sizeof(path), "%s/.sai-pool-state", pool->dir); + lws_snprintf(tmp, sizeof(tmp), "%s/.sai-tmp", pool->dir); + if (saib_pool_write_file(tmp, path, buf, (size_t)n)) + lwsl_err("%s: unable to write %s\n", __func__, path); + + /* + * Tasks replacing a corpus sub need to know where we pulled up to, + * see README-pool.md + */ + n = lws_snprintf(buf, sizeof(buf), "%llu\n", + (unsigned long long)pool->cursor[SAI_POOL_NS_CORPUS]); + lws_snprintf(path, sizeof(path), "%s/corpus/.sai-pool-seq", pool->dir); + saib_pool_write_file(tmp, path, buf, (size_t)n); +} + +/* + * Nspawns waiting to start until their pool is pulled + */ + +static void +saib_pool_release_waiters(saib_pool_t *pool, const char *why) +{ + lws_sul_cancel(&pool->sul_wait); + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + pool->waiters.head) { + struct sai_nspawn *ns = lws_container_of(d, struct sai_nspawn, + pool_wait_list); + + lws_dll2_remove(&ns->pool_wait_list); + + if (why) + saib_task_logf(ns->spm, ns, NULL, "Pool %s: %s, " + "starting anyway", pool->name, why); + + if (saib_spawn_script(ns)) { + saib_task_logf(ns->spm, ns, NULL, "Builder %s could " + "not start step %d, failing the task", + ns->sp->name, ns->task->build_step + 1); + saib_set_ns_state(ns, NSSTATE_FAILED); + } + + } lws_end_foreach_dll_safe(d, d1); +} + +static void +saib_pool_wait_timeout_cb(lws_sorted_usec_list_t *sul) +{ + saib_pool_t *pool = lws_container_of(sul, saib_pool_t, sul_wait); + + saib_pool_release_waiters(pool, "syncing is taking too long"); +} + +static void +saib_pool_sync_cb(lws_sorted_usec_list_t *sul) +{ + saib_pool_t *pool = lws_container_of(sul, saib_pool_t, sul_sync); + + saib_pool_sync_start(pool); +} + +static void +saib_pool_session_timeout_cb(lws_sorted_usec_list_t *sul) +{ + saib_pool_t *pool = lws_container_of(sul, saib_pool_t, sul_session); + + lwsl_warn("%s: pool %s sync took too long\n", __func__, pool->name); + if (pool->ss) + lws_ss_destroy(&pool->ss); +} + +/* + * Tx side of a sync + */ + +static saib_pool_rec_t * +saib_pool_rec_new(int type, int ns, const char *name, size_t name_len, + size_t data_len) +{ + saib_pool_rec_t *r; + + r = malloc(sizeof(*r) + SAI_POOL_REC_HDR_LEN + name_len + data_len); + if (!r) + return NULL; + memset(r, 0, sizeof(*r)); + + r->len = SAI_POOL_REC_HDR_LEN + name_len + data_len; + sai_pool_rec_hdr_write((uint8_t *)&r[1], type, ns, name_len, data_len); + if (name_len) { + memcpy((uint8_t *)&r[1] + SAI_POOL_REC_HDR_LEN, name, name_len); + lws_strnncpy(r->name, name, name_len, sizeof(r->name)); + } + + return r; +} + +static uint8_t * +saib_pool_rec_data(saib_pool_rec_t *r) +{ + sai_pool_rec_hdr_t h; + + sai_pool_rec_hdr_read((uint8_t *)&r[1], &h); + + return (uint8_t *)&r[1] + SAI_POOL_REC_HDR_LEN + h.name_len; +} + +static void +saib_pool_queue(saib_pool_ss_t *ps, saib_pool_rec_t *r) +{ + lws_dll2_add_tail(&r->list, &ps->txq); + if (lws_ss_request_tx(ps->ss)) + lwsl_warn("%s: request tx failed\n", __func__); +} + +static void +saib_pool_free_list(lws_dll2_owner_t *o) +{ + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, o->head) { + lws_dll2_remove(d); + free(lws_container_of(d, saib_pool_rec_t, list)); + } lws_end_foreach_dll_safe(d, d1); +} + +/* read a whole local file into a PUT record, if it's still there */ + +static saib_pool_rec_t * +saib_pool_put_rec(int ns, const char *name, const char *path, size_t max) +{ + saib_pool_rec_t *r; + struct stat s; + int fd; + + fd = lws_open(path, O_RDONLY +#if defined(WIN32) + | _O_BINARY +#endif + ); + if (fd < 0) + return NULL; + + if (fstat(fd, &s) || (uint64_t)s.st_size > max) { + if (!fstat(fd, &s)) + lwsl_warn("%s: %s is too big to send (%llu)\n", __func__, + path, (unsigned long long)s.st_size); + close(fd); + return NULL; + } + + r = saib_pool_rec_new(SAI_POOL_REC_PUT, ns, name, strlen(name), + (size_t)s.st_size); + if (r && s.st_size && + read(fd, saib_pool_rec_data(r), +#if defined(WIN32) + (unsigned int) +#endif + (size_t)s.st_size) != (ssize_t)s.st_size) { + free(r); + r = NULL; + } + close(fd); + + if (r) + lws_strncpy(r->path, path, sizeof(r->path)); + + return r; +} + +/* the sync is over when nothing is left to send or to hear back about */ + +static int +saib_pool_session_idle(saib_pool_ss_t *ps) +{ + return ps->pulled && !ps->pulls_left && !ps->offers_left && + !ps->txq.count && !ps->acks.count && !ps->puts.count; +} + +static lws_ss_state_return_t +saib_pool_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, + size_t *len, int *flags) +{ + saib_pool_ss_t *ps = (saib_pool_ss_t *)userobj; + saib_pool_t *pool = ps->pool; + lws_struct_serialize_t *js; + sai_pool_hello_t hello; + saib_pool_rec_t *r; + size_t w, n = 0; + + *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + + if (!ps->sent_auth) { + /* it lands on sai-server's /builder, so the link auth first */ + *len = (size_t)lws_snprintf((char *)buf, *len, + "{\"schema\":\"" SAI_LINKAUTH_SCHEMA + "\",\"secret\":\"%s\"}", + builder.link_key ? builder.link_key : ""); + ps->sent_auth = 1; + + return lws_ss_request_tx(ps->ss); + } + + if (!ps->sent_hello) { + memset(&hello, 0, sizeof(hello)); + lws_strncpy(hello.task_uuid, pool->task_uuid, + sizeof(hello.task_uuid)); + lws_strncpy(hello.nonce, pool->nonce, sizeof(hello.nonce)); + + js = lws_struct_json_serialize_create(lsm_schema_pool_hello, + LWS_ARRAY_SIZE(lsm_schema_pool_hello), 0, + &hello); + if (!js) + return LWSSSSRET_DESTROY_ME; + lws_struct_json_serialize(js, buf, *len, &w); + lws_struct_json_serialize_destroy(&js); + *len = w; + ps->sent_hello = 1; + + return lws_ss_request_tx(ps->ss); + } + + /* coalesce queued records into one message, up to what fits */ + + while (ps->txq.head && n < *len) { + size_t l; + + r = lws_container_of(ps->txq.head, saib_pool_rec_t, list); + l = r->len - r->ofs; + if (l > *len - n) + l = *len - n; + memcpy(buf + n, (uint8_t *)&r[1] + r->ofs, l); + r->ofs += l; + n += l; + + if (r->ofs == r->len) { + lws_dll2_remove(&r->list); + if (r->name[0]) + /* PUT and REPLACE hear back */ + lws_dll2_add_tail(&r->list, &ps->acks); + else + free(r); + } + } + + /* the next corpus file the server wants, now there's room */ + + while (!ps->txq.count && ps->puts.head) { + saib_pool_put_t *pp = lws_container_of(ps->puts.head, + saib_pool_put_t, list); + char path[384]; + + lws_dll2_remove(&pp->list); + lws_snprintf(path, sizeof(path), "%s/corpus/%s", pool->dir, + pp->name); + r = saib_pool_put_rec(SAI_POOL_NS_CORPUS, pp->name, path, + SAI_POOL_ENTRY_MAX); + if (!r) + /* it's gone since we offered it, or too big */ + continue; + r->path[0] = '\0'; /* we keep corpus files, of course */ + lws_dll2_add_tail(&r->list, &ps->txq); + } + + if (!n) { + if (ps->done) + return LWSSSSRET_DESTROY_ME; + return LWSSSSRET_TX_DONT_SEND; + } + + *len = n; + + if (ps->txq.count) + return lws_ss_request_tx(ps->ss); + + return LWSSSSRET_OK; +} + +/* + * The offer, findings and replace scans, once the pull is done + */ + +typedef struct { + saib_pool_ss_t *ps; + const char *sub; /* as we scan inside one */ + uint8_t *list; /* names being collected */ + size_t list_len; + size_t list_max; + uint64_t since; + int ns; + char overflow; +} saib_pool_scan_t; + +static int +saib_pool_offer_flush(saib_pool_scan_t *sc) +{ + saib_pool_rec_t *r; + + if (!sc->list_len) + return 0; + + r = saib_pool_rec_new(SAI_POOL_REC_OFFER, SAI_POOL_NS_CORPUS, NULL, 0, + sc->list_len); + if (!r) + return -1; + memcpy(saib_pool_rec_data(r), sc->list, sc->list_len); + saib_pool_queue(sc->ps, r); + sc->ps->offers_left++; + sc->list_len = 0; + + return 0; +} + +/* each file in corpus/<sub>/ or findings/<sub>/ */ + +static int +saib_pool_scan_file_cb(const char *dirpath, void *user, + struct lws_dir_entry *lde) +{ + saib_pool_scan_t *sc = (saib_pool_scan_t *)user; + char name[SAI_POOL_REC_NAME_MAX + 1], path[384]; + size_t fl = strlen(lde->name); + saib_pool_rec_t *r; + int nl; + + if (lde->type != LDOT_FILE) + return 0; + + nl = lws_snprintf(name, sizeof(name), "%s/%s", sc->sub, lde->name); + if (!sai_pool_entry_name_ok(sc->ns, name, (size_t)nl)) + /* not something we sync, eg, a temp file */ + return 0; + + lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); + + if (sc->ns == SAI_POOL_NS_FINDINGS) { + r = saib_pool_put_rec(SAI_POOL_NS_FINDINGS, name, path, + SAI_POOL_FINDING_MAX); + if (r) + /* the path stays set, we delete it once it's stored */ + saib_pool_queue(sc->ps, r); + + return 0; + } + + /* otherwise we're collecting the names in a sub for a replace */ + + if (sc->list_len + fl + 1 > sc->list_max) { + sc->overflow = 1; + return 1; + } + memcpy(sc->list + sc->list_len, lde->name, fl); + sc->list_len += fl; + sc->list[sc->list_len++] = '\n'; + + return 0; +} + +static int +saib_pool_offer_file_cb(const char *dirpath, void *user, + struct lws_dir_entry *lde) +{ + saib_pool_scan_t *sc = (saib_pool_scan_t *)user; + char name[SAI_POOL_REC_NAME_MAX + 1], path[384]; + struct stat s; + int nl; + + if (lde->type != LDOT_FILE) + return 0; + + nl = lws_snprintf(name, sizeof(name), "%s/%s", sc->sub, lde->name); + if (!sai_pool_entry_name_ok(SAI_POOL_NS_CORPUS, name, (size_t)nl)) + return 0; + + if (sc->since) { + lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); + if (stat(path, &s) || (uint64_t)s.st_mtime < sc->since) + /* nothing new */ + return 0; + } + + if (sc->list_len + (size_t)nl + 1 > SAI_POOL_OFFER_MAX && + saib_pool_offer_flush(sc)) + return 1; + + memcpy(sc->list + sc->list_len, name, (size_t)nl); + sc->list_len += (size_t)nl; + sc->list[sc->list_len++] = '\n'; + + return 0; +} + +/* each sub in corpus/ or findings/ */ + +static int +saib_pool_scan_sub_cb(const char *dirpath, void *user, + struct lws_dir_entry *lde) +{ + saib_pool_scan_t *sc = (saib_pool_scan_t *)user; + char path[384]; + + if (lde->type != LDOT_DIR || + !sai_pool_sub_ok(lde->name, strlen(lde->name))) + return 0; + + lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); + sc->sub = lde->name; + + return lws_dir(path, sc, sc->ns == SAI_POOL_NS_CORPUS ? + saib_pool_offer_file_cb : + saib_pool_scan_file_cb); +} + +/* + * A task that replaced the contents of corpus/<sub> leaves + * corpus/.sai-replace-<sub> holding the cursor it began from, see + * README-pool.md + */ + +static int +saib_pool_replace_marker_cb(const char *dirpath, void *user, + struct lws_dir_entry *lde) +{ + saib_pool_scan_t *sc = (saib_pool_scan_t *)user; + char path[384], sub[33], num[32]; + saib_pool_scan_t rsc; + saib_pool_rec_t *r; + uint64_t base = 0; + size_t sl; + int fd, n, i; + + if (lde->type != LDOT_FILE || strncmp(lde->name, ".sai-replace-", 13)) + return 0; + + sl = strlen(lde->name + 13); + if (!sai_pool_sub_ok(lde->name + 13, sl)) + return 0; + lws_strncpy(sub, lde->name + 13, sizeof(sub)); + + lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); + fd = lws_open(path, O_RDONLY); + if (fd < 0) + return 0; + n = (int)read(fd, num, sizeof(num) - 1); + close(fd); + if (n <= 0) + return 0; + for (i = 0; i < n && num[i] >= '0' && num[i] <= '9'; i++) + base = (base * 10) + (uint64_t)(num[i] - '0'); + if (!i) { + lwsl_warn("%s: ignoring %s, no cursor in it\n", __func__, path); + return 0; + } + + memset(&rsc, 0, sizeof(rsc)); + rsc.ps = sc->ps; + rsc.ns = SAI_POOL_NS_CORPUS; + rsc.sub = sub; + rsc.list_max = SAI_POOL_LIST_MAX - 8; + rsc.list = malloc(rsc.list_max); + if (!rsc.list) + return 1; + + { + char subdir[384]; + + lws_snprintf(subdir, sizeof(subdir), "%s/%s", dirpath, sub); + lws_dir(subdir, &rsc, saib_pool_scan_file_cb); + } + + if (rsc.overflow) { + lwsl_warn("%s: %s has too many files to replace\n", __func__, + sub); + free(rsc.list); + return 0; + } + + r = saib_pool_rec_new(SAI_POOL_REC_REPLACE, SAI_POOL_NS_CORPUS, sub, + strlen(sub), 8 + rsc.list_len); + if (r) { + sai_pool_u64_write(saib_pool_rec_data(r), base); + memcpy(saib_pool_rec_data(r) + 8, rsc.list, rsc.list_len); + lws_strncpy(r->path, path, sizeof(r->path)); + r->is_replace = 1; + r->replace_base = base; + saib_pool_queue(sc->ps, r); + } + free(rsc.list); + + return 0; +} + +static void +saib_pool_after_pull(saib_pool_ss_t *ps) +{ + saib_pool_t *pool = ps->pool; + saib_pool_scan_t sc; + char path[384]; + + ps->pulled = 1; + pool->last_pull = lws_now_usecs(); + saib_pool_state_save(pool); + saib_pool_release_waiters(pool, NULL); + + /* offer the server what appeared in corpus/ since the last sync */ + + ps->scan_start = (uint64_t)lws_now_secs(); + ps->full_offer = !pool->full_offer || + ps->scan_start - pool->full_offer > + SAIB_POOL_FULL_OFFER_SECS; + + memset(&sc, 0, sizeof(sc)); + sc.ps = ps; + sc.ns = SAI_POOL_NS_CORPUS; + if (!ps->full_offer && pool->offer_since > SAIB_POOL_OFFER_SLACK_SECS) + sc.since = pool->offer_since - SAIB_POOL_OFFER_SLACK_SECS; + sc.list = malloc(SAI_POOL_OFFER_MAX); + if (sc.list) { + lws_snprintf(path, sizeof(path), "%s/corpus", pool->dir); + lws_dir(path, &sc, saib_pool_scan_sub_cb); + saib_pool_offer_flush(&sc); + free(sc.list); + } + + /* send everything in findings/, deleting each once it's stored */ + + memset(&sc, 0, sizeof(sc)); + sc.ps = ps; + sc.ns = SAI_POOL_NS_FINDINGS; + lws_snprintf(path, sizeof(path), "%s/findings", pool->dir); + lws_dir(path, &sc, saib_pool_scan_sub_cb); + + /* any replacing tasks asked for */ + + memset(&sc, 0, sizeof(sc)); + sc.ps = ps; + lws_snprintf(path, sizeof(path), "%s/corpus", pool->dir); + lws_dir(path, &sc, saib_pool_replace_marker_cb); +} + +/* + * Rx side of a sync + */ + +static int +saib_pool_entry(saib_pool_t *pool, const sai_pool_rec_hdr_t *h, + const char *name, const uint8_t *data) +{ + char path[384], dir[384], tmp[384]; + const char *sl = memchr(name, '/', h->name_len); + + if ((h->ns != SAI_POOL_NS_CORPUS && h->ns != SAI_POOL_NS_KNOWN) || + !sai_pool_entry_name_ok(h->ns, name, h->name_len)) + return -1; + + lws_snprintf(dir, sizeof(dir), "%s/%s/%.*s", pool->dir, ns_dir[h->ns], + (int)lws_ptr_diff(sl, name), name); + lws_snprintf(path, sizeof(path), "%s/%s/%.*s", pool->dir, + ns_dir[h->ns], (int)h->name_len, name); + + if (h->type == SAI_POOL_REC_DEAD) { + unlink(path); + return 0; + } + + /* the server checked, but it's going in a file with that name */ + if (!sai_pool_content_matches(data, h->len, sl + 1)) { + lwsl_err("%s: content isn't %s\n", __func__, path); + return -1; + } + + if (mkdir(dir, 0755) && errno != EEXIST) + return -1; + + lws_snprintf(tmp, sizeof(tmp), "%s/.sai-tmp", pool->dir); + + return saib_pool_write_file(tmp, path, data, h->len); +} + +static int +saib_pool_ack(saib_pool_ss_t *ps, const char *name, size_t name_len) +{ + saib_pool_t *pool = ps->pool; + saib_pool_rec_t *r; + + if (!ps->acks.head) + return -1; + + /* the server answers in order */ + + r = lws_container_of(ps->acks.head, saib_pool_rec_t, list); + if (strlen(r->name) != name_len || memcmp(r->name, name, name_len)) { + lwsl_err("%s: ACK for %.*s, expected %s\n", __func__, + (int)name_len, name, r->name); + return -1; + } + lws_dll2_remove(&r->list); + + if (r->path[0]) + /* a finding it has now, or a replace marker it's dealt with */ + unlink(r->path); + + if (r->is_replace) { + /* + * Go back to where the replacing task started from, so we + * hear about what the server removed and what anyone added + * meanwhile, which the replace may have lost locally + */ + if (pool->cursor[SAI_POOL_NS_CORPUS] > r->replace_base) + pool->cursor[SAI_POOL_NS_CORPUS] = r->replace_base; + saib_pool_state_save(pool); + } + + free(r); + + return 0; +} + +static int +saib_pool_record(saib_pool_ss_t *ps, const sai_pool_rec_hdr_t *h, + const char *name, const uint8_t *data) +{ + saib_pool_t *pool = ps->pool; + const char *p, *end; + + switch (h->type) { + case SAI_POOL_REC_ENTRY: + case SAI_POOL_REC_DEAD: + return saib_pool_entry(pool, h, name, data); + + case SAI_POOL_REC_PULL_END: + if (h->ns > SAI_POOL_NS_KNOWN || h->len != 8 || !ps->pulls_left) + return -1; + pool->cursor[h->ns] = sai_pool_u64_read(data); + if (!--ps->pulls_left) + saib_pool_after_pull(ps); + return 0; + + case SAI_POOL_REC_WANT: + if (h->ns != SAI_POOL_NS_CORPUS || !ps->offers_left) + return -1; + ps->offers_left--; + + p = (const char *)data; + end = p + h->len; + while (p < end) { + const char *nl = memchr(p, '\n', + lws_ptr_diff_size_t(end, p)), + *e = nl ? nl : end; + size_t l = lws_ptr_diff_size_t(e, p); + saib_pool_put_t *pp; + + if (l && sai_pool_entry_name_ok(SAI_POOL_NS_CORPUS, p, l)) { + pp = lwsac_use_zero(&ps->ac_puts, sizeof(*pp), + 8192); + if (!pp) + return -1; + lws_strnncpy(pp->name, p, l, sizeof(pp->name)); + lws_dll2_add_tail(&pp->list, &ps->puts); + } + p = e + 1; + } + if (ps->puts.count && lws_ss_request_tx(ps->ss)) + return -1; + return 0; + + case SAI_POOL_REC_ACK: + return saib_pool_ack(ps, name, h->name_len); + } + + lwsl_err("%s: unexpected record type 0x%x\n", __func__, h->type); + + return -1; +} + +static lws_ss_state_return_t +saib_pool_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + saib_pool_ss_t *ps = (saib_pool_ss_t *)userobj; + sai_pool_rec_hdr_t h; + size_t ofs = 0, need; + + if (ps->rxb_len + len > ps->rxb_alloc) { + size_t na = ps->rxb_len + len + 4096; + uint8_t *nb; + + if (na > SAI_POOL_ENTRY_MAX + SAI_POOL_OFFER_MAX + 65536) + return LWSSSSRET_DESTROY_ME; + nb = realloc(ps->rxb, na); + if (!nb) + return LWSSSSRET_DESTROY_ME; + ps->rxb = nb; + ps->rxb_alloc = na; + } + + memcpy(ps->rxb + ps->rxb_len, buf, len); + ps->rxb_len += len; + + while (ps->rxb_len - ofs >= SAI_POOL_REC_HDR_LEN) { + sai_pool_rec_hdr_read(ps->rxb + ofs, &h); + + if (h.ns >= SAI_POOL_NS_COUNT || + h.name_len > SAI_POOL_REC_NAME_MAX || + h.len > sai_pool_rec_max(h.ns, h.type)) { + lwsl_err("%s: bad record from server\n", __func__); + return LWSSSSRET_DESTROY_ME; + } + + need = SAI_POOL_REC_HDR_LEN + h.name_len + h.len; + if (ps->rxb_len - ofs < need) + break; + + if (saib_pool_record(ps, &h, (const char *)ps->rxb + ofs + + SAI_POOL_REC_HDR_LEN, + ps->rxb + ofs + SAI_POOL_REC_HDR_LEN + + h.name_len)) + return LWSSSSRET_DESTROY_ME; + + ofs += need; + } + + if (ofs) { + memmove(ps->rxb, ps->rxb + ofs, ps->rxb_len - ofs); + ps->rxb_len -= ofs; + } + + if (saib_pool_session_idle(ps)) { + /* everything done, a successful sync */ + ps->pool->offer_since = ps->scan_start; + if (ps->full_offer) + ps->pool->full_offer = ps->scan_start; + saib_pool_state_save(ps->pool); + ps->done = 1; + + return lws_ss_request_tx(ps->ss); + } + + return LWSSSSRET_OK; +} + +static lws_ss_state_return_t +saib_pool_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + saib_pool_ss_t *ps = (saib_pool_ss_t *)userobj; + saib_pool_t *pool; + saib_pool_rec_t *r; + uint8_t since[8]; + int ns; + + switch (state) { + case LWSSSCS_CREATING: + ps->pool = (saib_pool_t *)ps->opaque_data; + break; + + case LWSSSCS_CONNECTED: + /* pull what happened since last time, in both namespaces */ + for (ns = SAI_POOL_NS_CORPUS; ns <= SAI_POOL_NS_KNOWN; ns++) { + sai_pool_u64_write(since, ps->pool->cursor[ns]); + r = saib_pool_rec_new(SAI_POOL_REC_PULL, ns, NULL, 0, 8); + if (!r) + return LWSSSSRET_DESTROY_ME; + memcpy(saib_pool_rec_data(r), since, 8); + lws_dll2_add_tail(&r->list, &ps->txq); + ps->pulls_left++; + } + return lws_ss_request_tx(ps->ss); + + case LWSSSCS_DISCONNECTED: + case LWSSSCS_ALL_RETRIES_FAILED: + case LWSSSCS_TIMEOUT: + return LWSSSSRET_DESTROY_ME; + + case LWSSSCS_DESTROYING: + pool = ps->pool; + + saib_pool_free_list(&ps->txq); + saib_pool_free_list(&ps->acks); + lws_dll2_owner_clear(&ps->puts); + lwsac_free(&ps->ac_puts); + free(ps->rxb); + + if (!pool) + break; + + pool->ss = NULL; + lws_sul_cancel(&pool->sul_session); + + if (!ps->done) + lwsl_warn("%s: pool %s sync didn't complete\n", + __func__, pool->name); + + if (pool->waiters.count) + /* nothing more is coming for them this time */ + saib_pool_release_waiters(pool, + "couldn't sync with the server"); + + /* decide when we sync next, if at all */ + + if (pool->resync) { + pool->resync = 0; + lws_sul_schedule(builder.context, 0, &pool->sul_sync, + saib_pool_sync_cb, 1); + } else + if (pool->users) + lws_sul_schedule(builder.context, 0, + &pool->sul_sync, + saib_pool_sync_cb, + SAIB_POOL_SYNC_INTERVAL_US); + else + if (!ps->done && + ++pool->final_tries < SAIB_POOL_FINAL_TRIES) + /* the last one after the task failed */ + lws_sul_schedule(builder.context, 0, + &pool->sul_sync, + saib_pool_sync_cb, + SAIB_POOL_SYNC_INTERVAL_US); + + saib_reassess_idle_situation(); + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +const lws_ss_info_t ssi_sai_pool = { + .handle_offset = offsetof(saib_pool_ss_t, ss), + .opaque_user_data_offset = offsetof(saib_pool_ss_t, opaque_data), + .rx = saib_pool_rx, + .tx = saib_pool_tx, + .state = saib_pool_state, + .user_alloc = sizeof(saib_pool_ss_t), + .streamtype = "sai_pool" +}; + +static void +saib_pool_sync_start(saib_pool_t *pool) +{ + lws_sul_cancel(&pool->sul_sync); + + if (pool->ss) { + /* one at a time, but have another go after this one */ + pool->resync = 1; + return; + } + + if (!pool->spm || !pool->spm->url || !pool->task_uuid[0]) + return; + + if (lws_ss_create(builder.context, 0, &ssi_sai_pool, pool, &pool->ss, + NULL, NULL)) { + lwsl_err("%s: unable to create pool sync stream\n", __func__); + pool->ss = NULL; + if (pool->users) + lws_sul_schedule(builder.context, 0, &pool->sul_sync, + saib_pool_sync_cb, + SAIB_POOL_SYNC_INTERVAL_US); + saib_pool_release_waiters(pool, "couldn't start syncing"); + return; + } + + if (lws_ss_set_metadata(pool->ss, "url", pool->spm->url, + strlen(pool->spm->url))) + lwsl_warn("%s: unable to set metadata\n", __func__); + + lws_sul_schedule(builder.context, 0, &pool->sul_session, + saib_pool_session_timeout_cb, SAIB_POOL_SESSION_MAX_US); + + if (lws_ss_client_connect(pool->ss)) + lwsl_warn("%s: connect failed\n", __func__); +} + +/* + * The task running in this nspawn names a pool: find or create our copy of + * it, and account for the nspawn using it + */ + +int +saib_pool_attach(struct sai_nspawn *ns) +{ + char path[300], pur[160], *p; + saib_pool_t *pool = NULL; + int n; + + if (!ns->task->pool[0] || ns->task->build_step < 2) + /* nothing to sync, or it's just git steps */ + return 0; + + if (!sai_pool_name_ok(ns->task->pool)) { + saib_task_logf(ns->spm, ns, NULL, "Ignoring bad pool name"); + return 0; + } + + /* one copy per server, repo and pool */ + + lws_snprintf(pur, sizeof(pur), "%s-%s-%s", ns->spm->name, + ns->task->repo_name, ns->task->pool); + lws_filename_purify_inplace(pur); + p = pur; + while ((p = strchr(p, '/'))) + *p++ = '_'; + + lws_start_foreach_dll(struct lws_dll2 *, d, builder.pool_owner.head) { + saib_pool_t *xp = lws_container_of(d, saib_pool_t, list); + + if (!strcmp(xp->key, pur)) { + pool = xp; + break; + } + + } lws_end_foreach_dll(d); + + if (!pool) { + pool = malloc(sizeof(*pool)); + if (!pool) + return -1; + memset(pool, 0, sizeof(*pool)); + + lws_strncpy(pool->name, ns->task->pool, sizeof(pool->name)); + lws_strncpy(pool->key, pur, sizeof(pool->key)); + lws_snprintf(pool->dir, sizeof(pool->dir), "%s/pools/%s", + builder.home, pur); + + lws_snprintf(path, sizeof(path), "%s/pools", builder.home); + if (mkdir(path, 0755) && errno != EEXIST) + goto bail; + if (mkdir(pool->dir, 0755) && errno != EEXIST) + goto bail; + for (n = 0; n < (int)LWS_ARRAY_SIZE(ns_dir); n++) { + lws_snprintf(path, sizeof(path), "%s/%s", pool->dir, + ns_dir[n]); + if (mkdir(path, 0755) && errno != EEXIST) + goto bail; + } + + saib_pool_state_load(pool); + lws_dll2_add_tail(&pool->list, &builder.pool_owner); + } + + /* the latest task to use it is who we sync on behalf of */ + + pool->spm = ns->spm; + lws_strncpy(pool->task_uuid, ns->task->uuid, sizeof(pool->task_uuid)); + lws_strncpy(pool->nonce, ns->task->art_up_nonce, sizeof(pool->nonce)); + pool->final_tries = 0; + + ns->pool = pool; + if (!pool->users++ && !pool->ss && !pool->sul_sync.list.owner) + lws_sul_schedule(builder.context, 0, &pool->sul_sync, + saib_pool_sync_cb, SAIB_POOL_SYNC_INTERVAL_US); + + return 0; + +bail: + saib_task_logf(ns->spm, ns, NULL, "Unable to create pool dir %s: " + "errno %d", pool->dir, errno); + free(pool); + + return -1; +} + +/* + * The step would like to start: if its pool hasn't been pulled lately, it + * waits for that. Returns 1 if the step was deferred, and will be spawned + * when the pull is done or given up on. + */ + +int +saib_pool_defer_spawn(struct sai_nspawn *ns) +{ + saib_pool_t *pool = ns->pool; + + if (!pool || (pool->last_pull && + lws_now_usecs() - pool->last_pull < SAIB_POOL_PULL_FRESH_US)) + return 0; + + saib_task_logf(ns->spm, ns, NULL, "Syncing pool %s before starting", + pool->name); + + lws_dll2_add_tail(&ns->pool_wait_list, &pool->waiters); + if (!pool->sul_wait.list.owner) + lws_sul_schedule(builder.context, 0, &pool->sul_wait, + saib_pool_wait_timeout_cb, + SAIB_POOL_WAIT_MAX_US); + + /* + * If a sync is already going, it hasn't pulled yet (or the pull would + * be fresh), so its pull releases us too + */ + if (!pool->ss) + saib_pool_sync_start(pool); + + return 1; +} + +/* + * The nspawn is being stopped while it waits for its pool: it never started, + * so it has to be failed here. Returns 1 if it was waiting. + */ + +int +saib_pool_waiter_abort(struct sai_nspawn *ns) +{ + if (lws_dll2_is_detached(&ns->pool_wait_list)) + return 0; + + lws_dll2_remove(&ns->pool_wait_list); + saib_set_ns_state(ns, NSSTATE_FAILED); + if (!ns->idle_yield) + ns->retcode = SAISPRF_TERMINATED; + + return 1; +} + +/* the nspawn is done with the pool, sync what it left there */ + +void +saib_pool_detach(struct sai_nspawn *ns) +{ + saib_pool_t *pool = ns->pool; + + if (!pool) + return; + + lws_dll2_remove(&ns->pool_wait_list); + ns->pool = NULL; + + if (--pool->users) + return; + + lws_sul_cancel(&pool->sul_sync); + saib_pool_sync_start(pool); +} + +/* what the task is told about its pool */ + +void +saib_pool_env(struct sai_nspawn *ns, char *buf, size_t len) +{ + saib_pool_t *pool = ns->pool; + + buf[0] = '\0'; + if (!pool) + return; + + lws_snprintf(buf, len, +#if defined(WIN32) + "set SAI_POOL_DIR=%s\\corpus\n" + "set SAI_POOL_KNOWN=%s\\known\n" + "set SAI_POOL_FINDINGS=%s\\findings\n", +#else + "export SAI_POOL_DIR=%s/corpus\n" + "export SAI_POOL_KNOWN=%s/known\n" + "export SAI_POOL_FINDINGS=%s/findings\n", +#endif + pool->dir, pool->dir, pool->dir); +} + +/* are we still syncing, or about to, so we shouldn't power off yet? */ + +int +saib_pool_busy(void) +{ + lws_start_foreach_dll(struct lws_dll2 *, d, builder.pool_owner.head) { + saib_pool_t *pool = lws_container_of(d, saib_pool_t, list); + + if (pool->ss || (!pool->users && pool->sul_sync.list.owner)) + return 1; + + } lws_end_foreach_dll(d); + + return 0; +} + +void +saib_pool_destroy_all(void) +{ + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + builder.pool_owner.head) { + saib_pool_t *pool = lws_container_of(d, saib_pool_t, list); + + lws_sul_cancel(&pool->sul_sync); + lws_sul_cancel(&pool->sul_wait); + lws_sul_cancel(&pool->sul_session); + if (pool->ss) { + /* its DESTROYING mustn't touch the pool after this */ + saib_pool_ss_t *ps = lws_ss_to_user_object(pool->ss); + + ps->pool = NULL; + lws_ss_destroy(&pool->ss); + } + lws_dll2_remove(&pool->list); + free(pool); + + } lws_end_foreach_dll_safe(d, d1); +} diff --git a/src/builder/b-power.c b/src/builder/b-power.c index 16f8ea8..1c49793 100644 --- a/src/builder/b-power.c +++ b/src/builder/b-power.c @@ -628,6 +628,12 @@ saib_reassess_idle_situation(void) in_use = 1; } + if (saib_pool_busy()) { + /* what the last tasks left in their pools isn't synced yet */ + lwsl_notice("%s: pools still syncing\n", __func__); + in_use = 1; + } + saib_power_event(in_use ? SAIB_PWR_EV_BUSY : SAIB_PWR_EV_IDLE); return 0; diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 0d774ac..4f827cb 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -202,6 +202,7 @@ struct sai_builder { lws_dll2_owner_t lsp_owner; /* list of lws_spawn_piped */ lws_dll2_owner_t jobdir_hold_owner; /* saib_jobdir_hold_t */ lws_dll2_owner_t shell_owner; /* list of sai_shell */ + lws_dll2_owner_t pool_owner; /* saib_pool_t, see b-pool.c */ struct lws_ss_handle *ss_stay; struct lws_ss_handle *ss_power_off; @@ -344,7 +345,8 @@ struct ws_capture_chunk { extern struct sai_builder builder; -extern const lws_ss_info_t ssi_sai_builder, ssi_sai_mirror, ssi_sai_artifact; +extern const lws_ss_info_t ssi_sai_builder, ssi_sai_mirror, ssi_sai_artifact, + ssi_sai_pool; extern const struct lws_protocols protocol_com_warmcat_sai; int saib_config_global(struct sai_builder *builder, const char *d); @@ -466,6 +468,29 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, void saib_sul_task_cancel(struct lws_sorted_usec_list *sul); +/* b-pool.c */ + +int +saib_pool_attach(struct sai_nspawn *ns); + +int +saib_pool_defer_spawn(struct sai_nspawn *ns); + +int +saib_pool_waiter_abort(struct sai_nspawn *ns); + +void +saib_pool_detach(struct sai_nspawn *ns); + +void +saib_pool_env(struct sai_nspawn *ns, char *buf, size_t len); + +int +saib_pool_busy(void); + +void +saib_pool_destroy_all(void); + int saib_suspender_get_pipe(void); diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index 8672c34..83dbc70 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -174,6 +174,24 @@ static const char * const default_ss_policy = "]" "}}," /* + * Ephemeral connections to the same server syncing a pool, + * see b-pool.c + */ + "{\"sai_pool\": {" + "\"endpoint\":" "\"${url}\"," + "\"port\":" "443," + "\"protocol\":" "\"ws\"," + "\"ws_subprotocol\":" "\"com-warmcat-sai\"," + "\"http_url\":" "\"\"," /* filled in by url */ + "\"tls\":" "true," + "\"opportunistic\":" "true," + "\"ws_binary\":" "true," + "\"retry\":" "\"default\"," + "\"metadata\": [" + "{\"url\": \"\"}" + "]" + "}}," + /* * Used to connect to sai-power to ask for power-off */ "{\"sai_power\": {" @@ -1072,6 +1090,7 @@ saib_app_run(int argc, const char **argv) } lws_end_foreach_dll_safe(mp, mp1); + saib_pool_destroy_all(); saib_config_destroy(&builder); saib_power_shutdown(); diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 75d264e..60b194d 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -458,6 +458,10 @@ saib_idletask_yield_all(void) xns->idle_yield = 1; + if (saib_pool_waiter_abort(xns)) + /* it hadn't started, waiting for its pool */ + continue; + if (!xns->op || !xns->op->lsp) /* * Nothing running to stop, eg, between @@ -597,6 +601,9 @@ saib_task_destroy(struct sai_nspawn *ns) lws_sul_cancel(&ns->sul_cleaner); lws_sul_cancel(&ns->sul_task_cancel); + /* sync what the task left in its pool, if it has one */ + saib_pool_detach(ns); + /* * Any artifact uploads still referencing us must go first... their * DESTROYING unlinks their temp file and accounts against @@ -1131,6 +1138,10 @@ saib_sul_task_cancel(struct lws_sorted_usec_list *sul) char s[64]; int n; + if (saib_pool_waiter_abort(ns)) + /* it never started, it was waiting for its pool */ + return; + if (!ns->op || !ns->op->lsp) return; @@ -1538,7 +1549,15 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, ns->user_cancel = 0; ns->spins = 0; - if (saib_spawn_script(ns)) { + if (saib_pool_attach(ns)) + goto bail; + + /* + * If the task has a pool we haven't pulled lately, the step is spawned + * once we did (or gave up), see b-pool.c + */ + + if (!saib_pool_defer_spawn(ns) && saib_spawn_script(ns)) { lwsl_err("%s: saib_spawn_script failed\n", __func__); goto bail; } diff --git a/src/common/c-pool.c b/src/common/c-pool.c new file mode 100644 index 0000000..20840f6 --- /dev/null +++ b/src/common/c-pool.c @@ -0,0 +1,214 @@ +/* + * Sai - ./src/common/c-pool.c + * + * Copyright (C) 2019 - 2026 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * Pool sync helpers shared by sai-builder and sai-server, see + * READMEs/README-pool.md. + * + * Everything a pool record names ends up as a path component on the builder, + * and in the server's db, so both sides check names with these before using + * them, whoever sent them. + */ + +#include <libwebsockets.h> +#include <string.h> + +#include "include/private.h" + +/* + * A pool name, from .sai.json: short, lowercase alnum plus - and _ + */ + +int +sai_pool_name_ok(const char *name) +{ + size_t n = 0; + + if (!name || !*name) + return 0; + + while (name[n]) { + char c = name[n]; + + if (!((c >= 'a' && c <= 'z') || (c >= '0' && c <= '9') || + c == '-' || c == '_')) + return 0; + if (++n > 32) + return 0; + } + + return 1; +} + +static int +sai_pool_safe_char(char c) +{ + return (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || + (c >= '0' && c <= '9') || c == '-' || c == '_' || c == '.'; +} + +/* + * The subdir part of an entry name, eg, "corpus-h2". It's a single path + * component that can't be hidden, ".." or anything else special. + */ + +int +sai_pool_sub_ok(const char *sub, size_t len) +{ + size_t n; + + if (!len || len > 32 || sub[0] == '.') + return 0; + + for (n = 0; n < len; n++) + if (!sai_pool_safe_char(sub[n])) + return 0; + + return 1; +} + +/* + * A whole entry name, "<sub>/<name>". In the content addressed namespaces the + * name is the 40 char lowercase hex sha1 of the content; findings can be named + * anything safe, that isn't hidden. + */ + +int +sai_pool_entry_name_ok(int ns, const char *name, size_t len) +{ + const char *sl = memchr(name, '/', len); + size_t sub_len, nl, n; + const char *nm; + + if (!sl) + return 0; + + sub_len = lws_ptr_diff_size_t(sl, name); + if (!sai_pool_sub_ok(name, sub_len)) + return 0; + + nm = sl + 1; + nl = len - sub_len - 1; + + if (ns == SAI_POOL_NS_FINDINGS) { + if (!nl || nl > 64 || nm[0] == '.') + return 0; + for (n = 0; n < nl; n++) + if (!sai_pool_safe_char(nm[n])) + return 0; + + return 1; + } + + if (nl != 40) + return 0; + + for (n = 0; n < nl; n++) + if (!((nm[n] >= '0' && nm[n] <= '9') || + (nm[n] >= 'a' && nm[n] <= 'f'))) + return 0; + + return 1; +} + +/* is the content really what its content addressed name says? */ + +int +sai_pool_content_matches(const uint8_t *data, size_t len, const char *sha1hex) +{ + uint8_t md[20]; + char hex[41]; + int n; + + lws_SHA1(data, len, md); + for (n = 0; n < 20; n++) + lws_snprintf(hex + (n * 2), 3, "%02x", md[n]); + + return !memcmp(hex, sha1hex, 40); +} + +void +sai_pool_rec_hdr_write(uint8_t *p, int type, int ns, size_t name_len, + size_t len) +{ + p[0] = (uint8_t)type; + p[1] = (uint8_t)ns; + p[2] = (uint8_t)(name_len >> 8); + p[3] = (uint8_t)name_len; + p[4] = (uint8_t)(len >> 24); + p[5] = (uint8_t)(len >> 16); + p[6] = (uint8_t)(len >> 8); + p[7] = (uint8_t)len; +} + +void +sai_pool_rec_hdr_read(const uint8_t *p, sai_pool_rec_hdr_t *h) +{ + h->type = p[0]; + h->ns = p[1]; + h->name_len = (uint16_t)((p[2] << 8) | p[3]); + h->len = ((uint32_t)p[4] << 24) | ((uint32_t)p[5] << 16) | + ((uint32_t)p[6] << 8) | (uint32_t)p[7]; +} + +/* the most data a record of this type and namespace may carry */ + +size_t +sai_pool_rec_max(int ns, int type) +{ + switch (type) { + case SAI_POOL_REC_PULL: + case SAI_POOL_REC_PULL_END: + return 8; + case SAI_POOL_REC_OFFER: + case SAI_POOL_REC_WANT: + return SAI_POOL_OFFER_MAX; + case SAI_POOL_REC_REPLACE: + return SAI_POOL_LIST_MAX; + case SAI_POOL_REC_PUT: + case SAI_POOL_REC_ENTRY: + return ns == SAI_POOL_NS_FINDINGS ? SAI_POOL_FINDING_MAX : + SAI_POOL_ENTRY_MAX; + } + + return 0; +} + +void +sai_pool_u64_write(uint8_t *p, uint64_t v) +{ + int n; + + for (n = 7; n >= 0; n--) { + p[n] = (uint8_t)v; + v >>= 8; + } +} + +uint64_t +sai_pool_u64_read(const uint8_t *p) +{ + uint64_t v = 0; + int n; + + for (n = 0; n < 8; n++) + v = (v << 8) | p[n]; + + return v; +} diff --git a/src/common/c-sqlite3.c b/src/common/c-sqlite3.c index 0bf57cd..d5b682d 100644 --- a/src/common/c-sqlite3.c +++ b/src/common/c-sqlite3.c @@ -132,6 +132,9 @@ sai_event_db_ensure_open(struct lws_context *cx, lws_dll2_owner_t *sqlite3_cache /* default 0 so every older task reads as not idle */ sqlite3_exec(*ppdb, "ALTER TABLE tasks ADD COLUMN idle integer default 0;", NULL, NULL, &err); if (err) sqlite3_free(err); + err = NULL; + sqlite3_exec(*ppdb, "ALTER TABLE tasks ADD COLUMN pool varchar;", NULL, NULL, &err); + if (err) sqlite3_free(err); } sai_sqlite3_statement(*ppdb, "CREATE UNIQUE INDEX IF NOT EXISTS idx_task_uuid ON tasks(uuid, run);", "create task index"); diff --git a/src/common/include/private.h b/src/common/include/private.h index 5831cd0..489ba1e 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -292,6 +292,12 @@ typedef struct { * does not count towards its event's state. See README-idle.md. */ int idle; + /* + * Optional name of the repo's shared pool the task works with, eg, + * "fuzz" for fuzzing corpora. The builder keeps it synced with the + * server while the task runs. See READMEs/README-pool.md. + */ + char pool[33]; } sai_task_t; struct saib_logproxy { @@ -305,6 +311,8 @@ struct saib_resproxy { struct sai_nspawn *ns; }; +struct saib_pool; + struct sai_nspawn { char inp[512]; char inp_vn[16]; @@ -338,6 +346,10 @@ struct sai_nspawn { /* builder: sai_artifact_t of uploads still in flight for this ns */ lws_dll2_owner_t artifact_owner; + /* builder: the task's pool, if any, and waiting for it to sync */ + struct saib_pool *pool; + lws_dll2_t pool_wait_list; + sai_plat_t *sp; /* the sai_plat */ struct sai_plat_server *spm; /* the sai plat / server with the ss / wsi */ @@ -1052,7 +1064,7 @@ extern const lws_struct_map_t lsm_schema_map_ta[1], lsm_schema_map_plat_simple[1], lsm_event[15], - lsm_task[33], + lsm_task[34], lsm_log[8], lsm_artifact[9], lsm_plat_list[1], @@ -1202,6 +1214,115 @@ int sai_sqlite3_statement(struct sqlite3 *pdb, const char *cmd, const char *desc); /* + * Pools: a repo's named sets of files that the builders running its tasks keep + * synced through sai-server, eg, fuzzing corpora. See READMEs/README-pool.md. + * + * A builder opens a connection to sai-server's /builder endpoint for each + * sync, proves the link key as usual, then sends a JSON "hello" naming the + * task it syncs for, with the task's artifact upload nonce, which decides the + * repo and pool. After that, both ways, it's a stream of binary records: + * + * u8 type, u8 ns, u16 name length, u32 data length (big endian), + * the name, then the data + * + * Names are "<sub>/<name>", eg, "corpus-h2/<sha1>", except as noted. + */ + +#define SAI_POOL_SCHEMA "com.warmcat.sai.pool" +#define SAI_POOL_REC_HDR_LEN 8 +#define SAI_POOL_REC_NAME_MAX 128 +/* the most one corpus or known entry may be */ +#define SAI_POOL_ENTRY_MAX (1024u * 1024u) +/* the most one finding may be */ +#define SAI_POOL_FINDING_MAX (8u * 1024u * 1024u) +/* + * The most an OFFER (and so the WANT answering it) may list, the sender splits + * longer lists over several; and the most a REPLACE's list of names to keep, + * which can't be split, may be + */ +#define SAI_POOL_OFFER_MAX (512u * 1024u) +#define SAI_POOL_LIST_MAX (8u * 1024u * 1024u) + +enum { + /* + * Both ways, and the entries are content addressed: the name is the + * lowercase hex sha1 of the content + */ + SAI_POOL_NS_CORPUS, + /* server -> builder only, content addressed, eg, known reproducers */ + SAI_POOL_NS_KNOWN, + /* builder -> server only, any safe name, eg, fuzzer findings */ + SAI_POOL_NS_FINDINGS, + + SAI_POOL_NS_COUNT +}; + +enum { + /* builder -> server */ + + SAI_POOL_REC_PULL = 1, /* data: u64 BE cursor we have up to */ + SAI_POOL_REC_OFFER, /* data: '\n'-separated names we have */ + SAI_POOL_REC_PUT, /* name, data: the content */ + SAI_POOL_REC_REPLACE, /* name: sub only, data: u64 BE base + * cursor then '\n'-separated names in + * sub to keep, see README-pool.md */ + + /* server -> builder */ + + SAI_POOL_REC_ENTRY = 0x81, /* name, data: the content */ + SAI_POOL_REC_DEAD, /* name: the entry was removed */ + SAI_POOL_REC_PULL_END, /* data: u64 BE cursor now */ + SAI_POOL_REC_WANT, /* data: '\n'-separated offered names + * the server doesn't have */ + SAI_POOL_REC_ACK, /* name: of the PUT, or sub of the + * REPLACE, that was stored */ +}; + +typedef struct sai_pool_hello { + char task_uuid[65]; + char nonce[33]; +} sai_pool_hello_t; + +typedef struct sai_pool_rec_hdr { + uint32_t len; + uint16_t name_len; + uint8_t type; + uint8_t ns; +} sai_pool_rec_hdr_t; + +extern const lws_struct_map_t lsm_pool_hello[2], lsm_schema_pool_hello[1]; + +/* src/common/c-pool.c */ + +int +sai_pool_name_ok(const char *name); + +int +sai_pool_sub_ok(const char *sub, size_t len); + +int +sai_pool_entry_name_ok(int ns, const char *name, size_t len); + +int +sai_pool_content_matches(const uint8_t *data, size_t len, const char *sha1hex); + +void +sai_pool_rec_hdr_write(uint8_t *p, int type, int ns, size_t name_len, + size_t len); + +void +sai_pool_rec_hdr_read(const uint8_t *p, sai_pool_rec_hdr_t *h); + +size_t +sai_pool_rec_max(int ns, int type); + +void +sai_pool_u64_write(uint8_t *p, uint64_t v); + +uint64_t +sai_pool_u64_read(const uint8_t *p); + +/* * c-conf.c: create the context and vhosts from an lwsws-style config dir. * info must be zeroed and passed through lws_cmdline_option_handle_builtin() * by the caller first. diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 7096516..e96c5c0 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -230,6 +230,7 @@ const lws_struct_map_t lsm_task[] = { LSM_SIGNED (sai_task_t, rebuildable, "rebuildable"), LSM_SIGNED (sai_task_t, run, "run"), LSM_SIGNED (sai_task_t, idle, "idle"), + LSM_CARRAY (sai_task_t, pool, "pool"), }; const lws_struct_map_t lsm_schema_json_map_task[] = { @@ -641,3 +642,14 @@ const lws_struct_map_t lsm_schema_pending_tasks[] = { LSM_SCHEMA(sai_platform_pending_tasks_t, NULL, lsm_pending_tasks, "com.warmcat.sai.power.pending_tasks"), }; + +/* builder -> server, the first thing on a pool sync connection after auth */ + +const lws_struct_map_t lsm_pool_hello[] = { + LSM_CARRAY (sai_pool_hello_t, task_uuid, "task_uuid"), + LSM_CARRAY (sai_pool_hello_t, nonce, "nonce"), +}; + +const lws_struct_map_t lsm_schema_pool_hello[] = { + LSM_SCHEMA (sai_pool_hello_t, NULL, lsm_pool_hello, SAI_POOL_SCHEMA), +}; diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index 1f7d7e0..cf8096e 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -12,12 +12,14 @@ set(SRCS s-task.c s-task-helpers.c s-idle.c + s-pool.c s-central.c s-ws-web.c s-webops.c s-resource.c s-watcher.c ../common/c-utils.c + ../common/c-pool.c ../common/c-conf.c ../common/c-sqlite3.c ../common/struct-metadata.c diff --git a/src/server/s-comms.c b/src/server/s-comms.c index c1bb921..14ac426 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -717,6 +717,8 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, else lwsl_wsi_user(wsi, "#### sai-server: CLOSED builder conn (no close payload) ####"); } + sais_pool_session_destroy(pss); + /* remove pss from vhd->builders (active connection list) */ lws_dll2_remove(&pss->same); diff --git a/src/server/s-notification.c b/src/server/s-notification.c index 0ed49d6..3226527 100644 --- a/src/server/s-notification.c +++ b/src/server/s-notification.c @@ -93,6 +93,7 @@ static const char * const saifile_paths[] = { "configurations.*.branches", "configurations.*.task_log_limit", "configurations.*.idle", + "configurations.*.pool", "configurations.*", }; @@ -112,6 +113,7 @@ enum enum_saifile_paths { LEJPNSAIF_CONFIGURATIONS_BRANCHES, LEJPNSAIF_CONFIGURATIONS_TASK_LOG_LIMIT, LEJPNSAIF_CONFIGURATIONS_IDLE, + LEJPNSAIF_CONFIGURATIONS_POOL, LEJPNSAIF_CONFIGURATIONS_NAME, }; @@ -296,6 +298,7 @@ sai_saifile_lejp_cb(struct lejp_ctx *ctx, char reason) sn->t.branches[0] = '\0'; sn->explicit_platforms[0] = '\0'; sn->idle_lanes = 0; + sn->t.pool[0] = '\0'; return 0; } @@ -777,6 +780,22 @@ insert_fail: sn->idle_lanes = SAI_IDLE_LANES_MAX; break; + case LEJPNSAIF_CONFIGURATIONS_POOL: + /* + * The repo's named pool the configuration's tasks work with, + * which builders keep synced with us while they run them. It + * becomes a path component on the builder and part of our db + * filename, so it has to be a plain short name. + */ + lws_strncpy(sn->t.pool, ctx->buf, sizeof(sn->t.pool)); + if (ctx->npos >= sizeof(sn->t.pool) || + !sai_pool_name_ok(sn->t.pool)) { + lwsl_notice("%s: rejecting pool name '%s'\n", + __func__, sn->t.pool); + return -1; + } + break; + case LEJPNSAIF_PLAT_BUILD: case LEJPNSAIF_PLAT_BUILD_STAGE: /* diff --git a/src/server/s-pool.c b/src/server/s-pool.c new file mode 100644 index 0000000..7db981e --- /dev/null +++ b/src/server/s-pool.c @@ -0,0 +1,870 @@ +/* + * Sai server - ./src/server/s-pool.c + * + * Copyright (C) 2019 - 2026 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * Pools, server side: see READMEs/README-pool.md + * + * Each repo's pool lives in its own sqlite db, as entries in an append-only + * log: every entry added, or removed, gets the next sequence number, so a + * builder only has to ask for what happened after the last sequence number it + * saw. Removed entries leave a "dead" row behind to tell builders about it. + * Findings builders send go in their own table. + * + * A pool sync connection arrives on /builder like any other, proves the link + * key, then sends a JSON hello naming a task and its artifact upload nonce. + * The task decides the repo and pool, so a builder can only sync the pool of + * a task it was given. After that it's binary records both ways. + */ + +#include <libwebsockets.h> +#include <string.h> + +#include "s-private.h" + +/* past this many live entries, a pool takes no more */ +#define SAIS_POOL_ENTRIES_MAX (2 * 1000 * 1000) +/* we queue pulled entries for sending until there's this much waiting */ +#define SAIS_POOL_TX_HIGHWATER (256 * 1024) +/* the most we send in one ws message */ +#define SAIS_POOL_TX_CHUNK (64 * 1024) + +typedef struct sais_pool_db { + lws_dll2_t list; /* vhd->pool_dbs */ + sqlite3 *pdb; + int refcount; + unsigned int live; /* live entries */ + char key[128]; /* "<repo>/<pool>" */ +} sais_pool_db_t; + +struct sais_pool_session { + sais_pool_db_t *db; + + struct lws_buflist *tx; /* records waiting to go */ + uint8_t *txb; /* LWS_PRE + chunk */ + + uint8_t *rxb; /* partial incoming record(s) */ + size_t rxb_len; + size_t rxb_alloc; + + /* pulls going on, per content addressed ns: from pos up to end */ + uint64_t pull_pos[2]; + uint64_t pull_end[2]; + char pulling[2]; + + char task_uuid[65]; +}; + +static const char * const pool_schema[] = { + "CREATE TABLE IF NOT EXISTS entries (" + "seq INTEGER PRIMARY KEY AUTOINCREMENT, " + "ns INTEGER NOT NULL, " + "sub VARCHAR(32) NOT NULL, " + "name VARCHAR(64) NOT NULL, " + "len INTEGER, " + "dead INTEGER NOT NULL DEFAULT 0, " + "added INTEGER, " + "blob BLOB, " + "UNIQUE(ns, sub, name));", + "CREATE INDEX IF NOT EXISTS idx_entries_ns_seq ON entries(ns, seq);", + "CREATE TABLE IF NOT EXISTS findings (" + "uid INTEGER PRIMARY KEY AUTOINCREMENT, " + "sub VARCHAR(32) NOT NULL, " + "name VARCHAR(64) NOT NULL, " + "len INTEGER, " + "received INTEGER, " + "task_uuid VARCHAR(65), " + "peer VARCHAR(48), " + "blob BLOB, " + "UNIQUE(sub, name));", + "PRAGMA journal_mode=WAL;", +}; + +static sais_pool_db_t * +sais_pool_db_get(struct vhd *vhd, const char *repo, const char *pool) +{ + char key[128], fn[256], saf[128], *p; + sais_pool_db_t *db; + sqlite3_stmt *sm; + size_t n; + + lws_snprintf(key, sizeof(key), "%s/%s", repo, pool); + + lws_start_foreach_dll(struct lws_dll2 *, d, vhd->pool_dbs.head) { + db = lws_container_of(d, sais_pool_db_t, list); + + if (!strcmp(db->key, key)) { + db->refcount++; + return db; + } + + } lws_end_foreach_dll(d); + + /* the repo name was checked at hook intake, but it's going in a path */ + + lws_snprintf(saf, sizeof(saf), "%s-%s", repo, pool); + lws_filename_purify_inplace(saf); + p = saf; + while ((p = strchr(p, '/'))) + *p++ = '_'; + + lws_snprintf(fn, sizeof(fn), "%s-pool-%s.sqlite3", + vhd->sqlite3_path_lhs, saf); + + db = malloc(sizeof(*db)); + if (!db) + return NULL; + memset(db, 0, sizeof(*db)); + lws_strncpy(db->key, key, sizeof(db->key)); + + if (sqlite3_open_v2(fn, &db->pdb, SQLITE_OPEN_READWRITE | + SQLITE_OPEN_CREATE, NULL) != SQLITE_OK) { + lwsl_err("%s: unable to open %s\n", __func__, fn); + sqlite3_close(db->pdb); + free(db); + return NULL; + } + + sqlite3_busy_timeout(db->pdb, SAI_SQLITE3_BUSY_TIMEOUT_MS); + + for (n = 0; n < LWS_ARRAY_SIZE(pool_schema); n++) + if (sai_sqlite3_statement(db->pdb, pool_schema[n], + "pool schema")) { + sqlite3_close(db->pdb); + free(db); + return NULL; + } + + if (sqlite3_prepare_v2(db->pdb, "SELECT count(*) FROM entries WHERE " + "dead = 0", -1, &sm, NULL) == SQLITE_OK) { + if (sqlite3_step(sm) == SQLITE_ROW) + db->live = (unsigned int)sqlite3_column_int(sm, 0); + sqlite3_finalize(sm); + } + + db->refcount = 1; + lws_dll2_add_tail(&db->list, &vhd->pool_dbs); + + lwsl_notice("%s: opened pool %s, %u entries\n", __func__, key, db->live); + + return db; +} + +static void +sais_pool_db_put(sais_pool_db_t *db) +{ + if (--db->refcount) + return; + + lws_dll2_remove(&db->list); + sqlite3_close(db->pdb); + free(db); +} + +/* the newest sequence number, ie, what a pull started now goes up to */ + +static uint64_t +sais_pool_seq_now(sais_pool_db_t *db) +{ + sqlite3_stmt *sm; + uint64_t v = 0; + + if (sqlite3_prepare_v2(db->pdb, "SELECT max(seq) FROM entries", -1, + &sm, NULL) != SQLITE_OK) + return 0; + if (sqlite3_step(sm) == SQLITE_ROW) + v = (uint64_t)sqlite3_column_int64(sm, 0); + sqlite3_finalize(sm); + + return v; +} + +/* 1 = live, 0 = never seen or dead, -1 = error */ + +static int +sais_pool_live(sais_pool_db_t *db, int ns, const char *sub, size_t sub_len, + const char *name, size_t name_len) +{ + sqlite3_stmt *sm; + int r = 0; + + if (sqlite3_prepare_v2(db->pdb, "SELECT 1 FROM entries WHERE ns = ? " + "AND sub = ? AND name = ? AND dead = 0", -1, + &sm, NULL) != SQLITE_OK) + return -1; + + sqlite3_bind_int(sm, 1, ns); + sqlite3_bind_text(sm, 2, sub, (int)sub_len, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 3, name, (int)name_len, SQLITE_TRANSIENT); + if (sqlite3_step(sm) == SQLITE_ROW) + r = 1; + sqlite3_finalize(sm); + + return r; +} + +/* + * (Re)write the entry as the newest thing in the log: any older row for it, + * dead or alive, goes, so the entry appears again with a new sequence number + * and builders hear about it on their next pull. blob NULL means make it a + * dead entry. + */ + +static int +sais_pool_log_entry(sais_pool_db_t *db, int ns, const char *sub, + size_t sub_len, const char *name, size_t name_len, + const uint8_t *blob, size_t len) +{ + sqlite3_stmt *sm; + int r; + + if (sqlite3_prepare_v2(db->pdb, "DELETE FROM entries WHERE ns = ? AND " + "sub = ? AND name = ?", -1, &sm, NULL) != SQLITE_OK) + return -1; + sqlite3_bind_int(sm, 1, ns); + sqlite3_bind_text(sm, 2, sub, (int)sub_len, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 3, name, (int)name_len, SQLITE_TRANSIENT); + r = sqlite3_step(sm); + sqlite3_finalize(sm); + if (r != SQLITE_DONE) + return -1; + + if (sqlite3_prepare_v2(db->pdb, "INSERT INTO entries (ns, sub, name, " + "len, dead, added, blob) VALUES (?,?,?,?,?,?,?)", + -1, &sm, NULL) != SQLITE_OK) + return -1; + sqlite3_bind_int(sm, 1, ns); + sqlite3_bind_text(sm, 2, sub, (int)sub_len, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 3, name, (int)name_len, SQLITE_TRANSIENT); + sqlite3_bind_int64(sm, 4, (sqlite3_int64)len); + sqlite3_bind_int(sm, 5, !blob); + sqlite3_bind_int64(sm, 6, (sqlite3_int64)lws_now_secs()); + if (blob) + sqlite3_bind_blob(sm, 7, blob, (int)len, SQLITE_TRANSIENT); + else + sqlite3_bind_null(sm, 7); + r = sqlite3_step(sm); + sqlite3_finalize(sm); + + return r == SQLITE_DONE ? 0 : -1; +} + +/* queue a record to go to the builder */ + +static int +sais_pool_queue(struct pss *pss, int type, int ns, const char *name, + size_t name_len, const uint8_t *data, size_t len) +{ + sais_pool_session_t *ps = pss->pool; + uint8_t hdr[SAI_POOL_REC_HDR_LEN]; + + sai_pool_rec_hdr_write(hdr, type, ns, name_len, len); + + if (lws_buflist_append_segment(&ps->tx, hdr, sizeof(hdr)) < 0 || + (name_len && lws_buflist_append_segment(&ps->tx, + (const uint8_t *)name, name_len) < 0) || + (len && lws_buflist_append_segment(&ps->tx, data, len) < 0)) { + lwsl_err("%s: tx OOM\n", __func__); + return -1; + } + + lws_callback_on_writable(pss->wsi); + + return 0; +} + +/* + * The builder offered us names it has in a content addressed ns... tell it + * which ones we want it to send + */ + +static int +sais_pool_offer(struct pss *pss, int ns, const uint8_t *data, size_t len) +{ + sais_pool_db_t *db = pss->pool->db; + const char *p = (const char *)data, *end = p + len; + uint8_t *want; + size_t wl = 0; + int r; + + want = malloc(len + 1); + if (!want) + return -1; + + while (p < end) { + const char *nl = memchr(p, '\n', lws_ptr_diff_size_t(end, p)), + *e = nl ? nl : end, *sl; + size_t l = lws_ptr_diff_size_t(e, p); + + if (l && sai_pool_entry_name_ok(ns, p, l)) { + sl = memchr(p, '/', l); + if (sais_pool_live(db, ns, p, lws_ptr_diff_size_t(sl, p), + sl + 1, + l - lws_ptr_diff_size_t(sl, p) - 1) == 0) { + memcpy(want + wl, p, l); + wl += l; + want[wl++] = '\n'; + } + } + + p = e + 1; + } + + r = sais_pool_queue(pss, SAI_POOL_REC_WANT, ns, NULL, 0, want, wl); + free(want); + + return r; +} + +static int +sais_pool_put(struct pss *pss, int ns, const char *name, size_t name_len, + const uint8_t *data, size_t len) +{ + sais_pool_session_t *ps = pss->pool; + const char *sl = memchr(name, '/', name_len); + size_t sub_len = lws_ptr_diff_size_t(sl, name); + sqlite3_stmt *sm; + int r; + + if (ns == SAI_POOL_NS_FINDINGS) { + if (sqlite3_prepare_v2(ps->db->pdb, "INSERT OR IGNORE INTO " + "findings (sub, name, len, received, " + "task_uuid, peer, blob) VALUES (?,?,?,?,?,?,?)", + -1, &sm, NULL) != SQLITE_OK) + return -1; + sqlite3_bind_text(sm, 1, name, (int)sub_len, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 2, sl + 1, (int)(name_len - sub_len - 1), + SQLITE_TRANSIENT); + sqlite3_bind_int64(sm, 3, (sqlite3_int64)len); + sqlite3_bind_int64(sm, 4, (sqlite3_int64)lws_now_secs()); + sqlite3_bind_text(sm, 5, ps->task_uuid, -1, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 6, pss->peer_ip, -1, SQLITE_TRANSIENT); + sqlite3_bind_blob(sm, 7, data, (int)len, SQLITE_TRANSIENT); + r = sqlite3_step(sm); + sqlite3_finalize(sm); + if (r != SQLITE_DONE) + return -1; + + lwsl_notice("%s: %s: finding %.*s from %s\n", __func__, + ps->db->key, (int)name_len, name, pss->peer_ip); + + goto ack; + } + + /* the corpus is content addressed, so we can check it's really that */ + + if (!sai_pool_content_matches(data, len, sl + 1)) { + lwsl_notice("%s: %s: content isn't %.*s\n", __func__, + ps->db->key, (int)name_len, name); + return -1; + } + + r = sais_pool_live(ps->db, ns, name, sub_len, sl + 1, + name_len - sub_len - 1); + if (r < 0) + return -1; + if (!r) { + if (ps->db->live >= SAIS_POOL_ENTRIES_MAX) { + lwsl_warn("%s: %s: full, dropping %.*s\n", __func__, + ps->db->key, (int)name_len, name); + goto ack; + } + + if (sais_pool_log_entry(ps->db, ns, name, sub_len, sl + 1, + name_len - sub_len - 1, data, len)) + return -1; + ps->db->live++; + } + +ack: + return sais_pool_queue(pss, SAI_POOL_REC_ACK, ns, name, name_len, + NULL, 0); +} + +/* + * The builder replaced the set of entries in a sub, eg, after minimizing a + * corpus: it tells us the sequence number it had pulled up to when it began, + * and the names it's keeping. Entries up to then that aren't on the list are + * removed. Anything added since then, by anyone, stays: they didn't know + * about it. + */ + +static int +sais_pool_replace(struct pss *pss, int ns, const char *sub, size_t sub_len, + const uint8_t *data, size_t len) +{ + sais_pool_db_t *db = pss->pool->db; + const char *p = (const char *)data + 8, *end = (const char *)data + len; + uint64_t base = sai_pool_u64_read(data); + struct lwsac *ac = NULL; + lws_dll2_owner_t owner; + unsigned int dropped = 0; + sqlite3_stmt *sm; + int r = -1; + + typedef struct { + lws_dll2_t list; + char name[41]; + } gone_t; + + lws_dll2_owner_clear(&owner); + + if (sai_sqlite3_statement(db->pdb, "CREATE TEMP TABLE IF NOT EXISTS " + "keep (name VARCHAR(64) PRIMARY KEY)", + "pool keep") || + sai_sqlite3_statement(db->pdb, "DELETE FROM keep", "pool keep")) + return -1; + + if (sqlite3_prepare_v2(db->pdb, "INSERT OR IGNORE INTO keep VALUES " + "(?)", -1, &sm, NULL) != SQLITE_OK) + return -1; + + sai_sqlite3_statement(db->pdb, "BEGIN", "pool replace"); + + while (p < end) { + const char *nl = memchr(p, '\n', lws_ptr_diff_size_t(end, p)), + *e = nl ? nl : end; + size_t l = lws_ptr_diff_size_t(e, p); + + if (l == 40) { + sqlite3_bind_text(sm, 1, p, 40, SQLITE_TRANSIENT); + /* + * A name that didn't make it onto the keep list would + * get removed below + */ + if (sqlite3_step(sm) != SQLITE_DONE) { + lwsl_err("%s: %s: unable to keep %.40s: %s\n", + __func__, db->key, p, + sqlite3_errmsg(db->pdb)); + sqlite3_finalize(sm); + goto bail; + } + sqlite3_reset(sm); + } + p = e + 1; + } + sqlite3_finalize(sm); + + /* collect first, we're going to be rewriting the rows */ + + if (sqlite3_prepare_v2(db->pdb, "SELECT name FROM entries WHERE ns = ? " + "AND sub = ? AND dead = 0 AND seq <= ? AND " + "name NOT IN (SELECT name FROM keep)", -1, &sm, + NULL) != SQLITE_OK) + goto bail; + sqlite3_bind_int(sm, 1, ns); + sqlite3_bind_text(sm, 2, sub, (int)sub_len, SQLITE_TRANSIENT); + sqlite3_bind_int64(sm, 3, (sqlite3_int64)base); + while (sqlite3_step(sm) == SQLITE_ROW) { + const char *n = (const char *)sqlite3_column_text(sm, 0); + gone_t *g; + + if (!n) + continue; + g = lwsac_use_zero(&ac, sizeof(*g), 16384); + if (!g) + break; + lws_strncpy(g->name, n, sizeof(g->name)); + lws_dll2_add_tail(&g->list, &owner); + } + sqlite3_finalize(sm); + + lws_start_foreach_dll(struct lws_dll2 *, d, owner.head) { + gone_t *g = lws_container_of(d, gone_t, list); + + if (sais_pool_log_entry(db, ns, sub, sub_len, g->name, + strlen(g->name), NULL, 0)) + goto bail; + dropped++; + + } lws_end_foreach_dll(d); + + db->live = db->live > dropped ? db->live - dropped : 0; + r = 0; + +bail: + sai_sqlite3_statement(db->pdb, r ? "ROLLBACK" : "COMMIT", + "pool replace"); + lwsac_free(&ac); + + if (r) + return r; + + lwsl_notice("%s: %s: replaced %.*s up to %llu, %u removed\n", __func__, + db->key, (int)sub_len, sub, (unsigned long long)base, + dropped); + + return sais_pool_queue(pss, SAI_POOL_REC_ACK, ns, sub, sub_len, + NULL, 0); +} + +/* one whole record from the builder: 0 if OK, -1 to drop the connection */ + +static int +sais_pool_record(struct pss *pss, const sai_pool_rec_hdr_t *h, + const char *name, const uint8_t *data) +{ + sais_pool_session_t *ps = pss->pool; + + switch (h->type) { + case SAI_POOL_REC_PULL: + if (h->ns > SAI_POOL_NS_KNOWN || h->len != 8 || h->name_len) + return -1; + ps->pull_pos[h->ns] = sai_pool_u64_read(data); + ps->pull_end[h->ns] = sais_pool_seq_now(ps->db); + ps->pulling[h->ns] = 1; + lws_callback_on_writable(pss->wsi); + return 0; + + case SAI_POOL_REC_OFFER: + if (h->ns != SAI_POOL_NS_CORPUS || h->name_len) + return -1; + return sais_pool_offer(pss, h->ns, data, h->len); + + case SAI_POOL_REC_PUT: + if ((h->ns != SAI_POOL_NS_CORPUS && + h->ns != SAI_POOL_NS_FINDINGS) || + !sai_pool_entry_name_ok(h->ns, name, h->name_len)) + return -1; + return sais_pool_put(pss, h->ns, name, h->name_len, data, + h->len); + + case SAI_POOL_REC_REPLACE: + if (h->ns != SAI_POOL_NS_CORPUS || h->len < 8 || + !sai_pool_sub_ok(name, h->name_len)) + return -1; + return sais_pool_replace(pss, h->ns, name, h->name_len, data, + h->len); + } + + lwsl_notice("%s: unexpected record type 0x%x\n", __func__, h->type); + + return -1; +} + +/* + * Bytes arrived on an established pool sync connection + */ + +int +sais_pool_rx(struct vhd *vhd, struct pss *pss, const uint8_t *buf, size_t len) +{ + sais_pool_session_t *ps = pss->pool; + sai_pool_rec_hdr_t h; + size_t ofs = 0, need; + + if (!len) + return 0; + + if (ps->rxb_len + len > ps->rxb_alloc) { + size_t na = ps->rxb_len + len + 4096; + uint8_t *nb; + + if (na > SAI_POOL_LIST_MAX + SAI_POOL_REC_HDR_LEN + + SAI_POOL_REC_NAME_MAX + 65536) { + lwsl_notice("%s: rx too large\n", __func__); + return -1; + } + nb = realloc(ps->rxb, na); + if (!nb) + return -1; + ps->rxb = nb; + ps->rxb_alloc = na; + } + + memcpy(ps->rxb + ps->rxb_len, buf, len); + ps->rxb_len += len; + + while (ps->rxb_len - ofs >= SAI_POOL_REC_HDR_LEN) { + sai_pool_rec_hdr_read(ps->rxb + ofs, &h); + + if (h.ns >= SAI_POOL_NS_COUNT || + h.name_len > SAI_POOL_REC_NAME_MAX || + h.len > sai_pool_rec_max(h.ns, h.type)) { + lwsl_notice("%s: bad record hdr type 0x%x, ns %u, " + "name %u, len %u\n", __func__, h.type, + h.ns, h.name_len, h.len); + return -1; + } + + need = SAI_POOL_REC_HDR_LEN + h.name_len + h.len; + if (ps->rxb_len - ofs < need) + break; /* the rest of it isn't here yet */ + + if (sais_pool_record(pss, &h, (const char *)ps->rxb + ofs + + SAI_POOL_REC_HDR_LEN, + ps->rxb + ofs + SAI_POOL_REC_HDR_LEN + + h.name_len)) + return -1; + + ofs += need; + } + + if (ofs) { + memmove(ps->rxb, ps->rxb + ofs, ps->rxb_len - ofs); + ps->rxb_len -= ofs; + } + + return 0; +} + +/* + * Queue up the next part of whatever the builder is pulling, as long as there + * isn't much waiting to go already + */ + +static int +sais_pool_pull_more(struct pss *pss) +{ + sais_pool_session_t *ps = pss->pool; + sqlite3_stmt *sm; + uint8_t cur[8]; + int ns, n; + + for (ns = 0; ns < 2; ns++) { + if (!ps->pulling[ns]) + continue; + + if (sqlite3_prepare_v2(ps->db->pdb, "SELECT seq, sub, name, " + "dead, blob FROM entries WHERE ns = ? AND " + "seq > ? AND seq <= ? ORDER BY seq LIMIT 64", + -1, &sm, NULL) != SQLITE_OK) + return -1; + sqlite3_bind_int(sm, 1, ns); + sqlite3_bind_int64(sm, 2, (sqlite3_int64)ps->pull_pos[ns]); + sqlite3_bind_int64(sm, 3, (sqlite3_int64)ps->pull_end[ns]); + + n = 0; + while (sqlite3_step(sm) == SQLITE_ROW) { + const char *sub = (const char *)sqlite3_column_text(sm, 1), + *name = (const char *)sqlite3_column_text(sm, 2); + char nm[SAI_POOL_REC_NAME_MAX]; + int nl; + + n++; + ps->pull_pos[ns] = (uint64_t)sqlite3_column_int64(sm, 0); + if (!sub || !name) + continue; + + nl = lws_snprintf(nm, sizeof(nm), "%s/%s", sub, name); + + if (sqlite3_column_int(sm, 3)) { + if (sais_pool_queue(pss, SAI_POOL_REC_DEAD, ns, + nm, (size_t)nl, NULL, 0)) + goto fail; + } else + if (sais_pool_queue(pss, SAI_POOL_REC_ENTRY, ns, + nm, (size_t)nl, + sqlite3_column_blob(sm, 4), + (size_t)sqlite3_column_bytes(sm, 4))) + goto fail; + + if (lws_buflist_total_len(&ps->tx) > + SAIS_POOL_TX_HIGHWATER) + break; + } + sqlite3_finalize(sm); + + if (!n) { + /* nothing left up to where it started */ + sai_pool_u64_write(cur, ps->pull_end[ns]); + ps->pulling[ns] = 0; + if (sais_pool_queue(pss, SAI_POOL_REC_PULL_END, ns, + NULL, 0, cur, sizeof(cur))) + return -1; + } + + /* one ns at a time */ + return 0; + } + + return 0; + +fail: + sqlite3_finalize(sm); + + return -1; +} + +int +sais_pool_tx(struct vhd *vhd, struct pss *pss) +{ + sais_pool_session_t *ps = pss->pool; + uint8_t *p; + size_t n; + + if (lws_buflist_total_len(&ps->tx) < SAIS_POOL_TX_HIGHWATER && + sais_pool_pull_more(pss)) + return -1; + + n = lws_buflist_next_segment_len(&ps->tx, &p); + if (!n) + return 0; + + /* coalesce what's waiting into one message, up to a chunk */ + + n = 0; + while (n < SAIS_POOL_TX_CHUNK) { + size_t l = lws_buflist_next_segment_len(&ps->tx, &p); + + if (!l) + break; + if (l > SAIS_POOL_TX_CHUNK - n) + l = SAIS_POOL_TX_CHUNK - n; + memcpy(ps->txb + LWS_PRE + n, p, l); + lws_buflist_use_segment(&ps->tx, l); + n += l; + } + + if (lws_write(pss->wsi, ps->txb + LWS_PRE, n, LWS_WRITE_BINARY) < (int)n) + return -1; + + if (lws_buflist_total_len(&ps->tx) || ps->pulling[0] || ps->pulling[1]) + lws_callback_on_writable(pss->wsi); + + return 0; +} + +/* + * A builder connection sent the pool hello: check it's for a real task with + * a pool that it has the nonce for, and turn the connection into a sync + * session for the task's repo and pool. -1 to drop the connection. + */ + +int +sais_pool_hello(struct vhd *vhd, struct pss *pss, const sai_pool_hello_t *hello) +{ + char event_uuid[33], esc[96], filt[128], repo[65], pool[33]; + struct lwsac *ac = NULL; + sais_pool_session_t *ps; + sqlite3 *pdb = NULL; + lws_dll2_owner_t o; + sai_event_t *e; + sai_task_t *t; + int n; + + if (pss->pool || sais_validate_id(hello->task_uuid, SAI_TASKID_LEN)) { + lwsl_notice("%s: bad pool hello\n", __func__); + return -1; + } + + sai_task_uuid_to_event_uuid(event_uuid, hello->task_uuid); + + lws_sql_purify(esc, event_uuid, sizeof(esc)); + lws_snprintf(filt, sizeof(filt), " and uuid='%s'", esc); + n = lws_struct_sq3_deserialize(vhd->server.pdb, filt, NULL, + lsm_schema_sq3_map_event, &o, &ac, 0, 1); + if (n < 0 || !o.head) + goto bail; + e = lws_container_of(o.head, sai_event_t, list); + lws_strncpy(repo, e->repo_name, sizeof(repo)); + + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) + goto bail; + + lws_sql_purify(esc, hello->task_uuid, sizeof(esc)); + lws_snprintf(filt, sizeof(filt), " and uuid='%s'", esc); + n = lws_struct_sq3_deserialize(pdb, filt, "run desc", + lsm_schema_sq3_map_task, &o, &ac, 0, 1); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); + if (n < 0 || !o.head) + goto bail; + t = lws_container_of(o.head, sai_task_t, list); + + /* both fixed 32-char hex in 33-byte arrays */ + if (lws_timingsafe_bcmp(t->art_up_nonce, hello->nonce, 32)) { + lwsl_notice("%s: pool hello nonce mismatch\n", __func__); + goto bail; + } + + if (!sai_pool_name_ok(t->pool) || sai_str_has_shell_metachars(repo)) { + lwsl_notice("%s: task %s has no usable pool\n", __func__, + hello->task_uuid); + goto bail; + } + lws_strncpy(pool, t->pool, sizeof(pool)); + lwsac_free(&ac); + + ps = malloc(sizeof(*ps)); + if (!ps) + return -1; + memset(ps, 0, sizeof(*ps)); + + ps->txb = malloc(LWS_PRE + SAIS_POOL_TX_CHUNK); + ps->db = sais_pool_db_get(vhd, repo, pool); + if (!ps->txb || !ps->db) { + if (ps->db) + sais_pool_db_put(ps->db); + free(ps->txb); + free(ps); + return -1; + } + lws_strncpy(ps->task_uuid, hello->task_uuid, sizeof(ps->task_uuid)); + + pss->pool = ps; + + /* + * It's not a builder platform connection: take it off that list, so + * nothing meant for builders gets queued on it, and drop anything that + * already was + */ + + lws_dll2_remove(&pss->same); + lws_dll2_add_tail(&pss->same, &vhd->pool_conns); + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + pss->viewer_state_owner.head) { + lws_dll2_remove(d); + free(lws_container_of(d, sai_viewer_state_t, list)); + } lws_end_foreach_dll_safe(d, d1); + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + pss->task_cancel_owner.head) { + lws_dll2_remove(d); + free(lws_container_of(d, sai_cancel_t, list)); + } lws_end_foreach_dll_safe(d, d1); + + lwsl_notice("%s: %s: sync from %s for %s\n", __func__, ps->db->key, + pss->peer_ip, ps->task_uuid); + + return 0; + +bail: + lwsac_free(&ac); + lwsl_notice("%s: refusing pool sync for %s\n", __func__, + hello->task_uuid); + + return -1; +} + +void +sais_pool_session_destroy(struct pss *pss) +{ + sais_pool_session_t *ps = pss->pool; + + if (!ps) + return; + + lws_buflist_destroy_all_segments(&ps->tx); + sais_pool_db_put(ps->db); + free(ps->txb); + free(ps->rxb); + free(ps); + pss->pool = NULL; +} diff --git a/src/server/s-private.h b/src/server/s-private.h index 5b31cd6..e30a4af 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -144,6 +144,9 @@ typedef struct sai_builder { struct vhd; +/* a pool sync connection's state, see s-pool.c */ +typedef struct sais_pool_session sais_pool_session_t; + struct pss { struct vhd *vhd; struct lws *wsi; @@ -159,6 +162,8 @@ struct pss { sqlite3 *pdb_artifact; sqlite3_blob *blob_artifact; + sais_pool_session_t *pool; /* it's a pool sync connection */ + lws_dll2_owner_t platform_owner; /* sai_platform_t builder offers */ lws_dll2_owner_t task_cancel_owner; /* sai_platform_t builder offers */ lws_dll2_owner_t rebuild_owner; @@ -297,6 +302,10 @@ struct vhd { lws_dll2_owner_t watcher_services; /* sai_watcher_service_t from config */ + /* pools, see s-pool.c */ + lws_dll2_owner_t pool_dbs; /* sais_pool_db_t, open pool dbs */ + lws_dll2_owner_t pool_conns; /* pss of pool sync connections */ + /* idle tasks, see s-idle.c */ lws_dll2_owner_t idle_budgets; /* sais_idle_budget_t */ lws_dll2_owner_t idle_hosts; /* sais_idle_host_t in ac_idle_hosts */ @@ -474,6 +483,18 @@ sais_idle_declined(struct vhd *vhd, sai_plat_t *sp, const char *task_uuid); void sais_idle_builder_gone(struct vhd *vhd, sai_plat_t *sp); +int +sais_pool_hello(struct vhd *vhd, struct pss *pss, const sai_pool_hello_t *hello); + +int +sais_pool_rx(struct vhd *vhd, struct pss *pss, const uint8_t *buf, size_t len); + +int +sais_pool_tx(struct vhd *vhd, struct pss *pss); + +void +sais_pool_session_destroy(struct pss *pss); + sai_plat_t * sais_builder_from_uuid(struct vhd *vhd, const char *hostname); sai_plat_t * diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index f3b8ded..ce03903 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -85,6 +85,8 @@ static const lws_struct_map_t lsm_schema_map_ba[] = { "com.warmcat.sai.ptydata"), LSM_SCHEMA (sai_active_shells_t, NULL, lsm_schema_active_shells, "com.warmcat.sai.active_shells"), + LSM_SCHEMA (sai_pool_hello_t, NULL, lsm_pool_hello, + SAI_POOL_SCHEMA), }; enum { @@ -97,6 +99,7 @@ enum { SAIM_WSSCH_BUILDER_METRIC, SAIM_WSSCH_BUILDER_PTYDATA, SAIM_WSSCH_BUILDER_ACTIVE_SHELLS, + SAIM_WSSCH_BUILDER_POOL_HELLO, }; static void @@ -865,6 +868,10 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b sais_metrics_db_init(vhd); + if (pss->pool) + /* after its hello, a pool sync connection is binary records */ + return sais_pool_rx(vhd, pss, buf, bl); + if (pss->bulk_binary_data) { lwsl_info("%s: bulk %d\n", __func__, (int)bl); m = (int)bl; @@ -1258,6 +1265,21 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b lwsac_free(&pss->a.ac); break; + case SAIM_WSSCH_BUILDER_POOL_HELLO: + /* + * This connection is for syncing a pool: from here on + * it carries binary records, including any that came + * after the hello in this buffer + */ + n = sais_pool_hello(vhd, pss, + (sai_pool_hello_t *)pss->a.dest); + lwsac_free(&pss->a.ac); + if (n) + return -1; + + return sais_pool_rx(vhd, pss, buf + bl - (unsigned int)m, + (size_t)m); + case SAIM_WSSCH_BUILDER_ACTIVE_SHELLS: { sai_active_shells_t *ash = (sai_active_shells_t *)pss->a.dest; @@ -1775,6 +1797,9 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, sai_task_t *task; size_t w; + if (pss->pool) + return sais_pool_tx(vhd, pss); + if (pss->viewer_state_owner.head) { /* * Pending viewer state message to send to a builder
Page fetched 0s ago, creation time: 12ms (vhost etag hits: 0%, cache hits: 0%)